fix(sse): buffer session events for late subscribers + clear on close
The Pi process produces events immediately after session creation, but the browser's SSE connection may not be open yet (React re-render delay, navigation). The hub discarded events with no subscribers, so the user saw a blank session. - Ring-buffer up to 200 events per session; replay on subscribe - hub.clear(id) on POST /sessions/:id/close frees memory
This commit is contained in:
@@ -62,7 +62,9 @@ export function sessionRoutes(
|
|||||||
return reply.code(200).send({ id });
|
return reply.code(200).send({ id });
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/close", async (req) => {
|
app.post("/sessions/:id/close", async (req) => {
|
||||||
d.mgr.teardown((req.params as any).id);
|
const id = (req.params as { id: string }).id;
|
||||||
|
d.mgr.teardown(id);
|
||||||
|
d.hub.clear(id);
|
||||||
return { closed: true };
|
return { closed: true };
|
||||||
});
|
});
|
||||||
app.get("/sessions/:id/events", (req, reply) => {
|
app.get("/sessions/:id/events", (req, reply) => {
|
||||||
|
|||||||
@@ -1,15 +1,50 @@
|
|||||||
type Send = (event: string, data: object) => void;
|
type Send = (event: string, data: object) => void;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Per-session ring buffer of recent events.
|
||||||
|
*
|
||||||
|
* When the browser opens the SSE connection slightly after session creation
|
||||||
|
* (React re-render, navigation, etc.), events produced by the Pi process in
|
||||||
|
* that gap would be lost. The buffer replays them to late subscribers.
|
||||||
|
*
|
||||||
|
* `system_event` with `event: "agent_end"` acts as the natural sentinel:
|
||||||
|
* once a subscriber sees it, the buffer for that session is safe to clear.
|
||||||
|
*/
|
||||||
|
const BUFFER_LIMIT = 200;
|
||||||
|
|
||||||
|
interface BufferedEvent { event: string; data: object }
|
||||||
|
|
||||||
export class SseHub {
|
export class SseHub {
|
||||||
private subs = new Map<string, Set<Send>>();
|
private subs = new Map<string, Set<Send>>();
|
||||||
|
private buffers = new Map<string, BufferedEvent[]>();
|
||||||
|
|
||||||
subscribe(sessionId: string, send: Send, pending?: object | null): () => void {
|
subscribe(sessionId: string, send: Send, pending?: object | null): () => void {
|
||||||
if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set());
|
if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set());
|
||||||
this.subs.get(sessionId)!.add(send);
|
this.subs.get(sessionId)!.add(send);
|
||||||
|
|
||||||
|
// Replay buffered events so late subscribers don't miss the session lifecycle.
|
||||||
|
const buf = this.buffers.get(sessionId);
|
||||||
|
if (buf) for (const { event, data } of buf) send(event, data);
|
||||||
|
|
||||||
// Match the live-event shape (hub.publish sends the full ClientEvent):
|
// Match the live-event shape (hub.publish sends the full ClientEvent):
|
||||||
// both carry { type, ui_request } so re-emit and live widgets are identical.
|
// both carry { type, ui_request } so re-emit and live widgets are identical.
|
||||||
if (pending) send("ui_request", { type: "ui_request", ui_request: pending });
|
if (pending) send("ui_request", { type: "ui_request", ui_request: pending });
|
||||||
return () => this.subs.get(sessionId)?.delete(send);
|
return () => this.subs.get(sessionId)?.delete(send);
|
||||||
}
|
}
|
||||||
|
|
||||||
publish(sessionId: string, event: string, data: object): void {
|
publish(sessionId: string, event: string, data: object): void {
|
||||||
|
// Buffer the event for late subscribers.
|
||||||
|
let buf = this.buffers.get(sessionId);
|
||||||
|
if (!buf) { buf = []; this.buffers.set(sessionId, buf); }
|
||||||
|
buf.push({ event, data });
|
||||||
|
if (buf.length > BUFFER_LIMIT) buf.splice(0, buf.length - BUFFER_LIMIT);
|
||||||
|
|
||||||
for (const s of this.subs.get(sessionId) ?? []) s(event, data);
|
for (const s of this.subs.get(sessionId) ?? []) s(event, data);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Clear buffer and subscribers for a finished session. */
|
||||||
|
clear(sessionId: string): void {
|
||||||
|
this.buffers.delete(sessionId);
|
||||||
|
this.subs.delete(sessionId);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user