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
47 changes: 40 additions & 7 deletions apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";

import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterRegistry.ts";
import {
Expand All @@ -81,7 +82,14 @@ import {

const DEFAULT_WAIT_TIMEOUT_MS = 10 * 60 * 1_000;
const MAX_WAIT_TIMEOUT_MS = 60 * 60 * 1_000;
const TASK_POLL_INTERVAL_MS = 50;
// Events that can make a delegated task terminal: the parent's task record,
// and the child's runs, nested tasks, and pending provider background work.
const TASK_WAKE_EVENTS = [
{ thread: "parent", eventType: "subagent.updated" },
{ thread: "child", eventType: "run.updated" },
{ thread: "child", eventType: "subagent.updated" },
{ thread: "child", eventType: "provider-thread.updated" },
] as const;
const DEFAULT_THREAD_LIST_LIMIT = 50;
const DEFAULT_THREAD_READ_LIMIT = 50;
const DEFAULT_THREAD_RUN_LIMIT = 10;
Expand Down Expand Up @@ -1298,14 +1306,39 @@ const make = Effect.gen(function* () {
return response;
});

// Re-read the task only when an event on the parent or child thread can
// change its status, instead of polling the projections every 50 ms.
const waitForTask = (scope: McpThreadInvocationScope, taskId: NodeId, timeoutMs: number) =>
Effect.gen(function* () {
while (true) {
const result = yield* readTask(scope, taskId, false, true);
if (isTerminalTaskStatus(result.status)) return result;
yield* Effect.sleep(Duration.millis(TASK_POLL_INTERVAL_MS));
}
}).pipe(Effect.timeoutOption(Duration.millis(timeoutMs)));
const streamError = (error: unknown) =>
failure(
"orchestration_error",
`Unable to watch delegated task ${taskId}: ${errorMessage(error)}`,
);
// Sequences are global, so one cursor taken before the first read
// replays anything either thread records after it.
const afterSequence = yield* threadManagement
.getThreadEventSequence(scope.thread.threadId)
.pipe(Effect.mapError(streamError));
const initial = yield* readTask(scope, taskId, false, true);
if (isTerminalTaskStatus(initial.status)) return Option.some(initial);
// One stream per event type, so transcript events never fill a buffer.
return yield* Stream.mergeAll(
TASK_WAKE_EVENTS.map(({ thread, eventType }) =>
threadManagement.streamStoredEventsFrom({
threadId: thread === "parent" ? scope.thread.threadId : initial.childThreadId,
afterSequence,
eventType,
}),
),
{ concurrency: "unbounded" },
).pipe(
Stream.mapError(streamError),
Stream.mapEffect(() => readTask(scope, taskId, false, true)),
Stream.filter((result) => isTerminalTaskStatus(result.status)),
Stream.runHead,
);
}).pipe(Effect.timeoutOption(Duration.millis(timeoutMs)), Effect.map(Option.flatten));

