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
3 changes: 3 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -819,6 +819,9 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
const nextText = Option.match(existingMessage, {
onNone: () => event.payload.text,
onSome: (message) => {
if (event.payload.reset === true) {
return event.payload.text;
}
if (event.payload.streaming) {
return `${message.text}${event.payload.text}`;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1703,7 +1703,93 @@ describe("ProviderCommandReactor", () => {
await waitFor(() => harness.interruptTurn.mock.calls.length === 1);
expect(harness.interruptTurn.mock.calls[0]?.[0]).toEqual({
threadId: "thread-1",
turnId: "turn-1",
});

const readModel = await harness.readModel();
const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(thread?.session?.status).toBe("ready");
expect(thread?.session?.activeTurnId).toBeNull();
});

it("settles a stuck running session when stop is pressed after the turn is already interrupted", async () => {
const harness = await createHarness();
const firstInterruptAt = "2026-01-01T00:00:00.000Z";
const repeatInterruptAt = "2026-01-01T00:00:01.000Z";
const threadId = ThreadId.make("thread-1");
const turnId = asTurnId("turn-1");

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-set-stuck"),
threadId,
session: {
threadId,
status: "running",
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: turnId,
lastError: null,
updatedAt: firstInterruptAt,
},
createdAt: firstInterruptAt,
}),
);

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.turn.interrupt",
commandId: CommandId.make("cmd-turn-interrupt-stuck"),
threadId,
turnId,
createdAt: firstInterruptAt,
}),
);

await waitFor(() => harness.interruptTurn.mock.calls.length === 1);

harness.interruptTurn.mockClear();

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("cmd-session-set-stuck-again"),
threadId,
session: {
threadId,
status: "running",
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: turnId,
lastError: null,
updatedAt: repeatInterruptAt,
},
createdAt: repeatInterruptAt,
}),
);

await Effect.runPromise(
harness.engine.dispatch({
type: "thread.turn.interrupt",
commandId: CommandId.make("cmd-turn-interrupt-stuck-repeat"),
threadId,
turnId,
createdAt: repeatInterruptAt,
}),
);

await waitFor(async () => {
const readModel = await harness.readModel();
const thread = readModel.threads.find((entry) => entry.id === threadId);
return thread?.session?.status === "ready" && thread?.session?.activeTurnId === null;
});

expect(harness.interruptTurn.mock.calls.length).toBe(0);
const readModel = await harness.readModel();
const thread = readModel.threads.find((entry) => entry.id === threadId);
expect(thread?.session?.status).toBe("ready");
expect(thread?.session?.activeTurnId).toBeNull();
});

