Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,13 @@ import type {
} from "@first-tree/shared";
import { encodeProviderRetryEventMessage, parseProviderRetryEventMessage } from "@first-tree/shared";
import { afterEach, describe, expect, it, vi } from "vitest";
import type { AgentHandler, HandlerFactory, SessionContext, SessionMessage } from "../runtime/handler.js";
import type {
AgentHandler,
DeliveryToken,
HandlerFactory,
SessionContext,
SessionMessage,
} from "../runtime/handler.js";
import type { DeliveryDecision, DeliveryRouteOwnership, DeliveryWork } from "../runtime/inbox-delivery-coordinator.js";
import type { SubprocessProbe } from "../runtime/process-tree-probe.js";
import { SessionManager } from "../runtime/session-manager.js";
Expand Down Expand Up @@ -54,6 +60,15 @@ type SessionManagerInternals = {
terminatePersistFailures: Set<string>;
awaitingResetFenceRelease: Set<string>;
pendingTeardowns: Map<string, Set<AgentHandler>>;
quarantinedSessions: Map<
string,
{
handler: AgentHandler;
generation: number;
reason: "operator_suspend_timeout";
routeTransitionInFlight: boolean;
}
>;
routeProducers: Map<string, Set<Promise<void>>>;
registry: SessionRegistry | null;
pendingQueue: Array<{ message: SessionMessage | null; chatId: string; deliveryKind: string }>;
Expand Down Expand Up @@ -3257,6 +3272,316 @@ describe("SessionManager edge coverage", () => {
expect(i.sessions.has(chatId)).toBe(false);
});

it("quarantines a timed-out operator suspend and recovers real inbox debt before a fresh handler", async () => {
vi.useFakeTimers();
let initialCtx: SessionContext | undefined;
let initialHead: SessionMessage | undefined;
const oldHandler = handler({
start: vi.fn().mockImplementation(async (message, ctx, token) => {
initialCtx = ctx;
initialHead = message;
token?.processingStarted(message);
return { sessionId: "established-session", route: { kind: "owned" as const, mode: "queued" as const } };
}),
inject: vi.fn().mockReturnValue({ kind: "owned", mode: "queued" } as const),
suspend: vi.fn(() => new Promise<void>(() => {})),
shutdown: vi.fn(() => new Promise<void>(() => {})),
});
let freshCtx: SessionContext | undefined;
let freshMessage: SessionMessage | undefined;
const freshHandler = handler({
resume: vi.fn().mockImplementation(async (message, _sessionId, ctx, token) => {
freshCtx = ctx;
freshMessage = message;
if (message) token?.processingStarted(message);
return { sessionId: "fresh-session", route: { kind: "owned" as const, mode: "queued" as const } };
}),
});
const ackEntry = vi.fn<(entryId: number) => Promise<void>>().mockResolvedValue(undefined);
const recoverChat = vi.fn<(chatId: string) => Promise<void>>().mockResolvedValue(undefined);
const onSessionEvent = vi.fn<(chatId: string, event: SessionEvent) => void>();
const sm = makeManager({ handlers: [oldHandler, freshHandler], ackEntry, recoverChat, onSessionEvent });
const i = internals(sm);
const chatId = "chat-timeout-recovery-before-fresh-handler";
const headEntry = mockEntry({ id: 9100, chatId, messageId: "msg-timeout-head" });
const queuedTailEntry = mockEntry({ id: 9101, chatId, messageId: "msg-timeout-queued-tail" });

await sm.dispatch(headEntry);
if (!initialCtx || !initialHead) throw new Error("initial route was not captured");
await sm.dispatch(queuedTailEntry);
expect(i.sessions.get(chatId)?.routeTransition).toBe(null);

await sm.handleCommand(chatId, "session:suspend");
const laterDispatch = sm.dispatch(mockEntry({ id: 9102, chatId, messageId: "msg-after-timeout" }));

await vi.advanceTimersByTimeAsync(29_999);
expect(freshHandler.resume).not.toHaveBeenCalled();

await vi.advanceTimersByTimeAsync(1);
await laterDispatch;

expect(ackEntry).toHaveBeenCalledWith(9100);
expect(ackEntry).not.toHaveBeenCalledWith(9101);
expect(recoverChat).toHaveBeenCalledTimes(1);
expect(freshHandler.resume).not.toHaveBeenCalled();
expect(oldHandler.shutdown).not.toHaveBeenCalled();
expect(i.pendingTeardowns.has(chatId)).toBe(false);
expect(i.quarantinedSessions.get(chatId)).toEqual(
expect.objectContaining({
handler: oldHandler,
reason: "operator_suspend_timeout",
routeTransitionInFlight: false,
}),
);
expect(onSessionEvent).toHaveBeenCalledWith(
chatId,
expect.objectContaining({
payload: expect.objectContaining({
message: expect.stringContaining('"routeTransitionInFlight":false'),
}),
}),
);

await sm.dispatch(queuedTailEntry);

const recoveryOrder = recoverChat.mock.invocationCallOrder[0];
const resumeOrder = vi.mocked(freshHandler.resume).mock.invocationCallOrder[0];
expect(recoveryOrder).toBeDefined();
expect(resumeOrder).toBeDefined();
expect(Number(recoveryOrder)).toBeLessThan(Number(resumeOrder));
expect(i.sessions.get(chatId)?.handler).toBe(freshHandler);
if (!freshCtx || !freshMessage) throw new Error("fresh route was not captured");
await freshCtx.finishTurn(freshMessage, { status: "success", terminal: true });

await sm.shutdown();
});

it("preserves terminal notice debt when operator suspend emits failure and then times out", async () => {
vi.useFakeTimers();
let initialCtx: SessionContext | undefined;
const oldHandler = handler({
start: vi.fn().mockImplementation(async (message, ctx, token) => {
initialCtx = ctx;
token?.processingStarted(message);
return { sessionId: "established-session", route: { kind: "owned" as const, mode: "queued" as const } };
}),
suspend: vi.fn().mockImplementation(() => {
if (!initialCtx) throw new Error("initial route context was not captured");
initialCtx.emitEvent({
kind: "error",
payload: {
source: "runtime",
message: encodeProviderRetryEventMessage({
event: "provider_failure_terminal",
provider: "codex",
scope: "provider_turn",
category: "credential",
reasonCode: "provider_credential_required",
replaySafety: "provider_entered",
userSeverity: "error",
messagePreview: "refresh token revoked while suspending",
}),
},
});
return new Promise<void>(() => {});
}),
shutdown: vi.fn(() => new Promise<void>(() => {})),
});
const ackEntry = vi.fn<(entryId: number) => Promise<void>>().mockResolvedValue(undefined);
const recoverChat = vi.fn<(chatId: string) => Promise<void>>().mockResolvedValue(undefined);
const sendMessage = vi.fn().mockResolvedValue({ id: "runtime-notice-after-timeout" });
const sdk = { ...mockSdk(), sendMessage } as unknown as FirstTreeHubSDK;
const sm = makeManager({ handlers: [oldHandler], ackEntry, recoverChat, sdk });
const i = internals(sm);
const chatId = "chat-timeout-terminal-notice-debt";
const headEntry = mockEntry({ id: 9110, chatId, messageId: "msg-timeout-terminal-notice" });

await sm.dispatch(headEntry);
await sm.handleCommand(chatId, "session:suspend");
await vi.advanceTimersByTimeAsync(30_000);
await i.sessions.get(chatId)?.suspending;

expect(ackEntry).not.toHaveBeenCalled();
expect(i.inboxDelivery.snapshot(chatId)).toMatchObject({
entries: [{ entryId: 9110, messageId: headEntry.message.id, phase: "owned" }],
recoveryDebt: "required",
});

// The first frame opens recovery; the server's redelivery then settles
// the retained notice debt without re-entering the provider.
await sm.dispatch(headEntry);
expect(recoverChat).toHaveBeenCalledTimes(1);
expect(sendMessage).not.toHaveBeenCalled();
expect(ackEntry).not.toHaveBeenCalled();

await sm.dispatch(headEntry);
await vi.waitFor(() => expect(ackEntry).toHaveBeenCalledTimes(1));

expect(sendMessage).toHaveBeenCalledTimes(1);
expect(ackEntry).toHaveBeenCalledWith(9110);
expect(oldHandler.start).toHaveBeenCalledTimes(1);
const noticeOrder = sendMessage.mock.invocationCallOrder[0];
const ackOrder = ackEntry.mock.invocationCallOrder[0];
if (noticeOrder === undefined || ackOrder === undefined) throw new Error("expected notice and ACK order");
expect(noticeOrder).toBeLessThan(ackOrder);
expect(oldHandler.shutdown).not.toHaveBeenCalled();

await sm.shutdown();
});

it("keeps later suspend and resume cycles healthy while the quarantined callback never settles", async () => {
vi.useFakeTimers();
const oldHandler = handler({
suspend: vi.fn(() => new Promise<void>(() => {})),
shutdown: vi.fn(() => new Promise<void>(() => {})),
});
const freshHandler = handler({
resume: vi
.fn()
.mockResolvedValue({ sessionId: "fresh-session", route: { kind: "owned" as const, mode: "queued" as const } }),
});
const sm = makeManager({ handlers: [freshHandler] });
const i = internals(sm);
const chatId = "chat-quarantine-repeated-resume";
i.sessions.set(chatId, makeSessionRecord(chatId, { handler: oldHandler, status: "active" }));
i._activeCount = 1;

await sm.handleCommand(chatId, "session:suspend");
await vi.advanceTimersByTimeAsync(30_000);
await i.sessions.get(chatId)?.suspending;
await sm.handleCommand(chatId, "session:resume");

expect(freshHandler.resume).toHaveBeenCalledTimes(1);
expect(i.sessions.get(chatId)?.handler).toBe(freshHandler);
expect(oldHandler.shutdown).not.toHaveBeenCalled();

await sm.handleCommand(chatId, "session:suspend");
await i.sessions.get(chatId)?.suspending;
await sm.handleCommand(chatId, "session:resume");

expect(freshHandler.resume).toHaveBeenCalledTimes(2);
expect(i.sessions.get(chatId)?.handler).toBe(freshHandler);
expect(i.quarantinedSessions.get(chatId)?.handler).toBe(oldHandler);
expect(oldHandler.shutdown).not.toHaveBeenCalled();

await sm.shutdown();
});

it("fences late output and route adoption from a quarantined generation", async () => {
vi.useFakeTimers();
let signalResumeStarted: (() => void) | undefined;
let resolveResume: (() => void) | undefined;
let resolveSuspend: (() => void) | undefined;
let oldCtx: SessionContext | undefined;
let oldToken: DeliveryToken | undefined;
let oldMessage: SessionMessage | undefined;
const resumeStarted = new Promise<void>((resolve) => {
signalResumeStarted = resolve;
});
const resumeGate = new Promise<void>((resolve) => {
resolveResume = resolve;
});
const suspendGate = new Promise<void>((resolve) => {
resolveSuspend = resolve;
});
const oldHandler = handler({
resume: vi.fn().mockImplementation(async (message, _sessionId, ctx, token) => {
oldCtx = ctx;
oldToken = token;
oldMessage = message;
token?.processingStarted(message);
signalResumeStarted?.();
await resumeGate;
return { sessionId: "late-session", route: { kind: "owned" as const, mode: "queued" as const } };
}),
suspend: vi.fn(() => suspendGate),
shutdown: vi.fn(() => new Promise<void>(() => {})),
});
const freshHandler = handler();
const ackEntry = vi.fn<(entryId: number) => Promise<void>>().mockResolvedValue(undefined);
const onSessionEvent = vi.fn<(chatId: string, event: SessionEvent) => void>();
const sm = makeManager({ handlers: [freshHandler], ackEntry, onSessionEvent });
const i = internals(sm);
const chatId = "chat-quarantined-late-generation";
i.sessions.set(
chatId,
makeSessionRecord(chatId, {
handler: oldHandler,
status: "suspended",
claudeSessionId: "existing-session",
}),
);

const resumeDispatch = sm.dispatch(mockEntry({ id: 9120, chatId, messageId: "msg-late-generation" }));
await resumeStarted;
await sm.handleCommand(chatId, "session:suspend");
await vi.advanceTimersByTimeAsync(30_000);
await i.sessions.get(chatId)?.suspending;

expect(i.quarantinedSessions.get(chatId)).toEqual(
expect.objectContaining({
handler: oldHandler,
routeTransitionInFlight: true,
}),
);
if (!oldCtx || !oldToken || !oldMessage) throw new Error("old route output handles were not captured");
const eventCount = onSessionEvent.mock.calls.length;
const ackCount = ackEntry.mock.calls.length;
const lastActivity = i.sessions.get(chatId)?.lastActivity;
const trigger = i.currentTrigger.get(chatId);

oldCtx.emitEvent({ kind: "assistant_text", payload: { text: "late output" } });
oldCtx.recordProviderActivity();
await oldCtx.forwardResult("late result");
await oldToken.complete(oldMessage, { status: "success", terminal: true });
resolveSuspend?.();

expect(onSessionEvent).toHaveBeenCalledTimes(eventCount);
expect(ackEntry).toHaveBeenCalledTimes(ackCount);
expect(i.sessions.get(chatId)?.lastActivity).toBe(lastActivity);
expect(i.currentTrigger.get(chatId)).toEqual(trigger);
expect(i.sessions.get(chatId)?.claudeSessionId).toBe("existing-session");

await sm.dispatch(mockEntry({ id: 9121, chatId, messageId: "msg-provider-admission-fenced" }));
expect(freshHandler.start).not.toHaveBeenCalled();
expect(freshHandler.resume).not.toHaveBeenCalled();

await sm.shutdown();
expect(oldHandler.shutdown).not.toHaveBeenCalled();

resolveResume?.();
await resumeDispatch;
expect(i.sessions.has(chatId)).toBe(false);
expect(i.pendingTeardowns.has(chatId)).toBe(false);
});

it("returns restart-required Reset failure and bounds manager shutdown after suspend timeout", async () => {
vi.useFakeTimers();
const oldHandler = handler({
suspend: vi.fn(() => new Promise<void>(() => {})),
shutdown: vi.fn(() => new Promise<void>(() => {})),
});
const sm = makeManager();
const i = internals(sm);
const chatId = "chat-quarantine-reset-restart-required";
i.sessions.set(chatId, makeSessionRecord(chatId, { handler: oldHandler, status: "active" }));
i._activeCount = 1;

await sm.handleCommand(chatId, "session:suspend");
await vi.advanceTimersByTimeAsync(30_000);
await i.sessions.get(chatId)?.suspending;

await expect(sm.handleCommand(chatId, "session:terminate", { resetRef: "ref-restart" })).rejects.toThrow(
"Reset blocked: operator_suspend_timeout; provider teardown was not confirmed",
);
expect(i.sessions.has(chatId)).toBe(true);
expect(i.quarantinedSessions.has(chatId)).toBe(true);
expect(oldHandler.shutdown).not.toHaveBeenCalled();

await expect(sm.shutdown()).resolves.toBeUndefined();
expect(oldHandler.shutdown).not.toHaveBeenCalled();
});

it("manager shutdown stops a handler after a completed failed suspend with a falsey rejection", async () => {
const targetHandler = handler({
suspend: vi.fn().mockRejectedValue(undefined),
Expand Down
Loading
Loading