fix: serialize session lifecycle transitions
This commit is contained in:
@@ -174,9 +174,16 @@ export class PiProcessManager {
|
||||
|
||||
teardown(id: string): void {
|
||||
const rt = this.runtimes.get(id);
|
||||
if (rt) {
|
||||
rt.child.kill();
|
||||
this.runtimes.delete(id);
|
||||
}
|
||||
if (rt) this.teardownIfCurrent(id, rt);
|
||||
}
|
||||
|
||||
/** Remove only the runtime identity the caller observed. */
|
||||
teardownIfCurrent(id: string, expected: SessionRuntime): boolean {
|
||||
if (this.runtimes.get(id) !== expected) return false;
|
||||
// Delete before signalling the child so its asynchronous exit cannot be mistaken for a
|
||||
// crash, and so a replacement installed by a later lifecycle operation is never targeted.
|
||||
this.runtimes.delete(id);
|
||||
expected.child.kill();
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -27,24 +27,27 @@ export function sessionRoutes(
|
||||
app: FastifyInstance,
|
||||
d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager },
|
||||
) {
|
||||
const resumeTails = new Map<string, Promise<void>>();
|
||||
const lifecycleTails = new Map<string, Promise<void>>();
|
||||
const boundRuntimes = new Map<
|
||||
string,
|
||||
ReturnType<PiProcessManager["createFor"]>
|
||||
>();
|
||||
const failurePersistenceClaimed = new WeakSet<
|
||||
ReturnType<PiProcessManager["createFor"]>
|
||||
>();
|
||||
|
||||
const withResumeLock = async <T>(id: string, work: () => Promise<T>): Promise<T> => {
|
||||
const previous = resumeTails.get(id) ?? Promise.resolve();
|
||||
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);
|
||||
resumeTails.set(id, tail);
|
||||
lifecycleTails.set(id, tail);
|
||||
await previous;
|
||||
try {
|
||||
return await work();
|
||||
} finally {
|
||||
release();
|
||||
if (resumeTails.get(id) === tail) resumeTails.delete(id);
|
||||
if (lifecycleTails.get(id) === tail) lifecycleTails.delete(id);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -61,7 +64,17 @@ export function sessionRoutes(
|
||||
// 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);
|
||||
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;
|
||||
await d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
||||
}).catch(() => undefined);
|
||||
}
|
||||
}
|
||||
d.hub.publish(id, e.type, e);
|
||||
if (
|
||||
@@ -91,14 +104,25 @@ export function sessionRoutes(
|
||||
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 {
|
||||
d.mgr.teardown(id);
|
||||
void d.tht.failSession(id, d.getSettings().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" });
|
||||
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 d.tht.failSession(id, d.getSettings().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" });
|
||||
});
|
||||
}
|
||||
})();
|
||||
};
|
||||
@@ -158,7 +182,7 @@ export function sessionRoutes(
|
||||
});
|
||||
app.post("/sessions/:id/resume", async (req, reply) => {
|
||||
const id = (req.params as any).id;
|
||||
return withResumeLock(id, async () => {
|
||||
return withSessionLifecycle(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);
|
||||
@@ -192,18 +216,31 @@ export function sessionRoutes(
|
||||
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 });
|
||||
}
|
||||
}
|
||||
|
||||
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
||||
try {
|
||||
if (existing) {
|
||||
boundRuntimes.delete(id);
|
||||
d.mgr.teardown(id);
|
||||
if (current) {
|
||||
if (boundRuntimes.get(id) === current) boundRuntimes.delete(id);
|
||||
d.mgr.teardownIfCurrent(id, current);
|
||||
}
|
||||
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);
|
||||
if (rt) {
|
||||
if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
|
||||
d.mgr.teardownIfCurrent(id, rt);
|
||||
}
|
||||
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
||||
}
|
||||
|
||||
@@ -217,14 +254,19 @@ export function sessionRoutes(
|
||||
});
|
||||
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 {
|
||||
return withSessionLifecycle(id, async () => {
|
||||
// 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);
|
||||
d.mgr.teardown(id);
|
||||
d.hub.clear(id);
|
||||
}
|
||||
return { closed: true };
|
||||
if (current) d.mgr.teardownIfCurrent(id, current);
|
||||
try {
|
||||
await d.tht.closeSession(id, d.getSettings().workspace);
|
||||
} finally {
|
||||
d.hub.clear(id);
|
||||
}
|
||||
return { closed: true };
|
||||
});
|
||||
});
|
||||
app.get("/sessions/:id/events", (req, reply) => {
|
||||
const id = (req.params as any).id;
|
||||
@@ -252,6 +294,11 @@ export function sessionRoutes(
|
||||
const off = d.hub.subscribe(id, send, {
|
||||
afterId,
|
||||
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);
|
||||
});
|
||||
@@ -273,11 +320,14 @@ 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();
|
||||
return withSessionLifecycle(id, async () => {
|
||||
const current = d.mgr.get(id);
|
||||
boundRuntimes.delete(id);
|
||||
if (current) d.mgr.teardownIfCurrent(id, current);
|
||||
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));
|
||||
}
|
||||
|
||||
@@ -1,5 +1,11 @@
|
||||
type Send = (event: string, data: object, id: number) => void;
|
||||
|
||||
interface Subscriber {
|
||||
send: Send;
|
||||
close: () => void;
|
||||
closed: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-session ring buffer of recent events.
|
||||
*
|
||||
@@ -17,6 +23,7 @@ interface BufferedEvent {
|
||||
interface SubscribeOptions {
|
||||
afterId?: number;
|
||||
pending?: object | null;
|
||||
close?: () => void;
|
||||
}
|
||||
|
||||
function descriptorId(value: unknown): string | null {
|
||||
@@ -32,7 +39,7 @@ function bufferedGateId(item: BufferedEvent): string | null {
|
||||
}
|
||||
|
||||
export class SseHub {
|
||||
private subs = new Map<string, Set<Send>>();
|
||||
private subs = new Map<string, Set<Subscriber>>();
|
||||
private buffers = new Map<string, BufferedEvent[]>();
|
||||
private lastIds = new Map<string, number>();
|
||||
|
||||
@@ -42,7 +49,8 @@ export class SseHub {
|
||||
options: SubscribeOptions = {},
|
||||
): () => void {
|
||||
if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set());
|
||||
this.subs.get(sessionId)!.add(send);
|
||||
const subscriber = { send, close: options.close ?? (() => undefined), closed: false };
|
||||
this.subs.get(sessionId)!.add(subscriber);
|
||||
|
||||
const afterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0
|
||||
? options.afterId ?? 0
|
||||
@@ -63,12 +71,14 @@ export class SseHub {
|
||||
send(item.event, item.data, item.id);
|
||||
}
|
||||
|
||||
return () => this.subs.get(sessionId)?.delete(send);
|
||||
return () => this.subs.get(sessionId)?.delete(subscriber);
|
||||
}
|
||||
|
||||
publish(sessionId: string, event: string, data: object): number {
|
||||
const item = this.buffer(sessionId, event, data);
|
||||
for (const send of this.subs.get(sessionId) ?? []) send(event, data, item.id);
|
||||
for (const subscriber of this.subs.get(sessionId) ?? []) {
|
||||
subscriber.send(event, data, item.id);
|
||||
}
|
||||
return item.id;
|
||||
}
|
||||
|
||||
@@ -88,6 +98,13 @@ export class SseHub {
|
||||
|
||||
/** Clear buffered/runtime bindings while retaining session event-id monotonicity. */
|
||||
clear(sessionId: string): void {
|
||||
// Snapshot because ending an HTTP response can synchronously/asynchronously unsubscribe it.
|
||||
// Mark before invoking callbacks so even a re-entrant clear cannot close a response twice.
|
||||
for (const subscriber of [...(this.subs.get(sessionId) ?? [])]) {
|
||||
if (subscriber.closed) continue;
|
||||
subscriber.closed = true;
|
||||
try { subscriber.close(); } catch { /* disconnect every remaining subscriber */ }
|
||||
}
|
||||
this.buffers.delete(sessionId);
|
||||
this.subs.delete(sessionId);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user