fix: make session resume atomic across restarts
This commit is contained in:
@@ -68,9 +68,14 @@ export class PiProcessManager {
|
||||
cwd: this.cfg.harnessDir,
|
||||
env,
|
||||
});
|
||||
// Log stderr for debugging (was silently drained)
|
||||
child.stderr.on("data", (d: Buffer) => console.error(`[pi:${sessionId}] stderr:`, d.toString().trim()));
|
||||
return child;
|
||||
try {
|
||||
// Log stderr for debugging (was silently drained)
|
||||
child.stderr.on("data", (d: Buffer) => console.error(`[pi:${sessionId}] stderr:`, d.toString().trim()));
|
||||
return child;
|
||||
} catch (error) {
|
||||
try { child.kill(); } catch { /* preserve the initialization error */ }
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
count(): number { return this.runtimes.size; }
|
||||
@@ -91,31 +96,39 @@ export class PiProcessManager {
|
||||
const author = o.author ?? "dev@local";
|
||||
const provider = canonicalPiProvider(o.provider ?? this.cfg.defaults.provider);
|
||||
const child = this.spawnFn(sessionId, author, provider);
|
||||
const rpc = new RpcClient(child);
|
||||
const bridge = new SessionBridge(rpc);
|
||||
const rt: SessionRuntime = { rpc, bridge, child };
|
||||
bridge.beginTurn();
|
||||
this.runtimes.set(sessionId, rt);
|
||||
// Identity-checked: a stale child's exit must not evict a newer runtime.
|
||||
// Expected teardowns (teardown()/respawn) delete the runtime from the map BEFORE the
|
||||
// exit event fires, so reaching this branch with `rt` still mapped means the child
|
||||
// died on its own: tell the client, or the UI spins forever waiting for a turn end.
|
||||
child.on("exit", (code) => {
|
||||
console.error(`[pi:${sessionId}] exited code=${code ?? "?"} mapped=${this.runtimes.get(sessionId) === rt}`);
|
||||
if (this.runtimes.get(sessionId) === rt) {
|
||||
rt.bridge.markFailed();
|
||||
this.runtimes.delete(sessionId);
|
||||
rt.bridge.emitClientEvent({
|
||||
type: "info",
|
||||
level: "error",
|
||||
text: `Pi process exited unexpectedly (code ${code ?? "?"})`,
|
||||
});
|
||||
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
|
||||
rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
|
||||
}
|
||||
});
|
||||
let rt: SessionRuntime | undefined;
|
||||
try {
|
||||
const rpc = new RpcClient(child);
|
||||
const bridge = new SessionBridge(rpc);
|
||||
const runtime: SessionRuntime = { rpc, bridge, child };
|
||||
rt = runtime;
|
||||
bridge.beginTurn();
|
||||
this.runtimes.set(sessionId, runtime);
|
||||
// Identity-checked: a stale child's exit must not evict a newer runtime.
|
||||
// Expected teardowns (teardown()/respawn) delete the runtime from the map BEFORE the
|
||||
// exit event fires, so reaching this branch with `rt` still mapped means the child
|
||||
// died on its own: tell the client, or the UI spins forever waiting for a turn end.
|
||||
child.on("exit", (code) => {
|
||||
console.error(`[pi:${sessionId}] exited code=${code ?? "?"} mapped=${this.runtimes.get(sessionId) === runtime}`);
|
||||
if (this.runtimes.get(sessionId) === runtime) {
|
||||
runtime.bridge.markFailed();
|
||||
this.runtimes.delete(sessionId);
|
||||
runtime.bridge.emitClientEvent({
|
||||
type: "info",
|
||||
level: "error",
|
||||
text: `Pi process exited unexpectedly (code ${code ?? "?"})`,
|
||||
});
|
||||
runtime.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
|
||||
runtime.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
|
||||
}
|
||||
});
|
||||
|
||||
return rt;
|
||||
return runtime;
|
||||
} catch (error) {
|
||||
if (rt && this.runtimes.get(sessionId) === rt) this.runtimes.delete(sessionId);
|
||||
try { child.kill(); } catch { /* preserve the initialization error */ }
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
/** Configure model and thinking. Safe to run alongside deterministic retrieval. */
|
||||
|
||||
+109
-36
@@ -10,6 +10,8 @@ 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.";
|
||||
|
||||
function eventCursor(...values: unknown[]): number {
|
||||
let cursor = 0;
|
||||
@@ -25,16 +27,58 @@ export function sessionRoutes(
|
||||
app: FastifyInstance,
|
||||
d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager },
|
||||
) {
|
||||
const resumeTails = new Map<string, Promise<void>>();
|
||||
const boundRuntimes = new Map<
|
||||
string,
|
||||
ReturnType<PiProcessManager["createFor"]>
|
||||
>();
|
||||
|
||||
const withResumeLock = async <T>(id: string, work: () => Promise<T>): Promise<T> => {
|
||||
const previous = resumeTails.get(id) ?? Promise.resolve();
|
||||
let release!: () => void;
|
||||
const gate = new Promise<void>((resolve) => { release = resolve; });
|
||||
const tail = previous.then(() => gate);
|
||||
resumeTails.set(id, tail);
|
||||
await previous;
|
||||
try {
|
||||
return await work();
|
||||
} finally {
|
||||
release();
|
||||
if (resumeTails.get(id) === tail) resumeTails.delete(id);
|
||||
}
|
||||
};
|
||||
|
||||
const info = (id: string, text: string, level = "info") =>
|
||||
d.hub.publish(id, "info", { type: "info", level, text });
|
||||
|
||||
const bindRuntime = (id: string, rt: ReturnType<PiProcessManager["createFor"]>) =>
|
||||
rt.bridge.onClientEvent((e) => {
|
||||
if (e.type === "system_event" && e.event === "session_failed") {
|
||||
void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
||||
}
|
||||
d.hub.publish(id, e.type, e);
|
||||
});
|
||||
const bindRuntime = (id: string, rt: ReturnType<PiProcessManager["createFor"]>) => {
|
||||
const previous = boundRuntimes.get(id);
|
||||
boundRuntimes.set(id, rt);
|
||||
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") {
|
||||
void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
||||
}
|
||||
d.hub.publish(id, e.type, e);
|
||||
if (
|
||||
e.type === "system_event"
|
||||
&& e.event === "agent_end"
|
||||
&& d.mgr.get(id) !== rt
|
||||
&& boundRuntimes.get(id) === rt
|
||||
) {
|
||||
boundRuntimes.delete(id);
|
||||
}
|
||||
});
|
||||
} catch (error) {
|
||||
if (previous) boundRuntimes.set(id, previous);
|
||||
else if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
|
||||
throw error;
|
||||
}
|
||||
};
|
||||
|
||||
const bootstrap = (
|
||||
id: string,
|
||||
@@ -114,42 +158,69 @@ export function sessionRoutes(
|
||||
});
|
||||
app.post("/sessions/:id/resume", async (req, reply) => {
|
||||
const id = (req.params as any).id;
|
||||
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 });
|
||||
return withResumeLock(id, async () => {
|
||||
// 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 manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null;
|
||||
if (manifest?.status === "finalized" || manifest?.archived) {
|
||||
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
||||
}
|
||||
const settings = d.getSettings();
|
||||
const ensure = await d.readiness.ensure(settings.workspace ?? "");
|
||||
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,
|
||||
thinking: saved?.thinking ?? settings.thinking,
|
||||
author: getUser(req).id,
|
||||
mode: "resume" as const,
|
||||
};
|
||||
if (existing) d.mgr.teardown(id);
|
||||
d.hub.clear(id);
|
||||
await d.tht.reopenSession(id, settings.workspace);
|
||||
const rt = d.mgr.createFor(id, options);
|
||||
bindRuntime(id, rt);
|
||||
info(id, "Resuming session");
|
||||
bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
||||
return reply.code(200).send({ id, alreadyActive: false });
|
||||
const manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null;
|
||||
if (manifest?.status === "finalized" || manifest?.archived) {
|
||||
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
||||
}
|
||||
const settings = d.getSettings();
|
||||
const ensure = await d.readiness.ensure(settings.workspace ?? "");
|
||||
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,
|
||||
thinking: saved?.thinking ?? settings.thinking,
|
||||
author: getUser(req).id,
|
||||
mode: "resume" as const,
|
||||
};
|
||||
|
||||
// Reopening is validation, not the transport commit point. Keep the old hub intact if
|
||||
// persistence cannot be reopened.
|
||||
try {
|
||||
await d.tht.reopenSession(id, settings.workspace);
|
||||
} catch {
|
||||
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
||||
}
|
||||
|
||||
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
||||
try {
|
||||
if (existing) {
|
||||
boundRuntimes.delete(id);
|
||||
d.mgr.teardown(id);
|
||||
}
|
||||
rt = d.mgr.createFor(id, options);
|
||||
bindRuntime(id, rt);
|
||||
} catch {
|
||||
// A created-but-unbound runtime is not usable. The old hub remains attached because
|
||||
// clear() has not happened yet.
|
||||
if (rt) d.mgr.teardown(id);
|
||||
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, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
||||
return reply.code(200).send({ id, alreadyActive: false });
|
||||
});
|
||||
});
|
||||
app.post("/sessions/:id/close", async (req) => {
|
||||
const id = (req.params as { id: string }).id;
|
||||
try {
|
||||
await d.tht.closeSession(id, d.getSettings().workspace);
|
||||
} finally {
|
||||
boundRuntimes.delete(id);
|
||||
d.mgr.teardown(id);
|
||||
d.hub.clear(id);
|
||||
}
|
||||
@@ -202,8 +273,10 @@ export function sessionRoutes(
|
||||
});
|
||||
app.delete("/sessions/:id", async (req, reply) => {
|
||||
const id = (req.params as any).id;
|
||||
boundRuntimes.delete(id);
|
||||
d.mgr.teardown(id); // drop any live runtime before deleting on disk
|
||||
await d.tht.deleteSession(id, d.getSettings().workspace);
|
||||
d.hub.forget(id);
|
||||
return reply.code(204).send();
|
||||
});
|
||||
app.get("/sessions/:id/documents", async (req) => d.tht.documents((req.params as any).id));
|
||||
|
||||
@@ -91,4 +91,10 @@ export class SseHub {
|
||||
this.buffers.delete(sessionId);
|
||||
this.subs.delete(sessionId);
|
||||
}
|
||||
|
||||
/** Permanently discard transport state for a deleted session id. */
|
||||
forget(sessionId: string): void {
|
||||
this.clear(sessionId);
|
||||
this.lastIds.delete(sessionId);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user