Files
ThothII/tools/tht/internal/pi/update.go
T

981 lines
37 KiB
Go

package pi
import (
"context"
"crypto/sha256"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"runtime"
"sort"
"strings"
"time"
"github.com/aritmolab/thothii/tools/tht/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")
)
// 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
RestartStatePath 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
Version 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)
}
// UpdateWithResolvedVersion resolves a requested version before Update obtains its lifecycle lock
// or contacts Docker. A registry failure therefore cannot drain sessions, write state, or mutate
// images and configuration.
func UpdateWithResolvedVersion(ctx context.Context, runner Runner, request Request, packageName string, registry RegistryClient) (Result, error) {
version, err := ResolveRequestedVersion(ctx, request.Version, packageName, registry)
if err != nil {
return Result{StatePath: request.StatePath}, err
}
request.Version = version
result, err := Update(ctx, runner, request)
result.Version = version
return result, err
}
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")
}
if request.RestartStatePath == "" {
request.RestartStatePath = pairedRestartStatePath(request.StatePath)
}
if err := validateRestartStatePaths(request.RestartStatePath, request.StatePath); err != nil {
return Result{StatePath: request.StatePath}, err
}
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 _, err := parseSemanticVersion(request.Version); err != nil {
return Result{StatePath: request.StatePath}, fmt.Errorf("%w: %v", ErrInvalidRequest, err)
}
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 err := prepareLifecycleMutation(request.RestartStatePath, hooks.removeFile); err != nil {
return Result{StatePath: request.StatePath}, err
}
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, restartStatePath string, confirm bool) (result Result, retErr error) {
return rollbackWithHooks(ctx, runner, statePath, restartStatePath, confirm, defaultLifecycleHooks)
}
func rollbackWithHooks(ctx context.Context, runner Runner, statePath, restartStatePath string, confirm bool, hooks lifecycleHooks) (result Result, retErr error) {
if restartStatePath == "" {
restartStatePath = pairedRestartStatePath(statePath)
}
if err := validateRestartStatePaths(restartStatePath, statePath); err != nil {
return Result{StatePath: statePath}, err
}
lock, err := acquireLock(statePath)
if err != nil {
return Result{StatePath: statePath}, err
}
defer lock.Release()
if !confirm {
return Result{StatePath: statePath}, ErrConfirmationRequired
}
if err := prepareLifecycleMutation(restartStatePath, hooks.removeFile); err != nil {
return Result{StatePath: statePath}, err
}
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:tht-" + 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
}