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
68 changes: 68 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2464,6 +2464,74 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => {
);
});

it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test-"))(
"OrchestrationProjectionPipeline pending turn cleanup",
(it) => {
it.effect("clears pending turn starts when startup reaches a terminal session state", () =>
Effect.gen(function* () {
const projectionPipeline = yield* OrchestrationProjectionPipeline;
const eventStore = yield* OrchestrationEventStore;
const sql = yield* SqlClient.SqlClient;

for (const [index, status] of (["error", "interrupted", "stopped"] as const).entries()) {
const threadId = ThreadId.make(`thread-terminal-${status}`);
const requestedAt = `2026-02-26T14:00:0${index}.000Z`;
yield* eventStore.append({
type: "thread.turn-start-requested",
eventId: EventId.make(`evt-terminal-pending-${status}`),
aggregateKind: "thread",
aggregateId: threadId,
occurredAt: requestedAt,
commandId: CommandId.make(`cmd-terminal-pending-${status}`),
causationEventId: null,
correlationId: CorrelationId.make(`cmd-terminal-pending-${status}`),
metadata: {},
payload: {
threadId,
messageId: MessageId.make(`message-terminal-${status}`),
runtimeMode: "approval-required",
createdAt: requestedAt,
},
});
yield* eventStore.append({
type: "thread.session-set",
eventId: EventId.make(`evt-terminal-session-${status}`),
aggregateKind: "thread",
aggregateId: threadId,
occurredAt: requestedAt,
commandId: CommandId.make(`cmd-terminal-session-${status}`),
causationEventId: null,
correlationId: CorrelationId.make(`cmd-terminal-session-${status}`),
metadata: {},
payload: {
threadId,
session: {
threadId,
status,
providerName: "codex",
runtimeMode: "approval-required",
activeTurnId: null,
lastError: status === "error" ? "startup failed" : null,
updatedAt: requestedAt,
},
},
});
}

yield* projectionPipeline.bootstrap;

const pendingRows = yield* sql<{ readonly threadId: string }>`
SELECT thread_id AS "threadId"
FROM projection_turns
WHERE turn_id IS NULL
AND state = 'pending'
`;
assert.deepEqual(pendingRows, []);
}),
);
},
);