it("starts a fresh session when only projected session state exists", async () => {
Expand Down
63 changes: 59 additions & 4 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -870,8 +870,8 @@ const make = Effect.gen(function* () {
if (!thread) {
return;
}
const hasSession = thread.session && thread.session.status !== "stopped";
if (!hasSession) {
const session = thread.session;
if (!session || session.status === "stopped") {
return yield* appendProviderFailureActivity({
threadId: event.payload.threadId,
kind: "provider.turn.interrupt.failed",
Expand All @@ -882,8 +882,63 @@ const make = Effect.gen(function* () {
});
}

// Orchestration turn ids are not provider turn ids, so interrupt by session.
yield* providerService.interruptTurn({ threadId: event.payload.threadId });
const now = event.payload.createdAt;
const turnId = event.payload.turnId;
const latestTurn = thread.latestTurn;
const sessionStillShowsActiveTurn =
session.status === "running" ||
session.status === "starting" ||
session.activeTurnId !== null;
const turnWasSettledBeforeThisInterrupt =
turnId !== undefined &&
latestTurn?.turnId === turnId &&
latestTurn.state !== "running" &&
latestTurn.completedAt !== null &&
latestTurn.completedAt < now;
const shouldSignalProvider =
!turnWasSettledBeforeThisInterrupt &&
(session.status === "running" ||
session.status === "starting" ||
session.activeTurnId !== null);

if (shouldSignalProvider) {
// Orchestration turn ids are not provider turn ids, so interrupt by session.
yield* providerService
.interruptTurn({
threadId: event.payload.threadId,
...(turnId !== undefined ? { turnId } : {}),
})
.pipe(
Effect.catchCause((cause) =>
appendProviderFailureActivity({
threadId: event.payload.threadId,
kind: "provider.turn.interrupt.failed",
summary: "Provider turn interrupt failed",
detail: Cause.pretty(cause),
turnId: turnId ?? null,
createdAt: now,
}),
),
);
}

// Projection marks the turn interrupted, but only provider runtime events
// (or an explicit session set) clear session.activeTurnId. Without this,
// stop leaves the thread "working" forever — especially when the provider
// process is already gone and interrupt could not reach it.
if (sessionStillShowsActiveTurn) {
yield* setThreadSession({
threadId: thread.id,
session: {
...session,
status: "ready",
activeTurnId: null,
lastError: null,
updatedAt: now,
},
createdAt: now,
});
}
});

const processApprovalResponseRequested = Effect.fn("processApprovalResponseRequested")(function* (
Expand Down
113 changes: 113 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2072,6 +2072,119 @@ describe("ProviderRuntimeIngestion", () => {
).toBe(" after approval");
});

it("resets a completed assistant segment when item.started arrives again after resume", async () => {
const harness = await createHarness({ serverSettings: { enableAssistantStreaming: true } });
const startedAt = "2026-03-28T08:00:00.000Z";
const resumedAt = "2026-03-28T08:00:10.000Z";

harness.emit({
type: "turn.started",
eventId: asEventId("evt-turn-started-resume-assistant-segment"),
provider: ProviderDriverKind.make("cursor"),
createdAt: startedAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-resume-assistant-segment"),
});
await waitForThread(
harness.readModel,
(thread) =>
thread.session?.status === "running" &&
thread.session?.activeTurnId === "turn-resume-assistant-segment",
);

const resumedItemId = asItemId("assistant:session-1:segment:3");
const resumedMessageId = "assistant:assistant:session-1:segment:3";

harness.emit({
type: "content.delta",
eventId: asEventId("evt-message-delta-resume-assistant-segment-initial"),
provider: ProviderDriverKind.make("cursor"),
createdAt: startedAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-resume-assistant-segment"),
itemId: resumedItemId,
payload: {
streamKind: "assistant_text",
delta: "stale summary before restart",
},
});
harness.emit({
type: "item.completed",
eventId: asEventId("evt-message-completed-resume-assistant-segment-initial"),
provider: ProviderDriverKind.make("cursor"),
createdAt: startedAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-resume-assistant-segment"),
itemId: resumedItemId,
payload: {
itemType: "assistant_message",
status: "completed",
},
});

await waitForThread(harness.readModel, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === resumedMessageId &&
!message.streaming &&
message.text === "stale summary before restart",
),
);

harness.emit({
type: "item.started",
eventId: asEventId("evt-message-started-resume-assistant-segment"),
provider: ProviderDriverKind.make("cursor"),
createdAt: resumedAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-resume-assistant-segment"),
itemId: resumedItemId,
payload: {
itemType: "assistant_message",
status: "inProgress",
},
});
harness.emit({
type: "content.delta",
eventId: asEventId("evt-message-delta-resume-assistant-segment-after-resume"),
provider: ProviderDriverKind.make("cursor"),
createdAt: resumedAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-resume-assistant-segment"),
itemId: resumedItemId,
payload: {
streamKind: "assistant_text",
delta: "fresh answer after resume",
},
});
harness.emit({
type: "item.completed",
eventId: asEventId("evt-message-completed-resume-assistant-segment-after-resume"),
provider: ProviderDriverKind.make("cursor"),
createdAt: resumedAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-resume-assistant-segment"),
itemId: resumedItemId,
payload: {
itemType: "assistant_message",
status: "completed",
},
});

const thread = await waitForThread(harness.readModel, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === resumedMessageId &&
!message.streaming &&
message.text === "fresh answer after resume",
),
);
expect(
thread.messages.find((message: ProviderRuntimeTestMessage) => message.id === resumedMessageId)
?.text,
).toBe("fresh answer after resume");
});

it("streams assistant deltas when thread.turn.start requests streaming mode", async () => {
const harness = await createHarness({ serverSettings: { enableAssistantStreaming: true } });
const now = "2026-01-01T00:00:00.000Z";
Expand Down
43 changes: 43 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,19 @@ function assistantSegmentMessageId(baseKey: string, segmentIndex: number): Messa
segmentIndex === 0 ? `assistant:${baseKey}` : `assistant:${baseKey}:segment:${segmentIndex}`,
);
}

function assistantMessageIdFromEvent(event: ProviderRuntimeEvent): MessageId {
return assistantSegmentMessageId(assistantSegmentBaseKeyFromEvent(event), 0);
}

function nextAssistantSegmentIndexFromItemId(itemId: string | undefined): number {
if (!itemId) {
return 1;
}
const match = itemId.match(/:segment:(\d+)$/);
return match ? Number(match[1]) + 1 : 1;
}

function buildContextWindowActivityPayload(
event: ProviderRuntimeEvent,
): ThreadTokenUsageSnapshot | undefined {
Expand Down Expand Up @@ -1365,6 +1378,36 @@ const make = Effect.gen(function* () {
const proposedPlanDelta =
event.type === "turn.proposed.delta" ? event.payload.delta : undefined;

if (event.type === "item.started" && event.payload.itemType === "assistant_message") {
const turnId = toTurnId(event.turnId);
if (turnId) {
const messageId = assistantMessageIdFromEvent(event);
const detailedThread = yield* getLoadedThreadDetail();
const existingMessage = detailedThread
? findMessageById(detailedThread.messages, messageId)
: undefined;
if (existingMessage && !existingMessage.streaming) {
yield* orchestrationEngine.dispatch({
type: "thread.message.assistant.reset",
commandId: yield* providerCommandId(event, "assistant-reset-on-item-started"),
threadId: thread.id,
messageId,
turnId,
createdAt: now,
});
}
yield* clearBufferedAssistantText(messageId);
yield* rememberAssistantMessageId(thread.id, turnId, messageId);
yield* setAssistantSegmentStateForTurn(thread.id, turnId, {
baseKey: assistantSegmentBaseKeyFromEvent(event),
nextSegmentIndex: nextAssistantSegmentIndexFromItemId(
event.itemId !== undefined ? String(event.itemId) : undefined,
),
activeMessageId: messageId,
});
}
}

if (assistantDelta && assistantDelta.length > 0) {
const turnId = toTurnId(event.turnId);
const assistantMessageId = yield* getOrCreateAssistantMessageId({
Expand Down
28 changes: 28 additions & 0 deletions apps/server/src/orchestration/decider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -654,6 +654,34 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"
};
}

case "thread.message.assistant.reset": {
yield* requireThread({
readModel,
command,
threadId: command.threadId,
});
return {
...(yield* withEventBase({
aggregateKind: "thread",
aggregateId: command.threadId,
occurredAt: command.createdAt,
commandId: command.commandId,
})),
type: "thread.message-sent",
payload: {
threadId: command.threadId,
messageId: command.messageId,
role: "assistant",
text: "",
turnId: command.turnId ?? null,
streaming: true,
reset: true,
createdAt: command.createdAt,
updatedAt: command.createdAt,
},
};
}

case "thread.proposed-plan.upsert": {
yield* requireThread({
readModel,
Expand Down
Loading