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
gamertan bf56dbce0f docs: publish Tend Compose continuity evidence
Export the reviewed allowlisted snapshot from private source commit 07c1655921f21ee5e4fc4d85639d199e8867b17d. This records the Docker Compose activation, schema-compatible rollback, and stateful migration resource findings from Observatory Preview 19 dogfooding.

AI-Assisted: OpenAI Codex
Signed-off-by: Cole Speelman <crspeelman@gmail.com>
2026-08-18 21:42:33 -04:00

1024 lines
35 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"`
LeaseCleanupPending bool `json:"lease_cleanup_pending,omitempty"`
}
type Status struct {
State *state.Record `json:"state,omitempty"`
Units map[string]bool `json:"units"`
StateInitialized bool `json:"state_initialized"`
}
type retainedCandidateError struct{ cause error }
func (e *retainedCandidateError) Error() string { return e.cause.Error() }
func (e *retainedCandidateError) Unwrap() error { return e.cause }
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 = ""
}
}
if cfg.Deployment.Strategy == "singleton_candidate" && operationID == "" {
return Report{}, errors.New("singleton activation requires a fresh operation identity")
}
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
}
leasePath := state.CandidateLeasePath(cfg.Deployment.StateFile)
var candidateLease state.CandidateLease
if cfg.Deployment.Strategy == "singleton_candidate" {
_, exists, leaseErr := loadCandidateLease(cfg)
if leaseErr != nil {
return Report{}, leaseErr
}
if record.CandidateRelease != "" || exists {
return Report{}, errors.New("singleton candidate lease is unresolved; run tend reconcile --json before another activation")
}
}
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 cfg.Deployment.Strategy == "singleton_candidate" {
candidateUnit, unitErr := singletonCandidateUnit(cfg.Service.Name, operationID)
if unitErr != nil {
return Report{}, unitErr
}
candidateLease = state.CandidateLease{SchemaVersion: state.CandidateLeaseSchemaVersion, Service: cfg.Service.Name, OperationID: operationID, Release: release, Unit: candidateUnit, Address: cfg.Deployment.Singleton.CandidateAddress, StartedAt: attemptAt}
}
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
return Report{}, err
}
if cfg.Deployment.Strategy == "singleton_candidate" {
if err := state.StoreCandidateLease(leasePath, cfg.Deployment.Root, candidateLease); err != nil {
failed := record
clearCandidateLease(&failed)
failed.LastAttemptOutcome = "failed"
failed.UpdatedAt = m.Now().UTC().Format(time.RFC3339)
if storeErr := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, failed); storeErr != nil {
return Report{}, errors.Join(err, storeErr)
}
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, candidateLease.Unit)
default:
err = errors.New("unsupported strategy")
}
if err != nil {
failed := record
var retained *retainedCandidateError
if !errors.As(err, &retained) {
if cleanupErr := state.RemoveCandidateLease(leasePath); cleanupErr != nil {
err = errors.Join(err, fmt.Errorf("candidate lease cleanup failed: %w", cleanupErr))
} else {
clearCandidateLease(&failed)
}
}
failed.LastAttemptOutcome = "failed"
failed.UpdatedAt = m.Now().UTC().Format(time.RFC3339)
if storeErr := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, failed); storeErr != nil {
err = errors.Join(err, storeErr)
}
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
}
report := Report{Validated: true, Mutation: "activated", Release: release, ActiveRelease: updated.ActiveRelease, PreviousRelease: updated.PreviousRelease, EventWarnings: eventWarnings}
if cfg.Deployment.Strategy == "singleton_candidate" {
if cleanupErr := state.RemoveCandidateLease(leasePath); cleanupErr != nil {
report.LeaseCleanupPending = true
}
}
return report, nil
}
func singletonCandidateUnit(service, operationID string) (string, error) {
if len(operationID) != 32 {
return "", errors.New("candidate operation identity is invalid")
}
for _, character := range operationID {
if !strings.ContainsRune("0123456789abcdef", character) {
return "", errors.New("candidate operation identity is invalid")
}
}
return service + "-tend-candidate-" + operationID[:12] + ".service", nil
}
func clearCandidateLease(record *state.Record) {
record.CandidateRelease = ""
}
func loadCandidateLease(cfg config.Config) (state.CandidateLease, bool, error) {
lease, err := state.LoadCandidateLease(state.CandidateLeasePath(cfg.Deployment.StateFile), cfg.Deployment.Root, cfg.Service.Name)
if err == nil {
return lease, true, nil
}
if os.IsNotExist(err) {
return state.CandidateLease{}, false, nil
}
return state.CandidateLease{}, false, fmt.Errorf("load singleton candidate lease: %w", err)
}
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: state.SchemaVersion, 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: state.SchemaVersion, 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: state.SchemaVersion, 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, candidateUnit string) (err error) {
if err = m.activateSingletonRelease(ctx, cfg, release, candidateUnit, true); err != nil {
return err
}
next := state.Record{SchemaVersion: state.SchemaVersion, 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")
}
leasePath := state.CandidateLeasePath(cfg.Deployment.StateFile)
if cfg.Deployment.Strategy == "singleton_candidate" {
_, exists, leaseErr := loadCandidateLease(cfg)
if leaseErr != nil {
return state.Record{}, leaseErr
}
if record.CandidateRelease != "" || exists {
return state.Record{}, errors.New("singleton candidate lease is unresolved; run tend reconcile --json before rollback")
}
}
started := m.Now()
identity, identityErr := m.ReadIdentity(record.PreviousRelease)
digest, digestErr := releaseDigest(record.PreviousRelease)
operationID := ""
if m.OperationID != nil {
operationID, err = m.OperationID()
}
if cfg.Deployment.Strategy == "singleton_candidate" && (err != nil || operationID == "") {
return state.Record{}, errors.New("singleton rollback requires a fresh operation identity")
}
candidateUnit := ""
rollbackAttempt := record
if cfg.Deployment.Strategy == "singleton_candidate" {
candidateUnit, err = singletonCandidateUnit(cfg.Service.Name, operationID)
if err != nil {
return state.Record{}, err
}
attemptAt := m.Now().UTC().Format(time.RFC3339)
rollbackAttempt.DesiredRelease = record.PreviousRelease
rollbackAttempt.CandidateRelease = record.PreviousRelease
rollbackAttempt.LastAttemptRelease = record.PreviousRelease
rollbackAttempt.LastAttemptOutcome = "running"
rollbackAttempt.LastAttemptAt = attemptAt
rollbackAttempt.UpdatedAt = attemptAt
if err = state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, rollbackAttempt); err != nil {
return state.Record{}, err
}
lease := state.CandidateLease{SchemaVersion: state.CandidateLeaseSchemaVersion, Service: cfg.Service.Name, OperationID: operationID, Release: record.PreviousRelease, Unit: candidateUnit, Address: cfg.Deployment.Singleton.CandidateAddress, StartedAt: attemptAt}
if err = state.StoreCandidateLease(leasePath, cfg.Deployment.Root, lease); err != nil {
failed := rollbackAttempt
clearCandidateLease(&failed)
failed.LastAttemptOutcome = "failed"
failed.UpdatedAt = m.Now().UTC().Format(time.RFC3339)
if storeErr := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, failed); storeErr != nil {
return state.Record{}, errors.Join(err, storeErr)
}
return state.Record{}, err
}
}
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, candidateUnit)
default:
err = errors.New("unsupported strategy")
}
if err != nil {
if cfg.Deployment.Strategy == "singleton_candidate" {
failed := rollbackAttempt
var retained *retainedCandidateError
if !errors.As(err, &retained) {
if cleanupErr := state.RemoveCandidateLease(leasePath); cleanupErr != nil {
err = errors.Join(err, fmt.Errorf("candidate lease cleanup failed: %w", cleanupErr))
} else {
clearCandidateLease(&failed)
}
}
failed.LastAttemptOutcome = "failed"
failed.UpdatedAt = m.Now().UTC().Format(time.RFC3339)
if storeErr := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, failed); storeErr != nil {
err = errors.Join(err, storeErr)
}
}
emit("failed")
return state.Record{}, err
}
emit("succeeded")
updated, loadErr := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if loadErr != nil {
return state.Record{}, loadErr
}
if cfg.Deployment.Strategy == "singleton_candidate" {
if cleanupErr := state.RemoveCandidateLease(leasePath); cleanupErr != nil {
return updated, fmt.Errorf("rollback succeeded but candidate lease cleanup is pending; run tend reconcile --json: %w", cleanupErr)
}
}
return updated, nil
}
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: state.SchemaVersion, 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, candidateUnit string) (err error) {
if err = m.activateSingletonRelease(ctx, cfg, record.PreviousRelease, candidateUnit, false); err != nil {
return err
}
next := state.Record{SchemaVersion: state.SchemaVersion, 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, candidateUnit string, checkMarkers bool) (err error) {
single := *cfg.Deployment.Singleton
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 = &retainedCandidateError{cause: 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}
lease, exists, leaseErr := loadCandidateLease(cfg)
if leaseErr != nil {
return Status{}, leaseErr
}
if exists {
units = append(units, lease.Unit)
} else if result.StateInitialized && record.CandidateRelease != "" {
return Status{}, errors.New("candidate release has no operation-scoped lease; run tend reconcile --json")
}
}
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}
if cfg.Deployment.Strategy == "singleton_candidate" {
lease, exists, leaseErr := loadCandidateLease(cfg)
if leaseErr != nil {
return nil, leaseErr
}
if exists {
protected[lease.Release] = true
} else if record.CandidateRelease != "" {
protected[record.CandidateRelease] = 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: state.SchemaVersion, 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
}