it.effect("restores pending turn-start metadata across projection pipeline restart", () =>
Effect.gen(function* () {
const { dbPath } = yield* ServerConfig;
Expand Down
9 changes: 9 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1059,6 +1059,15 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti
case "thread.session-set": {
const turnId = event.payload.session.activeTurnId;
if (turnId === null || event.payload.session.status !== "running") {
if (
event.payload.session.status === "error" ||
event.payload.session.status === "stopped" ||
event.payload.session.status === "interrupted"
) {
yield* projectionTurnRepository.deletePendingTurnStartByThreadId({
threadId: event.payload.threadId,
});
}
// Leaving the "running" session status is the turn-end signal:
// settle still-running turns so their duration reflects the whole
// turn rather than the last assistant message.
Expand Down
125 changes: 123 additions & 2 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import {
TurnId,
} from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Deferred from "effect/Deferred";
import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
import * as ManagedRuntime from "effect/ManagedRuntime";
Expand Down Expand Up @@ -145,6 +146,9 @@ describe("ProviderCommandReactor", () => {
readonly threadModelSelection?: ModelSelection;
readonly sessionModelSwitch?: "unsupported" | "in-session";
readonly requiresNewThreadForModelChange?: boolean;
readonly startSessionEffect?: (
session: ProviderSession,
) => Effect.Effect<ProviderSession, ProviderAdapterRequestError>;
}) {
const now = "2026-01-01T00:00:00.000Z";
const baseDir =
Expand All @@ -159,6 +163,7 @@ describe("ProviderCommandReactor", () => {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5-codex",
};
const startSessionEffect = input?.startSessionEffect;
const startSession = vi.fn((_: unknown, input: unknown) => {
const sessionIndex = nextSessionIndex++;
const resumeCursor =
Expand Down Expand Up @@ -212,8 +217,13 @@ describe("ProviderCommandReactor", () => {
createdAt: now,
updatedAt: now,
};
runtimeSessions.push(session);
return Effect.succeed(session);
return (startSessionEffect?.(session) ?? Effect.succeed(session)).pipe(
Effect.tap((startedSession) =>
Effect.sync(() => {
runtimeSessions.push(startedSession);
}),
),
);
});
const sendTurn = vi.fn((_: unknown) =>
Effect.succeed({
Expand Down Expand Up @@ -463,9 +473,120 @@ describe("ProviderCommandReactor", () => {
const readModel = await harness.readModel();
const thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(thread?.session?.threadId).toBe("thread-1");
expect(thread?.session?.status).toBe("starting");
expect(thread?.session?.runtimeMode).toBe("approval-required");
});

effectIt.effect("projects starting before a slow provider session finishes", () =>
Effect.gen(function* () {
const releaseStart = yield* Deferred.make<void>();
const harness = yield* Effect.promise(() =>
createHarness({
startSessionEffect: (session) => Deferred.await(releaseStart).pipe(Effect.as(session)),
}),
);
const now = "2026-01-01T00:00:00.000Z";

yield* harness.engine.dispatch({
type: "thread.turn.start",
commandId: CommandId.make("cmd-turn-start-slow-provider"),
threadId: ThreadId.make("thread-1"),
message: {
messageId: asMessageId("user-message-slow-provider"),
role: "user",
text: "start slowly",
attachments: [],
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
createdAt: now,
});

yield* Effect.promise(() => waitFor(() => harness.startSession.mock.calls.length === 1));
const duringStartup = yield* Effect.promise(() => harness.readModel());
expect(
duringStartup.threads.find((entry) => entry.id === ThreadId.make("thread-1"))?.session
?.status,
).toBe("starting");
expect(harness.sendTurn).not.toHaveBeenCalled();

yield* Deferred.succeed(releaseStart, undefined);
yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1));
}),
);

effectIt.effect("settles a failed provider startup and allows a clean retry", () =>
Effect.gen(function* () {
let failStartup = true;
const harness = yield* Effect.promise(() =>
createHarness({
startSessionEffect: (session) =>
failStartup
? Effect.fail(
new ProviderAdapterRequestError({
provider: "codex",
method: "thread.start",
detail: "deterministic startup failure",
}),
)
: Effect.succeed(session),
}),
);
const now = "2026-01-01T00:00:00.000Z";

yield* harness.engine.dispatch({
type: "thread.turn.start",
commandId: CommandId.make("cmd-turn-start-provider-failure"),
threadId: ThreadId.make("thread-1"),
message: {
messageId: asMessageId("user-message-provider-failure"),
role: "user",
text: "fail once",
attachments: [],
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
createdAt: now,
});

yield* Effect.promise(() =>
waitFor(async () => {
const readModel = await harness.readModel();
return (
readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"))?.session
?.status === "error"
);
}),
);
let readModel = yield* Effect.promise(() => harness.readModel());
let thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(thread?.session?.lastError).toContain("deterministic startup failure");
expect(harness.sendTurn).not.toHaveBeenCalled();

failStartup = false;
yield* harness.engine.dispatch({
type: "thread.turn.start",
commandId: CommandId.make("cmd-turn-start-provider-retry"),
threadId: ThreadId.make("thread-1"),
message: {
messageId: asMessageId("user-message-provider-retry"),
role: "user",
text: "retry",
attachments: [],
},
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
createdAt: "2026-01-01T00:00:01.000Z",
});

yield* Effect.promise(() => waitFor(() => harness.sendTurn.mock.calls.length === 1));
readModel = yield* Effect.promise(() => harness.readModel());
thread = readModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(thread?.session?.status).toBe("starting");
expect(thread?.session?.lastError).toBeNull();
}),
);

it("generates a thread title on the first turn", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
44 changes: 34 additions & 10 deletions apps/server/src/orchestration/Layers/ProviderCommandReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -288,15 +288,20 @@ const make = Effect.gen(function* () {
readonly createdAt: string;
}) {
const thread = yield* resolveThread(input.threadId);
const session = thread?.session;
if (!session) {
if (!thread) {
return;
}
const session = thread.session;
yield* setThreadSession({
threadId: input.threadId,
session: {
...session,
status: session.status === "stopped" ? "stopped" : "ready",
...(session ?? {
threadId: input.threadId,
providerName: null,
providerInstanceId: thread.modelSelection.instanceId,
runtimeMode: thread.runtimeMode,
}),
status: session?.status === "stopped" ? "stopped" : "error",
activeTurnId: null,
lastError: input.detail,
updatedAt: input.createdAt,
Expand Down Expand Up @@ -354,6 +359,7 @@ const make = Effect.gen(function* () {
createdAt: string,
options?: {
readonly modelSelection?: ModelSelection;
readonly pendingTurnStart?: boolean;
},
) {
const thread = yield* resolveThread(threadId);
Expand Down Expand Up @@ -428,6 +434,22 @@ const make = Effect.gen(function* () {
});
}
const preferredProvider: ProviderDriverKind = desiredDriverKind;
if (options?.pendingTurnStart === true && thread.session?.status !== "running") {
yield* setThreadSession({
threadId,
session: {
threadId,
status: "starting",
providerName: activeSession?.provider ?? preferredProvider,
providerInstanceId: activeSession?.providerInstanceId ?? desiredInstanceId,
runtimeMode: desiredRuntimeMode,
activeTurnId: null,
lastError: null,
updatedAt: createdAt,
},
createdAt,
});
}
if (thread.session !== null) {
yield* rejectStartedThreadModelChangeIfRequired({
threadId,
Expand Down Expand Up @@ -498,7 +520,10 @@ const make = Effect.gen(function* () {
threadId,
session: {
threadId,
status: mapProviderSessionStatusToOrchestrationStatus(session.status),
status:
options?.pendingTurnStart === true && session.status === "ready"
? "starting"
: mapProviderSessionStatusToOrchestrationStatus(session.status),
providerName: session.provider,
providerInstanceId: session.providerInstanceId,
runtimeMode: desiredRuntimeMode,
Expand Down Expand Up @@ -597,11 +622,10 @@ const make = Effect.gen(function* () {
new Error(`Thread '${input.threadId}' was not found in read model.`),
);
}
yield* ensureSessionForThread(
input.threadId,
input.createdAt,
input.modelSelection !== undefined ? { modelSelection: input.modelSelection } : {},
);
yield* ensureSessionForThread(input.threadId, input.createdAt, {
...(input.modelSelection !== undefined ? { modelSelection: input.modelSelection } : {}),
pendingTurnStart: true,
});
if (input.modelSelection !== undefined) {
threadModelSelections.set(input.threadId, input.modelSelection);
}
Expand Down
Loading
Loading