fix: acknowledge pi maintenance barrier
This commit is contained in:
+20
-4
@@ -15,7 +15,7 @@ import { settingsRoutes, effectiveSettings } from "./routes/settings.js";
|
|||||||
import { createPiModelLister } from "./pi/list-models.js";
|
import { createPiModelLister } from "./pi/list-models.js";
|
||||||
import { loadSettings, type Settings } from "./settings/settings-store.js";
|
import { loadSettings, type Settings } from "./settings/settings-store.js";
|
||||||
import { ReadinessManager } from "./runtime/readiness-manager.js";
|
import { ReadinessManager } from "./runtime/readiness-manager.js";
|
||||||
import { createMaintenanceGate } from "./runtime/maintenance-gate.js";
|
import { MaintenanceBarrier } from "./runtime/maintenance-gate.js";
|
||||||
import { WorkspaceRegistry } from "./workspaces/registry.js";
|
import { WorkspaceRegistry } from "./workspaces/registry.js";
|
||||||
import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js";
|
import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js";
|
||||||
import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js";
|
import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js";
|
||||||
@@ -33,8 +33,7 @@ export interface BuildAppDeps {
|
|||||||
workspaceRegistry?: WorkspaceRegistry;
|
workspaceRegistry?: WorkspaceRegistry;
|
||||||
workspaceDiagnoser?: WorkspaceDiagnoser;
|
workspaceDiagnoser?: WorkspaceDiagnoser;
|
||||||
workspaceRuntimeSupport?: (workspace: WorkspaceDescriptor) => boolean;
|
workspaceRuntimeSupport?: (workspace: WorkspaceDescriptor) => boolean;
|
||||||
/** Returns true while a host maintenance transaction is preventing new runtimes. */
|
maintenanceBarrier?: MaintenanceBarrier;
|
||||||
maintenanceGate?: () => boolean;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance {
|
export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance {
|
||||||
@@ -97,12 +96,27 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
|||||||
app.get("/health", async () => ({ status: "ok" }));
|
app.get("/health", async () => ({ status: "ok" }));
|
||||||
app.get("/health/dwh", async () => tht.dbPing());
|
app.get("/health/dwh", async () => tht.dbPing());
|
||||||
app.get("/me", async (req) => getPrincipal(req));
|
app.get("/me", async (req) => getPrincipal(req));
|
||||||
|
const maintenanceBarrier = deps?.maintenanceBarrier ?? new MaintenanceBarrier();
|
||||||
sessionRoutes(app, {
|
sessionRoutes(app, {
|
||||||
mgr, tht: tht as ThtRunner, hub, getSettings, readiness, listModels, workspaceRegistry,
|
mgr, tht: tht as ThtRunner, hub, getSettings, readiness, listModels, workspaceRegistry,
|
||||||
dwhPrecheck: config.dwhPrecheck,
|
dwhPrecheck: config.dwhPrecheck,
|
||||||
legacyWorkspaceMode: config.legacyWorkspaceMode,
|
legacyWorkspaceMode: config.legacyWorkspaceMode,
|
||||||
workspaceRuntimeSupport,
|
workspaceRuntimeSupport,
|
||||||
maintenanceGate: deps?.maintenanceGate ?? createMaintenanceGate(config.maintenanceFile),
|
maintenanceBarrier,
|
||||||
|
});
|
||||||
|
app.post("/internal/maintenance/activate", async (req, reply) => {
|
||||||
|
if (!isLoopback(req.ip) || req.principal?.subject !== "thothctl-maintenance") return reply.code(403).send({ error: "loopback maintenance control required" });
|
||||||
|
await maintenanceBarrier.activate();
|
||||||
|
return maintenanceBarrier.status();
|
||||||
|
});
|
||||||
|
app.post("/internal/maintenance/deactivate", async (req, reply) => {
|
||||||
|
if (!isLoopback(req.ip) || req.principal?.subject !== "thothctl-maintenance") return reply.code(403).send({ error: "loopback maintenance control required" });
|
||||||
|
maintenanceBarrier.deactivate();
|
||||||
|
return maintenanceBarrier.status();
|
||||||
|
});
|
||||||
|
app.get("/internal/maintenance/status", async (req, reply) => {
|
||||||
|
if (!isLoopback(req.ip) || req.principal?.subject !== "thothctl-maintenance") return reply.code(403).send({ error: "loopback maintenance control required" });
|
||||||
|
return maintenanceBarrier.status();
|
||||||
});
|
});
|
||||||
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
|
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
|
||||||
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
|
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
|
||||||
@@ -111,3 +125,5 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
|||||||
|
|
||||||
return app;
|
return app;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function isLoopback(ip: string): boolean { return ip === "127.0.0.1" || ip === "::1" || ip === "::ffff:127.0.0.1"; }
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
|||||||
import type { ListModelsFn } from "./meta.js";
|
import type { ListModelsFn } from "./meta.js";
|
||||||
import type { WorkspaceRegistry } from "../workspaces/registry.js";
|
import type { WorkspaceRegistry } from "../workspaces/registry.js";
|
||||||
import type { WorkspaceDescriptor } from "../workspaces/schema.js";
|
import type { WorkspaceDescriptor } from "../workspaces/schema.js";
|
||||||
|
import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js";
|
||||||
|
|
||||||
const BOOTSTRAP_FAILURE_MESSAGE =
|
const BOOTSTRAP_FAILURE_MESSAGE =
|
||||||
"Session startup failed. Check configuration and connectivity, then Resume the session.";
|
"Session startup failed. Check configuration and connectivity, then Resume the session.";
|
||||||
@@ -37,8 +38,7 @@ export function sessionRoutes(
|
|||||||
legacyWorkspaceMode?: boolean;
|
legacyWorkspaceMode?: boolean;
|
||||||
/** Fail-closed installation/runtime transport capability check. */
|
/** Fail-closed installation/runtime transport capability check. */
|
||||||
workspaceRuntimeSupport: (workspace: WorkspaceDescriptor) => boolean;
|
workspaceRuntimeSupport: (workspace: WorkspaceDescriptor) => boolean;
|
||||||
/** Host-controlled admission guard. Existing runtimes deliberately continue. */
|
maintenanceBarrier: MaintenanceBarrier;
|
||||||
maintenanceGate: () => boolean;
|
|
||||||
},
|
},
|
||||||
) {
|
) {
|
||||||
const lifecycleTails = new Map<string, Promise<void>>();
|
const lifecycleTails = new Map<string, Promise<void>>();
|
||||||
@@ -81,6 +81,14 @@ export function sessionRoutes(
|
|||||||
code: "maintenance",
|
code: "maintenance",
|
||||||
error: "Session admission is temporarily paused for maintenance. Try again shortly.",
|
error: "Session admission is temporarily paused for maintenance. Try again shortly.",
|
||||||
});
|
});
|
||||||
|
const admissionLeases = new WeakMap<object, () => void>();
|
||||||
|
app.addHook("preHandler", async (req, reply) => {
|
||||||
|
if (req.method !== "POST" || !(req.url === "/sessions" || /^\/sessions\/[^/]+\/resume(?:\?|$)/.test(req.url))) return;
|
||||||
|
const release = d.maintenanceBarrier.acquire();
|
||||||
|
if (!release) return maintenanceReply(reply);
|
||||||
|
admissionLeases.set(req, release);
|
||||||
|
});
|
||||||
|
app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); });
|
||||||
|
|
||||||
/** Include retained historical descriptors so removed workspaces remain resumable. */
|
/** Include retained historical descriptors so removed workspaces remain resumable. */
|
||||||
const sessionRevisions = async () => {
|
const sessionRevisions = async () => {
|
||||||
@@ -279,7 +287,6 @@ export function sessionRoutes(
|
|||||||
});
|
});
|
||||||
|
|
||||||
app.post("/sessions", async (req, reply) => {
|
app.post("/sessions", async (req, reply) => {
|
||||||
if (d.maintenanceGate()) return maintenanceReply(reply);
|
|
||||||
const b = req.body as {
|
const b = req.body as {
|
||||||
question: string; name?: string; workspace?: string; workspaceId?: string;
|
question: string; name?: string; workspace?: string; workspaceId?: string;
|
||||||
provider?: string; model?: string; thinking?: string;
|
provider?: string; model?: string; thinking?: string;
|
||||||
@@ -510,7 +517,6 @@ export function sessionRoutes(
|
|||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/resume", async (req, reply) => {
|
app.post("/sessions/:id/resume", async (req, reply) => {
|
||||||
if (d.maintenanceGate()) return maintenanceReply(reply);
|
|
||||||
const id = (req.params as any).id;
|
const id = (req.params as any).id;
|
||||||
const principal = getPrincipal(req);
|
const principal = getPrincipal(req);
|
||||||
return withSessionLifecycle(id, async () => {
|
return withSessionLifecycle(id, async () => {
|
||||||
|
|||||||
@@ -1,10 +1,27 @@
|
|||||||
import { existsSync } from "node:fs";
|
/** An in-process admission barrier. A lease spans the complete create/resume decision. */
|
||||||
|
export class MaintenanceBarrier {
|
||||||
|
private active = false;
|
||||||
|
private admissions = 0;
|
||||||
|
private waiters: (() => void)[] = [];
|
||||||
|
|
||||||
/**
|
acquire(): (() => void) | undefined {
|
||||||
* A host-side lifecycle transaction creates this marker before it checks for active sessions.
|
if (this.active) return undefined;
|
||||||
* The marker is intentionally only an admission gate: it must never terminate existing Pi
|
this.admissions += 1;
|
||||||
* processes or make their in-flight work unavailable.
|
let released = false;
|
||||||
*/
|
return () => {
|
||||||
export function createMaintenanceGate(markerFile: string): () => boolean {
|
if (released) return;
|
||||||
return () => existsSync(markerFile);
|
released = true;
|
||||||
|
this.admissions -= 1;
|
||||||
|
if (this.admissions === 0) this.waiters.splice(0).forEach((resolve) => resolve());
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
async activate(): Promise<void> {
|
||||||
|
this.active = true;
|
||||||
|
if (this.admissions === 0) return;
|
||||||
|
await new Promise<void>((resolve) => this.waiters.push(resolve));
|
||||||
|
}
|
||||||
|
|
||||||
|
deactivate(): void { this.active = false; }
|
||||||
|
status(): { active: boolean; admissions: number } { return { active: this.active, admissions: this.admissions }; }
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,19 @@
|
|||||||
|
import { test, expect } from "vitest";
|
||||||
|
import { MaintenanceBarrier } from "../src/runtime/maintenance-gate.js";
|
||||||
|
|
||||||
|
test("activation waits for an in-flight admission lease and rejects later admissions", async () => {
|
||||||
|
const gate = new MaintenanceBarrier();
|
||||||
|
const release = gate.acquire();
|
||||||
|
expect(release).toBeTypeOf("function");
|
||||||
|
let acknowledged = false;
|
||||||
|
const activation = gate.activate().then(() => { acknowledged = true; });
|
||||||
|
await Promise.resolve();
|
||||||
|
expect(acknowledged).toBe(false);
|
||||||
|
expect(gate.acquire()).toBeUndefined();
|
||||||
|
release?.();
|
||||||
|
await activation;
|
||||||
|
expect(acknowledged).toBe(true);
|
||||||
|
expect(gate.status()).toEqual({ active: true, admissions: 0 });
|
||||||
|
gate.deactivate();
|
||||||
|
expect(gate.acquire()).toBeTypeOf("function");
|
||||||
|
});
|
||||||
@@ -7,6 +7,7 @@ import { tmpdir } from "node:os";
|
|||||||
import { buildApp as buildRealApp } from "../src/app.js";
|
import { buildApp as buildRealApp } from "../src/app.js";
|
||||||
import { loadConfig } from "../src/config.js";
|
import { loadConfig } from "../src/config.js";
|
||||||
import { SseHub } from "../src/sse/sse-hub.js";
|
import { SseHub } from "../src/sse/sse-hub.js";
|
||||||
|
import { MaintenanceBarrier } from "../src/runtime/maintenance-gate.js";
|
||||||
|
|
||||||
const FAKE = path.resolve("../harness/tests/fake_pi/fake_pi_rpc.mjs");
|
const FAKE = path.resolve("../harness/tests/fake_pi/fake_pi_rpc.mjs");
|
||||||
const SCRIPT = path.resolve("../harness/tests/fake_pi/scripts/f1_disambiguation.json");
|
const SCRIPT = path.resolve("../harness/tests/fake_pi/scripts/f1_disambiguation.json");
|
||||||
@@ -76,8 +77,10 @@ test("upstream requests without a principal fail before a Pi runtime can be crea
|
|||||||
});
|
});
|
||||||
|
|
||||||
test("maintenance rejects new and resumed session admission without interrupting running sessions", async () => {
|
test("maintenance rejects new and resumed session admission without interrupting running sessions", async () => {
|
||||||
|
const maintenanceBarrier = new MaintenanceBarrier();
|
||||||
|
await maintenanceBarrier.activate();
|
||||||
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
||||||
maintenanceGate: () => true,
|
maintenanceBarrier,
|
||||||
thtRunner: { withPrincipal: () => ({ sessionShow: async () => ({ id: "open", status: "open" }) }) } as any,
|
thtRunner: { withPrincipal: () => ({ sessionShow: async () => ({ id: "open", status: "open" }) }) } as any,
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -94,7 +97,7 @@ test("maintenance rejects new and resumed session admission without interrupting
|
|||||||
expect(resume.json()).toEqual(create.json());
|
expect(resume.json()).toEqual(create.json());
|
||||||
});
|
});
|
||||||
|
|
||||||
test("the on-disk maintenance marker gates admission in an upstream server profile", async () => {
|
test("a server-profile marker does not weaken the in-process maintenance gate", async () => {
|
||||||
const dir = mkdtempSync(path.join(tmpdir(), "tht-maintenance-"));
|
const dir = mkdtempSync(path.join(tmpdir(), "tht-maintenance-"));
|
||||||
try {
|
try {
|
||||||
const marker = path.join(dir, "maintenance.json");
|
const marker = path.join(dir, "maintenance.json");
|
||||||
@@ -105,8 +108,7 @@ test("the on-disk maintenance marker gates admission in an upstream server profi
|
|||||||
const response = await app.inject({
|
const response = await app.inject({
|
||||||
method: "POST", url: "/sessions", headers: aliceHeaders, payload: { question: "q" },
|
method: "POST", url: "/sessions", headers: aliceHeaders, payload: { question: "q" },
|
||||||
});
|
});
|
||||||
expect(response.statusCode).toBe(503);
|
expect(response.statusCode).not.toBe(503);
|
||||||
expect(response.json().code).toBe("maintenance");
|
|
||||||
} finally {
|
} finally {
|
||||||
rmSync(dir, { recursive: true, force: true });
|
rmSync(dir, { recursive: true, force: true });
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -172,7 +172,9 @@ func piCommand(ctx context.Context, installation config.Installation, runner com
|
|||||||
if err := pi.Configure(ctx, controlled, defaults); err != nil {
|
if err := pi.Configure(ctx, controlled, defaults); err != nil {
|
||||||
return piFailure(stderr, err, secretValues)
|
return piFailure(stderr, err, secretValues)
|
||||||
}
|
}
|
||||||
fmt.Fprintln(stdout, "Pi defaults applied and read back. Put credentials only in PI_AUTH_FILE (/home/thoth/.pi/agent/auth.json, mode 0600); never pass credentials to thothctl.")
|
authFile, _ := installation.EnvironmentValue("PI_AUTH_FILE")
|
||||||
|
if authFile == "" { authFile = "the host path declared by PI_AUTH_FILE" }
|
||||||
|
fmt.Fprintf(stdout, "Pi defaults applied and read back. Put credentials only in PI_AUTH_FILE=%s (mode 0600); expected variables/files are PI_AUTH_FILE and /home/thoth/.pi/agent/auth.json. Never pass credentials to thothctl.\n", authFile)
|
||||||
return 0
|
return 0
|
||||||
case "update":
|
case "update":
|
||||||
request, err := parsePiUpdateArgs(args[1:], filepath.Join(installation.ProjectDirectory, ".thothctl", "update-state.json"))
|
request, err := parsePiUpdateArgs(args[1:], filepath.Join(installation.ProjectDirectory, ".thothctl", "update-state.json"))
|
||||||
|
|||||||
@@ -164,6 +164,16 @@ func (i Installation) SecretFiles() ([]string, error) {
|
|||||||
return files, nil
|
return files, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// EnvironmentValue returns one declared installation value without exposing dotenv parsing to
|
||||||
|
// callers. It is used only for operator-visible file locations, never for secret content.
|
||||||
|
func (i Installation) EnvironmentValue(name string) (string, error) {
|
||||||
|
contents, err := safeio.ReadCanonicalRegular(i.EnvFile, maxEnvironmentFileBytes)
|
||||||
|
if err != nil { return "", errors.New("installation environment could not be read") }
|
||||||
|
values, err := parseComposeDotenv(contents)
|
||||||
|
if err != nil { return "", errors.New("installation environment could not be read") }
|
||||||
|
return values[name], nil
|
||||||
|
}
|
||||||
|
|
||||||
func parseComposeDotenv(contents []byte) (map[string]string, error) {
|
func parseComposeDotenv(contents []byte) (map[string]string, error) {
|
||||||
dotenvParseMu.Lock()
|
dotenvParseMu.Lock()
|
||||||
defer dotenvParseMu.Unlock()
|
defer dotenvParseMu.Unlock()
|
||||||
|
|||||||
@@ -30,7 +30,7 @@ var internalIdentityHeaders = []string{
|
|||||||
|
|
||||||
// Configure changes the backend's real installation settings through a core-side helper. It
|
// 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.
|
// deliberately has no secret or endpoint input: external endpoints remain Compose-owned.
|
||||||
func Configure(ctx context.Context, runner Runner, value Defaults) error {
|
func Configure(ctx context.Context, runner Runner, value Defaults) (retErr error) {
|
||||||
if !choicePattern.MatchString(value.Provider) || !choicePattern.MatchString(value.Model) {
|
if !choicePattern.MatchString(value.Provider) || !choicePattern.MatchString(value.Model) {
|
||||||
return errors.New("provider and model must be supported identifiers")
|
return errors.New("provider and model must be supported identifiers")
|
||||||
}
|
}
|
||||||
@@ -63,10 +63,22 @@ func Configure(ctx context.Context, runner Runner, value Defaults) error {
|
|||||||
if !found {
|
if !found {
|
||||||
return errors.New("provider/model is not in Pi options")
|
return errors.New("provider/model is not in Pi options")
|
||||||
}
|
}
|
||||||
result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/settings/settings-cli.js", "--provider", value.Provider, "--model", value.Model, "--thinking", value.Thinking)
|
|
||||||
if err != nil { return commandError("Pi installation settings write", result, err) }
|
|
||||||
settingsArgs := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...)
|
settingsArgs := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...)
|
||||||
settingsArgs = append(settingsArgs, "http://127.0.0.1:8787/settings")
|
settingsArgs = append(settingsArgs, "http://127.0.0.1:8787/settings")
|
||||||
|
oldResult, err := runCompose(ctx, runner, settingsArgs...)
|
||||||
|
if err != nil { return commandError("Pi installation settings capture", oldResult, err) }
|
||||||
|
var old Defaults
|
||||||
|
if json.Unmarshal([]byte(oldResult.Stdout), &old) != nil || old.Provider == "" || old.Model == "" || old.Thinking == "" { return errors.New("Pi installation settings capture is invalid") }
|
||||||
|
wrote := false
|
||||||
|
defer func() {
|
||||||
|
if retErr != nil && wrote {
|
||||||
|
result, restoreErr := runCompose(context.Background(), runner, "exec", "-T", "core", "node", "/app/backend/dist/settings/settings-cli.js", "--provider", old.Provider, "--model", old.Model, "--thinking", old.Thinking)
|
||||||
|
if restoreErr != nil || result.ExitCode != 0 { retErr = fmt.Errorf("%w; previous Pi settings could not be restored: recovery required", retErr) }
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
result, err := runCompose(ctx, runner, "exec", "-T", "core", "node", "/app/backend/dist/settings/settings-cli.js", "--provider", value.Provider, "--model", value.Model, "--thinking", value.Thinking)
|
||||||
|
if err != nil { return commandError("Pi installation settings write", result, err) }
|
||||||
|
wrote = true
|
||||||
settings, err := runCompose(ctx, runner, settingsArgs...)
|
settings, err := runCompose(ctx, runner, settingsArgs...)
|
||||||
if err != nil { return commandError("Pi installation settings read-back", settings, err) }
|
if err != nil { return commandError("Pi installation settings read-back", settings, err) }
|
||||||
var saved Defaults
|
var saved Defaults
|
||||||
|
|||||||
@@ -79,6 +79,9 @@ func readState(path string) (State, error) {
|
|||||||
if state.Version != stateFileVersion || state.Previous.ID == "" || state.Previous.Reference == "" || state.Previous.MountFingerprint == "" {
|
if state.Version != stateFileVersion || state.Previous.ID == "" || state.Previous.Reference == "" || state.Previous.MountFingerprint == "" {
|
||||||
return State{}, errors.New("update recovery state is incomplete")
|
return State{}, errors.New("update recovery state is incomplete")
|
||||||
}
|
}
|
||||||
|
if mountFingerprint(state.Previous.Mounts) != state.Previous.MountFingerprint || (state.Candidate.ID != "" && mountFingerprint(state.Candidate.Mounts) != state.Candidate.MountFingerprint) {
|
||||||
|
return State{}, errors.New("update recovery state mount fingerprint is invalid")
|
||||||
|
}
|
||||||
return state, nil
|
return state, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -144,9 +147,10 @@ func acquireLock(statePath string) (*updateLock, error) {
|
|||||||
return nil, errors.New("could not create Pi update recovery directory")
|
return nil, errors.New("could not create Pi update recovery directory")
|
||||||
}
|
}
|
||||||
path := statePath + ".lock"
|
path := statePath + ".lock"
|
||||||
if err := os.Mkdir(path, 0o700); err != nil {
|
file, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o600)
|
||||||
|
if err != nil {
|
||||||
if errors.Is(err, os.ErrExist) {
|
if errors.Is(err, os.ErrExist) {
|
||||||
if reclaimDeadLocalLock(path) {
|
if reclaimDeadLocalLock(path, statePath) {
|
||||||
return acquireLock(statePath)
|
return acquireLock(statePath)
|
||||||
}
|
}
|
||||||
return nil, ErrLockHeld
|
return nil, ErrLockHeld
|
||||||
@@ -154,28 +158,36 @@ func acquireLock(statePath string) (*updateLock, error) {
|
|||||||
return nil, errors.New("could not acquire Pi update lock")
|
return nil, errors.New("could not acquire Pi update lock")
|
||||||
}
|
}
|
||||||
host, err := os.Hostname()
|
host, err := os.Hostname()
|
||||||
if err != nil { _ = os.Remove(path); return nil, errors.New("could not identify Pi update lock owner") }
|
if err != nil { _ = file.Close(); _ = os.Remove(path); return nil, errors.New("could not identify Pi update lock owner") }
|
||||||
owner := lockOwner{PID: os.Getpid(), Host: host, StartedAt: time.Now().UTC(), Transaction: fmt.Sprintf("%d-%d", os.Getpid(), time.Now().UnixNano())}
|
owner := lockOwner{PID: os.Getpid(), Host: host, StartedAt: time.Now().UTC(), Transaction: fmt.Sprintf("%d-%d", os.Getpid(), time.Now().UnixNano())}
|
||||||
contents, err := json.Marshal(owner)
|
contents, err := json.Marshal(owner)
|
||||||
if err != nil { _ = os.Remove(path); return nil, errors.New("could not record Pi update lock owner") }
|
if err != nil { _ = file.Close(); _ = os.Remove(path); return nil, errors.New("could not record Pi update lock owner") }
|
||||||
if err := writeFileDurably(filepath.Join(path, "owner.json"), ".owner-", append(contents, '\n')); err != nil {
|
if _, err := file.Write(append(contents, '\n')); err != nil || file.Sync() != nil || file.Close() != nil {
|
||||||
_ = os.Remove(path)
|
_ = file.Close(); _ = os.Remove(path)
|
||||||
return nil, errors.New("could not record Pi update lock owner")
|
return nil, errors.New("could not record Pi update lock owner")
|
||||||
}
|
}
|
||||||
return &updateLock{path: path}, nil
|
return &updateLock{path: path}, nil
|
||||||
}
|
}
|
||||||
func (l *updateLock) Release() { _ = os.Remove(filepath.Join(l.path, "owner.json")); _ = os.Remove(l.path) }
|
func (l *updateLock) Release() { _ = os.Remove(l.path) }
|
||||||
|
|
||||||
// reclaimDeadLocalLock is deliberately conservative: a malformed, remote, or merely old lock
|
// reclaimDeadLocalLock is deliberately conservative: a malformed, remote, or merely old lock
|
||||||
// is recovery-required. Only a process we can prove is gone on this machine is reclaimed.
|
// is recovery-required. Only a process we can prove is gone on this machine is reclaimed.
|
||||||
func reclaimDeadLocalLock(path string) bool {
|
func reclaimDeadLocalLock(path, statePath string) bool {
|
||||||
contents, err := os.ReadFile(filepath.Join(path, "owner.json"))
|
info, err := os.Stat(path)
|
||||||
if err != nil { return false }
|
if err != nil || time.Since(info.ModTime()) < 5*time.Minute || !hasPendingRecoveryState(statePath) { return false }
|
||||||
|
contents, err := os.ReadFile(path)
|
||||||
|
if err != nil { return os.Remove(path) == nil }
|
||||||
var owner lockOwner
|
var owner lockOwner
|
||||||
if json.Unmarshal(contents, &owner) != nil || owner.PID <= 0 || owner.Host == "" { return false }
|
if json.Unmarshal(contents, &owner) != nil || owner.PID <= 0 || owner.Host == "" { return false }
|
||||||
host, err := os.Hostname()
|
host, err := os.Hostname()
|
||||||
if err != nil || owner.Host != host { return false }
|
if err != nil || owner.Host != host { return false }
|
||||||
if processAlive(owner.PID) { return false }
|
if processAlive(owner.PID) { return false }
|
||||||
if err := os.Remove(filepath.Join(path, "owner.json")); err != nil { return false }
|
|
||||||
return os.Remove(path) == nil
|
return os.Remove(path) == nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func hasPendingRecoveryState(path string) bool {
|
||||||
|
contents, err := os.ReadFile(path); if err != nil { return false }
|
||||||
|
var state State
|
||||||
|
if json.Unmarshal(contents, &state) != nil { return false }
|
||||||
|
return state.Phase != PhaseVerified && state.Phase != PhaseRolledBack && state.Phase != PhaseNoop
|
||||||
|
}
|
||||||
|
|||||||
@@ -166,7 +166,7 @@ func Update(ctx context.Context, runner Runner, request Request) (result Result,
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Rollback restores the image recorded in durable update state. It is safe for interrupted runs.
|
// Rollback restores the image recorded in durable update state. It is safe for interrupted runs.
|
||||||
func Rollback(ctx context.Context, runner Runner, statePath string, confirm bool) (Result, error) {
|
func Rollback(ctx context.Context, runner Runner, statePath string, confirm bool) (result Result, retErr error) {
|
||||||
lock, err := acquireLock(statePath)
|
lock, err := acquireLock(statePath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Result{StatePath: statePath}, err
|
return Result{StatePath: statePath}, err
|
||||||
@@ -175,6 +175,16 @@ func Rollback(ctx context.Context, runner Runner, statePath string, confirm bool
|
|||||||
if !confirm {
|
if !confirm {
|
||||||
return Result{StatePath: statePath}, ErrConfirmationRequired
|
return Result{StatePath: statePath}, ErrConfirmationRequired
|
||||||
}
|
}
|
||||||
|
if err := setMaintenance(ctx, runner, true); err != nil { return Result{StatePath: statePath}, err }
|
||||||
|
defer func() {
|
||||||
|
if clearErr := setMaintenance(context.Background(), runner, false); clearErr != nil {
|
||||||
|
result = Result{Phase: PhaseFailed, StatePath: statePath}
|
||||||
|
if retErr == nil { retErr = errors.New("maintenance admission gate could not be cleared: recovery required")
|
||||||
|
} else { retErr = fmt.Errorf("%w; maintenance admission gate could not be cleared: recovery required", retErr) }
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
if active, err := activeSessions(ctx, runner); err != nil { return Result{StatePath: statePath}, err
|
||||||
|
} else if active { return Result{StatePath: statePath}, ErrActiveSessions }
|
||||||
state, err := readState(statePath)
|
state, err := readState(statePath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Result{StatePath: statePath}, err
|
return Result{StatePath: statePath}, err
|
||||||
@@ -232,30 +242,34 @@ func canonicalDigestReference(value string) (string, error) {
|
|||||||
return reference.FamiliarString(canonical), nil
|
return reference.FamiliarString(canonical), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// The command text is fixed; no operator input or host path is interpolated into the core shell.
|
|
||||||
// The marker lives alongside SETTINGS_FILE's named/bind-mounted directory and is read by backend.
|
|
||||||
func setMaintenance(ctx context.Context, runner Runner, enabled bool) error {
|
func setMaintenance(ctx context.Context, runner Runner, enabled bool) error {
|
||||||
command := "mkdir -p /data/settings && : > /data/settings/maintenance.json && chmod 600 /data/settings/maintenance.json"
|
path := "deactivate"
|
||||||
if !enabled { command = "rm -f /data/settings/maintenance.json" }
|
if enabled { path = "activate" }
|
||||||
result, err := runCompose(ctx, runner, "exec", "-T", "core", "sh", "-ceu", command)
|
args := append([]string{"exec", "-T", "core", "curl", "-fsS", "-X", "POST"}, internalIdentityHeaders...)
|
||||||
|
args = append(args, "http://127.0.0.1:8787/internal/maintenance/"+path)
|
||||||
|
result, err := runCompose(ctx, runner, args...)
|
||||||
if err != nil { return commandError("maintenance admission gate", result, err) }
|
if err != nil { return commandError("maintenance admission gate", result, err) }
|
||||||
|
var status struct { Active bool `json:"active"`; Admissions int `json:"admissions"` }
|
||||||
|
if json.Unmarshal([]byte(result.Stdout), &status) != nil || status.Active != enabled || status.Admissions != 0 { return errors.New("maintenance admission gate did not acknowledge a quiescent state") }
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func activeSessions(ctx context.Context, runner Runner) (bool, error) {
|
func activeSessions(ctx context.Context, runner Runner) (bool, error) {
|
||||||
result, err := runCompose(ctx, runner, "exec", "-T", "core", "tht", "session", "list", "--json")
|
args := append([]string{"exec", "-T", "core", "curl", "-fsS"}, internalIdentityHeaders...)
|
||||||
|
args = append(args, "http://127.0.0.1:8787/sessions?scope=all")
|
||||||
|
result, err := runCompose(ctx, runner, args...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return false, commandError("active-session check", result, err)
|
return false, commandError("active-session check", result, err)
|
||||||
}
|
}
|
||||||
var sessions []struct {
|
var payload struct { Sessions []struct {
|
||||||
Status string `json:"status"`
|
Status string `json:"status"`
|
||||||
Archived bool `json:"archived"`
|
Archived bool `json:"archived"`
|
||||||
}
|
} `json:"sessions"` }
|
||||||
if err := json.Unmarshal([]byte(result.Stdout), &sessions); err != nil {
|
if err := json.Unmarshal([]byte(result.Stdout), &payload); err != nil {
|
||||||
return false, errors.New("active-session check returned invalid session data")
|
return false, errors.New("active-session check returned invalid session data")
|
||||||
}
|
}
|
||||||
for _, session := range sessions {
|
for _, session := range payload.Sessions {
|
||||||
if !session.Archived && session.Status == "open" {
|
if !session.Archived && session.Status != "finalized" && session.Status != "closed" {
|
||||||
return true, nil
|
return true, nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -132,7 +132,9 @@ func TestUpdateRequiresConfirmationAndDrainsActiveSessions(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("Update() with drain error = %v", err)
|
t.Fatalf("Update() with drain error = %v", err)
|
||||||
}
|
}
|
||||||
assertCalled(t, fake.calls, "compose exec -T core tht session list --json")
|
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")
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestRollbackRestoresInterruptedOrPreviouslyRecordedState(t *testing.T) {
|
func TestRollbackRestoresInterruptedOrPreviouslyRecordedState(t *testing.T) {
|
||||||
@@ -162,7 +164,7 @@ func TestRollbackRestoresInterruptedOrPreviouslyRecordedState(t *testing.T) {
|
|||||||
func TestUpdateRefusesToOverwriteInterruptedRecoveryState(t *testing.T) {
|
func TestUpdateRefusesToOverwriteInterruptedRecoveryState(t *testing.T) {
|
||||||
fake := newFakeRunner()
|
fake := newFakeRunner()
|
||||||
statePath := filepath.Join(t.TempDir(), "state.json")
|
statePath := filepath.Join(t.TempDir(), "state.json")
|
||||||
writeStateForTest(t, statePath, State{Phase: PhaseRecreated, Previous: Image{ID: "sha256:old", Reference: "thothii-core:local", Volumes: []string{"settings"}, MountFingerprint: "recorded"}})
|
writeStateForTest(t, statePath, State{Phase: PhaseRecreated, Previous: Image{ID: "sha256:old", Reference: "thothii-core:local", Volumes: []string{"settings"}, MountFingerprint: mountFingerprint(nil)}})
|
||||||
_, err := Update(context.Background(), fake, Request{StatePath: statePath, Version: "0.81.0", Source: BuildSource, Confirm: true})
|
_, err := Update(context.Background(), fake, Request{StatePath: statePath, Version: "0.81.0", Source: BuildSource, Confirm: true})
|
||||||
if !errors.Is(err, ErrInterruptedUpdate) {
|
if !errors.Is(err, ErrInterruptedUpdate) {
|
||||||
t.Fatalf("Update() error = %v, want interrupted update error", err)
|
t.Fatalf("Update() error = %v, want interrupted update error", err)
|
||||||
@@ -242,12 +244,16 @@ func (f *fakeRunner) Run(_ context.Context, args []string, _ io.Reader) (compose
|
|||||||
return compose.Result{Stdout: f.mountsJSON}, nil
|
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
|
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, "tht session list --json"):
|
case strings.Contains(call, "/internal/maintenance/activate"):
|
||||||
|
return compose.Result{Stdout: `{"active":true,"admissions":0}`}, nil
|
||||||
|
case strings.Contains(call, "/internal/maintenance/deactivate"):
|
||||||
|
return compose.Result{Stdout: `{"active":false,"admissions":0}`}, nil
|
||||||
|
case strings.Contains(call, "/sessions?scope=all"):
|
||||||
if f.activeSessions {
|
if f.activeSessions {
|
||||||
f.activeSessions = false
|
f.activeSessions = false
|
||||||
return compose.Result{Stdout: `[{"status":"open","archived":false}]`}, nil
|
return compose.Result{Stdout: `{"sessions":[{"status":"open","archived":false}]}`}, nil
|
||||||
}
|
}
|
||||||
return compose.Result{Stdout: "[]"}, nil
|
return compose.Result{Stdout: `{"sessions":[]}`}, nil
|
||||||
case strings.Contains(call, "compose build"):
|
case strings.Contains(call, "compose build"):
|
||||||
f.built = true
|
f.built = true
|
||||||
f.version = "0.81.0"
|
f.version = "0.81.0"
|
||||||
|
|||||||
Reference in New Issue
Block a user