814 lines
37 KiB
TypeScript
814 lines
37 KiB
TypeScript
import type { FastifyInstance } from "fastify";
|
|
import type { PiProcessManager } from "../pi/pi-process-manager.js";
|
|
import type { SessionRow, ThtRunner } from "../tht/tht-runner.js";
|
|
import type { SseHub } from "../sse/sse-hub.js";
|
|
import type { Settings } from "../settings/settings-store.js";
|
|
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";
|
|
import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js";
|
|
import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js";
|
|
import { hasPermission, isPrincipalContext, requirePermission } from "../auth/authorization.js";
|
|
|
|
const BOOTSTRAP_FAILURE_MESSAGE =
|
|
"Session startup failed. Check configuration and connectivity, then Resume the session.";
|
|
const READINESS_FAILURE_MESSAGE =
|
|
"Session services are not ready. Check configuration and connectivity, then try again.";
|
|
const RESUME_FAILURE_MESSAGE =
|
|
"Session could not be resumed. Check configuration and connectivity, then try again.";
|
|
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,
|
|
d: {
|
|
mgr: PiProcessManager; tht: ThtRunner; hub: SseHub;
|
|
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;
|
|
/** Explicit loopback-only compatibility path for old clients that send `workspace`. */
|
|
legacyWorkspaceMode?: boolean;
|
|
/** Fail-closed installation/runtime transport capability check. */
|
|
workspaceRuntimeSupport: (workspace: WorkspaceDescriptor) => boolean;
|
|
maintenanceBarrier: MaintenanceBarrier;
|
|
},
|
|
) {
|
|
const lifecycleTails = new Map<string, Promise<void>>();
|
|
const boundRuntimes = new Map<
|
|
string,
|
|
ReturnType<PiProcessManager["createFor"]>
|
|
>();
|
|
const runtimeManifestReaders = new WeakMap<
|
|
ReturnType<PiProcessManager["createFor"]>,
|
|
() => Promise<any>
|
|
>();
|
|
const failurePersistenceClaimed = new WeakSet<
|
|
ReturnType<PiProcessManager["createFor"]>
|
|
>();
|
|
|
|
const withSessionLifecycle = async <T>(id: string, work: () => Promise<T>): Promise<T> => {
|
|
const previous = lifecycleTails.get(id) ?? Promise.resolve();
|
|
let release!: () => void;
|
|
const gate = new Promise<void>((resolve) => { release = resolve; });
|
|
const tail = previous.then(() => gate);
|
|
lifecycleTails.set(id, tail);
|
|
await previous;
|
|
try {
|
|
return await work();
|
|
} finally {
|
|
release();
|
|
if (lifecycleTails.get(id) === tail) lifecycleTails.delete(id);
|
|
}
|
|
};
|
|
|
|
const info = (id: string, text: string, level = "info") =>
|
|
d.hub.publish(id, "info", { type: "info", level, text });
|
|
|
|
const runnerFor = (principal: PrincipalContext): any => {
|
|
const runner = d.tht as any;
|
|
return typeof runner.withPrincipal === "function" ? runner.withPrincipal(principal) : runner;
|
|
};
|
|
const ownershipPrincipal = (principal: PrincipalContext, permission: "session.read_all" | "session.manage_all") => ({
|
|
...principal,
|
|
// The harness transition remains isAdmin-based, but only the relevant all-session
|
|
// permission can enable its RLS bypass.
|
|
isAdmin: hasPermission(principal, permission),
|
|
});
|
|
|
|
const optionsWithRuntimeConfig = (runner: any, workspaceConfigPath: string | undefined, options: any) => (
|
|
workspaceConfigPath && typeof runner.acquireWorkspaceRuntime === "function"
|
|
? { ...options, runtimeConfig: runner.acquireWorkspaceRuntime(workspaceConfigPath) }
|
|
: options
|
|
);
|
|
|
|
const maintenanceReply = (reply: any) => reply.code(503).send({
|
|
code: "maintenance",
|
|
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("preHandler", async (req, reply) => {
|
|
const pathname = req.url.split("?", 1)[0];
|
|
if (pathname === "/runtime/prewarm" || pathname === "/sessions" || pathname.startsWith("/sessions/")) {
|
|
const principal = requirePermission(req, reply, "session.use");
|
|
if (!isPrincipalContext(principal)) return principal;
|
|
}
|
|
});
|
|
app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); });
|
|
|
|
/** Include retained historical descriptors so removed workspaces remain resumable. */
|
|
const sessionRevisions = async () => {
|
|
const registry = d.workspaceRegistry as Partial<WorkspaceRegistry>;
|
|
if (typeof registry.listRetainedSnapshots === "function") {
|
|
return await registry.listRetainedSnapshots();
|
|
}
|
|
return await d.workspaceRegistry.list();
|
|
};
|
|
|
|
const isNotFound = (error: unknown) =>
|
|
/not found|non trovata|inesistente|404/i.test(error instanceof Error ? error.message : String(error));
|
|
|
|
type LocatedSession = {
|
|
manifest: any;
|
|
workspaceConfigPath: string;
|
|
workspace?: WorkspaceDescriptor;
|
|
};
|
|
|
|
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, permission: "session.read_all" | "session.manage_all" = "session.read_all",
|
|
): Promise<LocatedSession | undefined> => {
|
|
const runner = runnerFor(ownershipPrincipal(principal, permission));
|
|
// Dependency-injected runners in legacy route tests may model only the mutation under test.
|
|
if (typeof runner.sessionShow !== "function") return { manifest: {}, workspaceConfigPath: "" };
|
|
const legacySession = async (): Promise<LocatedSession | undefined> => {
|
|
try {
|
|
const manifest = await runner.sessionShow(id);
|
|
return manifest && !manifest.workspace_id && !manifest.workspace_revision
|
|
? { manifest, workspaceConfigPath: "" }
|
|
: undefined;
|
|
} catch (error) {
|
|
if (isNotFound(error)) return undefined;
|
|
throw error;
|
|
}
|
|
};
|
|
let revisions: Awaited<ReturnType<typeof d.workspaceRegistry.list>>;
|
|
try {
|
|
revisions = await sessionRevisions();
|
|
} catch (registryError) {
|
|
// Sessions created before revision pinning still live under the installation's legacy
|
|
// default config. Keep that compatibility path available when a fresh installation has
|
|
// no registry snapshot yet; a pinned session remains fail-closed below.
|
|
const legacy = await legacySession();
|
|
if (legacy) return legacy;
|
|
throw registryError;
|
|
}
|
|
for (const revision of revisions) {
|
|
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 await legacySession();
|
|
};
|
|
|
|
/** 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 pinned = await d.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision);
|
|
const workspace = validateOperationalWorkspace(pinned.workspace);
|
|
return {
|
|
...located,
|
|
workspace,
|
|
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, permission: "session.read_all" | "session.manage_all" = "session.read_all",
|
|
): Promise<LocatedSession | undefined> => await locateSession(principal, id, permission);
|
|
|
|
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"]>,
|
|
): Promise<boolean> => {
|
|
if (d.mgr.get(id) !== rt || boundRuntimes.get(id) !== rt) return false;
|
|
const readManifest = runtimeManifestReaders.get(rt);
|
|
if (!readManifest || (await readManifest())?.status !== "finalized") return false;
|
|
if (!d.mgr.teardownIfCurrent(id, rt)) return false;
|
|
if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
|
|
return true;
|
|
};
|
|
|
|
const bindRuntime = (
|
|
id: string, rt: ReturnType<PiProcessManager["createFor"]>, runner: any, workspace?: string,
|
|
) => {
|
|
const previous = boundRuntimes.get(id);
|
|
boundRuntimes.set(id, rt);
|
|
if (typeof runner.sessionShow === "function") {
|
|
runtimeManifestReaders.set(rt, () => runner.sessionShow(id, workspace));
|
|
}
|
|
try {
|
|
rt.bridge.onClientEvent((e) => {
|
|
// Child termination is asynchronous. Ignore queued events from a runtime once a newer
|
|
// identity is bound or the session is explicitly closed/deleted. The active identity
|
|
// remains bound after an unexpected exit so its public failure events still reach SSE.
|
|
if (boundRuntimes.get(id) !== rt) return;
|
|
if (e.type === "system_event" && e.event === "session_failed") {
|
|
if (!failurePersistenceClaimed.has(rt)) {
|
|
failurePersistenceClaimed.add(rt);
|
|
void withSessionLifecycle(id, async () => {
|
|
// agent_end can release the old binding before this queued work acquires the lock.
|
|
// Undefined means no replacement; a different identity means Resume won and the
|
|
// old failure must not touch its manifest.
|
|
const bound = boundRuntimes.get(id);
|
|
if (bound !== undefined && bound !== rt) return;
|
|
// Best-effort by design (crash path), but a storage outage must be
|
|
// visible server-side: the manifest stays open and resume remains legal.
|
|
await runner.failSession(id, workspace).catch((error: unknown) => {
|
|
console.error(`[session:${id}] failSession persistence failed:`, error);
|
|
});
|
|
}).catch(() => undefined);
|
|
}
|
|
}
|
|
d.hub.publish(id, e.type, e);
|
|
if (e.type === "system_event" && e.event === "agent_end") {
|
|
if (d.mgr.get(id) !== rt && boundRuntimes.get(id) === rt) {
|
|
boundRuntimes.delete(id);
|
|
} else if (typeof runner.sessionShow === "function") {
|
|
void withSessionLifecycle(id, () => releaseIfFinalized(id, rt)).catch((error: unknown) => {
|
|
console.error(`[session:${id}] terminal runtime cleanup failed:`, error);
|
|
});
|
|
}
|
|
}
|
|
});
|
|
} catch (error) {
|
|
if (previous) boundRuntimes.set(id, previous);
|
|
else if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
|
|
throw error;
|
|
}
|
|
};
|
|
|
|
const bootstrap = (
|
|
id: string,
|
|
rt: ReturnType<PiProcessManager["createFor"]>,
|
|
runner: any,
|
|
workspace: string | undefined,
|
|
configure: Promise<void>,
|
|
retrieval: Promise<void> | null,
|
|
start: () => void,
|
|
) => {
|
|
void (async () => {
|
|
try {
|
|
if (retrieval) info(id, "Preparing retrieval context");
|
|
await Promise.all([configure, retrieval]);
|
|
if (d.mgr.get(id) !== rt) return;
|
|
info(id, "Starting model");
|
|
if (d.mgr.get(id) !== rt) return;
|
|
start();
|
|
} catch (error) {
|
|
// The client only ever sees the generic BOOTSTRAP_FAILURE_MESSAGE; log the real
|
|
// cause server-side so failures (e.g. an unreachable DWH/vector host behind a
|
|
// dropped VPN) are diagnosable from the backend console instead of silent.
|
|
console.error(`[pi:${id}] bootstrap failed:`, error instanceof Error ? error.message : error);
|
|
void withSessionLifecycle(id, async () => {
|
|
// A bootstrap continuation can settle after Close/Delete or after a replacement was
|
|
// installed. Claim only the runtime identity that actually failed; holding the same
|
|
// lifecycle lock through persistence prevents a Resume from becoming that failure's
|
|
// accidental target.
|
|
if (d.mgr.get(id) !== rt || !d.mgr.teardownIfCurrent(id, rt)) return;
|
|
if (!failurePersistenceClaimed.has(rt)) {
|
|
failurePersistenceClaimed.add(rt);
|
|
await runner.failSession(id, workspace).catch(() => undefined);
|
|
}
|
|
rt.bridge.emitClientEvent({ type: "info", level: "error", text: BOOTSTRAP_FAILURE_MESSAGE });
|
|
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
|
|
rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
|
|
});
|
|
}
|
|
})();
|
|
};
|
|
|
|
app.post("/runtime/prewarm", async (req, reply) => {
|
|
let settings: Settings;
|
|
try { settings = await d.getSettings(getPrincipal(req)); } catch { return storageFailure(reply); }
|
|
const workspace = settings.workspace ?? "";
|
|
void d.readiness.ensure(workspace, getPrincipal(req)).catch(() => undefined);
|
|
return reply.code(202).send({ status: "warming" });
|
|
});
|
|
|
|
app.post("/sessions", async (req, reply) => {
|
|
const b = req.body as {
|
|
question: string; name?: string; workspace?: 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);
|
|
// `workspace` was the legacy request field before browser-local registry preferences.
|
|
// It is available only through the explicit loopback-only compatibility mode; every normal
|
|
// new session must resolve and pin an immutable registry revision.
|
|
const legacyWorkspaceRequest = d.legacyWorkspaceMode
|
|
&& typeof b.workspace === "string" && b.workspace.length > 0;
|
|
const requestedWorkspaceId = b.workspaceId ?? (legacyWorkspaceRequest ? undefined : s.workspace);
|
|
if (!requestedWorkspaceId && !legacyWorkspaceRequest) {
|
|
return reply.code(409).send({
|
|
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
|
|
code: "workspace_revision_unavailable",
|
|
});
|
|
}
|
|
let revisionLease: Awaited<ReturnType<WorkspaceRegistry["acquireSessionRevision"]>> | undefined;
|
|
let manifestPersisted = false;
|
|
try {
|
|
let workspaceConfigPath: string | undefined;
|
|
let workspaceId: string | undefined;
|
|
let workspaceRevision: string | undefined;
|
|
let workspaceDescriptor: WorkspaceDescriptor | 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 (!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;
|
|
workspaceDescriptor = resolved.workspace;
|
|
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, workspaceDescriptor,
|
|
);
|
|
if (!ensure.ok) return reply.code(503).send({
|
|
error: READINESS_FAILURE_MESSAGE,
|
|
...(ensure.code ? { code: ensure.code } : {}),
|
|
});
|
|
// 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 {
|
|
({ 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 runtimeOptions = options;
|
|
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
|
try {
|
|
runtimeOptions = optionsWithRuntimeConfig(runner, workspaceConfigPath, options);
|
|
rt = d.mgr.createFor(id, runtimeOptions);
|
|
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, runtimeOptions),
|
|
runner.searchPack(b.question, id, workspaceConfigPath),
|
|
() => d.mgr.start(id, rt, runtimeOptions),
|
|
);
|
|
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",
|
|
);
|
|
});
|
|
}
|
|
}
|
|
});
|
|
app.get("/sessions", async (req, reply) => {
|
|
const principal = getPrincipal(req);
|
|
const scope = (req.query as { scope?: string }).scope ?? "mine";
|
|
if (scope !== "mine" && scope !== "all") return reply.code(400).send({ error: "scope must be mine or all" });
|
|
if (scope === "all") {
|
|
const allPrincipal = requirePermission(req, reply, "session.read_all");
|
|
if (!isPrincipalContext(allPrincipal)) return allPrincipal;
|
|
}
|
|
try {
|
|
// Admin RLS is deliberately disabled for a normal 'mine' listing.
|
|
const scopedPrincipal = scope === "mine"
|
|
? { ...principal, isAdmin: false }
|
|
: ownershipPrincipal(principal, "session.read_all");
|
|
const runner = runnerFor(scopedPrincipal);
|
|
const revisions = await sessionRevisions();
|
|
const lists = await Promise.all(revisions
|
|
.map((revision) => runner.sessionList(revision.snapshotPath) as Promise<SessionRow[]>));
|
|
const sessions = new Map<string, SessionRow>();
|
|
for (const row of lists.flat()) {
|
|
if (!sessions.has(row.id)) sessions.set(row.id, row);
|
|
}
|
|
const list = [...sessions.values()];
|
|
// Only an administrator-visible complete list (or the single local principal) is safe
|
|
// input for retention. A remote per-user view can never discard another principal's pin.
|
|
const reconcileSnapshotRetention = (d.workspaceRegistry as Partial<WorkspaceRegistry>).reconcileSnapshotRetention;
|
|
const hasCompleteRetentionView = (scope === "all" || principal.issuer === "local")
|
|
&& hasPermission(principal, "session.read_all");
|
|
if (hasCompleteRetentionView && typeof reconcileSnapshotRetention === "function") {
|
|
const retained = [...new Set(list
|
|
.filter((row) => row.status !== "finalized" && !row.archived && typeof row.workspace_revision === "string")
|
|
.map((row) => row.workspace_revision!))];
|
|
await reconcileSnapshotRetention.call(d.workspaceRegistry, retained);
|
|
}
|
|
// 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
|
|
// never spawns a runtime.
|
|
return list.map((row) => ({ ...row, active: d.mgr.get(row.id) !== undefined }));
|
|
} catch { return storageFailure(reply); }
|
|
});
|
|
app.get("/sessions/:id", async (req, reply) => {
|
|
const principal = getPrincipal(req);
|
|
try {
|
|
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 (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 {
|
|
if (!await authorize(principal, id, "session.manage_all")) 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)) {
|
|
return reply.code(409).send({ error: "risposta non corrispondente al gate in attesa" });
|
|
}
|
|
return reply.code(204).send();
|
|
});
|
|
app.post("/sessions/:id/steer", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
const principal = getPrincipal(req);
|
|
try {
|
|
if (!await authorize(principal, id, "session.manage_all")) 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);
|
|
return reply.code(204).send();
|
|
});
|
|
app.post("/sessions/:id/resume", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
const principal = getPrincipal(req);
|
|
return withSessionLifecycle(id, async () => {
|
|
let settings: Settings;
|
|
let located: LocatedSession | undefined;
|
|
try {
|
|
located = await locateSession(principal, id, "session.manage_all");
|
|
} catch { return storageFailure(reply); }
|
|
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.
|
|
if (manifest?.status === "finalized" || manifest?.archived) {
|
|
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
|
}
|
|
const saved = manifest as {
|
|
provider?: string; model?: string; thinking?: string;
|
|
workspace_id?: string; workspace_revision?: string;
|
|
};
|
|
let workspaceConfigPath: string;
|
|
let workspaceDescriptor: WorkspaceDescriptor | undefined;
|
|
try {
|
|
const resolved = await resolveSessionWorkspace(located);
|
|
workspaceConfigPath = resolved.workspaceConfigPath;
|
|
workspaceDescriptor = resolved.workspace;
|
|
}
|
|
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);
|
|
if (existing) {
|
|
const state = existing.bridge.turnState();
|
|
if (state === "running" || state === "waiting") {
|
|
return reply.code(200).send({ id, alreadyActive: true });
|
|
}
|
|
}
|
|
const ensure = await d.readiness.ensure(
|
|
workspaceConfigPath ?? "", principal, workspaceDescriptor,
|
|
);
|
|
if (!ensure.ok) return reply.code(503).send({
|
|
error: READINESS_FAILURE_MESSAGE,
|
|
...(ensure.code ? { code: ensure.code } : {}),
|
|
});
|
|
const options = {
|
|
provider: saved?.provider,
|
|
model: saved?.model,
|
|
thinking: saved?.thinking ?? settings.thinking,
|
|
author: principal.displayName ?? principal.subject,
|
|
principal,
|
|
mode: "resume" as const,
|
|
};
|
|
let runtimeOptions = options;
|
|
|
|
// Reopening is validation, not the transport commit point. Keep the old hub intact if
|
|
// persistence cannot be reopened.
|
|
try {
|
|
await runner.reopenSession(id, workspaceConfigPath);
|
|
} catch {
|
|
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
|
}
|
|
|
|
// Re-check immediately before the replacement commit. A lifecycle operation that ran
|
|
// before this request acquired the lock may have changed or removed the runtime.
|
|
const current = d.mgr.get(id);
|
|
if (current) {
|
|
const state = current.bridge.turnState();
|
|
if (state === "running" || state === "waiting") {
|
|
return reply.code(200).send({ id, alreadyActive: true });
|
|
}
|
|
}
|
|
|
|
for (const stoppedId of d.mgr.teardownForPrincipal?.(principal) ?? []) {
|
|
boundRuntimes.delete(stoppedId);
|
|
}
|
|
|
|
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
|
try {
|
|
if (current) {
|
|
if (boundRuntimes.get(id) === current) boundRuntimes.delete(id);
|
|
d.mgr.teardownIfCurrent(id, current);
|
|
}
|
|
runtimeOptions = optionsWithRuntimeConfig(runner, workspaceConfigPath, options);
|
|
rt = d.mgr.createFor(id, runtimeOptions);
|
|
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.
|
|
if (rt) {
|
|
if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
|
|
d.mgr.teardownIfCurrent(id, rt);
|
|
}
|
|
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
|
}
|
|
|
|
// Commit the replacement only after reopen + runtime creation/binding succeeded, and
|
|
// immediately before the first event produced by the new Resume.
|
|
d.hub.clear(id);
|
|
info(id, "Resuming session");
|
|
bootstrap(
|
|
id, rt, runner, workspaceConfigPath,
|
|
d.mgr.configure(rt, runtimeOptions), null,
|
|
() => d.mgr.start(id, rt, runtimeOptions),
|
|
);
|
|
return reply.code(200).send({ id, alreadyActive: false });
|
|
});
|
|
});
|
|
app.post("/sessions/:id/close", async (req, reply) => {
|
|
const id = (req.params as { id: string }).id;
|
|
const principal = getPrincipal(req);
|
|
return withSessionLifecycle(id, async () => {
|
|
let session: LocatedSession | undefined;
|
|
try {
|
|
session = await authorize(principal, id, "session.manage_all");
|
|
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, 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
|
|
// emptied here; only delete discards the id counter.
|
|
d.hub.clear(id);
|
|
}
|
|
return { closed: true };
|
|
});
|
|
});
|
|
app.get("/sessions/:id/events", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
const principal = getPrincipal(req);
|
|
try {
|
|
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);
|
|
// reply.raw.writeHead bypasses Fastify's CORS hook. Only a cookie-authenticated request
|
|
// from the exact configured public origin receives credentialed SSE CORS headers.
|
|
const origin = typeof req.headers.origin === "string" ? req.headers.origin : undefined;
|
|
let isConfiguredOrigin = false;
|
|
try {
|
|
isConfiguredOrigin = req.authPublicOrigin !== undefined
|
|
&& origin !== undefined
|
|
&& new URL(origin).origin === req.authPublicOrigin;
|
|
} catch {
|
|
isConfiguredOrigin = false;
|
|
}
|
|
const corsHeaders = isConfiguredOrigin
|
|
? {
|
|
"Access-Control-Allow-Origin": req.authPublicOrigin,
|
|
"Access-Control-Allow-Credentials": "true",
|
|
}
|
|
: {};
|
|
reply.raw.writeHead(200, {
|
|
"Content-Type": "text/event-stream",
|
|
"Cache-Control": "no-cache",
|
|
"X-Accel-Buffering": "no",
|
|
Connection: "keep-alive",
|
|
...corsHeaders,
|
|
});
|
|
// Send the handshake immediately. Without this, Node waits for the first event body and
|
|
// proxies/clients cannot establish an idle SSE subscription or inspect its headers.
|
|
reply.raw.flushHeaders();
|
|
const send = (event: string, data: object, eventId: string) =>
|
|
reply.raw.write(`id: ${eventId}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
|
|
const off = d.hub.subscribe(id, send, {
|
|
// Raw cursor candidates: the hub parses "<generation>:<seq>" and treats any
|
|
// other-generation (or legacy numeric) cursor as stale → replay from the start.
|
|
after: [req.headers["last-event-id"], (req.query as { lastEventId?: unknown }).lastEventId],
|
|
pending: rt?.bridge.pendingWidget() ?? null,
|
|
// clear()/forget() end every old transport so native EventSource reconnects with its
|
|
// Last-Event-ID instead of remaining attached to a subscriber callback that no longer exists.
|
|
close: () => {
|
|
if (!reply.raw.writableEnded) reply.raw.end();
|
|
},
|
|
});
|
|
req.raw.on("close", off);
|
|
});
|
|
app.post("/sessions/:id/rename", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
const principal = getPrincipal(req);
|
|
try {
|
|
const session = await authorize(principal, id, "session.manage_all");
|
|
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 session = await authorize(principal, id, "session.manage_all");
|
|
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 session = await authorize(principal, id, "session.manage_all");
|
|
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 session = await authorize(principal, id, "session.manage_all");
|
|
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 session: LocatedSession | undefined;
|
|
try {
|
|
session = await authorize(principal, id, "session.manage_all");
|
|
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, session.workspaceConfigPath);
|
|
} catch (error) {
|
|
return lifecycleFailure(reply, error);
|
|
}
|
|
d.hub.forget(id);
|
|
return reply.code(204).send();
|
|
});
|
|
});
|
|
app.get("/sessions/:id/documents", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
const principal = getPrincipal(req);
|
|
try {
|
|
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); }
|
|
});
|
|
}
|