diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 6f1819a672b7..9f67cd0898b8 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -8,16 +8,20 @@ import { type ModelSelection, NodeId, type OrchestrationV2Run, + type OrchestrationV2ThreadProjection, ProjectId, ProviderDriverKind, ProviderInstanceId, ProviderThreadId, + RunAttemptId, RunId, + RuntimeRequestId, ThreadId, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import * as Stream from "effect/Stream"; import * as CheckpointStore from "../checkpointing/CheckpointStore.ts"; @@ -38,6 +42,7 @@ import * as Orchestrator from "./Orchestrator.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import { continueRestartedRun } from "./RestartContinuation.ts"; import * as RuntimeLayer from "./runtimeLayer.ts"; +import { makeSubagentChildThread } from "./SubagentProjection.ts"; import * as ProviderTurnStartServiceTestkit from "./ProviderTurnStartService.testkit.ts"; const layerPlatformTest = Layer.merge( @@ -288,6 +293,220 @@ const seedParentWithTerminalTask = (input: { }); }); +// A parent with a running turn and two running delegated tasks: the first +// spawned by that turn, the second by an earlier one, so each belongs to a +// different completion cohort. `block` records a pending request on a child +// and returns the parent sequence to wait from for the notice it queues. +const seedParentWithBlockingTasks = (name: string) => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const sink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make(`thread:${name}`); + const runId = RunId.make(`run:${name}`); + const rootNodeId = NodeId.make(`node:${name}-root`); + const taskId = NodeId.make(`node:${name}-task`); + const secondTaskId = NodeId.make(`node:${name}-second`); + const childThreadId = ThreadId.make(`thread:${name}-child`); + const secondChildThreadId = ThreadId.make(`thread:${name}-second-child`); + yield* seedParentWithTerminalTask({ + threadId, + runId, + projectId: ProjectId.make(`project:${name}`), + rootNodeId, + taskId, + deliveryState: "claimed", + completionWake: "always", + now, + }); + const parent = yield* orchestrator.getThreadProjection(threadId); + const parentRun = parent.runs.find((run) => run.id === runId)!; + const providerThreadId = parentRun.providerThreadId!; + const attemptId = RunAttemptId.make(`attempt:${name}`); + const runningTask = (id: NodeId, child: ThreadId, title: string, taskRunId = runId) => [ + { + id: EventId.make(`event:${name}:child:${id}`), + type: "thread.created" as const, + threadId: child, + driver, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: makeSubagentChildThread({ + parentThread: parent.thread, + childThreadId: child, + parentNodeId: id, + activeProviderThreadId: null, + providerInstanceId: modelSelection.instanceId, + modelSelection, + title, + now, + createdBy: "agent", + creationSource: "server", + }), + }, + { + id: EventId.make(`event:${name}:task:${id}`), + type: "subagent.updated" as const, + threadId, + runId: taskRunId, + nodeId: id, + driver, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { + ...parent.subagents[0]!, + id, + runId: taskRunId, + title, + childThreadId: child, + status: "running" as const, + result: null, + completedAt: null, + completionDelivery: undefined, + }, + }, + ]; + // Two running children, and a parent run that Stop can interrupt. + yield* sink.write({ + commandId: CommandId.make(`command:${name}:seed`), + events: [ + ...runningTask(taskId, childThreadId, "Load test lane"), + // Spawned by an earlier parent turn, so Stop's cohort barrier alone + // would not reach its notices. + ...runningTask( + secondTaskId, + secondChildThreadId, + "Second lane", + RunId.make(`run:${name}-earlier`), + ), + { + id: EventId.make(`event:${name}:root`), + type: "node.updated", + threadId, + runId, + nodeId: rootNodeId, + occurredAt: now, + payload: { + id: rootNodeId, + threadId, + runId, + parentNodeId: null, + rootNodeId, + kind: "root_turn", + status: "running", + countsForRun: true, + providerThreadId, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + id: EventId.make(`event:${name}:attempt`), + type: "run-attempt.updated", + threadId, + runId, + nodeId: rootNodeId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { + id: attemptId, + runId, + attemptOrdinal: 1, + rootNodeId, + providerInstanceId: modelSelection.instanceId, + providerThreadId, + providerTurnId: null, + reason: "initial", + status: "running", + startedAt: now, + completedAt: null, + }, + }, + { + id: EventId.make(`event:${name}:active-attempt`), + type: "run.updated", + threadId, + runId, + nodeId: rootNodeId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { ...parentRun, activeAttemptId: attemptId }, + }, + ], + }); + + // Records a pending request on a child. Returns the parent sequence to + // wait from for the notice that the blocked-child reactor queues. + const block = ( + child: ThreadId, + nodeId: NodeId, + id: string, + kind: "user_input" | "command" | "auth_refresh", + attempt = 0, + ) => + Effect.gen(function* () { + const afterSequence = yield* sink.latestSequence({ threadId }); + yield* sink.write({ + commandId: CommandId.make(`command:${name}:${id}:${attempt}`), + events: [ + { + id: EventId.make(`event:${name}:${id}:${attempt}`), + type: "runtime-request.updated", + threadId: child, + nodeId, + occurredAt: now, + payload: { + id: RuntimeRequestId.make(id), + nodeId, + providerTurnId: null, + nativeRequestRef: null, + kind, + status: "pending", + responseCapability: { type: "message" }, + createdAt: now, + resolvedAt: null, + }, + }, + ], + }); + return afterSequence; + }); + const nextNotice = (afterSequence: number) => + sink.stream({ threadId, afterSequence, eventType: "message.updated" }).pipe( + Stream.filter( + (stored) => + stored.event.type === "message.updated" && + stored.event.payload.notification?.source.kind === "delegated_task", + ), + Stream.runHead, + Effect.map((stored) => + Option.isSome(stored) && stored.value.event.type === "message.updated" + ? stored.value.event.payload + : undefined, + ), + ); + const noticeRunStatus = ( + projection: OrchestrationV2ThreadProjection, + messageId: MessageId | undefined, + ) => projection.runs.find((run) => run.userMessageId === messageId)?.status; + + return { + threadId, + runId, + taskId, + secondTaskId, + childThreadId, + secondChildThreadId, + block, + nextNotice, + noticeRunStatus, + }; + }); + it.layer(layerTest)("delegated completion delivery repairs", (it) => { it.effect("acceptance batches pending siblings without acknowledging their results", () => Effect.gen(function* () { @@ -474,6 +693,202 @@ it.layer(layerTest)("delegated completion delivery repairs", (it) => { }), ); + it.effect("tells the parent once when its child blocks on a request", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const sink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const { + threadId, + runId, + taskId, + secondTaskId, + childThreadId, + secondChildThreadId, + block, + nextNotice, + noticeRunStatus, + } = yield* seedParentWithBlockingTasks("delegated-task-blocked"); + + const question = yield* nextNotice( + yield* block(childThreadId, taskId, "request:question", "user_input"), + ); + assert.deepEqual(question?.notification, { + source: { kind: "delegated_task", taskIds: [taskId] }, + outcome: "updated", + summary: "Load test lane is waiting for an answer to a question", + }); + assert.include(question?.text ?? "", "t3_pending_request_respond"); + assert.include(question?.text ?? "", String(childThreadId)); + + // A repeated update for the same request must not queue a second notice. + // The reactor handles events in order, so the approval notice below + // arrives only after the repeat has been processed. + yield* block(childThreadId, taskId, "request:question", "user_input", 1); + const approval = yield* nextNotice( + yield* block(childThreadId, taskId, "request:approval", "command"), + ); + assert.equal( + approval?.notification?.summary, + "Load test lane is waiting for approval (command)", + ); + assert.include(approval?.text ?? "", "Only the user can resolve this"); + + const signIn = yield* nextNotice( + yield* block(childThreadId, taskId, "request:sign-in", "auth_refresh"), + ); + assert.equal( + signIn?.notification?.summary, + "Load test lane is waiting for the user to sign in again", + ); + assert.include(signIn?.text ?? "", "Only the user can resolve this"); + + // A rejected earlier attempt must not swallow a later update's notice. + yield* sink.commitRejectedCommand({ + commandId: CommandId.make("command:delegated-task-blocked:request:retry"), + threadId, + commandType: "message.dispatch", + rejectedAt: now, + error: "Thread has a pending merge-back transfer.", + }); + const retried = yield* nextNotice( + yield* block(childThreadId, taskId, "request:retry", "user_input"), + ); + assert.equal( + retried?.notification?.summary, + "Load test lane is waiting for an answer to a question", + ); + + const blocked = yield* orchestrator.getThreadProjection(threadId); + assert.equal( + blocked.messages.filter( + (message) => + message.notification?.source.kind === "delegated_task" && + message.notification.source.taskIds.includes(taskId), + ).length, + 4, + ); + // The child is not finished: no result is published to the parent. + assert.equal(blocked.subagents.find((row) => row.id === taskId)?.status, "running"); + assert.isFalse(blocked.contextTransfers.some((row) => row.type === "subagent_result")); + assert.equal(noticeRunStatus(blocked, question?.id), "queued"); + assert.equal(noticeRunStatus(blocked, approval?.id), "queued"); + + // Once the parent drops a task (task_cancel), its queued notices must + // not start a parent turn. + yield* orchestrator.dispatch({ + type: "delegated_task.completion-delivery.dispose", + commandId: CommandId.make("command:delegated-task-blocked:dispose"), + parentThreadId: threadId, + taskId, + }); + const disposed = yield* orchestrator.getThreadProjection(threadId); + assert.equal(noticeRunStatus(disposed, question?.id), "cancelled"); + assert.equal(noticeRunStatus(disposed, approval?.id), "cancelled"); + assert.equal(noticeRunStatus(disposed, signIn?.id), "cancelled"); + assert.equal(noticeRunStatus(disposed, retried?.id), "cancelled"); + + // Stop must not let any queued notice start a new parent turn either. + const second = yield* nextNotice( + yield* block(secondChildThreadId, secondTaskId, "request:second", "user_input"), + ); + assert.equal( + noticeRunStatus(yield* orchestrator.getThreadProjection(threadId), second?.id), + "queued", + ); + yield* orchestrator.dispatch({ + type: "run.interrupt", + commandId: CommandId.make("command:delegated-task-blocked:stop-parent"), + threadId, + runId, + reason: "User stopped the parent.", + holdQueue: true, + }); + assert.equal( + noticeRunStatus(yield* orchestrator.getThreadProjection(threadId), second?.id), + "cancelled", + ); + }), + ); + + it.effect("a plain interrupt drops only its own cohort's blocked-task notices", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const { + threadId, + runId, + taskId, + secondTaskId, + childThreadId, + secondChildThreadId, + block, + nextNotice, + noticeRunStatus, + } = yield* seedParentWithBlockingTasks("delegated-task-blocked-plain-interrupt"); + const noticesFor = ( + projection: OrchestrationV2ThreadProjection, + id: NodeId, + ): ReadonlyArray => + projection.messages + .filter( + (message) => + message.notification?.source.kind === "delegated_task" && + message.notification.source.taskIds.includes(id), + ) + .map((message) => message.id); + + const stopped = yield* nextNotice( + yield* block(childThreadId, taskId, "request:plain-interrupt-own", "user_input"), + ); + const kept = yield* nextNotice( + yield* block(secondChildThreadId, secondTaskId, "request:plain-interrupt-other", "command"), + ); + const before = yield* orchestrator.getThreadProjection(threadId); + assert.equal(noticeRunStatus(before, stopped?.id), "queued"); + assert.equal(noticeRunStatus(before, kept?.id), "queued"); + + // Interrupt without holding the queue: only the interrupted run's own + // cohort is dropped, and the rest of the queue stays. + yield* orchestrator.dispatch({ + type: "run.interrupt", + commandId: CommandId.make("command:delegated-task-blocked-plain-interrupt:interrupt"), + threadId, + runId, + reason: "User interrupted the parent.", + }); + const interrupted = yield* orchestrator.getThreadProjection(threadId); + assert.equal(noticeRunStatus(interrupted, stopped?.id), "cancelled"); + assert.notEqual(noticeRunStatus(interrupted, kept?.id), "cancelled"); + assert.equal( + interrupted.subagents.find((row) => row.id === taskId)?.completionDelivery?.state, + "disposed", + ); + assert.notEqual( + interrupted.subagents.find((row) => row.id === secondTaskId)?.completionDelivery?.state, + "disposed", + ); + + // A later request from the dropped cohort queues nothing, while the + // other cohort still notifies. The reactor handles events in order, so + // the second notice arrives only after the first request was handled. + yield* block(childThreadId, taskId, "request:plain-interrupt-own-later", "user_input"); + const keptLater = yield* nextNotice( + yield* block( + secondChildThreadId, + secondTaskId, + "request:plain-interrupt-other-later", + "user_input", + ), + ); + assert.isDefined(keptLater); + const after = yield* orchestrator.getThreadProjection(threadId); + assert.deepEqual(noticesFor(after, taskId), [stopped!.id]); + assert.deepEqual(noticesFor(after, secondTaskId), [kept!.id, keptLater!.id]); + assert.notEqual(noticeRunStatus(after, kept?.id), "cancelled"); + assert.notEqual(noticeRunStatus(after, keptLater?.id), "cancelled"); + }), + ); + it.effect("builds completion text and metadata from the same live cohort", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index de100d291b8f..5db6f5311a38 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -37,6 +37,7 @@ import { type OrchestrationV2ProviderTurn, type OrchestrationV2Run, type OrchestrationV2RunAttempt, + type OrchestrationV2RuntimeRequest, type OrchestrationV2ThreadShell, type OrchestrationV2ThreadShellSnapshot, type OrchestrationV2StoredEvent, @@ -1943,6 +1944,15 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio occurredAt: now, payload: updatedTask, }); + // The parent has the outcome or dropped the task, so a queued notice + // that it is blocked is stale. + yield* cancelQueuedBlockedTaskNotices({ + command, + events, + projection, + now, + taskIds: [task.id], + }); const parentRun = task.runId === null @@ -2024,6 +2034,39 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } }); + /** + * Cancels queued notices that a delegated task is blocked on a request. + * They are server-owned like completion deliveries, so the same Stop and + * disposal barriers must keep them from starting a parent turn. + */ + const cancelQueuedBlockedTaskNotices = (input: { + readonly command: OrchestrationV2ServerCommand; + readonly events: Ref.Ref>; + readonly projection: Pick< + OrchestrationV2ThreadProjection, + "runs" | "messages" | "nodes" | "attempts" + >; + readonly now: DateTime.Utc; + readonly taskIds?: ReadonlyArray; + }) => + Effect.forEach( + input.projection.runs.filter((run) => { + if (run.status !== "queued") return false; + const message = input.projection.messages.find( + (candidate) => candidate.id === run.userMessageId, + ); + const source = message?.notification?.source; + return ( + message?.delegatedCompletion === undefined && + source?.kind === "delegated_task" && + (input.taskIds === undefined || + source.taskIds.some((taskId) => input.taskIds!.includes(taskId))) + ); + }), + (run) => emitQueuedRunCancellation({ ...input, run }), + { discard: true }, + ); + const disposeDelegatedCompletionCohort = (input: { readonly command: OrchestrationV2ServerCommand; readonly events: Ref.Ref>; @@ -2105,6 +2148,15 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio now: input.now, }); } + if (input.cancelQueuedDelivery !== false) { + yield* cancelQueuedBlockedTaskNotices({ + command: input.command, + events: input.events, + projection: input.projection, + now: input.now, + taskIds: tasks.map((task) => task.id), + }); + } }); const disposeAllDelegatedCompletionCohorts = (input: { @@ -8468,6 +8520,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio now: input.now, }); } + // Stop means nothing new starts, whichever cohort a notice belongs to. + yield* cancelQueuedBlockedTaskNotices({ + command: input.command, + events: input.events, + projection: yield* getProjectionWithPendingEvents(thread.id, input.events), + now: input.now, + }); }); const dispatchBackgroundWorkSettle = ( @@ -10573,6 +10632,98 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ); + /** + * Tells the parent when its app-owned delegated child blocks on a request. + * The child keeps running and its eventual result is still delivered; this + * only queues a notice so the parent can answer a question or ask the user. + * The command id is per request, so replays and repeated updates are no-ops. + */ + const notifyParentOfBlockedChild = ( + childThreadId: ThreadId, + request: OrchestrationV2RuntimeRequest, + ) => + Effect.gen(function* () { + const parentThreadId = yield* appOwnedSubagentParentThreadId(childThreadId); + if (parentThreadId === undefined) return; + // Hold the parent lock from the eligibility check through dispatch so a + // concurrent Stop, disposal, or archive cannot slip in between. + yield* threadDispatch.withLock( + parentThreadId, + dispatchBlockedChildNotice(parentThreadId, childThreadId, request), + ); + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("Failed to notify parent of blocked delegated task", { + childThreadId, + requestId: request.id, + cause, + }), + ), + ); + + const dispatchBlockedChildNotice = ( + parentThreadId: ThreadId, + childThreadId: ThreadId, + request: OrchestrationV2RuntimeRequest, + ) => + Effect.gen(function* () { + const parent = yield* projectionStore.getThreadRecords(parentThreadId, ["subagents"]); + const task = parent.subagents.find( + (candidate) => + candidate.origin === "app_owned" && candidate.childThreadId === childThreadId, + ); + if ( + task === undefined || + isTerminalDelegatedTaskStatus(task.status) || + task.completionDelivery?.state === "disposed" || + parent.thread.archivedAt !== null || + parent.thread.deletedAt !== null + ) { + return; + } + const label = task.title?.trim() || "Delegated task"; + const need = + request.kind === "user_input" + ? "is waiting for an answer to a question" + : request.kind === "auth_refresh" + ? "is waiting for the user to sign in again" + : `is waiting for approval (${request.kind})`; + const action = + request.kind === "user_input" + ? `Read it with t3_pending_request_read (threadId ${childThreadId}, requestId ${request.id}) and answer with t3_pending_request_respond, or ask the user.` + : "Only the user can resolve this. Ask them to open the delegated task's thread."; + // One notice per request: an accepted receipt means it was delivered. + // A rejected one (say, a pending merge-back) must not suppress a later + // update for the same request, so each rejection moves to a new id. + let commandId = CommandId.make(`command:delegated-task-blocked:${request.id}`); + for (let attempt = 1; ; attempt++) { + const receipt = yield* commandReceipts.getByCommandId(commandId); + if (Option.isNone(receipt)) break; + if (receipt.value.status !== "rejected") return; + commandId = CommandId.make(`command:delegated-task-blocked:${request.id}:${attempt}`); + } + const messageId = yield* idAllocator.allocate.message({ + threadId: parentThreadId, + ordinal: (yield* projectionStore.getMessageCount(parentThreadId)) + 1, + }); + yield* dispatchWithReceiptEffect({ + type: "message.dispatch", + commandId, + threadId: parentThreadId, + messageId, + text: `Delegated task ${task.id} ${need}. ${action} It may already be resolved; its result still arrives when it finishes.`, + notification: { + source: { kind: "delegated_task", taskIds: [task.id] }, + outcome: "updated", + summary: `${label} ${need}`, + }, + attachments: [], + dispatchMode: { type: "queue_after_active" }, + createdBy: "agent", + creationSource: "server", + }); + }); + // Historical terminal events are already represented by the projections // below. Replaying the full event table on every server start delays live // queue promotion in proportion to the lifetime size of the database. @@ -10596,8 +10747,22 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio Effect.forkDetach, ); + yield* eventSink + .stream({ afterSequence: terminalEventsAfterSequence, eventType: "runtime-request.updated" }) + .pipe( + Stream.runForEach((stored) => + stored.event.type === "runtime-request.updated" && + stored.event.threadId !== undefined && + stored.event.payload.status === "pending" && + stored.event.payload.kind !== "dynamic_tool_call" + ? notifyParentOfBlockedChild(stored.event.threadId, stored.event.payload) + : Effect.void, + ), + Effect.forkDetach, + ); + // Settles child results and completion deliveries whose runs ended without - // the listener above: before this boot, or in runtime reconciliation, which + // the terminal-run listener above: before this boot, or in runtime reconciliation, which // it skips. Startup runs this after reconciliation and before the effect // worker. Queue recovery instead holds unstarted runs until an explicit // queue.resume command arrives.