fix: harden embedded Pi lifecycle recovery

This commit is contained in:
2026-08-04 23:14:47 +02:00
parent a368889838
commit 935bb1db0e
15 changed files with 572 additions and 101 deletions
+6 -6
View File
@@ -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",
+2
View File
@@ -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",
});
}
+26 -9
View File
@@ -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));
}
}
+32 -5
View File
@@ -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 });
}
+13 -3
View File
@@ -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 });
}
+1 -1
View File
@@ -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 ./
+31 -12
View File
@@ -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.
+8 -2
View File
@@ -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'
+9 -5
View File
@@ -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
)
+15 -11
View File
@@ -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=
+8 -4
View File
@@ -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")
}
@@ -61,8 +61,50 @@ 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
restoreDurabilityFailure bool
settings Defaults
settingsExist bool
settingsRaw []byte
@@ -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") {
@@ -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}
}
+126 -27
View File
@@ -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 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)
@@ -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)
}
@@ -421,6 +498,7 @@ func setMaintenance(ctx context.Context, runner Runner, enabled bool) error {
type MaintenanceState struct {
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 {
+210
View File
@@ -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,6 +738,14 @@ 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 {
@@ -611,12 +753,20 @@ func newFakeRunner() *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"},
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)