feat: pin sessions to workspace revisions

This commit is contained in:
2026-08-04 05:13:39 +02:00
parent 301db4bd85
commit bc39730b58
3 changed files with 234 additions and 83 deletions
+111 -79
View File
@@ -73,22 +73,64 @@ export function sessionRoutes(
const isNotFound = (error: unknown) =>
/not found|non trovata|inesistente|404/i.test(error instanceof Error ? error.message : String(error));
/** RLS makes a foreign session indistinguishable from a missing one. */
const authorize = async (principal: PrincipalContext, id: string, workspace?: string): Promise<any | undefined> => {
type LocatedSession = { manifest: any; workspaceConfigPath: string };
const workspaceRevisionUnavailable = () => Object.assign(
new Error("workspace revision unavailable"), { code: "workspace_revision_unavailable" },
);
const unavailableWorkspaceReply = (reply: any) => reply.code(409).send({
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
code: "workspace_revision_unavailable",
});
/**
* Find a session by asking every active registry snapshot, never by using the installation
* default. `tht` applies RLS for the supplied principal, so a foreign ID remains a 404.
*/
const locateSession = async (principal: PrincipalContext, id: string): Promise<LocatedSession | undefined> => {
const runner = runnerFor(principal);
// Dependency-injected runners in legacy route tests may model only the mutation under test.
if (typeof runner.sessionShow !== "function") return {
manifest: {}, workspaceConfigPath: (await d.workspaceRegistry.list())[0]?.snapshotPath ?? "",
};
const revisions = await d.workspaceRegistry.list();
for (const revision of revisions) {
if (revision.state !== "operational") continue;
try {
const manifest = await runner.sessionShow(id, revision.snapshotPath);
if (manifest) return { manifest, workspaceConfigPath: revision.snapshotPath };
} catch (error) {
if (isNotFound(error)) continue;
throw error;
}
}
return undefined;
};
/** Read the durable pinned descriptor only after the owner-visible manifest is located. */
const resolveSessionWorkspace = async (located: LocatedSession): Promise<LocatedSession> => {
const saved = located.manifest as { workspace_id?: string; workspace_revision?: string };
if (!saved.workspace_id || !saved.workspace_revision) return located;
try {
const runner = runnerFor(principal);
// Dependency-injected runners in legacy route tests may model only the mutation under
// test. Production ThtRunner always exposes sessionShow; keep that test seam harmless.
if (typeof runner.sessionShow !== "function") return {};
const manifest = await runner.sessionShow(id, workspace);
return manifest ?? undefined;
} catch (error) {
if (isNotFound(error)) return undefined;
throw error;
const pinned = await d.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision);
return { ...located, workspaceConfigPath: pinned.workspaceConfigPath ?? (pinned as any).revision?.snapshotPath };
} catch {
throw workspaceRevisionUnavailable();
}
};
/** RLS makes a foreign session indistinguishable from a missing one. */
const authorize = async (principal: PrincipalContext, id: string): Promise<LocatedSession | undefined> => {
const located = await locateSession(principal, id);
return located && await resolveSessionWorkspace(located);
};
const storageFailure = (reply: any) => reply.code(503).send({ error: "session storage is unavailable" });
const lifecycleFailure = (reply: any, error: unknown) =>
(error as { code?: string } | undefined)?.code === "workspace_revision_unavailable"
? unavailableWorkspaceReply(reply)
: storageFailure(reply);
const releaseIfFinalized = async (
id: string, rt: ReturnType<PiProcessManager["createFor"]>,
@@ -323,10 +365,14 @@ export function sessionRoutes(
if (scope !== "mine" && scope !== "all") return reply.code(400).send({ error: "scope must be mine or all" });
if (scope === "all" && !principal.isAdmin) return reply.code(403).send({ error: "admin scope required" });
try {
const settings = await d.getSettings(principal);
// Admin RLS is deliberately disabled for a normal 'mine' listing.
const scopedPrincipal = scope === "mine" ? { ...principal, isAdmin: false } : principal;
const list: SessionRow[] = await runnerFor(scopedPrincipal).sessionList(settings.workspace);
const runner = runnerFor(scopedPrincipal);
const revisions = await d.workspaceRegistry.list();
const lists = await Promise.all(revisions
.filter((revision) => revision.state === "operational")
.map((revision) => runner.sessionList(revision.snapshotPath) as Promise<SessionRow[]>));
const list = lists.flat();
// Annotate each row with whether a live Pi runtime is currently bound. The client
// opens an `active` session straight into its live view (reconnecting to its pending
// gate), while a cold session keeps its explicit Resume affordance — so a mere click
@@ -337,21 +383,20 @@ export function sessionRoutes(
app.get("/sessions/:id", async (req, reply) => {
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
const manifest = await authorize(principal, (req.params as any).id, settings.workspace);
if (!manifest) return reply.code(404).send({ error: "session not found" });
const session = await authorize(principal, (req.params as any).id);
if (!session) return reply.code(404).send({ error: "session not found" });
const manifest = session.manifest;
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); }
} catch (error) { return lifecycleFailure(reply, error); }
});
app.post("/sessions/:id/response", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
if (!await authorize(principal, id)) return reply.code(404).send({ error: "session not found" });
} catch (error) { return lifecycleFailure(reply, error); }
const rt = d.mgr.get(id);
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
if (!rt.bridge.respond((req.body as any).ui_response)) {
@@ -363,9 +408,8 @@ export function sessionRoutes(
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
if (!await authorize(principal, id)) return reply.code(404).send({ error: "session not found" });
} catch (error) { return lifecycleFailure(reply, error); }
const rt = d.mgr.get(id);
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
rt.bridge.steer((req.body as any).text);
@@ -376,12 +420,12 @@ export function sessionRoutes(
const principal = getPrincipal(req);
return withSessionLifecycle(id, async () => {
let settings: Settings;
let manifest: any;
let located: LocatedSession | undefined;
try {
settings = await d.getSettings(principal);
manifest = await authorize(principal, id, settings.workspace);
located = await locateSession(principal, id);
} catch { return storageFailure(reply); }
if (!manifest) return reply.code(404).send({ error: "session not found" });
if (!located) return reply.code(404).send({ error: "session not found" });
const manifest = located.manifest;
const runner = runnerFor(principal);
// Read-only contract FIRST: finalized or archived sessions never attempt compatibility
// resolution, even when their historical snapshot was subsequently pruned.
@@ -392,21 +436,10 @@ export function sessionRoutes(
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",
});
}
}
let workspaceConfigPath: string;
try { workspaceConfigPath = (await resolveSessionWorkspace(located)).workspaceConfigPath; }
catch { return unavailableWorkspaceReply(reply); }
try { settings = await d.getSettings(principal); } catch { return storageFailure(reply); }
// This check belongs inside the per-session lock: a preceding cold Resume may have
// installed a running runtime while this request was waiting.
const existing = d.mgr.get(id);
@@ -479,20 +512,20 @@ export function sessionRoutes(
const id = (req.params as { id: string }).id;
const principal = getPrincipal(req);
return withSessionLifecycle(id, async () => {
let settings: Settings;
let session: LocatedSession | undefined;
try {
settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
} catch (error) { return lifecycleFailure(reply, error); }
// Invalidate the live generation before persistence can yield. Otherwise its deferred
// bootstrap may start Pi while Close is already in progress.
const current = d.mgr.get(id);
boundRuntimes.delete(id);
if (current) d.mgr.teardownIfCurrent(id, current);
try {
await runnerFor(principal).closeSession(id, settings.workspace);
} catch {
return storageFailure(reply);
await runnerFor(principal).closeSession(id, session.workspaceConfigPath);
} catch (error) {
return lifecycleFailure(reply, error);
} finally {
// clear, NOT forget: a closed session can be reopened, and the per-session seq
// monotonicity is what keeps a browser's old cursor detectable. The buffer is
@@ -506,9 +539,8 @@ export function sessionRoutes(
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
if (!await authorize(principal, id)) return reply.code(404).send({ error: "session not found" });
} catch (error) { return lifecycleFailure(reply, error); }
const rt = d.mgr.get(id);
// Add CORS headers manually: reply.raw.writeHead bypasses Fastify's onSend hooks
// (where @fastify/cors injects headers), so we must set them explicitly here.
@@ -543,58 +575,58 @@ export function sessionRoutes(
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).setName(id, (req.body as any).name, settings.workspace);
} catch { return storageFailure(reply); }
const session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).setName(id, (req.body as any).name, session.workspaceConfigPath);
} catch (error) { return lifecycleFailure(reply, error); }
return reply.code(204).send();
});
app.post("/sessions/:id/group", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).setGroup(id, (req.body as any).group, settings.workspace);
} catch { return storageFailure(reply); }
const session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).setGroup(id, (req.body as any).group, session.workspaceConfigPath);
} catch (error) { return lifecycleFailure(reply, error); }
return reply.code(204).send();
});
app.post("/sessions/:id/archive", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).archive(id, settings.workspace);
} catch { return storageFailure(reply); }
const session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).archive(id, session.workspaceConfigPath);
} catch (error) { return lifecycleFailure(reply, error); }
return reply.code(204).send();
});
app.post("/sessions/:id/unarchive", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).unarchive(id, settings.workspace);
} catch { return storageFailure(reply); }
const session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).unarchive(id, session.workspaceConfigPath);
} catch (error) { return lifecycleFailure(reply, error); }
return reply.code(204).send();
});
app.delete("/sessions/:id", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
return withSessionLifecycle(id, async () => {
let settings: Settings;
let session: LocatedSession | undefined;
try {
settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
} catch (error) { return lifecycleFailure(reply, error); }
const current = d.mgr.get(id);
boundRuntimes.delete(id);
if (current) d.mgr.teardownIfCurrent(id, current);
try {
await runnerFor(principal).deleteSession(id, settings.workspace);
} catch {
return storageFailure(reply);
await runnerFor(principal).deleteSession(id, session.workspaceConfigPath);
} catch (error) {
return lifecycleFailure(reply, error);
}
d.hub.forget(id);
return reply.code(204).send();
@@ -604,9 +636,9 @@ export function sessionRoutes(
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
return await runnerFor(principal).documents(id, settings.workspace);
} catch { return storageFailure(reply); }
const session = await authorize(principal, id);
if (!session) return reply.code(404).send({ error: "session not found" });
return await runnerFor(principal).documents(id, session.workspaceConfigPath);
} catch (error) { return lifecycleFailure(reply, error); }
});
}