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>
857 lines
28 KiB
Go
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
|
|
}
|