fix(bridge): forward Pi agent_end so the spinner stops at workflow completion
The FE derived 'working' purely as activeSession && !pendingWidget, so the
final workflow turn — the only one that ends without a follow-up gate —
left the spinner on forever (observed live: 21592s after F8 approve).
- SessionBridge maps Pi's agent_end -> SSE system_event {event: agent_end}
- PiProcessManager notifies the client (info error + synthetic agent_end)
when the child dies unexpectedly; expected teardowns stay silent
- sessionStore tracks agentActive (on: user entry/text_delta/ui_request,
off: agent_end); AppShell working now requires it; resume sets it
optimistically
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -34,6 +34,11 @@ export class SessionBridge {
|
|||||||
this.fan({ type: "text_delta", text: m.text ?? "" });
|
this.fan({ type: "text_delta", text: m.text ?? "" });
|
||||||
} else if (m.type === "system_event") {
|
} else if (m.type === "system_event") {
|
||||||
this.fan(m as ClientEvent);
|
this.fan(m as ClientEvent);
|
||||||
|
} else if (m.type === "agent_end") {
|
||||||
|
// Fine turno di Pi (pi-agent-core agent-loop): e' l'unico segnale che il turno e'
|
||||||
|
// concluso. Senza inoltrarlo, il FE resta "working" per sempre quando il turno
|
||||||
|
// finisce senza un gate successivo (ultimo step del workflow).
|
||||||
|
this.fan({ type: "system_event", event: "agent_end" });
|
||||||
}
|
}
|
||||||
// altri method nativi (setStatus/setWidget) e altri eventi Pi ignorati in MVP
|
// altri method nativi (setStatus/setWidget) e altri eventi Pi ignorati in MVP
|
||||||
});
|
});
|
||||||
@@ -43,6 +48,9 @@ export class SessionBridge {
|
|||||||
|
|
||||||
onClientEvent(cb: (e: ClientEvent) => void): void { this.cbs.add(cb); }
|
onClientEvent(cb: (e: ClientEvent) => void): void { this.cbs.add(cb); }
|
||||||
|
|
||||||
|
/** Eventi generati dal backend stesso (es. exit inatteso del child Pi), non da Pi. */
|
||||||
|
emitClientEvent(e: ClientEvent): void { this.fan(e); }
|
||||||
|
|
||||||
respond(uiResponse: object & { id: string }): void {
|
respond(uiResponse: object & { id: string }): void {
|
||||||
// Correla sull'id RPC di Pi; `value` porta l'uiResponse (con l'id del descriptor) cosi'
|
// Correla sull'id RPC di Pi; `value` porta l'uiResponse (con l'id del descriptor) cosi'
|
||||||
// il check interno del gate (resp.id === descriptor.id) regge.
|
// il check interno del gate (resp.id === descriptor.id) regge.
|
||||||
|
|||||||
@@ -70,7 +70,20 @@ export class PiProcessManager {
|
|||||||
const rt: SessionRuntime = { rpc, bridge, child };
|
const rt: SessionRuntime = { rpc, bridge, child };
|
||||||
this.runtimes.set(sessionId, rt);
|
this.runtimes.set(sessionId, rt);
|
||||||
// Identity-checked: a stale child's exit must not evict a newer runtime.
|
// Identity-checked: a stale child's exit must not evict a newer runtime.
|
||||||
child.on("exit", () => { if (this.runtimes.get(sessionId) === rt) this.runtimes.delete(sessionId); });
|
// 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) => {
|
||||||
|
if (this.runtimes.get(sessionId) === rt) {
|
||||||
|
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: "agent_end" });
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
const provider = o.provider ?? this.cfg.defaults.provider;
|
const provider = o.provider ?? this.cfg.defaults.provider;
|
||||||
const model = o.model ?? this.cfg.defaults.model;
|
const model = o.model ?? this.cfg.defaults.model;
|
||||||
|
|||||||
@@ -80,6 +80,33 @@ function recordingChild() {
|
|||||||
return ch;
|
return ch;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
test("un exit INATTESO del child notifica il client (info error + agent_end)", async () => {
|
||||||
|
const cfg = loadConfig({});
|
||||||
|
const child = recordingChild();
|
||||||
|
const mgr = new PiProcessManager(cfg, { spawnFn: () => child as any });
|
||||||
|
const rt = await mgr.spawnFor("crash-id", {});
|
||||||
|
const seen: any[] = [];
|
||||||
|
rt.bridge.onClientEvent((e) => seen.push(e));
|
||||||
|
child.emit("exit", 137);
|
||||||
|
expect(seen).toEqual([
|
||||||
|
{ type: "info", level: "error", text: expect.stringContaining("137") },
|
||||||
|
{ type: "system_event", event: "agent_end" },
|
||||||
|
]);
|
||||||
|
expect(mgr.count()).toBe(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("l'exit dopo teardown NON emette eventi al client (uscita attesa)", async () => {
|
||||||
|
const cfg = loadConfig({});
|
||||||
|
const child = recordingChild();
|
||||||
|
const mgr = new PiProcessManager(cfg, { spawnFn: () => child as any });
|
||||||
|
const rt = await mgr.spawnFor("stop-id", {});
|
||||||
|
const seen: any[] = [];
|
||||||
|
rt.bridge.onClientEvent((e) => seen.push(e));
|
||||||
|
mgr.teardown("stop-id");
|
||||||
|
child.emit("exit", 0);
|
||||||
|
expect(seen).toEqual([]);
|
||||||
|
});
|
||||||
|
|
||||||
test("spawnFor resume mode sends /riprendi-sessione <id>", async () => {
|
test("spawnFor resume mode sends /riprendi-sessione <id>", async () => {
|
||||||
const cfg = loadConfig({}); // no provider/model/thinking -> no rpc.request handshakes
|
const cfg = loadConfig({}); // no provider/model/thinking -> no rpc.request handshakes
|
||||||
const child = recordingChild();
|
const child = recordingChild();
|
||||||
|
|||||||
@@ -54,6 +54,17 @@ test("respond correla sull'id RPC di Pi (non sull'id del descriptor) e azzera il
|
|||||||
expect(b.pendingWidget()).toBeNull();
|
expect(b.pendingWidget()).toBeNull();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("agent_end di Pi diventa un system_event agent_end per il FE", () => {
|
||||||
|
const { rpc, fire } = fakeRpc();
|
||||||
|
const b = new SessionBridge(rpc);
|
||||||
|
const seen: any[] = [];
|
||||||
|
b.onClientEvent((e) => seen.push(e));
|
||||||
|
// Pi emette agent_end alla fine di ogni prompt (pi-agent-core agent-loop); e' il solo
|
||||||
|
// segnale di fine turno: senza mapparlo il FE non puo' mai uscire dallo stato "working".
|
||||||
|
fire({ type: "agent_end", messages: [] });
|
||||||
|
expect(seen).toEqual([{ type: "system_event", event: "agent_end" }]);
|
||||||
|
});
|
||||||
|
|
||||||
test("steer invia un comando steer", () => {
|
test("steer invia un comando steer", () => {
|
||||||
const { rpc, sent } = fakeRpc();
|
const { rpc, sent } = fakeRpc();
|
||||||
new SessionBridge(rpc).steer("considera solo il 2024");
|
new SessionBridge(rpc).steer("considera solo il 2024");
|
||||||
|
|||||||
@@ -76,6 +76,8 @@ export function AppShell() {
|
|||||||
// working spinner shows straight away; the backend calls run after.
|
// working spinner shows straight away; the backend calls run after.
|
||||||
setPanelSession(null);
|
setPanelSession(null);
|
||||||
resetSession();
|
resetSession();
|
||||||
|
// Optimistic: the resume POST is about to hand the ball to the harness.
|
||||||
|
setAgentActive(true);
|
||||||
setActiveSessionId(id);
|
setActiveSessionId(id);
|
||||||
try {
|
try {
|
||||||
await resumeSession(id);
|
await resumeSession(id);
|
||||||
@@ -152,13 +154,17 @@ export function AppShell() {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
// The harness "holds the ball" whenever a session is live and no widget is
|
// The harness "holds the ball" whenever a session is live, no widget is waiting on
|
||||||
// waiting on the human; once a widget appears, input is back in the user's court.
|
// the human, AND the Pi turn is still in flight (agentActive). Without the last
|
||||||
|
// condition the final workflow step — which ends with no follow-up gate — would
|
||||||
|
// leave the spinner on forever.
|
||||||
const pendingWidget = useSessionStore((s) => s.pendingWidget);
|
const pendingWidget = useSessionStore((s) => s.pendingWidget);
|
||||||
const resetSession = useSessionStore((s) => s.resetSession);
|
const resetSession = useSessionStore((s) => s.resetSession);
|
||||||
const setPhase = useSessionStore((s) => s.setPhase);
|
const setPhase = useSessionStore((s) => s.setPhase);
|
||||||
|
const setAgentActive = useSessionStore((s) => s.setAgentActive);
|
||||||
const lastSystemEvent = useSessionStore((s) => s.lastSystemEvent);
|
const lastSystemEvent = useSessionStore((s) => s.lastSystemEvent);
|
||||||
const working = Boolean(activeSessionId) && !pendingWidget;
|
const agentActive = useSessionStore((s) => s.agentActive);
|
||||||
|
const working = Boolean(activeSessionId) && !pendingWidget && agentActive;
|
||||||
// Processing time counts only while the harness works, not while a finalized
|
// Processing time counts only while the harness works, not while a finalized
|
||||||
// session sits idle or a gate awaits the reviewer (pendingWidget).
|
// session sits idle or a gate awaits the reviewer (pendingWidget).
|
||||||
const running = working && !finalized;
|
const running = working && !finalized;
|
||||||
|
|||||||
@@ -39,6 +39,35 @@ test("ui_request without a phase keeps the existing currentPhase", () => {
|
|||||||
expect(useSessionStore.getState().currentPhase).toBe("F2");
|
expect(useSessionStore.getState().currentPhase).toBe("F2");
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("agentActive lifecycle: off by default, on with activity, off on agent_end", () => {
|
||||||
|
const st = useSessionStore.getState();
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(false);
|
||||||
|
// Any streamed text means the turn is alive (also covers reattaching mid-turn).
|
||||||
|
st.applyEvent({ type: "text_delta", text: "Analisi" });
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(true);
|
||||||
|
// Pi's end-of-turn signal (mapped by the backend) releases the working state.
|
||||||
|
st.applyEvent({ type: "system_event", event: "agent_end" });
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("setLastUserEntry hands the ball back to the harness (agentActive on)", () => {
|
||||||
|
useSessionStore.getState().setLastUserEntry({ kind: "choice", text: "approve" });
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("setAgentActive drives the flag directly (resume optimistic spin)", () => {
|
||||||
|
useSessionStore.getState().setAgentActive(true);
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(true);
|
||||||
|
useSessionStore.getState().resetSession();
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a ui_request keeps agentActive on (turn blocked on the gate, not ended)", () => {
|
||||||
|
const st = useSessionStore.getState();
|
||||||
|
st.applyEvent({ type: "ui_request", ui_request: { id: "u1", widget: "select" } });
|
||||||
|
expect(useSessionStore.getState().agentActive).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
test("pushToast appends an error toast", () => {
|
test("pushToast appends an error toast", () => {
|
||||||
useSessionStore.getState().pushToast({ level: "error", text: "boom" });
|
useSessionStore.getState().pushToast({ level: "error", text: "boom" });
|
||||||
expect(useSessionStore.getState().toasts.at(-1)).toEqual({ level: "error", text: "boom" });
|
expect(useSessionStore.getState().toasts.at(-1)).toEqual({ level: "error", text: "boom" });
|
||||||
|
|||||||
@@ -15,12 +15,17 @@ interface SessionState {
|
|||||||
lastSystemEvent: StreamEvent | null;
|
lastSystemEvent: StreamEvent | null;
|
||||||
currentPhase: string | null;
|
currentPhase: string | null;
|
||||||
phaseError: string | null;
|
phaseError: string | null;
|
||||||
|
// True while a Pi turn is in flight. Turned off by the backend-forwarded `agent_end`
|
||||||
|
// system event — the only end-of-turn signal on the final workflow step, where no
|
||||||
|
// follow-up gate arrives to release the spinner.
|
||||||
|
agentActive: boolean;
|
||||||
applyEvent: (e: StreamEvent) => void;
|
applyEvent: (e: StreamEvent) => void;
|
||||||
clearPending: () => void;
|
clearPending: () => void;
|
||||||
resetSession: () => void;
|
resetSession: () => void;
|
||||||
setPhase: (phase: string | null) => void;
|
setPhase: (phase: string | null) => void;
|
||||||
pushToast: (toast: { level: string; text: string }) => void;
|
pushToast: (toast: { level: string; text: string }) => void;
|
||||||
setLastUserEntry: (e: { kind: "input" | "choice"; text: string }) => void;
|
setLastUserEntry: (e: { kind: "input" | "choice"; text: string }) => void;
|
||||||
|
setAgentActive: (v: boolean) => void;
|
||||||
}
|
}
|
||||||
|
|
||||||
const empty = {
|
const empty = {
|
||||||
@@ -32,6 +37,7 @@ const empty = {
|
|||||||
lastSystemEvent: null,
|
lastSystemEvent: null,
|
||||||
currentPhase: null as string | null,
|
currentPhase: null as string | null,
|
||||||
phaseError: null as string | null,
|
phaseError: null as string | null,
|
||||||
|
agentActive: false,
|
||||||
};
|
};
|
||||||
|
|
||||||
export const useSessionStore = create<SessionState>((set) => ({
|
export const useSessionStore = create<SessionState>((set) => ({
|
||||||
@@ -41,6 +47,8 @@ export const useSessionStore = create<SessionState>((set) => ({
|
|||||||
if (e.type === "ui_request")
|
if (e.type === "ui_request")
|
||||||
return {
|
return {
|
||||||
pendingWidget: e.ui_request,
|
pendingWidget: e.ui_request,
|
||||||
|
// The turn is blocked on the gate, not ended: keep the agent marked active.
|
||||||
|
agentActive: true,
|
||||||
currentPhase: e.ui_request.phase
|
currentPhase: e.ui_request.phase
|
||||||
? e.ui_request.phase.split("_")[0]
|
? e.ui_request.phase.split("_")[0]
|
||||||
: st.currentPhase,
|
: st.currentPhase,
|
||||||
@@ -52,19 +60,25 @@ export const useSessionStore = create<SessionState>((set) => ({
|
|||||||
const last = t.at(-1);
|
const last = t.at(-1);
|
||||||
if (last && !st.pendingWidget) t[t.length - 1] = { role: "assistant", text: last.text + e.text };
|
if (last && !st.pendingWidget) t[t.length - 1] = { role: "assistant", text: last.text + e.text };
|
||||||
else t.push({ role: "assistant", text: e.text });
|
else t.push({ role: "assistant", text: e.text });
|
||||||
return { transcript: t };
|
// Streamed text means the turn is alive (also covers reattaching mid-turn).
|
||||||
|
return { transcript: t, agentActive: true };
|
||||||
}
|
}
|
||||||
if (e.type === "info") {
|
if (e.type === "info") {
|
||||||
const stepMessages = [...st.stepMessages, { level: e.level ?? "info", text: e.text }];
|
const stepMessages = [...st.stepMessages, { level: e.level ?? "info", text: e.text }];
|
||||||
// An error during the active phase marks that phase red until the next gate.
|
// An error during the active phase marks that phase red until the next gate.
|
||||||
return e.level === "error" ? { stepMessages, phaseError: st.currentPhase } : { stepMessages };
|
return e.level === "error" ? { stepMessages, phaseError: st.currentPhase } : { stepMessages };
|
||||||
}
|
}
|
||||||
if (e.type === "system_event") return { lastSystemEvent: e };
|
if (e.type === "system_event")
|
||||||
|
return e.event === "agent_end"
|
||||||
|
? { lastSystemEvent: e, agentActive: false }
|
||||||
|
: { lastSystemEvent: e };
|
||||||
return {};
|
return {};
|
||||||
}),
|
}),
|
||||||
clearPending: () => set({ pendingWidget: null }),
|
clearPending: () => set({ pendingWidget: null }),
|
||||||
resetSession: () => set({ ...empty }),
|
resetSession: () => set({ ...empty }),
|
||||||
setPhase: (phase) => set({ currentPhase: phase }),
|
setPhase: (phase) => set({ currentPhase: phase }),
|
||||||
pushToast: (toast) => set((st) => ({ toasts: [...st.toasts, toast] })),
|
pushToast: (toast) => set((st) => ({ toasts: [...st.toasts, toast] })),
|
||||||
setLastUserEntry: (e) => set({ lastUserEntry: e, stepMessages: [] }),
|
// A user entry (question, gate choice, steer) hands the ball back to the harness.
|
||||||
|
setLastUserEntry: (e) => set({ lastUserEntry: e, stepMessages: [], agentActive: true }),
|
||||||
|
setAgentActive: (v) => set({ agentActive: v }),
|
||||||
}));
|
}));
|
||||||
|
|||||||
Reference in New Issue
Block a user