/**
* A scheduled task the caller may change: one whose modes are no broader
Expand Down
54 changes: 50 additions & 4 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2291,6 +2291,12 @@ describe("CodexAdapterV2 post-settle continuation", () => {
(event): event is Extract<ProviderAdapterV2Event, { type: "message.updated" }> =>
event.type === "message.updated" && event.message.role === "assistant",
);
const assistantTurnItems = (events: ReadonlyArray<ProviderAdapterV2Event>) =>
events.flatMap((event) =>
event.type === "turn_item.updated" && event.turnItem.type === "assistant_message"
? [event.turnItem]
: [],
);

it.effect("keeps an asynchronous Codex question actionable after the turn completes", () =>
Effect.scoped(
Expand Down Expand Up @@ -3427,6 +3433,49 @@ describe("CodexAdapterV2 post-settle continuation", () => {
),
);

it.effect("streams text on the turn item and sends the message once, when it completes", () =>
Effect.scoped(
Effect.gen(function* () {
const transcript = finalAnswerTranscript("codex-streamed-message-once", [
{ id: "answer", text: "CODEX_RECOVERY_OK", streamed: true, completionDelayMs: 100 },
]);
const harness = yield* makeCodexReplayHarness(transcript);
const now = yield* DateTime.now;

yield* harness.runtime.startTurn(
makeCodexTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now,
attemptId: RunAttemptId.make("attempt-codex-streamed-message-once"),
text: "Reply with the requested recovery marker.",
}),
);
yield* Effect.yieldNow;
yield* TestClock.adjust("50 millis");
yield* awaitUntil(
() => assistantTurnItems(harness.events).length === 1,
"streamed turn item",
);
assert.deepEqual(
assistantTurnItems(harness.events).map(({ text, streaming }) => ({ text, streaming })),
[{ text: "CODEX_RECOVERY_OK", streaming: true }],
);
assert.deepEqual(assistantMessages(harness.events), []);

yield* TestClock.adjust("50 millis");
yield* awaitUntil(() => harness.terminalEvents().length === 1, "root turn terminal");
assert.deepEqual(
assistantMessages(harness.events).map(({ message }) => ({
text: message.text,
streaming: message.streaming,
})),
[{ text: "CODEX_RECOVERY_OK", streaming: false }],
);
}).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
),
);

it.effect("suppresses a later streamed duplicate final answer", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down Expand Up @@ -3578,10 +3627,7 @@ describe("CodexAdapterV2 post-settle continuation", () => {
yield* TestClock.adjust("50 millis");
yield* Effect.yieldNow;

assert.equal(
new Set(assistantMessages(harness.events).map((event) => event.message.id)).size,
1,
);
assert.equal(new Set(assistantTurnItems(harness.events).map((item) => item.id)).size, 1);

yield* TestClock.adjust("50 millis");
yield* awaitUntil(() => harness.terminalEvents().length === 1, "root turn terminal");
Expand Down
15 changes: 10 additions & 5 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3180,11 +3180,16 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
driver: CODEX_PROVIDER,
node: artifacts.node,
});
yield* emitProviderEvent({
type: "message.updated",
driver: CODEX_PROVIDER,
message: artifacts.message,
});
// Clients render streaming text from the turn item, so the
// message only carries the final text. flushTurn completes every
// buffered item when a turn ends, interrupted or not.
if (update.completed) {
yield* emitProviderEvent({
type: "message.updated",
driver: CODEX_PROVIDER,
message: artifacts.message,
});
}
yield* emitProviderEvent({
type: "turn_item.updated",
driver: CODEX_PROVIDER,
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -316,6 +316,8 @@ export interface OrchestratorV2Shape {
readonly streamStoredEventsFrom: (input?: {
readonly threadId?: ThreadId;
readonly afterSequence?: number;
/** Keep only this type, before the bounded buffer retains anything. */
readonly eventType?: OrchestrationV2DomainEvent["type"];
}) => Stream.Stream<OrchestrationV2StoredEvent, OrchestratorV2Error>;
readonly streamDomainEvents: Stream.Stream<OrchestrationV2DomainEvent, OrchestratorV2Error>;
}
Expand Down
15 changes: 15 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3762,6 +3762,21 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => {
(yield* projectionStore.getThreadProjection(sourceThreadId)).messages,
);
const targetAfterRollback = yield* projectionStore.getThreadProjection(targetThreadId);
// Both shell reads drop the rolled-back run from the source count, while
// the fork keeps the prefix it inherited.
const shellSnapshot = yield* projectionStore.getShellSnapshot();
for (const shell of [
yield* projectionStore.getThreadShell(sourceThreadId),
shellSnapshot.threads.find((thread) => thread.id === sourceThreadId),
]) {
assert.equal(shell?.itemCount, 2);
}
for (const shell of [
yield* projectionStore.getThreadShell(targetThreadId),
shellSnapshot.threads.find((thread) => thread.id === targetThreadId),
]) {
assert.equal(shell?.visibleItemCount, targetAfterRollback.visibleTurnItems.length);
}
const forwardPage = yield* projectionStore.getTimelinePage(targetThreadId, {
view: "activity",
limit: 2,
Expand Down
53 changes: 42 additions & 11 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1550,6 +1550,25 @@ function itemCountThroughRun(input: {
return count;
}

/**
* Threads that some shell row forks from at a run. Only these need run
* ordinals and per-run item counts, to count the inherited prefix, so a thread
* nobody forked from never pays for a scan of its run history.
*/
function shellForkSourceIds(
rows: ReadonlyArray<{ readonly forked_from_run_source_thread_id: string | null }>,
): ReadonlyArray<ThreadId> {
return [
...new Set(
rows.flatMap((row) =>
row.forked_from_run_source_thread_id === null
? []
: [ThreadId.make(row.forked_from_run_source_thread_id)],
),
),
];
}

function visibleItemCountForShell(input: {
readonly threadId: ThreadId;
readonly statesByThreadId: ReadonlyMap<ThreadId, ShellThreadState>;
Expand Down Expand Up @@ -4938,13 +4957,19 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
AND plan.kind = 'proposed_plan'
AND plan.status = 'active'
) AS has_actionable_proposed_plan,
-- Count per run on the covering (thread_id, run_id) index, then
-- look up each run once, instead of one run lookup per item.
(
SELECT COUNT(*)
FROM orchestration_v2_projection_turn_items i
SELECT COALESCE(SUM(per_run.item_count), 0)
FROM (
SELECT i.run_id, COUNT(*) AS item_count
FROM orchestration_v2_projection_turn_items i
WHERE i.thread_id = t.thread_id
GROUP BY i.run_id
) per_run
LEFT JOIN orchestration_v2_projection_runs r
ON r.run_id = i.run_id
WHERE i.thread_id = t.thread_id
AND (i.run_id IS NULL OR r.status <> 'rolled_back')
ON r.run_id = per_run.run_id
WHERE per_run.run_id IS NULL OR r.status <> 'rolled_back'
) AS item_count,
(
SELECT COUNT(*)
Expand Down Expand Up @@ -5426,14 +5451,15 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
}
const threadRows = [...rowsByThreadId.values()];
const threadIds = [...rowsByThreadId.keys()];
const forkSourceIds = shellForkSourceIds(threadRows);
const readForThreadIds = <A>(
read: (ids: ReadonlyArray<ThreadId>) => Effect.Effect<ReadonlyArray<A>, unknown>,
) =>
threadIds.length === 0 ? Effect.succeed([] as ReadonlyArray<A>) : read(threadIds);
ids: ReadonlyArray<ThreadId> = threadIds,
) => (ids.length === 0 ? Effect.succeed([] as ReadonlyArray<A>) : read(ids));
const [runRows, itemCountRows, sequenceRows, providerThreadRows, pendingTurnItemRows] =
yield* Effect.all([
readForThreadIds(selectShellRunRows),
readForThreadIds(selectShellRunItemCounts),
readForThreadIds(selectShellRunRows, forkSourceIds),
readForThreadIds(selectShellRunItemCounts, forkSourceIds),
sql<{ readonly snapshot_sequence: number | null }>`
SELECT MAX(sequence) AS snapshot_sequence
FROM orchestration_events
Expand Down Expand Up @@ -5520,10 +5546,15 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
}

const threadIds = [...rowsByThreadId.keys()];
const forkSourceIds = shellForkSourceIds([...rowsByThreadId.values()]);
const [runRows, itemCountRows, providerThreadRows, pendingTurnItemRows] =
yield* Effect.all([
selectShellRunRows(threadIds),
selectShellRunItemCounts(threadIds),
forkSourceIds.length === 0
? Effect.succeed([] as ReadonlyArray<ShellRunRow>)
: selectShellRunRows(forkSourceIds),
forkSourceIds.length === 0
? Effect.succeed([] as ReadonlyArray<ShellRunItemCountRow>)
: selectShellRunItemCounts(forkSourceIds),
selectShellProviderThreadRows(threadIds),
selectShellPendingTurnItemRows(threadIds),
]);
Expand Down
Loading
Loading