fix: finalize durable Pi lifecycle

This commit is contained in:
2026-08-04 21:34:08 +02:00
parent 5ba2821a1b
commit a368889838
20 changed files with 1016 additions and 147 deletions
+18 -4
View File
@@ -111,12 +111,26 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
maintenanceBarrier,
});
app.post("/internal/maintenance/activate", async (req, reply) => {
await maintenanceBarrier.activate();
return maintenanceBarrier.status();
try {
await maintenanceBarrier.activate();
return maintenanceBarrier.status();
} catch {
return reply.code(500).send({
...maintenanceBarrier.status(),
error: "maintenance activation durability was not acknowledged",
});
}
});
app.post("/internal/maintenance/deactivate", async (req, reply) => {
maintenanceBarrier.deactivate();
return maintenanceBarrier.status();
try {
maintenanceBarrier.deactivate();
return maintenanceBarrier.status();
} catch {
return reply.code(500).send({
...maintenanceBarrier.status(),
error: "maintenance deactivation durability was not acknowledged",
});
}
});
app.get("/internal/maintenance/status", async (req, reply) => {
return maintenanceBarrier.status();
+52 -15
View File
@@ -10,17 +10,25 @@ import {
} from "node:fs";
import { dirname } from "node:path";
export interface MaintenanceDurability {
syncDirectory(directory: string): void;
}
/** A durable admission barrier. A lease spans the complete create/resume decision. */
export class MaintenanceBarrier {
private active: boolean;
private admissions = 0;
private waiters: (() => void)[] = [];
constructor(private readonly markerFile?: string) {
constructor(
private readonly markerFile?: string,
private readonly durability: MaintenanceDurability = defaultDurability,
) {
this.active = markerFile === undefined ? false : existsSync(markerFile);
}
acquire(): (() => void) | undefined {
this.reconcileActive();
if (this.active) return undefined;
this.admissions += 1;
let released = false;
@@ -33,17 +41,44 @@ export class MaintenanceBarrier {
}
async activate(): Promise<void> {
this.persistMarker();
this.active = true;
if (this.admissions === 0) return;
await new Promise<void>((resolve) => this.waiters.push(resolve));
let persistError: unknown;
if (this.markerFile === undefined) {
this.active = true;
} else {
try {
this.persistMarker();
} catch (error) {
persistError = error;
} finally {
this.reconcileActive();
}
}
if (!this.active) throw persistError;
if (this.admissions > 0) {
await new Promise<void>((resolve) => this.waiters.push(resolve));
}
if (persistError !== undefined) throw persistError;
}
deactivate(): void {
this.removeMarker();
this.active = false;
if (this.markerFile === undefined) {
this.active = false;
return;
}
try {
this.removeMarker();
} finally {
this.reconcileActive();
}
}
status(): { active: boolean; admissions: number } {
this.reconcileActive();
return { active: this.active, admissions: this.admissions };
}
private reconcileActive(): void {
if (this.markerFile !== undefined) this.active = existsSync(this.markerFile);
}
status(): { active: boolean; admissions: number } { return { active: this.active, admissions: this.admissions }; }
private persistMarker(): void {
if (!this.markerFile) return;
@@ -59,7 +94,7 @@ export class MaintenanceBarrier {
}
try {
renameSync(temporary, this.markerFile);
syncDirectory(directory);
this.durability.syncDirectory(directory);
} catch (error) {
try { unlinkSync(temporary); } catch { /* already renamed or best-effort cleanup */ }
throw error;
@@ -69,12 +104,14 @@ export class MaintenanceBarrier {
private removeMarker(): void {
if (!this.markerFile || !existsSync(this.markerFile)) return;
unlinkSync(this.markerFile);
syncDirectory(dirname(this.markerFile));
this.durability.syncDirectory(dirname(this.markerFile));
}
}
function syncDirectory(directory: string): void {
if (process.platform === "win32") return;
const fd = openSync(directory, "r");
try { fsyncSync(fd); } finally { closeSync(fd); }
}
const defaultDurability: MaintenanceDurability = {
syncDirectory(directory: string): void {
if (process.platform === "win32") return;
const fd = openSync(directory, "r");
try { fsyncSync(fd); } finally { closeSync(fd); }
},
};
+27 -9
View File
@@ -1,7 +1,15 @@
/* Core-side, non-interactive installation-default writer used only through compose exec.
* It accepts no credentials and writes the same SETTINGS_FILE consumed by session creation. */
import { readFileSync } from "node:fs";
import { loadConfig } from "../config.js";
import { loadSettings, saveSettings, type Settings } from "./settings-store.js";
import {
captureSettingsSnapshot,
loadSettings,
restoreSettingsSnapshot,
saveSettings,
type Settings,
type SettingsSnapshot,
} from "./settings-store.js";
const choice = /^[A-Za-z0-9][A-Za-z0-9._/-]{0,127}$/;
@@ -15,15 +23,25 @@ function value(args: string[], flag: string): string {
try {
const args = process.argv.slice(2);
if (args.length !== 6) throw new Error("only provider, model, and thinking may be configured");
const provider = value(args, "--provider");
const model = value(args, "--model");
const thinking = value(args, "--thinking");
if (!choice.test(provider) || !choice.test(model)) throw new Error("invalid provider or model");
if (!["low", "medium", "high"].includes(thinking)) throw new Error("invalid thinking level");
const cfg = loadConfig(process.env);
const next: Settings = { ...loadSettings(cfg), provider, model, thinking };
saveSettings(cfg, next);
if (args.length === 1 && args[0] === "--snapshot") {
process.stdout.write(`${JSON.stringify(captureSettingsSnapshot(cfg))}\n`);
} else if (args.length === 1 && args[0] === "--restore") {
const parsed = JSON.parse(readFileSync(0, "utf8")) as Partial<SettingsSnapshot>;
if (Object.keys(parsed).some((key) => key !== "exists" && key !== "rawBase64")) {
throw new Error("invalid settings snapshot");
}
restoreSettingsSnapshot(cfg, parsed as SettingsSnapshot);
} else {
if (args.length !== 6) throw new Error("only provider, model, and thinking may be configured");
const provider = value(args, "--provider");
const model = value(args, "--model");
const thinking = value(args, "--thinking");
if (!choice.test(provider) || !choice.test(model)) throw new Error("invalid provider or model");
if (!["low", "medium", "high"].includes(thinking)) throw new Error("invalid thinking level");
const next: Settings = { ...loadSettings(cfg), provider, model, thinking };
saveSettings(cfg, next);
}
} catch (error) {
process.stderr.write(`settings-cli: ${error instanceof Error ? error.message : "invalid configuration"}\n`);
process.exitCode = 2;
+44
View File
@@ -22,6 +22,11 @@ export interface SettingsDurability {
syncDirectory(directory: string): void;
}
export interface SettingsSnapshot {
exists: boolean;
rawBase64: string;
}
/**
* Settings files are installation defaults only. Personal workspace/model/thinking choices
* belong to the browser and must never be written back here by request handlers.
@@ -39,6 +44,45 @@ export function loadSettings(cfg: AppConfig): Settings {
}
}
/** Capture exact file existence and bytes so host-side configuration can compensate losslessly. */
export function captureSettingsSnapshot(cfg: AppConfig): SettingsSnapshot {
try {
return { exists: true, rawBase64: readFileSync(cfg.settingsFile).toString("base64") };
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") {
return { exists: false, rawBase64: "" };
}
throw error;
}
}
/** Restore a previously captured settings file exactly, including the clean absent state. */
export function restoreSettingsSnapshot(
cfg: AppConfig,
snapshot: SettingsSnapshot,
durability: SettingsDurability = defaultDurability,
): void {
if (typeof snapshot.exists !== "boolean" || typeof snapshot.rawBase64 !== "string") {
throw new Error("invalid settings snapshot");
}
const raw = Buffer.from(snapshot.rawBase64, "base64");
if (raw.toString("base64") !== snapshot.rawBase64 || (!snapshot.exists && raw.length !== 0)) {
throw new Error("invalid settings snapshot");
}
const directory = dirname(cfg.settingsFile);
mkdirSync(directory, { recursive: true });
if (snapshot.exists) {
replaceSettingsFile(cfg.settingsFile, raw);
} else {
try {
unlinkSync(cfg.settingsFile);
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error;
}
}
durability.syncDirectory(directory);
}
/** Persist settings (pretty JSON). Creates the parent directory if needed. */
export function saveSettings(
cfg: AppConfig,
+56 -1
View File
@@ -1,5 +1,5 @@
import { test, expect } from "vitest";
import { existsSync, mkdtempSync, rmSync } from "node:fs";
import { existsSync, mkdtempSync, rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { MaintenanceBarrier } from "../src/runtime/maintenance-gate.js";
@@ -45,3 +45,58 @@ test("durable activation survives recreation and deactivation removes the marker
rmSync(directory, { recursive: true, force: true });
}
});
test("activation reconciles active state when directory fsync fails after marker rename", async () => {
const directory = mkdtempSync(join(tmpdir(), "tht-maintenance-activate-fsync-"));
const marker = join(directory, "maintenance.json");
try {
const gate = new MaintenanceBarrier(marker, {
syncDirectory() { throw new Error("injected post-rename fsync failure"); },
});
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.acquire()).toBeUndefined();
const recovered = new MaintenanceBarrier(marker);
expect(recovered.status()).toEqual({ active: true, admissions: 0 });
expect(recovered.acquire()).toBeUndefined();
} 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");
let failSync = false;
try {
const gate = new MaintenanceBarrier(marker, {
syncDirectory() {
if (failSync) throw new Error("injected post-remove fsync failure");
},
});
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");
} finally {
rmSync(directory, { recursive: true, force: true });
}
});
test("status and admission recover from a persistent marker even when memory started inactive", () => {
const directory = mkdtempSync(join(tmpdir(), "tht-maintenance-reconcile-"));
const marker = join(directory, "maintenance.json");
try {
const gate = new MaintenanceBarrier(marker);
expect(gate.status().active).toBe(false);
writeFileSync(marker, '{"version":1,"active":true}\n', { mode: 0o600 });
expect(gate.status()).toEqual({ active: true, admissions: 0 });
expect(gate.acquire()).toBeUndefined();
} finally {
rmSync(directory, { recursive: true, force: true });
}
});
+32
View File
@@ -152,6 +152,38 @@ test.each(["none", "upstream"] as const)(
},
);
test("maintenance endpoints report marker-derived state after post-rename and post-remove fsync failures", async () => {
const dir = mkdtempSync(path.join(tmpdir(), "tht-maintenance-endpoint-fsync-"));
const marker = path.join(dir, "maintenance.json");
let failSync = true;
try {
const maintenanceBarrier = new MaintenanceBarrier(marker, {
syncDirectory() {
if (failSync) throw new Error("injected maintenance fsync failure");
},
});
const app = buildApp(loadConfig({
AUTH_MODE: "none",
THT_HARNESS_DIR: "../harness",
THT_MAINTENANCE_FILE: marker,
}), { thtRunner: {} as any, maintenanceBarrier });
const activated = await app.inject({ method: "POST", url: "/internal/maintenance/activate" });
expect(activated.statusCode).toBe(500);
expect(activated.json()).toMatchObject({ active: true, admissions: 0 });
failSync = false;
expect((await app.inject({ method: "GET", url: "/internal/maintenance/status" })).json())
.toEqual({ active: true, admissions: 0 });
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 });
} finally {
rmSync(dir, { recursive: true, force: true });
}
});
test("admin all-sessions response matches the authenticated lifecycle wire fixture", async () => {
const fixture = JSON.parse(readFileSync(
path.join(import.meta.dirname, "fixtures", "sessions-scope-all.json"),
+29 -2
View File
@@ -1,8 +1,13 @@
import { test, expect } from "vitest";
import { closeSync, fsyncSync, mkdtempSync, openSync, rmSync, writeFileSync } from "node:fs";
import { closeSync, existsSync, fsyncSync, mkdtempSync, openSync, readFileSync, rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { loadSettings, saveSettings } from "../src/settings/settings-store.js";
import {
captureSettingsSnapshot,
loadSettings,
restoreSettingsSnapshot,
saveSettings,
} from "../src/settings/settings-store.js";
import { loadConfig } from "../src/config.js";
function cfgWith(file: string) {
@@ -70,3 +75,25 @@ test("saveSettings restores the previous file when post-rename directory durabil
rmSync(dir, { recursive: true, force: true });
}
});
test("settings snapshots restore exact absent and empty-file states", () => {
const dir = mkdtempSync(join(tmpdir(), "tht-settings-snapshot-"));
try {
const file = join(dir, "settings.json");
const cfg = cfgWith(file);
const absent = captureSettingsSnapshot(cfg);
expect(absent).toEqual({ exists: false, rawBase64: "" });
saveSettings(cfg, { provider: "new", model: "model", thinking: "high" });
restoreSettingsSnapshot(cfg, absent);
expect(existsSync(file)).toBe(false);
writeFileSync(file, Buffer.alloc(0), { mode: 0o600 });
const empty = captureSettingsSnapshot(cfg);
expect(empty.exists).toBe(true);
saveSettings(cfg, { provider: "new", model: "model", thinking: "high" });
restoreSettingsSnapshot(cfg, empty);
expect(readFileSync(file)).toEqual(Buffer.alloc(0));
} finally {
rmSync(dir, { recursive: true, force: true });
}
});
@@ -66,6 +66,8 @@ test("CLI accepts an explicit valid ID when a legacy filename contains dots", as
test("declares a durable isolated registry volume and only read-only Git credential mounts", () => {
const compose = readFileSync(new URL("../../compose.yaml", import.meta.url), "utf8");
const development = readFileSync(new URL("../../docker-compose.dev.yml", import.meta.url), "utf8");
const gitHttps = readFileSync(new URL("../../deploy/compose.git-https.yaml", import.meta.url), "utf8");
const gitSsh = readFileSync(new URL("../../deploy/compose.git-ssh.yaml", import.meta.url), "utf8");
const dockerfile = readFileSync(new URL("../../docker/core.Dockerfile", import.meta.url), "utf8");
const smoke = readFileSync(new URL("../../scripts/workspace-registry-smoke.sh", import.meta.url), "utf8");
@@ -73,12 +75,14 @@ test("declares a durable isolated registry volume and only read-only Git credent
expect(source).toContain("THT_WORKSPACE_REGISTRY_ROOT: /data/workspace-registry");
expect(source).toContain("THT_WORKSPACE_GIT_REMOTE: ${THT_WORKSPACE_GIT_REMOTE:?set THT_WORKSPACE_GIT_REMOTE}");
expect(source).toContain("workspace-registry:/data/workspace-registry");
expect(source).toMatch(/workspace-registry-git-credentials:ro/);
expect(source).toMatch(/workspace-registry-git-ca:ro/);
expect(source).toMatch(/workspace-registry-git-ssh-key:ro/);
expect(source).toMatch(/workspace-registry-git-known-hosts:ro/);
}
expect(dockerfile).toMatch(/mkdir -p \/data\/workspace-registry && chown -R thoth:thoth \/data\/workspace-registry/);
expect(compose).not.toMatch(/workspace-registry-git-(?:credentials|ca|ssh-key|known-hosts):ro/);
expect(gitHttps).toMatch(/workspace-registry-git-credentials:ro/);
expect(gitHttps).toMatch(/workspace-registry-git-ca:ro/);
expect(gitSsh).toMatch(/workspace-registry-git-ssh-key:ro/);
expect(gitSsh).toMatch(/workspace-registry-git-known-hosts:ro/);
expect(dockerfile).toMatch(/mkdir -p[^\n]*\/data\/workspace-registry/);
expect(dockerfile).toMatch(/chown -R thoth:thoth \/home\/thoth\/\.pi \/data/);
expect(smoke).toContain('core_remote="/fixtures/offline.git"');
expect(smoke).toContain('"degraded":true');
expect(smoke).toContain('core_remote="/fixtures/remote.git"');
+53 -16
View File
@@ -13,10 +13,12 @@ thothctl --installation /absolute/path/thothii-installation.yaml pi logs
thothctl --installation /absolute/path/thothii-installation.yaml pi configure
```
`status` executes the image-bundled `pi --version`. `doctor` and `test` require a healthy core,
a successful Pi smoke, valid settings, and an exact selected provider/model pair from the
backend's available model entries. `pi check` remains an alias for `pi test`. Logs are always a
bounded, sanitized 200-line snapshot; there is no follow mode.
`status` executes the image-bundled `pi --version`. `doctor` compares that value with both the
container's `PI_VERSION` contract and the `io.thothii.pi.version` image label; a merely nonempty
version is not sufficient. `doctor` and `test` also require a healthy core, a successful Pi smoke,
valid settings, and an exact selected provider/model pair from the backend's available model
entries. `pi check` remains an alias for `pi test`. Logs are always a bounded, sanitized 200-line
snapshot; there is no follow mode.
On a TTY, `pi configure` presents numbered provider, model, and thinking choices. Providers and
models come from the backend's closed model list, and the model choices are restricted to the
@@ -27,10 +29,29 @@ thothctl --installation /absolute/path/thothii-installation.yaml pi configure \
--provider zai --model glm-5.2 --thinking medium
```
The helper snapshots the previous settings, applies the new values atomically, and verifies the
readback and rendered-configuration digest. A helper, readback, or digest failure triggers an
attempted restore followed by another readback. The command reports the actual host path from
`PI_AUTH_FILE`; credentials remain in that protected host file and must never be passed as flags.
The helper snapshots exact settings-file existence and raw bytes, applies the new values atomically,
and verifies the readback and rendered-configuration digest. A helper, readback, or digest failure
restores those exact bytes when the prior file existed; on a clean installation it removes the new
file and verifies the absent/default state. Empty prior files are supported. The command reports
the actual host path from `PI_AUTH_FILE`; credentials remain in that protected host file and must
never be passed as flags.
## Supported Compose entry points and current image
Use `thothctl start`, `stop`, `status`, `logs`, and `doctor` for ordinary installation lifecycle
operations. All `thothctl` Compose commands automatically include the installation-specific
durable selector when it exists:
```text
<projectDirectory>/.thothctl/<installation-id>/current-image.yaml
```
This selector is part of the supported installation state: it keeps a verified Pi image selected
across a fresh `thothctl` process, stop/start, reconcile, and source checkout whose base image is
digest-pinned. Do not delete or hand-edit it. Direct raw `docker compose` lifecycle commands bypass
this protection and are unsupported. Advanced documented Compose rendering must use
`scripts/compose-with-preflight.sh` and include the same selector with `-f` when present; connector
secret overrides must never bypass that preflight wrapper.
## Updating Pi
@@ -63,7 +84,11 @@ their work, `--drain` makes the command poll the authenticated bare-array
The configured `core.image` is never retagged or mutated. Each installation transaction creates
unique candidate and previous tags, including when two installations share a configured tag or
the configured image is digest-pinned. A temporary lifecycle-only Compose override selects those
tags for build, recreate, and rollback. Terminal success removes the override.
tags for build, recreate, and rollback. Candidate build, pull, or candidate-tag failures happen
before `mutation_started` and therefore never recreate or roll back core. After verification, the
temporary candidate selector is atomically promoted to the durable current-image override.
Rollback atomically promotes the previous selector. Terminal cleanup removes only transaction
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
@@ -75,8 +100,8 @@ persistence-mount fingerprint.
Recovery state and lock diagnostics live under:
```text
<projectDirectory>/.thothctl/update-state.json
<projectDirectory>/.thothctl/update-state.json.lock.owner.json
<projectDirectory>/.thothctl/<installation-id>/update-state.json
<projectDirectory>/.thothctl/<installation-id>/update-state.json.lock.owner.json
```
The recovery file is mode `0600` and records transaction-scoped image identities, mount
@@ -103,12 +128,24 @@ thothctl --installation /absolute/path/thothii-installation.yaml pi maintenance
thothctl --installation /absolute/path/thothii-installation.yaml pi maintenance recover --yes
```
`maintenance recover` refuses a pending transaction. For terminal or absent recovery state, it
removes a stale lifecycle override, verifies the running installation when the gate is active, and
only then removes the durable marker and reopens admission. If rollback or recovery fails, leave
the marker in place, preserve `update-state.json`, repair the reported Docker/configuration issue,
and rerun rollback or maintenance recovery.
`maintenance recover` completes an interrupted verified-image promotion, safely finalizes a
preparation interrupted before core mutation, and refuses other pending mutations. For terminal or
absent recovery state, it removes only a stale transaction override, verifies the running
installation when the gate is active, and only then removes the durable marker and reopens
admission. It never removes `current-image.yaml`. If rollback or recovery fails, leave the marker
in place, preserve `update-state.json`, repair the reported Docker/configuration issue, and rerun
rollback or maintenance recovery.
Missing confirmation, invalid arguments, active sessions, and an interrupted transaction exit
`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
`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.
+3 -3
View File
@@ -189,7 +189,7 @@ func piCommand(ctx context.Context, installation config.Installation, runner com
fmt.Fprintf(stdout, "Pi defaults applied and read back. Provider credentials remain only in the host file %s (mode 0600). Never pass credentials to thothctl.\n", authFile)
return 0
case "update":
request, err := parsePiUpdateArgs(args[1:], filepath.Join(installation.ProjectDirectory, ".thothctl", "update-state.json"))
request, err := parsePiUpdateArgs(args[1:], installation.UpdateStatePath())
if err != nil {
return commandUsageError(stderr, err.Error())
}
@@ -207,7 +207,7 @@ func piCommand(ctx context.Context, installation config.Installation, runner com
if len(args) != 2 || args[1] != "--yes" {
return commandUsageError(stderr, "pi rollback requires --yes")
}
result, err := pi.Rollback(ctx, controlled, filepath.Join(installation.ProjectDirectory, ".thothctl", "update-state.json"), true)
result, err := pi.Rollback(ctx, controlled, installation.UpdateStatePath(), true)
if err != nil {
return piFailure(stderr, err, secretValues)
}
@@ -223,7 +223,7 @@ func piCommand(ctx context.Context, installation config.Installation, runner com
return 0
}
if len(args) == 3 && args[1] == "recover" && args[2] == "--yes" {
statePath := filepath.Join(installation.ProjectDirectory, ".thothctl", "update-state.json")
statePath := installation.UpdateStatePath()
if err := pi.RecoverMaintenance(ctx, controlled, statePath, true); err != nil {
return piFailure(stderr, err, secretValues)
}
+28 -1
View File
@@ -12,6 +12,7 @@ import (
"testing"
"github.com/aritmolab/thothii/tools/thothctl/internal/compose"
"github.com/aritmolab/thothii/tools/thothctl/internal/config"
"github.com/aritmolab/thothii/tools/thothctl/internal/pi"
"github.com/aritmolab/thothii/tools/thothctl/internal/testsupport"
)
@@ -50,7 +51,8 @@ func TestResolvePiConfigureRequiresExplicitFlagsWithoutTTY(t *testing.T) {
func TestPiLifecycleContractErrorsExitTwo(t *testing.T) {
for _, lifecycleErr := range []error{pi.ErrActiveSessions, pi.ErrInterruptedUpdate} {
var stderr bytes.Buffer
if code := piFailure(&stderr, lifecycleErr, nil); code != 2 {
wrapped := fmt.Errorf("automatic rollback succeeded: %w", lifecycleErr)
if code := piFailure(&stderr, wrapped, nil); code != 2 {
t.Errorf("piFailure(%v) = %d, want 2", lifecycleErr, code)
}
}
@@ -373,6 +375,28 @@ func TestRunStatusUsesStableComposeArguments(t *testing.T) {
}
}
func TestRunStartAutomaticallyUsesTheDurableCurrentImageOverride(t *testing.T) {
fixture := newCLIFixture(t, "SAFE_VALUE=1\n")
fixture.setEnvironment(t)
installation, err := config.Load(fixture.installationPath)
if err != nil {
t.Fatal(err)
}
currentImage := installation.CurrentImageOverridePath()
if err := os.MkdirAll(filepath.Dir(currentImage), 0o700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(currentImage, []byte("services:\n core:\n image: thothii-core:verified\n"), 0o600); err != nil {
t.Fatal(err)
}
var stdout, stderr bytes.Buffer
if code := run(context.Background(), []string{"--installation", fixture.installationPath, "start"}, &stdout, &stderr); code != 0 {
t.Fatalf("start exit = %d, stderr = %s", code, stderr.String())
}
assertInvocationContains(t, fixture.invocations(t), "-f", currentImage, "up", "--detach", "--remove-orphans")
}
func TestRunExplainsWhenDockerIsNotAvailable(t *testing.T) {
fixture := newCLIFixture(t, "SAFE_VALUE=1\n")
fixture.setEnvironment(t)
@@ -591,8 +615,11 @@ printf '%s\n' -- >> "$THOTHCTL_FAKE_ARGS"
case " $* " in
*" config --format json "*) printf '%s\n' '{"volumes":{"settings":{}},"services":{"core":{"image":"thothii-core:local","environment":{"THT_LLM_URL":"https://llm.example.invalid"}}}}' ;;
*" ps --format json "*) printf '%s\n' '[{"Service":"core","State":"running","Health":"healthy"},{"Service":"frontend","State":"running","Health":"healthy"}]' ;;
*"io.thothii.pi.version"*) printf '%s\n' '0.80.3' ;;
*"PI_VERSION"*) printf '%s\n' '0.80.3' ;;
*" pi --version "*) printf '%s\n' '0.80.3' ;;
*"/models "*) printf '%s\n' '{"models":[{"provider":"provider","id":"model"}]}' ;;
*"settings-cli.js --snapshot"*) printf '%s\n' '{"exists":false,"rawBase64":""}' ;;
*"/settings "*) printf '%s\n' '{"provider":"provider","model":"model","thinking":"medium"}' ;;
*"/internal/maintenance/status "*) printf '%s\n' '{"active":true,"admissions":0}' ;;
*" logs "*) printf '%s\n' "$THOTHCTL_FAKE_LOG" ;;
+4 -2
View File
@@ -1,6 +1,8 @@
module github.com/aritmolab/thothii/tools/thothctl
go 1.24
go 1.24.0
toolchain go1.24.13
require gopkg.in/yaml.v3 v3.0.1
@@ -9,7 +11,7 @@ require (
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.22.0
golang.org/x/sys v0.41.0
)
require github.com/opencontainers/go-digest v1.0.0 // indirect
+2
View File
@@ -24,6 +24,8 @@ 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=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+33 -3
View File
@@ -101,6 +101,13 @@ func Load(path string) (Installation, error) {
return Installation{}, err
}
}
if info, err := os.Lstat(installation.CurrentImageOverridePath()); err == nil {
if !info.Mode().IsRegular() {
return Installation{}, errors.New("installation current-image override must be a regular file")
}
} else if !errors.Is(err, os.ErrNotExist) {
return Installation{}, errors.New("installation current-image override could not be inspected")
}
return installation, nil
}
@@ -111,7 +118,26 @@ func (i Installation) ComposeFiles() []string {
filepath.Join(i.ProjectDirectory, "compose.yaml"),
filepath.Join(i.ProjectDirectory, "deploy", "compose."+i.Profile+".yaml"),
}
return append(files, i.Overrides...)
files = append(files, i.Overrides...)
currentImage := i.CurrentImageOverridePath()
if info, err := os.Lstat(currentImage); err == nil && info.Mode().IsRegular() {
files = append(files, currentImage)
}
return files
}
// ControlDirectory contains state that is private to one installation descriptor, even when
// multiple installations intentionally share one source checkout.
func (i Installation) ControlDirectory() string {
return filepath.Join(i.ProjectDirectory, ".thothctl", i.ProjectName())
}
func (i Installation) CurrentImageOverridePath() string {
return filepath.Join(i.ControlDirectory(), "current-image.yaml")
}
func (i Installation) UpdateStatePath() string {
return filepath.Join(i.ControlDirectory(), "update-state.json")
}
// ProjectName is stable for one installation and avoids collisions between different checkouts.
@@ -168,9 +194,13 @@ func (i Installation) SecretFiles() ([]string, error) {
// callers. It is used only for operator-visible file locations, never for secret content.
func (i Installation) EnvironmentValue(name string) (string, error) {
contents, err := safeio.ReadCanonicalRegular(i.EnvFile, maxEnvironmentFileBytes)
if err != nil { return "", errors.New("installation environment could not be read") }
if err != nil {
return "", errors.New("installation environment could not be read")
}
values, err := parseComposeDotenv(contents)
if err != nil { return "", errors.New("installation environment could not be read") }
if err != nil {
return "", errors.New("installation environment could not be read")
}
return values[name], nil
}
@@ -51,6 +51,51 @@ func TestLoadSelectsServerComposeFiles(t *testing.T) {
assertStringsEqual(t, installation.ComposeFiles(), want)
}
func TestComposeArgsAutomaticallyIncludeTheInstallationCurrentImageOverride(t *testing.T) {
t.Parallel()
installationPath, _, _, _ := writeInstallation(t, "local")
seed, err := Load(installationPath)
if err != nil {
t.Fatal(err)
}
currentImage := seed.CurrentImageOverridePath()
if err := os.MkdirAll(filepath.Dir(currentImage), 0o700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(currentImage, []byte("services:\n core:\n image: candidate\n"), 0o600); err != nil {
t.Fatal(err)
}
installation, err := Load(installationPath)
if err != nil {
t.Fatal(err)
}
args := installation.ComposeArgs("up", "--detach")
want := []string{"-f", currentImage, "up", "--detach"}
if !containsSequence(args, want) {
t.Fatalf("ComposeArgs() = %#v, want durable override immediately before command", args)
}
}
func TestInstallationControlPathsAreIsolatedForDescriptorsSharingOneCheckout(t *testing.T) {
projectDirectory := t.TempDir()
first := Installation{Path: filepath.Join(t.TempDir(), installationFileName), ProjectDirectory: projectDirectory}
second := Installation{Path: filepath.Join(t.TempDir(), installationFileName), ProjectDirectory: projectDirectory}
if first.CurrentImageOverridePath() == second.CurrentImageOverridePath() {
t.Fatalf("shared-checkout installations reused %q", first.CurrentImageOverridePath())
}
for _, installation := range []Installation{first, second} {
if filepath.Dir(filepath.Dir(installation.CurrentImageOverridePath())) != filepath.Join(projectDirectory, ".thothctl") {
t.Fatalf("current-image path %q is not installation-specific under .thothctl", installation.CurrentImageOverridePath())
}
if filepath.Dir(installation.UpdateStatePath()) != filepath.Dir(installation.CurrentImageOverridePath()) {
t.Fatalf("state %q and selector %q do not share one installation control directory", installation.UpdateStatePath(), installation.CurrentImageOverridePath())
}
}
}
func TestLoadRejectsRelativeInstallationPaths(t *testing.T) {
t.Parallel()
@@ -100,3 +145,22 @@ func assertStringsEqual(t *testing.T, got, want []string) {
}
}
}
func containsSequence(values, wanted []string) bool {
for start := range values {
if len(values)-start < len(wanted) {
continue
}
matched := true
for offset := range wanted {
if values[start+offset] != wanted[offset] {
matched = false
break
}
}
if matched {
return true
}
}
return false
}
+94 -26
View File
@@ -1,8 +1,10 @@
package pi
import (
"bytes"
"context"
"crypto/sha256"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
@@ -26,6 +28,11 @@ type ModelOption struct {
ID string `json:"id"`
}
type settingsFileSnapshot struct {
Exists bool `json:"exists"`
RawBase64 string `json:"rawBase64"`
}
var internalIdentityHeaders = []string{
"-H", "x-thoth-principal-issuer: thothctl",
"-H", "x-thoth-principal-subject: thothctl-maintenance",
@@ -59,40 +66,34 @@ func Configure(ctx context.Context, runner Runner, value Defaults) error {
if !found {
return errors.New("provider/model is not in Pi options")
}
settingsArgs := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...)
settingsArgs = append(settingsArgs, "http://127.0.0.1:8787/settings")
oldResult, err := runCompose(ctx, runner, settingsArgs...)
old, err := captureSettingsFile(ctx, runner)
if err != nil {
return commandError("Pi installation settings capture", oldResult, err)
return err
}
var old Defaults
if json.Unmarshal([]byte(oldResult.Stdout), &old) != nil || old.Provider == "" || old.Model == "" || old.Thinking == "" {
return errors.New("Pi installation settings capture is invalid")
oldEffective, err := readEffectiveSettings(ctx, runner)
if err != nil {
return err
}
restore := func(cause error) error {
result, restoreErr := writeDefaults(context.Background(), runner, old)
if restoreErr != nil {
return fmt.Errorf("%w; previous Pi settings could not be restored: recovery required", cause)
}
if result.ExitCode != 0 {
return fmt.Errorf("%w; previous Pi settings could not be restored: recovery required", cause)
}
verified, readErr := readDefaults(context.Background(), runner, settingsArgs)
if readErr != nil || verified != old {
if restoreErr := restoreSettingsFile(context.Background(), runner, old); restoreErr != nil {
return fmt.Errorf("%w; previous Pi settings restoration could not be verified: recovery required", cause)
}
restoredEffective, restoreErr := readEffectiveSettings(context.Background(), runner)
if restoreErr != nil || !bytes.Equal(restoredEffective, oldEffective) {
return fmt.Errorf("%w; previous effective Pi settings could not be verified: recovery required", cause)
}
return cause
}
result, err := writeDefaults(ctx, runner, value)
if err != nil {
return restore(commandError("Pi installation settings write", result, err))
}
settings, err := runCompose(ctx, runner, settingsArgs...)
settings, err := readEffectiveSettings(ctx, runner)
if err != nil {
return restore(commandError("Pi installation settings read-back", settings, err))
return restore(err)
}
var saved Defaults
if json.Unmarshal([]byte(settings.Stdout), &saved) != nil || saved.Provider != value.Provider || saved.Model != value.Model || saved.Thinking != value.Thinking {
if json.Unmarshal(settings, &saved) != nil || saved.Provider != value.Provider || saved.Model != value.Model || saved.Thinking != value.Thinking {
return restore(errors.New("Pi installation settings read-back did not match requested provider, model, and thinking"))
}
after, err := renderedCore(ctx, runner)
@@ -130,16 +131,55 @@ func writeDefaults(ctx context.Context, runner Runner, value Defaults) (compose.
return runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/settings/settings-cli.js", "--provider", value.Provider, "--model", value.Model, "--thinking", value.Thinking)
}
func readDefaults(ctx context.Context, runner Runner, args []string) (Defaults, error) {
func captureSettingsFile(ctx context.Context, runner Runner) (settingsFileSnapshot, error) {
result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/settings/settings-cli.js", "--snapshot")
if err != nil {
return settingsFileSnapshot{}, commandError("Pi installation settings snapshot", result, err)
}
var snapshot settingsFileSnapshot
if json.Unmarshal([]byte(result.Stdout), &snapshot) != nil {
return settingsFileSnapshot{}, errors.New("Pi installation settings snapshot is invalid")
}
raw, decodeErr := base64.StdEncoding.DecodeString(snapshot.RawBase64)
if decodeErr != nil || base64.StdEncoding.EncodeToString(raw) != snapshot.RawBase64 || (!snapshot.Exists && len(raw) != 0) {
return settingsFileSnapshot{}, errors.New("Pi installation settings snapshot is invalid")
}
return snapshot, nil
}
func restoreSettingsFile(ctx context.Context, runner Runner, snapshot settingsFileSnapshot) error {
payload, err := json.Marshal(snapshot)
if err != nil {
return errors.New("Pi installation settings snapshot could not be encoded")
}
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 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")
}
func readEffectiveSettings(ctx context.Context, runner Runner) ([]byte, error) {
args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...)
args = append(args, "http://127.0.0.1:8787/settings")
result, err := runCompose(ctx, runner, args...)
if err != nil {
return Defaults{}, commandError("Pi installation settings restoration read-back", result, err)
return nil, commandError("Pi installation settings read-back", result, err)
}
var value Defaults
if json.Unmarshal([]byte(result.Stdout), &value) != nil {
return Defaults{}, errors.New("Pi installation settings restoration read-back is invalid")
var settings map[string]json.RawMessage
if json.Unmarshal([]byte(result.Stdout), &settings) != nil || settings == nil {
return nil, errors.New("Pi installation settings read-back is invalid")
}
return value, nil
canonical, err := json.Marshal(settings)
if err != nil {
return nil, errors.New("Pi installation settings read-back could not be normalized")
}
return canonical, nil
}
// Runner is the narrow, shell-free command boundary shared with thothctl.
@@ -165,9 +205,17 @@ func Doctor(ctx context.Context, runner Runner) error {
if _, err := renderedCore(ctx, runner); err != nil {
return err
}
if _, err := Status(ctx, runner); err != nil {
actual, err := Status(ctx, runner)
if err != nil {
return err
}
expected, label, err := expectedVersions(ctx, runner)
if err != nil {
return err
}
if actual != expected || actual != label {
return errors.New("Pi version does not match the image PI_VERSION and io.thothii.pi.version contract")
}
for _, check := range [][]string{
{"exec", "-T", "core", "sh", "-ceu", "test -w /home/thoth/.pi"},
{"exec", "-T", "core", "sh", "-ceu", "test -r /home/thoth/.pi/agent/auth.json"},
@@ -181,6 +229,26 @@ func Doctor(ctx context.Context, runner Runner) error {
return Test(ctx, runner)
}
func expectedVersions(ctx context.Context, runner Runner) (string, string, error) {
environment, err := runCompose(ctx, runner, "exec", "-T", "core", "sh", "-ceu", `printf '%s\n' "${PI_VERSION:-}"`)
if err != nil {
return "", "", commandError("Pi expected-version check", environment, err)
}
container, err := runCompose(ctx, runner, "ps", "-q", "core")
if err != nil || strings.TrimSpace(container.Stdout) == "" {
return "", "", commandError("Pi image-label check", container, err)
}
label, err := runner.Run(ctx, []string{"inspect", "--format", `{{ index .Config.Labels "io.thothii.pi.version" }}`, strings.TrimSpace(container.Stdout)}, nil)
if err != nil {
return "", "", commandError("Pi image-label check", label, err)
}
expectedValue, labelValue := strings.TrimSpace(environment.Stdout), strings.TrimSpace(label.Stdout)
if expectedValue == "" || labelValue == "" {
return "", "", errors.New("Pi image expected-version contract is empty")
}
return expectedValue, labelValue, nil
}
// Test performs the pre-Task-8 composite smoke through core's private loopback endpoint.
func Test(ctx context.Context, runner Runner) error {
if _, err := Status(ctx, runner); err != nil {
+99 -10
View File
@@ -2,6 +2,7 @@ package pi
import (
"context"
"encoding/base64"
"encoding/json"
"errors"
"io"
@@ -16,15 +17,36 @@ func TestDoctorRequiresExternalEndpointAuthPiStateAndHealth(t *testing.T) {
if err := Doctor(context.Background(), fake); err != nil {
t.Fatalf("Doctor() error = %v", err)
}
for _, command := range []string{"pi --version", "test -w /home/thoth/.pi", "test -r /home/thoth/.pi/agent/auth.json", "/health"} {
for _, command := range []string{"pi --version", "PI_VERSION", "io.thothii.pi.version", "test -w /home/thoth/.pi", "test -r /home/thoth/.pi/agent/auth.json", "/health"} {
assertCalled(t, fake.calls, command)
}
}
func TestDoctorRejectsActualEnvironmentAndImageLabelVersionMismatches(t *testing.T) {
for _, mismatch := range []string{"actual", "environment", "label"} {
t.Run(mismatch, func(t *testing.T) {
fake := newFakeRunner()
switch mismatch {
case "actual":
fake.version = "0.80.2"
case "environment":
fake.expectedVersion = "0.80.2"
case "label":
fake.labelVersion = "0.80.2"
}
if err := Doctor(context.Background(), fake); err == nil || !strings.Contains(err.Error(), "version") {
t.Fatalf("Doctor() error = %v, want expected-version mismatch", err)
}
})
}
}
func TestConfigureRestoresAndVerifiesOldSettingsAfterEveryPostSnapshotFailure(t *testing.T) {
for _, failure := range []string{"helper", "readback", "digest"} {
t.Run(failure, func(t *testing.T) {
fake := &configureRunner{failure: failure, settings: Defaults{Provider: "old", Model: "old-model", Thinking: "low"}}
old := Defaults{Provider: "old", Model: "old-model", Thinking: "low"}
raw, _ := json.Marshal(old)
fake := &configureRunner{failure: failure, settings: old, settingsExist: true, settingsRaw: raw}
err := Configure(context.Background(), fake, Defaults{Provider: "new", Model: "new-model", Thinking: "high"})
if err == nil {
t.Fatal("Configure() error = nil, want injected failure")
@@ -32,12 +54,8 @@ func TestConfigureRestoresAndVerifiesOldSettingsAfterEveryPostSnapshotFailure(t
if fake.settings != (Defaults{Provider: "old", Model: "old-model", Thinking: "low"}) {
t.Fatalf("settings after failure = %#v, want old snapshot", fake.settings)
}
minimumReads := 3
if failure == "helper" {
minimumReads = 2
}
if fake.settingsReads < minimumReads {
t.Fatalf("settings read count = %d, want capture/failure reads plus verified restore", fake.settingsReads)
if !fake.settingsExist || string(fake.settingsRaw) != string(raw) {
t.Fatalf("settings raw snapshot after failure = exists:%t raw:%q, want %q", fake.settingsExist, fake.settingsRaw, raw)
}
})
}
@@ -46,11 +64,14 @@ func TestConfigureRestoresAndVerifiesOldSettingsAfterEveryPostSnapshotFailure(t
type configureRunner struct {
failure string
settings Defaults
settingsExist bool
settingsRaw []byte
settingsReads int
configReads int
writes int
}
func (f *configureRunner) Run(_ context.Context, args []string, _ io.Reader) (compose.Result, error) {
func (f *configureRunner) Run(_ context.Context, args []string, stdin io.Reader) (compose.Result, error) {
call := strings.Join(args, " ")
switch {
case strings.Contains(call, "config --format json"):
@@ -62,21 +83,50 @@ func (f *configureRunner) Run(_ context.Context, args []string, _ io.Reader) (co
return compose.Result{Stdout: `{"services":{"core":{"image":"thothii-core:local","environment":{"THT_LLM_URL":"` + endpoint + `"}}}}`}, nil
case strings.Contains(call, "/models"):
return compose.Result{Stdout: `{"models":[{"provider":"old","id":"old-model"},{"provider":"new","id":"new-model"}]}`}, nil
case strings.Contains(call, "settings-cli.js --snapshot"):
raw := f.settingsRaw
payload := map[string]any{"exists": f.settingsExist, "rawBase64": base64.StdEncoding.EncodeToString(raw)}
contents, _ := json.Marshal(payload)
return compose.Result{Stdout: string(contents)}, nil
case strings.Contains(call, "settings-cli.js --restore"):
var payload struct {
Exists bool `json:"exists"`
RawBase64 string `json:"rawBase64"`
}
contents, _ := io.ReadAll(stdin)
if json.Unmarshal(contents, &payload) != nil {
return compose.Result{ExitCode: 2}, errors.New("invalid restore payload")
}
f.settingsExist = payload.Exists
f.settingsRaw, _ = base64.StdEncoding.DecodeString(payload.RawBase64)
f.settings = Defaults{Provider: "old", Model: "old-model", Thinking: "low"}
if payload.Exists {
_ = json.Unmarshal(f.settingsRaw, &f.settings)
}
return compose.Result{}, nil
case strings.Contains(call, "settings-cli.js"):
if strings.Contains(call, "--provider new") {
f.settings = Defaults{Provider: "new", Model: "new-model", Thinking: "high"}
f.settingsExist = true
f.settingsRaw, _ = json.MarshalIndent(f.settings, "", " ")
f.writes++
if f.failure == "helper" {
return compose.Result{ExitCode: 17}, errors.New("injected helper failure")
}
} else {
f.settings = Defaults{Provider: "old", Model: "old-model", Thinking: "low"}
f.settingsExist = true
f.settingsRaw, _ = json.Marshal(f.settings)
}
return compose.Result{}, nil
case strings.Contains(call, "/settings"):
f.settingsReads++
if f.failure == "readback" && f.settingsReads == 2 {
if f.failure == "readback" && f.settings.Provider == "new" {
return compose.Result{Stdout: `{}`}, nil
}
if !f.settingsExist {
return compose.Result{Stdout: `{"provider":"old","model":"old-model","thinking":"low"}`}, nil
}
contents, _ := json.Marshal(f.settings)
return compose.Result{Stdout: string(contents)}, nil
default:
@@ -84,6 +134,45 @@ func (f *configureRunner) Run(_ context.Context, args []string, _ io.Reader) (co
}
}
func TestConfigureAllowsAFirstRunWithoutAnExistingSettingsFile(t *testing.T) {
fake := &configureRunner{}
if err := Configure(context.Background(), fake, Defaults{Provider: "new", Model: "new-model", Thinking: "high"}); err != nil {
t.Fatalf("Configure() clean install error = %v", err)
}
if !fake.settingsExist || fake.writes != 1 || fake.settings.Provider != "new" {
t.Fatalf("clean settings = exists:%t writes:%d value:%#v", fake.settingsExist, fake.writes, fake.settings)
}
}
func TestConfigureCompensationRestoresAbsentAndExactEmptyPriorFiles(t *testing.T) {
for _, prior := range []struct {
name string
exists bool
raw []byte
}{
{name: "absent"},
{name: "empty", exists: true, raw: []byte{}},
{name: "exact raw", exists: true, raw: []byte("{\n \"workspace\": \"kept\",\n \"provider\": \"old\",\n \"model\": \"old-model\",\n \"thinking\": \"low\"\n}\n")},
} {
t.Run(prior.name, func(t *testing.T) {
fake := &configureRunner{failure: "digest", settingsExist: prior.exists, settingsRaw: append([]byte{}, prior.raw...), settings: Defaults{Provider: "old", Model: "old-model", Thinking: "low"}}
err := Configure(context.Background(), fake, Defaults{Provider: "new", Model: "new-model", Thinking: "high"})
if err == nil {
t.Fatal("Configure() error = nil, want compensated digest failure")
}
if fake.writes != 1 {
t.Fatalf("settings writes = %d, want selected values written before compensation", fake.writes)
}
if fake.settingsExist != prior.exists || string(fake.settingsRaw) != string(prior.raw) {
t.Fatalf("restored exists/raw = %t/%q, want %t/%q", fake.settingsExist, fake.settingsRaw, prior.exists, prior.raw)
}
if fake.settingsReads < 3 {
t.Fatalf("settings reads = %d, want prior effective state, requested readback, and restored default verification", fake.settingsReads)
}
})
}
}
func TestConfigureValidatesBackendModelOptionsWritesRealCoreSettingsAndUsesUpstreamIdentity(t *testing.T) {
fake := newFakeRunner()
if err := Configure(context.Background(), fake, Defaults{Provider: "provider", Model: "model", Thinking: "medium"}); err != nil {
+11 -9
View File
@@ -15,7 +15,7 @@ import (
"github.com/gofrs/flock"
)
const stateFileVersion = 3
const stateFileVersion = 4
// Phase describes the durable point reached by a Pi update.
type Phase string
@@ -24,6 +24,7 @@ const (
PhasePreflight Phase = "preflight"
PhaseBuilding Phase = "building"
PhaseRecreated Phase = "recreated"
PhasePromoting Phase = "promoting"
PhaseVerified Phase = "verified"
PhaseRolledBack Phase = "rolled_back"
PhaseFailed Phase = "failed"
@@ -59,14 +60,15 @@ type Target struct {
// State is recovery metadata stored below the installation project. It never stores environment
// values, secret paths, credentials, or command output.
type State struct {
Version int `json:"version"`
Transaction string `json:"transaction"`
Phase Phase `json:"phase"`
UpdatedAt time.Time `json:"updated_at"`
Target Target `json:"target,omitempty"`
Previous Image `json:"previous"`
Candidate Image `json:"candidate,omitempty"`
Error string `json:"error,omitempty"`
Version int `json:"version"`
Transaction string `json:"transaction"`
Phase Phase `json:"phase"`
UpdatedAt time.Time `json:"updated_at"`
Target Target `json:"target,omitempty"`
Previous Image `json:"previous"`
Candidate Image `json:"candidate,omitempty"`
MutationStarted bool `json:"mutation_started,omitempty"`
Error string `json:"error,omitempty"`
}
func readState(path string) (State, error) {
+158 -20
View File
@@ -95,10 +95,14 @@ func updateWithHooks(ctx context.Context, runner Runner, request Request, hooks
}
request.Image = canonical
}
if old, err := readState(request.StatePath); err == nil && old.Phase != PhaseVerified && old.Phase != PhaseRolledBack && old.Phase != PhaseNoop {
if old, err := readState(request.StatePath); err == nil && stateNeedsRecovery(old) {
return Result{StatePath: request.StatePath}, ErrInterruptedUpdate
} else if err != nil && !errors.Is(err, os.ErrNotExist) {
return Result{StatePath: request.StatePath}, err
} else if err == nil && !old.MutationStarted {
if cleanupErr := hooks.removeFile(lifecycleOverridePath(request.StatePath, old.Transaction)); cleanupErr != nil {
return Result{StatePath: request.StatePath}, errors.New("safe prior preparation state could not be cleaned up")
}
}
if err := setMaintenance(ctx, runner, true); err != nil {
return Result{StatePath: request.StatePath}, err
@@ -183,24 +187,30 @@ func updateWithHooks(ctx context.Context, runner Runner, request Request, hooks
lifecycle := composeOverrideRunner{Runner: runner, path: overridePath}
state.Phase = PhaseBuilding
clearMaintenance = false
if err := hooks.writeState(request.StatePath, state); err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
if err := prepareCandidate(ctx, lifecycle, request, candidateReference); err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
running, err = activeSessions(ctx, runner)
if err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
if running {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, ErrActiveSessions, hooks)
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, ErrActiveSessions, hooks)
return result, retErr
}
state.MutationStarted = true
if err := hooks.writeState(request.StatePath, state); err != nil {
state.MutationStarted = false
result, retErr, clearMaintenance = failPreparation(request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
clearMaintenance = false
if err := recreateCore(ctx, lifecycle); err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
return result, retErr
@@ -223,18 +233,48 @@ func updateWithHooks(ctx context.Context, runner Runner, request Request, hooks
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
state.Phase, state.Error = PhaseVerified, ""
state.Phase, state.Error = PhasePromoting, ""
if err := hooks.writeState(request.StatePath, state); err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
if err := hooks.removeFile(overridePath); err != nil {
return Result{Phase: PhaseFailed, StatePath: request.StatePath}, errors.New("verified update override cleanup failed: maintenance recovery required")
if err := promoteLifecycleOverride(overridePath, currentImageOverridePath(request.StatePath), candidateReference); err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
state.Phase = PhaseVerified
if err := hooks.writeState(request.StatePath, state); err != nil {
result, retErr, clearMaintenance = compensate(ctx, runner, request.StatePath, overridePath, state, err, hooks)
return result, retErr
}
clearMaintenance = true
return Result{Phase: PhaseVerified, StatePath: request.StatePath}, nil
}
func stateNeedsRecovery(state State) bool {
switch state.Phase {
case PhaseVerified, PhaseRolledBack, PhaseNoop:
return false
case PhaseFailed:
return state.MutationStarted
default:
return state.MutationStarted
}
}
func failPreparation(statePath, overridePath string, state State, cause error, hooks lifecycleHooks) (Result, error, bool) {
state.Phase = PhaseFailed
state.MutationStarted = false
state.Error = "candidate preparation failed before core mutation"
writeErr := hooks.writeState(statePath, state)
removeErr := hooks.removeFile(overridePath)
message := "candidate preparation failed before core mutation"
if writeErr != nil || removeErr != nil {
message += "; safe preparation cleanup was incomplete"
}
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("%s: %w", message, cause), true
}
// Rollback restores the image recorded in durable update state. It is safe for interrupted runs.
func Rollback(ctx context.Context, runner Runner, statePath string, confirm bool) (result Result, retErr error) {
return rollbackWithHooks(ctx, runner, statePath, confirm, defaultLifecycleHooks)
@@ -290,13 +330,13 @@ func rollbackWithHooks(ctx context.Context, runner Runner, statePath string, con
}
return Result{Phase: PhaseFailed, StatePath: statePath}, err
}
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")
}
state.Phase, state.Error = PhaseRolledBack, ""
if err := hooks.writeState(statePath, state); err != nil {
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("rollback restored the core but recovery state could not be persisted")
}
if err := hooks.removeFile(overridePath); err != nil {
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("rollback override cleanup failed: maintenance recovery required")
}
clearMaintenance = true
return Result{Phase: PhaseRolledBack, StatePath: statePath}, nil
}
@@ -323,14 +363,14 @@ 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 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
}
state.Phase, state.Error = PhaseRolledBack, ""
if writeErr := hooks.writeState(statePath, state); writeErr != nil {
return Result{Phase: PhaseFailed, StatePath: statePath}, fmt.Errorf("previous core image was restored but recovery state write failed: recovery required"), false
}
if err := hooks.removeFile(overridePath); err != nil {
return Result{Phase: PhaseFailed, StatePath: statePath}, errors.New("previous core image was restored but override cleanup failed: recovery required"), false
}
return Result{Phase: PhaseRolledBack, StatePath: statePath}, fmt.Errorf("update failed; previous core image was restored"), true
return Result{Phase: PhaseRolledBack, StatePath: statePath}, fmt.Errorf("update failed; previous core image was restored: %w", cause), true
}
func sourceValue(request Request) string {
@@ -600,6 +640,10 @@ func lifecycleOverridePath(statePath, transaction string) string {
return filepath.Join(filepath.Dir(statePath), "pi-lifecycle-"+transaction+".yaml")
}
func currentImageOverridePath(statePath string) string {
return filepath.Join(filepath.Dir(statePath), "current-image.yaml")
}
func writeLifecycleOverride(path, image string) error {
quoted, err := json.Marshal(image)
if err != nil {
@@ -612,6 +656,41 @@ func writeLifecycleOverride(path, image string) error {
return nil
}
func promoteLifecycleOverride(source, destination, expectedImage string) error {
if err := durableReplace(source, destination, filepath.Dir(destination)); err != nil {
selected, readErr := readLifecycleOverride(destination)
if readErr == nil && selected == expectedImage {
return nil
}
return errors.New("lifecycle image override could not be promoted durably")
}
selected, err := readLifecycleOverride(destination)
if err != nil || selected != expectedImage {
return errors.New("promoted lifecycle image override could not be verified")
}
return nil
}
func readLifecycleOverride(path string) (string, error) {
contents, err := os.ReadFile(path)
if err != nil {
return "", err
}
for _, line := range strings.Split(string(contents), "\n") {
line = strings.TrimSpace(line)
if !strings.HasPrefix(line, "image:") {
continue
}
encoded := strings.TrimSpace(strings.TrimPrefix(line, "image:"))
var image string
if json.Unmarshal([]byte(encoded), &image) != nil || image == "" || strings.ContainsAny(image, "\r\n") {
return "", errors.New("lifecycle image override is invalid")
}
return image, nil
}
return "", errors.New("lifecycle image override has no core image")
}
// RecoverMaintenance clears a stale durable gate only after the running core and terminal
// recovery metadata prove that no rollback is still required.
func RecoverMaintenance(ctx context.Context, runner Runner, statePath string, confirm bool) error {
@@ -625,11 +704,26 @@ func RecoverMaintenance(ctx context.Context, runner Runner, statePath string, co
defer lock.Release()
state, stateErr := readState(statePath)
if stateErr == nil {
if state.Phase != PhaseVerified && state.Phase != PhaseRolledBack && state.Phase != PhaseNoop {
transactionOverride := lifecycleOverridePath(statePath, state.Transaction)
switch {
case state.Phase == PhasePromoting:
if err := recoverPromotion(ctx, runner, statePath, transactionOverride, &state); err != nil {
return err
}
case !state.MutationStarted && state.Phase != PhaseVerified && state.Phase != PhaseRolledBack && state.Phase != PhaseNoop:
state.Phase, state.Error = PhaseFailed, "candidate preparation interrupted before core mutation"
if err := writeState(statePath, state); err != nil {
return errors.New("maintenance recovery could not finalize safe preparation state")
}
if err := durableRemove(transactionOverride); err != nil {
return errors.New("maintenance recovery could not remove the safe preparation override")
}
case stateNeedsRecovery(state):
return ErrInterruptedUpdate
}
if err := durableRemove(lifecycleOverridePath(statePath, state.Transaction)); err != nil {
return errors.New("maintenance recovery could not remove the lifecycle override")
default:
if err := durableRemove(transactionOverride); err != nil {
return errors.New("maintenance recovery could not remove the lifecycle override")
}
}
} else if !errors.Is(stateErr, os.ErrNotExist) {
return stateErr
@@ -647,6 +741,50 @@ func RecoverMaintenance(ctx context.Context, runner Runner, statePath string, co
return setMaintenance(ctx, runner, false)
}
func recoverPromotion(ctx context.Context, runner Runner, statePath, transactionOverride string, state *State) error {
currentOverride := currentImageOverridePath(statePath)
selected, currentErr := readLifecycleOverride(currentOverride)
if currentErr != nil || selected != state.Candidate.Reference {
pending, pendingErr := readLifecycleOverride(transactionOverride)
if pendingErr != nil || pending != state.Candidate.Reference {
if currentErr == nil && selected == state.Previous.Reference {
if err := verifyRestoredCurrent(ctx, runner, state.Previous); err != nil {
return ErrInterruptedUpdate
}
state.Phase, state.Error = PhaseRolledBack, ""
return writeState(statePath, *state)
}
return ErrInterruptedUpdate
}
if err := promoteLifecycleOverride(transactionOverride, currentOverride, state.Candidate.Reference); err != nil {
return err
}
}
if err := verifyCandidate(ctx, runner, state.Target.Version, state.Previous); err != nil {
return err
}
state.Phase, state.Error = PhaseVerified, ""
return writeState(statePath, *state)
}
func verifyRestoredCurrent(ctx context.Context, runner Runner, previous Image) error {
configured, err := renderedCore(ctx, runner)
if err != nil {
return err
}
after, err := runningImage(ctx, runner, configured.Reference)
if err != nil {
return err
}
if after.ID != previous.ID || configured.ConfigurationSHA != previous.ConfigurationSHA || !sameMounts(previous.Mounts, after.Mounts) {
return errors.New("running core does not match the durable previous-image selector")
}
if err := Doctor(ctx, runner); err != nil {
return err
}
return Test(ctx, runner)
}
func sameStrings(left, right []string) bool {
left, right = append([]string(nil), left...), append([]string(nil), right...)
sort.Strings(left)
+200 -21
View File
@@ -54,6 +54,9 @@ func TestUpdateBuildsPinnedVersionRecreatesOnlyCoreAndPersistsRecoveryState(t *t
if got := string(readStateBytes(t, result.StatePath)); strings.Contains(got, "llm.example.invalid") {
t.Fatalf("state = %q, want an endpoint-free configuration digest", got)
}
if selected := readSelectorReference(t, currentImageOverridePath(result.StatePath)); selected != fake.buildReference {
t.Fatalf("durable selector = %q, want verified candidate %q", selected, fake.buildReference)
}
}
func TestUpdateUsesATransactionScopedComposeOverrideWithoutMutatingTheConfiguredImage(t *testing.T) {
@@ -69,16 +72,46 @@ func TestUpdateUsesATransactionScopedComposeOverrideWithoutMutatingTheConfigured
if matches, err := filepath.Glob(filepath.Join(filepath.Dir(statePath), "pi-lifecycle-*.yaml")); err != nil || len(matches) != 0 {
t.Fatalf("terminal lifecycle overrides = %v, error = %v; want none", matches, err)
}
if _, err := os.Stat(currentImageOverridePath(statePath)); err != nil {
t.Fatalf("durable current-image override missing: %v", err)
}
}
func TestSuccessfulUpdateAndRollbackRemainSelectedOnFreshRecreate(t *testing.T) {
fake := newFakeRunner()
statePath := filepath.Join(t.TempDir(), ".thothctl", "update-state.json")
if _, err := Update(context.Background(), fake, Request{StatePath: statePath, Version: "0.81.0", Source: BuildSource, Confirm: true}); err != nil {
t.Fatal(err)
}
fake.currentImage = "sha256:old"
if err := recreateCore(context.Background(), composeOverrideRunner{Runner: fake, path: currentImageOverridePath(statePath)}); err != nil {
t.Fatal(err)
}
if fake.currentImage != "sha256:candidate" {
t.Fatalf("fresh recreate image = %q, want verified candidate", fake.currentImage)
}
if _, err := Rollback(context.Background(), fake, statePath, true); err != nil {
t.Fatal(err)
}
fake.currentImage = "sha256:candidate"
if err := recreateCore(context.Background(), composeOverrideRunner{Runner: fake, path: currentImageOverridePath(statePath)}); err != nil {
t.Fatal(err)
}
if fake.currentImage != "sha256:old" {
t.Fatalf("fresh recreate after rollback image = %q, want previous image", fake.currentImage)
}
}
func TestTwoInstallationsSharingAConfiguredTagUseDifferentLifecycleTags(t *testing.T) {
first, second := newFakeRunner(), newFakeRunner()
firstPath := filepath.Join(t.TempDir(), "one", "state.json")
secondPath := filepath.Join(t.TempDir(), "two", "state.json")
for _, item := range []struct {
fake *fakeRunner
path string
}{
{first, filepath.Join(t.TempDir(), "one", "state.json")},
{second, filepath.Join(t.TempDir(), "two", "state.json")},
{first, firstPath},
{second, secondPath},
} {
if _, err := Update(context.Background(), item.fake, Request{StatePath: item.path, Version: "0.81.0", Source: BuildSource, Confirm: true}); err != nil {
t.Fatal(err)
@@ -87,15 +120,28 @@ func TestTwoInstallationsSharingAConfiguredTagUseDifferentLifecycleTags(t *testi
if first.buildReference == second.buildReference {
t.Fatalf("installations reused lifecycle tag %q", first.buildReference)
}
firstSelector := readSelectorReference(t, currentImageOverridePath(firstPath))
secondSelector := readSelectorReference(t, currentImageOverridePath(secondPath))
if firstSelector == secondSelector || firstSelector != first.buildReference || secondSelector != second.buildReference {
t.Fatalf("installation selectors = %q / %q, want isolated lifecycle references", firstSelector, secondSelector)
}
}
func TestDigestPinnedConfiguredImageIsNeverUsedAsARollbackTagTarget(t *testing.T) {
fake := newFakeRunner()
fake.configuredImage = "registry.example.invalid/core@sha256:" + strings.Repeat("b", 64)
fake.tags = map[string]string{fake.configuredImage: "sha256:old"}
fake.fail = "health"
_, _ = Update(context.Background(), fake, Request{StatePath: filepath.Join(t.TempDir(), "state.json"), Version: "0.81.0", Source: BuildSource, Confirm: true})
statePath := filepath.Join(t.TempDir(), ".thothctl", "state.json")
if _, err := Update(context.Background(), fake, Request{StatePath: statePath, Version: "0.81.0", Source: BuildSource, Confirm: true}); err != nil {
t.Fatal(err)
}
if _, err := Rollback(context.Background(), fake, statePath, true); err != nil {
t.Fatal(err)
}
assertNotCalled(t, fake.calls, "image tag sha256:old "+fake.configuredImage)
if selected := readSelectorReference(t, currentImageOverridePath(statePath)); !strings.Contains(selected, "-previous") {
t.Fatalf("rollback selector = %q, want transaction previous tag for digest-pinned base", selected)
}
}
func TestMaintenanceLostResponsesAreResolvedByStatusAndEveryRecreateStartsGated(t *testing.T) {
@@ -189,7 +235,7 @@ func TestUpdateRollsBackAfterPostRecreateFailures(t *testing.T) {
}
func TestEveryRecoveryStateWriteFailureIsHandledTransactionally(t *testing.T) {
for failAt := 1; failAt <= 4; failAt++ {
for failAt := 1; failAt <= 6; failAt++ {
t.Run(fmt.Sprintf("write-%d", failAt), func(t *testing.T) {
fake := newFakeRunner()
writes := 0
@@ -210,8 +256,14 @@ func TestEveryRecoveryStateWriteFailureIsHandledTransactionally(t *testing.T) {
if fake.currentImage != "sha256:old" {
t.Fatalf("current image = %q, want restored previous", fake.currentImage)
}
if failAt > 1 && result.Phase != PhaseRolledBack {
t.Fatalf("phase = %q, want rolled_back", result.Phase)
wantPhase := Phase("")
if failAt == 2 || failAt == 3 {
wantPhase = PhaseFailed
} else if failAt >= 4 {
wantPhase = PhaseRolledBack
}
if result.Phase != wantPhase {
t.Fatalf("phase = %q, want %q for write %d", result.Phase, wantPhase, failAt)
}
if fake.maintenance {
t.Fatal("maintenance remained active after proven stable recovery")
@@ -227,7 +279,7 @@ func TestCompensationWriteFailureKeepsMaintenanceActiveForExplicitRecovery(t *te
hooks := defaultLifecycleHooks
hooks.writeState = func(path string, state State) error {
writes++
if writes == 4 {
if writes == 5 {
return errors.New("injected compensation state write failure")
}
return writeState(path, state)
@@ -264,12 +316,15 @@ func TestRecoverMaintenanceClearsOnlyAfterTerminalStateAndVerifiedSmoke(t *testi
fake.maintenance = true
statePath := filepath.Join(t.TempDir(), "state.json")
previous := stateImageForTest(t, fake)
state := State{Transaction: "recover-test", Phase: PhaseVerified, Previous: previous}
state := State{Transaction: "recover-test", Phase: PhaseRolledBack, MutationStarted: true, Previous: previous}
writeStateForTest(t, statePath, state)
overridePath := lifecycleOverridePath(statePath, state.Transaction)
if err := writeLifecycleOverride(overridePath, previous.Reference); err != nil {
t.Fatal(err)
}
if err := writeLifecycleOverride(currentImageOverridePath(statePath), previous.Reference); err != nil {
t.Fatal(err)
}
if err := RecoverMaintenance(context.Background(), fake, statePath, true); err != nil {
t.Fatalf("RecoverMaintenance() error = %v", err)
@@ -280,6 +335,9 @@ func TestRecoverMaintenanceClearsOnlyAfterTerminalStateAndVerifiedSmoke(t *testi
if _, err := os.Stat(overridePath); !errors.Is(err, os.ErrNotExist) {
t.Fatalf("lifecycle override still exists: %v", err)
}
if selected := readSelectorReference(t, currentImageOverridePath(statePath)); selected != previous.Reference {
t.Fatalf("maintenance cleanup changed durable selector to %q", selected)
}
assertCalled(t, fake.calls, "/models")
assertCalled(t, fake.calls, "/settings")
}
@@ -288,7 +346,7 @@ func TestRecoverMaintenanceRefusesPendingTransaction(t *testing.T) {
fake := newFakeRunner()
fake.maintenance = true
statePath := filepath.Join(t.TempDir(), "state.json")
writeStateForTest(t, statePath, State{Phase: PhaseRecreated, Previous: stateImageForTest(t, fake)})
writeStateForTest(t, statePath, State{Phase: PhaseRecreated, MutationStarted: true, Previous: stateImageForTest(t, fake)})
err := RecoverMaintenance(context.Background(), fake, statePath, true)
if !errors.Is(err, ErrInterruptedUpdate) {
@@ -299,6 +357,43 @@ func TestRecoverMaintenanceRefusesPendingTransaction(t *testing.T) {
}
}
func TestRecoverMaintenanceCompletesAnInterruptedDurablePromotion(t *testing.T) {
fake := newFakeRunner()
fake.maintenance = true
fake.currentImage = "sha256:candidate"
fake.version = "0.81.0"
fake.expectedVersion = "0.81.0"
fake.labelVersion = "0.81.0"
statePath := filepath.Join(t.TempDir(), ".thothctl", "update-state.json")
previous := stateImageForTest(t, newFakeRunner())
candidate := previous
candidate.ID = "sha256:candidate"
candidate.Reference = "thothii-core:thothctl-recover-candidate"
fake.tags[candidate.Reference] = candidate.ID
state := State{
Transaction: "promotion-recovery",
Phase: PhasePromoting,
MutationStarted: true,
Target: Target{Version: "0.81.0", Source: string(BuildSource)},
Previous: previous,
Candidate: candidate,
}
writeStateForTest(t, statePath, state)
if err := writeLifecycleOverride(lifecycleOverridePath(statePath, state.Transaction), candidate.Reference); err != nil {
t.Fatal(err)
}
if err := RecoverMaintenance(context.Background(), fake, statePath, true); err != nil {
t.Fatalf("RecoverMaintenance() promotion error = %v", err)
}
if selected := readSelectorReference(t, currentImageOverridePath(statePath)); selected != candidate.Reference {
t.Fatalf("recovered selector = %q, want %q", selected, candidate.Reference)
}
if recovered, err := readState(statePath); err != nil || recovered.Phase != PhaseVerified {
t.Fatalf("recovered state = %+v, %v; want verified", recovered, err)
}
}
func TestRollbackFinalStateWriteFailureKeepsMaintenanceAndOverrideForRecovery(t *testing.T) {
fake := newFakeRunner()
statePath := filepath.Join(t.TempDir(), "state.json")
@@ -320,13 +415,13 @@ func TestRollbackFinalStateWriteFailureKeepsMaintenanceAndOverrideForRecovery(t
if !fake.maintenance {
t.Fatal("maintenance was cleared without durable rollback finalization")
}
if _, err := os.Stat(lifecycleOverridePath(statePath, "rollback-test")); err != nil {
t.Fatalf("recovery override was not preserved: %v", err)
if selected := readSelectorReference(t, currentImageOverridePath(statePath)); selected != previous.Reference {
t.Fatalf("durable rollback selector = %q, want %q", selected, previous.Reference)
}
}
func TestUpdateDoesNotRecreateWhenPreflightOrBuildFails(t *testing.T) {
for _, failure := range []string{"preflight", "build"} {
func TestUpdateDoesNotRecreateWhenPreflightFails(t *testing.T) {
for _, failure := range []string{"preflight"} {
t.Run(failure, func(t *testing.T) {
fake := newFakeRunner()
fake.fail = failure
@@ -337,16 +432,67 @@ func TestUpdateDoesNotRecreateWhenPreflightOrBuildFails(t *testing.T) {
if failure == "preflight" && result.Phase == PhaseRolledBack {
t.Fatalf("preflight failure unexpectedly rolled back: %+v", result)
}
if failure == "build" && result.Phase != PhaseRolledBack {
t.Fatalf("candidate build failure must compensate: %+v", result)
assertNotCalled(t, fake.calls, "force-recreate")
})
}
}
func TestCandidateBuildAndPullFailuresRemainPreMutationAndNeverRecreateCore(t *testing.T) {
for _, testCase := range []struct {
name string
source Source
image string
failure string
}{
{name: "build", source: BuildSource, failure: "build"},
{name: "pull", source: PullSource, image: "registry.example.invalid/core@sha256:" + strings.Repeat("a", 64), failure: "pull"},
{name: "candidate tag", source: PullSource, image: "registry.example.invalid/core@sha256:" + strings.Repeat("b", 64), failure: "tag"},
} {
t.Run(testCase.name, func(t *testing.T) {
fake := newFakeRunner()
fake.fail = testCase.failure
statePath := filepath.Join(t.TempDir(), ".thothctl", "update-state.json")
result, err := Update(context.Background(), fake, Request{StatePath: statePath, Version: "0.81.0", Source: testCase.source, Image: testCase.image, Confirm: true})
if err == nil {
t.Fatal("Update() error = nil, want preparation failure")
}
if failure == "preflight" {
assertNotCalled(t, fake.calls, "force-recreate")
if result.Phase != PhaseFailed {
t.Fatalf("phase = %q, want safe failed preparation", result.Phase)
}
state, stateErr := readState(statePath)
if stateErr != nil {
t.Fatal(stateErr)
}
if state.MutationStarted {
t.Fatal("preparation failure recorded mutationStarted")
}
assertNotCalled(t, fake.calls, "force-recreate")
if fake.maintenance {
t.Fatal("maintenance remained active after safe preparation failure")
}
})
}
}
func TestSuccessfulCompensationPreservesTheOriginalTypedCause(t *testing.T) {
for _, cause := range []error{ErrActiveSessions, ErrInterruptedUpdate} {
fake := newFakeRunner()
fake.maintenance = true
statePath := filepath.Join(t.TempDir(), ".thothctl", "update-state.json")
state := State{Transaction: "typed-cause", Phase: PhaseRecreated, MutationStarted: true, Previous: stateImageForTest(t, fake)}
result, err, clear := compensate(context.Background(), fake, statePath, lifecycleOverridePath(statePath, state.Transaction), state, cause, defaultLifecycleHooks)
if result.Phase != PhaseRolledBack || !clear {
t.Fatalf("compensation = %+v, clear=%t; want successful rollback", result, clear)
}
if !errors.Is(err, cause) {
t.Fatalf("compensation error = %v, want errors.Is(..., %v)", err, cause)
}
if !strings.Contains(err.Error(), "previous core image was restored") {
t.Fatalf("compensation error = %v, want rollback-success report", err)
}
}
}
func TestUpdateRequiresConfirmationAndDrainsActiveSessions(t *testing.T) {
fake := newFakeRunner()
_, err := Update(context.Background(), fake, Request{StatePath: filepath.Join(t.TempDir(), "state.json"), Version: "0.81.0", Source: BuildSource})
@@ -402,7 +548,7 @@ func TestRollbackRestoresInterruptedOrPreviouslyRecordedState(t *testing.T) {
func TestUpdateRefusesToOverwriteInterruptedRecoveryState(t *testing.T) {
fake := newFakeRunner()
statePath := filepath.Join(t.TempDir(), "state.json")
writeStateForTest(t, statePath, State{Phase: PhaseRecreated, Previous: Image{ID: "sha256:old", Reference: "thothii-core:local", MountFingerprint: mountFingerprint(nil)}})
writeStateForTest(t, statePath, State{Phase: PhaseRecreated, MutationStarted: true, Previous: Image{ID: "sha256:old", Reference: "thothii-core:local", MountFingerprint: mountFingerprint(nil)}})
_, err := Update(context.Background(), fake, Request{StatePath: statePath, Version: "0.81.0", Source: BuildSource, Confirm: true})
if !errors.Is(err, ErrInterruptedUpdate) {
t.Fatalf("Update() error = %v, want interrupted update error", err)
@@ -441,6 +587,8 @@ type fakeRunner struct {
calls []string
fail string
version string
expectedVersion string
labelVersion string
activeSessions bool
built bool
currentImage string
@@ -460,7 +608,7 @@ type fakeRunner struct {
func newFakeRunner() *fakeRunner {
return &fakeRunner{
version: "0.80.3", currentImage: "sha256:old", configuredImage: "thothii-core:local",
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"},
}
@@ -478,6 +626,12 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose
if f.fail == "build" && containsArg(args, "build") {
return compose.Result{ExitCode: 1}, errors.New("build token=secret")
}
if f.fail == "pull" && len(args) > 0 && args[0] == "pull" {
return compose.Result{ExitCode: 1}, errors.New("pull token=secret")
}
if f.fail == "tag" && len(args) >= 4 && args[0] == "image" && args[1] == "tag" && strings.Contains(args[3], "-candidate") {
return compose.Result{ExitCode: 1}, errors.New("tag token=secret")
}
if f.fail == "health" && f.built && strings.Contains(call, "curl -fsS http://127.0.0.1:8787/health") {
return compose.Result{ExitCode: 1}, errors.New("health token=secret")
}
@@ -501,6 +655,8 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose
return compose.Result{Stdout: "core-container\n"}, nil
case strings.Contains(call, "inspect --format {{.Image}}"):
return compose.Result{Stdout: f.currentImage + "\n"}, nil
case strings.Contains(call, "io.thothii.pi.version"):
return compose.Result{Stdout: f.labelVersion + "\n"}, nil
case strings.Contains(call, "inspect --format {{json .Mounts}}"):
if f.fail == "mount-drift" && f.currentImage == "sha256:candidate" {
return compose.Result{Stdout: `[{"Type":"volume","Name":"wrong-settings","Source":"wrong-settings","Destination":"/data/settings","RW":true}]`}, nil
@@ -580,6 +736,8 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose
return compose.Result{}, nil
case strings.Contains(call, "pi --version"):
return compose.Result{Stdout: f.version + "\n"}, nil
case strings.Contains(call, "PI_VERSION"):
return compose.Result{Stdout: f.expectedVersion + "\n"}, nil
case strings.Contains(call, "/models"):
if f.modelsWire != "" {
return compose.Result{Stdout: f.modelsWire}, nil
@@ -604,7 +762,7 @@ func containsArg(args []string, wanted string) bool {
func selectedCoreReference(args []string, fallback string) string {
for index := 0; index+1 < len(args); index++ {
if args[index] != "-f" || !strings.Contains(filepath.Base(args[index+1]), "pi-lifecycle-") {
if args[index] != "-f" || (!strings.Contains(filepath.Base(args[index+1]), "pi-lifecycle-") && filepath.Base(args[index+1]) != "current-image.yaml") {
continue
}
contents, err := os.ReadFile(args[index+1])
@@ -626,6 +784,27 @@ func selectedCoreReference(args []string, fallback string) string {
return fallback
}
func readSelectorReference(t *testing.T, path string) string {
t.Helper()
contents, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
for _, line := range strings.Split(string(contents), "\n") {
line = strings.TrimSpace(line)
if !strings.HasPrefix(line, "image:") {
continue
}
value := strings.TrimSpace(strings.TrimPrefix(line, "image:"))
if decoded, err := strconv.Unquote(value); err == nil {
return decoded
}
return value
}
t.Fatalf("selector %s has no image", path)
return ""
}
func callIndex(calls []string, contains string) int {
for index, call := range calls {
if strings.Contains(call, contains) {