feat(auth): integrate authentication with installation lifecycle

This commit is contained in:
2026-08-17 23:33:51 +02:00
parent 0d0c15b4b8
commit e8b9995ed0
21 changed files with 672 additions and 234 deletions
-15
View File
@@ -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;
+123
View File
@@ -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<Array<Pick<SessionRow, "status" | "archived">>> {
const registry = new WorkspaceRegistry(config.workspaceRegistry);
const revisions = await registry.listRetainedSnapshots();
const runner = operatorRunner(config);
const sessions = new Map<string, SessionRow>();
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<unknown> {
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<void> {
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;
});
}
+2 -2
View File
@@ -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<void> {
const config = loadConfig(process.env);
const app = buildApp(config) as AppWithAuthSessionStore;
const sessions = app.thothiiAuthSessionStore;
if (sessions) {
await sessions.prune();
+7 -1
View File
@@ -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}`;
}
+3 -3
View File
@@ -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");
}
});
+34
View File
@@ -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");
});
});
+61 -7
View File
@@ -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'
}
+8 -8
View File
@@ -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
+22 -13
View File
@@ -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)
+15 -7
View File
@@ -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
}
+76 -45
View File
@@ -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)
+83 -15
View File
@@ -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
}
+165 -23
View File
@@ -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
},
+6 -3
View File
@@ -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
}
+5 -5
View File
@@ -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
+5 -18
View File
@@ -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)
}
+17 -12
View File
@@ -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"))
}
+7 -7
View File
@@ -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")
})
}
}
+5 -11
View File
@@ -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)
}
+25 -36
View File
@@ -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
}
+3 -3
View File
@@ -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{}