Skip to content
Closed
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
140 changes: 140 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4680,6 +4680,146 @@ describe("ProviderRuntimeIngestion", () => {
expect(activity?.payload).toMatchObject({ requestId: "message-compact" });
});

it("keeps background liveness through a provider resume until the follow-up turn starts", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
const threadId = asThreadId("thread-1");
const provider = ProviderDriverKind.make("claudeAgent");

await harness.emitAndDrain([
{
type: "task.started",
eventId: asEventId("evt-resume-task-started"),
provider,
createdAt: now,
threadId,
payload: {
taskId: "subagent-1",
taskType: "local_agent",
description: "Review the diff",
},
},
]);
expect((await harness.readThreadShell()).backgroundLiveness).toBe("working");

await harness.emitAndDrain([
{
type: "turn.completed",
eventId: asEventId("evt-resume-turn-completed"),
provider,
createdAt: now,
threadId,
turnId: asTurnId("turn-parent"),
payload: { state: "completed" },
},
]);
expect((await harness.readThreadShell()).session?.status).toBe("ready");
expect((await harness.readThreadShell()).backgroundLiveness).toBe("working");

await harness.emitAndDrain([
{
type: "task.completed",
eventId: asEventId("evt-resume-task-completed"),
provider,
createdAt: now,
threadId,
payload: {
taskId: "subagent-1",
status: "completed",
resumesProvider: true,
},
},
]);
const held = await harness.readThreadShell();
expect(held.session?.status).toBe("ready");
expect(held.backgroundLiveness).toBe("working");

await harness.emitAndDrain([
{
type: "turn.started",
eventId: asEventId("evt-resume-follow-up"),
provider,
createdAt: "2026-01-01T00:00:02.000Z",
threadId,
turnId: asTurnId("turn-follow-up"),
},
]);
const resumed = await harness.readThreadShell();
expect(resumed.session?.status).toBe("running");
expect(resumed.backgroundLiveness).toBeNull();
});

it("clears background liveness when a completion will not resume the provider", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
const threadId = asThreadId("thread-1");
const provider = ProviderDriverKind.make("claudeAgent");

await harness.emitAndDrain([
{
type: "task.started",
eventId: asEventId("evt-settle-task-started"),
provider,
createdAt: now,
threadId,
payload: {
taskId: "monitor-1",
taskType: "local_bash",
description: "Watch the checks",
},
},
{
type: "task.completed",
eventId: asEventId("evt-settle-task-completed"),
provider,
createdAt: now,
threadId,
payload: {
taskId: "monitor-1",
status: "completed",
},
},
]);
expect((await harness.readThreadShell()).backgroundLiveness).toBeNull();
});

it("drops a provider-resume hold when the session errors", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
const threadId = asThreadId("thread-1");
const provider = ProviderDriverKind.make("claudeAgent");

await harness.emitAndDrain([
{
type: "task.completed",
eventId: asEventId("evt-error-task-completed"),
provider,
createdAt: now,
threadId,
payload: {
taskId: "subagent-1",
status: "completed",
resumesProvider: true,
},
},
]);
expect((await harness.readThreadShell()).backgroundLiveness).toBe("working");

await harness.emitAndDrain([
{
type: "session.state.changed",
eventId: asEventId("evt-resume-session-error"),
provider,
createdAt: now,
threadId,
payload: { state: "error", reason: "provider crashed" },
},
]);
const failed = await harness.readThreadShell();
expect(failed.session?.status).toBe("error");
expect(failed.backgroundLiveness).toBeNull();
});

