This repository has been archived on 2026-08-19. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
tend/internal/deploy/deploy.go
T
gamertan 9d9fc83dd0 feat: publish Tend v0.2 Preview 2 source
Export the reviewed allowlisted snapshot from private source commit 8aab3db43f35e6a49aa497f45d73701b13fc9f32 and tree 992132ea4703437dc13ffdbb04a077816c02caf9. This includes routed singleton continuity, deployment evidence, strict schema-2 configuration, restricted transport, and the independently compilable public-tree guard.

AI-Assisted: OpenAI Codex
Signed-off-by: Cole Speelman <crspeelman@gmail.com>
2026-08-18 06:40:58 -04:00

857 lines
28 KiB
Go

// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"time"
"gamertan.com/tend/internal/config"
"gamertan.com/tend/internal/eventlog"
"gamertan.com/tend/internal/state"
)
type Request struct {
Artifact string
SHA256 string
ApprovedSHA256 string
Activate bool
}
type Report struct {
Validated bool `json:"validated"`
Mutation string `json:"mutation"`
Release string `json:"release,omitempty"`
ActiveRelease string `json:"active_release,omitempty"`
PreviousRelease string `json:"previous_release,omitempty"`
EventWarnings int `json:"event_warnings,omitempty"`
}
type Status struct {
State *state.Record `json:"state,omitempty"`
Units map[string]bool `json:"units"`
StateInitialized bool `json:"state_initialized"`
}
type Manager struct {
Operator Operator
Now func() time.Time
Prepare func(config.Config, string, string, string) (string, error)
Inspect func(config.Config, string, string, string) error
ReadIdentity func(string) (releaseIdentity, error)
OperationID func() (string, error)
AppendEvent func(string, eventlog.Event) error
Sleep func(context.Context, time.Duration) error
}
func NewManager(operator Operator) Manager {
return Manager{Operator: operator, Now: time.Now, Prepare: prepareRelease, Inspect: inspectArtifact, ReadIdentity: readReleaseIdentity, OperationID: eventlog.OperationID, AppendEvent: eventlog.Append, Sleep: sleepContext}
}
func (m Manager) Deploy(ctx context.Context, cfg config.Config, request Request) (Report, error) {
if err := cfg.Validate(); err != nil {
return Report{}, err
}
if !request.Activate {
if err := m.Inspect(cfg, request.Artifact, request.SHA256, request.ApprovedSHA256); err != nil {
return Report{}, err
}
return Report{Validated: true, Mutation: "none"}, nil
}
lock, err := acquireLock(cfg.Deployment.LockFile)
if err != nil {
return Report{}, err
}
defer lock.Close()
release, err := m.Prepare(cfg, request.Artifact, request.SHA256, request.ApprovedSHA256)
if err != nil {
return Report{}, err
}
started := m.Now()
eventWarnings := 0
identity, identityErr := m.ReadIdentity(release)
if identityErr != nil {
eventWarnings++
}
operationID := ""
if m.OperationID != nil {
operationID, err = m.OperationID()
if err != nil {
eventWarnings++
operationID = ""
}
}
emit := func(phase, slot, outcome string) {
if m.AppendEvent == nil || identityErr != nil || operationID == "" {
return
}
event := eventlog.Event{Version: eventlog.Version, OperationID: operationID, Service: cfg.Service.Name, ArtifactDigest: request.ApprovedSHA256, Commit: identity.Commit, ReleaseVersion: identity.Version, Phase: phase, Slot: slot, DurationMillis: max(0, m.Now().Sub(started).Milliseconds()), Outcome: outcome, ObservedAt: m.Now().UTC().Format(time.RFC3339Nano)}
if eventErr := m.AppendEvent(cfg.Deployment.EventLog, event); eventErr != nil {
eventWarnings++
}
}
record, err := loadOrBootstrap(cfg, m.Now())
if err != nil {
return Report{}, err
}
attemptAt := m.Now().UTC().Format(time.RFC3339)
record.DesiredRelease = release
record.CandidateRelease = release
record.LastAttemptRelease = release
record.LastAttemptOutcome = "running"
record.LastAttemptAt = attemptAt
record.UpdatedAt = attemptAt
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
return Report{}, err
}
emit("candidate", inactiveSlot(cfg, record), "running")
switch cfg.Deployment.Strategy {
case "blue_green":
err = m.deployBlueGreen(ctx, cfg, record, release)
case "singleton_candidate":
err = m.deploySingleton(ctx, cfg, record, release)
default:
err = errors.New("unsupported strategy")
}
if err != nil {
failed := record
failed.CandidateRelease = ""
failed.LastAttemptOutcome = "failed"
failed.UpdatedAt = m.Now().UTC().Format(time.RFC3339)
_ = state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, failed)
emit("activation", inactiveSlot(cfg, record), "failed")
return Report{EventWarnings: eventWarnings}, err
}
emit("activation", inactiveSlot(cfg, record), "succeeded")
updated, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
return Report{}, err
}
return Report{Validated: true, Mutation: "activated", Release: release, ActiveRelease: updated.ActiveRelease, PreviousRelease: updated.PreviousRelease, EventWarnings: eventWarnings}, nil
}
type releaseIdentity struct {
Version string `json:"version"`
Commit string `json:"commit"`
}
func readReleaseIdentity(release string) (releaseIdentity, error) {
b, err := os.ReadFile(filepath.Join(release, "RELEASE.json"))
if err != nil {
return releaseIdentity{}, fmt.Errorf("read installed release identity: %w", err)
}
if len(b) > 1<<20 {
return releaseIdentity{}, errors.New("installed release identity is too large")
}
var identity releaseIdentity
if err := json.Unmarshal(b, &identity); err != nil {
return releaseIdentity{}, errors.New("decode installed release identity")
}
if identity.Version == "" || identity.Commit == "" {
return releaseIdentity{}, errors.New("installed release identity is incomplete")
}
return identity, nil
}
func inactiveSlot(cfg config.Config, record state.Record) string {
if cfg.Deployment.Strategy == "singleton_candidate" {
return "singleton"
}
if record.ActiveSlot == "blue" {
return "green"
}
return "blue"
}
func sleepContext(ctx context.Context, duration time.Duration) error {
timer := time.NewTimer(duration)
defer timer.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-timer.C:
return nil
}
}
func loadOrBootstrap(cfg config.Config, now time.Time) (state.Record, error) {
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err == nil {
return record, nil
}
if !os.IsNotExist(err) {
return state.Record{}, err
}
switch cfg.Deployment.Strategy {
case "blue_green":
slot := cfg.Deployment.BlueGreen.BootstrapActive
release, err := resolveReleaseLink(cfg.Deployment.Root, slotConfig(*cfg.Deployment.BlueGreen, slot).Link)
if err != nil {
return state.Record{}, fmt.Errorf("bootstrap active slot: %w", err)
}
return state.Record{SchemaVersion: 1, Strategy: cfg.Deployment.Strategy, DesiredRelease: release, ActiveSlot: slot, ActiveRelease: release, UpdatedAt: now.UTC().Format(time.RFC3339)}, nil
case "singleton_candidate":
release, err := resolveReleaseLink(cfg.Deployment.Root, cfg.Deployment.Singleton.CurrentLink)
if err != nil {
return state.Record{}, fmt.Errorf("bootstrap singleton: %w", err)
}
return state.Record{SchemaVersion: 1, Strategy: cfg.Deployment.Strategy, DesiredRelease: release, ActiveSlot: "singleton", ActiveRelease: release, UpdatedAt: now.UTC().Format(time.RFC3339)}, nil
}
return state.Record{}, errors.New("unsupported strategy")
}
func (m Manager) deployBlueGreen(ctx context.Context, cfg config.Config, record state.Record, release string) (err error) {
bg := *cfg.Deployment.BlueGreen
inactive := "blue"
if record.ActiveSlot == "blue" {
inactive = "green"
}
slot := slotConfig(bg, inactive)
oldInactive, oldErr := resolveReleaseLink(cfg.Deployment.Root, slot.Link)
if oldErr != nil && !os.IsNotExist(oldErr) {
return oldErr
}
oldHandler, err := os.ReadFile(bg.CaddyHandler)
if err != nil {
return fmt.Errorf("read current Caddy handler: %w", err)
}
handlerChanged := false
linkChanged := false
defer func() {
if err == nil {
return
}
if handlerChanged {
_ = atomicWrite(bg.CaddyHandler, oldHandler, 0o644)
_ = m.Operator.ValidateCaddy(ctx, bg.CaddyConfig)
_ = m.Operator.ReloadCaddy(ctx)
}
if linkChanged {
if oldErr == nil {
_ = replaceSymlink(slot.Link, oldInactive)
_ = m.Operator.Restart(ctx, slot.Unit)
} else {
_ = removeSymlink(slot.Link)
_ = m.Operator.Stop(ctx, slot.Unit)
}
}
}()
if err = replaceSymlink(slot.Link, release); err != nil {
return err
}
linkChanged = true
if err = m.Operator.Restart(ctx, slot.Unit); err != nil {
return err
}
if err = m.probeAll(ctx, cfg, slot.Address); err != nil {
return fmt.Errorf("candidate failed: %w", err)
}
handler, err := renderHandler(bg.CaddyHandlerTemplate, slot.Address)
if err != nil {
return err
}
if err = atomicWrite(bg.CaddyHandler, handler, 0o644); err != nil {
return err
}
handlerChanged = true
if err = m.Operator.ValidateCaddy(ctx, bg.CaddyConfig); err != nil {
return fmt.Errorf("Caddy validation failed: %w", err)
}
if err = m.Operator.ReloadCaddy(ctx); err != nil {
return fmt.Errorf("Caddy reload failed: %w", err)
}
if err = m.probeAll(ctx, cfg, slot.Address); err != nil {
return fmt.Errorf("post-activation smoke failed: %w", err)
}
if err = m.probePublic(ctx, cfg, true); err != nil {
return fmt.Errorf("public-origin smoke failed: %w", err)
}
previous := slotConfig(bg, record.ActiveSlot)
if err = m.continuityWindow(ctx, cfg, previous.Address, true); err != nil {
return fmt.Errorf("activation continuity failed: %w", err)
}
next := state.Record{SchemaVersion: 1, Strategy: cfg.Deployment.Strategy, DesiredRelease: release, ActiveSlot: inactive, ActiveRelease: release, PreviousSlot: record.ActiveSlot, PreviousRelease: record.ActiveRelease, LastAttemptRelease: release, LastAttemptOutcome: "succeeded", LastAttemptAt: record.LastAttemptAt, UpdatedAt: m.Now().UTC().Format(time.RFC3339)}
if err = state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, next); err != nil {
return err
}
return nil
}
func (m Manager) deploySingleton(ctx context.Context, cfg config.Config, record state.Record, release string) (err error) {
if err = m.activateSingletonRelease(ctx, cfg, release, true); err != nil {
return err
}
next := state.Record{SchemaVersion: 1, Strategy: cfg.Deployment.Strategy, DesiredRelease: release, ActiveSlot: "singleton", ActiveRelease: release, PreviousSlot: "singleton", PreviousRelease: record.ActiveRelease, LastAttemptRelease: release, LastAttemptOutcome: "succeeded", LastAttemptAt: record.LastAttemptAt, UpdatedAt: m.Now().UTC().Format(time.RFC3339)}
if err = state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, next); err != nil {
return err
}
return nil
}
func (m Manager) Rollback(ctx context.Context, cfg config.Config) (state.Record, error) {
lock, err := acquireLock(cfg.Deployment.LockFile)
if err != nil {
return state.Record{}, err
}
defer lock.Close()
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
return state.Record{}, err
}
if record.PreviousRelease == "" {
return state.Record{}, errors.New("no previous release is recorded")
}
started := m.Now()
identity, identityErr := m.ReadIdentity(record.PreviousRelease)
digest, digestErr := releaseDigest(record.PreviousRelease)
operationID := ""
if m.OperationID != nil {
operationID, _ = m.OperationID()
}
emit := func(outcome string) {
if m.AppendEvent == nil || identityErr != nil || digestErr != nil || operationID == "" {
return
}
event := eventlog.Event{Version: eventlog.Version, OperationID: operationID, Service: cfg.Service.Name, ArtifactDigest: digest, Commit: identity.Commit, ReleaseVersion: identity.Version, Phase: "rollback", Slot: record.PreviousSlot, DurationMillis: max(0, m.Now().Sub(started).Milliseconds()), Outcome: outcome, ObservedAt: m.Now().UTC().Format(time.RFC3339Nano)}
_ = m.AppendEvent(cfg.Deployment.EventLog, event)
}
emit("running")
switch cfg.Deployment.Strategy {
case "blue_green":
err = m.rollbackBlueGreen(ctx, cfg, record)
case "singleton_candidate":
err = m.rollbackSingleton(ctx, cfg, record)
default:
err = errors.New("unsupported strategy")
}
if err != nil {
emit("failed")
return state.Record{}, err
}
emit("succeeded")
return state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
}
func releaseDigest(release string) (string, error) {
name := filepath.Base(release)
if !strings.HasPrefix(name, "sha256-") {
return "", errors.New("release is not content addressed")
}
digest := strings.TrimPrefix(name, "sha256-")
if len(digest) != 64 {
return "", errors.New("release digest is invalid")
}
for _, character := range digest {
if !strings.ContainsRune("0123456789abcdef", character) {
return "", errors.New("release digest is invalid")
}
}
return digest, nil
}
func (m Manager) rollbackBlueGreen(ctx context.Context, cfg config.Config, record state.Record) (err error) {
bg := *cfg.Deployment.BlueGreen
slot := slotConfig(bg, record.PreviousSlot)
if active, checkErr := m.Operator.IsActive(ctx, slot.Unit); checkErr != nil {
return checkErr
} else if !active {
if err = m.Operator.Restart(ctx, slot.Unit); err != nil {
return err
}
}
if err = m.probeHealthReadiness(ctx, cfg, slot.Address); err != nil {
return err
}
oldHandler, err := os.ReadFile(bg.CaddyHandler)
if err != nil {
return err
}
changed := false
defer func() {
if err != nil && changed {
_ = atomicWrite(bg.CaddyHandler, oldHandler, 0o644)
_ = m.Operator.ValidateCaddy(ctx, bg.CaddyConfig)
_ = m.Operator.ReloadCaddy(ctx)
}
}()
handler, err := renderHandler(bg.CaddyHandlerTemplate, slot.Address)
if err != nil {
return err
}
if err = atomicWrite(bg.CaddyHandler, handler, 0o644); err != nil {
return err
}
changed = true
if err = m.Operator.ValidateCaddy(ctx, bg.CaddyConfig); err != nil {
return err
}
if err = m.Operator.ReloadCaddy(ctx); err != nil {
return err
}
if err = m.probeHealthReadiness(ctx, cfg, slot.Address); err != nil {
return err
}
if err = m.probePublic(ctx, cfg, false); err != nil {
return err
}
next := state.Record{SchemaVersion: 1, Strategy: record.Strategy, DesiredRelease: record.PreviousRelease, ActiveSlot: record.PreviousSlot, ActiveRelease: record.PreviousRelease, PreviousSlot: record.ActiveSlot, PreviousRelease: record.ActiveRelease, LastAttemptRelease: record.PreviousRelease, LastAttemptOutcome: "rolled_back", LastAttemptAt: m.Now().UTC().Format(time.RFC3339), UpdatedAt: m.Now().UTC().Format(time.RFC3339)}
return state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, next)
}
func (m Manager) rollbackSingleton(ctx context.Context, cfg config.Config, record state.Record) (err error) {
if err = m.activateSingletonRelease(ctx, cfg, record.PreviousRelease, false); err != nil {
return err
}
next := state.Record{SchemaVersion: 1, Strategy: record.Strategy, DesiredRelease: record.PreviousRelease, ActiveSlot: "singleton", ActiveRelease: record.PreviousRelease, PreviousSlot: "singleton", PreviousRelease: record.ActiveRelease, LastAttemptRelease: record.PreviousRelease, LastAttemptOutcome: "rolled_back", LastAttemptAt: m.Now().UTC().Format(time.RFC3339), UpdatedAt: m.Now().UTC().Format(time.RFC3339)}
return state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, next)
}
// activateSingletonRelease keeps public traffic on a proven process while the
// installed fixed-address unit changes release. The transient candidate first
// receives traffic, remains healthy through the handoff, and is stopped only
// after Caddy points back to the verified installed unit.
func (m Manager) activateSingletonRelease(ctx context.Context, cfg config.Config, release string, checkMarkers bool) (err error) {
single := *cfg.Deployment.Singleton
candidateUnit := cfg.Service.Name + "-tend-candidate.service"
env := map[string]string{single.ListenEnv: single.CandidateAddress}
binary := filepath.Join(release, cfg.Build.Binary)
if err = m.Operator.StartCandidate(ctx, candidateUnit, binary, cfg.Service.EnvironmentFile, env); err != nil {
return err
}
stopCandidate := true
defer func() {
if !stopCandidate {
return
}
stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = m.Operator.Stop(stopCtx, candidateUnit)
}()
probeLocal := m.probeHealthReadiness
if checkMarkers {
probeLocal = m.probeAll
}
if err = probeLocal(ctx, cfg, single.CandidateAddress); err != nil {
return fmt.Errorf("candidate failed: %w", err)
}
oldHandler, err := os.ReadFile(single.CaddyHandler)
if err != nil {
return fmt.Errorf("read current Caddy handler: %w", err)
}
oldCurrent, err := resolveReleaseLink(cfg.Deployment.Root, single.CurrentLink)
if err != nil {
return err
}
oldPrevious, previousErr := resolveReleaseLink(cfg.Deployment.Root, single.PreviousLink)
handlerChanged := false
currentChanged := false
previousChanged := false
defer func() {
if err == nil {
return
}
recoveryErr := error(nil)
if currentChanged {
if restoreErr := replaceSymlink(single.CurrentLink, oldCurrent); restoreErr != nil {
recoveryErr = errors.Join(recoveryErr, restoreErr)
} else if restoreErr = m.Operator.Restart(ctx, single.Unit); restoreErr != nil {
recoveryErr = errors.Join(recoveryErr, restoreErr)
} else if restoreErr = m.probeHealthReadiness(ctx, cfg, single.Address); restoreErr != nil {
recoveryErr = errors.Join(recoveryErr, restoreErr)
}
}
if previousChanged {
var restoreErr error
if previousErr == nil {
restoreErr = replaceSymlink(single.PreviousLink, oldPrevious)
} else {
restoreErr = removeSymlink(single.PreviousLink)
}
recoveryErr = errors.Join(recoveryErr, restoreErr)
}
if handlerChanged && recoveryErr == nil {
if restoreErr := atomicWrite(single.CaddyHandler, oldHandler, 0o644); restoreErr != nil {
recoveryErr = errors.Join(recoveryErr, restoreErr)
} else if restoreErr = m.Operator.ValidateCaddy(ctx, single.CaddyConfig); restoreErr != nil {
recoveryErr = errors.Join(recoveryErr, restoreErr)
} else if restoreErr = m.Operator.ReloadCaddy(ctx); restoreErr != nil {
recoveryErr = errors.Join(recoveryErr, restoreErr)
}
}
if recoveryErr != nil && handlerChanged {
stopCandidate = false
err = errors.Join(err, fmt.Errorf("singleton recovery incomplete; candidate remains routed for operator recovery: %w", recoveryErr))
}
}()
candidateHandler, err := renderHandler(single.CaddyHandlerTemplate, single.CandidateAddress)
if err != nil {
return err
}
if err = atomicWrite(single.CaddyHandler, candidateHandler, 0o644); err != nil {
return err
}
handlerChanged = true
if err = m.Operator.ValidateCaddy(ctx, single.CaddyConfig); err != nil {
return fmt.Errorf("candidate Caddy validation failed: %w", err)
}
if err = m.Operator.ReloadCaddy(ctx); err != nil {
return fmt.Errorf("candidate Caddy reload failed: %w", err)
}
if err = m.probePublic(ctx, cfg, checkMarkers); err != nil {
return fmt.Errorf("candidate public-origin smoke failed: %w", err)
}
if err = replaceSymlink(single.PreviousLink, oldCurrent); err != nil {
return err
}
previousChanged = true
if err = replaceSymlink(single.CurrentLink, release); err != nil {
return err
}
currentChanged = true
if err = m.Operator.Restart(ctx, single.Unit); err != nil {
return err
}
if err = probeLocal(ctx, cfg, single.Address); err != nil {
return fmt.Errorf("post-activation smoke failed: %w", err)
}
installedHandler, err := renderHandler(single.CaddyHandlerTemplate, single.Address)
if err != nil {
return err
}
if err = atomicWrite(single.CaddyHandler, installedHandler, 0o644); err != nil {
return err
}
if err = m.Operator.ValidateCaddy(ctx, single.CaddyConfig); err != nil {
return fmt.Errorf("installed Caddy validation failed: %w", err)
}
if err = m.Operator.ReloadCaddy(ctx); err != nil {
return fmt.Errorf("installed Caddy reload failed: %w", err)
}
if err = m.probePublic(ctx, cfg, checkMarkers); err != nil {
return fmt.Errorf("public-origin smoke failed: %w", err)
}
if err = m.continuityWindow(ctx, cfg, single.CandidateAddress, checkMarkers); err != nil {
return fmt.Errorf("activation continuity failed: %w", err)
}
return nil
}
func (m Manager) Status(ctx context.Context, cfg config.Config) (Status, error) {
result := Status{Units: map[string]bool{}}
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err == nil {
result.State = &record
result.StateInitialized = true
} else if !os.IsNotExist(err) {
return Status{}, err
}
units := []string{}
if cfg.Deployment.Strategy == "blue_green" {
units = []string{cfg.Deployment.BlueGreen.Blue.Unit, cfg.Deployment.BlueGreen.Green.Unit}
} else {
units = []string{cfg.Deployment.Singleton.Unit}
}
for _, unit := range units {
active, err := m.Operator.IsActive(ctx, unit)
if err != nil {
return Status{}, err
}
result.Units[unit] = active
}
return result, nil
}
func (m Manager) Prune(cfg config.Config, keep int, apply bool) ([]string, error) {
if keep < 2 || keep > 100 {
return nil, errors.New("keep must be between 2 and 100")
}
lock, err := acquireLock(cfg.Deployment.LockFile)
if err != nil {
return nil, err
}
defer lock.Close()
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
return nil, err
}
entries, err := os.ReadDir(filepath.Join(cfg.Deployment.Root, "releases"))
if err != nil {
return nil, err
}
type candidate struct {
name, path string
mod time.Time
}
items := []candidate{}
protected := map[string]bool{record.ActiveRelease: true, record.PreviousRelease: true}
for _, entry := range entries {
if !entry.IsDir() || entry.Type()&os.ModeSymlink != 0 || !strings.HasPrefix(entry.Name(), "sha256-") {
continue
}
path := filepath.Join(cfg.Deployment.Root, "releases", entry.Name())
if protected[path] {
continue
}
info, err := entry.Info()
if err != nil {
return nil, err
}
items = append(items, candidate{entry.Name(), path, info.ModTime()})
}
sort.Slice(items, func(i, j int) bool { return items[i].mod.After(items[j].mod) })
retained := keep - 2
if retained < 0 {
retained = 0
}
if retained > len(items) {
retained = len(items)
}
items = items[retained:]
paths := make([]string, 0, len(items))
for _, item := range items {
paths = append(paths, item.path)
if apply {
if err := removeRelease(item.path, cfg.Deployment.Root); err != nil {
return paths, err
}
}
}
return paths, nil
}
func (m Manager) probeAll(ctx context.Context, cfg config.Config, address string) error {
checks := append([]config.Smoke{{Path: cfg.Deployment.HealthPath}, {Path: cfg.Deployment.ReadinessPath}}, cfg.Deployment.Smoke...)
return m.probe(ctx, cfg, address, checks)
}
func (m Manager) probeHealthReadiness(ctx context.Context, cfg config.Config, address string) error {
checks := []config.Smoke{{Path: cfg.Deployment.HealthPath}, {Path: cfg.Deployment.ReadinessPath}}
return m.probe(ctx, cfg, address, checks)
}
func (m Manager) probePublic(ctx context.Context, cfg config.Config, checkMarkers bool) error {
timeout := time.Duration(cfg.Deployment.CandidateTimeoutSecs) * time.Second
for _, check := range cfg.Deployment.PublicSmoke {
attempt, cancel := context.WithTimeout(ctx, timeout)
contains := check.Contains
if !checkMarkers {
contains = ""
}
err := m.Operator.ProbeURL(attempt, check.URL, contains)
cancel()
if err != nil {
return err
}
}
return nil
}
func (m Manager) continuityWindow(ctx context.Context, cfg config.Config, previousAddress string, checkMarkers bool) error {
steps := cfg.Deployment.ActivationWindowSecs * 4
if steps < 1 {
steps = 1
}
for step := 0; step < steps; step++ {
if err := m.probePublic(ctx, cfg, checkMarkers); err != nil {
return err
}
if previousAddress != "" {
if err := m.probeHealthReadiness(ctx, cfg, previousAddress); err != nil {
return fmt.Errorf("previous slot lost continuity: %w", err)
}
}
if step+1 < steps {
sleep := m.Sleep
if sleep == nil {
sleep = sleepContext
}
if err := sleep(ctx, 250*time.Millisecond); err != nil {
return err
}
}
}
return nil
}
func (m Manager) probe(ctx context.Context, cfg config.Config, address string, checks []config.Smoke) error {
timeout := time.Duration(cfg.Deployment.CandidateTimeoutSecs) * time.Second
for _, check := range checks {
deadline := m.Now().Add(timeout)
var last error
for {
attempt, cancel := context.WithTimeout(ctx, 2*time.Second)
last = m.Operator.Probe(attempt, address, cfg.Service.AllowedHost, check.Path, check.Contains)
cancel()
if last == nil {
break
}
if !m.Now().Before(deadline) {
return last
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(200 * time.Millisecond):
}
}
}
return nil
}
func slotConfig(bg config.BlueGreen, name string) config.Slot {
if name == "blue" {
return bg.Blue
}
return bg.Green
}
func resolveReleaseLink(root, link string) (string, error) {
info, err := os.Lstat(link)
if err != nil {
return "", err
}
if info.Mode()&os.ModeSymlink == 0 {
return "", errors.New("release pointer is not a symlink")
}
target, err := os.Readlink(link)
if err != nil {
return "", err
}
if !filepath.IsAbs(target) {
target = filepath.Join(filepath.Dir(link), target)
}
target = filepath.Clean(target)
probe := state.Record{SchemaVersion: 1, Strategy: "singleton_candidate", ActiveSlot: "singleton", ActiveRelease: target, UpdatedAt: time.Unix(1, 0).UTC().Format(time.RFC3339)}
if err := probe.Validate(root, "singleton_candidate"); err != nil {
return "", err
}
targetInfo, err := os.Lstat(target)
if err != nil {
return "", err
}
if !targetInfo.IsDir() || targetInfo.Mode()&os.ModeSymlink != 0 {
return "", errors.New("release target must be a real directory")
}
return target, nil
}
func replaceSymlink(link, target string) error {
if info, err := os.Lstat(link); err == nil && info.Mode()&os.ModeSymlink == 0 {
return errors.New("refusing to replace non-symlink release pointer")
} else if err != nil && !os.IsNotExist(err) {
return err
}
if err := os.MkdirAll(filepath.Dir(link), 0o755); err != nil {
return err
}
stage, err := os.MkdirTemp(filepath.Dir(link), ".tend-link-")
if err != nil {
return err
}
defer os.RemoveAll(stage)
tmp := filepath.Join(stage, "next")
if err := os.Symlink(target, tmp); err != nil {
return err
}
return os.Rename(tmp, link)
}
func removeSymlink(path string) error {
info, err := os.Lstat(path)
if os.IsNotExist(err) {
return nil
}
if err != nil {
return err
}
if info.Mode()&os.ModeSymlink == 0 {
return errors.New("refusing to remove non-symlink")
}
return os.Remove(path)
}
func renderHandler(templatePath, address string) ([]byte, error) {
b, err := os.ReadFile(templatePath)
if err != nil {
return nil, err
}
const marker = "{{UPSTREAM}}"
if bytes := strings.Count(string(b), marker); bytes != 1 {
return nil, errors.New("Caddy handler template must contain exactly one upstream marker")
}
return []byte(strings.Replace(string(b), marker, address, 1)), nil
}
func atomicWrite(path string, data []byte, mode os.FileMode) error {
var existing os.FileInfo
if info, err := os.Lstat(path); err == nil {
if info.Mode()&os.ModeSymlink != 0 || !info.Mode().IsRegular() {
return errors.New("refusing to replace non-regular or symlink file")
}
existing = info
} else if !os.IsNotExist(err) {
return err
}
identity := identityFor(existing, mode)
dir := filepath.Dir(path)
tmp, err := os.CreateTemp(dir, ".tend-write-")
if err != nil {
return err
}
name := tmp.Name()
ok := false
defer func() {
_ = tmp.Close()
if !ok {
_ = os.Remove(name)
}
}()
if err := applyIdentity(tmp, identity); err != nil {
return err
}
if _, err := tmp.Write(data); err != nil {
return err
}
if err := tmp.Sync(); err != nil {
return err
}
if err := tmp.Close(); err != nil {
return err
}
if err := os.Rename(name, path); err != nil {
return err
}
ok = true
return syncDirectory(dir)
}
func syncDirectory(path string) error {
dir, err := os.Open(path)
if err != nil {
return err
}
defer dir.Close()
return dir.Sync()
}
func removeRelease(path, root string) error {
releases := filepath.Join(root, "releases")
rel, err := filepath.Rel(releases, path)
if err != nil || rel == "." || rel == ".." || strings.ContainsRune(rel, filepath.Separator) || !strings.HasPrefix(rel, "sha256-") {
return errors.New("unsafe prune target")
}
info, err := os.Lstat(path)
if err != nil {
return err
}
if !info.IsDir() || info.Mode()&os.ModeSymlink != 0 {
return errors.New("prune target is not a real release directory")
}
return os.RemoveAll(path)
}
func Marshal(value any) ([]byte, error) {
b, err := json.MarshalIndent(value, "", " ")
if err != nil {
return nil, err
}
return append(b, '\n'), nil
}