diff --git a/backend/src/auth/auth.ts b/backend/src/auth/auth.ts index b2633022..e73008ac 100644 --- a/backend/src/auth/auth.ts +++ b/backend/src/auth/auth.ts @@ -82,12 +82,6 @@ export function authenticateSession(deps: AuthDependencies): preHandlerHookHandl const snapshot = captureAuthConfigSnapshot(request, deps.authentication); if (isPublicRoute(request)) return; - const operator = loopbackMaintenancePrincipal(request); - if (operator) { - request.principal = operator; - return; - } - if (legacy) { await legacy(request, reply); if (reply.sent || !STATE_CHANGING_METHODS.has(request.method)) return; @@ -144,15 +138,6 @@ export function authenticateSession(deps: AuthDependencies): preHandlerHookHandl }; } -function loopbackMaintenancePrincipal(request: FastifyRequest): PrincipalContext | undefined { - if (request.ip !== "127.0.0.1" && request.ip !== "::1" && request.ip !== "::ffff:127.0.0.1") return undefined; - if (singleHeader(request.headers["x-thoth-principal-issuer"]) !== "tht" - || singleHeader(request.headers["x-thoth-principal-subject"]) !== "tht-maintenance" - || singleHeader(request.headers["x-thoth-principal-display-name"]) !== "Tht maintenance" - || singleHeader(request.headers["x-thoth-is-admin"]) !== "1") return undefined; - return upstreamPrincipal(request.headers); -} - export function requireCsrf(request: FastifyRequest, reply: FastifyReply): true | FastifyReply { const expectedOrigin = request.authPublicOrigin; const token = request.authSessionToken; diff --git a/backend/src/operator-command.ts b/backend/src/operator-command.ts new file mode 100644 index 00000000..326ca2fd --- /dev/null +++ b/backend/src/operator-command.ts @@ -0,0 +1,123 @@ +/** + * Installation-scoped host operations that must not impersonate an HTTP administrator. + * This command runs only through `docker compose exec core`; it never accepts credentials, + * headers, paths, or arbitrary code from the caller. + */ +import { pathToFileURL } from "node:url"; +import { join } from "node:path"; +import { loadConfig, type AppConfig } from "./config.js"; +import { rolesToPermissions } from "./auth/config.js"; +import type { PrincipalContext } from "./auth/principal.js"; +import { createPiModelLister } from "./pi/list-models.js"; +import { createPiManagement } from "./pi/management.js"; +import { effectiveSettings } from "./routes/settings.js"; +import { MaintenanceBarrier } from "./runtime/maintenance-gate.js"; +import { loadSettings } from "./settings/settings-store.js"; +import { ThtRunner, type SessionRow } from "./tht/tht-runner.js"; +import { WorkspaceRegistry } from "./workspaces/registry.js"; +import { WorkspaceSecretStore } from "./workspaces/secret-store.js"; + +type OperatorAction = "maintenance-activate" | "maintenance-deactivate" | "maintenance-status" + | "session-inventory" | "workflow-doctor" | "pi-options" | "pi-test" | "effective-settings"; + +const lifecyclePrincipal: PrincipalContext = { + issuer: "tht-operator-command", + subject: "installation-lifecycle", + displayName: "Installation lifecycle", + roles: ["admin"], + permissions: rolesToPermissions(["admin"]), + isAdmin: true, +}; + +function operatorRunner(config: AppConfig): ThtRunner { + const workspaceSecretStore = new WorkspaceSecretStore({ + root: config.workspaceSecretStoreRoot, + runtimeRoot: config.workspaceSecretRuntimeRoot, + installationId: config.workspaceRegistry.installationId, + }); + return new ThtRunner({ + thtBin: config.thtBin, + harnessDir: config.harnessDir, + configPath: process.env.THT_CONFIG ?? "config/tht.yaml", + dataRoot: config.dataRoot, + runtimeSnapshotRoot: join(config.workspaceRegistry.root, "snapshots", "runtime"), + secretRoots: config.workspaceRegistry.secretRoots, + secretsFile: config.secretsFile, + secretFiles: config.secretFiles, + workspaceSecretStore, + semanticRuntime: { + internalQdrantUrl: config.internalQdrantUrl, + internalEmbeddingUrl: config.internalEmbeddingUrl, + internalEmbeddingModel: config.internalEmbeddingModel, + internalEmbeddingDimensions: config.internalEmbeddingDimensions, + }, + }).withPrincipal(lifecyclePrincipal); +} + +async function sessionInventory(config: AppConfig): Promise>> { + const registry = new WorkspaceRegistry(config.workspaceRegistry); + const revisions = await registry.listRetainedSnapshots(); + const runner = operatorRunner(config); + const sessions = new Map(); + for (const revision of revisions) { + for (const session of await runner.sessionList(revision.snapshotPath)) sessions.set(session.id, session); + } + return [...sessions.values()].map(({ status, archived }) => ({ status, archived: archived === true })); +} + +async function workflowDiagnostics(config: AppConfig): Promise<{ ready: true; workspaces: number }> { + const registry = new WorkspaceRegistry(config.workspaceRegistry); + const revisions = await registry.listRetainedSnapshots(); + if (revisions.length === 0) throw new Error("workflow diagnostics unavailable"); + const runner = operatorRunner(config); + for (const revision of revisions) { + const result = await runner.run(["doctor", "--json"], revision.snapshotPath); + let payload: unknown; + try { + payload = JSON.parse(result.stdout); + } catch { + throw new Error("workflow diagnostics failed"); + } + if ( + result.code !== 0 || !payload || typeof payload !== "object" + || (payload as { ok?: unknown }).ok !== true + ) throw new Error("workflow diagnostics failed"); + } + return { ready: true, workspaces: revisions.length }; +} + +export async function runOperatorAction( + action: OperatorAction, + config: AppConfig, +): Promise { + if (action.startsWith("maintenance-")) { + const barrier = new MaintenanceBarrier(config.maintenanceFile); + if (action === "maintenance-activate") await barrier.activate(); + if (action === "maintenance-deactivate") barrier.deactivate(); + return barrier.status(); + } + if (action === "session-inventory") return await sessionInventory(config); + if (action === "workflow-doctor") return await workflowDiagnostics(config); + if (action === "effective-settings") return effectiveSettings(config, loadSettings(config)); + const service = createPiManagement(config, { listModels: createPiModelLister(config) }); + if (action === "pi-options") return await service.options(); + if (action === "pi-test") return await service.test(); + throw new Error("unsupported operator action"); +} + +async function main(): Promise { + const action = process.argv[2] as OperatorAction | undefined; + if (!action || ![ + "maintenance-activate", "maintenance-deactivate", "maintenance-status", "session-inventory", + "workflow-doctor", "pi-options", "pi-test", "effective-settings", + ].includes(action)) throw new Error("invalid operator action"); + const result = await runOperatorAction(action, loadConfig(process.env)); + process.stdout.write(`${JSON.stringify(result)}\n`); +} + +if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) { + void main().catch(() => { + process.stderr.write("operator command failed\n"); + process.exitCode = 2; + }); +} diff --git a/backend/src/server.ts b/backend/src/server.ts index 06bf5f7f..726aedb3 100644 --- a/backend/src/server.ts +++ b/backend/src/server.ts @@ -1,10 +1,10 @@ import { buildApp, type AppWithAuthSessionStore } from "./app.js"; import { loadConfig } from "./config.js"; import { formatStartupFailure } from "./startup-error.js"; -const config = loadConfig(process.env); -const app = buildApp(config) as AppWithAuthSessionStore; async function start(): Promise { + const config = loadConfig(process.env); + const app = buildApp(config) as AppWithAuthSessionStore; const sessions = app.thothiiAuthSessionStore; if (sessions) { await sessions.prune(); diff --git a/backend/src/startup-error.ts b/backend/src/startup-error.ts index d04ea6ce..feb49ed8 100644 --- a/backend/src/startup-error.ts +++ b/backend/src/startup-error.ts @@ -4,9 +4,15 @@ const STARTUP_CAUSES = new Set([ "workspace_registry_invalid", ]); +const EXACT_CAUSES = new Map([ + ["authentication configuration is invalid", "auth_config_invalid"], + ["authentication session state is invalid", "auth_session_store_invalid"], + ["workspace registry configuration is invalid", "workspace_registry_invalid"], +]); + /** Return one bounded machine cause; never include the original error text or stack. */ export function formatStartupFailure(error: unknown): string { const message = error instanceof Error ? error.message : ""; - const cause = STARTUP_CAUSES.has(message) ? message : "startup_unknown"; + const cause = STARTUP_CAUSES.has(message) ? message : EXACT_CAUSES.get(message) ?? "startup_unknown"; return `backend startup failed: ${cause}`; } diff --git a/backend/test/auth.test.ts b/backend/test/auth.test.ts index 3bd5332c..531d0243 100644 --- a/backend/test/auth.test.ts +++ b/backend/test/auth.test.ts @@ -169,7 +169,7 @@ test("the session boundary exposes only exact health and authentication protocol expect((await app.inject({ method: "GET", url: "/auth/configured" })).statusCode).toBe(401); }); -test("the session boundary retains the exact loopback tht maintenance identity in configured auth modes", async () => { +test("loopback maintenance headers can never mint an administrator in configured auth modes", async () => { const app = Fastify(); app.addHook("preHandler", authenticateSession({ mode: "local" })); app.get("/private", async (request) => getPrincipal(request)); @@ -183,8 +183,8 @@ test("the session boundary retains the exact loopback tht maintenance identity i for (const method of ["GET", "POST"] as const) { const response = await app.inject({ method, url: "/private", headers, remoteAddress: "127.0.0.1" }); - expect(response.statusCode).toBe(200); - expect(response.json()).toMatchObject({ issuer: "tht", subject: "tht-maintenance", isAdmin: true }); + expect(response.statusCode).toBe(503); + expect(response.body).not.toContain("tht-maintenance"); } }); diff --git a/backend/test/startup-error.test.ts b/backend/test/startup-error.test.ts index bf1685bd..8faed34c 100644 --- a/backend/test/startup-error.test.ts +++ b/backend/test/startup-error.test.ts @@ -1,3 +1,7 @@ +import { spawnSync } from "node:child_process"; +import { chmodSync, mkdirSync, mkdtempSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join, resolve } from "node:path"; import { describe, expect, it } from "vitest"; import { formatStartupFailure } from "../src/startup-error.js"; @@ -25,4 +29,34 @@ describe("formatStartupFailure", () => { expect(formatted).not.toContain(leaked); } }); + + it("sanitizes synchronous configuration failures from the real server subprocess", () => { + const root = mkdtempSync(join(tmpdir(), "thothii-startup-secret-")); + const secret = "startup-password-do-not-log"; + const authDirectory = join(root, `auth-${secret}`); + mkdirSync(authDirectory, { mode: 0o700 }); + const authFile = join(authDirectory, "auth.yaml"); + writeFileSync(authFile, `version: 1\nmode: local\npassword: ${secret}\n`, { mode: 0o600 }); + chmodSync(authDirectory, 0o700); + const entrypoint = resolve(process.cwd(), "src/server.ts"); + + const result = spawnSync(process.execPath, ["--import", "tsx", entrypoint], { + cwd: process.cwd(), + encoding: "utf8", + timeout: 15_000, + env: { + ...process.env, + NODE_ENV: "test", + THT_AUTH_CONFIG_FILE: authFile, + THT_AUTH_STATE_ROOT: join(root, "auth-state"), + }, + }); + + expect(result.status).toBe(1); + expect(result.stdout).toBe(""); + expect(result.stderr.trim()).toBe("backend startup failed: auth_config_invalid"); + expect(`${result.stdout}${result.stderr}`).not.toContain(secret); + expect(`${result.stdout}${result.stderr}`).not.toContain(authDirectory); + expect(`${result.stdout}${result.stderr}`).not.toContain("server.ts"); + }); }); diff --git a/scripts/unified-deployment-smoke.sh b/scripts/unified-deployment-smoke.sh index 8370e1a2..93e14e54 100755 --- a/scripts/unified-deployment-smoke.sh +++ b/scripts/unified-deployment-smoke.sh @@ -226,7 +226,7 @@ task13_write_fixture_files() { printf '%s' "task13-runtime-password-$TASK13_RUN_ID" >"$TASK13_SESSION_RUNTIME_PASSWORD" printf '%s' "$TASK13_AUTH_PASSWORD" >"$TASK13_AUTH_PASSWORD_FILE" mkdir -p "$TASK13_AUTH_ROOT" - chmod 0644 "$TASK13_PI_AUTH" + chmod 0600 "$TASK13_PI_AUTH" chmod 0700 "$TASK13_AUTH_ROOT" chmod 0600 "$TASK13_SECRETS" "$TASK13_SESSION_RUNTIME_PASSWORD" "$TASK13_AUTH_PASSWORD_FILE" @@ -417,7 +417,7 @@ task13_write_server_fixture_files() { VEFTSzEzLURJU1BPU0FCTEUtU0VTU0lPTi1DQQ== -----END CERTIFICATE----- EOF - chmod 0644 "$TASK13_PI_AUTH" "$TASK13_SESSION_CA" + chmod 0600 "$TASK13_PI_AUTH" "$TASK13_SESSION_CA" chmod 0600 "$TASK13_SECRETS" "$TASK13_SESSION_RUNTIME_PASSWORD" \ "$TASK13_SESSION_MIGRATOR_PASSWORD_FILE" @@ -949,13 +949,23 @@ PY --fail --silent --show-error --cookie "$cookie_after" "http://$frontend/api/me")" node -e 'const value=JSON.parse(process.argv[1]); if(value.issuer!=="local"||value.session?.remembered!==true) process.exit(1)' "$me" \ || task13_fail "post-backup browser session was not authenticated" - task13_compose_logged "seed pending OIDC state excluded from restore" exec -T core sh -ceu \ - 'printf %s "{}" > /data/auth/oidc/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.json && chmod 0600 /data/auth/oidc/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.json && test "$(find /data/auth/sessions -type f | wc -l | tr -d " ")" -ge 2' + task13_compose_logged "seed valid pending OIDC state excluded from restore" exec -T core node --input-type=module -e ' + const { loadConfig } = await import("/app/backend/dist/config.js"); + const { createFileAuthSessionStore } = await import("/app/backend/dist/auth/session-store.js"); + const config = loadConfig(process.env); + const revision = config.authentication.current().revision; + const store = createFileAuthSessionStore(config.authStateRoot); + await store.createOidcState({ + nonce: "n".repeat(43), codeVerifier: "v".repeat(43), returnTo: "/", + authConfigRevision: revision, issuer: "https://pending.task13.invalid", + browserTransactionDigest: "d".repeat(64), browserTransactionTransport: "loopback_http", + }); + ' + task13_compose_logged "verify pre-restore authentication runtime population" exec -T core sh -ceu \ + 'test "$(find /data/auth/sessions -type f | wc -l | tr -d " ")" -ge 2 && test -n "$(find /data/auth/oidc -type f -name "*.json" -print -quit)"' - task13_compose_logged "stop local stack for restore" stop - task13_run_logged "perform real production restore" "$TASK13_THT" \ + task13_run_logged "perform real running local production restore" "$TASK13_THT" \ --installation "$TASK13_INSTALLATION" restore "$archive" --yes - task13_compose_start_logged "start restored local stack" up --detach --wait --wait-timeout 120 status_before="$(curl --connect-timeout "$TASK13_CURL_CONNECT_TIMEOUT" --max-time "$TASK13_CURL_MAX_TIME" \ --silent --output /dev/null --write-out '%{http_code}' --cookie "$cookie_before" "http://$frontend/api/me")" status_after="$(curl --connect-timeout "$TASK13_CURL_CONNECT_TIMEOUT" --max-time "$TASK13_CURL_MAX_TIME" \ @@ -972,6 +982,49 @@ PY task13_create_admin_session } +task13_assert_server_oidc_restore_verification() { + local archive frontend status diagnostics + archive="$TASK13_TMP/server-oidc-restore-source.zip" + frontend="$(task13_frontend_address)" + task13_compose_logged "seed valid server authentication runtime excluded from restore" exec -T core node --input-type=module -e ' + const { loadConfig } = await import("/app/backend/dist/config.js"); + const { createFileAuthSessionStore } = await import("/app/backend/dist/auth/session-store.js"); + const config = loadConfig(process.env); + const revision = config.authentication.current().revision; + const store = createFileAuthSessionStore(config.authStateRoot); + await store.create({ + principal: { issuer: "https://task13-fake-oidc:9443/application/o/task13/", subject: "restore-browser", roles: ["user"], permissions: ["session.use"], isAdmin: false }, + method: "oidc", remembered: false, authConfigRevision: revision, + idleTtlMs: 60_000, absoluteTtlMs: 120_000, + }); + await store.createOidcState({ + nonce: "n".repeat(43), codeVerifier: "v".repeat(43), returnTo: "/", + authConfigRevision: revision, issuer: "https://task13-fake-oidc:9443/application/o/task13/", + browserTransactionDigest: "d".repeat(64), browserTransactionTransport: "https", + }); + ' + task13_compose_logged "stop server stack for OIDC restore" stop + task13_run_logged "create real server default-custody backup" "$TASK13_THT" \ + --installation "$TASK13_INSTALLATION" backup --output "$archive" + task13_run_logged "perform real stopped OIDC production restore verification" "$TASK13_THT" \ + --installation "$TASK13_INSTALLATION" restore "$archive" --yes + task13_compose_start_logged "start restored server stack" up --detach --wait --wait-timeout 120 core frontend + status="$(curl --connect-timeout "$TASK13_CURL_CONNECT_TIMEOUT" --max-time "$TASK13_CURL_MAX_TIME" \ + --silent --output /dev/null --write-out '%{http_code}' "http://$frontend/api/me")" + [[ "$status" == 401 ]] || task13_fail "OIDC restore did not require browser reauthentication" + task13_compose_logged "verify private empty server authentication state" exec -T core sh -ceu ' + test "$(stat -c %a /data/auth)" = 700 + test "$(stat -c %a /data/auth/sessions)" = 700 + test "$(stat -c %a /data/auth/oidc)" = 700 + test "$(stat -c %u /data/auth)" = "$(id -u)" + test -z "$(find /data/auth/sessions /data/auth/oidc -mindepth 1 -print -quit)" + ' + diagnostics="$TASK13_TMP/server-auth-diagnostics-after-restore.json" + "$TASK13_THT" --installation "$TASK13_INSTALLATION" auth check --json >"$diagnostics" + node -e 'const value=JSON.parse(require("fs").readFileSync(process.argv[1], "utf8")); if(value.mode!=="oidc"||value.ready!==true||value.checks.length!==1||value.checks[0].code!=="auth_ready") process.exit(1)' "$diagnostics" \ + || task13_fail "restored server did not retain strict fake-provider OIDC diagnostics" +} + task13_assert_runtime() { local frontend expected_pi actual_pi core_id frontend="$(task13_frontend_address)" @@ -2137,6 +2190,7 @@ task13_server_smoke_main() { task13_assert_project_ownership task13_assert_built_image_ownership task13_assert_server_runtime + task13_assert_server_oidc_restore_verification printf 'Task 13 Linux server deployment smoke passed.\n' } diff --git a/tools/tht/cmd/tht/main_test.go b/tools/tht/cmd/tht/main_test.go index c5099280..8be1bfe1 100644 --- a/tools/tht/cmd/tht/main_test.go +++ b/tools/tht/cmd/tht/main_test.go @@ -1504,16 +1504,16 @@ case " $* " in *"io.thothii.pi.version"*) printf '%s\n' "${THT_FAKE_PI_VERSION:-0.80.3}" ;; *"PI_VERSION"*) printf '%s\n' "${THT_FAKE_PI_VERSION:-0.80.3}" ;; *" pi --version "*) printf '%s\n' "${THT_FAKE_PI_VERSION:-0.80.3}" ;; - *"/pi-management/options "*) printf '%s\n' '{"providers":["provider"],"models":[{"provider":"provider","id":"model"}],"reasoning":["low","medium","high"]}' ;; + *"operator-command.js pi-options "*) printf '%s\n' '{"providers":["provider"],"models":[{"provider":"provider","id":"model"}],"reasoning":["low","medium","high"]}' ;; *"settings-cli.js --snapshot"*) printf '%s\n' '{"exists":false,"rawBase64":""}' ;; - *"/settings "*) printf '%s\n' '{"provider":"provider","model":"model","thinking":"medium"}' ;; - *"/pi-management/test "*) printf '%s\n' '{"ready":true}' ;; - *" tht doctor --json"*) printf '%s\n' '{"ok":true}' ;; + *"operator-command.js effective-settings "*) printf '%s\n' '{"provider":"provider","model":"model","thinking":"medium"}' ;; + *"operator-command.js pi-test "*) printf '%s\n' '{"ready":true}' ;; + *"operator-command.js workflow-doctor "*) printf '%s\n' '{"ready":true,"workspaces":1}' ;; *"dist/auth/diagnostic-command.js --json"*) printf '%s\n' '{"ready":true,"mode":"oidc","checks":[{"level":"info","code":"auth_ready","message":"Authentication is ready."}]}' ;; - *"/sessions?scope="*) printf '%s\n' '[]' ;; - *"/internal/maintenance/activate "*) printf '%s\n' '{"active":true,"admissions":0}' ;; - *"/internal/maintenance/deactivate "*) printf '%s\n' '{"active":false,"admissions":0}' ;; - *"/internal/maintenance/status "*) printf '%s\n' '{"active":true,"admissions":0}' ;; + *"operator-command.js session-inventory "*) printf '%s\n' '[]' ;; + *"operator-command.js maintenance-activate "*) printf '%s\n' '{"active":true,"admissions":0,"recoveryRequired":false}' ;; + *"operator-command.js maintenance-deactivate "*) printf '%s\n' '{"active":false,"admissions":0,"recoveryRequired":false}' ;; + *"operator-command.js maintenance-status "*) printf '%s\n' '{"active":true,"admissions":0,"recoveryRequired":false}' ;; *" logs "*) printf '%s\n' "$THT_FAKE_LOG" ;; *" run --rm --no-deps --no-TTY "*" workspace-maintenance "*) printf '%s\n' "$THT_FAKE_WORKSPACE_RESULT" ;; esac diff --git a/tools/tht/internal/backup/create.go b/tools/tht/internal/backup/create.go index 1d7f9ce9..4015bdd0 100644 --- a/tools/tht/internal/backup/create.go +++ b/tools/tht/internal/backup/create.go @@ -43,6 +43,15 @@ var requiredVolumes = []string{ "embedding-models", } +var requiredServerVolumes = []string{"qdrant-data", "embedding-models"} + +func requiredBackupVolumes(installation config.Installation) []string { + if installation.Profile == "server" { + return requiredServerVolumes + } + return requiredVolumes +} + const ( helperImage = "busybox:1.36.1" drainPollInterval = time.Second @@ -165,7 +174,7 @@ func createWithDependencies(ctx context.Context, installation config.Installatio if err != nil { return Result{}, err } - volumes, err := inspectRequiredVolumes(ctx, dependencies.runner, rendered) + volumes, err := inspectRequiredVolumes(ctx, installation, dependencies.runner, rendered) if err != nil { return Result{}, err } @@ -408,7 +417,7 @@ func renderedConfiguration(ctx context.Context, installation config.Installation if json.Unmarshal([]byte(result.Stdout), &rendered) != nil { return renderedCompose{}, errors.New("Docker Compose returned invalid backup configuration") } - for _, logical := range requiredVolumes { + for _, logical := range requiredBackupVolumes(installation) { if rendered.Volumes[logical].Name == "" { return renderedCompose{}, fmt.Errorf("required backup volume %q is not configured", logical) } @@ -416,9 +425,10 @@ func renderedConfiguration(ctx context.Context, installation config.Installation return rendered, nil } -func inspectRequiredVolumes(ctx context.Context, runner archiveRunner, rendered renderedCompose) ([]VolumeMetadata, error) { - names := make([]string, 0, len(requiredVolumes)) - for _, logical := range requiredVolumes { +func inspectRequiredVolumes(ctx context.Context, installation config.Installation, runner archiveRunner, rendered renderedCompose) ([]VolumeMetadata, error) { + required := requiredBackupVolumes(installation) + names := make([]string, 0, len(required)) + for _, logical := range required { names = append(names, rendered.Volumes[logical].Name) } result, err := runner.Run(ctx, append([]string{"volume", "inspect"}, names...), nil) @@ -445,8 +455,8 @@ func inspectRequiredVolumes(ctx context.Context, runner archiveRunner, rendered Labels map[string]string }{item.Name, item.Driver, item.Labels} } - volumes := make([]VolumeMetadata, 0, len(requiredVolumes)) - for _, logical := range requiredVolumes { + volumes := make([]VolumeMetadata, 0, len(required)) + for _, logical := range required { name := rendered.Volumes[logical].Name item, ok := byName[name] if !ok || item.Driver == "" { @@ -614,13 +624,12 @@ func normalizeBackupServiceStates(states []backupServiceState) ([]backupServiceS } func maintenance(ctx context.Context, installation config.Installation, runner archiveRunner, activate bool) error { - action := "deactivate" + action := "maintenance-deactivate" if activate { - action = "activate" + action = "maintenance-activate" } result, err := runner.Run(ctx, installation.ComposeArgs( - "exec", "-T", "core", "curl", "-fsS", "--max-time", "5", "-X", "POST", - "http://127.0.0.1:8787/internal/maintenance/"+action, + "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", action, ), nil) if err != nil { return dockerError("change maintenance admissions", result, err) @@ -631,8 +640,8 @@ func maintenance(ctx context.Context, installation config.Installation, runner a func waitForNoActiveSessions(ctx context.Context, installation config.Installation, runner archiveRunner, drain bool, sleep func(time.Duration)) error { for attempt := 0; attempt < maxDrainPolls; attempt++ { result, err := runner.Run(ctx, installation.ComposeArgs( - "exec", "-T", "core", "curl", "-fsS", "--max-time", "5", - "http://127.0.0.1:8787/sessions?scope="+runner.SessionInventoryScope(), + "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", + "session-inventory", runner.SessionInventoryScope(), ), nil) if err != nil { return dockerError("inspect active sessions", result, err) diff --git a/tools/tht/internal/backup/create_test.go b/tools/tht/internal/backup/create_test.go index 67957540..1364c7fa 100644 --- a/tools/tht/internal/backup/create_test.go +++ b/tools/tht/internal/backup/create_test.go @@ -501,6 +501,13 @@ func TestCreateIncludesServerPreservationRootsWithoutRecursingIntoBackupRoot(t * t.Fatal(err) } archive := readFixtureArchive(t, output) + serverVolumes := make([]string, 0, len(archive.manifest.Volumes)) + for _, volume := range archive.manifest.Volumes { + serverVolumes = append(serverVolumes, volume.LogicalName) + } + if strings.Join(serverVolumes, ",") != "embedding-models,qdrant-data" { + t.Fatalf("server backup volumes = %v, want only infrastructure named volumes", serverVolumes) + } all := bytes.Join(mapValues(archive.files), nil) for _, value := range []string{"session-state", "pi-state", "workspace-registry"} { if !bytes.Contains(all, []byte(value)) { @@ -659,7 +666,7 @@ func (runner *fakeBackupRunner) Run(_ context.Context, args []string, _ io.Reade switch { case strings.Contains(command, " config --format json"): volumes := map[string]map[string]string{} - for _, logical := range requiredTestVolumes { + for _, logical := range requiredBackupVolumes(runner.installation) { volumes[logical] = map[string]string{"name": runner.installation.ProjectName() + "_" + logical} } payload := map[string]any{ @@ -672,8 +679,9 @@ func (runner *fakeBackupRunner) Run(_ context.Context, args []string, _ io.Reade encoded, _ := json.Marshal(payload) return compose.Result{Stdout: string(encoded)}, nil case strings.HasPrefix(command, "volume inspect "): - items := make([]map[string]any, 0, len(requiredTestVolumes)) - for _, logical := range requiredTestVolumes { + required := requiredBackupVolumes(runner.installation) + items := make([]map[string]any, 0, len(required)) + for _, logical := range required { items = append(items, map[string]any{ "Name": runner.installation.ProjectName() + "_" + logical, "Driver": "local", @@ -703,15 +711,15 @@ func (runner *fakeBackupRunner) Run(_ context.Context, args []string, _ io.Reade lines = append(lines, fmt.Sprintf(`{"Service":%q,"State":%q}`, service, states[service])) } return compose.Result{Stdout: strings.Join(lines, "\n")}, nil - case strings.Contains(command, "/internal/maintenance/status"): + case strings.Contains(command, "operator-command.js maintenance-status"): return compose.Result{Stdout: fmt.Sprintf(`{"active":%t,"admissions":0,"recoveryRequired":false}`, runner.maintenance)}, nil - case strings.Contains(command, "/internal/maintenance/activate"): + case strings.Contains(command, "operator-command.js maintenance-activate"): runner.maintenance = true return compose.Result{Stdout: `{"active":true,"admissions":0,"recoveryRequired":false}`}, nil - case strings.Contains(command, "/internal/maintenance/deactivate"): + case strings.Contains(command, "operator-command.js maintenance-deactivate"): runner.maintenance = false return compose.Result{Stdout: `{"active":false,"admissions":0,"recoveryRequired":false}`}, nil - case strings.Contains(command, "/sessions?scope="): + case strings.Contains(command, "operator-command.js session-inventory"): if len(runner.sessionResponses) == 0 { return compose.Result{Stdout: `[]`}, nil } diff --git a/tools/tht/internal/backup/restore.go b/tools/tht/internal/backup/restore.go index ffca69de..6ff619e5 100644 --- a/tools/tht/internal/backup/restore.go +++ b/tools/tht/internal/backup/restore.go @@ -20,8 +20,10 @@ type RestoreRequest struct { Drain bool } -// RestoreResult records the retained recovery point and the final service state. +// RestoreResult records the final service state. Secret-bearing recovery checkpoints are always +// destroyed internally and are therefore never exposed to callers. type RestoreResult struct { + // Deprecated: always empty. Recovery checkpoints are private transaction internals. Checkpoint string Restarted bool Verified bool @@ -33,6 +35,9 @@ type restoreVerify func(context.Context, config.Installation, archiveRunner) err type restoreDependencies struct { preflight func(context.Context, config.Installation, PreflightRequest) (PreflightResult, error) checkpoint func(context.Context, config.Installation, CreateRequest) (Result, error) + prepareRecovery func(context.Context, config.Installation, string) (PreflightResult, error) + recover func(context.Context, config.Installation, PreflightResult, bool) error + cleanupCheckpoint func(string) error acquireLock func(config.Installation) (restoreLock, error) runner archiveRunner sleep func(duration time.Duration) @@ -56,7 +61,7 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati if request.Archive == "" { return RestoreResult{}, errors.New("restore archive is required") } - if deps.preflight == nil || deps.checkpoint == nil || deps.acquireLock == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == nil || deps.resetAuthenticationState == nil || deps.verify == nil { + if deps.preflight == nil || deps.checkpoint == nil || deps.prepareRecovery == nil || deps.recover == nil || deps.cleanupCheckpoint == nil || deps.acquireLock == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == nil || deps.resetAuthenticationState == nil || deps.verify == nil { return RestoreResult{}, errors.New("restore dependencies are incomplete") } preflight, err := deps.preflight(ctx, installation, PreflightRequest{Archive: request.Archive, Confirm: true, AllowExternalSecrets: true}) @@ -69,11 +74,21 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati return RestoreResult{}, err } - checkpoint, err := deps.checkpoint(ctx, installation, CreateRequest{}) + checkpoint, err := deps.checkpoint(ctx, installation, CreateRequest{IncludeSecrets: true, Confirm: true}) if err != nil { return RestoreResult{}, fmt.Errorf("create recovery checkpoint: %w", err) } - result.Checkpoint = checkpoint.Path + recovery, err := deps.prepareRecovery(ctx, installation, checkpoint.Path) + if err != nil { + cleanupErr := deps.cleanupCheckpoint(checkpoint.Path) + return RestoreResult{}, errors.Join(fmt.Errorf("validate recovery checkpoint: %w", err), cleanupErr) + } + defer func() { + if cleanupErr := deps.cleanupCheckpoint(checkpoint.Path); cleanupErr != nil { + resultErr = errors.Join(resultErr, fmt.Errorf("destroy recovery checkpoint: %w", cleanupErr)) + } + }() + defer recovery.CloseArchive() lock, err := deps.acquireLock(installation) if err != nil { return result, err @@ -91,7 +106,9 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati mutated := false defer func() { if resultErr != nil && mutated { - _ = runCompose(context.Background(), installation, deps.runner, "stop") + if recoveryErr := deps.recover(context.Background(), installation, recovery, wasRunning); recoveryErr != nil { + resultErr = errors.Join(resultErr, fmt.Errorf("restore recovery checkpoint: %w", recoveryErr)) + } } }() if wasRunning { @@ -105,42 +122,9 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati return result, err } } - reader, err := zip.NewReader(archive, preflight.ArchiveSize) - if err != nil { - return result, fmt.Errorf("read verified restore archive: %w", err) - } - members := make(map[string]*zip.File, len(reader.File)) - for _, member := range reader.File { - members[member.Name] = member - } - for _, entry := range preflight.Entries { - member := members[entry.Path] - if member == nil { - return result, fmt.Errorf("verified archive is missing %q", entry.Path) - } - stream, openErr := member.Open() - if openErr != nil { - return result, fmt.Errorf("open verified archive member %q: %w", entry.Path, openErr) - } - mutated = true - var restoreErr error - if entry.Kind == EntryVolume { - volume, found := restoreVolumeMetadata(preflight.Manifest, entry.LogicalName) - if !found { - _ = stream.Close() - return result, errors.New("verified volume metadata is incomplete") - } - restoreErr = deps.restoreVolume(ctx, installation, volume, stream) - } else { - restoreErr = deps.restoreFile(ctx, installation, entry, stream) - } - closeErr := stream.Close() - if restoreErr != nil { - return result, restoreErr - } - if closeErr != nil { - return result, closeErr - } + mutated = true + if err := restoreVerifiedEntries(ctx, installation, preflight, archive, deps.restoreFile, deps.restoreVolume); err != nil { + return result, err } if err := deps.resetAuthenticationState(ctx, installation, deps.runner); err != nil { return result, fmt.Errorf("reset authentication state: %w", err) @@ -164,6 +148,53 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati return result, nil } +func restoreVerifiedEntries( + ctx context.Context, + installation config.Installation, + preflight PreflightResult, + archive io.ReaderAt, + restoreFile func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error, + restoreVolume func(context.Context, config.Installation, VolumeMetadata, io.Reader) error, +) error { + reader, err := zip.NewReader(archive, preflight.ArchiveSize) + if err != nil { + return fmt.Errorf("read verified restore archive: %w", err) + } + members := make(map[string]*zip.File, len(reader.File)) + for _, member := range reader.File { + members[member.Name] = member + } + for _, entry := range preflight.Entries { + member := members[entry.Path] + if member == nil { + return fmt.Errorf("verified archive is missing %q", entry.Path) + } + stream, openErr := member.Open() + if openErr != nil { + return fmt.Errorf("open verified archive member %q: %w", entry.Path, openErr) + } + var restoreErr error + if entry.Kind == EntryVolume { + volume, found := restoreVolumeMetadata(preflight.Manifest, entry.LogicalName) + if !found { + _ = stream.Close() + return errors.New("verified volume metadata is incomplete") + } + restoreErr = restoreVolume(ctx, installation, volume, stream) + } else { + restoreErr = restoreFile(ctx, installation, entry, stream) + } + closeErr := stream.Close() + if restoreErr != nil { + return restoreErr + } + if closeErr != nil { + return closeErr + } + } + return nil +} + func restoreVolumeMetadata(manifest Manifest, logicalName string) (VolumeMetadata, bool) { if logicalName == "" { return VolumeMetadata{}, false @@ -177,12 +208,12 @@ func restoreVolumeMetadata(manifest Manifest, logicalName string) (VolumeMetadat } // resetAuthenticationState clears browser sessions and pending OIDC transactions without touching -// installation-global auth.yaml or users.yaml. The command runs as the unprivileged core user so -// the recreated state root is private to the service on both the local volume and server /data bind. +// installation-global auth.yaml or users.yaml. A root-scoped one-shot repairs ownership and mode +// before clearing children; links and malformed roots are rejected before any recursive removal. func resetAuthenticationState(ctx context.Context, installation config.Installation, runner archiveRunner) error { result, err := runner.Run(ctx, installation.ComposeArgs( - "run", "--rm", "--no-deps", "--no-TTY", "--entrypoint", "sh", "core", "-ceu", - "find /data/auth -mindepth 1 -maxdepth 1 -exec rm -rf -- {} + && install -d -m 0700 /data/auth /data/auth/sessions /data/auth/oidc && test -z \"$(find /data/auth/sessions /data/auth/oidc -mindepth 1 -print -quit)\"", + "run", "--rm", "--no-deps", "--no-TTY", "--user", "0:0", "--entrypoint", "sh", "core", "-ceu", + "test ! -L /data/auth && { test ! -e /data/auth || test -d /data/auth; } && install -d -o 10001 -g 10001 -m 0700 /data/auth && find /data/auth -mindepth 1 -maxdepth 1 -exec rm -rf -- {} + && install -d -o 10001 -g 10001 -m 0700 /data/auth/sessions /data/auth/oidc && chmod 0700 /data/auth /data/auth/sessions /data/auth/oidc && chown 10001:10001 /data/auth /data/auth/sessions /data/auth/oidc && test -z \"$(find /data/auth/sessions /data/auth/oidc -mindepth 1 -print -quit)\"", ), nil) if err != nil { return dockerError("reset authentication state", result, err) diff --git a/tools/tht/internal/backup/restore_host.go b/tools/tht/internal/backup/restore_host.go index e970992b..33dc45bb 100644 --- a/tools/tht/internal/backup/restore_host.go +++ b/tools/tht/internal/backup/restore_host.go @@ -5,7 +5,6 @@ import ( "crypto/rand" "crypto/sha256" "encoding/hex" - "encoding/json" "errors" "fmt" "io" @@ -15,6 +14,7 @@ import ( "strings" "time" + "github.com/aritmolab/thothii/tools/tht/internal/authconfig" "github.com/aritmolab/thothii/tools/tht/internal/compose" "github.com/aritmolab/thothii/tools/tht/internal/config" "github.com/aritmolab/thothii/tools/tht/internal/doctor" @@ -26,7 +26,7 @@ import ( func productionRestoreDependencies(installation config.Installation) restoreDependencies { runner := hostRunner{runner: compose.NewRunner(""), binary: "docker", profile: installation.Profile} - return restoreDependencies{ + deps := restoreDependencies{ preflight: func(ctx context.Context, target config.Installation, request PreflightRequest) (PreflightResult, error) { return Preflight(ctx, target, request, PreflightDependencies{ FreeBytes: restoreFreeBytes, @@ -47,6 +47,7 @@ func productionRestoreDependencies(installation config.Installation) restoreDepe request.Output = path return Create(ctx, target, request) }, + cleanupCheckpoint: cleanupRecoveryCheckpoint, acquireLock: func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) }, @@ -71,6 +72,46 @@ func productionRestoreDependencies(installation config.Installation) restoreDepe "workspace": verifyRestoreWorkspace, }, } + deps.prepareRecovery = func(ctx context.Context, target config.Installation, path string) (PreflightResult, error) { + return deps.preflight(ctx, target, PreflightRequest{Archive: path, Confirm: true, AllowExternalSecrets: true}) + } + deps.recover = func(ctx context.Context, target config.Installation, recovery PreflightResult, wasRunning bool) error { + return recoverRestoreTransaction(ctx, target, recovery, wasRunning, deps) + } + return deps +} + +func cleanupRecoveryCheckpoint(path string) error { + if err := safeio.RemoveCanonicalPrivateRegular(path); err != nil { + return errors.New("private recovery checkpoint could not be destroyed safely") + } + return nil +} + +func recoverRestoreTransaction(ctx context.Context, installation config.Installation, recovery PreflightResult, wasRunning bool, deps restoreDependencies) error { + var resultErr error + if err := runCompose(ctx, installation, deps.runner, "stop"); err != nil { + resultErr = errors.Join(resultErr, err) + } + archive, err := recovery.RevalidateArchive() + if err != nil { + return errors.Join(resultErr, err) + } + if err := restoreVerifiedEntries(ctx, installation, recovery, archive, deps.restoreFile, deps.restoreVolume); err != nil { + return errors.Join(resultErr, err) + } + // Recovery restores only configuration/secret files and durable application volumes. Runtime + // authentication state is deliberately reset again so neither sessions nor OIDC transactions + // survive a failed restore attempt. + if err := deps.resetAuthenticationState(ctx, installation, deps.runner); err != nil { + return errors.Join(resultErr, err) + } + if wasRunning { + if err := composeStartAndVerify(ctx, installation, deps.runner); err != nil { + return errors.Join(resultErr, err) + } + } + return resultErr } func restoreCheckpointPath(installation config.Installation, now time.Time) (string, error) { @@ -131,14 +172,14 @@ func validateRestoreReference(installation config.Installation, entry Entry) err } func validateRestoreVolumes(ctx context.Context, installation config.Installation, manifest Manifest, runner archiveRunner) error { - if len(manifest.Volumes) != len(requiredVolumes) { + if len(manifest.Volumes) != len(requiredBackupVolumes(installation)) { return errors.New("backup volume set is incomplete") } rendered, err := renderedConfiguration(ctx, installation, runner) if err != nil { return err } - current, err := inspectRequiredVolumes(ctx, runner, rendered) + current, err := inspectRequiredVolumes(ctx, installation, runner, rendered) if err != nil { return err } @@ -313,6 +354,13 @@ func verifyRestoreDoctor(ctx context.Context, installation config.Installation, if configErr != nil { return dockerError("verify restored Compose configuration", result, configErr) } + diagnostics, diagnosticErr := authconfig.Check(ctx, installation, runner, false, false) + if diagnosticErr != nil || !diagnostics.Ready { + if diagnosticErr != nil { + return diagnosticErr + } + return errors.New("authentication diagnostics did not pass after restore") + } return nil } report, err := doctor.Run(ctx, installation, runner) @@ -320,7 +368,16 @@ func verifyRestoreDoctor(ctx context.Context, installation config.Installation, return err } if !report.OK { - return errors.New("aggregate doctor did not pass after restore") + failed := make([]string, 0, len(report.Checks)) + for _, check := range report.Checks { + if check.Status == doctor.StatusFailed { + failed = append(failed, check.Name) + } + } + if len(failed) == 0 { + return errors.New("aggregate doctor did not pass after restore") + } + return fmt.Errorf("aggregate doctor failed checks: %s", strings.Join(failed, ",")) } return nil } @@ -328,25 +385,36 @@ func verifyRestoreDoctor(ctx context.Context, installation config.Installation, func verifyRestorePi(ctx context.Context, installation config.Installation, runner archiveRunner) error { running, err := restoreVerificationRunning(ctx, installation, runner) if err != nil || !running { - return err + if err != nil { + return err + } + result, runErr := runner.Run(ctx, installation.ComposeArgs( + "run", "--rm", "--no-deps", "--no-TTY", "core", "pi", "--version", + ), nil) + if runErr != nil || strings.TrimSpace(result.Stdout) == "" { + if runErr == nil { + runErr = errors.New("Pi version probe returned no version") + } + return dockerError("verify restored Pi runtime", result, runErr) + } + return nil } return pi.Doctor(ctx, compose.InstallationRunner{Installation: installation, Runner: runner}) } func verifyRestoreWorkspace(ctx context.Context, installation config.Installation, runner archiveRunner) error { running, err := restoreVerificationRunning(ctx, installation, runner) - if err != nil || !running { + if err != nil { return err } - result, err := runner.Run(ctx, installation.ComposeArgs( - "exec", "-T", "core", "curl", "-fsS", "--max-time", "5", "http://127.0.0.1:8787/workspaces", - ), nil) - if err != nil { - return dockerError("inspect restored workspaces", result, err) + command := []string{"exec", "-T", "core", "node", "-e"} + if !running { + command = []string{"run", "--rm", "--no-deps", "--no-TTY", "core", "node", "-e"} } - var workspaces []json.RawMessage - if json.Unmarshal([]byte(result.Stdout), &workspaces) != nil { - return errors.New("restored workspace inspection returned invalid JSON") + command = append(command, `const fs=require("node:fs");const p="/data/workspace-registry/state/active.json";const s=JSON.parse(fs.readFileSync(p,"utf8"));const h=/^[0-9a-f]{40}$/;if(!h.test(s.head)||!Array.isArray(s.revisions)||s.revisions.some(r=>!r||typeof r.id!=="string"||!r.id||!h.test(r.commit)||!h.test(r.blob)))process.exit(1);for(const r of s.revisions)fs.accessSync("/data/workspace-registry/snapshots/"+r.commit+"/"+r.id+".yaml",fs.constants.R_OK)`) + result, err := runner.Run(ctx, installation.ComposeArgs(command...), nil) + if err != nil { + return dockerError("validate restored workspace registry", result, err) } return nil } diff --git a/tools/tht/internal/backup/restore_test.go b/tools/tht/internal/backup/restore_test.go index 8dae6f3d..34270c0c 100644 --- a/tools/tht/internal/backup/restore_test.go +++ b/tools/tht/internal/backup/restore_test.go @@ -6,6 +6,7 @@ import ( "context" "errors" "io" + "os" "path/filepath" "strings" "testing" @@ -122,6 +123,14 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi checkpointRequest = request return Result{Path: "/tmp/checkpoint.zip"}, nil } + deps.prepareRecovery = func(context.Context, config.Installation, string) (PreflightResult, error) { + events = append(events, "prepare-recovery") + return PreflightResult{}, nil + } + deps.cleanupCheckpoint = func(string) error { + events = append(events, "cleanup-checkpoint") + return nil + } deps.acquireLock = func(config.Installation) (restoreLock, error) { events = append(events, "lock") return fakeRestoreLock{release: func() { events = append(events, "unlock") }}, nil @@ -142,13 +151,13 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi if err != nil { t.Fatal(err) } - if result.Checkpoint != "/tmp/checkpoint.zip" || result.Restarted || !result.Verified { + if result.Checkpoint != "" || result.Restarted || !result.Verified { t.Fatalf("Restore() result = %#v", result) } - if checkpointRequest.IncludeSecrets || checkpointRequest.Confirm { - t.Fatalf("checkpoint request = %#v, want non-secret unconfirmed checkpoint", checkpointRequest) + if !checkpointRequest.IncludeSecrets || !checkpointRequest.Confirm { + t.Fatalf("checkpoint request = %#v, want private confirmed secret-aware checkpoint", checkpointRequest) } - if got, want := events, []string{"checkpoint", "lock", "file:configuration/operator.env", "health", "doctor", "pi", "workspace", "unlock"}; !equalStrings(got, want) { + if got, want := events, []string{"checkpoint", "prepare-recovery", "lock", "file:configuration/operator.env", "health", "doctor", "pi", "workspace", "unlock", "cleanup-checkpoint"}; !equalStrings(got, want) { t.Fatalf("restore events = %v, want %v", got, want) } } @@ -196,9 +205,11 @@ func TestResetAuthenticationStateCreatesOnlyPrivateEmptyStateDirectories(t *test } joined := strings.Join(runner.args, "\x00") for _, required := range []string{ - "run", "--rm", "--no-deps", "--no-TTY", "--entrypoint", "sh", "core", "-ceu", + "run", "--rm", "--no-deps", "--no-TTY", "--user", "0:0", "--entrypoint", "sh", "core", "-ceu", + "test ! -L /data/auth", + "install -d -o 10001 -g 10001 -m 0700 /data/auth", "find /data/auth -mindepth 1 -maxdepth 1 -exec rm -rf -- {} +", - "install -d -m 0700 /data/auth /data/auth/sessions /data/auth/oidc", + "install -d -o 10001 -g 10001 -m 0700 /data/auth/sessions /data/auth/oidc", "find /data/auth/sessions /data/auth/oidc -mindepth 1 -print -quit", } { if !strings.Contains(joined, required) { @@ -207,6 +218,16 @@ func TestResetAuthenticationStateCreatesOnlyPrivateEmptyStateDirectories(t *test } } +func TestResetAuthenticationStateFailsClosedForDamagedRoot(t *testing.T) { + installation := preflightTestInstallation(t) + runner := &authenticationStateResetRunner{result: compose.Result{ExitCode: 1}, err: errors.New("damaged auth root")} + + err := resetAuthenticationState(context.Background(), installation, runner) + if err == nil || !strings.Contains(err.Error(), "reset authentication state") { + t.Fatalf("resetAuthenticationState() error = %v, want bounded fail-closed error", err) + } +} + func TestRestorePreflightFailureDoesNotMutateTarget(t *testing.T) { installation := preflightTestInstallation(t) runner := newBackupRunner(installation, true) @@ -259,7 +280,7 @@ func TestRestoreCheckpointFailureDoesNotMutateTarget(t *testing.T) { } } -func TestRestoreFileFailureStopsMutatedTargetAndRetainsCheckpoint(t *testing.T) { +func TestRestoreFileFailureRollsBackSecretAwareCheckpointBeforeCleanup(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) @@ -268,7 +289,17 @@ func TestRestoreFileFailureStopsMutatedTargetAndRetainsCheckpoint(t *testing.T) deps.checkpoint = func(context.Context, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/recovery.zip"}, nil } + var events []string + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + events = append(events, "recover") + return nil + } + deps.cleanupCheckpoint = func(string) error { + events = append(events, "cleanup") + return nil + } deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + events = append(events, "mutate") return fileErr } @@ -276,43 +307,145 @@ func TestRestoreFileFailureStopsMutatedTargetAndRetainsCheckpoint(t *testing.T) if !errors.Is(err, fileErr) { t.Fatalf("restore error = %v, want file error", err) } - if result.Checkpoint != "/tmp/recovery.zip" { - t.Fatalf("recovery checkpoint = %q, want retained path", result.Checkpoint) + if result.Checkpoint != "" { + t.Fatalf("recovery checkpoint = %q, want no retained secret-bearing path", result.Checkpoint) } - if runner.running || runner.stopCount != 2 { - t.Fatalf("mutated target was not stopped: running=%t stops=%d", runner.running, runner.stopCount) + if got, want := events, []string{"mutate", "recover", "cleanup"}; !equalStrings(got, want) { + t.Fatalf("failure recovery events = %v, want %v", got, want) } } -func TestRestoreStartFailureStopsRunningTarget(t *testing.T) { +func TestRestoreFailureAfterAuthenticationMutationRollsBackAndClearsRuntimeState(t *testing.T) { + installation := preflightTestInstallation(t) + archive := restoreArchive(t) + runner := newBackupRunner(installation, false) + deps := restoreTestDependencies(t, runner) + resetErr := errors.New("authentication reset failed") + var events []string + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + events = append(events, "auth-config-mutated") + return nil + } + deps.resetAuthenticationState = func(context.Context, config.Installation, archiveRunner) error { + events = append(events, "auth-runtime-reset-failed") + return resetErr + } + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + events = append(events, "secret-aware-recovery-and-reauth-reset") + return nil + } + deps.cleanupCheckpoint = func(string) error { + events = append(events, "checkpoint-cleaned") + return nil + } + + _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) + if !errors.Is(err, resetErr) { + t.Fatalf("restore error = %v, want reset failure", err) + } + if got, want := events, []string{"auth-config-mutated", "auth-runtime-reset-failed", "secret-aware-recovery-and-reauth-reset", "checkpoint-cleaned"}; !equalStrings(got, want) { + t.Fatalf("post-auth recovery events = %v, want %v", got, want) + } +} + +func TestRestoreCleanupFailureDoesNotSuppressRollback(t *testing.T) { + installation := preflightTestInstallation(t) + archive := restoreArchive(t) + deps := restoreTestDependencies(t, newBackupRunner(installation, false)) + mutationErr := errors.New("post-mutation failure") + cleanupErr := errors.New("checkpoint cleanup failure") + recovered := false + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return mutationErr } + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { recovered = true; return nil } + deps.cleanupCheckpoint = func(string) error { return cleanupErr } + + _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) + if !recovered || !errors.Is(err, mutationErr) || !errors.Is(err, cleanupErr) { + t.Fatalf("restore error = %v recovered=%t, want joined mutation/cleanup failure after rollback", err, recovered) + } +} + +func TestRecoveryCheckpointCleanupRejectsSymlinksAndPermissiveFiles(t *testing.T) { + root := t.TempDir() + target := filepath.Join(root, "checkpoint.zip") + if err := os.WriteFile(target, []byte("secret"), 0o600); err != nil { + t.Fatal(err) + } + link := filepath.Join(root, "checkpoint-link.zip") + if err := os.Symlink(target, link); err != nil { + t.Skipf("symlink unavailable: %v", err) + } + if err := cleanupRecoveryCheckpoint(link); err == nil { + t.Fatal("cleanupRecoveryCheckpoint accepted a symlink") + } + if _, err := os.Stat(target); err != nil { + t.Fatalf("symlink target was changed: %v", err) + } + if err := os.Chmod(target, 0o644); err != nil { + t.Fatal(err) + } + if err := cleanupRecoveryCheckpoint(target); err == nil { + t.Fatal("cleanupRecoveryCheckpoint accepted a permissive secret checkpoint") + } +} + +func TestRecoveryCheckpointCleanupRemovesOnlyPrivateRegularFile(t *testing.T) { + root, err := filepath.EvalSymlinks(t.TempDir()) + if err != nil { + t.Fatal(err) + } + checkpoint := filepath.Join(root, "checkpoint.zip") + if err := os.WriteFile(checkpoint, []byte("secret checkpoint"), 0o600); err != nil { + t.Fatal(err) + } + if err := cleanupRecoveryCheckpoint(checkpoint); err != nil { + t.Fatal(err) + } + if _, err := os.Lstat(checkpoint); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("private checkpoint still exists after cleanup: %v", err) + } +} + +func TestRestoreStartFailureRecoversPreviouslyRunningTarget(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) backingRunner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, failStartRestoreRunner{fakeBackupRunner: backingRunner}) + recovered := false + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + recovered = true + backingRunner.running = true + return nil + } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if err == nil || !strings.Contains(err.Error(), "start refused") { t.Fatalf("restore error = %v, want restart failure", err) } - if backingRunner.running || backingRunner.stopCount != 2 || backingRunner.startCount != 0 { - t.Fatalf("failed restart left target available: running=%t stops=%d starts=%d", backingRunner.running, backingRunner.stopCount, backingRunner.startCount) + if !recovered || !backingRunner.running { + t.Fatalf("failed restart was not rolled back: recovered=%t running=%t", recovered, backingRunner.running) } } -func TestRestoreVerificationFailureStopsRunningTarget(t *testing.T) { +func TestRestoreVerificationFailureRecoversPreviouslyRunningTarget(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) verificationErr := errors.New("Pi is unavailable") deps.verify["pi"] = func(context.Context, config.Installation, archiveRunner) error { return verificationErr } + recovered := false + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + recovered = true + return nil + } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, verificationErr) { t.Fatalf("restore error = %v, want verification failure", err) } - if runner.running || runner.stopCount != 2 || runner.startCount != 1 { - t.Fatalf("verification failure left target available: running=%t stops=%d starts=%d", runner.running, runner.stopCount, runner.startCount) + if !recovered || !runner.running { + t.Fatalf("verification failure was not rolled back: recovered=%t running=%t", recovered, runner.running) } } @@ -361,11 +494,15 @@ type fakeRestoreLock struct { release func() } -type authenticationStateResetRunner struct{ args []string } +type authenticationStateResetRunner struct { + args []string + result compose.Result + err error +} func (runner *authenticationStateResetRunner) Run(_ context.Context, args []string, _ io.Reader) (compose.Result, error) { runner.args = append([]string(nil), args...) - return compose.Result{}, nil + return runner.result, runner.err } func (runner *authenticationStateResetRunner) Stream(context.Context, []string, io.Reader, io.Writer) (compose.Result, error) { @@ -402,10 +539,15 @@ func restoreTestDependencies(t *testing.T, runner archiveRunner) restoreDependen checkpoint: func(context.Context, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/default-checkpoint.zip"}, nil }, - acquireLock: func(config.Installation) (restoreLock, error) { return fakeRestoreLock{}, nil }, - runner: runner, - sleep: func(time.Duration) {}, - restoreFile: func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return nil }, + prepareRecovery: func(context.Context, config.Installation, string) (PreflightResult, error) { + return PreflightResult{}, nil + }, + recover: func(context.Context, config.Installation, PreflightResult, bool) error { return nil }, + cleanupCheckpoint: func(string) error { return nil }, + acquireLock: func(config.Installation) (restoreLock, error) { return fakeRestoreLock{}, nil }, + runner: runner, + sleep: func(time.Duration) {}, + restoreFile: func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return nil }, restoreVolume: func(context.Context, config.Installation, VolumeMetadata, io.Reader) error { return nil }, diff --git a/tools/tht/internal/doctor/report.go b/tools/tht/internal/doctor/report.go index c492eba2..0f8d0843 100644 --- a/tools/tht/internal/doctor/report.go +++ b/tools/tht/internal/doctor/report.go @@ -293,15 +293,18 @@ func authenticationCheck(ctx context.Context, installation config.Installation, } func workflowCheck(ctx context.Context, installation config.Installation, runner Runner, secrets []string, add func(string, string, string)) { - result, err := runner.Run(ctx, installation.ComposeArgs("exec", "-T", "core", "tht", "doctor", "--json"), nil) + result, err := runner.Run(ctx, installation.ComposeArgs( + "exec", "-T", "core", "node", "dist/operator-command.js", "workflow-doctor", + ), nil) if err != nil { add("workflow", StatusFailed, commandDetail("container-local workflow doctor", result, err, secrets)) return } var payload struct { - OK bool `json:"ok"` + Ready bool `json:"ready"` + Workspaces int `json:"workspaces"` } - if json.Unmarshal([]byte(result.Stdout), &payload) != nil || !payload.OK { + if json.Unmarshal([]byte(result.Stdout), &payload) != nil || !payload.Ready || payload.Workspaces < 1 { add("workflow", StatusFailed, "container-local workflow doctor returned an invalid or failing report") return } diff --git a/tools/tht/internal/doctor/report_test.go b/tools/tht/internal/doctor/report_test.go index 1d844d04..695954ad 100644 --- a/tools/tht/internal/doctor/report_test.go +++ b/tools/tht/internal/doctor/report_test.go @@ -125,8 +125,8 @@ func TestRunUsesOnlyContainerLocalWorkflowAndPiDiagnosticsWhenCoreRuns(t *testin if !strings.Contains(calls, "exec -T core node dist/auth/diagnostic-command.js --json") { t.Fatalf("Run() calls = %s, want core-local authentication diagnostic", calls) } - if !strings.Contains(calls, "exec -T core tht doctor --json") { - t.Fatalf("Run() calls = %s, want core-local workflow doctor", calls) + if !strings.Contains(calls, "exec -T core node dist/operator-command.js workflow-doctor") { + t.Fatalf("Run() calls = %s, want registry-bound core-local workflow doctor", calls) } if strings.Contains(calls, ".venv") || strings.Contains(calls, "python") { t.Fatalf("Run() calls = %s, must not require host Python", calls) @@ -271,11 +271,11 @@ func (r *doctorRunner) Run(_ context.Context, args []string, _ io.Reader) (compo return compose.Result{Stdout: healthyServices}, nil } return compose.Result{Stdout: r.services}, nil - case strings.Contains(call, "tht doctor --json"): + case strings.Contains(call, "operator-command.js workflow-doctor"): if r.workflowFailure != "" { return compose.Result{Stderr: r.workflowFailure, ExitCode: 23}, errors.New("workflow failed") } - return compose.Result{Stdout: `{"ok":true,"checks":[]}`}, nil + return compose.Result{Stdout: `{"ready":true,"workspaces":1}`}, nil case strings.Contains(call, "dist/auth/diagnostic-command.js --json"): return compose.Result{Stdout: `{"ready":true,"mode":"oidc","checks":[{"level":"info","code":"auth_ready","message":"Authentication is ready."}]}`}, nil case strings.Contains(call, "workspace-registry/state/active.json"): @@ -287,7 +287,7 @@ func (r *doctorRunner) Run(_ context.Context, args []string, _ io.Reader) (compo return compose.Result{Stdout: "0.80.3\n"}, nil case strings.Contains(call, "test -w /home/thoth/.pi") || strings.Contains(call, "test -r /home/thoth/.pi/agent/auth.json") || strings.Contains(call, "/health"): return compose.Result{Stdout: `{"ready":true}`}, nil - case strings.Contains(call, "/pi-management/test"): + case strings.Contains(call, "operator-command.js pi-test"): return compose.Result{Stdout: `{"ready":true}`}, nil case strings.Contains(call, "ps -q core"): return compose.Result{Stdout: "core-id\n"}, nil diff --git a/tools/tht/internal/pi/commands.go b/tools/tht/internal/pi/commands.go index def58e40..107c6f5b 100644 --- a/tools/tht/internal/pi/commands.go +++ b/tools/tht/internal/pi/commands.go @@ -39,13 +39,6 @@ type settingsFileSnapshot struct { RawBase64 string `json:"rawBase64"` } -var internalIdentityHeaders = []string{ - "-H", "x-thoth-principal-issuer: tht", - "-H", "x-thoth-principal-subject: tht-maintenance", - "-H", "x-thoth-principal-display-name: Tht maintenance", - "-H", "x-thoth-is-admin: 1", -} - // Configure changes the backend's real installation settings through a core-side helper. It // deliberately has no secret or endpoint input: external endpoints remain Compose-owned. func Configure(ctx context.Context, runner Runner, value Defaults) error { @@ -127,9 +120,7 @@ func ConfigurationOptions(ctx context.Context, runner Runner) ([]ModelOption, er } func configurationOptions(ctx context.Context, runner Runner) (piOptions, error) { - args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...) - args = append(args, "http://127.0.0.1:8787/pi-management/options") - result, err := runCompose(ctx, runner, args...) + result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", "pi-options") if err != nil { return piOptions{}, commandError("Pi options check", result, err) } @@ -204,9 +195,7 @@ func restoreSettingsFile(ctx context.Context, runner Runner, snapshot settingsFi } 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...) + result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", "effective-settings") if err != nil { return nil, commandError("Pi installation settings read-back", result, err) } @@ -288,15 +277,13 @@ func expectedVersions(ctx context.Context, runner Runner) (string, string, error return expectedValue, labelValue, nil } -// Test retains the direct image-version signal, then delegates all Pi configuration/provider smoke -// validation to core's dedicated, admin-only Pi Management endpoint. +// Test retains the direct image-version signal, then invokes the same Pi management service via a +// scoped core-side command. No HTTP principal or privileged header exists on this path. func Test(ctx context.Context, runner Runner) error { if _, err := Status(ctx, runner); err != nil { return err } - args := append([]string{"exec", "-T", "core", "curl", "-fsS", "-X", "POST"}, internalIdentityHeaders...) - args = append(args, "http://127.0.0.1:8787/pi-management/test") - smoke, err := runCompose(ctx, runner, args...) + smoke, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", "pi-test") if err != nil { return commandError("Pi smoke check", smoke, err) } diff --git a/tools/tht/internal/pi/commands_test.go b/tools/tht/internal/pi/commands_test.go index bd96a006..e4644748 100644 --- a/tools/tht/internal/pi/commands_test.go +++ b/tools/tht/internal/pi/commands_test.go @@ -103,6 +103,7 @@ func TestConfigurePreservesTypedRecoveryRequiredErrorFromSettingsRestore(t *test } type configureRunner struct { + calls []string failure string restoreDurabilityFailure bool settings Defaults @@ -115,6 +116,7 @@ type configureRunner struct { func (f *configureRunner) Run(_ context.Context, args []string, stdin io.Reader) (compose.Result, error) { call := strings.Join(args, " ") + f.calls = append(f.calls, call) switch { case strings.Contains(call, "config --format json"): f.configReads++ @@ -123,7 +125,7 @@ func (f *configureRunner) Run(_ context.Context, args []string, stdin io.Reader) endpoint = "https://drift.example.invalid" } return compose.Result{Stdout: `{"services":{"core":{"image":"thothii-core:local","environment":{"THT_LLM_URL":"` + endpoint + `"}}}}`}, nil - case strings.Contains(call, "/pi-management/options"): + case strings.Contains(call, "operator-command.js pi-options"): return compose.Result{Stdout: `{"providers":["old","new"],"models":[{"provider":"old","id":"old-model"},{"provider":"new","id":"new-model"}],"reasoning":["low","medium","high"]}`}, nil case strings.Contains(call, "settings-cli.js --snapshot"): raw := f.settingsRaw @@ -164,7 +166,7 @@ func (f *configureRunner) Run(_ context.Context, args []string, stdin io.Reader) f.settingsRaw, _ = json.Marshal(f.settings) } return compose.Result{}, nil - case strings.Contains(call, "/settings"): + case strings.Contains(call, "operator-command.js effective-settings"): f.settingsReads++ if f.failure == "readback" && f.settings.Provider == "new" { return compose.Result{Stdout: `{}`}, nil @@ -220,18 +222,21 @@ func TestConfigureCompensationRestoresAbsentAndExactEmptyPriorFiles(t *testing.T // Catches tht reading the legacy public model route instead of the admin-only closed Pi // Management choices before it writes shared installation defaults. -func TestConfigureLoadsDedicatedClosedOptionsWritesRealCoreSettingsAndUsesUpstreamIdentity(t *testing.T) { - fake := newFakeRunner() - if err := Configure(context.Background(), fake, Defaults{Provider: "provider", Model: "model", Thinking: "medium"}); err != nil { +func TestConfigureUsesScopedCoreCommandWithoutMintingAnHTTPIdentity(t *testing.T) { + fake := &configureRunner{ + settings: Defaults{Provider: "old", Model: "old-model", Thinking: "low"}, + settingsExist: true, + settingsRaw: []byte(`{"provider":"old","model":"old-model","thinking":"low"}`), + } + if err := Configure(context.Background(), fake, Defaults{Provider: "new", Model: "new-model", Thinking: "high"}); err != nil { t.Fatal(err) } - assertCalled(t, fake.calls, "/pi-management/options") - assertCalled(t, fake.calls, "node /app/backend/dist/settings/settings-cli.js --provider provider --model model --thinking medium") - assertCalled(t, fake.calls, "x-thoth-principal-subject: tht-maintenance") - if got := strings.Join(fake.calls, "\n"); strings.Contains(got, "pi-defaults.json") || strings.Contains(got, "secret") { + assertCalled(t, fake.calls, "operator-command.js pi-options") + assertCalled(t, fake.calls, "node /app/backend/dist/settings/settings-cli.js --provider new --model new-model --thinking high") + if got := strings.Join(fake.calls, "\n"); strings.Contains(got, "x-thoth-principal") || strings.Contains(got, "x-thoth-is-admin") || strings.Contains(got, "pi-defaults.json") || strings.Contains(got, "secret") { t.Fatalf("commands=%q", got) } - if err := Configure(context.Background(), fake, Defaults{Provider: "provider", Model: "unknown", Thinking: "medium"}); err == nil { + if err := Configure(context.Background(), fake, Defaults{Provider: "new", Model: "unknown", Thinking: "medium"}); err == nil { t.Fatal("expected unknown model rejection") } } @@ -243,10 +248,10 @@ func TestTestUsesDedicatedSmokeEndpointAndIndependentImageVersionProbe(t *testin if err := Test(context.Background(), fake); err != nil { t.Fatalf("Test() error = %v", err) } - for _, command := range []string{"pi --version", "/pi-management/test", "x-thoth-principal-subject: tht-maintenance"} { + for _, command := range []string{"pi --version", "operator-command.js pi-test"} { assertCalled(t, fake.calls, command) } - for _, legacy := range []string{"/health", "/models", "/settings"} { + for _, legacy := range []string{"/pi-management/test", "x-thoth-principal", "x-thoth-is-admin"} { if strings.Contains(strings.Join(fake.calls, "\n"), legacy) { t.Fatalf("Pi smoke invoked legacy endpoint %q: %s", legacy, strings.Join(fake.calls, "\n")) } diff --git a/tools/tht/internal/pi/restart_test.go b/tools/tht/internal/pi/restart_test.go index 369158ec..3ccc2754 100644 --- a/tools/tht/internal/pi/restart_test.go +++ b/tools/tht/internal/pi/restart_test.go @@ -149,8 +149,8 @@ func TestRestartActivationFailureClearsPreMutationMaintenance(t *testing.T) { if fake.recreated { t.Fatal("pre-mutation activation failure recreated core") } - assertCalled(t, fake.calls, "/internal/maintenance/activate") - assertCalled(t, fake.calls, "/internal/maintenance/deactivate") + assertCalled(t, fake.calls, "operator-command.js maintenance-activate") + assertCalled(t, fake.calls, "operator-command.js maintenance-deactivate") }) } } @@ -319,17 +319,17 @@ func TestRecoverLifecycleMaintenanceVerifiesAndClearsRestartState(t *testing.T) if err != nil || updateState.Phase != PhaseVerified { t.Fatalf("update recovery state = %+v, %v; want verified image rollback metadata", updateState, err) } - deactivate := callIndex(fake.calls, "/internal/maintenance/deactivate") - lastVerification := lastCallIndexBefore(fake.calls, "/pi-management/test", deactivate) + deactivate := callIndex(fake.calls, "operator-command.js maintenance-deactivate") + lastVerification := lastCallIndexBefore(fake.calls, "operator-command.js pi-test", deactivate) if deactivate < 0 || lastVerification < 0 { t.Fatalf("calls = %v; want verification before maintenance deactivation", fake.calls) } verificationCount := 0 for index := 0; index < deactivate; index++ { - if strings.Contains(fake.calls[index], "/internal/maintenance/deactivate") { + if strings.Contains(fake.calls[index], "operator-command.js maintenance-deactivate") { t.Fatalf("maintenance reopened before combined verification: %v", fake.calls) } - if strings.Contains(fake.calls[index], "/pi-management/test") { + if strings.Contains(fake.calls[index], "operator-command.js pi-test") { verificationCount++ } } @@ -373,7 +373,7 @@ func TestRecoverLifecycleMaintenanceRejectsMalformedRestartState(t *testing.T) { if _, stateErr := os.Stat(restartStatePath); stateErr != nil { t.Fatalf("malformed restart state was removed: %v", stateErr) } - assertNotCalled(t, fake.calls, "/internal/maintenance/deactivate") + assertNotCalled(t, fake.calls, "operator-command.js maintenance-deactivate") }) } } diff --git a/tools/tht/internal/pi/update.go b/tools/tht/internal/pi/update.go index 84785ed1..ba518a09 100644 --- a/tools/tht/internal/pi/update.go +++ b/tools/tht/internal/pi/update.go @@ -473,13 +473,11 @@ func canonicalDigestReference(value string) (string, error) { } func setMaintenance(ctx context.Context, runner Runner, enabled bool) error { - path := "deactivate" + action := "maintenance-deactivate" if enabled { - path = "activate" + action = "maintenance-activate" } - args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...) - args = append(args, "-X", "POST", "http://127.0.0.1:8787/internal/maintenance/"+path) - result, err := runCompose(ctx, runner, args...) + result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", action) status, valid := parseMaintenanceStatus(result.Stdout) if err == nil && valid && status.Active == enabled && status.Admissions == 0 && !status.RecoveryRequired { return nil @@ -518,9 +516,7 @@ func parseMaintenanceStatus(value string) (MaintenanceState, bool) { } func MaintenanceStatus(ctx context.Context, runner Runner) (MaintenanceState, error) { - args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...) - args = append(args, "http://127.0.0.1:8787/internal/maintenance/status") - result, err := runCompose(ctx, runner, args...) + result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", "maintenance-status") if err != nil { return MaintenanceState{}, commandError("maintenance status check", result, err) } @@ -549,9 +545,7 @@ func activeSessions(ctx context.Context, runner Runner) (bool, error) { scope = requested } } - args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...) - args = append(args, "http://127.0.0.1:8787/sessions?scope="+scope) - result, err := runCompose(ctx, runner, args...) + result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", "session-inventory", scope) if err != nil { return false, commandError("active-session check", result, err) } diff --git a/tools/tht/internal/pi/update_test.go b/tools/tht/internal/pi/update_test.go index 79ef8426..f6fe6310 100644 --- a/tools/tht/internal/pi/update_test.go +++ b/tools/tht/internal/pi/update_test.go @@ -48,7 +48,7 @@ func TestActiveSessionsUsesInstallationScopedInventory(t *testing.T) { if active, err := activeSessions(context.Background(), runner); err != nil || active { t.Fatalf("activeSessions(%s) = %t, %v; want false, nil", scope, active, err) } - if len(runner.calls) != 1 || !strings.Contains(runner.calls[0], "/sessions?scope="+scope) { + if len(runner.calls) != 1 || !strings.Contains(runner.calls[0], "operator-command.js session-inventory "+scope) { t.Fatalf("activeSessions(%s) call = %v; want installation-scoped inventory", scope, runner.calls) } } @@ -250,7 +250,7 @@ func TestMaintenanceTransportLossAfterBackendRestartRequiresRecovery(t *testing. if fake.backendRestarts != 1 { t.Fatalf("backend restarts = %d, want 1", fake.backendRestarts) } - assertCalled(t, fake.calls, "/internal/maintenance/status") + assertCalled(t, fake.calls, "operator-command.js maintenance-status") for index, active := range fake.maintenanceAtRecreate { if !active { t.Fatalf("recreate %d started without durable maintenance", index+1) @@ -344,8 +344,8 @@ func TestCompensationReactivatesMaintenanceAndRescansBeforeRollback(t *testing.T t.Fatalf("calls %v contain no rollback", fake.calls) } recreate := callIndex(fake.calls, "force-recreate core") - activate := lastCallIndexBefore(fake.calls, "/internal/maintenance/activate", rollback) - scan := lastCallIndexBefore(fake.calls, "/sessions?scope=all", rollback) + activate := lastCallIndexBefore(fake.calls, "operator-command.js maintenance-activate", rollback) + scan := lastCallIndexBefore(fake.calls, "operator-command.js session-inventory all", rollback) if activate <= recreate || scan <= recreate { t.Fatalf("calls %v do not reactivate/confirm and rescan after candidate recreate before rollback", fake.calls) } @@ -599,7 +599,7 @@ func TestRecoverMaintenanceClearsOnlyAfterTerminalStateAndVerifiedSmoke(t *testi if selected := readSelectorReference(t, currentImageOverridePath(statePath)); selected != previous.Reference { t.Fatalf("maintenance cleanup changed durable selector to %q", selected) } - assertCalled(t, fake.calls, "/pi-management/test") + assertCalled(t, fake.calls, "operator-command.js pi-test") } func TestRecoverMaintenanceRefusesPendingTransaction(t *testing.T) { @@ -774,12 +774,12 @@ func TestUpdateRequiresConfirmationAndDrainsActiveSessions(t *testing.T) { if err != nil { t.Fatalf("Update() with drain error = %v", err) } - assertCalled(t, fake.calls, "http://127.0.0.1:8787/sessions?scope=all") - assertCalled(t, fake.calls, "/internal/maintenance/activate") - assertCalled(t, fake.calls, "/internal/maintenance/deactivate") + assertCalled(t, fake.calls, "operator-command.js session-inventory all") + assertCalled(t, fake.calls, "operator-command.js maintenance-activate") + assertCalled(t, fake.calls, "operator-command.js maintenance-deactivate") } -func TestMaintenanceControlUsesExactLoopbackOperatorIdentity(t *testing.T) { +func TestMaintenanceControlUsesOnlyTheScopedNonNetworkCommand(t *testing.T) { fake := newFakeRunner() if err := setMaintenance(context.Background(), fake, true); err != nil { t.Fatal(err) @@ -788,26 +788,15 @@ func TestMaintenanceControlUsesExactLoopbackOperatorIdentity(t *testing.T) { if _, err := MaintenanceStatus(context.Background(), fake); err != nil { t.Fatal(err) } - for _, path := range []string{"/internal/maintenance/activate", "/internal/maintenance/status"} { - found := false - for _, call := range fake.calls { - if !strings.Contains(call, path) { - continue - } - found = true - for _, header := range []string{ - "x-thoth-principal-issuer: tht", - "x-thoth-principal-subject: tht-maintenance", - "x-thoth-principal-display-name: Tht maintenance", - "x-thoth-is-admin: 1", - } { - if !strings.Contains(call, header) { - t.Fatalf("maintenance call %q lacks %q", call, header) - } - } + joined := strings.Join(fake.calls, "\n") + for _, action := range []string{"operator-command.js maintenance-activate", "operator-command.js maintenance-status"} { + if !strings.Contains(joined, action) { + t.Fatalf("maintenance action %q was not made: %s", action, joined) } - if !found { - t.Fatalf("maintenance call %q was not made", path) + } + for _, forbidden := range []string{"curl", "x-thoth-principal", "x-thoth-is-admin"} { + if strings.Contains(joined, forbidden) { + t.Fatalf("maintenance commands retained HTTP identity material %q: %s", forbidden, joined) } } } @@ -1074,7 +1063,7 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose if f.fail == "version" && (f.built || f.recreated) && strings.Contains(call, "pi --version") && strings.Contains(call, "exec") { return compose.Result{ExitCode: 1}, errors.New("version token=secret") } - if f.fail == "smoke" && (f.built || f.recreated) && strings.Contains(call, "127.0.0.1:8787/pi-management/test") { + if f.fail == "smoke" && (f.built || f.recreated) && strings.Contains(call, "operator-command.js pi-test") { return compose.Result{ExitCode: 1}, errors.New("smoke token=secret") } switch { @@ -1106,7 +1095,7 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose return compose.Result{Stdout: f.mountsJSON}, nil } 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"): + case strings.Contains(call, "operator-command.js maintenance-activate"): f.maintenance = true if f.fail == "maintenance-activate-durability" || f.fail == "maintenance-activate-durability-without-status-flag" { return compose.Result{ExitCode: 22}, errors.New("maintenance activation durability was not acknowledged") @@ -1119,7 +1108,7 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose return compose.Result{ExitCode: 52}, errors.New("lost activation response") } return compose.Result{Stdout: `{"active":true,"admissions":0}`}, nil - case strings.Contains(call, "/internal/maintenance/deactivate"): + case strings.Contains(call, "operator-command.js maintenance-deactivate"): if f.fail == "maintenance-clear" { return compose.Result{ExitCode: 53}, errors.New("maintenance clear failure") } @@ -1138,7 +1127,7 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose 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"): + case strings.Contains(call, "operator-command.js maintenance-status"): if f.fail == "maintenance-proof" && f.recreated { return compose.Result{Stdout: fmt.Sprintf(`{"active":%t,"admissions":0,"recoveryRequired":true}`, f.maintenance)}, nil } @@ -1149,7 +1138,7 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose 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"): + case strings.Contains(call, "operator-command.js session-inventory all"): if f.sessionsWire != "" { return compose.Result{Stdout: f.sessionsWire}, nil } @@ -1230,12 +1219,12 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose 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, "/pi-management/options"): + case strings.Contains(call, "operator-command.js pi-options"): if f.piManagementOptionsWire != "" { return compose.Result{Stdout: f.piManagementOptionsWire}, nil } return compose.Result{Stdout: `{"providers":["provider"],"models":[{"id":"model","provider":"provider"}],"reasoning":["low","medium","high"]}`}, nil - case strings.Contains(call, "/pi-management/test"): + case strings.Contains(call, "operator-command.js pi-test"): if f.currentImage == "sha256:old" { f.restoredProofComplete = true } @@ -1248,7 +1237,7 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose return compose.Result{Stdout: f.modelsWire}, nil } return compose.Result{Stdout: `{"models":[{"id":"model","provider":"provider"}]}`}, nil - case strings.Contains(call, "/settings"): + case strings.Contains(call, "operator-command.js effective-settings"): if f.currentImage == "sha256:old" { f.restoredProofComplete = true } diff --git a/tools/tht/internal/setup/run_test.go b/tools/tht/internal/setup/run_test.go index d5a81279..bd99fa83 100644 --- a/tools/tht/internal/setup/run_test.go +++ b/tools/tht/internal/setup/run_test.go @@ -331,8 +331,8 @@ func setupStage(args []string) (string, compose.Result) { return "doctor config", compose.Result{Stdout: renderedSetupConfig} case strings.Contains(joined, "exec -T core node dist/auth/diagnostic-command.js --json"): return "authentication", compose.Result{Stdout: `{"ready":true,"mode":"oidc","checks":[{"level":"info","code":"auth_ready","message":"Authentication is ready."}]}`} - case strings.Contains(joined, "exec -T core tht doctor --json"): - return "workflow doctor", compose.Result{Stdout: `{"ok":true,"components":{}}`} + case strings.Contains(joined, "exec -T core node dist/operator-command.js workflow-doctor"): + return "workflow doctor", compose.Result{Stdout: `{"ready":true,"workspaces":1}`} case strings.Contains(joined, "exec -T core curl -fsS --max-time 5 http://127.0.0.1:8787/health"): return "core HTTP", compose.Result{} case strings.Contains(joined, "exec -T frontend wget -q -T 5 -O /dev/null http://127.0.0.1:8080/"): @@ -345,7 +345,7 @@ func setupStage(args []string) (string, compose.Result) { return "pi doctor", compose.Result{Stdout: "0.80.3\n"} case strings.Contains(joined, "pi --version") || strings.Contains(joined, "PI_VERSION") || strings.Contains(joined, "test -w") || strings.Contains(joined, "test -r") || strings.Contains(joined, "127.0.0.1:8787/health"): return "pi doctor", compose.Result{Stdout: "0.80.3\n"} - case strings.Contains(joined, "/pi-management/test"): + case strings.Contains(joined, "operator-command.js pi-test"): return "pi doctor", compose.Result{Stdout: `{"ready":true}`} default: return "unexpected: " + joined, compose.Result{}