feat: pin sessions to workspace revisions
This commit is contained in:
@@ -7,6 +7,7 @@ import { getPrincipal } from "../auth/auth.js";
|
||||
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";
|
||||
|
||||
const BOOTSTRAP_FAILURE_MESSAGE =
|
||||
"Session startup failed. Check configuration and connectivity, then Resume the session.";
|
||||
@@ -18,6 +19,8 @@ const DWH_UNREACHABLE_MESSAGE =
|
||||
"Cannot start a session: the database is unreachable. Check the VPN connection and try again.";
|
||||
const MODEL_UNAVAILABLE_MESSAGE =
|
||||
"Selected model is unavailable. Check Pi authentication and model settings, then try again.";
|
||||
const WORKSPACE_REVISION_UNAVAILABLE_MESSAGE =
|
||||
"Session workspace configuration is unavailable. Check configuration and try again.";
|
||||
|
||||
export function sessionRoutes(
|
||||
app: FastifyInstance,
|
||||
@@ -26,6 +29,7 @@ export function sessionRoutes(
|
||||
getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||
readiness: ReadinessManager;
|
||||
listModels: ListModelsFn;
|
||||
workspaceRegistry: WorkspaceRegistry;
|
||||
/** Local-only guard: probe DWH reachability before creating a session (run-stack.sh). */
|
||||
dwhPrecheck?: boolean;
|
||||
},
|
||||
@@ -195,28 +199,61 @@ export function sessionRoutes(
|
||||
});
|
||||
|
||||
app.post("/sessions", async (req, reply) => {
|
||||
const b = req.body as { question: string; name?: string };
|
||||
const b = req.body as {
|
||||
question: string; name?: string; workspaceId?: string;
|
||||
provider?: string; model?: string; thinking?: string;
|
||||
};
|
||||
const principal = getPrincipal(req);
|
||||
let s: Settings;
|
||||
try { s = await d.getSettings(principal); } catch { return storageFailure(reply); }
|
||||
const runner = runnerFor(principal);
|
||||
let workspaceConfigPath = s.workspace;
|
||||
let workspaceId: string | undefined;
|
||||
let workspaceRevision: string | undefined;
|
||||
let allowedModels: readonly string[] | undefined;
|
||||
if (b.workspaceId) {
|
||||
try {
|
||||
const resolved = await d.workspaceRegistry.read(b.workspaceId);
|
||||
if (resolved.revision.state !== "operational") {
|
||||
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(s.workspace ?? "", principal);
|
||||
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(s.workspace);
|
||||
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 (s.provider && s.model) {
|
||||
if (provider && model) {
|
||||
let available: Awaited<ReturnType<ListModelsFn>>;
|
||||
try {
|
||||
available = await d.listModels();
|
||||
@@ -227,7 +264,7 @@ export function sessionRoutes(
|
||||
});
|
||||
}
|
||||
const selectedAvailable = available.some(
|
||||
(candidate) => candidate.provider === s.provider && candidate.id === s.model,
|
||||
(candidate) => candidate.provider === provider && candidate.id === model,
|
||||
);
|
||||
if (!selectedAvailable) {
|
||||
return reply.code(503).send({
|
||||
@@ -236,19 +273,18 @@ export function sessionRoutes(
|
||||
});
|
||||
}
|
||||
}
|
||||
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
||||
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
||||
// 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, workspace: s.workspace,
|
||||
provider: s.provider, model: s.model, thinking: s.thinking,
|
||||
question: b.question, name: b.name, workspaceConfigPath,
|
||||
workspaceId, workspaceRevision, provider, model, thinking,
|
||||
}));
|
||||
} catch { return storageFailure(reply); }
|
||||
const options = {
|
||||
provider: s.provider,
|
||||
model: s.model,
|
||||
thinking: s.thinking,
|
||||
provider, model, thinking,
|
||||
author: principal.displayName ?? principal.subject,
|
||||
principal,
|
||||
question: b.question,
|
||||
@@ -256,22 +292,22 @@ export function sessionRoutes(
|
||||
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
||||
try {
|
||||
rt = d.mgr.createFor(id, options);
|
||||
bindRuntime(id, rt, runner, s.workspace);
|
||||
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, s.workspace).catch((persistenceError: unknown) => {
|
||||
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, s.workspace, d.mgr.configure(rt, options),
|
||||
runner.searchPack(b.question, id, s.workspace),
|
||||
id, rt, runner, workspaceConfigPath, d.mgr.configure(rt, options),
|
||||
runner.searchPack(b.question, id, workspaceConfigPath),
|
||||
() => d.mgr.start(id, rt, options),
|
||||
);
|
||||
return { id };
|
||||
@@ -298,7 +334,10 @@ export function sessionRoutes(
|
||||
try {
|
||||
const settings = await d.getSettings(principal);
|
||||
const manifest = await authorize(principal, (req.params as any).id, settings.workspace);
|
||||
return manifest ?? reply.code(404).send({ error: "session not found" });
|
||||
if (!manifest) return reply.code(404).send({ error: "session not found" });
|
||||
return (!manifest.workspace_id || !manifest.workspace_revision)
|
||||
? { ...manifest, warning: "Legacy session: this session is not pinned to a workspace revision." }
|
||||
: manifest;
|
||||
} catch { return storageFailure(reply); }
|
||||
});
|
||||
app.post("/sessions/:id/response", async (req, reply) => {
|
||||
@@ -339,6 +378,25 @@ export function sessionRoutes(
|
||||
} catch { return storageFailure(reply); }
|
||||
if (!manifest) return reply.code(404).send({ error: "session not found" });
|
||||
const runner = runnerFor(principal);
|
||||
const saved = manifest as {
|
||||
provider?: string; model?: string; thinking?: string;
|
||||
workspace_id?: string; workspace_revision?: string;
|
||||
};
|
||||
let workspaceConfigPath = settings.workspace;
|
||||
const warning = !saved.workspace_id || !saved.workspace_revision
|
||||
? "Legacy session: this session is not pinned to a workspace revision."
|
||||
: undefined;
|
||||
if (!warning) {
|
||||
try {
|
||||
const pinned = await d.workspaceRegistry.readPinned(saved.workspace_id!, saved.workspace_revision!);
|
||||
workspaceConfigPath = pinned.workspaceConfigPath ?? (pinned as any).revision?.snapshotPath;
|
||||
} catch {
|
||||
return reply.code(409).send({
|
||||
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
|
||||
code: "workspace_revision_unavailable",
|
||||
});
|
||||
}
|
||||
}
|
||||
// Read-only contract FIRST: a finalized/archived session must refuse resume even
|
||||
// when a lingering runtime still looks active — the manifest is the truth.
|
||||
if (manifest?.status === "finalized" || manifest?.archived) {
|
||||
@@ -353,9 +411,8 @@ export function sessionRoutes(
|
||||
return reply.code(200).send({ id, alreadyActive: true });
|
||||
}
|
||||
}
|
||||
const ensure = await d.readiness.ensure(settings.workspace ?? "", principal);
|
||||
const ensure = await d.readiness.ensure(workspaceConfigPath ?? "", principal);
|
||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
|
||||
const options = {
|
||||
provider: saved?.provider,
|
||||
model: saved?.model,
|
||||
@@ -368,7 +425,7 @@ export function sessionRoutes(
|
||||
// Reopening is validation, not the transport commit point. Keep the old hub intact if
|
||||
// persistence cannot be reopened.
|
||||
try {
|
||||
await runner.reopenSession(id, settings.workspace);
|
||||
await runner.reopenSession(id, workspaceConfigPath);
|
||||
} catch {
|
||||
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
||||
}
|
||||
@@ -394,7 +451,7 @@ export function sessionRoutes(
|
||||
d.mgr.teardownIfCurrent(id, current);
|
||||
}
|
||||
rt = d.mgr.createFor(id, options);
|
||||
bindRuntime(id, rt, runner, settings.workspace);
|
||||
bindRuntime(id, rt, runner, workspaceConfigPath);
|
||||
} catch {
|
||||
// A created-but-unbound runtime is not usable. The old hub remains attached because
|
||||
// clear() has not happened yet.
|
||||
@@ -409,7 +466,7 @@ export function sessionRoutes(
|
||||
// immediately before the first event produced by the new Resume.
|
||||
d.hub.clear(id);
|
||||
info(id, "Resuming session");
|
||||
bootstrap(id, rt, runner, settings.workspace, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
||||
bootstrap(id, rt, runner, workspaceConfigPath, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
||||
return reply.code(200).send({ id, alreadyActive: false });
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { FastifyInstance } from "fastify";
|
||||
import type { AppConfig } from "../config.js";
|
||||
import { loadSettings, saveSettings, type Settings } from "../settings/settings-store.js";
|
||||
import type { Settings } from "../settings/settings-store.js";
|
||||
import { listWorkspaces, type ListModelsFn } from "./meta.js";
|
||||
import { getPrincipal } from "../auth/auth.js";
|
||||
import type { PrincipalContext } from "../auth/principal.js";
|
||||
@@ -9,10 +9,10 @@ import type { PrincipalContext } from "../auth/principal.js";
|
||||
export function effectiveSettings(cfg: AppConfig, stored: Settings): Settings {
|
||||
const workspaces = listWorkspaces(cfg.harnessDir);
|
||||
return {
|
||||
workspace: stored.workspace ?? (workspaces[0]?.name),
|
||||
provider: stored.provider ?? cfg.defaults.provider,
|
||||
model: stored.model ?? cfg.defaults.model,
|
||||
thinking: stored.thinking ?? cfg.defaults.thinking,
|
||||
workspace: workspaces[0]?.name ?? stored.workspace,
|
||||
provider: cfg.defaults.provider ?? stored.provider,
|
||||
model: cfg.defaults.model ?? stored.model,
|
||||
thinking: cfg.defaults.thinking ?? stored.thinking,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -21,7 +21,6 @@ export function settingsRoutes(
|
||||
deps: {
|
||||
cfg: AppConfig; listModels: ListModelsFn;
|
||||
getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||
saveSettings: (principal: PrincipalContext, settings: Settings) => Promise<void>;
|
||||
},
|
||||
): void {
|
||||
app.get("/settings", async (req, reply) => {
|
||||
@@ -50,15 +49,10 @@ export function settingsRoutes(
|
||||
});
|
||||
}
|
||||
}
|
||||
const next: Settings = {
|
||||
workspace: b.workspace,
|
||||
provider: b.provider,
|
||||
model: b.model,
|
||||
thinking: b.thinking,
|
||||
};
|
||||
try {
|
||||
await deps.saveSettings(getPrincipal(req), next);
|
||||
return effectiveSettings(deps.cfg, next);
|
||||
// Retain this endpoint as a validating compatibility surface for older clients, but do
|
||||
// not write anonymous users' choices to shared server storage.
|
||||
return await deps.getSettings(getPrincipal(req));
|
||||
} catch {
|
||||
return reply.code(503).send({ error: "settings storage is unavailable" });
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user