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 @@ -2005,6 +2005,100 @@ describe("ProviderCommandReactor", () => {
}),
);

effectIt.effect(
"keeps exact admission pending while a provider reports a running prior turn",
() =>
Effect.gen(function* () {
const sendEntered = yield* Deferred.make<void>();
const releaseSend = yield* Deferred.make<void>();
const requestId = CommandId.make("cmd-running-provider-exact-admission");
const messageId = asMessageId("message-running-provider-exact-admission");
const harness = yield* Effect.promise(() =>
createHarness({
startSessionEffect: (session) => Effect.succeed({ ...session, status: "running" }),
sendTurnEffect: () =>
Deferred.succeed(sendEntered, undefined).pipe(
Effect.andThen(Deferred.await(releaseSend)),
Effect.as({ threadId: ThreadId.make("thread-1"), turnId: asTurnId("new-turn") }),
),
}),
);

yield* harness.engine.dispatch({
type: "thread.turn.start",
commandId: requestId,
threadId: ThreadId.make("thread-1"),
message: { messageId, role: "user", text: "follow up", attachments: [] },
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
createdAt: isoAt(0),
});
yield* Deferred.await(sendEntered);
let thread = (yield* Effect.promise(() => harness.readModel())).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(thread?.session?.status).toBe("starting");
expect(thread?.session?.pendingTurnRequestId).toBe(requestId);
expect(thread?.session?.pendingTurnMessageId).toBe(messageId);

yield* Deferred.succeed(releaseSend, undefined);
yield* Effect.promise(() =>
waitFor(async () => {
const thread = (await harness.readModel()).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
return thread?.session?.activeTurnId === asTurnId("new-turn");
}),
);
thread = (yield* Effect.promise(() => harness.readModel())).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(thread?.session?.status).toBe("running");
expect(thread?.session?.activeTurnRequestId).toBe(requestId);
expect(thread?.session?.pendingTurnRequestId).toBeUndefined();
}),
);

for (const status of ["error", "closed"] as const) {
effectIt.effect(`rejects exact admission when the provider session is ${status}`, () =>
Effect.gen(function* () {
const requestId = CommandId.make(`cmd-terminal-provider-${status}`);
const harness = yield* Effect.promise(() =>
createHarness({
startSessionEffect: (session) => Effect.succeed({ ...session, status }),
}),
);
yield* harness.engine.dispatch({
type: "thread.turn.start",
commandId: requestId,
threadId: ThreadId.make("thread-1"),
message: {
messageId: asMessageId(`message-terminal-provider-${status}`),
role: "user",
text: "follow up",
attachments: [],
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
createdAt: isoAt(0),
});
yield* Effect.promise(() =>
waitFor(async () => {
const thread = (await harness.readModel()).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
return thread?.session?.status === "error";
}),
);
const thread = (yield* Effect.promise(() => harness.readModel())).threads.find(
(entry) => entry.id === ThreadId.make("thread-1"),
);
expect(thread?.session?.failedTurnRequestId).toBe(requestId);
expect(harness.sendTurn).not.toHaveBeenCalled();
}),
);
}

effectIt.effect("does not disarm when raw turn start ingestion never accepts the CAS", () =>
Effect.gen(function* () {
const testClock = yield* TestClock.make();
Expand Down
15 changes: 14 additions & 1 deletion apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1048,10 +1048,23 @@ const make = Effect.gen(function* () {
detail: `Provider session '${session.threadId}' started without a provider instance or session incarnation id.`,
});
}
if (
options?.pendingTurnStart === true &&
(session.status === "error" || session.status === "closed")
) {
return yield* new ProviderAdapterRequestError({
provider: providerErrorLabel(session.provider),
method: "thread.turn.start",
detail: `Provider session '${session.threadId}' is ${session.status}; it cannot accept a new turn.`,
});
}
const sessionBinding: OrchestrationSession = {
threadId,
// The provider can still report the previous turn as running while
// a new exact admission is reserved. The reservation owns the
// projected lifecycle until its matching turn.started is accepted.
status:
options?.pendingTurnStart === true && session.status === "ready"
options?.pendingTurnStart === true
? "starting"
: mapProviderSessionStatusToOrchestrationStatus(session.status),
providerName: session.provider,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1663,6 +1663,27 @@ describe("ProviderRuntimeIngestion", () => {
expect(thread?.session?.pendingTurnRequestId).toBe(requestId);
expect(thread?.messages.some((message) => message.role === "assistant")).toBe(false);

// The prior turn can finish on the same provider session after the new
// admission is reserved. Its terminal event must not clear that pending
// admission before the matching turn.started arrives.
await harness.emitAndDrain([
{
type: "turn.completed",
eventId: asEventId("evt-runtime-session-b-prior-turn-completed"),
provider: ProviderDriverKind.make("codex"),
providerInstanceId: ProviderInstanceId.make("codex"),
threadId,
turnId: asTurnId("turn-runtime-session-b-prior"),
admissionRequestId: CommandId.make("cmd-runtime-session-b-prior"),
sessionIncarnationId: sessionB,
createdAt: "2026-01-01T00:00:02.500Z",
payload: { state: "completed" },
},
]);
thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId);
expect(thread?.session?.status).toBe("starting");
expect(thread?.session?.pendingTurnRequestId).toBe(requestId);

const currentTurnId = asTurnId("turn-runtime-session-b");
harness.emit({
type: "turn.started",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,12 @@ it.layer(NodeServices.layer)("session lifecycle CAS decider", (it) => {
readModel: makeReadModel(current),
});
expect(stale).toEqual([]);

const runningWithPendingRequest = yield* decideOrchestrationCommand({
command: { ...command, session: { ...acceptedSession, status: "running" } },
readModel: makeReadModel(current),
});
expect(runningWithPendingRequest).toEqual([]);
}),
);

Expand Down
1 change: 1 addition & 0 deletions apps/server/src/orchestration/decider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2629,6 +2629,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"
thread.session.activeTurnId !== null ||
command.session.pendingTurnRequestId !== command.requestId ||
command.session.pendingTurnMessageId !== command.messageId ||
command.session.status !== "starting" ||
command.session.providerInstanceId !== command.modelSelection.instanceId ||
command.session.runtimeMode !== command.runtimeMode ||
command.session.sessionIncarnationId === undefined ||
Expand Down
Loading