it("projects Codex task lifecycle chunks into thread activities", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
28 changes: 28 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1778,6 +1778,12 @@ const make = Effect.gen(function* () {
},
);

/**
* Fold one provider runtime event into the thread: session lifecycle,
* activities, and in-memory background liveness. A waking task completion
* keeps that liveness working until the follow-up turn starts, so the
* shell does not read as ready in the gap.
*/
const processRuntimeEvent = (event: ProviderRuntimeEvent) =>
Effect.gen(function* () {
if (
Expand Down Expand Up @@ -2504,6 +2510,7 @@ const make = Effect.gen(function* () {
taskType?: string;
status?: string;
agentId?: string;
resumesProvider?: boolean;
};
threadBackgroundLiveness.recordTaskLiveness({
threadId: thread.id,
Expand All @@ -2519,9 +2526,30 @@ const make = Effect.gen(function* () {
: event.type === "task.updated"
? "updated"
: "completed",
awaitsProviderResume:
event.type === "task.completed" && payload.resumesProvider === true,
});
break;
}
case "turn.started":
// The follow-up turn is running. Session status keeps the thread
// working until that turn settles, so the resume hold can drop.
if (shouldApplyThreadLifecycle) {
threadBackgroundLiveness.releaseProviderResume(thread.id);
}
break;
case "turn.aborted":
if (shouldApplyThreadLifecycle) {
threadBackgroundLiveness.releaseProviderResume(thread.id);
}
break;
case "session.state.changed":
// Ready is the gap itself. Release only when the session can no
// longer resume, so a failure is not pinned on Working.
if (event.payload.state === "error" || event.payload.state === "stopped") {
threadBackgroundLiveness.releaseProviderResume(thread.id);
}
break;
case "session.exited":
threadBackgroundLiveness.clearThreadLiveness(thread.id);
break;
Expand Down
97 changes: 97 additions & 0 deletions apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -226,4 +226,101 @@ describe("ThreadBackgroundLiveness", () => {
a.clearThreadLiveness("t");
expect(a.getThreadBackgroundLiveness("t")).toBeNull();
});

it("holds working after a waking completion until the provider resumes", () => {
const liveness = ThreadBackgroundLiveness.make();
const threadId = "thread-resume";
liveness.recordTaskLiveness({
threadId,
taskId: "subagent",
taskType: "local_agent",
status: undefined,
kind: "started",
});
liveness.recordTaskLiveness({
threadId,
taskId: "monitor",
taskType: "local_bash",
status: undefined,
kind: "started",
});
liveness.recordTaskLiveness({
threadId,
taskId: "subagent",
taskType: "local_agent",
status: "completed",
kind: "completed",
awaitsProviderResume: true,
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("monitoring");
liveness.recordTaskLiveness({
threadId,
taskId: "monitor",
taskType: "local_bash",
status: "completed",
kind: "completed",
awaitsProviderResume: true,
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working");
liveness.releaseProviderResume(threadId);
expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull();
});

it("does not hold a completion that will not resume the provider", () => {
const liveness = ThreadBackgroundLiveness.make();
const threadId = "thread-settle";
liveness.recordTaskLiveness({
threadId,
taskId: "monitor",
taskType: "local_bash",
status: undefined,
kind: "started",
});
liveness.recordTaskLiveness({
threadId,
taskId: "monitor",
taskType: "local_bash",
status: "completed",
kind: "completed",
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull();
});

it("drops a resume hold when the thread's background work is cleared", () => {
const liveness = ThreadBackgroundLiveness.make();
liveness.recordTaskLiveness({
threadId: "thread",
taskId: "subagent",
taskType: "local_agent",
status: "completed",
kind: "completed",
awaitsProviderResume: true,
});
expect(liveness.getThreadBackgroundLiveness("thread")).toBe("working");
liveness.clearThreadLiveness("thread");
expect(liveness.getThreadBackgroundLiveness("thread")).toBeNull();
});

it("releases only the resume hold and leaves a task that started during the handoff", () => {
const liveness = ThreadBackgroundLiveness.make();
const threadId = "thread-handoff";
liveness.recordTaskLiveness({
threadId,
taskId: "subagent",
taskType: "local_agent",
status: "completed",
kind: "completed",
awaitsProviderResume: true,
});
liveness.recordTaskLiveness({
threadId,
taskId: "monitor",
taskType: "local_bash",
status: undefined,
kind: "started",
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("working");
liveness.releaseProviderResume(threadId);
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("monitoring");
});
});
Loading
Loading