fix: harden workspace activation and snapshot retention

This commit is contained in:
2026-08-04 09:25:06 +02:00
parent 3b23cf3714
commit e4fdbed864
14 changed files with 474 additions and 103 deletions
+11
View File
@@ -18,6 +18,8 @@ import { ReadinessManager } from "./runtime/readiness-manager.js";
import { WorkspaceRegistry } from "./workspaces/registry.js";
import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js";
import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js";
import { resolveRuntimeBindings, supportsSessionRuntime } from "./workspaces/bindings.js";
import type { WorkspaceDescriptor } from "./workspaces/schema.js";
export interface BuildAppDeps {
thtRunner?: ThtRunner;
@@ -29,6 +31,7 @@ export interface BuildAppDeps {
hub?: SseHub;
workspaceRegistry?: WorkspaceRegistry;
workspaceDiagnoser?: WorkspaceDiagnoser;
workspaceRuntimeSupport?: (workspace: WorkspaceDescriptor) => boolean;
}
export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance {
@@ -55,6 +58,13 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
const workspaceRegistry = deps?.workspaceRegistry ?? new WorkspaceRegistry(config.workspaceRegistry);
const workspaceDiagnoser = deps?.workspaceDiagnoser
?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs);
const workspaceRuntimeSupport = deps?.workspaceRuntimeSupport ?? ((workspace: WorkspaceDescriptor) => (
supportsSessionRuntime(resolveRuntimeBindings(
workspace,
process.env,
config.workspaceRegistry.secretRoots,
))
));
const readiness = deps?.readiness ?? new ReadinessManager(
tht as ThtRunner,
Math.round(config.ollamaEnsureTimeoutMs / 1000),
@@ -88,6 +98,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
mgr, tht: tht as ThtRunner, hub, getSettings, readiness, listModels, workspaceRegistry,
dwhPrecheck: config.dwhPrecheck,
legacyWorkspaceMode: config.legacyWorkspaceMode,
workspaceRuntimeSupport,
});
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
+127 -90
View File
@@ -8,6 +8,7 @@ import type { PrincipalContext } from "../auth/principal.js";
import type { ReadinessManager } from "../runtime/readiness-manager.js";
import type { ListModelsFn } from "./meta.js";
import type { WorkspaceRegistry } from "../workspaces/registry.js";
import type { WorkspaceDescriptor } from "../workspaces/schema.js";
const BOOTSTRAP_FAILURE_MESSAGE =
"Session startup failed. Check configuration and connectivity, then Resume the session.";
@@ -34,6 +35,8 @@ export function sessionRoutes(
dwhPrecheck?: boolean;
/** Explicit loopback-only compatibility path for old clients that send `workspace`. */
legacyWorkspaceMode?: boolean;
/** Fail-closed installation/runtime transport capability check. */
workspaceRuntimeSupport: (workspace: WorkspaceDescriptor) => boolean;
},
) {
const lifecycleTails = new Map<string, Promise<void>>();
@@ -289,110 +292,144 @@ export function sessionRoutes(
code: "workspace_revision_unavailable",
});
}
let workspaceConfigPath: string | undefined;
let workspaceId: string | undefined;
let workspaceRevision: string | undefined;
let allowedModels: readonly string[] | undefined;
if (requestedWorkspaceId) {
try {
const resolved = await d.workspaceRegistry.read(requestedWorkspaceId);
if (resolved.revision.state !== "operational") {
let revisionLease: Awaited<ReturnType<WorkspaceRegistry["acquireSessionRevision"]>> | undefined;
let manifestPersisted = false;
try {
let workspaceConfigPath: string | undefined;
let workspaceId: string | undefined;
let workspaceRevision: string | undefined;
let allowedModels: readonly string[] | undefined;
if (requestedWorkspaceId) {
try {
const registry = d.workspaceRegistry as Partial<WorkspaceRegistry>;
const resolved = typeof registry.acquireSessionRevision === "function"
? await registry.acquireSessionRevision.call(d.workspaceRegistry, requestedWorkspaceId)
: await d.workspaceRegistry.read(requestedWorkspaceId);
if ("markPersisted" in resolved && "abort" in resolved) {
revisionLease = resolved as Awaited<ReturnType<WorkspaceRegistry["acquireSessionRevision"]>>;
}
if (resolved.revision.state !== "operational") {
return reply.code(409).send({
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
code: "workspace_revision_unavailable",
});
}
if (!d.workspaceRuntimeSupport(resolved.workspace)) {
return reply.code(409).send({
error: "This workspace transport is not available to runtime sessions.",
code: "workspace_not_activatable",
});
}
workspaceConfigPath = resolved.revision.snapshotPath;
workspaceId = resolved.revision.id;
workspaceRevision = resolved.revision.commit;
allowedModels = resolved.workspace.llm_policy.allowed;
} catch {
return reply.code(409).send({
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
code: "workspace_revision_unavailable",
});
}
workspaceConfigPath = resolved.revision.snapshotPath;
workspaceId = resolved.revision.id;
workspaceRevision = resolved.revision.commit;
allowedModels = resolved.workspace.llm_policy.allowed;
} catch {
return reply.code(409).send({
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
code: "workspace_revision_unavailable",
});
}
}
const provider = b.provider ?? s.provider;
const model = b.model ?? s.model;
const thinking = b.thinking ?? s.thinking;
if (allowedModels && provider && model && !allowedModels.includes(`${provider}/${model}`)) {
return reply.code(400).send({ error: "Selected model is not allowed by this workspace." });
}
// A persisted session is resumable without keeping Pi alive. New work replaces every
// runtime owned by this principal, while runtimes belonging to other users remain intact.
// Optional chaining preserves the deliberately narrow manager stubs used by route tests.
for (const id of d.mgr.teardownForPrincipal?.(principal) ?? []) boundRuntimes.delete(id);
const ensure = await d.readiness.ensure(workspaceConfigPath ?? "", principal);
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
// Local-only: verify the DWH is reachable BEFORE creating the session, so a dropped
// VPN surfaces as an up-front alert instead of a session that spawns Pi and then dies
// in bootstrap retrieval. `code` lets the client show a specific message.
if (d.dwhPrecheck) {
const ping = await runner.dbPing(workspaceConfigPath);
if (!ping.ok) {
console.error(`[dwh-precheck] refusing new session — DWH unreachable: ${ping.detail}`);
return reply.code(503).send({ error: DWH_UNREACHABLE_MESSAGE, code: "dwh_unreachable" });
const provider = b.provider ?? s.provider;
const model = b.model ?? s.model;
const thinking = b.thinking ?? s.thinking;
if (allowedModels && provider && model && !allowedModels.includes(`${provider}/${model}`)) {
return reply.code(400).send({ error: "Selected model is not allowed by this workspace." });
}
}
if (provider && model) {
let available: Awaited<ReturnType<ListModelsFn>>;
// A persisted session is resumable without keeping Pi alive. New work replaces every
// runtime owned by this principal, while runtimes belonging to other users remain intact.
// Optional chaining preserves the deliberately narrow manager stubs used by route tests.
for (const id of d.mgr.teardownForPrincipal?.(principal) ?? []) boundRuntimes.delete(id);
const ensure = await d.readiness.ensure(workspaceConfigPath ?? "", principal);
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
// Local-only: verify the DWH is reachable BEFORE creating the session, so a dropped
// VPN surfaces as an up-front alert instead of a session that spawns Pi and then dies
// in bootstrap retrieval. `code` lets the client show a specific message.
if (d.dwhPrecheck) {
const ping = await runner.dbPing(workspaceConfigPath);
if (!ping.ok) {
console.error(`[dwh-precheck] refusing new session — DWH unreachable: ${ping.detail}`);
return reply.code(503).send({ error: DWH_UNREACHABLE_MESSAGE, code: "dwh_unreachable" });
}
}
if (provider && model) {
let available: Awaited<ReturnType<ListModelsFn>>;
try {
available = await d.listModels();
} catch {
return reply.code(503).send({
error: MODEL_UNAVAILABLE_MESSAGE,
code: "model_unavailable",
});
}
const selectedAvailable = available.some(
(candidate) => candidate.provider === provider && candidate.id === model,
);
if (!selectedAvailable) {
return reply.code(503).send({
error: MODEL_UNAVAILABLE_MESSAGE,
code: "model_unavailable",
});
}
}
// Browser choices are copied to the persisted manifest together with the immutable
// registry snapshot. The legacy fallback stays available for sessions created before
// the browser-local preference migration.
let id: string;
try {
available = await d.listModels();
} catch {
return reply.code(503).send({
error: MODEL_UNAVAILABLE_MESSAGE,
code: "model_unavailable",
({ id } = await runner.sessionNew({
question: b.question, name: b.name, workspaceConfigPath,
workspaceId, workspaceRevision, provider, model, thinking,
}));
manifestPersisted = true;
if (revisionLease) {
await revisionLease.markPersisted().catch((error: unknown) => {
console.error(
`[session:${id}] revision lease hand-off failed:`,
error instanceof Error ? error.message : "unknown error",
);
});
}
} catch { return storageFailure(reply); }
const options = {
provider, model, thinking,
author: principal.displayName ?? principal.subject,
principal,
question: b.question,
};
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
try {
rt = d.mgr.createFor(id, options);
bindRuntime(id, rt, runner, workspaceConfigPath);
} catch (error) {
if (rt) d.mgr.teardownIfCurrent(id, rt);
console.error(
`[pi:${id}] runtime construction failed:`,
error instanceof Error ? error.message : "unknown error",
);
await runner.failSession(id, workspaceConfigPath).catch((persistenceError: unknown) => {
console.error(`[session:${id}] failSession persistence failed:`, persistenceError);
});
return reply.code(503).send({ error: BOOTSTRAP_FAILURE_MESSAGE });
}
const selectedAvailable = available.some(
(candidate) => candidate.provider === provider && candidate.id === model,
info(id, "Session created");
bootstrap(
id, rt, runner, workspaceConfigPath, d.mgr.configure(rt, options),
runner.searchPack(b.question, id, workspaceConfigPath),
() => d.mgr.start(id, rt, options),
);
if (!selectedAvailable) {
return reply.code(503).send({
error: MODEL_UNAVAILABLE_MESSAGE,
code: "model_unavailable",
return { id };
} finally {
if (revisionLease && !manifestPersisted) {
await revisionLease.abort().catch((error: unknown) => {
console.error(
"[session] revision lease cleanup failed:",
error instanceof Error ? error.message : "unknown error",
);
});
}
}
// Browser choices are copied to the persisted manifest together with the immutable
// registry snapshot. The legacy fallback stays available for sessions created before
// the browser-local preference migration.
let id: string;
try {
({ id } = await runner.sessionNew({
question: b.question, name: b.name, workspaceConfigPath,
workspaceId, workspaceRevision, provider, model, thinking,
}));
} catch { return storageFailure(reply); }
const options = {
provider, model, thinking,
author: principal.displayName ?? principal.subject,
principal,
question: b.question,
};
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
try {
rt = d.mgr.createFor(id, options);
bindRuntime(id, rt, runner, workspaceConfigPath);
} catch (error) {
if (rt) d.mgr.teardownIfCurrent(id, rt);
console.error(
`[pi:${id}] runtime construction failed:`,
error instanceof Error ? error.message : "unknown error",
);
await runner.failSession(id, workspaceConfigPath).catch((persistenceError: unknown) => {
console.error(`[session:${id}] failSession persistence failed:`, persistenceError);
});
return reply.code(503).send({ error: BOOTSTRAP_FAILURE_MESSAGE });
}
info(id, "Session created");
bootstrap(
id, rt, runner, workspaceConfigPath, d.mgr.configure(rt, options),
runner.searchPack(b.question, id, workspaceConfigPath),
() => d.mgr.start(id, rt, options),
);
return { id };
});
app.get("/sessions", async (req, reply) => {
const principal = getPrincipal(req);
+8
View File
@@ -143,3 +143,11 @@ export function resolveRuntimeBindings(
embedding: resolveBinding(workspace, "EMBEDDING", env, secretRoots),
};
}
/**
* SSH bindings are currently probe-only: diagnostics owns a short-lived tunnel, while the
* session runtime has no tunnel owner. Keep activation fail-closed until that lifecycle exists.
*/
export function supportsSessionRuntime(bindings: RuntimeBindings): boolean {
return bindings.dwh.transport !== "ssh_tunnel" && bindings.vector.transport !== "ssh_tunnel";
}
+13 -1
View File
@@ -535,7 +535,9 @@ function diagnosticError(code: WorkspaceErrorCode, field?: string): Diagnostic {
? "Installation binding is missing or invalid."
: code === "semantic_index_incompatible"
? "Semantic index metadata is incompatible with this workspace."
: "Connector diagnostic failed.",
: code === "workspace_not_activatable"
? "This transport can be tested, but it is not available to runtime sessions."
: "Connector diagnostic failed.",
};
}
@@ -851,6 +853,16 @@ export function createWorkspaceDiagnoser(
}
}
// The concrete SSH adapter deliberately owns only a bounded diagnostic tunnel and closes it
// in `finally`. Until a session runtime owns an equivalent long-lived tunnel, a successful
// probe is connectivity evidence only and must never be advertised as activatable.
if (
(bindings.dwh.transport === "ssh_tunnel" || bindings.vector.transport === "ssh_tunnel")
&& !diagnostics.some((diagnostic) => diagnostic.level === "error")
) {
diagnostics.push(diagnosticError("workspace_not_activatable"));
}
return {
activatable: !diagnostics.some((diagnostic) => diagnostic.level === "error"),
diagnostics,
+139 -1
View File
@@ -28,6 +28,15 @@ export interface WorkspaceRevision {
state: "operational" | "migration_required";
}
export interface SessionRevisionLease {
workspace: WorkspaceDescriptor;
revision: WorkspaceRevision;
/** Mark the manifest durable; retention removes the lease only after observing that manifest. */
markPersisted(): Promise<void>;
/** Remove a lease for a session that failed before its manifest was durable. */
abort(): Promise<void>;
}
export type PublishWorkspaceRequest =
| { action: "create"; workspace: CanonicalWorkspace; baseCommit: string }
| { action: "update"; workspace: CanonicalWorkspace; baseCommit: string; baseBlob: string }
@@ -56,6 +65,14 @@ interface SnapshotManifest extends ActiveState {
files: Record<string, string>;
}
interface RevisionLeaseRecord {
version: 1;
token: string;
workspaceId: string;
commit: string;
state: "creating" | "persisted";
}
type LegacyWorkspaceRevision = Omit<WorkspaceRevision, "state">;
interface LegacyActiveState {
@@ -178,6 +195,56 @@ export class WorkspaceRegistry {
}
}
/**
* Resolve the active revision and create its cross-process retention lease under the same
* repository lock. The lease bridges the interval before `session_manifest.yaml` is durable.
*/
async acquireSessionRevision(id: string): Promise<SessionRevisionLease> {
await this.repository.ensureLayout();
return await this.lock.run(async () => {
const state = await this.activeState();
const revision = state.revisions.find((candidate) => candidate.id === id);
if (!revision) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable");
let workspace: WorkspaceDescriptor;
try {
workspace = parseWorkspaceYaml(await readFile(revision.snapshotPath, "utf8"));
} catch (error) {
throw workspaceError(error);
}
const token = randomUUID();
const record: RevisionLeaseRecord = {
version: 1,
token,
workspaceId: id,
commit: revision.commit,
state: "creating",
};
const path = await this.writeRevisionLease(record, true);
let localState: RevisionLeaseRecord["state"] | "aborted" = "creating";
return {
workspace,
revision,
markPersisted: async () => {
if (localState === "persisted") return;
if (localState === "aborted") throw new WorkspaceRegistryError(
"workspace_invalid", "Workspace revision lease is unavailable",
);
await this.lock.run(async () => {
await this.replaceRevisionLease(path, { ...record, state: "persisted" });
});
localState = "persisted";
},
abort: async () => {
if (localState !== "creating") return;
await this.lock.run(async () => { await rm(path, { force: true }); });
localState = "aborted";
},
};
});
}
/** Read a retained immutable snapshot for a session pinned to a historical commit. */
async readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }> {
const snapshotPath = this.snapshotPath(safeCommit(commit), id);
@@ -195,9 +262,12 @@ export class WorkspaceRegistry {
* a partial, per-user list could otherwise remove another user's resumable workspace pin.
*/
async reconcileSnapshotRetention(referencedCommits: readonly string[]): Promise<void> {
const retained = new Set(referencedCommits.map(safeCommit));
const manifestReferences = new Set(referencedCommits.map(safeCommit));
const retained = new Set(manifestReferences);
await this.repository.ensureLayout();
await this.lock.run(async () => {
const leases = await this.revisionLeases();
for (const { record } of leases) retained.add(record.commit);
retained.add((await this.activeState()).head);
const entries = await readdir(this.repository.snapshotsPath, { withFileTypes: true });
for (const entry of entries) {
@@ -210,9 +280,77 @@ export class WorkspaceRegistry {
if (!current.isDirectory() || current.isSymbolicLink()) continue;
await rm(path, { recursive: true, force: true });
}
// A persisted lease is handed off only when this exact authoritative scan has observed a
// manifest pin for its commit. A stale scan therefore keeps the lease and cannot prune it.
for (const { path, record } of leases) {
if (record.state === "persisted" && manifestReferences.has(record.commit)) {
await rm(path, { force: true });
}
}
});
}
private revisionLeaseDirectory(): string {
return join(this.repository.statePath, "revision-leases");
}
private async writeRevisionLease(record: RevisionLeaseRecord, exclusive: boolean): Promise<string> {
const directory = this.revisionLeaseDirectory();
await mkdir(directory, { recursive: true, mode: 0o700 });
const path = join(directory, `${record.token}.json`);
await writeFile(path, JSON.stringify(record), {
encoding: "utf8",
mode: 0o600,
flush: true,
...(exclusive ? { flag: "wx" } : {}),
});
return path;
}
private async replaceRevisionLease(path: string, record: RevisionLeaseRecord): Promise<void> {
const staging = `${path}.staging-${randomUUID()}`;
try {
await writeFile(staging, JSON.stringify(record), {
encoding: "utf8", mode: 0o600, flag: "wx", flush: true,
});
await rename(staging, path);
} catch (error) {
await rm(staging, { force: true });
throw error;
}
}
private async revisionLeases(): Promise<Array<{ path: string; record: RevisionLeaseRecord }>> {
const directory = this.revisionLeaseDirectory();
await mkdir(directory, { recursive: true, mode: 0o700 });
const entries = await readdir(directory, { withFileTypes: true });
const leases: Array<{ path: string; record: RevisionLeaseRecord }> = [];
for (const entry of entries) {
if (!entry.isFile() || entry.isSymbolicLink() || !/^[0-9a-f-]{36}\.json$/.test(entry.name)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision lease is invalid");
}
const path = join(directory, entry.name);
let record: RevisionLeaseRecord;
try {
record = JSON.parse(await readFile(path, "utf8")) as RevisionLeaseRecord;
} catch {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision lease is invalid");
}
if (
record.version !== 1
|| `${record.token}.json` !== entry.name
|| !/^[0-9a-f-]{36}$/.test(record.token)
|| !/^[a-z][a-z0-9-]{2,62}$/.test(record.workspaceId)
|| !/^[0-9a-f]{40}$/.test(record.commit)
|| (record.state !== "creating" && record.state !== "persisted")
) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision lease is invalid");
}
leases.push({ path, record });
}
return leases;
}
/**
* Publish canonical YAML and derived public documentation as one optimistic Git revision.
* The browser never provides paths or generated artifacts; those are derived server-side.