949 lines
36 KiB
Go
949 lines
36 KiB
Go
package pi
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"regexp"
|
|
"runtime"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/aritmolab/thothii/tools/thothctl/internal/compose"
|
|
"github.com/distribution/reference"
|
|
)
|
|
|
|
var (
|
|
ErrConfirmationRequired = errors.New("update requires --yes after reviewing the planned Pi version")
|
|
ErrActiveSessions = errors.New("active sessions must be drained before updating Pi; use --drain only after they are complete")
|
|
ErrInterruptedUpdate = errors.New("a previous Pi update is incomplete; run pi rollback --yes before starting another update")
|
|
ErrInvalidRequest = errors.New("invalid Pi lifecycle request")
|
|
versionPattern = regexp.MustCompile(`^[0-9]+(?:\.[0-9]+){1,3}(?:[-+][0-9A-Za-z.-]+)?$`)
|
|
)
|
|
|
|
// Source chooses whether the candidate is built from this checkout or pulled from an immutable image.
|
|
type Source string
|
|
|
|
const (
|
|
BuildSource Source = "build"
|
|
PullSource Source = "pull"
|
|
)
|
|
|
|
// Request contains only non-secret operator inputs.
|
|
type Request struct {
|
|
StatePath string
|
|
Version string
|
|
Source Source
|
|
Image string
|
|
Confirm bool
|
|
Drain bool
|
|
}
|
|
|
|
// Result summarizes the completed, failed, or recovered transaction without command output.
|
|
type Result struct {
|
|
Phase Phase
|
|
StatePath string
|
|
}
|
|
|
|
type lifecycleHooks struct {
|
|
writeState func(string, State) error
|
|
removeFile func(string) error
|
|
sleep func(time.Duration)
|
|
}
|
|
|
|
var defaultLifecycleHooks = lifecycleHooks{
|
|
writeState: writeState,
|
|
removeFile: durableRemove,
|
|
sleep: time.Sleep,
|
|
}
|
|
|
|
// Update performs a recoverable core-only Pi update using the default Compose command layout.
|
|
func Update(ctx context.Context, runner Runner, request Request) (result Result, retErr error) {
|
|
return updateWithHooks(ctx, runner, request, defaultLifecycleHooks)
|
|
}
|
|
|
|
func updateWithHooks(ctx context.Context, runner Runner, request Request, hooks lifecycleHooks) (result Result, retErr error) {
|
|
if request.StatePath == "" {
|
|
return Result{}, errors.New("update state path is required")
|
|
}
|
|
lock, err := acquireLock(request.StatePath)
|
|
if err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
defer lock.Release()
|
|
if !request.Confirm {
|
|
return Result{StatePath: request.StatePath}, ErrConfirmationRequired
|
|
}
|
|
if !versionPattern.MatchString(request.Version) {
|
|
return Result{StatePath: request.StatePath}, fmt.Errorf("%w: Pi version must be an explicit pinned version", ErrInvalidRequest)
|
|
}
|
|
if request.Source == "" {
|
|
return Result{StatePath: request.StatePath}, fmt.Errorf("%w: Pi update requires an explicit source: build or pull", ErrInvalidRequest)
|
|
}
|
|
if request.Source != BuildSource && request.Source != PullSource {
|
|
return Result{StatePath: request.StatePath}, fmt.Errorf("%w: Pi update source must be build or pull", ErrInvalidRequest)
|
|
}
|
|
if request.Source == PullSource {
|
|
canonical, err := canonicalDigestReference(request.Image)
|
|
if err != nil {
|
|
return Result{StatePath: request.StatePath}, fmt.Errorf("%w: %v", ErrInvalidRequest, err)
|
|
}
|
|
request.Image = canonical
|
|
}
|
|
if old, err := readState(request.StatePath); err == nil && stateNeedsRecovery(old) {
|
|
return Result{StatePath: request.StatePath}, ErrInterruptedUpdate
|
|
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
|
|
return Result{StatePath: request.StatePath}, err
|
|
} else if err == nil && !old.MutationStarted {
|
|
if cleanupErr := hooks.removeFile(lifecycleOverridePath(request.StatePath, old.Transaction)); cleanupErr != nil {
|
|
return Result{StatePath: request.StatePath}, errors.New("safe prior preparation state could not be cleaned up")
|
|
}
|
|
}
|
|
if err := setMaintenance(ctx, runner, true); err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
clearMaintenance := true
|
|
defer func() {
|
|
if !clearMaintenance {
|
|
return
|
|
}
|
|
if clearErr := setMaintenance(context.Background(), runner, false); clearErr != nil {
|
|
result = Result{Phase: PhaseFailed, StatePath: request.StatePath}
|
|
retErr = errors.Join(retErr, fmt.Errorf("maintenance admission gate could not be cleared: %w", clearErr))
|
|
}
|
|
}()
|
|
|
|
if err := waitForInactiveSessions(ctx, runner, request.Drain, hooks.sleep); err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
if err := Doctor(ctx, runner); err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
|
|
currentVersion, err := Status(ctx, runner)
|
|
if err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
if currentVersion == request.Version {
|
|
return Result{Phase: PhaseNoop, StatePath: request.StatePath}, nil
|
|
}
|
|
configured, err := renderedCore(ctx, runner)
|
|
if err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
previous, err := runningImage(ctx, runner, configured.Reference)
|
|
if err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
previous.ConfigurationSHA = configured.ConfigurationSHA
|
|
transaction := lifecycleTransaction(request.StatePath)
|
|
previous.Reference = lifecycleImageTag(transaction, "previous")
|
|
candidateReference := lifecycleImageTag(transaction, "candidate")
|
|
if err := tagImage(ctx, runner, previous.ID, previous.Reference, "previous Pi image pin"); err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
state := State{
|
|
Transaction: transaction,
|
|
Phase: PhasePreflight,
|
|
Target: Target{Version: request.Version, Source: sourceValue(request)},
|
|
Previous: previous,
|
|
Candidate: Image{Reference: candidateReference},
|
|
}
|
|
if err := hooks.writeState(request.StatePath, state); err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
overridePath := lifecycleOverridePath(request.StatePath, transaction)
|
|
if err := writeLifecycleOverride(overridePath, candidateReference); err != nil {
|
|
return Result{StatePath: request.StatePath}, err
|
|
}
|
|
lifecycle := composeOverrideRunner{Runner: runner, path: overridePath}
|
|
|
|
state.Phase = PhaseBuilding
|
|
if err := hooks.writeState(request.StatePath, state); err != nil {
|
|
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
if err := prepareCandidate(ctx, lifecycle, request, candidateReference); err != nil {
|
|
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
running, err := activeSessions(ctx, runner)
|
|
if err != nil {
|
|
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
if running {
|
|
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, ErrActiveSessions, hooks)
|
|
return result, retErr
|
|
}
|
|
state.MutationStarted = true
|
|
if err := hooks.writeState(request.StatePath, state); err != nil {
|
|
state.MutationStarted = false
|
|
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
clearMaintenance = false
|
|
if err := recreateCore(ctx, lifecycle); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
if err := ensureMaintenance(ctx, lifecycle); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
state.Phase = PhaseRecreated
|
|
state.Candidate, err = runningImage(ctx, lifecycle, candidateReference)
|
|
if err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
if err := hooks.writeState(request.StatePath, state); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
if err := verifyCandidate(ctx, lifecycle, request.Version, previous); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
state.Phase, state.Error = PhasePromoting, ""
|
|
if err := hooks.writeState(request.StatePath, state); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
if err := promoteLifecycleOverride(overridePath, currentImageOverridePath(request.StatePath), candidateReference); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
state.Phase = PhaseVerified
|
|
if err := hooks.writeState(request.StatePath, state); err != nil {
|
|
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
|
|
return result, retErr
|
|
}
|
|
clearMaintenance = true
|
|
return Result{Phase: PhaseVerified, StatePath: request.StatePath}, nil
|
|
}
|
|
|
|
func stateNeedsRecovery(state State) bool {
|
|
switch state.Phase {
|
|
case PhaseVerified, PhaseRolledBack, PhaseNoop:
|
|
return false
|
|
case PhaseFailed:
|
|
return state.MutationStarted
|
|
default:
|
|
return state.MutationStarted
|
|
}
|
|
}
|
|
|
|
func failPreparation(statePath, overridePath string, state State, cause error, hooks lifecycleHooks) (Result, error, bool) {
|
|
state.Phase = PhaseFailed
|
|
state.MutationStarted = false
|
|
state.Error = "candidate preparation failed before core mutation"
|
|
writeErr := hooks.writeState(statePath, state)
|
|
removeErr := hooks.removeFile(overridePath)
|
|
message := "candidate preparation failed before core mutation"
|
|
if writeErr != nil || removeErr != nil {
|
|
message += "; safe preparation cleanup was incomplete"
|
|
}
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("%s: %w", message, cause), true
|
|
}
|
|
|
|
// Rollback restores the image recorded in durable update state. It is safe for interrupted runs.
|
|
func Rollback(ctx context.Context, runner Runner, statePath string, confirm bool) (result Result, retErr error) {
|
|
return rollbackWithHooks(ctx, runner, statePath, confirm, defaultLifecycleHooks)
|
|
}
|
|
|
|
func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, confirm bool, hooks lifecycleHooks) (result Result, retErr error) {
|
|
lock, err := acquireLock(statePath)
|
|
if err != nil {
|
|
return Result{StatePath: statePath}, err
|
|
}
|
|
defer lock.Release()
|
|
if !confirm {
|
|
return Result{StatePath: statePath}, ErrConfirmationRequired
|
|
}
|
|
maintenanceErr := ensureMaintenance(ctx, runner)
|
|
clearMaintenance := maintenanceErr == nil
|
|
defer func() {
|
|
if !clearMaintenance {
|
|
return
|
|
}
|
|
if clearErr := setMaintenance(context.Background(), runner, false); clearErr != nil {
|
|
result = Result{Phase: PhaseFailed, StatePath: statePath}
|
|
retErr = errors.Join(retErr, fmt.Errorf("maintenance admission gate could not be cleared: %w", clearErr))
|
|
}
|
|
}()
|
|
if maintenanceErr == nil {
|
|
if active, err := activeSessions(ctx, runner); err != nil {
|
|
return Result{StatePath: statePath}, err
|
|
} else if active {
|
|
return Result{StatePath: statePath}, ErrActiveSessions
|
|
}
|
|
}
|
|
state, err := readState(statePath)
|
|
if err != nil {
|
|
if maintenanceErr == nil {
|
|
clearMaintenance = false
|
|
}
|
|
return Result{StatePath: statePath}, err
|
|
}
|
|
overridePath := lifecycleOverridePath(statePath, state.Transaction)
|
|
if err := writeLifecycleOverride(overridePath, state.Previous.Reference); err != nil {
|
|
clearMaintenance = false
|
|
return Result{StatePath: statePath}, err
|
|
}
|
|
lifecycle := composeOverrideRunner{Runner: runner, path: overridePath}
|
|
if maintenanceErr != nil {
|
|
stopped, stopErr := coreIsStopped(ctx, runner)
|
|
if stopErr != nil || !stopped {
|
|
return Result{StatePath: statePath}, maintenanceErr
|
|
}
|
|
if err := persistMaintenanceWithoutLiveCore(ctx, lifecycle); err != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, err
|
|
}
|
|
}
|
|
clearMaintenance = false
|
|
if err := restore(ctx, lifecycle, state.Previous); err != nil {
|
|
state.Phase, state.Error = PhaseFailed, "rollback failed"
|
|
if writeErr := hooks.writeState(statePath, state); writeErr != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("rollback failed and recovery state could not be persisted")
|
|
}
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, err
|
|
}
|
|
if active, err := activeSessions(ctx, runner); err != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, err
|
|
} else if active {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, ErrActiveSessions
|
|
}
|
|
if err := promoteLifecycleOverride(overridePath, currentImageOverridePath(statePath), state.Previous.Reference); err != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("rollback restored the core but durable current-image promotion failed: %w", err)
|
|
}
|
|
state.Phase, state.Error = PhaseRolledBack, ""
|
|
if err := hooks.writeState(statePath, state); err != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("rollback restored the core but recovery state could not be persisted")
|
|
}
|
|
clearMaintenance = true
|
|
return Result{Phase: PhaseRolledBack, StatePath: statePath}, nil
|
|
}
|
|
|
|
func compensate(ctx context.Context, runner Runner, statePath, overridePath string, state State, cause error, hooks lifecycleHooks) (Result, error, bool) {
|
|
if err := writeLifecycleOverride(overridePath, state.Previous.Reference); err != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("update failed and rollback override could not be prepared: recovery required"), false
|
|
}
|
|
lifecycle := composeOverrideRunner{Runner: runner, path: overridePath}
|
|
if err := ensureMaintenance(context.Background(), runner); err != nil {
|
|
stopped, stopErr := coreIsStopped(context.Background(), runner)
|
|
if stopErr != nil || !stopped {
|
|
state.Phase, state.Error = PhaseFailed, "maintenance recovery failed"
|
|
_ = hooks.writeState(statePath, state)
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("update failed and maintenance could not be reactivated: recovery required"), false
|
|
}
|
|
if markerErr := persistMaintenanceWithoutLiveCore(context.Background(), lifecycle); markerErr != nil {
|
|
state.Phase, state.Error = PhaseFailed, "maintenance recovery failed"
|
|
_ = hooks.writeState(statePath, state)
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("update failed and durable maintenance could not be established: recovery required"), false
|
|
}
|
|
} else if active, err := activeSessions(context.Background(), runner); err != nil || active {
|
|
state.Phase, state.Error = PhaseFailed, "rollback inventory failed"
|
|
_ = hooks.writeState(statePath, state)
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("update failed and rollback inventory is not quiescent: recovery required"), false
|
|
}
|
|
if restoreErr := restore(ctx, lifecycle, state.Previous); restoreErr != nil {
|
|
state.Phase, state.Error = PhaseFailed, "candidate verification and automatic rollback failed"
|
|
if writeErr := hooks.writeState(statePath, state); writeErr != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("update failed and rollback proof failed; recovery state could not be persisted"), false
|
|
}
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("update failed; automatic rollback also failed: recovery required"), false
|
|
}
|
|
if active, err := activeSessions(context.Background(), runner); err != nil || active {
|
|
state.Phase, state.Error = PhaseFailed, "restored rollback inventory failed"
|
|
_ = hooks.writeState(statePath, state)
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("previous core image was restored but rollback inventory is not quiescent: recovery required"), false
|
|
}
|
|
if err := promoteLifecycleOverride(overridePath, currentImageOverridePath(statePath), state.Previous.Reference); err != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("previous core image was restored but durable selector promotion failed: %w", err), false
|
|
}
|
|
state.Phase, state.Error = PhaseRolledBack, ""
|
|
if writeErr := hooks.writeState(statePath, state); writeErr != nil {
|
|
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("previous core image was restored but recovery state write failed: recovery required"), false
|
|
}
|
|
return Result{Phase: PhaseRolledBack, StatePath: statePath}, fmt.Errorf("update failed; previous core image was restored: %w", cause), true
|
|
}
|
|
|
|
func coreIsStopped(ctx context.Context, runner Runner) (bool, error) {
|
|
result, err := runCompose(ctx, runner, "ps", "--status", "running", "-q", "core")
|
|
if err != nil {
|
|
return false, commandError("core running-state check", result, err)
|
|
}
|
|
return strings.TrimSpace(result.Stdout) == "", nil
|
|
}
|
|
|
|
const maintenanceMarkerScript = `
|
|
const fs = require("node:fs");
|
|
const path = require("node:path");
|
|
const marker = process.env.THT_MAINTENANCE_FILE;
|
|
if (!marker) throw new Error("THT_MAINTENANCE_FILE is required");
|
|
const directory = path.dirname(marker);
|
|
fs.mkdirSync(directory, { recursive: true });
|
|
const temporary = marker + ".rollback-" + process.pid + "-" + Date.now();
|
|
let file;
|
|
try {
|
|
file = fs.openSync(temporary, "wx", 0o600);
|
|
fs.writeFileSync(file, "{\"version\":1,\"active\":true}\n", "utf8");
|
|
fs.fsyncSync(file);
|
|
fs.closeSync(file);
|
|
file = undefined;
|
|
fs.renameSync(temporary, marker);
|
|
const directoryFile = fs.openSync(directory, "r");
|
|
try { fs.fsyncSync(directoryFile); } finally { fs.closeSync(directoryFile); }
|
|
} catch (error) {
|
|
if (file !== undefined) try { fs.closeSync(file); } catch {}
|
|
try { fs.unlinkSync(temporary); } catch {}
|
|
throw error;
|
|
}
|
|
`
|
|
|
|
func persistMaintenanceWithoutLiveCore(ctx context.Context, runner Runner) error {
|
|
result, err := runCompose(ctx, runner,
|
|
"run", "--rm", "--no-deps", "--entrypoint", "node", "core", "-e", maintenanceMarkerScript,
|
|
)
|
|
if err != nil {
|
|
return commandError("durable maintenance recovery", result, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func sourceValue(request Request) string {
|
|
if request.Source == PullSource {
|
|
return request.Image
|
|
}
|
|
return string(BuildSource)
|
|
}
|
|
|
|
func canonicalDigestReference(value string) (string, error) {
|
|
if strings.Contains(value, "://") || strings.ContainsAny(value, "?#") || strings.Contains(value, "@") && strings.Contains(strings.Split(value, "@")[0], ":") && strings.Contains(strings.Split(value, "@")[0], "//") {
|
|
return "", errors.New("pulled Pi image must be a credential-free canonical sha256 digest reference")
|
|
}
|
|
parsed, err := reference.ParseAnyReference(value)
|
|
if err != nil {
|
|
return "", errors.New("pulled Pi image must be a valid canonical sha256 digest reference")
|
|
}
|
|
canonical, ok := parsed.(reference.Canonical)
|
|
if !ok || canonical.Digest().Algorithm().String() != "sha256" || len(canonical.Digest().Encoded()) != 64 {
|
|
return "", errors.New("pulled Pi image must use an immutable sha256 digest")
|
|
}
|
|
return reference.FamiliarString(canonical), nil
|
|
}
|
|
|
|
func setMaintenance(ctx context.Context, runner Runner, enabled bool) error {
|
|
path := "deactivate"
|
|
if enabled {
|
|
path = "activate"
|
|
}
|
|
args := []string{"exec", "-T", "core", "curl", "-fsS", "-X", "POST", "http://127.0.0.1:8787/internal/maintenance/" + path}
|
|
result, err := runCompose(ctx, runner, args...)
|
|
status, valid := parseMaintenanceStatus(result.Stdout)
|
|
if err == nil && valid && status.Active == enabled && status.Admissions == 0 && !status.RecoveryRequired {
|
|
return nil
|
|
}
|
|
// Status identifies the safest immediate state after an ambiguous response. It cannot
|
|
// acknowledge durability for an operation whose command returned an error.
|
|
observed, statusErr := MaintenanceStatus(ctx, runner)
|
|
if valid && status.RecoveryRequired {
|
|
return recoveryRequired("maintenance durability was explicitly not acknowledged", err)
|
|
}
|
|
if statusErr == nil && observed.RecoveryRequired {
|
|
return recoveryRequired("maintenance durability was explicitly not acknowledged", err)
|
|
}
|
|
if err != nil {
|
|
if statusErr == nil && observed.Active == enabled && observed.Admissions == 0 {
|
|
return recoveryRequired("maintenance durability was not acknowledged after a failed command", err)
|
|
}
|
|
return commandError("maintenance admission gate", result, err)
|
|
}
|
|
if statusErr == nil && observed.Active == enabled && observed.Admissions == 0 {
|
|
return nil
|
|
}
|
|
return errors.New("maintenance admission gate did not acknowledge a quiescent state")
|
|
}
|
|
|
|
type MaintenanceState struct {
|
|
Active bool `json:"active"`
|
|
Admissions int `json:"admissions"`
|
|
RecoveryRequired bool `json:"recoveryRequired"`
|
|
}
|
|
|
|
func parseMaintenanceStatus(value string) (MaintenanceState, bool) {
|
|
var status MaintenanceState
|
|
err := json.Unmarshal([]byte(value), &status)
|
|
return status, err == nil && status.Admissions >= 0
|
|
}
|
|
|
|
func MaintenanceStatus(ctx context.Context, runner Runner) (MaintenanceState, error) {
|
|
result, err := runCompose(ctx, runner, "exec", "-T", "core", "curl", "-fsS", "http://127.0.0.1:8787/internal/maintenance/status")
|
|
if err != nil {
|
|
return MaintenanceState{}, commandError("maintenance status check", result, err)
|
|
}
|
|
status, valid := parseMaintenanceStatus(result.Stdout)
|
|
if !valid {
|
|
return MaintenanceState{}, errors.New("maintenance status check returned invalid data")
|
|
}
|
|
return status, nil
|
|
}
|
|
|
|
func ensureMaintenance(ctx context.Context, runner Runner) error {
|
|
status, err := MaintenanceStatus(ctx, runner)
|
|
if err == nil && status.Active && status.Admissions == 0 && !status.RecoveryRequired {
|
|
return nil
|
|
}
|
|
if err == nil && status.RecoveryRequired {
|
|
return recoveryRequired("maintenance durability was explicitly not acknowledged", nil)
|
|
}
|
|
return setMaintenance(ctx, runner, true)
|
|
}
|
|
|
|
func activeSessions(ctx context.Context, runner Runner) (bool, error) {
|
|
scope := "all"
|
|
if scoped, ok := runner.(interface{ SessionInventoryScope() string }); ok {
|
|
if requested := scoped.SessionInventoryScope(); requested == "mine" || requested == "all" {
|
|
scope = requested
|
|
}
|
|
}
|
|
args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...)
|
|
args = append(args, "http://127.0.0.1:8787/sessions?scope="+scope)
|
|
result, err := runCompose(ctx, runner, args...)
|
|
if err != nil {
|
|
return false, commandError("active-session check", result, err)
|
|
}
|
|
var payload []struct {
|
|
Status string `json:"status"`
|
|
Archived bool `json:"archived"`
|
|
}
|
|
if err := json.Unmarshal([]byte(result.Stdout), &payload); err != nil {
|
|
return false, errors.New("active-session check returned invalid session data")
|
|
}
|
|
for _, session := range payload {
|
|
if !session.Archived && session.Status != "finalized" && session.Status != "closed" {
|
|
return true, nil
|
|
}
|
|
}
|
|
return false, nil
|
|
}
|
|
|
|
func waitForInactiveSessions(ctx context.Context, runner Runner, drain bool, sleep func(time.Duration)) error {
|
|
running, err := activeSessions(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !running {
|
|
return nil
|
|
}
|
|
if !drain {
|
|
return ErrActiveSessions
|
|
}
|
|
for attempts := 0; attempts < 30; attempts++ {
|
|
running, err = activeSessions(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !running {
|
|
return nil
|
|
}
|
|
sleep(time.Second)
|
|
}
|
|
return ErrActiveSessions
|
|
}
|
|
|
|
func runningImage(ctx context.Context, runner Runner, reference string) (Image, error) {
|
|
container, err := runCompose(ctx, runner, "ps", "-q", "core")
|
|
if err != nil || strings.TrimSpace(container.Stdout) == "" {
|
|
return Image{}, commandError("running core image check", container, err)
|
|
}
|
|
id := strings.TrimSpace(container.Stdout)
|
|
image, err := runner.Run(ctx, []string{"inspect", "--format", "{{.Image}}", id}, nil)
|
|
if err != nil || strings.TrimSpace(image.Stdout) == "" {
|
|
return Image{}, commandError("running core image check", image, err)
|
|
}
|
|
mounts, err := runner.Run(ctx, []string{"inspect", "--format", "{{json .Mounts}}", id}, nil)
|
|
if err != nil {
|
|
return Image{}, commandError("core volume check", mounts, err)
|
|
}
|
|
var raw []struct {
|
|
Type string `json:"Type"`
|
|
Name string `json:"Name"`
|
|
Source string `json:"Source"`
|
|
Destination string `json:"Destination"`
|
|
RW bool `json:"RW"`
|
|
Mode string `json:"Mode"`
|
|
Propagation string `json:"Propagation"`
|
|
Driver string `json:"Driver"`
|
|
}
|
|
if err := json.Unmarshal([]byte(mounts.Stdout), &raw); err != nil {
|
|
return Image{}, errors.New("core returned invalid persistence mount data")
|
|
}
|
|
if len(raw) == 0 {
|
|
return Image{}, errors.New("core has no persistence mounts to preserve")
|
|
}
|
|
contract := make([]Mount, 0, len(raw))
|
|
for _, mount := range raw {
|
|
if mount.Type == "" || mount.Source == "" || mount.Destination == "" {
|
|
return Image{}, errors.New("core returned incomplete persistence mount data")
|
|
}
|
|
contract = append(contract, Mount{Type: mount.Type, Name: mount.Name, SourceSHA256: mountSourceHash(mount.Source), SourceAliases: mountSourceAliases(mount.Type, mount.Source, runtime.GOOS), Destination: mount.Destination, RW: mount.RW, Options: strings.Join([]string{mount.Mode, mount.Propagation, mount.Driver}, "\x00")})
|
|
}
|
|
return Image{ID: strings.TrimSpace(image.Stdout), Reference: reference, Mounts: contract, MountFingerprint: mountFingerprint(contract)}, nil
|
|
}
|
|
|
|
func prepareCandidate(ctx context.Context, runner Runner, request Request, candidateReference string) error {
|
|
if request.Source == BuildSource {
|
|
result, err := runCompose(ctx, runner, "build", "--pull", "--build-arg", "PI_VERSION="+request.Version, "core")
|
|
if err != nil {
|
|
return commandError("Pi image build", result, err)
|
|
}
|
|
return nil
|
|
}
|
|
pull, err := runner.Run(ctx, []string{"pull", request.Image}, nil)
|
|
if err != nil {
|
|
return commandError("Pi image pull", pull, err)
|
|
}
|
|
return tagImage(ctx, runner, request.Image, candidateReference, "Pi image tag")
|
|
}
|
|
|
|
func recreateCore(ctx context.Context, runner Runner) error {
|
|
result, err := runCompose(ctx, runner, "up", "--detach", "--wait", "--wait-timeout", "45", "--no-deps", "--force-recreate", "core")
|
|
if err != nil {
|
|
return commandError("core recreation", result, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func verifyCandidate(ctx context.Context, runner Runner, wanted string, previous Image) error {
|
|
health, err := runCompose(ctx, runner, "exec", "-T", "core", "curl", "-fsS", "http://127.0.0.1:8787/health")
|
|
if err != nil {
|
|
return commandError("core health check", health, err)
|
|
}
|
|
if err := verifyCandidateVersionIdentity(ctx, runner, wanted); err != nil {
|
|
return err
|
|
}
|
|
if err := Test(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
configured, err := renderedCore(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if configured.ConfigurationSHA != previous.ConfigurationSHA {
|
|
return errors.New("external endpoint configuration changed during Pi update")
|
|
}
|
|
after, err := runningImage(ctx, runner, configured.Reference)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !sameMounts(previous.Mounts, after.Mounts) {
|
|
return errors.New("core persistence mount contract changed during Pi update")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func verifyCandidateVersionIdentity(ctx context.Context, runner Runner, wanted string) error {
|
|
executable, err := Status(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
environment, label, err := expectedVersions(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if executable != wanted || environment != wanted || label != wanted {
|
|
return errors.New("candidate Pi executable, PI_VERSION, and image label do not all match the requested pinned version")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func restore(ctx context.Context, runner Runner, previous Image) error {
|
|
if err := tagImage(ctx, runner, previous.ID, previous.Reference, "rollback image restore"); err != nil {
|
|
return err
|
|
}
|
|
if err := recreateCore(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
if err := ensureMaintenance(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
configured, err := renderedCore(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
after, err := runningImage(ctx, runner, configured.Reference)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if after.ID != previous.ID {
|
|
return errors.New("rollback core image does not match recorded previous image")
|
|
}
|
|
if configured.ConfigurationSHA != previous.ConfigurationSHA {
|
|
return errors.New("external endpoint configuration drift prevents rollback proof")
|
|
}
|
|
if !sameMounts(previous.Mounts, after.Mounts) {
|
|
return errors.New("core persistence mount contract changed during rollback")
|
|
}
|
|
if err := Doctor(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
if err := Test(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func tagImage(ctx context.Context, runner Runner, source, target, label string) error {
|
|
result, err := runner.Run(ctx, []string{"image", "tag", source, target}, nil)
|
|
if err != nil {
|
|
return commandError(label, result, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type composeOverrideRunner struct {
|
|
Runner
|
|
path string
|
|
}
|
|
|
|
func (r composeOverrideRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) {
|
|
if len(args) > 0 && args[0] == "compose" {
|
|
withOverride := append([]string{"compose", "-f", r.path}, args[1:]...)
|
|
return r.Runner.Run(ctx, withOverride, stdin)
|
|
}
|
|
return r.Runner.Run(ctx, args, stdin)
|
|
}
|
|
|
|
func lifecycleTransaction(statePath string) string {
|
|
value := fmt.Sprintf("%s\x00%d\x00%d", filepath.Clean(statePath), os.Getpid(), time.Now().UnixNano())
|
|
sum := sha256.Sum256([]byte(value))
|
|
return fmt.Sprintf("%x", sum[:8])
|
|
}
|
|
|
|
func lifecycleImageTag(transaction, role string) string {
|
|
return "thothii-core:thothctl-" + transaction + "-" + role
|
|
}
|
|
|
|
func lifecycleOverridePath(statePath, transaction string) string {
|
|
if transaction == "" {
|
|
transaction = "recovery"
|
|
}
|
|
return filepath.Join(filepath.Dir(statePath), "pi-lifecycle-"+transaction+".yaml")
|
|
}
|
|
|
|
func currentImageOverridePath(statePath string) string {
|
|
return filepath.Join(filepath.Dir(statePath), "current-image.yaml")
|
|
}
|
|
|
|
func writeLifecycleOverride(path, image string) error {
|
|
quoted, err := json.Marshal(image)
|
|
if err != nil {
|
|
return errors.New("lifecycle image override could not be encoded")
|
|
}
|
|
contents := []byte(`services:
|
|
core:
|
|
image: ` + string(quoted) + `
|
|
workspace-maintenance:
|
|
image: ` + string(quoted) + `
|
|
`)
|
|
if err := writeFileDurably(path, ".pi-lifecycle-", contents); err != nil {
|
|
return errors.New("lifecycle image override could not be written durably")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func promoteLifecycleOverride(source, destination, expectedImage string) error {
|
|
return promoteLifecycleOverrideWith(source, destination, expectedImage, durableReplace)
|
|
}
|
|
|
|
func promoteLifecycleOverrideWith(
|
|
source, destination, expectedImage string,
|
|
replace func(string, string, string) error,
|
|
) error {
|
|
if err := replace(source, destination, filepath.Dir(destination)); err != nil {
|
|
selected, readErr := readLifecycleOverride(destination)
|
|
if readErr == nil && selected == expectedImage {
|
|
return recoveryRequired("lifecycle image override changed but durability was not acknowledged", err)
|
|
}
|
|
return recoveryRequired("lifecycle image override could not be promoted durably", err)
|
|
}
|
|
selected, err := readLifecycleOverride(destination)
|
|
if err != nil || selected != expectedImage {
|
|
return errors.New("promoted lifecycle image override could not be verified")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func readLifecycleOverride(path string) (string, error) {
|
|
contents, err := os.ReadFile(path)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
for _, line := range strings.Split(string(contents), "\n") {
|
|
line = strings.TrimSpace(line)
|
|
if !strings.HasPrefix(line, "image:") {
|
|
continue
|
|
}
|
|
encoded := strings.TrimSpace(strings.TrimPrefix(line, "image:"))
|
|
var image string
|
|
if json.Unmarshal([]byte(encoded), &image) != nil || image == "" || strings.ContainsAny(image, "\r\n") {
|
|
return "", errors.New("lifecycle image override is invalid")
|
|
}
|
|
return image, nil
|
|
}
|
|
return "", errors.New("lifecycle image override has no core image")
|
|
}
|
|
|
|
// RecoverMaintenance clears a stale durable gate only after the running core and terminal
|
|
// recovery metadata prove that no rollback is still required.
|
|
func RecoverMaintenance(ctx context.Context, runner Runner, statePath string, confirm bool) error {
|
|
if !confirm {
|
|
return ErrConfirmationRequired
|
|
}
|
|
lock, err := acquireLock(statePath)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer lock.Release()
|
|
return recoverMaintenanceLocked(ctx, runner, statePath)
|
|
}
|
|
|
|
func recoverMaintenanceLocked(ctx context.Context, runner Runner, statePath string) error {
|
|
state, stateErr := readState(statePath)
|
|
if stateErr == nil {
|
|
transactionOverride := lifecycleOverridePath(statePath, state.Transaction)
|
|
switch {
|
|
case state.Phase == PhasePromoting:
|
|
if err := recoverPromotion(ctx, runner, statePath, transactionOverride, &state); err != nil {
|
|
return err
|
|
}
|
|
case !state.MutationStarted && state.Phase != PhaseVerified && state.Phase != PhaseRolledBack && state.Phase != PhaseNoop:
|
|
state.Phase, state.Error = PhaseFailed, "candidate preparation interrupted before core mutation"
|
|
if err := writeState(statePath, state); err != nil {
|
|
return errors.New("maintenance recovery could not finalize safe preparation state")
|
|
}
|
|
if err := durableRemove(transactionOverride); err != nil {
|
|
return errors.New("maintenance recovery could not remove the safe preparation override")
|
|
}
|
|
case stateNeedsRecovery(state):
|
|
return ErrInterruptedUpdate
|
|
default:
|
|
if err := durableRemove(transactionOverride); err != nil {
|
|
return errors.New("maintenance recovery could not remove the lifecycle override")
|
|
}
|
|
}
|
|
} else if !errors.Is(stateErr, os.ErrNotExist) {
|
|
return stateErr
|
|
}
|
|
status, err := MaintenanceStatus(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !status.Active {
|
|
return nil
|
|
}
|
|
if err := Doctor(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
return setMaintenance(ctx, runner, false)
|
|
}
|
|
|
|
func recoverPromotion(ctx context.Context, runner Runner, statePath, transactionOverride string, state *State) error {
|
|
currentOverride := currentImageOverridePath(statePath)
|
|
selected, currentErr := readLifecycleOverride(currentOverride)
|
|
if currentErr != nil || selected != state.Candidate.Reference {
|
|
pending, pendingErr := readLifecycleOverride(transactionOverride)
|
|
if pendingErr != nil || pending != state.Candidate.Reference {
|
|
if currentErr == nil && selected == state.Previous.Reference {
|
|
if err := verifyRestoredCurrent(ctx, runner, state.Previous); err != nil {
|
|
return ErrInterruptedUpdate
|
|
}
|
|
state.Phase, state.Error = PhaseRolledBack, ""
|
|
return writeState(statePath, *state)
|
|
}
|
|
return ErrInterruptedUpdate
|
|
}
|
|
if err := promoteLifecycleOverride(transactionOverride, currentOverride, state.Candidate.Reference); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if err := verifyCandidate(ctx, runner, state.Target.Version, state.Previous); err != nil {
|
|
return err
|
|
}
|
|
state.Phase, state.Error = PhaseVerified, ""
|
|
return writeState(statePath, *state)
|
|
}
|
|
|
|
func verifyRestoredCurrent(ctx context.Context, runner Runner, previous Image) error {
|
|
configured, err := renderedCore(ctx, runner)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
after, err := runningImage(ctx, runner, configured.Reference)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if after.ID != previous.ID || configured.ConfigurationSHA != previous.ConfigurationSHA || !sameMounts(previous.Mounts, after.Mounts) {
|
|
return errors.New("running core does not match the durable previous-image selector")
|
|
}
|
|
if err := Doctor(ctx, runner); err != nil {
|
|
return err
|
|
}
|
|
return Test(ctx, runner)
|
|
}
|
|
|
|
func sameStrings(left, right []string) bool {
|
|
left, right = append([]string(nil), left...), append([]string(nil), right...)
|
|
sort.Strings(left)
|
|
sort.Strings(right)
|
|
return strings.Join(left, "\x00") == strings.Join(right, "\x00")
|
|
}
|
|
|
|
func sameMounts(left, right []Mount) bool {
|
|
if len(left) != len(right) {
|
|
return false
|
|
}
|
|
identityWithoutSource := func(m Mount) string {
|
|
return m.Type + "\x00" + m.Name + "\x00" + m.Destination + "\x00" + fmt.Sprint(m.RW) + "\x00" + m.Options
|
|
}
|
|
sourceMatches := func(a, b Mount) bool {
|
|
if a.SourceSHA256 == b.SourceSHA256 {
|
|
return true
|
|
}
|
|
for _, alias := range a.SourceAliases {
|
|
if alias == b.SourceSHA256 {
|
|
return true
|
|
}
|
|
}
|
|
for _, alias := range b.SourceAliases {
|
|
if alias == a.SourceSHA256 {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
matched := make([]bool, len(right))
|
|
for _, candidate := range left {
|
|
found := false
|
|
for index, observed := range right {
|
|
if matched[index] || identityWithoutSource(candidate) != identityWithoutSource(observed) || !sourceMatches(candidate, observed) {
|
|
continue
|
|
}
|
|
matched[index], found = true, true
|
|
break
|
|
}
|
|
if !found {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|