// 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 }