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>
This commit is contained in:
2026-08-18 21:42:33 -04:00
commit bf56dbce0f
83 changed files with 8555 additions and 0 deletions
+1023
View File
@@ -0,0 +1,1023 @@
// 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
}
+591
View File
@@ -0,0 +1,591 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"time"
"gamertan.com/tend/internal/config"
"gamertan.com/tend/internal/eventlog"
"gamertan.com/tend/internal/state"
)
const testCandidateUnit = "example-site-tend-candidate-dddddddddddd.service"
func storeTestCandidateLease(t *testing.T, cfg config.Config, release, startedAt string) {
t.Helper()
lease := state.CandidateLease{SchemaVersion: state.CandidateLeaseSchemaVersion, Service: cfg.Service.Name, OperationID: strings.Repeat("d", 32), Release: release, Unit: testCandidateUnit, Address: cfg.Deployment.Singleton.CandidateAddress, StartedAt: startedAt}
if err := state.StoreCandidateLease(state.CandidateLeasePath(cfg.Deployment.StateFile), cfg.Deployment.Root, lease); err != nil {
t.Fatal(err)
}
}
type fakeOperator struct {
failReload bool
failReloadAt int
reloads int
failRestartUnit string
failRestartOnce bool
rejectMarkers bool
active map[string]bool
starts, stops, restarts []string
probes []string
publicProbes []string
failPublic bool
failPublicAfter int
candidateEnvironment map[string]string
candidateFile string
}
func (f *fakeOperator) Restart(_ context.Context, unit string) error {
f.restarts = append(f.restarts, unit)
if unit == f.failRestartUnit {
if f.failRestartOnce {
f.failRestartUnit = ""
}
return errors.New("injected restart failure")
}
f.active[unit] = true
return nil
}
func (f *fakeOperator) Stop(_ context.Context, unit string) error {
f.stops = append(f.stops, unit)
f.active[unit] = false
return nil
}
func (f *fakeOperator) IsActive(_ context.Context, unit string) (bool, error) {
return f.active[unit], nil
}
func (f *fakeOperator) StartCandidate(_ context.Context, unit, binary, environmentFile string, env map[string]string) error {
if !filepath.IsAbs(binary) || !filepath.IsAbs(environmentFile) || len(env) == 0 {
return errors.New("bad candidate")
}
f.starts = append(f.starts, unit)
f.candidateFile = environmentFile
f.candidateEnvironment = make(map[string]string, len(env))
for key, value := range env {
f.candidateEnvironment[key] = value
}
f.active[unit] = true
return nil
}
func (f *fakeOperator) ProbeURL(_ context.Context, value, contains string) error {
f.publicProbes = append(f.publicProbes, value)
if f.failPublic || (f.failPublicAfter > 0 && len(f.publicProbes) >= f.failPublicAfter) {
return errors.New("injected public smoke failure")
}
if f.rejectMarkers && contains != "" {
return errors.New("unexpected future-release smoke marker")
}
return nil
}
func TestPublicSmokeFailureRestoresBlueGreenHandlerAndSlot(t *testing.T) {
cfg, old, fresh := baseConfig(t, "blue_green")
handler := filepath.Join(cfg.Deployment.Root, "handler.caddy")
template := filepath.Join(cfg.Deployment.Root, "handler.template")
original := []byte("reverse_proxy 127.0.0.1:8090\n")
_ = os.WriteFile(handler, original, 0o644)
_ = os.WriteFile(template, []byte("reverse_proxy {{UPSTREAM}}\n"), 0o644)
blue := filepath.Join(cfg.Deployment.Root, "slots", "blue")
green := filepath.Join(cfg.Deployment.Root, "slots", "green")
_ = replaceSymlink(blue, old)
_ = replaceSymlink(green, old)
cfg.Deployment.BlueGreen = &config.BlueGreen{CaddyConfig: filepath.Join(cfg.Deployment.Root, "Caddyfile"), CaddyHandler: handler, CaddyHandlerTemplate: template, BootstrapActive: "blue", Blue: config.Slot{Unit: "example-blue.service", Address: "127.0.0.1:8090", Link: blue}, Green: config.Slot{Unit: "example-green.service", Address: "127.0.0.1:8091", Link: green}}
operator := &fakeOperator{active: map[string]bool{}, failPublic: true}
if _, err := manager(operator, fresh).Deploy(context.Background(), cfg, Request{Activate: true}); err == nil {
t.Fatal("expected public smoke failure")
}
body, _ := os.ReadFile(handler)
if string(body) != string(original) {
t.Fatalf("handler not restored: %q", body)
}
target, err := resolveReleaseLink(cfg.Deployment.Root, green)
if err != nil || target != old {
t.Fatalf("green=%q err=%v", target, err)
}
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
t.Fatal(err)
}
if record.ActiveRelease != old || record.LastAttemptOutcome != "failed" || record.CandidateRelease != "" || record.LastAttemptRelease != fresh {
t.Fatalf("failed attempt state=%+v", record)
}
}
func TestContinuityFailureRestoresBlueGreenRoute(t *testing.T) {
cfg, old, fresh := baseConfig(t, "blue_green")
handler := filepath.Join(cfg.Deployment.Root, "handler.caddy")
template := filepath.Join(cfg.Deployment.Root, "handler.template")
original := []byte("reverse_proxy 127.0.0.1:8090\n")
_ = os.WriteFile(handler, original, 0o644)
_ = os.WriteFile(template, []byte("reverse_proxy {{UPSTREAM}}\n"), 0o644)
blue := filepath.Join(cfg.Deployment.Root, "slots", "blue")
green := filepath.Join(cfg.Deployment.Root, "slots", "green")
_ = replaceSymlink(blue, old)
_ = replaceSymlink(green, old)
cfg.Deployment.BlueGreen = &config.BlueGreen{CaddyConfig: filepath.Join(cfg.Deployment.Root, "Caddyfile"), CaddyHandler: handler, CaddyHandlerTemplate: template, BootstrapActive: "blue", Blue: config.Slot{Unit: "example-blue.service", Address: "127.0.0.1:8090", Link: blue}, Green: config.Slot{Unit: "example-green.service", Address: "127.0.0.1:8091", Link: green}}
operator := &fakeOperator{active: map[string]bool{}, failPublicAfter: 3}
if _, err := manager(operator, fresh).Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)}); err == nil || !strings.Contains(err.Error(), "continuity") {
t.Fatalf("expected continuity failure, got %v", err)
}
body, _ := os.ReadFile(handler)
if string(body) != string(original) {
t.Fatalf("handler not restored: %q", body)
}
target, err := resolveReleaseLink(cfg.Deployment.Root, green)
if err != nil || target != old {
t.Fatalf("green=%q err=%v", target, err)
}
}
func (f *fakeOperator) ValidateCaddy(context.Context, string) error { return nil }
func (f *fakeOperator) ReloadCaddy(context.Context) error {
f.reloads++
if f.failReload || (f.failReloadAt > 0 && f.reloads == f.failReloadAt) {
return errors.New("injected reload failure")
}
return nil
}
func (f *fakeOperator) Probe(_ context.Context, address, host, path, contains string) error {
f.probes = append(f.probes, address+path)
if f.rejectMarkers && contains != "" {
return errors.New("unexpected future-release smoke marker")
}
return nil
}
func baseConfig(t *testing.T, strategy string) (config.Config, string, string) {
t.Helper()
root := filepath.Join(t.TempDir(), "service")
if err := os.MkdirAll(filepath.Join(root, "releases"), 0o755); err != nil {
t.Fatal(err)
}
old := filepath.Join(root, "releases", "sha256-"+strings.Repeat("c", 64))
fresh := filepath.Join(root, "releases", "sha256-"+strings.Repeat("a", 64))
for _, dir := range []string{old, fresh} {
if err := os.Mkdir(dir, 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(dir, "app"), []byte("x"), 0o755); err != nil {
t.Fatal(err)
}
}
cfg := config.Config{SchemaVersion: 2, Service: config.Service{Name: "example-site", AllowedHost: "example.test", EnvironmentFile: "/etc/tend/environment/example-site.env"}, Build: config.Build{Package: "./cmd/site", Binary: "app", Branch: "main"}, Deployment: config.Deployment{Strategy: strategy, Root: root, LockFile: filepath.Join(root, "deploy.lock"), StateFile: filepath.Join(root, "state.json"), EventLog: filepath.Join(root, "deployment-events.jsonl"), HealthPath: "/healthz", ReadinessPath: "/readyz", CandidateTimeoutSecs: 2, ActivationWindowSecs: 1, Smoke: []config.Smoke{{Path: "/", Contains: "Example"}}, PublicSmoke: []config.PublicSmoke{{URL: "https://example.test/", Contains: "Example"}}}}
return cfg, old, fresh
}
func manager(operator Operator, fresh string) Manager {
return Manager{Operator: operator, Now: func() time.Time { return time.Unix(100, 0).UTC() }, Prepare: func(config.Config, string, string, string) (string, error) { return fresh, nil }, Inspect: func(config.Config, string, string, string) error { return nil }, ReadIdentity: func(string) (releaseIdentity, error) {
return releaseIdentity{Version: "v0.2.0-preview.1", Commit: strings.Repeat("b", 40)}, nil
}, OperationID: func() (string, error) { return strings.Repeat("d", 32), nil }, Sleep: func(context.Context, time.Duration) error { return nil }}
}
func singletonSettings(t *testing.T, cfg config.Config, currentRelease string) (*config.Singleton, string) {
t.Helper()
handler := filepath.Join(cfg.Deployment.Root, "handler.caddy")
template := filepath.Join(cfg.Deployment.Root, "handler.template")
if err := os.WriteFile(handler, []byte("reverse_proxy 127.0.0.1:8092\n"), 0o640); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(template, []byte("reverse_proxy {{UPSTREAM}}\n"), 0o644); err != nil {
t.Fatal(err)
}
current := filepath.Join(cfg.Deployment.Root, "current")
if err := replaceSymlink(current, currentRelease); err != nil {
t.Fatal(err)
}
return &config.Singleton{
Unit: "example-site.service", Address: "127.0.0.1:8092", CandidateAddress: "127.0.0.1:18092", ListenEnv: "EXAMPLE_LISTEN",
CurrentLink: current, PreviousLink: filepath.Join(cfg.Deployment.Root, "previous"), CaddyConfig: filepath.Join(cfg.Deployment.Root, "Caddyfile"),
CaddyHandler: handler, CaddyHandlerTemplate: template,
}, handler
}
func TestBlueGreenActivationAndRollback(t *testing.T) {
cfg, old, fresh := baseConfig(t, "blue_green")
handler := filepath.Join(cfg.Deployment.Root, "handler.caddy")
template := filepath.Join(cfg.Deployment.Root, "handler.template")
if err := os.WriteFile(handler, []byte("reverse_proxy 127.0.0.1:8090\n"), 0o640); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(template, []byte("reverse_proxy {{UPSTREAM}}\n"), 0o644); err != nil {
t.Fatal(err)
}
blue := filepath.Join(cfg.Deployment.Root, "slots", "blue")
green := filepath.Join(cfg.Deployment.Root, "slots", "green")
if err := replaceSymlink(blue, old); err != nil {
t.Fatal(err)
}
if err := replaceSymlink(green, old); err != nil {
t.Fatal(err)
}
cfg.Deployment.BlueGreen = &config.BlueGreen{CaddyConfig: filepath.Join(cfg.Deployment.Root, "Caddyfile"), CaddyHandler: handler, CaddyHandlerTemplate: template, BootstrapActive: "blue", Blue: config.Slot{Unit: "example-blue.service", Address: "127.0.0.1:8090", Link: blue}, Green: config.Slot{Unit: "example-green.service", Address: "127.0.0.1:8091", Link: green}}
if err := cfg.Validate(); err != nil {
t.Fatal(err)
}
operator := &fakeOperator{active: map[string]bool{"example-blue.service": true, "example-green.service": true}}
m := manager(operator, fresh)
var events []eventlog.Event
m.AppendEvent = func(_ string, event eventlog.Event) error {
events = append(events, event)
return nil
}
report, err := m.Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)})
if err != nil {
t.Fatal(err)
}
if report.ActiveRelease != fresh || report.PreviousRelease != old {
t.Fatalf("report=%+v", report)
}
deployed, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
t.Fatal(err)
}
if deployed.DesiredRelease != fresh || deployed.CandidateRelease != "" || deployed.LastAttemptRelease != fresh || deployed.LastAttemptOutcome != "succeeded" {
t.Fatalf("deployment identity state=%+v", deployed)
}
if len(events) != 2 || events[0].Phase != "candidate" || events[0].Outcome != "running" || events[1].Phase != "activation" || events[1].Outcome != "succeeded" || events[0].OperationID != events[1].OperationID {
t.Fatalf("events=%+v", events)
}
if info, err := os.Stat(handler); err != nil || info.Mode().Perm() != 0o640 {
t.Fatalf("handler mode=%v err=%v", info.Mode().Perm(), err)
}
if target, err := resolveReleaseLink(cfg.Deployment.Root, green); err != nil || target != fresh {
t.Fatalf("green=%q err=%v", target, err)
}
operator.rejectMarkers = true
operator.probes = nil
record, err := m.Rollback(context.Background(), cfg)
if err != nil {
t.Fatal(err)
}
if record.ActiveRelease != old || record.PreviousRelease != fresh {
t.Fatalf("rollback=%+v", record)
}
if len(events) != 4 || events[2].Phase != "rollback" || events[2].Outcome != "running" || events[3].Phase != "rollback" || events[3].Outcome != "succeeded" || events[2].OperationID != events[3].OperationID || events[2].ArtifactDigest != strings.Repeat("c", 64) {
t.Fatalf("rollback events=%+v", events)
}
if len(operator.probes) != 4 {
t.Fatalf("rollback probes=%#v", operator.probes)
}
for _, probe := range operator.probes {
if !strings.HasSuffix(probe, "/healthz") && !strings.HasSuffix(probe, "/readyz") {
t.Fatalf("rollback applied future-release smoke checks: %#v", operator.probes)
}
}
}
func TestDeploymentEvidenceCanNeverBlockActivationOrRollback(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, old)
m := manager(&fakeOperator{active: map[string]bool{}}, fresh)
appendCalls := 0
m.AppendEvent = func(string, eventlog.Event) error {
appendCalls++
return errors.New("injected event failure")
}
report, err := m.Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)})
if err != nil {
t.Fatal(err)
}
if report.EventWarnings != 2 || report.ActiveRelease != fresh {
t.Fatalf("report=%+v", report)
}
if _, err := m.Rollback(context.Background(), cfg); err != nil {
t.Fatal(err)
}
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil || record.ActiveRelease != old || appendCalls != 4 {
t.Fatalf("record=%+v appends=%d err=%v", record, appendCalls, err)
}
}
func TestDeploymentEvidenceIdentityIsBestEffort(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, old)
m := manager(&fakeOperator{active: map[string]bool{}}, fresh)
m.ReadIdentity = func(string) (releaseIdentity, error) {
return releaseIdentity{}, errors.New("injected identity failure")
}
report, err := m.Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)})
if err != nil || report.EventWarnings != 1 || report.ActiveRelease != fresh {
t.Fatalf("report=%+v err=%v", report, err)
}
}
func TestSingletonOperationIdentityIsRequiredBeforeCandidateStart(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, old)
operator := &fakeOperator{active: map[string]bool{}}
m := manager(operator, fresh)
m.OperationID = func() (string, error) { return "", errors.New("injected entropy failure") }
if _, err := m.Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)}); err == nil || !strings.Contains(err.Error(), "operation identity") {
t.Fatalf("expected operation identity refusal, got %v", err)
}
if len(operator.starts) != 0 {
t.Fatalf("candidate started without an operation identity: %#v", operator.starts)
}
}
func TestSingletonUnresolvedLeaseBlocksReplacementBeforeCandidateStart(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, old)
at := time.Unix(100, 0).UTC().Format(time.RFC3339)
record := state.Record{SchemaVersion: state.SchemaVersion, Strategy: cfg.Deployment.Strategy, DesiredRelease: fresh, CandidateRelease: fresh, ActiveSlot: "singleton", ActiveRelease: old, LastAttemptRelease: fresh, LastAttemptOutcome: "failed", LastAttemptAt: at, UpdatedAt: at}
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
t.Fatal(err)
}
storeTestCandidateLease(t, cfg, fresh, at)
operator := &fakeOperator{active: map[string]bool{testCandidateUnit: true}}
if _, err := manager(operator, fresh).Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)}); err == nil || !strings.Contains(err.Error(), "run tend reconcile") {
t.Fatalf("expected unresolved lease refusal, got %v", err)
}
if len(operator.starts) != 0 || !operator.active[testCandidateUnit] {
t.Fatalf("existing candidate was disturbed: starts=%#v active=%#v", operator.starts, operator.active)
}
}
func TestBlueGreenCaddyFailureRestoresHandlerAndSlot(t *testing.T) {
cfg, old, fresh := baseConfig(t, "blue_green")
handler := filepath.Join(cfg.Deployment.Root, "handler.caddy")
template := filepath.Join(cfg.Deployment.Root, "handler.template")
original := []byte("reverse_proxy 127.0.0.1:8090\n")
_ = os.WriteFile(handler, original, 0o644)
_ = os.WriteFile(template, []byte("reverse_proxy {{UPSTREAM}}\n"), 0o644)
blue := filepath.Join(cfg.Deployment.Root, "slots", "blue")
green := filepath.Join(cfg.Deployment.Root, "slots", "green")
_ = replaceSymlink(blue, old)
_ = replaceSymlink(green, old)
cfg.Deployment.BlueGreen = &config.BlueGreen{CaddyConfig: filepath.Join(cfg.Deployment.Root, "Caddyfile"), CaddyHandler: handler, CaddyHandlerTemplate: template, BootstrapActive: "blue", Blue: config.Slot{Unit: "example-blue.service", Address: "127.0.0.1:8090", Link: blue}, Green: config.Slot{Unit: "example-green.service", Address: "127.0.0.1:8091", Link: green}}
operator := &fakeOperator{active: map[string]bool{}, failReload: true}
m := manager(operator, fresh)
if _, err := m.Deploy(context.Background(), cfg, Request{Activate: true}); err == nil {
t.Fatal("expected failure")
}
body, _ := os.ReadFile(handler)
if string(body) != string(original) {
t.Fatalf("handler not restored: %q", body)
}
target, err := resolveReleaseLink(cfg.Deployment.Root, green)
if err != nil || target != old {
t.Fatalf("green=%q err=%v", target, err)
}
failed, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
t.Fatal(err)
}
if failed.ActiveRelease != old || failed.LastAttemptOutcome != "failed" || failed.CandidateRelease != "" {
t.Fatalf("failed state=%+v", failed)
}
}
func TestSingletonRestartFailureRestoresPointers(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
settings, handler := singletonSettings(t, cfg, old)
cfg.Deployment.Singleton = settings
operator := &fakeOperator{active: map[string]bool{}, failRestartUnit: "example-site.service", failRestartOnce: true}
m := manager(operator, fresh)
if _, err := m.Deploy(context.Background(), cfg, Request{Activate: true}); err == nil {
t.Fatal("expected failure")
}
target, err := resolveReleaseLink(cfg.Deployment.Root, cfg.Deployment.Singleton.CurrentLink)
if err != nil || target != old {
t.Fatalf("current=%q err=%v", target, err)
}
if _, err := os.Lstat(cfg.Deployment.Singleton.PreviousLink); !os.IsNotExist(err) {
t.Fatal("previous pointer was not restored")
}
body, err := os.ReadFile(handler)
if err != nil || string(body) != "reverse_proxy 127.0.0.1:8092\n" {
t.Fatalf("handler=%q err=%v", body, err)
}
if info, err := os.Stat(handler); err != nil || info.Mode().Perm() != 0o640 {
t.Fatalf("handler mode=%v err=%v", info.Mode().Perm(), err)
}
if operator.active[testCandidateUnit] {
t.Fatal("candidate was not stopped after successful restoration")
}
}
func TestStateRecordsSuccessfulActivation(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
settings, handler := singletonSettings(t, cfg, old)
cfg.Deployment.Singleton = settings
operator := &fakeOperator{active: map[string]bool{}}
m := manager(operator, fresh)
if _, err := m.Deploy(context.Background(), cfg, Request{Activate: true}); err != nil {
t.Fatal(err)
}
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
t.Fatal(err)
}
if record.ActiveRelease != fresh || record.PreviousRelease != old || record.CandidateRelease != "" {
t.Fatalf("state=%+v", record)
}
if _, err := os.Lstat(state.CandidateLeasePath(cfg.Deployment.StateFile)); !os.IsNotExist(err) {
t.Fatalf("candidate lease was not removed: %v", err)
}
if len(operator.starts) != 1 || operator.starts[0] != testCandidateUnit {
t.Fatalf("candidate starts=%#v", operator.starts)
}
if operator.candidateFile != cfg.Service.EnvironmentFile || len(operator.candidateEnvironment) != 1 || operator.candidateEnvironment["EXAMPLE_LISTEN"] != "127.0.0.1:18092" {
t.Fatalf("candidate file=%q environment=%#v", operator.candidateFile, operator.candidateEnvironment)
}
body, err := os.ReadFile(handler)
if err != nil || string(body) != "reverse_proxy 127.0.0.1:8092\n" {
t.Fatalf("handler=%q err=%v", body, err)
}
if operator.reloads != 2 || operator.active[testCandidateUnit] {
t.Fatalf("reloads=%d active=%#v", operator.reloads, operator.active)
}
reconciliation, err := m.Reconcile(context.Background(), cfg)
if err != nil || reconciliation.Mutation != "none" || reconciliation.Disposition != "settled" || reconciliation.Observed.CandidateLease || !reconciliation.Consistent {
t.Fatalf("reconciliation=%+v err=%v", reconciliation, err)
}
}
func TestSingletonContinuityFailureRestoresHandlerPointersAndService(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
settings, handler := singletonSettings(t, cfg, old)
cfg.Deployment.Singleton = settings
operator := &fakeOperator{active: map[string]bool{"example-site.service": true}, failPublicAfter: 4}
m := manager(operator, fresh)
if _, err := m.Deploy(context.Background(), cfg, Request{Activate: true}); err == nil || !strings.Contains(err.Error(), "continuity") {
t.Fatalf("expected continuity failure, got %v", err)
}
current, err := resolveReleaseLink(cfg.Deployment.Root, cfg.Deployment.Singleton.CurrentLink)
if err != nil || current != old {
t.Fatalf("current=%q err=%v", current, err)
}
body, err := os.ReadFile(handler)
if err != nil || string(body) != "reverse_proxy 127.0.0.1:8092\n" {
t.Fatalf("handler=%q err=%v", body, err)
}
if operator.active[testCandidateUnit] {
t.Fatal("candidate was not stopped after continuity restoration")
}
}
func TestSingletonCaddyReloadFailuresRestorePriorRoute(t *testing.T) {
for _, reload := range []int{1, 2} {
t.Run(fmt.Sprintf("reload-%d", reload), func(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
settings, handler := singletonSettings(t, cfg, old)
cfg.Deployment.Singleton = settings
operator := &fakeOperator{active: map[string]bool{"example-site.service": true}, failReloadAt: reload}
if _, err := manager(operator, fresh).Deploy(context.Background(), cfg, Request{Activate: true}); err == nil || !strings.Contains(err.Error(), "Caddy reload failed") {
t.Fatalf("expected Caddy reload failure, got %v", err)
}
current, err := resolveReleaseLink(cfg.Deployment.Root, settings.CurrentLink)
if err != nil || current != old {
t.Fatalf("current=%q err=%v", current, err)
}
body, err := os.ReadFile(handler)
if err != nil || string(body) != "reverse_proxy 127.0.0.1:8092\n" {
t.Fatalf("handler=%q err=%v", body, err)
}
if operator.active[testCandidateUnit] {
t.Fatal("candidate was not stopped after route restoration")
}
})
}
}
func TestSingletonIncompleteRecoveryKeepsProvenCandidateRunning(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
settings, _ := singletonSettings(t, cfg, old)
cfg.Deployment.Singleton = settings
operator := &fakeOperator{active: map[string]bool{"example-site.service": true}, failReload: true}
_, err := manager(operator, fresh).Deploy(context.Background(), cfg, Request{Activate: true})
if err == nil || !strings.Contains(err.Error(), "candidate remains routed for operator recovery") {
t.Fatalf("expected explicit incomplete recovery, got %v", err)
}
if !operator.active[testCandidateUnit] {
t.Fatal("proven candidate was stopped despite incomplete route restoration")
}
record, loadErr := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
lease, leaseErr := state.LoadCandidateLease(state.CandidateLeasePath(cfg.Deployment.StateFile), cfg.Deployment.Root, cfg.Service.Name)
if loadErr != nil || leaseErr != nil || record.CandidateRelease != fresh || lease.OperationID != strings.Repeat("d", 32) || lease.Unit != testCandidateUnit {
t.Fatalf("retained state=%+v err=%v", record, loadErr)
}
}
func TestReconcileReportsRetainedCandidateWithoutMutation(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
settings, handler := singletonSettings(t, cfg, old)
cfg.Deployment.Singleton = settings
at := time.Unix(100, 0).UTC().Format(time.RFC3339)
record := state.Record{SchemaVersion: state.SchemaVersion, Strategy: cfg.Deployment.Strategy, DesiredRelease: fresh, CandidateRelease: fresh, ActiveSlot: "singleton", ActiveRelease: old, LastAttemptRelease: fresh, LastAttemptOutcome: "failed", LastAttemptAt: at, UpdatedAt: at}
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
t.Fatal(err)
}
storeTestCandidateLease(t, cfg, fresh, at)
candidateHandler, err := renderHandler(settings.CaddyHandlerTemplate, settings.CandidateAddress)
if err != nil {
t.Fatal(err)
}
if err := os.WriteFile(handler, candidateHandler, 0o640); err != nil {
t.Fatal(err)
}
operator := &fakeOperator{active: map[string]bool{settings.Unit: true, testCandidateUnit: true}}
reconciliation, err := manager(operator, fresh).Reconcile(context.Background(), cfg)
if err != nil || reconciliation.Mutation != "none" || !reconciliation.Observed.CandidateLease || reconciliation.Observed.CandidateUnitActive == nil || !*reconciliation.Observed.CandidateUnitActive || reconciliation.Observed.HandlerFileTarget != "candidate" || reconciliation.Disposition != "retained_candidate_handler_file" {
t.Fatalf("reconciliation=%+v err=%v", reconciliation, err)
}
if len(operator.stops) != 0 || len(operator.restarts) != 0 || operator.reloads != 0 {
t.Fatalf("reconcile mutated services: stops=%#v restarts=%#v reloads=%d", operator.stops, operator.restarts, operator.reloads)
}
}
func TestSingletonRollbackRetainsOperationLeaseWhenRecoveryIsIncomplete(t *testing.T) {
cfg, old, fresh := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, old)
operator := &fakeOperator{active: map[string]bool{cfg.Deployment.Singleton.Unit: true}}
m := manager(operator, fresh)
if _, err := m.Deploy(context.Background(), cfg, Request{Activate: true, ApprovedSHA256: strings.Repeat("a", 64)}); err != nil {
t.Fatal(err)
}
operator.failReload = true
if _, err := m.Rollback(context.Background(), cfg); err == nil || !strings.Contains(err.Error(), "candidate remains routed for operator recovery") {
t.Fatalf("expected retained rollback candidate, got %v", err)
}
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err != nil {
t.Fatal(err)
}
lease, leaseErr := state.LoadCandidateLease(state.CandidateLeasePath(cfg.Deployment.StateFile), cfg.Deployment.Root, cfg.Service.Name)
if record.ActiveRelease != fresh || record.CandidateRelease != old || leaseErr != nil || lease.OperationID != strings.Repeat("d", 32) || lease.Unit != testCandidateUnit || record.LastAttemptOutcome != "failed" {
t.Fatalf("rollback state=%+v", record)
}
if !operator.active[testCandidateUnit] {
t.Fatal("rollback candidate was stopped despite incomplete recovery")
}
}
func TestPruneProtectsOperationScopedCandidateRelease(t *testing.T) {
cfg, active, candidate := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, active)
at := time.Unix(100, 0).UTC().Format(time.RFC3339)
record := state.Record{SchemaVersion: state.SchemaVersion, Strategy: cfg.Deployment.Strategy, DesiredRelease: candidate, CandidateRelease: candidate, ActiveSlot: "singleton", ActiveRelease: active, LastAttemptRelease: candidate, LastAttemptOutcome: "failed", LastAttemptAt: at, UpdatedAt: at}
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
t.Fatal(err)
}
storeTestCandidateLease(t, cfg, candidate, at)
removed, err := manager(&fakeOperator{active: map[string]bool{}}, candidate).Prune(cfg, 2, true)
if err != nil {
t.Fatal(err)
}
if len(removed) != 0 {
t.Fatalf("candidate release was selected for pruning: %#v", removed)
}
if info, err := os.Stat(candidate); err != nil || !info.IsDir() {
t.Fatalf("candidate release was not preserved: %v", err)
}
}
+32
View File
@@ -0,0 +1,32 @@
//go:build linux
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"errors"
"os"
"syscall"
)
type fileLock struct{ file *os.File }
func acquireLock(path string) (*fileLock, error) {
file, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0o600)
if err != nil {
return nil, err
}
if err := syscall.Flock(int(file.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil {
_ = file.Close()
return nil, errors.New("deployment lock is held")
}
return &fileLock{file: file}, nil
}
func (l *fileLock) Close() error {
if l == nil || l.file == nil {
return nil
}
_ = syscall.Flock(int(l.file.Fd()), syscall.LOCK_UN)
return l.file.Close()
}
+30
View File
@@ -0,0 +1,30 @@
// SPDX-License-Identifier: AGPL-3.0-only
//go:build linux
package deploy
import (
"path/filepath"
"testing"
)
func TestHostWideLockSerializesIndependentServices(t *testing.T) {
path := filepath.Join(t.TempDir(), "tend-deploy.lock")
first, err := acquireLock(path)
if err != nil {
t.Fatal(err)
}
defer first.Close()
if second, err := acquireLock(path); err == nil {
_ = second.Close()
t.Fatal("second service acquired the shared activation lock")
}
if err = first.Close(); err != nil {
t.Fatal(err)
}
third, err := acquireLock(path)
if err != nil {
t.Fatal(err)
}
_ = third.Close()
}
+14
View File
@@ -0,0 +1,14 @@
//go:build !linux
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import "errors"
type fileLock struct{}
func acquireLock(string) (*fileLock, error) {
return nil, errors.New("deployment mutations require Linux")
}
func (*fileLock) Close() error { return nil }
+125
View File
@@ -0,0 +1,125 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"sort"
"strings"
"time"
"gamertan.com/tend/internal/process"
)
type Operator interface {
Restart(context.Context, string) error
Stop(context.Context, string) error
IsActive(context.Context, string) (bool, error)
StartCandidate(context.Context, string, string, string, map[string]string) error
ValidateCaddy(context.Context, string) error
ReloadCaddy(context.Context) error
Probe(context.Context, string, string, string, string) error
ProbeURL(context.Context, string, string) error
}
type SystemOperator struct {
Runner process.Runner
Timeout time.Duration
}
func (o SystemOperator) Restart(ctx context.Context, unit string) error {
_, err := o.Runner.Run(ctx, "/", nil, "systemctl", "restart", unit)
return err
}
func (o SystemOperator) Stop(ctx context.Context, unit string) error {
_, err := o.Runner.Run(ctx, "/", nil, "systemctl", "stop", unit)
return err
}
func (o SystemOperator) IsActive(ctx context.Context, unit string) (bool, error) {
out, err := o.Runner.Run(ctx, "/", nil, "systemctl", "is-active", unit)
if err != nil {
if strings.TrimSpace(string(out)) == "inactive" || strings.TrimSpace(string(out)) == "failed" {
return false, nil
}
return false, err
}
return strings.TrimSpace(string(out)) == "active", nil
}
func (o SystemOperator) StartCandidate(ctx context.Context, unit, binary, environmentFile string, env map[string]string) error {
args := []string{
"--unit", unit, "--collect",
"--property=DynamicUser=yes", "--property=NoNewPrivileges=yes",
"--property=PrivateDevices=yes", "--property=PrivateTmp=yes",
"--property=ProtectClock=yes", "--property=ProtectControlGroups=yes",
"--property=ProtectHome=yes", "--property=ProtectHostname=yes",
"--property=ProtectKernelLogs=yes", "--property=ProtectKernelModules=yes",
"--property=ProtectKernelTunables=yes", "--property=ProtectSystem=strict",
"--property=RestrictAddressFamilies=AF_INET AF_INET6 AF_UNIX",
"--property=RestrictNamespaces=yes", "--property=RestrictRealtime=yes",
"--property=RestrictSUIDSGID=yes", "--property=LockPersonality=yes",
"--property=MemoryDenyWriteExecute=yes", "--property=CapabilityBoundingSet=",
"--property=AmbientCapabilities=",
"--property=EnvironmentFile=" + environmentFile,
}
keys := make([]string, 0, len(env))
for key := range env {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
args = append(args, "--setenv", key+"="+env[key])
}
args = append(args, "--", binary)
_, err := o.Runner.Run(ctx, "/", nil, "systemd-run", args...)
return err
}
func (o SystemOperator) ValidateCaddy(ctx context.Context, path string) error {
_, err := o.Runner.Run(ctx, "/", nil, "caddy", "validate", "--config", path, "--adapter", "caddyfile")
return err
}
func (o SystemOperator) ReloadCaddy(ctx context.Context) error {
_, err := o.Runner.Run(ctx, "/", nil, "systemctl", "reload", "caddy.service")
return err
}
func (o SystemOperator) Probe(ctx context.Context, address, host, path, contains string) error {
u := url.URL{Scheme: "http", Host: address, Path: path}
return o.probeRequest(ctx, u.String(), host, contains)
}
func (o SystemOperator) ProbeURL(ctx context.Context, value, contains string) error {
u, err := url.Parse(value)
if err != nil || u.Scheme != "https" || u.Host == "" || u.User != nil || u.Fragment != "" {
return errors.New("public probe URL is invalid")
}
return o.probeRequest(ctx, u.String(), "", contains)
}
func (o SystemOperator) probeRequest(ctx context.Context, value, host, contains string) error {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, value, nil)
if err != nil {
return err
}
if host != "" {
req.Host = host
}
client := &http.Client{Timeout: o.Timeout, CheckRedirect: func(*http.Request, []*http.Request) error { return errors.New("redirect refused") }}
response, err := client.Do(req)
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode != http.StatusOK {
return fmt.Errorf("probe returned HTTP %d", response.StatusCode)
}
body, err := io.ReadAll(io.LimitReader(response.Body, 1<<20))
if err != nil {
return err
}
if contains != "" && !strings.Contains(string(body), contains) {
return errors.New("probe response omitted required marker")
}
return nil
}
+49
View File
@@ -0,0 +1,49 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"context"
"reflect"
"strings"
"testing"
)
type recordingRunner struct {
name string
args []string
}
func (r *recordingRunner) Run(_ context.Context, _ string, _ map[string]string, name string, args ...string) ([]byte, error) {
r.name = name
r.args = append([]string(nil), args...)
return nil, nil
}
func TestStartCandidateUsesArgumentVectorAndHardenedUnit(t *testing.T) {
runner := &recordingRunner{}
operator := SystemOperator{Runner: runner}
env := map[string]string{"Z_ENV": "safe value", "A_ENV": "first"}
if err := operator.StartCandidate(context.Background(), "example-tend-candidate.service", "/opt/example/releases/sha256-a/app", "/etc/tend/environment/example.env", env); err != nil {
t.Fatal(err)
}
if runner.name != "systemd-run" {
t.Fatalf("command=%q", runner.name)
}
required := []string{"--property=DynamicUser=yes", "--property=NoNewPrivileges=yes", "--property=ProtectSystem=strict", "--property=MemoryDenyWriteExecute=yes", "--property=CapabilityBoundingSet=", "--property=EnvironmentFile=/etc/tend/environment/example.env", "--setenv", "A_ENV=first", "--setenv", "Z_ENV=safe value", "--", "/opt/example/releases/sha256-a/app"}
cursor := 0
for _, arg := range runner.args {
if cursor < len(required) && arg == required[cursor] {
cursor++
}
}
if cursor != len(required) {
t.Fatalf("arguments omitted ordered security boundary: %#v", runner.args)
}
if reflect.DeepEqual(runner.args, []string{"sh", "-c"}) {
t.Fatal("candidate command used a shell")
}
if strings.Contains(strings.Join(runner.args, "\n"), "SUPER_SECRET") {
t.Fatal("candidate arguments exposed a secret value")
}
}
+38
View File
@@ -0,0 +1,38 @@
//go:build linux
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"os"
"syscall"
)
type fileIdentity struct {
mode os.FileMode
uid, gid int
owned bool
}
func identityFor(info os.FileInfo, fallback os.FileMode) fileIdentity {
identity := fileIdentity{mode: fallback}
if info == nil {
return identity
}
identity.mode = info.Mode().Perm()
if stat, ok := info.Sys().(*syscall.Stat_t); ok {
identity.uid = int(stat.Uid)
identity.gid = int(stat.Gid)
identity.owned = true
}
return identity
}
func applyIdentity(file *os.File, identity fileIdentity) error {
if identity.owned {
if err := file.Chown(identity.uid, identity.gid); err != nil {
return err
}
}
return file.Chmod(identity.mode)
}
+17
View File
@@ -0,0 +1,17 @@
//go:build !linux
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import "os"
type fileIdentity struct{ mode os.FileMode }
func identityFor(info os.FileInfo, fallback os.FileMode) fileIdentity {
if info != nil {
return fileIdentity{mode: info.Mode().Perm()}
}
return fileIdentity{mode: fallback}
}
func applyIdentity(file *os.File, identity fileIdentity) error { return file.Chmod(identity.mode) }
+270
View File
@@ -0,0 +1,270 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"bytes"
"context"
"errors"
"os"
"gamertan.com/tend/internal/config"
"gamertan.com/tend/internal/state"
)
type Finding struct {
Code string `json:"code"`
Severity string `json:"severity"`
Message string `json:"message"`
}
type ObservedReleaseIdentity struct {
Version string `json:"version"`
Commit string `json:"commit"`
}
type ObservedState struct {
ActiveRelease string `json:"active_release,omitempty"`
PreviousRelease string `json:"previous_release,omitempty"`
ActiveIdentity *ObservedReleaseIdentity `json:"active_identity,omitempty"`
Units map[string]bool `json:"units"`
CandidateLease bool `json:"candidate_lease"`
CandidateUnit string `json:"candidate_unit,omitempty"`
CandidateUnitActive *bool `json:"candidate_unit_active,omitempty"`
LegacyCandidateUnit string `json:"legacy_candidate_unit,omitempty"`
LegacyCandidateActive *bool `json:"legacy_candidate_unit_active,omitempty"`
RouteHandlerMatches bool `json:"route_handler_matches"`
HandlerFileTarget string `json:"handler_file_target,omitempty"`
}
type Reconciliation struct {
Service string `json:"service"`
Strategy string `json:"strategy"`
Mutation string `json:"mutation"`
Consistent bool `json:"consistent"`
StateInitialized bool `json:"state_initialized"`
State *state.Record `json:"state,omitempty"`
Observed ObservedState `json:"observed"`
Disposition string `json:"disposition"`
Findings []Finding `json:"findings"`
}
// Reconcile observes configured units, release pointers, installed release
// identity, and the imported Caddy handler. It never acquires the deployment
// lock or mutates service state; proposed repairs remain an operator decision.
func (m Manager) Reconcile(ctx context.Context, cfg config.Config) (Reconciliation, error) {
if err := cfg.Validate(); err != nil {
return Reconciliation{}, err
}
if m.Operator == nil || m.ReadIdentity == nil {
return Reconciliation{}, errors.New("reconciliation dependencies are unavailable")
}
report := Reconciliation{
Service: cfg.Service.Name, Strategy: cfg.Deployment.Strategy, Mutation: "none",
Observed: ObservedState{Units: map[string]bool{}}, Findings: []Finding{},
}
add := func(code, severity, message string) {
report.Findings = append(report.Findings, Finding{Code: code, Severity: severity, Message: message})
}
record, err := state.Load(cfg.Deployment.StateFile, cfg.Deployment.Root, cfg.Deployment.Strategy)
if err == nil {
report.State = &record
report.StateInitialized = true
} else if os.IsNotExist(err) {
add("state_uninitialized", "warning", "Tend has no validated state record for this service.")
} else {
add("state_invalid", "error", "The Tend state record could not be validated.")
}
activeSlot := ""
activeUnit := ""
activeAddress := ""
handler := ""
template := ""
var candidateLease *state.CandidateLease
switch cfg.Deployment.Strategy {
case "singleton_candidate":
single := *cfg.Deployment.Singleton
activeSlot = "singleton"
activeUnit = single.Unit
activeAddress = single.Address
handler, template = single.CaddyHandler, single.CaddyHandlerTemplate
report.Observed.ActiveRelease = observeReleaseLink(cfg, single.CurrentLink, true, add)
report.Observed.PreviousRelease = observeReleaseLink(cfg, single.PreviousLink, false, add)
legacyCandidateUnit := cfg.Service.Name + "-tend-candidate.service"
candidateUnit := legacyCandidateUnit
lease, leaseErr := state.LoadCandidateLease(state.CandidateLeasePath(cfg.Deployment.StateFile), cfg.Deployment.Root, cfg.Service.Name)
if leaseErr == nil {
candidateLease = &lease
report.Observed.CandidateLease = true
candidateUnit = lease.Unit
} else if !os.IsNotExist(leaseErr) {
report.Observed.CandidateLease = true
candidateUnit = ""
add("candidate_lease_invalid", "error", "The operation-scoped candidate lease could not be validated.")
} else if report.State != nil && report.State.CandidateRelease != "" {
report.Observed.CandidateLease = true
candidateUnit = ""
add("legacy_candidate_lease", "error", "The state records a candidate release without an operation-scoped lease and requires manual review.")
}
if candidateUnit != "" {
report.Observed.CandidateUnit = candidateUnit
candidateActive, candidateErr := m.Operator.IsActive(ctx, candidateUnit)
if candidateErr != nil {
add("candidate_unit_unobservable", "error", "The transient candidate unit state could not be observed.")
} else {
report.Observed.CandidateUnitActive = &candidateActive
report.Observed.Units[candidateUnit] = candidateActive
}
}
if candidateUnit != legacyCandidateUnit {
report.Observed.LegacyCandidateUnit = legacyCandidateUnit
legacyActive, legacyErr := m.Operator.IsActive(ctx, legacyCandidateUnit)
if legacyErr != nil {
add("legacy_candidate_unit_unobservable", "error", "The legacy fixed candidate unit state could not be observed.")
} else {
report.Observed.LegacyCandidateActive = &legacyActive
report.Observed.Units[legacyCandidateUnit] = legacyActive
if legacyActive {
add("legacy_candidate_unit_active", "error", "A legacy fixed-name candidate remains active beside an operation-scoped lease.")
}
}
}
case "blue_green":
blueGreen := *cfg.Deployment.BlueGreen
handler, template = blueGreen.CaddyHandler, blueGreen.CaddyHandlerTemplate
activeSlot = blueGreen.BootstrapActive
if report.State != nil {
activeSlot = report.State.ActiveSlot
}
active := slotConfig(blueGreen, activeSlot)
previousName := "blue"
if activeSlot == "blue" {
previousName = "green"
}
previous := slotConfig(blueGreen, previousName)
activeUnit, activeAddress = active.Unit, active.Address
report.Observed.ActiveRelease = observeReleaseLink(cfg, active.Link, true, add)
report.Observed.PreviousRelease = observeReleaseLink(cfg, previous.Link, false, add)
for _, slot := range []config.Slot{blueGreen.Blue, blueGreen.Green} {
observeUnit(ctx, m.Operator, slot.Unit, report.Observed.Units, add)
}
}
if _, exists := report.Observed.Units[activeUnit]; !exists {
observeUnit(ctx, m.Operator, activeUnit, report.Observed.Units, add)
}
if active, observed := report.Observed.Units[activeUnit]; activeUnit != "" && observed && !active {
add("active_unit_inactive", "error", "The configured active service unit is not active.")
}
if report.Observed.ActiveRelease != "" {
identity, identityErr := m.ReadIdentity(report.Observed.ActiveRelease)
if identityErr != nil {
add("active_identity_unreadable", "error", "The observed active release identity could not be validated.")
} else {
report.Observed.ActiveIdentity = &ObservedReleaseIdentity{Version: identity.Version, Commit: identity.Commit}
}
}
expectedHandler, renderErr := renderHandler(template, activeAddress)
actualHandler, readErr := os.ReadFile(handler)
if renderErr != nil || readErr != nil {
add("route_handler_unreadable", "error", "The configured Caddy handler or its template could not be validated.")
} else {
report.Observed.RouteHandlerMatches = bytes.Equal(expectedHandler, actualHandler)
if report.Observed.RouteHandlerMatches {
report.Observed.HandlerFileTarget = "installed"
} else if cfg.Deployment.Strategy == "singleton_candidate" {
candidateAddress := cfg.Deployment.Singleton.CandidateAddress
if candidateLease != nil {
candidateAddress = candidateLease.Address
}
candidateHandler, candidateErr := renderHandler(template, candidateAddress)
if candidateErr == nil && bytes.Equal(candidateHandler, actualHandler) {
report.Observed.HandlerFileTarget = "candidate"
} else {
report.Observed.HandlerFileTarget = "other"
}
}
if !report.Observed.RouteHandlerMatches {
add("route_handler_drift", "error", "The installed Caddy handler does not match the configured active upstream.")
}
}
if report.State != nil {
if report.State.ActiveSlot != activeSlot || report.State.ActiveRelease != report.Observed.ActiveRelease {
add("active_release_drift", "error", "Recorded active state does not match the observed active release pointer.")
}
if report.State.PreviousRelease != report.Observed.PreviousRelease {
add("previous_release_drift", "warning", "Recorded rollback state does not match the observed previous release pointer.")
}
if cfg.Deployment.Strategy == "singleton_candidate" {
leased := report.Observed.CandidateLease && candidateLease != nil
candidateActive := report.Observed.CandidateUnitActive != nil && *report.Observed.CandidateUnitActive
if candidateLease != nil && report.State.CandidateRelease == "" {
add("candidate_lease_without_running_attempt", "error", "An operation-scoped candidate lease remains after the recorded attempt settled.")
}
if candidateLease != nil && report.State.CandidateRelease != "" && candidateLease.Release != report.State.CandidateRelease {
add("candidate_lease_release_mismatch", "error", "The operation-scoped candidate lease does not match the recorded candidate release.")
}
if leased && report.Observed.CandidateUnitActive != nil && !candidateActive {
add("inactive_candidate_lease", "error", "State retains a candidate lease but its operation-scoped unit is inactive.")
}
if !report.Observed.CandidateLease && candidateActive {
add("unleased_candidate_active", "error", "A transient candidate unit is active without a matching running attempt.")
}
if leased && candidateActive && report.Observed.HandlerFileTarget == "candidate" {
add("retained_candidate_routed", "warning", "The retained candidate appears in the handler file; do not stop it before establishing another healthy route.")
}
if leased && candidateActive && report.Observed.HandlerFileTarget != "candidate" {
add("candidate_active_not_routed", "warning", "The leased candidate is active but the handler file does not target it; cleanup remains an explicit reviewed operation.")
}
}
}
switch {
case !report.StateInitialized:
report.Disposition = "state_uninitialized"
case report.Observed.CandidateLease && report.Observed.CandidateUnit == "":
report.Disposition = "legacy_or_invalid_candidate_lease"
case report.Observed.CandidateLease && report.Observed.CandidateUnitActive != nil && *report.Observed.CandidateUnitActive && report.Observed.HandlerFileTarget == "candidate":
report.Disposition = "retained_candidate_handler_file"
case report.Observed.CandidateLease && report.Observed.CandidateUnitActive != nil && *report.Observed.CandidateUnitActive:
report.Disposition = "candidate_active_not_in_handler_file"
case report.Observed.CandidateLease:
report.Disposition = "inactive_candidate_lease"
case len(report.Findings) == 0:
report.Disposition = "settled"
default:
report.Disposition = "manual_review_required"
}
report.Consistent = len(report.Findings) == 0
return report, nil
}
func observeReleaseLink(cfg config.Config, link string, required bool, add func(string, string, string)) string {
release, err := resolveReleaseLink(cfg.Deployment.Root, link)
if err == nil {
return release
}
if !required && os.IsNotExist(err) {
return ""
}
code := "previous_pointer_unreadable"
message := "The configured previous release pointer could not be validated."
severity := "warning"
if required {
code = "active_pointer_unreadable"
message = "The configured active release pointer could not be validated."
severity = "error"
}
add(code, severity, message)
return ""
}
func observeUnit(ctx context.Context, operator Operator, unit string, units map[string]bool, add func(string, string, string)) {
active, err := operator.IsActive(ctx, unit)
if err != nil {
add("unit_unobservable", "error", "A configured service unit state could not be observed.")
return
}
units[unit] = active
}
+182
View File
@@ -0,0 +1,182 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"context"
"os"
"path/filepath"
"strings"
"testing"
"time"
"gamertan.com/tend/internal/config"
"gamertan.com/tend/internal/state"
)
func writeReleaseIdentity(t *testing.T, release, version, commit string) {
t.Helper()
body := `{"version":"` + version + `","commit":"` + commit + `"}`
if err := os.WriteFile(filepath.Join(release, "RELEASE.json"), []byte(body), 0o644); err != nil {
t.Fatal(err)
}
}
func storeSingletonState(t *testing.T, cfg config.Config, active, previous string) {
t.Helper()
record := state.Record{
SchemaVersion: state.SchemaVersion,
Strategy: "singleton_candidate",
DesiredRelease: active,
ActiveSlot: "singleton",
ActiveRelease: active,
PreviousRelease: previous,
LastAttemptRelease: active,
LastAttemptOutcome: "succeeded",
LastAttemptAt: time.Unix(100, 0).UTC().Format(time.RFC3339),
UpdatedAt: time.Unix(100, 0).UTC().Format(time.RFC3339),
}
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
t.Fatal(err)
}
}
func reconciliationManager(operator Operator) Manager {
return Manager{Operator: operator, ReadIdentity: readReleaseIdentity}
}
func TestReconcileReportsHealthySingletonWithoutMutation(t *testing.T) {
cfg, active, previous := baseConfig(t, "singleton_candidate")
cfg.Deployment.Singleton, _ = singletonSettings(t, cfg, active)
if err := replaceSymlink(cfg.Deployment.Singleton.PreviousLink, previous); err != nil {
t.Fatal(err)
}
commit := strings.Repeat("a", 40)
writeReleaseIdentity(t, active, "v0.2.0-preview.2", commit)
storeSingletonState(t, cfg, active, previous)
operator := &fakeOperator{active: map[string]bool{
cfg.Deployment.Singleton.Unit: true,
cfg.Service.Name + "-tend-candidate.service": false,
}}
report, err := reconciliationManager(operator).Reconcile(context.Background(), cfg)
if err != nil {
t.Fatal(err)
}
if !report.Consistent || report.Mutation != "none" || !report.StateInitialized || len(report.Findings) != 0 {
t.Fatalf("report=%+v", report)
}
if report.Observed.ActiveRelease != active || report.Observed.PreviousRelease != previous || !report.Observed.RouteHandlerMatches {
t.Fatalf("observed=%+v", report.Observed)
}
if report.Observed.ActiveIdentity == nil || report.Observed.ActiveIdentity.Version != "v0.2.0-preview.2" || report.Observed.ActiveIdentity.Commit != commit {
t.Fatalf("identity=%+v", report.Observed.ActiveIdentity)
}
if !report.Observed.Units[cfg.Deployment.Singleton.Unit] || report.Observed.CandidateUnitActive == nil || *report.Observed.CandidateUnitActive {
t.Fatalf("units=%#v candidate=%v", report.Observed.Units, report.Observed.CandidateUnitActive)
}
}
func TestReconcileReportsHealthyBlueGreenDeployment(t *testing.T) {
cfg, active, previous := baseConfig(t, "blue_green")
handler := filepath.Join(cfg.Deployment.Root, "handler.caddy")
template := filepath.Join(cfg.Deployment.Root, "handler.template")
if err := os.WriteFile(handler, []byte("reverse_proxy 127.0.0.1:8090\n"), 0o640); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(template, []byte("reverse_proxy {{UPSTREAM}}\n"), 0o644); err != nil {
t.Fatal(err)
}
blue := filepath.Join(cfg.Deployment.Root, "slots", "blue")
green := filepath.Join(cfg.Deployment.Root, "slots", "green")
if err := replaceSymlink(blue, active); err != nil {
t.Fatal(err)
}
if err := replaceSymlink(green, previous); err != nil {
t.Fatal(err)
}
cfg.Deployment.BlueGreen = &config.BlueGreen{
CaddyConfig: filepath.Join(cfg.Deployment.Root, "Caddyfile"), CaddyHandler: handler,
CaddyHandlerTemplate: template, BootstrapActive: "blue",
Blue: config.Slot{Unit: "example-blue.service", Address: "127.0.0.1:8090", Link: blue},
Green: config.Slot{Unit: "example-green.service", Address: "127.0.0.1:8091", Link: green},
}
writeReleaseIdentity(t, active, "v0.2.0-preview.2", strings.Repeat("c", 40))
record := state.Record{
SchemaVersion: state.SchemaVersion, Strategy: "blue_green", DesiredRelease: active,
ActiveSlot: "blue", ActiveRelease: active, PreviousSlot: "green", PreviousRelease: previous,
LastAttemptRelease: active, LastAttemptOutcome: "succeeded",
LastAttemptAt: time.Unix(100, 0).UTC().Format(time.RFC3339), UpdatedAt: time.Unix(100, 0).UTC().Format(time.RFC3339),
}
if err := state.Store(cfg.Deployment.StateFile, cfg.Deployment.Root, record); err != nil {
t.Fatal(err)
}
operator := &fakeOperator{active: map[string]bool{
"example-blue.service": true, "example-green.service": true,
}}
report, err := reconciliationManager(operator).Reconcile(context.Background(), cfg)
if err != nil {
t.Fatal(err)
}
if !report.Consistent || report.Mutation != "none" || report.Observed.ActiveRelease != active || report.Observed.PreviousRelease != previous || !report.Observed.RouteHandlerMatches {
t.Fatalf("report=%+v", report)
}
if report.Observed.CandidateUnit != "" || report.Observed.CandidateUnitActive != nil {
t.Fatalf("unexpected candidate observation=%+v", report.Observed)
}
}
func TestReconcileExplainsDriftWithoutRepairingIt(t *testing.T) {
cfg, recorded, observed := baseConfig(t, "singleton_candidate")
var handler string
cfg.Deployment.Singleton, handler = singletonSettings(t, cfg, recorded)
if err := replaceSymlink(cfg.Deployment.Singleton.PreviousLink, observed); err != nil {
t.Fatal(err)
}
storeSingletonState(t, cfg, recorded, observed)
if err := replaceSymlink(cfg.Deployment.Singleton.CurrentLink, observed); err != nil {
t.Fatal(err)
}
writeReleaseIdentity(t, observed, "v0.2.0-preview.3", strings.Repeat("b", 40))
candidateHandler, err := renderHandler(cfg.Deployment.Singleton.CaddyHandlerTemplate, cfg.Deployment.Singleton.CandidateAddress)
if err != nil {
t.Fatal(err)
}
if err := os.WriteFile(handler, candidateHandler, 0o640); err != nil {
t.Fatal(err)
}
beforeHandler, err := os.ReadFile(handler)
if err != nil {
t.Fatal(err)
}
operator := &fakeOperator{active: map[string]bool{
cfg.Deployment.Singleton.Unit: false,
cfg.Service.Name + "-tend-candidate.service": true,
}}
report, err := reconciliationManager(operator).Reconcile(context.Background(), cfg)
if err != nil {
t.Fatal(err)
}
if report.Consistent || report.Mutation != "none" {
t.Fatalf("report=%+v", report)
}
codes := map[string]bool{}
for _, finding := range report.Findings {
codes[finding.Code] = true
}
for _, code := range []string{"active_release_drift", "active_unit_inactive", "route_handler_drift", "unleased_candidate_active"} {
if !codes[code] {
t.Fatalf("missing %s in %#v", code, report.Findings)
}
}
target, err := resolveReleaseLink(cfg.Deployment.Root, cfg.Deployment.Singleton.CurrentLink)
if err != nil || target != observed {
t.Fatalf("current=%q err=%v", target, err)
}
afterHandler, err := os.ReadFile(handler)
if err != nil || string(afterHandler) != string(beforeHandler) {
t.Fatalf("handler changed err=%v", err)
}
}
+318
View File
@@ -0,0 +1,318 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"archive/tar"
"compress/gzip"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"regexp"
"sort"
"strings"
"gamertan.com/tend/internal/config"
"gamertan.com/tend/internal/packager"
)
const maxArtifactSize int64 = 512 << 20
func prepareRelease(cfg config.Config, artifact, expected, approved string) (string, error) {
if err := checkArtifact(artifact, expected, approved); err != nil {
return "", err
}
if err := ensureTree(cfg.Deployment.Root); err != nil {
return "", err
}
releases := filepath.Join(cfg.Deployment.Root, "releases")
release := filepath.Join(releases, "sha256-"+expected)
if info, err := os.Lstat(release); err == nil {
if !info.IsDir() || info.Mode()&os.ModeSymlink != 0 {
return "", errors.New("existing release is not a directory")
}
if err := validateRelease(cfg, release); err != nil {
return "", err
}
return release, nil
} else if !os.IsNotExist(err) {
return "", err
}
stage := filepath.Join(releases, ".tend-stage-"+expected)
if err := os.Mkdir(stage, 0o700); err != nil {
return "", err
}
ok := false
defer func() {
if !ok {
_ = os.RemoveAll(stage)
}
}()
if err := extractArtifact(artifact, stage); err != nil {
return "", err
}
if err := validateRelease(cfg, stage); err != nil {
return "", err
}
if err := os.Chmod(stage, 0o755); err != nil {
return "", err
}
if err := os.Rename(stage, release); err != nil {
return "", err
}
ok = true
return release, nil
}
func inspectArtifact(cfg config.Config, artifact, expected, approved string) error {
if err := checkArtifact(artifact, expected, approved); err != nil {
return err
}
stage, err := os.MkdirTemp("", "tend-inspect-")
if err != nil {
return err
}
defer os.RemoveAll(stage)
if err := extractArtifact(artifact, stage); err != nil {
return err
}
return validateRelease(cfg, stage)
}
func checkArtifact(artifact, expected, approved string) error {
if expected != approved || len(expected) != 64 {
return errors.New("artifact digest was not explicitly approved")
}
if _, err := hex.DecodeString(expected); err != nil {
return errors.New("artifact digest is not hexadecimal")
}
if !filepath.IsAbs(artifact) || filepath.Clean(artifact) != artifact {
return errors.New("artifact path must be a clean absolute path")
}
info, err := os.Lstat(artifact)
if err != nil {
return err
}
if !info.Mode().IsRegular() || info.Mode()&os.ModeSymlink != 0 || info.Size() > maxArtifactSize {
return errors.New("artifact must be a bounded regular file")
}
actual, err := fileSHA(artifact)
if err != nil {
return err
}
if actual != expected {
return errors.New("artifact digest does not match")
}
return nil
}
func ensureTree(root string) error {
if os.Geteuid() != 0 && strings.HasPrefix(root, "/opt/") {
return errors.New("deployment under /opt requires root")
}
if err := rejectSymlinkAncestors(root); err != nil {
return err
}
if err := os.MkdirAll(filepath.Join(root, "releases"), 0o755); err != nil {
return err
}
return rejectSymlinkAncestors(filepath.Join(root, "releases"))
}
func rejectSymlinkAncestors(path string) error {
clean := filepath.Clean(path)
parts := strings.Split(strings.TrimPrefix(clean, string(filepath.Separator)), string(filepath.Separator))
current := string(filepath.Separator)
for _, part := range parts {
current = filepath.Join(current, part)
info, err := os.Lstat(current)
if os.IsNotExist(err) {
continue
}
if err != nil {
return err
}
if info.Mode()&os.ModeSymlink != 0 {
return fmt.Errorf("symlink ancestor refused: %s", current)
}
if !info.IsDir() {
return fmt.Errorf("non-directory ancestor refused: %s", current)
}
}
return nil
}
func extractArtifact(artifact, stage string) error {
file, err := os.Open(artifact)
if err != nil {
return err
}
defer file.Close()
gz, err := gzip.NewReader(file)
if err != nil {
return err
}
defer gz.Close()
tr := tar.NewReader(io.LimitReader(gz, maxArtifactSize))
files := 0
var total int64
for {
header, err := tr.Next()
if errors.Is(err, io.EOF) {
break
}
if err != nil {
return err
}
clean := filepath.Clean(filepath.FromSlash(header.Name))
parts := strings.Split(clean, string(filepath.Separator))
if len(parts) == 1 && header.Typeflag == tar.TypeDir {
continue
}
if len(parts) != 2 || parts[0] != "bundle" || parts[1] == "" || parts[1] == "." || parts[1] == ".." {
return fmt.Errorf("unsafe archive path %q", header.Name)
}
if header.Typeflag != tar.TypeReg || header.Size < 0 {
return errors.New("archive may contain only regular files")
}
files++
total += header.Size
if files > 16 || total > maxArtifactSize {
return errors.New("artifact exceeds extraction bounds")
}
target := filepath.Join(stage, parts[1])
mode := os.FileMode(0o644)
if !strings.HasSuffix(parts[1], ".json") && parts[1] != "SHA256SUMS" {
mode = 0o755
}
out, err := os.OpenFile(target, os.O_CREATE|os.O_EXCL|os.O_WRONLY, mode)
if err != nil {
return err
}
if _, err := io.CopyN(out, tr, header.Size); err != nil {
_ = out.Close()
return err
}
// OpenFile modes are filtered through the caller's umask. Tend is
// commonly invoked by a root account with umask 0077, while release
// binaries must remain executable by their dedicated service users.
// Reapply the validated, name-derived mode explicitly before the file
// becomes part of an immutable release.
if err := out.Chmod(mode); err != nil {
_ = out.Close()
return err
}
if err := out.Sync(); err != nil {
_ = out.Close()
return err
}
if err := out.Close(); err != nil {
return err
}
}
if files < 5 {
return errors.New("artifact is incomplete")
}
return nil
}
func validateRelease(cfg config.Config, release string) error {
manifestPath := filepath.Join(release, "RELEASE.json")
b, err := os.ReadFile(manifestPath)
if err != nil {
return err
}
dec := json.NewDecoder(strings.NewReader(string(b)))
dec.DisallowUnknownFields()
var manifest packager.Manifest
if err := dec.Decode(&manifest); err != nil {
return err
}
var trailing any
if err := dec.Decode(&trailing); !errors.Is(err, io.EOF) {
return errors.New("release manifest contains trailing data")
}
if manifest.SchemaVersion != 1 || manifest.Service != cfg.Service.Name || manifest.Binary != cfg.Build.Binary || manifest.GOOS != "linux" || manifest.GOARCH != "amd64" || manifest.CGOEnabled {
return errors.New("release manifest does not match configuration")
}
if matched, _ := regexp.MatchString(`^[0-9a-f]{40}$`, manifest.Commit); !matched {
return errors.New("release manifest commit is invalid")
}
binary := filepath.Join(release, cfg.Build.Binary)
sum, err := fileSHA(binary)
if err != nil {
return err
}
if sum != manifest.BinarySHA256 {
return errors.New("release binary digest does not match manifest")
}
if err := verifySums(release); err != nil {
return err
}
entries, err := os.ReadDir(release)
if err != nil {
return err
}
allowed := map[string]bool{cfg.Build.Binary: true, "BUILDINFO.json": true, "RELEASE.json": true, "SBOM.spdx.json": true, "SHA256SUMS": true}
if len(entries) != len(allowed) {
return errors.New("release contains unexpected files")
}
for _, entry := range entries {
if !allowed[entry.Name()] || !entry.Type().IsRegular() {
return fmt.Errorf("unexpected release entry %s", entry.Name())
}
}
return nil
}
func verifySums(release string) error {
b, err := os.ReadFile(filepath.Join(release, "SHA256SUMS"))
if err != nil {
return err
}
lines := strings.Split(strings.TrimSpace(string(b)), "\n")
if len(lines) != 4 {
return errors.New("SHA256SUMS must cover four release files")
}
seen := map[string]bool{}
for _, line := range lines {
fields := strings.Fields(line)
if len(fields) != 2 || len(fields[0]) != 64 {
return errors.New("malformed SHA256SUMS")
}
name := fields[1]
if filepath.Base(name) != name || seen[name] {
return errors.New("unsafe or duplicate checksum entry")
}
seen[name] = true
actual, err := fileSHA(filepath.Join(release, name))
if err != nil {
return err
}
if actual != fields[0] {
return fmt.Errorf("checksum mismatch for %s", name)
}
}
required := []string{"BUILDINFO.json", "RELEASE.json", "SBOM.spdx.json"}
sort.Strings(required)
for _, name := range required {
if !seen[name] {
return fmt.Errorf("checksum omitted %s", name)
}
}
return nil
}
func fileSHA(path string) (string, error) {
f, err := os.Open(path)
if err != nil {
return "", err
}
defer f.Close()
h := sha256.New()
if _, err := io.Copy(h, f); err != nil {
return "", err
}
return hex.EncodeToString(h.Sum(nil)), nil
}
@@ -0,0 +1,75 @@
// SPDX-License-Identifier: AGPL-3.0-only
//go:build linux
package deploy
import (
"archive/tar"
"compress/gzip"
"os"
"path/filepath"
"syscall"
"testing"
)
func TestExtractArtifactAppliesReleaseModesUnderRestrictiveUmask(t *testing.T) {
dir := t.TempDir()
artifact := filepath.Join(dir, "release.tar.gz")
file, err := os.OpenFile(artifact, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
if err != nil {
t.Fatal(err)
}
gz := gzip.NewWriter(file)
tw := tar.NewWriter(gz)
entries := map[string][]byte{
"bundle/app": []byte("executable"),
"bundle/BUILDINFO.json": []byte("{}"),
"bundle/RELEASE.json": []byte("{}"),
"bundle/SBOM.spdx.json": []byte("{}"),
"bundle/SHA256SUMS": []byte("checksums"),
}
for name, body := range entries {
if err := tw.WriteHeader(&tar.Header{Name: name, Typeflag: tar.TypeReg, Mode: 0o600, Size: int64(len(body))}); err != nil {
t.Fatal(err)
}
if _, err := tw.Write(body); err != nil {
t.Fatal(err)
}
}
if err := tw.Close(); err != nil {
t.Fatal(err)
}
if err := gz.Close(); err != nil {
t.Fatal(err)
}
if err := file.Close(); err != nil {
t.Fatal(err)
}
oldUmask := syscall.Umask(0o077)
t.Cleanup(func() { syscall.Umask(oldUmask) })
stage := filepath.Join(dir, "stage")
if err := os.Mkdir(stage, 0o700); err != nil {
t.Fatal(err)
}
if err := extractArtifact(artifact, stage); err != nil {
t.Fatal(err)
}
for name, want := range map[string]os.FileMode{
"app": 0o755,
"BUILDINFO.json": 0o644,
"RELEASE.json": 0o644,
"SBOM.spdx.json": 0o644,
"SHA256SUMS": 0o644,
} {
info, err := os.Stat(filepath.Join(stage, name))
if err != nil {
t.Fatal(err)
}
if got := info.Mode().Perm(); got != want {
t.Fatalf("%s mode=%#o want=%#o", name, got, want)
}
}
}
+41
View File
@@ -0,0 +1,41 @@
// SPDX-License-Identifier: AGPL-3.0-only
package deploy
import (
"archive/tar"
"compress/gzip"
"os"
"path/filepath"
"testing"
)
func writeHostileArchive(t *testing.T, name string, typeflag byte) {
t.Helper()
path := filepath.Join(t.TempDir(), "bad.tar.gz")
file, err := os.Create(path)
if err != nil {
t.Fatal(err)
}
gz := gzip.NewWriter(file)
tw := tar.NewWriter(gz)
body := []byte("x")
if err := tw.WriteHeader(&tar.Header{Name: name, Typeflag: typeflag, Mode: 0o644, Size: int64(len(body))}); err != nil {
t.Fatal(err)
}
if typeflag == tar.TypeReg {
_, _ = tw.Write(body)
}
_ = tw.Close()
_ = gz.Close()
_ = file.Close()
stage := filepath.Join(t.TempDir(), "stage")
_ = os.Mkdir(stage, 0o700)
if err := extractArtifact(path, stage); err == nil {
t.Fatalf("accepted hostile entry %q type %d", name, typeflag)
}
}
func TestExtractionRejectsTraversalAndLinks(t *testing.T) {
writeHostileArchive(t, "bundle/../../escape", tar.TypeReg)
writeHostileArchive(t, "bundle/link", tar.TypeSymlink)
}