From 935bb1db0e9248b3d4b453a74f85c696563d4201 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 4 Aug 2026 23:14:47 +0200 Subject: [PATCH] fix: harden embedded Pi lifecycle recovery --- backend/package-lock.json | 12 +- backend/src/app.ts | 2 + backend/src/runtime/maintenance-gate.ts | 35 ++- backend/test/maintenance-gate.test.ts | 37 +++- backend/test/routes-sessions.test.ts | 16 +- docker/thothctl.Dockerfile | 2 +- docs/contracts/thothctl-pi.md | 43 ++-- scripts/test-thothctl-build-contract.sh | 10 +- tools/thothctl/go.mod | 14 +- tools/thothctl/go.sum | 26 ++- tools/thothctl/internal/pi/commands.go | 12 +- tools/thothctl/internal/pi/commands_test.go | 59 ++++- tools/thothctl/internal/pi/recovery_error.go | 24 +++ tools/thothctl/internal/pi/update.go | 167 ++++++++++++--- tools/thothctl/internal/pi/update_test.go | 214 ++++++++++++++++++- 15 files changed, 572 insertions(+), 101 deletions(-) create mode 100644 tools/thothctl/internal/pi/recovery_error.go diff --git a/backend/package-lock.json b/backend/package-lock.json index bb100dc7..dbc16ac6 100644 --- a/backend/package-lock.json +++ b/backend/package-lock.json @@ -1455,9 +1455,9 @@ } }, "node_modules/fast-uri": { - "version": "3.1.2", - "resolved": "https://registry.npmjs.org/fast-uri/-/fast-uri-3.1.2.tgz", - "integrity": "sha512-rVjf7ArG3LTk+FS6Yw81V1DLuZl1bRbNrev6Tmd/9RaroeeRRJhAt7jg/6YFxbvAQXUCavSoZhPPj6oOx+5KjQ==", + "version": "3.1.5", + "resolved": "https://registry.npmjs.org/fast-uri/-/fast-uri-3.1.5.tgz", + "integrity": "sha512-gHwA1O9LDIcKunMKhObS/HimwtehO1nPUECKAu5TpKgaO19fcWEl4bliWe1jWxVFvIXztJjjQ4L8XQ1EU9f7Jw==", "funding": [ { "type": "github", @@ -1529,9 +1529,9 @@ } }, "node_modules/find-my-way": { - "version": "9.6.0", - "resolved": "https://registry.npmjs.org/find-my-way/-/find-my-way-9.6.0.tgz", - "integrity": "sha512-Zf4Xve4RymLl7NgaavNebZ01joJ8MfVerOG43wy7SHLO+r+K0C6d/SE0BiR7AV5V1VOCFlOP7ecdo+I4qmiHrQ==", + "version": "9.7.0", + "resolved": "https://registry.npmjs.org/find-my-way/-/find-my-way-9.7.0.tgz", + "integrity": "sha512-f2JHn75x2JlwUwLenZypgczR7YWMb/uO9BvUXtus+JMgkbIkLADd38cI4EiV+OQqrGo1Zlq6V8wnqMJ8e62wUQ==", "license": "MIT", "dependencies": { "fast-deep-equal": "^3.1.3", diff --git a/backend/src/app.ts b/backend/src/app.ts index f543e240..085ef8f6 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -117,6 +117,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc } catch { return reply.code(500).send({ ...maintenanceBarrier.status(), + code: "maintenance_durability_failed", error: "maintenance activation durability was not acknowledged", }); } @@ -128,6 +129,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc } catch { return reply.code(500).send({ ...maintenanceBarrier.status(), + code: "maintenance_durability_failed", error: "maintenance deactivation durability was not acknowledged", }); } diff --git a/backend/src/runtime/maintenance-gate.ts b/backend/src/runtime/maintenance-gate.ts index 97ab888c..569c79da 100644 --- a/backend/src/runtime/maintenance-gate.ts +++ b/backend/src/runtime/maintenance-gate.ts @@ -11,6 +11,8 @@ import { import { dirname } from "node:path"; export interface MaintenanceDurability { + writeFile?(descriptor: number, contents: string): void; + syncFile?(descriptor: number): void; syncDirectory(directory: string): void; } @@ -19,6 +21,7 @@ export class MaintenanceBarrier { private active: boolean; private admissions = 0; private waiters: (() => void)[] = []; + private recoveryRequired = false; constructor( private readonly markerFile?: string, @@ -44,11 +47,14 @@ export class MaintenanceBarrier { let persistError: unknown; if (this.markerFile === undefined) { this.active = true; + this.recoveryRequired = false; } else { try { this.persistMarker(); + this.recoveryRequired = false; } catch (error) { persistError = error; + this.recoveryRequired = true; } finally { this.reconcileActive(); } @@ -63,17 +69,24 @@ export class MaintenanceBarrier { deactivate(): void { if (this.markerFile === undefined) { this.active = false; + this.recoveryRequired = false; return; } try { this.removeMarker(); + this.recoveryRequired = false; + } catch (error) { + try { this.persistMarker(); } catch { /* marker existence is reconciled below */ } + this.recoveryRequired = true; + throw error; } finally { this.reconcileActive(); } } - status(): { active: boolean; admissions: number } { + status(): { active: boolean; admissions: number; recoveryRequired?: true } { this.reconcileActive(); - return { active: this.active, admissions: this.admissions }; + const status = { active: this.active, admissions: this.admissions }; + return this.recoveryRequired ? { ...status, recoveryRequired: true } : status; } private reconcileActive(): void { @@ -86,24 +99,28 @@ export class MaintenanceBarrier { mkdirSync(directory, { recursive: true }); const temporary = `${this.markerFile}.tmp-${process.pid}-${Date.now()}`; const fd = openSync(temporary, "wx", 0o600); + let closed = false; try { - writeFileSync(fd, '{"version":1,"active":true}\n', "utf8"); - fsyncSync(fd); - } finally { + const contents = '{"version":1,"active":true}\n'; + if (this.durability.writeFile) this.durability.writeFile(fd, contents); + else writeFileSync(fd, contents, "utf8"); + if (this.durability.syncFile) this.durability.syncFile(fd); + else fsyncSync(fd); closeSync(fd); - } - try { + closed = true; renameSync(temporary, this.markerFile); this.durability.syncDirectory(directory); } catch (error) { + if (!closed) try { closeSync(fd); } catch { /* preserve the original failure */ } try { unlinkSync(temporary); } catch { /* already renamed or best-effort cleanup */ } throw error; } } private removeMarker(): void { - if (!this.markerFile || !existsSync(this.markerFile)) return; - unlinkSync(this.markerFile); + if (!this.markerFile) return; + if (existsSync(this.markerFile)) unlinkSync(this.markerFile); + else if (!this.recoveryRequired) return; this.durability.syncDirectory(dirname(this.markerFile)); } } diff --git a/backend/test/maintenance-gate.test.ts b/backend/test/maintenance-gate.test.ts index 9704d80c..b086a983 100644 --- a/backend/test/maintenance-gate.test.ts +++ b/backend/test/maintenance-gate.test.ts @@ -1,5 +1,5 @@ import { test, expect } from "vitest"; -import { existsSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { existsSync, mkdtempSync, readdirSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { MaintenanceBarrier } from "../src/runtime/maintenance-gate.js"; @@ -55,7 +55,7 @@ test("activation reconciles active state when directory fsync fails after marker }); await expect(gate.activate()).rejects.toThrow(/post-rename fsync failure/); expect(existsSync(marker)).toBe(true); - expect(gate.status()).toEqual({ active: true, admissions: 0 }); + expect(gate.status()).toEqual({ active: true, admissions: 0, recoveryRequired: true }); expect(gate.acquire()).toBeUndefined(); const recovered = new MaintenanceBarrier(marker); @@ -66,6 +66,33 @@ test("activation reconciles active state when directory fsync fails after marker } }); +test.each(["write", "fsync"] as const)( + "activation cleans its marker temp file after a pre-rename %s failure", + async (failure) => { + const directory = mkdtempSync(join(tmpdir(), "tht-maintenance-temp-cleanup-")); + const marker = join(directory, "maintenance.json"); + try { + const gate = new MaintenanceBarrier(marker, { + writeFile(descriptor, contents) { + if (failure === "write") throw new Error("injected marker write failure"); + writeFileSync(descriptor, contents, "utf8"); + }, + syncFile() { + if (failure === "fsync") throw new Error("injected marker fsync failure"); + }, + syncDirectory() {}, + }); + + await expect(gate.activate()).rejects.toThrow(`injected marker ${failure} failure`); + expect(existsSync(marker)).toBe(false); + expect(readdirSync(directory).filter((entry) => entry.includes(".tmp-"))).toEqual([]); + expect(gate.status()).toEqual({ active: false, admissions: 0, recoveryRequired: true }); + } finally { + rmSync(directory, { recursive: true, force: true }); + } + }, +); + test("deactivation reconciles inactive state when directory fsync fails after marker removal", async () => { const directory = mkdtempSync(join(tmpdir(), "tht-maintenance-deactivate-fsync-")); const marker = join(directory, "maintenance.json"); @@ -79,9 +106,9 @@ test("deactivation reconciles inactive state when directory fsync fails after ma await gate.activate(); failSync = true; expect(() => gate.deactivate()).toThrow(/post-remove fsync failure/); - expect(existsSync(marker)).toBe(false); - expect(gate.status()).toEqual({ active: false, admissions: 0 }); - expect(gate.acquire()).toBeTypeOf("function"); + expect(existsSync(marker)).toBe(true); + expect(gate.status()).toEqual({ active: true, admissions: 0, recoveryRequired: true }); + expect(gate.acquire()).toBeUndefined(); } finally { rmSync(directory, { recursive: true, force: true }); } diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index fb1f83d6..8b8e2e8b 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -170,15 +170,25 @@ test("maintenance endpoints report marker-derived state after post-rename and po const activated = await app.inject({ method: "POST", url: "/internal/maintenance/activate" }); expect(activated.statusCode).toBe(500); - expect(activated.json()).toMatchObject({ active: true, admissions: 0 }); + expect(activated.json()).toMatchObject({ + active: true, + admissions: 0, + recoveryRequired: true, + code: "maintenance_durability_failed", + }); failSync = false; expect((await app.inject({ method: "GET", url: "/internal/maintenance/status" })).json()) - .toEqual({ active: true, admissions: 0 }); + .toEqual({ active: true, admissions: 0, recoveryRequired: true }); failSync = true; const deactivated = await app.inject({ method: "POST", url: "/internal/maintenance/deactivate" }); expect(deactivated.statusCode).toBe(500); - expect(deactivated.json()).toMatchObject({ active: false, admissions: 0 }); + expect(deactivated.json()).toMatchObject({ + active: true, + admissions: 0, + recoveryRequired: true, + code: "maintenance_durability_failed", + }); } finally { rmSync(dir, { recursive: true, force: true }); } diff --git a/docker/thothctl.Dockerfile b/docker/thothctl.Dockerfile index 8af20f39..5fd51bb8 100644 --- a/docker/thothctl.Dockerfile +++ b/docker/thothctl.Dockerfile @@ -1,4 +1,4 @@ -FROM golang:1.24@sha256:d2d2bc1c84f7e60d7d2438a3836ae7d0c847f4888464e7ec9ba3a1339a1ee804 AS build +FROM golang:1.26.5-bookworm@sha256:1ecb7edf62a0408027bd5729dfd6b1b8766e578e8df93995b225dfd0944eb651 AS build WORKDIR /src/tools/thothctl COPY tools/thothctl/go.mod tools/thothctl/go.sum ./ diff --git a/docs/contracts/thothctl-pi.md b/docs/contracts/thothctl-pi.md index 98cb9a7b..a5306635 100644 --- a/docs/contracts/thothctl-pi.md +++ b/docs/contracts/thothctl-pi.md @@ -75,7 +75,10 @@ Before inventory, update activates the durable maintenance gate. Activation writ all leases. A recreated candidate reads that marker at startup and therefore starts gated. The loopback-only control endpoints cannot be reached through the frontend proxy and do not depend on the configured authentication principal mode. Lost activation/deactivation responses are resolved -by querying gate status. +by querying gate status only when the original result is unknown. An explicit file or directory +durability failure is never converted to success by matching readback: the control API reports +`maintenance_durability_failed`, keeps or restores the safest durable marker state, and requires +recovery. Open, unarchived sessions stop an update. After an operator has completed or otherwise drained their work, `--drain` makes the command poll the authenticated bare-array @@ -91,9 +94,10 @@ Rollback atomically promotes the previous selector. Terminal cleanup removes onl files and never deletes the durable selector. Only `core` is recreated, with `--no-deps --force-recreate`; `frontend` is not recreated and no -volume-replacement flags are used. Verification checks health, exact requested Pi version, the -provider/model/settings smoke, unchanged non-secret rendered configuration, and the complete -persistence-mount fingerprint. +volume-replacement flags are used. Verification checks health; exact requested Pi version at all +three declared boundaries (the candidate executable, `PI_VERSION` environment, and +`org.opencontainers.image.version` image label); the provider/model/settings smoke; unchanged +non-secret rendered configuration; and the complete persistence-mount fingerprint. ## Recovery, rollback, and maintenance cleanup @@ -114,6 +118,11 @@ Any post-candidate failure explicitly confirms or reactivates maintenance and re before compensation. Automatic rollback selects the transaction's previous image through the lifecycle override and clears maintenance only after the previous image, configuration, mounts, health, Pi smoke, and terminal recovery write are verified. Ambiguous compensation remains gated. +If the candidate core is stopped and cannot serve the maintenance endpoint, rollback proves that +state with Compose and writes the marker through a one-off previous-image `core` container sharing +the settings volume. It does not require the failed candidate, a host Node runtime, or the Docker +socket inside a container. The restored core is then recreated, verified, and rescanned before the +gate can open. For an interrupted transaction, first run: @@ -140,12 +149,22 @@ Missing confirmation, invalid arguments, active sessions, and an interrupted tra `2`. Docker and verification failures exit nonzero with concise, redacted guidance. Direct read-only/log commands preserve the original Docker child exit code. -## Go 1.24 dependency security boundary +## Go dependency security boundary -`golang.org/x/sys` is a direct dependency for the Windows durable-replace implementation. The -newest release compatible with the required Go 1.24 toolchain is pinned (`v0.41.0`). Releases -`v0.42.0` through `v0.44.0` require Go 1.25, so `v0.44.0` cannot be selected without changing the -product toolchain contract. Govulncheck reports GO-2026-5024 at module level for `v0.41.0`, but no -thothctl call trace reaches the vulnerable `windows.NewNTUnicodeString`; thothctl calls only -`UTF16PtrFromString` and `MoveFileEx` in that package. Upgrade to at least `v0.44.0` together with -the planned Go 1.25-or-newer toolchain migration. +The supported toolchain is Go `1.26.5`, released 2026-07-07, with module language version +`1.26.0`. The Docker builder is pinned by both patch tag and the multi-platform manifest-list +digest: + +```text +golang:1.26.5-bookworm@sha256:1ecb7edf62a0408027bd5729dfd6b1b8766e578e8df93995b225dfd0944eb651 +``` + +That manifest provides both `linux/amd64` and `linux/arm64/v8` builders. Go's official release +history is the authority for the patch level (`https://go.dev/doc/devel/release`); the Docker +Official Image is the authority for the builder (`https://hub.docker.com/_/golang`). +`golang.org/x/sys`, used by the Windows durable-replace implementation, is pinned to `v0.47.0`. +The directly used `github.com/sirupsen/logrus` is pinned to `v1.9.1`, which removes +GO-2025-4188 from the imported package set. The build contract verifies the exact toolchain, +dependencies, digest, and all five supported target builds (Windows amd64, Darwin amd64/arm64, and +Linux amd64/arm64). `go mod verify`, tests including the race detector, `go vet`, and +`govulncheck` are release gates. diff --git a/scripts/test-thothctl-build-contract.sh b/scripts/test-thothctl-build-contract.sh index cdd20775..7c785dc3 100755 --- a/scripts/test-thothctl-build-contract.sh +++ b/scripts/test-thothctl-build-contract.sh @@ -4,12 +4,18 @@ set -euo pipefail repository_root=$(cd "$(dirname "$0")/.." && pwd) dockerfile="$repository_root/docker/thothctl.Dockerfile" builder_image=$(awk '$1 == "FROM" && $3 == "AS" && $4 == "build" { print $2; exit }' "$dockerfile") +expected_builder='golang:1.26.5-bookworm@sha256:1ecb7edf62a0408027bd5729dfd6b1b8766e578e8df93995b225dfd0944eb651' -if [[ ! "$builder_image" =~ ^golang:1\.24@sha256:[0-9a-f]{64}$ ]]; then - echo "thothctl builder must use a readable golang:1.24 tag with an immutable digest" >&2 +if [[ "$builder_image" != "$expected_builder" ]]; then + echo "thothctl builder must pin golang:1.26.5-bookworm by the approved multi-platform digest" >&2 exit 1 fi +grep -qx 'go 1.26.0' "$repository_root/tools/thothctl/go.mod" +grep -qx 'toolchain go1.26.5' "$repository_root/tools/thothctl/go.mod" +grep -Eq '^[[:space:]]*github.com/sirupsen/logrus v1\.9\.1$' "$repository_root/tools/thothctl/go.mod" +grep -Eq '^[[:space:]]*golang.org/x/sys v0\.47\.0$' "$repository_root/tools/thothctl/go.mod" + manifest=$(docker buildx imagetools inspect "$builder_image") printf '%s\n' "$manifest" | grep -Eq 'Platform:[[:space:]]+linux/amd64' printf '%s\n' "$manifest" | grep -Eq 'Platform:[[:space:]]+linux/arm64' diff --git a/tools/thothctl/go.mod b/tools/thothctl/go.mod index ad542ecf..c1663424 100644 --- a/tools/thothctl/go.mod +++ b/tools/thothctl/go.mod @@ -1,8 +1,8 @@ module github.com/aritmolab/thothii/tools/thothctl -go 1.24.0 +go 1.26.0 -toolchain go1.24.13 +toolchain go1.26.5 require gopkg.in/yaml.v3 v3.0.1 @@ -10,8 +10,12 @@ require ( github.com/compose-spec/compose-go/v2 v2.14.0 github.com/distribution/reference v0.6.0 github.com/gofrs/flock v0.12.1 - github.com/sirupsen/logrus v1.9.0 - golang.org/x/sys v0.41.0 + github.com/sirupsen/logrus v1.9.1 + golang.org/x/sys v0.47.0 ) -require github.com/opencontainers/go-digest v1.0.0 // indirect +require ( + github.com/kr/text v0.2.0 // indirect + github.com/opencontainers/go-digest v1.0.0 // indirect + github.com/rogpeppe/go-internal v1.15.0 // indirect +) diff --git a/tools/thothctl/go.sum b/tools/thothctl/go.sum index 2810f68c..1b6db1ef 100644 --- a/tools/thothctl/go.sum +++ b/tools/thothctl/go.sum @@ -1,5 +1,6 @@ github.com/compose-spec/compose-go/v2 v2.14.0 h1:uaJeo5B3+OVlu+Rx2qLBcAdXPEUUzm5nQrRiGJafRAQ= github.com/compose-spec/compose-go/v2 v2.14.0/go.mod h1:ZU6zlcweCZKyiB7BVfCizQT9XmkEIMFE+PRZydVcsZg= +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -9,25 +10,28 @@ github.com/gofrs/flock v0.12.1 h1:MTLVXXHf8ekldpJk3AKicLij9MdwOWkZ+a/jHHZby9E= github.com/gofrs/flock v0.12.1/go.mod h1:9zxTsyu5xtJ9DK+1tFZyibEV7y3uwDxPPfbxeeHCoD0= github.com/google/go-cmp v0.5.9 h1:O2Tfq5qg4qc4AmwVlvv0oLiVAGB7enBSJ2x2DqQFi38= github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/sirupsen/logrus v1.9.0 h1:trlNQbNUG3OdDrDil03MCb1H2o9nJ1x4/5LYw7byDE0= -github.com/sirupsen/logrus v1.9.0/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= +github.com/rogpeppe/go-internal v1.15.0 h1:D0RCU5rMAp+SpgkiNdrjfJ+LX4J1M32V2NeCY7EJ6hc= +github.com/rogpeppe/go-internal v1.15.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs= +github.com/sirupsen/logrus v1.9.1 h1:Ou41VVR3nMWWmTiEUnj0OlsgOSCUFgsPAOl6jRIcVtQ= +github.com/sirupsen/logrus v1.9.1/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -github.com/stretchr/testify v1.8.4 h1:CcVxjf3Q8PM0mHUKJCdn+eZZtm5yQwehR5yeSVQQcUk= -github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= +github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg= +github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.5.0 h1:MUK/U/4lj1t1oPg0HfuXDN/Z1wv31ZJ/YcPiGccS4DU= -golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.22.0 h1:RI27ohtqKCnwULzJLqkv897zojh5/DwS/ENaMzUOaWI= -golang.org/x/sys v0.22.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= -golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/tools/thothctl/internal/pi/commands.go b/tools/thothctl/internal/pi/commands.go index b5057683..a4bbd926 100644 --- a/tools/thothctl/internal/pi/commands.go +++ b/tools/thothctl/internal/pi/commands.go @@ -76,7 +76,7 @@ func Configure(ctx context.Context, runner Runner, value Defaults) error { } restore := func(cause error) error { if restoreErr := restoreSettingsFile(context.Background(), runner, old); restoreErr != nil { - return fmt.Errorf("%w; previous Pi settings restoration could not be verified: recovery required", cause) + return fmt.Errorf("%w; previous Pi settings restoration could not be verified: %w", cause, restoreErr) } restoredEffective, restoreErr := readEffectiveSettings(context.Background(), runner) if restoreErr != nil || !bytes.Equal(restoredEffective, oldEffective) { @@ -155,12 +155,16 @@ func restoreSettingsFile(ctx context.Context, runner Runner, snapshot settingsFi args := []string{"compose", "exec", "-T", "core", "node", "/app/backend/dist/settings/settings-cli.js", "--restore"} result, restoreErr := runner.Run(ctx, args, bytes.NewReader(payload)) verified, verifyErr := captureSettingsFile(ctx, runner) + if restoreErr != nil { + cause := commandError("Pi installation settings restore", result, restoreErr) + if verifyErr == nil && verified == snapshot { + return recoveryRequired("previous Pi settings bytes were restored but durability was not acknowledged", cause) + } + return cause + } if verifyErr == nil && verified == snapshot { return nil } - if restoreErr != nil { - return commandError("Pi installation settings restore", result, restoreErr) - } return errors.New("Pi installation settings restore did not reproduce the exact prior file state") } diff --git a/tools/thothctl/internal/pi/commands_test.go b/tools/thothctl/internal/pi/commands_test.go index 8c4358b0..c5bd36bd 100644 --- a/tools/thothctl/internal/pi/commands_test.go +++ b/tools/thothctl/internal/pi/commands_test.go @@ -61,14 +61,56 @@ func TestConfigureRestoresAndVerifiesOldSettingsAfterEveryPostSnapshotFailure(t } } +func TestSettingsRestoreDoesNotMaskExplicitDurabilityFailureWithMatchingReadback(t *testing.T) { + old := Defaults{Provider: "old", Model: "old-model", Thinking: "low"} + raw, _ := json.Marshal(old) + fake := &configureRunner{ + failure: "restore-durability", + settings: Defaults{Provider: "new", Model: "new-model", Thinking: "high"}, + settingsExist: true, + settingsRaw: []byte(`{"provider":"new","model":"new-model","thinking":"high"}`), + } + snapshot := settingsFileSnapshot{Exists: true, RawBase64: base64.StdEncoding.EncodeToString(raw)} + + err := restoreSettingsFile(context.Background(), fake, snapshot) + + var recovery interface{ RecoveryRequired() bool } + if err == nil || !errors.As(err, &recovery) || !recovery.RecoveryRequired() { + t.Fatalf("restore error = %v; want typed recovery-required result", err) + } + if !fake.settingsExist || string(fake.settingsRaw) != string(raw) || fake.settings != old { + t.Fatalf("restored state = exists:%t raw:%q value:%#v; want exact old bytes", fake.settingsExist, fake.settingsRaw, fake.settings) + } +} + +func TestConfigurePreservesTypedRecoveryRequiredErrorFromSettingsRestore(t *testing.T) { + old := Defaults{Provider: "old", Model: "old-model", Thinking: "low"} + raw, _ := json.Marshal(old) + fake := &configureRunner{ + failure: "helper", + restoreDurabilityFailure: true, + settings: old, + settingsExist: true, + settingsRaw: raw, + } + + err := Configure(context.Background(), fake, Defaults{Provider: "new", Model: "new-model", Thinking: "high"}) + + var recovery interface{ RecoveryRequired() bool } + if err == nil || !errors.As(err, &recovery) || !recovery.RecoveryRequired() { + t.Fatalf("Configure() error = %v; want typed recovery-required result", err) + } +} + type configureRunner struct { - failure string - settings Defaults - settingsExist bool - settingsRaw []byte - settingsReads int - configReads int - writes int + failure string + restoreDurabilityFailure bool + settings Defaults + settingsExist bool + settingsRaw []byte + settingsReads int + configReads int + writes int } func (f *configureRunner) Run(_ context.Context, args []string, stdin io.Reader) (compose.Result, error) { @@ -103,6 +145,9 @@ func (f *configureRunner) Run(_ context.Context, args []string, stdin io.Reader) if payload.Exists { _ = json.Unmarshal(f.settingsRaw, &f.settings) } + if f.failure == "restore-durability" || f.restoreDurabilityFailure { + return compose.Result{ExitCode: 2}, errors.New("injected post-rename directory fsync failure") + } return compose.Result{}, nil case strings.Contains(call, "settings-cli.js"): if strings.Contains(call, "--provider new") { diff --git a/tools/thothctl/internal/pi/recovery_error.go b/tools/thothctl/internal/pi/recovery_error.go new file mode 100644 index 00000000..943d75ab --- /dev/null +++ b/tools/thothctl/internal/pi/recovery_error.go @@ -0,0 +1,24 @@ +package pi + +// RecoveryRequiredError marks a result whose immediate state may be safe but whose durability +// was explicitly not acknowledged. Callers must not report success or clear maintenance. +type RecoveryRequiredError struct { + Operation string + Cause error +} + +func (e *RecoveryRequiredError) Error() string { + return e.Operation + ": recovery required" +} + +func (e *RecoveryRequiredError) Unwrap() error { + return e.Cause +} + +func (e *RecoveryRequiredError) RecoveryRequired() bool { + return true +} + +func recoveryRequired(operation string, cause error) error { + return &RecoveryRequiredError{Operation: operation, Cause: cause} +} diff --git a/tools/thothctl/internal/pi/update.go b/tools/thothctl/internal/pi/update.go index 2fae5b6e..4ce191a7 100644 --- a/tools/thothctl/internal/pi/update.go +++ b/tools/thothctl/internal/pi/update.go @@ -289,10 +289,8 @@ func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, con if !confirm { return Result{StatePath: statePath}, ErrConfirmationRequired } - if err := setMaintenance(ctx, runner, true); err != nil { - return Result{StatePath: statePath}, err - } - clearMaintenance := true + maintenanceErr := ensureMaintenance(ctx, runner) + clearMaintenance := maintenanceErr == nil defer func() { if !clearMaintenance { return @@ -306,14 +304,18 @@ func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, con } } }() - if active, err := activeSessions(ctx, runner); err != nil { - return Result{StatePath: statePath}, err - } else if active { - return Result{StatePath: statePath}, ErrActiveSessions + 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 { - clearMaintenance = false + if maintenanceErr == nil { + clearMaintenance = false + } return Result{StatePath: statePath}, err } overridePath := lifecycleOverridePath(statePath, state.Transaction) @@ -321,8 +323,17 @@ func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, con clearMaintenance = false return Result{StatePath: statePath}, err } - clearMaintenance = false 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 { @@ -330,8 +341,13 @@ func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, con } 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}, errors.New("rollback restored the core but durable current-image promotion failed: recovery required") + 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 { @@ -342,20 +358,27 @@ func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, con } func compensate(ctx context.Context, runner Runner, statePath, overridePath string, state State, cause error, hooks lifecycleHooks) (Result, error, bool) { - if err := ensureMaintenance(context.Background(), runner); err != nil { - 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 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 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 { @@ -363,8 +386,13 @@ func compensate(ctx context.Context, runner Runner, statePath, overridePath stri } 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}, errors.New("previous core image was restored but durable selector promotion failed: recovery required"), false + 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 { @@ -373,6 +401,49 @@ func compensate(ctx context.Context, runner Runner, statePath, overridePath stri 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 @@ -403,15 +474,21 @@ func setMaintenance(ctx context.Context, runner Runner, enabled bool) error { 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 { + if err == nil && valid && status.Active == enabled && status.Admissions == 0 && !status.RecoveryRequired { return nil } // A core recreate or transport interruption may lose only the response. Resolve ambiguity by // reading the durable gate state before deciding that operator recovery is required. observed, statusErr := MaintenanceStatus(ctx, runner) if statusErr == nil && observed.Active == enabled && observed.Admissions == 0 { + if observed.RecoveryRequired { + return recoveryRequired("maintenance durability was explicitly not acknowledged", err) + } return nil } + if valid && status.RecoveryRequired { + return recoveryRequired("maintenance durability was explicitly not acknowledged", err) + } if err != nil { return commandError("maintenance admission gate", result, err) } @@ -419,8 +496,9 @@ func setMaintenance(ctx context.Context, runner Runner, enabled bool) error { } type MaintenanceState struct { - Active bool `json:"active"` - Admissions int `json:"admissions"` + Active bool `json:"active"` + Admissions int `json:"admissions"` + RecoveryRequired bool `json:"recoveryRequired"` } func parseMaintenanceStatus(value string) (MaintenanceState, bool) { @@ -443,9 +521,12 @@ func MaintenanceStatus(ctx context.Context, runner Runner) (MaintenanceState, er func ensureMaintenance(ctx context.Context, runner Runner) error { status, err := MaintenanceStatus(ctx, runner) - if err == nil && status.Active && status.Admissions == 0 { + 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) } @@ -539,13 +620,9 @@ func verifyCandidate(ctx context.Context, runner Runner, wanted string, previous if err != nil { return commandError("core health check", health, err) } - version, err := Status(ctx, runner) - if err != nil { + if err := verifyCandidateVersionIdentity(ctx, runner, wanted); err != nil { return err } - if version != wanted { - return errors.New("candidate Pi version does not match requested pinned version") - } if err := Test(ctx, runner); err != nil { return err } @@ -566,6 +643,21 @@ func verifyCandidate(ctx context.Context, runner Runner, wanted string, previous 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 @@ -657,12 +749,19 @@ func writeLifecycleOverride(path, image string) error { } func promoteLifecycleOverride(source, destination, expectedImage string) error { - if err := durableReplace(source, destination, filepath.Dir(destination)); err != nil { + 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 nil + return recoveryRequired("lifecycle image override changed but durability was not acknowledged", err) } - return errors.New("lifecycle image override could not be promoted durably") + return recoveryRequired("lifecycle image override could not be promoted durably", err) } selected, err := readLifecycleOverride(destination) if err != nil || selected != expectedImage { diff --git a/tools/thothctl/internal/pi/update_test.go b/tools/thothctl/internal/pi/update_test.go index 2c160a98..3f3391ce 100644 --- a/tools/thothctl/internal/pi/update_test.go +++ b/tools/thothctl/internal/pi/update_test.go @@ -102,6 +102,38 @@ func TestSuccessfulUpdateAndRollbackRemainSelectedOnFreshRecreate(t *testing.T) } } +func TestSelectorPromotionDoesNotMaskPostRenameDirectoryFsyncFailure(t *testing.T) { + directory := t.TempDir() + source := filepath.Join(directory, "candidate.yaml") + destination := filepath.Join(directory, "current-image.yaml") + if err := writeLifecycleOverride(source, "thothii-core:candidate"); err != nil { + t.Fatal(err) + } + if err := writeLifecycleOverride(destination, "thothii-core:old"); err != nil { + t.Fatal(err) + } + + err := promoteLifecycleOverrideWith( + source, + destination, + "thothii-core:candidate", + func(source, destination, _ string) error { + if err := os.Rename(source, destination); err != nil { + return err + } + return errors.New("injected post-rename directory fsync failure") + }, + ) + + var recovery interface{ RecoveryRequired() bool } + if err == nil || !errors.As(err, &recovery) || !recovery.RecoveryRequired() { + t.Fatalf("promotion error = %v; want typed recovery-required result", err) + } + if selected := readSelectorReference(t, destination); selected != "thothii-core:candidate" { + t.Fatalf("immediate selector = %q, want landed candidate bytes", selected) + } +} + func TestTwoInstallationsSharingAConfiguredTagUseDifferentLifecycleTags(t *testing.T) { first, second := newFakeRunner(), newFakeRunner() firstPath := filepath.Join(t.TempDir(), "one", "state.json") @@ -162,6 +194,21 @@ func TestMaintenanceLostResponsesAreResolvedByStatusAndEveryRecreateStartsGated( } } +func TestMaintenanceReconciliationDoesNotMaskExplicitDurabilityFailure(t *testing.T) { + fake := newFakeRunner() + fake.fail = "maintenance-activate-durability" + + err := setMaintenance(context.Background(), fake, true) + + var recovery interface{ RecoveryRequired() bool } + if err == nil || !errors.As(err, &recovery) || !recovery.RecoveryRequired() { + t.Fatalf("maintenance error = %v; want typed recovery-required result", err) + } + if !fake.maintenance { + t.Fatal("safe marker state was not retained after activation durability failure") + } +} + func TestCompensationReactivatesMaintenanceAndRescansBeforeRollback(t *testing.T) { fake := newFakeRunner() fake.fail = "version" @@ -179,6 +226,63 @@ func TestCompensationReactivatesMaintenanceAndRescansBeforeRollback(t *testing.T } } +func TestAutomaticRollbackSurvivesADeadCandidateCore(t *testing.T) { + fake := newFakeRunner() + fake.fail = "dead-candidate" + statePath := filepath.Join(t.TempDir(), ".thothctl", "update-state.json") + + result, err := Update(context.Background(), fake, Request{ + StatePath: statePath, + Version: "0.81.0", + Source: BuildSource, + Confirm: true, + }) + + if err == nil || result.Phase != PhaseRolledBack { + t.Fatalf("Update() = %+v, %v; want automatic rollback after dead candidate", result, err) + } + if fake.currentImage != "sha256:old" || !fake.coreRunning { + t.Fatalf("restored core = image:%q running:%t; want previous running image", fake.currentImage, fake.coreRunning) + } + if fake.execFailuresWhileStopped == 0 { + t.Fatal("fake did not exercise candidate exec failure") + } + assertCalled(t, fake.calls, "run --rm --no-deps --entrypoint node") + assertMaintenanceClearedAfterRestoredProof(t, fake) +} + +func TestManualRollbackSurvivesADeadCandidateCore(t *testing.T) { + fake := newFakeRunner() + statePath := filepath.Join(t.TempDir(), ".thothctl", "update-state.json") + previous := stateImageForTest(t, fake) + previous.Reference = "thothii-core:thothctl-dead-candidate-previous" + fake.tags[previous.Reference] = previous.ID + writeStateForTest(t, statePath, State{ + Transaction: "dead-candidate", + Phase: PhaseRecreated, + MutationStarted: true, + Previous: previous, + }) + fake.currentImage = "sha256:candidate" + fake.version = "0.81.0" + fake.coreRunning = false + fake.maintenance = false + + result, err := Rollback(context.Background(), fake, statePath, true) + + if err != nil || result.Phase != PhaseRolledBack { + t.Fatalf("Rollback() = %+v, %v; want restored previous core", result, err) + } + if fake.currentImage != "sha256:old" || !fake.coreRunning { + t.Fatalf("restored core = image:%q running:%t; want previous running image", fake.currentImage, fake.coreRunning) + } + if fake.execFailuresWhileStopped == 0 { + t.Fatal("fake did not exercise candidate exec failure") + } + assertCalled(t, fake.calls, "run --rm --no-deps --entrypoint node") + assertMaintenanceClearedAfterRestoredProof(t, fake) +} + func TestUpdatePullsOnlyDigestPinnedSource(t *testing.T) { fake := newFakeRunner() digest := "registry.example.invalid/thothii-core@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" @@ -234,6 +338,36 @@ func TestUpdateRollsBackAfterPostRecreateFailures(t *testing.T) { } } +func TestCandidateVerificationRejectsEveryDeclaredVersionBoundaryMismatch(t *testing.T) { + for _, boundary := range []string{"executable", "environment", "image-label"} { + t.Run(boundary, func(t *testing.T) { + fake := newFakeRunner() + switch boundary { + case "executable": + fake.candidateVersion = "0.80.9" + case "environment": + fake.candidateExpectedVersion = "0.80.9" + case "image-label": + fake.candidateLabelVersion = "0.80.9" + } + + result, err := Update(context.Background(), fake, Request{ + StatePath: filepath.Join(t.TempDir(), "state.json"), + Version: "0.81.0", + Source: BuildSource, + Confirm: true, + }) + + if err == nil || result.Phase != PhaseRolledBack { + t.Fatalf("Update() = %+v, %v; want rollback for candidate %s mismatch", result, err, boundary) + } + if fake.currentImage != "sha256:old" { + t.Fatalf("current image = %q, want restored previous image", fake.currentImage) + } + }) + } +} + func TestEveryRecoveryStateWriteFailureIsHandledTransactionally(t *testing.T) { for failAt := 1; failAt <= 6; failAt++ { t.Run(fmt.Sprintf("write-%d", failAt), func(t *testing.T) { @@ -604,19 +738,35 @@ type fakeRunner struct { dropMaintenanceAfterCandidate bool modelsWire string rollbackPrepared bool + coreRunning bool + execFailuresWhileStopped int + maintenanceHelperImages []string + maintenanceClearImages []string + restoredProofComplete bool + candidateVersion string + candidateExpectedVersion string + candidateLabelVersion string } func newFakeRunner() *fakeRunner { return &fakeRunner{ version: "0.80.3", expectedVersion: "0.80.3", labelVersion: "0.80.3", currentImage: "sha256:old", configuredImage: "thothii-core:local", - tags: map[string]string{"thothii-core:local": "sha256:old"}, - imageVersions: map[string]string{"sha256:old": "0.80.3"}, + tags: map[string]string{"thothii-core:local": "sha256:old"}, + imageVersions: map[string]string{"sha256:old": "0.80.3"}, + coreRunning: true, + candidateVersion: "0.81.0", + candidateExpectedVersion: "0.81.0", + candidateLabelVersion: "0.81.0", } } func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose.Result, error) { call := strings.Join(args, " ") f.calls = append(f.calls, call) + if !f.coreRunning && containsArg(args, "exec") { + f.execFailuresWhileStopped++ + return compose.Result{ExitCode: 1}, errors.New("core service is not running") + } if f.built && f.fail != "compensation" && strings.Contains(call, "image tag sha256:old") { f.fail = "" } @@ -651,6 +801,11 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose endpoint = "https://drift.example.invalid" } return compose.Result{Stdout: `{"services":{"core":{"image":"` + selectedCoreReference(args, f.configuredImage) + `","environment":{"THT_LLM_URL":"` + endpoint + `"}}}}`}, nil + case strings.Contains(call, "ps --status running -q core"): + if f.coreRunning { + return compose.Result{Stdout: "core-container\n"}, nil + } + return compose.Result{}, nil case strings.Contains(call, "ps -q core"): return compose.Result{Stdout: "core-container\n"}, nil case strings.Contains(call, "inspect --format {{.Image}}"): @@ -667,6 +822,9 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose return compose.Result{Stdout: `[{"Type":"volume","Name":"settings","Source":"settings","Destination":"/data/settings","RW":true},{"Type":"volume","Name":"pi-state","Source":"pi-state","Destination":"/home/thoth/.pi","RW":true},{"Type":"volume","Name":"sessions","Source":"sessions","Destination":"/data/sessions","RW":true},{"Type":"volume","Name":"workspace-registry","Source":"workspace-registry","Destination":"/data/workspace-registry","RW":true}]`}, nil case strings.Contains(call, "/internal/maintenance/activate"): f.maintenance = true + if f.fail == "maintenance-activate-durability" { + return compose.Result{ExitCode: 22}, errors.New("maintenance activation durability was not acknowledged") + } if f.lostMaintenanceResponse == "activate" { f.lostMaintenanceResponse = "" return compose.Result{ExitCode: 52}, errors.New("lost activation response") @@ -677,12 +835,16 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose return compose.Result{ExitCode: 53}, errors.New("maintenance clear failure") } f.maintenance = false + f.maintenanceClearImages = append(f.maintenanceClearImages, f.currentImage) if f.lostMaintenanceResponse == "deactivate" { f.lostMaintenanceResponse = "" return compose.Result{ExitCode: 52}, errors.New("lost deactivation response") } return compose.Result{Stdout: `{"active":false,"admissions":0}`}, nil case strings.Contains(call, "/internal/maintenance/status"): + if f.fail == "maintenance-activate-durability" { + return compose.Result{Stdout: fmt.Sprintf(`{"active":%t,"admissions":0,"recoveryRequired":true}`, f.maintenance)}, nil + } return compose.Result{Stdout: fmt.Sprintf(`{"active":%t,"admissions":0}`, f.maintenance)}, nil case strings.Contains(call, "/sessions?scope=all"): if f.sessionsWire != "" { @@ -693,6 +855,17 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose return compose.Result{Stdout: `[{"status":"open","archived":false}]`}, nil } return compose.Result{Stdout: `[]`}, nil + case containsArg(args, "run") && containsArg(args, "--entrypoint") && containsArg(args, "node"): + reference := selectedCoreReference(args, f.configuredImage) + f.maintenanceHelperImages = append(f.maintenanceHelperImages, reference) + if reference == "" { + return compose.Result{ExitCode: 1}, errors.New("maintenance helper has no selected image") + } + if f.tags[reference] != "sha256:old" { + return compose.Result{ExitCode: 1}, errors.New("maintenance helper did not select the previous image") + } + f.maintenance = true + return compose.Result{}, nil case containsArg(args, "build"): f.built = true f.buildReference = selectedCoreReference(args, f.configuredImage) @@ -723,6 +896,21 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose if version, ok := f.imageVersions[f.currentImage]; ok { f.version = version } + if f.currentImage == "sha256:candidate" { + f.version = f.candidateVersion + f.expectedVersion = f.candidateExpectedVersion + f.labelVersion = f.candidateLabelVersion + } else if f.currentImage == "sha256:old" { + f.expectedVersion = "0.80.3" + f.labelVersion = "0.80.3" + } + f.coreRunning = true + if f.currentImage == "sha256:candidate" { + f.restoredProofComplete = false + } + if f.fail == "dead-candidate" && f.currentImage == "sha256:candidate" { + f.coreRunning = false + } if f.dropMaintenanceAfterCandidate && f.currentImage == "sha256:candidate" { f.maintenance = false f.dropMaintenanceAfterCandidate = false @@ -744,6 +932,9 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose } return compose.Result{Stdout: `{"models":[{"id":"model","provider":"provider"}]}`}, nil case strings.Contains(call, "/settings"): + if f.currentImage == "sha256:old" { + f.restoredProofComplete = true + } return compose.Result{Stdout: `{"provider":"provider","model":"model","thinking":"medium"}`}, nil case strings.Contains(call, "/health"): return compose.Result{Stdout: `{"status":"ok"}`}, nil @@ -840,6 +1031,25 @@ func assertNotCalled(t *testing.T, calls []string, prohibited string) { } } } + +func assertMaintenanceClearedAfterRestoredProof(t *testing.T, fake *fakeRunner) { + t.Helper() + if fake.maintenance { + t.Fatal("maintenance remained active after restored-core proof") + } + if !fake.restoredProofComplete { + t.Fatal("maintenance cleared before restored settings smoke completed") + } + if len(fake.maintenanceClearImages) == 0 { + t.Fatal("maintenance was never durably cleared") + } + for _, image := range fake.maintenanceClearImages { + if image != "sha256:old" { + t.Fatalf("maintenance cleared while image %q was selected; want previous image", image) + } + } +} + func readStateBytes(t *testing.T, path string) []byte { t.Helper() contents, err := os.ReadFile(path)