From 0eec4bd634866e9dde8944e0d6746ad868117924 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 29 Aug 2026 19:27:54 -0700 Subject: [PATCH 1/9] feat(mcp): run scheduled tasks immediately --- apps/server/src/mcp/OrchestratorMcpService.ts | 75 +- ...OrchestratorMcpToolkit.integration.test.ts | 95 ++- .../src/mcp/toolkits/orchestrator/handlers.ts | 6 + .../src/mcp/toolkits/orchestrator/tools.ts | 17 + .../src/orchestration-v2/Orchestrator.ts | 111 ++- .../orchestration-v2/ThreadLaunchService.ts | 6 + .../ThreadManagementService.ts | 5 + .../provider/T3OrchestrationInstructions.ts | 1 + .../ScheduledTaskService.test.ts | 489 ++++++++++++ .../scheduledTasks/ScheduledTaskService.ts | 752 +++++++++++++----- .../orchestrator-mcp-server.md | 21 +- docs/user/scheduled-tasks.md | 12 + packages/contracts/src/orchestrationV2.ts | 9 + packages/contracts/src/orchestratorMcp.ts | 29 + packages/shared/src/t3McpToolPresentation.ts | 1 + 15 files changed, 1441 insertions(+), 188 deletions(-) create mode 100644 apps/server/src/scheduledTasks/ScheduledTaskService.test.ts create mode 100644 docs/user/scheduled-tasks.md diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index ea0abc1e9785..8e5ed5dd2009 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -21,6 +21,8 @@ import { type OrchestratorMcpDeleteScheduledTaskResult, type OrchestratorMcpListScheduledTasksResult, type OrchestratorMcpRuntimeMode, + type OrchestratorMcpRunScheduledTaskNowInput, + type OrchestratorMcpRunScheduledTaskNowResult, type OrchestratorMcpScheduledTask, type OrchestratorMcpScheduleTaskInput, type OrchestratorMcpScheduleTaskResult, @@ -71,7 +73,10 @@ import { ThreadManagementService, } from "../orchestration-v2/ThreadManagementService.ts"; import { ProviderRegistry } from "../provider/Services/ProviderRegistry.ts"; -import { ScheduledTaskService } from "../scheduledTasks/ScheduledTaskService.ts"; +import { + type ScheduledTaskManualRunError, + ScheduledTaskService, +} from "../scheduledTasks/ScheduledTaskService.ts"; import type { McpInvocationScope } from "./McpInvocationContext.ts"; const DEFAULT_WAIT_TIMEOUT_MS = 10 * 60 * 1_000; @@ -126,6 +131,10 @@ export interface OrchestratorMcpServiceShape { scope: McpInvocationScope, input: OrchestratorMcpDeleteScheduledTaskInput, ) => Effect.Effect; + readonly runScheduledTaskNow: ( + scope: McpInvocationScope, + input: OrchestratorMcpRunScheduledTaskNowInput, + ) => Effect.Effect; readonly listThreads: ( scope: McpInvocationScope, input: OrchestratorMcpThreadListInput, @@ -181,6 +190,22 @@ function errorMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } +function scheduledTaskManualRunFailure(error: ScheduledTaskManualRunError): OrchestratorMcpFailure { + switch (error._tag) { + case "ScheduledTaskManualRunNotFoundError": + case "ScheduledTaskManualRunScopeError": + return failure("task_not_found", error.message); + case "ScheduledTaskManualRunRuntimeCeilingError": + return failure("runtime_mode_escalation_denied", error.message); + case "ScheduledTaskManualRunInteractionCeilingError": + return failure("interaction_mode_escalation_denied", error.message); + case "ScheduledTaskManualRunConflictError": + return failure("invalid_request", error.message); + case "ScheduledTaskError": + return failure("orchestration_error", error.message); + } +} + /** * Workspace strategy for a scheduled task created/updated over MCP: bound runs * post into the existing thread (the strategy is unused, keep root); unbound @@ -1135,6 +1160,54 @@ const make = Effect.gen(function* () { ); return { scheduledTaskId: existing.id, deleted: true }; }), + runScheduledTaskNow: (scope, input) => + Effect.gen(function* () { + yield* requireCapability(scope); + const parent = yield* loadProjection(scope.threadId); + const commandId = stableCommandId({ + scope, + requestKey: input.clientRequestId, + operation: "run-scheduled-task-now", + }); + const operation = `run-scheduled-task-now:${input.scheduledTaskId}`; + const result = yield* scheduledTasks + .runNowIdempotent({ + id: input.scheduledTaskId, + commandId, + messageId: stableOperationMessageId({ + scope, + requestKey: input.clientRequestId, + operation, + }), + unboundThreadId: stableThreadId({ + scope, + requestKey: `${input.clientRequestId}:${input.scheduledTaskId}`, + index: 0, + }), + projectId: parent.thread.projectId, + policyCeiling: { + callerThreadId: scope.threadId, + runtimeMode: parent.thread.runtimeMode, + interactionMode: parent.thread.interactionMode, + }, + }) + .pipe(Effect.mapError(scheduledTaskManualRunFailure)); + return { + scheduledTaskId: input.scheduledTaskId, + threadId: result.threadId, + messageId: result.messageId, + runId: result.runId, + status: result.status, + replayed: result.replayed, + receipt: { + commandId: result.receipt.commandId, + acceptedAt: DateTime.formatIso(result.receipt.acceptedAt), + resultSequence: result.receipt.resultSequence, + }, + nextRunAt: result.task?.nextRunAt ?? null, + runCount: result.task?.runCount ?? null, + } satisfies OrchestratorMcpRunScheduledTaskNowResult; + }), capabilities: (scope) => Effect.gen(function* () { yield* requireCapability(scope); diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index d9e06a5bacb6..ba5ae8c23cd1 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -25,6 +25,7 @@ import { type ProviderOptionDescriptor, ProviderThreadId, ProviderTurnId, + RunId, type ScheduledTask, ScheduledTaskId, type ScheduledTaskUpsertInput, @@ -71,7 +72,10 @@ import { materializeReplayTranscriptWorkspace, } from "../orchestration-v2/testkit/ReplayTranscriptNdjson.ts"; import { makeProviderRegistryLayer } from "../provider/testUtils/providerRegistryMock.ts"; -import { ScheduledTaskService } from "../scheduledTasks/ScheduledTaskService.ts"; +import { + type ScheduledTaskManualRunResult, + ScheduledTaskService, +} from "../scheduledTasks/ScheduledTaskService.ts"; import * as McpHttpServer from "./McpHttpServer.ts"; import * as McpInvocationContext from "./McpInvocationContext.ts"; import { delegatedTaskRun, hasPendingChildRuns } from "./OrchestratorMcpService.ts"; @@ -590,6 +594,9 @@ describe("orchestrator MCP toolkit", () => { // In-memory ScheduledTaskService stub so the schedule/list/update/ // delete tools can be exercised without SQL/launch wiring. const scheduledStore = yield* Ref.make>([]); + const scheduledManualRuns = yield* Ref.make>( + [], + ); const scheduledTaskStubLayer = Layer.succeed( ScheduledTaskService, ScheduledTaskService.of({ @@ -611,6 +618,41 @@ describe("orchestrator MCP toolkit", () => { all.filter((candidate) => candidate.id !== input.id), ).pipe(Effect.as({ id: input.id })), runNow: () => Effect.die("ScheduledTaskService.runNow is unused in this test"), + runNowIdempotent: (input) => + Effect.gen(function* () { + const replay = (yield* Ref.get(scheduledManualRuns)).find( + (candidate) => candidate.receipt.commandId === input.commandId, + ); + if (replay !== undefined) { + return { ...replay, task: null, replayed: true }; + } + const task = (yield* Ref.get(scheduledStore)).find( + (candidate) => candidate.id === input.id, + ); + if (task === undefined) { + return yield* Effect.die("Scheduled task is missing from the toolkit stub"); + } + const threadId = task.threadId ?? input.unboundThreadId; + const result = { + task: { ...task, runCount: task.runCount + 1 }, + threadId, + messageId: input.messageId, + runId: RunId.make(`run:${input.messageId}`), + status: "queued" as const, + replayed: false, + receipt: { + commandId: input.commandId, + threadId, + commandType: "message.dispatch", + acceptedAt: DateTime.makeUnsafe("2026-07-01T09:00:00.000Z"), + resultSequence: 23, + status: "accepted" as const, + error: null, + }, + } satisfies ScheduledTaskManualRunResult; + yield* Ref.update(scheduledManualRuns, (runs) => [...runs, result]); + return result; + }), }), ); const testLayer = McpHttpServer.OrchestratorToolkitRegistrationLive.pipe( @@ -1262,6 +1304,22 @@ describe("orchestrator MCP toolkit", () => { const scheduleTool = server.tools.find(({ tool }) => tool.name === "schedule_task"); expect(scheduleTool?.tool.annotations?.destructiveHint).toBe(true); + const runScheduledTaskNowTool = server.tools.find( + ({ tool }) => tool.name === "run_scheduled_task_now", + ); + expect(runScheduledTaskNowTool?.tool.inputSchema).toMatchObject({ + type: "object", + properties: { + scheduledTaskId: expect.any(Object), + clientRequestId: expect.any(Object), + }, + required: expect.arrayContaining(["scheduledTaskId", "clientRequestId"]), + }); + expect(runScheduledTaskNowTool?.tool.annotations).toMatchObject({ + destructiveHint: true, + idempotentHint: true, + openWorldHint: true, + }); const scheduleCall = yield* invoke("schedule_task", { prompt: "wake up in this thread and say hello", schedule: { type: "interval", everyMs: 60_000 }, @@ -1308,6 +1366,25 @@ describe("orchestrator MCP toolkit", () => { enabled: false, }); + const scheduledRunNowCall = yield* invoke("run_scheduled_task_now", { + scheduledTaskId, + clientRequestId: "run-scheduled-hello-1", + }); + expect(scheduledRunNowCall.isError).toBe(false); + expect(scheduledRunNowCall.structuredContent).toMatchObject({ + scheduledTaskId, + threadId: parentThreadId, + messageId: expect.stringContaining("run-scheduled-hello-1"), + runId: expect.any(String), + status: "queued", + replayed: false, + receipt: { + commandId: expect.stringContaining("run-scheduled-hello-1"), + resultSequence: 23, + }, + runCount: 1, + }); + // delete_scheduled_task removes it entirely. const scheduledDeleteCall = yield* invoke("delete_scheduled_task", { scheduledTaskId }); expect(scheduledDeleteCall.isError).toBe(false); @@ -1317,6 +1394,22 @@ describe("orchestrator MCP toolkit", () => { }); expect(yield* Ref.get(scheduledStore)).toHaveLength(0); + const scheduledRunNowReplay = yield* invoke("run_scheduled_task_now", { + scheduledTaskId, + clientRequestId: "run-scheduled-hello-1", + }); + expect(scheduledRunNowReplay.isError).toBe(false); + expect(scheduledRunNowReplay.structuredContent).toMatchObject({ + scheduledTaskId, + threadId: parentThreadId, + messageId: scheduledRunNowCall.structuredContent?.messageId, + runId: scheduledRunNowCall.structuredContent?.runId, + replayed: true, + receipt: scheduledRunNowCall.structuredContent?.receipt, + nextRunAt: null, + runCount: null, + }); + // OpenCode 1.15 has emitted this exact nested-object-as-JSON-string // shape. Decode it at the MCP boundary rather than failing a task // the model otherwise specified correctly. diff --git a/apps/server/src/mcp/toolkits/orchestrator/handlers.ts b/apps/server/src/mcp/toolkits/orchestrator/handlers.ts index c72ba2da07a2..f912890ba256 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/handlers.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/handlers.ts @@ -54,6 +54,12 @@ const handlers = { const service = yield* OrchestratorMcpService; return yield* service.deleteScheduledTask(scope, input); }), + run_scheduled_task_now: (input) => + Effect.gen(function* () { + const scope = yield* McpInvocationContext; + const service = yield* OrchestratorMcpService; + return yield* service.runScheduledTaskNow(scope, input); + }), create_threads: (input) => Effect.gen(function* () { const scope = yield* McpInvocationContext; diff --git a/apps/server/src/mcp/toolkits/orchestrator/tools.ts b/apps/server/src/mcp/toolkits/orchestrator/tools.ts index 2a884c6de109..a989e256c494 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/tools.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/tools.ts @@ -11,6 +11,8 @@ import { OrchestratorMcpListScheduledTasksResult, OrchestratorMcpScheduleTaskInput, OrchestratorMcpScheduleTaskResult, + OrchestratorMcpRunScheduledTaskNowInput, + OrchestratorMcpRunScheduledTaskNowResult, OrchestratorMcpTaskCancelInput, OrchestratorMcpTaskCancelResult, OrchestratorMcpUpdateScheduledTaskInput, @@ -143,6 +145,20 @@ export const DeleteScheduledTaskTool = Tool.make("delete_scheduled_task", { .annotate(Tool.Title, "Delete a scheduled task") .annotate(Tool.Destructive, true); +export const RunScheduledTaskNowTool = Tool.make("run_scheduled_task_now", { + description: + "Run one scheduled task immediately through the app scheduler. scheduledTaskId comes from list_scheduled_tasks. clientRequestId is required: reuse the same value only to retry this exact manual run. Success means the scheduled prompt was durably accepted for dispatch, not that the agent turn finished. This does not change the schedule definition or enabled state, but it records a run and advances nextRunAt using the current schedule.", + parameters: OrchestratorMcpRunScheduledTaskNowInput, + success: OrchestratorMcpRunScheduledTaskNowResult, + failure: OrchestratorMcpFailure, + failureMode: "return", + dependencies, +}) + .annotate(Tool.Title, "Run a scheduled task now") + .annotate(Tool.Destructive, true) + .annotate(Tool.Idempotent, true) + .annotate(Tool.OpenWorld, true); + export const CreateThreadsTool = Tool.make("create_threads", { description: "Create one or more ORDINARY TOP-LEVEL T3 conversations. This is not delegation and does not create child agents/subagents. If the user asks for agents, subagents, workers, delegation, or parallel help, call delegate_task once per child insteadβ€”even when selecting different providers. Use create_threads only when the user explicitly asks for separate/new/top-level threads or conversations. Each entry may override provider, model, options, runtime mode, and interaction mode; omitted settings inherit.", @@ -258,6 +274,7 @@ export const OrchestratorToolkit = Toolkit.make( ListScheduledTasksTool, UpdateScheduledTaskTool, DeleteScheduledTaskTool, + RunScheduledTaskNowTool, CreateThreadsTool, ThreadStartTool, ThreadListTool, diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 88fe7598325f..585a23b20844 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -41,7 +41,7 @@ import * as Stream from "effect/Stream"; import { CheckpointServiceV2 } from "./CheckpointService.ts"; import { CommandPolicyV2 } from "./CommandPolicy.ts"; -import { CommandReceiptStoreV2 } from "./CommandReceiptStore.ts"; +import { CommandReceiptStoreV2, type CommandReceiptStoreV2Shape } from "./CommandReceiptStore.ts"; import { ContextHandoffServiceV2 } from "./ContextHandoffService.ts"; import { EventSinkV2 } from "./EventSink.ts"; import type { OrchestrationEffectRequestV2, PendingOrchestrationEffectV2 } from "./EffectOutbox.ts"; @@ -181,6 +181,7 @@ export interface OrchestratorV2Shape { readonly dispatch: ( command: OrchestrationV2Command, ) => Effect.Effect; + readonly getCommandReceipt: CommandReceiptStoreV2Shape["getByCommandId"]; readonly getThreadProjection: ( threadId: ThreadId, ) => Effect.Effect; @@ -241,6 +242,23 @@ function isNativeMaintenanceCommand(message: { ); } +function runtimeModeRank(mode: OrchestrationV2AppThread["runtimeMode"]): number { + switch (mode) { + case "approval-required": + return 0; + case "auto-accept-edits": + return 1; + case "auto": + return 2; + case "full-access": + return 3; + } +} + +function interactionModeRank(mode: OrchestrationV2AppThread["interactionMode"]): number { + return mode === "plan" ? 0 : 1; +} + function commandThreadId(command: OrchestrationV2Command): ThreadId { switch (command.type) { case "thread.create": @@ -639,6 +657,60 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return withId; }); + const enforcePolicyCeiling = Effect.fn("orchestrationV2.enforcePolicyCeiling")(function* (input: { + readonly command: Extract< + OrchestrationV2Command, + { readonly type: "thread.create" | "message.dispatch" } + >; + readonly projectId: OrchestrationV2AppThread["projectId"]; + readonly runtimeMode: OrchestrationV2AppThread["runtimeMode"]; + readonly interactionMode: OrchestrationV2AppThread["interactionMode"]; + }) { + const ceiling = input.command.policyCeiling; + if (ceiling === undefined) return; + const caller = yield* projectionStore.getThreadProjection(ceiling.callerThreadId).pipe( + Effect.mapError( + (cause) => + new OrchestratorProjectionError({ + threadId: ceiling.callerThreadId, + cause, + }), + ), + ); + if ( + caller.thread.deletedAt !== null || + caller.thread.archivedAt !== null || + caller.thread.projectId !== input.projectId + ) { + return yield* new OrchestratorDispatchError({ + commandId: input.command.commandId, + commandType: input.command.type, + cause: `Caller thread ${ceiling.callerThreadId} does not authorize this target project.`, + }); + } + if ( + runtimeModeRank(input.runtimeMode) > runtimeModeRank(ceiling.runtimeMode) || + runtimeModeRank(input.runtimeMode) > runtimeModeRank(caller.thread.runtimeMode) + ) { + return yield* new OrchestratorDispatchError({ + commandId: input.command.commandId, + commandType: input.command.type, + cause: `Target runtime mode ${input.runtimeMode} exceeds the caller ceiling.`, + }); + } + if ( + interactionModeRank(input.interactionMode) > interactionModeRank(ceiling.interactionMode) || + interactionModeRank(input.interactionMode) > + interactionModeRank(caller.thread.interactionMode) + ) { + return yield* new OrchestratorDispatchError({ + commandId: input.command.commandId, + commandType: input.command.type, + cause: `Target interaction mode ${input.interactionMode} exceeds the caller ceiling.`, + }); + } + }); + const getProjectionWithPendingEvents = ( threadId: ThreadId, events: Ref.Ref>, @@ -1349,6 +1421,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio "orchestration_v2.driver": command.modelSelection.instanceId, }); + yield* enforcePolicyCeiling({ + command, + projectId: command.projectId, + runtimeMode: command.runtimeMode, + interactionMode: command.interactionMode, + }); + const now = yield* DateTime.now; const emitEvent = emit(events, command); const thread: OrchestrationV2AppThread = { @@ -2888,6 +2967,12 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } } + yield* enforcePolicyCeiling({ + command, + projectId: projection.thread.projectId, + runtimeMode: projection.thread.runtimeMode, + interactionMode: projection.thread.interactionMode, + }); if (projection.thread.settledOverride !== null) { const now = yield* DateTime.now; const thread: OrchestrationV2AppThread = { @@ -7183,8 +7268,28 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } satisfies OrchestratorV2DispatchResult; }); + const dispatchLockKeys = (command: OrchestrationV2Command): ReadonlyArray => { + const targetThreadId = commandThreadId(command); + const keys = + (command.type === "thread.create" || command.type === "message.dispatch") && + command.policyCeiling !== undefined + ? [targetThreadId, command.policyCeiling.callerThreadId] + : [targetThreadId]; + return [...new Set(keys)].toSorted((left, right) => (left < right ? -1 : left > right ? 1 : 0)); + }; + + const withDispatchLocks = ( + keys: ReadonlyArray, + effect: Effect.Effect, + ): Effect.Effect => { + const [key, ...remaining] = keys; + return key === undefined + ? effect + : threadDispatch.withLock(key, withDispatchLocks(remaining, effect)); + }; + const dispatchWithReceipt = (command: OrchestrationV2Command) => - threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command)); + withDispatchLocks(dispatchLockKeys(command), dispatchWithReceiptEffect(command)); const handleTerminalRun = (stored: OrchestrationV2StoredEvent) => Effect.gen(function* () { @@ -7335,6 +7440,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return OrchestratorV2.of({ resumeQueuedRuns, dispatch: dispatchWithReceipt, + getCommandReceipt: commandReceipts.getByCommandId, getThreadProjection: (threadId) => projectionStore .getThreadProjection(threadId) @@ -7443,6 +7549,7 @@ export const layerUnavailable: Layer.Layer = Layer.succeed( cause: "Orchestration V2 live runtime is not configured.", }), ), + getCommandReceipt: () => Effect.die("Orchestration V2 live runtime is not configured."), getThreadProjection: (threadId) => Effect.fail( new OrchestratorProjectionError({ diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index a7695f527450..d7faf53aee76 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -5,6 +5,7 @@ import { type ModelSelection, type OrchestrationV2Actor, type OrchestrationV2CreationSource, + type OrchestrationV2PolicyCeiling, type OrchestrationV2ThreadProjection, type ProviderInteractionMode, ProjectId, @@ -65,6 +66,7 @@ export interface ThreadLaunchInput { readonly modelSelection: ModelSelection; readonly runtimeMode: RuntimeMode; readonly interactionMode: ProviderInteractionMode; + readonly policyCeiling?: OrchestrationV2PolicyCeiling; readonly workspaceStrategy: ThreadLaunchWorkspaceStrategy; readonly initialMessage?: ThreadLaunchInitialMessage; readonly createdBy: OrchestrationV2Actor; @@ -485,6 +487,9 @@ export const make = Effect.gen(function* () { modelSelection: input.modelSelection, runtimeMode: input.runtimeMode, interactionMode: input.interactionMode, + ...(input.policyCeiling === undefined + ? {} + : { policyCeiling: input.policyCeiling }), branch: initialBranch, worktreePath: initialWorktreePath, createdBy: input.createdBy, @@ -527,6 +532,7 @@ export const make = Effect.gen(function* () { attachments: input.initialMessage.attachments, ...(input.generateTitle === true ? { titleSeed: input.title } : {}), modelSelection: input.modelSelection, + ...(input.policyCeiling === undefined ? {} : { policyCeiling: input.policyCeiling }), dispatchMode: { type: "defer_start" }, createdBy: input.createdBy, creationSource: input.creationSource, diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.ts b/apps/server/src/orchestration-v2/ThreadManagementService.ts index 258c553b705e..41685058249e 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.ts @@ -7,6 +7,7 @@ import { type OrchestrationV2Command, type OrchestrationV2ConversationMessage, type OrchestrationV2CreationSource, + type OrchestrationV2PolicyCeiling, type OrchestrationV2Run, type OrchestrationV2ThreadShellSnapshot, type OrchestrationV2ThreadProjection, @@ -105,6 +106,7 @@ export interface ThreadManagementSendInput { readonly text: string; readonly attachments: ReadonlyArray; readonly modelSelection?: ModelSelection; + readonly policyCeiling?: OrchestrationV2PolicyCeiling; readonly mode: ThreadManagementSendMode; readonly createdBy: OrchestrationV2Actor; readonly creationSource: OrchestrationV2CreationSource; @@ -271,6 +273,7 @@ export interface ThreadManagementServiceShape { readonly dispatch: ( command: OrchestrationV2Command, ) => Effect.Effect; + readonly getCommandReceipt: OrchestratorV2["Service"]["getCommandReceipt"]; readonly getThreadProjection: ( threadId: ThreadId, ) => Effect.Effect; @@ -518,6 +521,7 @@ const make = Effect.gen(function* () { text: input.text, attachments: input.attachments, ...(input.modelSelection === undefined ? {} : { modelSelection: input.modelSelection }), + ...(input.policyCeiling === undefined ? {} : { policyCeiling: input.policyCeiling }), dispatchMode, createdBy: input.createdBy, creationSource: input.creationSource, @@ -659,6 +663,7 @@ const make = Effect.gen(function* () { return ThreadManagementService.of({ ensureLegacyTranscript, dispatch, + getCommandReceipt: orchestrator.getCommandReceipt, getThreadProjection, getCheckpointContext, getThreadSnapshot, diff --git a/apps/server/src/provider/T3OrchestrationInstructions.ts b/apps/server/src/provider/T3OrchestrationInstructions.ts index 83816ddf626f..f112619a9e4e 100644 --- a/apps/server/src/provider/T3OrchestrationInstructions.ts +++ b/apps/server/src/provider/T3OrchestrationInstructions.ts @@ -7,6 +7,7 @@ The \`t3-code\` MCP server provides app-owned orchestration. Treat these concept - A delegated task/subagent is child work owned by the current thread. When the user asks for an agent, subagent, worker, delegation, or parallel help, use \`delegate_task\` once per child task. This remains true when targeting a different provider. Use \`orchestrator_capabilities\` to discover provider/model IDs, retain each returned \`taskId\`, and use \`task_status\` or \`task_cancel\` to manage it. The returned \`childThreadId\` is backing storage for the subagent; do not replace delegation with ordinary thread creation. - \`create_threads\` and \`t3_thread_start\` create ordinary top-level T3 conversations. Use them only when the user explicitly asks for separate/new/top-level threads or conversations. Never use them merely because the user said "subagent" or requested parallel delegated work. - \`schedule_task\` creates persistent recurring work in the app scheduler. Pass \`schedule\` as a structured object, never as JSON text: \`{"type":"interval","everyMs":3600000}\` for an interval, or \`{"type":"fixed_time","timeOfDay":"09:00","weekdays":[1,2,3,4,5]}\` for a wall-clock schedule. By default runs return to the current thread; set \`bindToCurrentThread=false\` only when the user wants a fresh thread for every run. After scheduling, report the returned cadence and next run time. +- \`run_scheduled_task_now\` asks the scheduler to dispatch one existing task immediately without changing its definition or enabled state. Supply \`scheduledTaskId\` and a stable \`clientRequestId\`; reuse that key only for an exact retry. An accepted result identifies the target run but does not mean the agent turn has finished. Tool names may include an MCP prefix (for example \`mcp__t3-code__delegate_task\`); the semantics are the same. Keep polling/wait loops bounded, do not duplicate active work, and use stable \`clientRequestId\` values when retrying mutations. `; diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts new file mode 100644 index 000000000000..56450bc90572 --- /dev/null +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -0,0 +1,489 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, describe, expect, it } from "@effect/vitest"; +import { + CommandId, + MessageId, + ProjectId, + ProviderDriverKind, + ProviderInstanceId, + type ScheduledTaskUpsertInput, + ScheduledTaskId, + ThreadId, +} from "@t3tools/contracts"; +import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Stream from "effect/Stream"; +import * as TestClock from "effect/testing/TestClock"; + +import * as GitWorkflow from "../git/GitWorkflowService.ts"; +import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; +import * as ProjectService from "../project/ProjectService.ts"; +import * as ProjectSetupScriptRunner from "../project/ProjectSetupScriptRunner.ts"; +import { makeProviderRegistryLayer } from "../provider/testUtils/providerRegistryMock.ts"; +import * as ServerSettings from "../serverSettings.ts"; +import * as TextGeneration from "../textGeneration/TextGeneration.ts"; +import { CodexProviderCapabilitiesV2 } from "../orchestration-v2/Adapters/CodexAdapterV2.ts"; +import * as CommandReceiptStore from "../orchestration-v2/CommandReceiptStore.ts"; +import * as IdAllocator from "../orchestration-v2/IdAllocator.ts"; +import { OrchestratorV2 } from "../orchestration-v2/Orchestrator.ts"; +import type { ProviderAdapterV2Shape } from "../orchestration-v2/ProviderAdapter.ts"; +import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterRegistry.ts"; +import * as ThreadLaunch from "../orchestration-v2/ThreadLaunchService.ts"; +import * as ThreadManagement from "../orchestration-v2/ThreadManagementService.ts"; +import { makeOrchestratorV2ReplayLayerWithRegistry } from "../orchestration-v2/testkit/ProviderReplayHarness.ts"; +import { + ScheduledTaskManualRunConflictError, + ScheduledTaskManualRunRuntimeCeilingError, + ScheduledTaskService, + layer as scheduledTaskLayer, +} from "./ScheduledTaskService.ts"; + +const projectId = ProjectId.make("project:scheduled-run-now"); +const otherProjectId = ProjectId.make("project:scheduled-run-now-other"); +const callerThreadId = ThreadId.make("thread:scheduled-run-now-caller"); +const boundThreadId = ThreadId.make("thread:scheduled-run-now-bound"); +const modelSelection = { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5.6-codex", +} as const; +const project = { + id: projectId, + title: "Scheduled run project", + workspaceRoot: "/repo", + repositoryIdentity: null, + faviconPath: null, + defaultModelSelection: modelSelection, + defaultThreadEnvMode: null, + scripts: [], + createdAt: "2026-08-29T00:00:00.000Z", + updatedAt: "2026-08-29T00:00:00.000Z", + deletedAt: null, +} as const; + +const adapter = { + instanceId: modelSelection.instanceId, + driver: ProviderDriverKind.make("codex"), + getCapabilities: () => Effect.succeed(CodexProviderCapabilitiesV2), + planSelectionTransition: () => Effect.succeed({ type: "apply_on_next_turn" as const }), + openSession: () => Effect.die("provider execution is disabled in scheduler tests"), +} as ProviderAdapterV2Shape; + +interface HarnessOptions { + readonly getProject?: ProjectService.ProjectService["Service"]["getById"]; + readonly providerAdapter?: ProviderAdapterV2Shape | null; +} + +function makeHarness(options: HarnessOptions = {}) { + const database = SqlitePersistenceMemory; + const adapterRegistry = ProviderAdapterRegistry.makeLayer( + options.providerAdapter === null ? [] : [options.providerAdapter ?? adapter], + ); + const orchestrator = makeOrchestratorV2ReplayLayerWithRegistry( + { name: "scheduled-run-now" }, + adapterRegistry, + { databaseLayer: database, runEffectWorker: false }, + ); + const threads = ThreadManagement.layer.pipe(Layer.provide(orchestrator)); + const receipts = CommandReceiptStore.layer.pipe(Layer.provide(database)); + const externalServices = Layer.mergeAll( + Layer.succeed( + ProjectService.ProjectService, + ProjectService.ProjectService.of({ + create: () => Effect.die("unused"), + bootstrap: () => Effect.die("unused"), + update: () => Effect.die("unused"), + delete: () => Effect.die("unused"), + getById: + options.getProject ?? + ((id) => Effect.succeed(id === projectId ? Option.some(project) : Option.none())), + getByWorkspaceRoot: () => Effect.succeed(Option.some(project)), + snapshot: Effect.die("unused"), + }), + ), + Layer.mock(GitWorkflow.GitWorkflowService)({ + createWorktree: () => Effect.die("unused"), + renameBranch: () => Effect.die("unused"), + fetchRemote: () => Effect.die("unused"), + removeWorktree: () => Effect.die("unused"), + resolveRemoteTrackingCommit: () => Effect.die("unused"), + }), + Layer.succeed(ProjectSetupScriptRunner.ProjectSetupScriptRunner, { + runForThread: () => Effect.succeed({ status: "no-script" as const }), + }), + Layer.mock(TextGeneration.TextGeneration)({ + generateThreadTitle: () => Effect.die("unused"), + generateBranchName: () => Effect.die("unused"), + }), + ServerSettings.layerTest(), + makeProviderRegistryLayer(), + ); + const launches = ThreadLaunch.layer.pipe( + Layer.provide(Layer.mergeAll(externalServices, threads, receipts, IdAllocator.layer)), + ); + const scheduler = scheduledTaskLayer.pipe( + Layer.provide(Layer.mergeAll(database, launches, threads)), + ); + return Layer.mergeAll(database, orchestrator, threads, launches, scheduler).pipe( + Layer.provideMerge(NodeServices.layer), + ); +} + +function createThread(input: { + readonly commandId: string; + readonly threadId: ThreadId; + readonly project?: ProjectId; + readonly runtimeMode?: "approval-required" | "auto-accept-edits" | "auto" | "full-access"; + readonly interactionMode?: "plan" | "default"; +}) { + return Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make(input.commandId), + threadId: input.threadId, + projectId: input.project ?? projectId, + title: "Scheduled run test thread", + modelSelection, + runtimeMode: input.runtimeMode ?? "full-access", + interactionMode: input.interactionMode ?? "default", + branch: null, + worktreePath: "/repo", + createdBy: "user", + creationSource: "web", + }); + }); +} + +function taskInput(input: { + readonly id: ScheduledTaskId; + readonly threadId: ThreadId | null; + readonly project?: ProjectId; + readonly runtimeMode?: "approval-required" | "auto-accept-edits" | "auto" | "full-access"; + readonly interactionMode?: "plan" | "default"; +}): ScheduledTaskUpsertInput { + return { + id: input.id, + title: `Task ${input.id}`, + prompt: "Run the scheduled maintenance check.", + enabled: true, + schedule: { type: "interval", everyMs: 60_000 }, + projectId: input.project ?? projectId, + threadId: input.threadId, + workspaceStrategy: { type: "root" }, + modelSelection, + runtimeMode: input.runtimeMode ?? "full-access", + interactionMode: input.interactionMode ?? "default", + createdBy: "agent", + creationSource: "mcp", + }; +} + +function manualRunInput(input: { + readonly taskId: ScheduledTaskId; + readonly key: string; + readonly unboundThreadId?: ThreadId; + readonly project?: ProjectId; + readonly callerRuntimeMode?: "approval-required" | "auto-accept-edits" | "auto" | "full-access"; + readonly callerInteractionMode?: "plan" | "default"; +}) { + return { + id: input.taskId, + commandId: CommandId.make(`command:scheduled-run-now:${input.key}`), + messageId: MessageId.make(`message:scheduled-run-now:${input.taskId}:${input.key}`), + unboundThreadId: + input.unboundThreadId ?? + ThreadId.make(`thread:scheduled-run-now:${input.taskId}:${input.key}`), + projectId: input.project ?? projectId, + policyCeiling: { + callerThreadId, + runtimeMode: input.callerRuntimeMode ?? "full-access", + interactionMode: input.callerInteractionMode ?? "default", + }, + } as const; +} + +describe("ScheduledTaskService.runNowIdempotent", () => { + it.layer(makeHarness(), { timeout: "30 seconds" })( + "real bound and unbound launches return durable accepted runs", + (it) => { + it.effect("dispatches through ThreadManagement and ThreadLaunch", () => + Effect.gen(function* () { + const scheduler = yield* ScheduledTaskService; + const threads = yield* ThreadManagement.ThreadManagementService; + yield* createThread({ commandId: "command:caller:create", threadId: callerThreadId }); + yield* createThread({ commandId: "command:bound:create", threadId: boundThreadId }); + + const boundTaskId = ScheduledTaskId.make("scheduled-task:bound-real"); + yield* scheduler.upsert(taskInput({ id: boundTaskId, threadId: boundThreadId })); + const bound = yield* scheduler.runNowIdempotent( + manualRunInput({ taskId: boundTaskId, key: "bound-real" }), + ); + expect(bound).toMatchObject({ + task: { + id: boundTaskId, + enabled: true, + runCount: 1, + lastRunStatus: "succeeded", + nextRunAt: expect.any(String), + }, + threadId: boundThreadId, + status: "starting", + replayed: false, + receipt: { commandType: "message.dispatch", status: "accepted" }, + }); + const boundProjection = yield* threads.getThreadProjection(boundThreadId); + expect( + boundProjection.messages.find((message) => message.id === bound.messageId)?.text, + ).toBe( + "[Triggered by schedule task: Task scheduled-task:bound-real]\n\nRun the scheduled maintenance check.", + ); + + const unboundTaskId = ScheduledTaskId.make("scheduled-task:unbound-real"); + const newThreadId = ThreadId.make("thread:scheduled-run-now:unbound-real"); + yield* scheduler.upsert(taskInput({ id: unboundTaskId, threadId: null })); + const unbound = yield* scheduler.runNowIdempotent( + manualRunInput({ + taskId: unboundTaskId, + key: "unbound-real", + unboundThreadId: newThreadId, + }), + ); + expect(unbound).toMatchObject({ + task: { id: unboundTaskId, runCount: 1, lastRunStatus: "succeeded" }, + threadId: newThreadId, + status: "preparing", + replayed: false, + receipt: { commandType: "message.dispatch", status: "accepted" }, + }); + const unboundProjection = yield* threads.getThreadProjection(newThreadId); + expect(unboundProjection.messages).toHaveLength(1); + expect(unboundProjection.runs.find((run) => run.id === unbound.runId)?.status).toBe( + "preparing", + ); + }), + ); + }, + ); + + it.layer(makeHarness(), { timeout: "30 seconds" })("accepted retry behavior", (it) => { + it.effect("replays after task deletion and rejects the same key for another task", () => + Effect.gen(function* () { + const scheduler = yield* ScheduledTaskService; + const threads = yield* ThreadManagement.ThreadManagementService; + yield* createThread({ commandId: "command:caller:retry", threadId: callerThreadId }); + yield* createThread({ commandId: "command:bound:retry", threadId: boundThreadId }); + + const firstTaskId = ScheduledTaskId.make("scheduled-task:replay-first"); + const firstInput = manualRunInput({ taskId: firstTaskId, key: "accepted-replay" }); + yield* scheduler.upsert(taskInput({ id: firstTaskId, threadId: boundThreadId })); + const accepted = yield* scheduler.runNowIdempotent(firstInput); + yield* scheduler.delete({ id: firstTaskId }); + const replayed = yield* scheduler.runNowIdempotent(firstInput); + expect(replayed).toEqual({ ...accepted, task: null, replayed: true }); + expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); + + const otherTaskId = ScheduledTaskId.make("scheduled-task:replay-other"); + yield* scheduler.upsert(taskInput({ id: otherTaskId, threadId: boundThreadId })); + const conflict = yield* Effect.flip( + scheduler.runNowIdempotent({ + ...manualRunInput({ taskId: otherTaskId, key: "accepted-replay" }), + commandId: firstInput.commandId, + }), + ); + expect(conflict).toBeInstanceOf(ScheduledTaskManualRunConflictError); + expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); + }), + ); + }); + + it.layer(makeHarness({ providerAdapter: null }), { timeout: "30 seconds" })( + "partial launch failure behavior", + (it) => { + it.effect("does not duplicate an unbound thread after provider dispatch rejection", () => + Effect.gen(function* () { + const scheduler = yield* ScheduledTaskService; + const orchestrator = yield* OrchestratorV2; + yield* createThread({ commandId: "command:caller:partial", threadId: callerThreadId }); + const taskId = ScheduledTaskId.make("scheduled-task:partial-launch"); + const targetThreadId = ThreadId.make("thread:scheduled-run-now:partial-launch"); + const input = manualRunInput({ + taskId, + key: "partial-launch", + unboundThreadId: targetThreadId, + }); + yield* scheduler.upsert(taskInput({ id: taskId, threadId: null })); + + const firstFailure = yield* Effect.flip(scheduler.runNowIdempotent(input)); + expect(firstFailure._tag).toBe("ScheduledTaskError"); + expect(yield* orchestrator.getThreadShell(targetThreadId)).toMatchObject({ + id: targetThreadId, + projectId, + }); + expect((yield* scheduler.list()).tasks.find((task) => task.id === taskId)).toMatchObject({ + runCount: 1, + lastRunStatus: "failed", + }); + + const retryFailure = yield* Effect.flip(scheduler.runNowIdempotent(input)); + expect(retryFailure._tag).toBe("ScheduledTaskError"); + const shells = (yield* orchestrator.getShellSnapshot()).threads; + expect(shells).toHaveLength(2); + expect(shells).toEqual( + expect.arrayContaining([ + expect.objectContaining({ id: targetThreadId }), + expect.objectContaining({ id: callerThreadId }), + ]), + ); + expect((yield* scheduler.list()).tasks.find((task) => task.id === taskId)).toMatchObject({ + runCount: 1, + lastRunStatus: "failed", + }); + }), + ); + }, + ); + + it.effect("serializes active admission and rejects a caller downgrade before V2 acceptance", () => + Effect.gen(function* () { + const lookupEntered = yield* Deferred.make(); + const allowLookup = yield* Deferred.make(); + let lookupCount = 0; + const harness = makeHarness({ + getProject: (id) => { + lookupCount += 1; + return lookupCount === 1 + ? Deferred.succeed(lookupEntered, undefined).pipe( + Effect.andThen(Deferred.await(allowLookup)), + Effect.as(id === projectId ? Option.some(project) : Option.none()), + ) + : Effect.succeed(id === projectId ? Option.some(project) : Option.none()); + }, + }); + + yield* Effect.gen(function* () { + const scheduler = yield* ScheduledTaskService; + const orchestrator = yield* OrchestratorV2; + yield* createThread({ commandId: "command:caller:race", threadId: callerThreadId }); + const taskId = ScheduledTaskId.make("scheduled-task:admission-race"); + const targetThreadId = ThreadId.make("thread:scheduled-run-now:admission-race"); + yield* scheduler.upsert(taskInput({ id: taskId, threadId: null })); + const firstInput = manualRunInput({ + taskId, + key: "admission-race-first", + unboundThreadId: targetThreadId, + }); + const firstFiber = yield* scheduler.runNowIdempotent(firstInput).pipe(Effect.forkChild); + yield* Deferred.await(lookupEntered); + + const overlap = yield* Effect.flip( + scheduler.runNowIdempotent(manualRunInput({ taskId, key: "admission-race-overlap" })), + ); + expect(overlap).toBeInstanceOf(ScheduledTaskManualRunConflictError); + + yield* orchestrator.dispatch({ + type: "thread.runtime-mode.set", + commandId: CommandId.make("command:caller:race:downgrade"), + threadId: callerThreadId, + runtimeMode: "approval-required", + }); + yield* Deferred.succeed(allowLookup, undefined); + const failure = yield* Effect.flip(Fiber.join(firstFiber)); + expect(failure._tag).toBe("ScheduledTaskError"); + expect(yield* orchestrator.getThreadShell(targetThreadId)).toBeNull(); + const stored = (yield* scheduler.list()).tasks.find((task) => task.id === taskId); + expect(stored).toMatchObject({ runCount: 1, lastRunStatus: "failed" }); + }).pipe(Effect.provide(harness), Effect.scoped); + }), + ); + + it.effect("rejects a manual run while the recurring scheduler owns admission", () => + Effect.gen(function* () { + const lookupEntered = yield* Deferred.make(); + const allowLookup = yield* Deferred.make(); + let lookupCount = 0; + const harness = makeHarness({ + getProject: (id) => { + lookupCount += 1; + return lookupCount === 1 + ? Deferred.succeed(lookupEntered, undefined).pipe( + Effect.andThen(Deferred.await(allowLookup)), + Effect.as(id === projectId ? Option.some(project) : Option.none()), + ) + : Effect.succeed(id === projectId ? Option.some(project) : Option.none()); + }, + }); + + yield* Effect.gen(function* () { + const scheduler = yield* ScheduledTaskService; + yield* createThread({ commandId: "command:caller:recurring", threadId: callerThreadId }); + const taskId = ScheduledTaskId.make("scheduled-task:recurring-overlap"); + yield* scheduler.upsert(taskInput({ id: taskId, threadId: null })); + const completed = yield* scheduler.subscribeList().pipe( + Stream.filter((snapshot) => + snapshot.tasks.some( + (task) => + task.id === taskId && task.lastRunStatus === "succeeded" && task.runCount === 1, + ), + ), + Stream.runHead, + Effect.forkChild, + ); + + yield* TestClock.adjust(Duration.seconds(65)); + yield* Deferred.await(lookupEntered); + const overlap = yield* Effect.flip( + scheduler.runNowIdempotent(manualRunInput({ taskId, key: "recurring-overlap-manual" })), + ); + expect(overlap).toBeInstanceOf(ScheduledTaskManualRunConflictError); + yield* Deferred.succeed(allowLookup, undefined); + const snapshot = yield* Fiber.join(completed); + expect(Option.getOrThrow(snapshot)?.tasks.find((task) => task.id === taskId)).toMatchObject( + { + lastRunStatus: "succeeded", + runCount: 1, + }, + ); + }).pipe(Effect.provide(harness), Effect.scoped); + }), + ); + + it.layer(makeHarness(), { timeout: "30 seconds" })("fresh scope and ceiling checks", (it) => { + it.effect("rejects a higher-mode bound target and a task from another project", () => + Effect.gen(function* () { + const scheduler = yield* ScheduledTaskService; + yield* createThread({ + commandId: "command:caller:limited", + threadId: callerThreadId, + runtimeMode: "approval-required", + }); + yield* createThread({ commandId: "command:bound:higher", threadId: boundThreadId }); + + const highTaskId = ScheduledTaskId.make("scheduled-task:higher-mode"); + yield* scheduler.upsert(taskInput({ id: highTaskId, threadId: boundThreadId })); + const ceilingFailure = yield* Effect.flip( + scheduler.runNowIdempotent( + manualRunInput({ + taskId: highTaskId, + key: "higher-mode", + callerRuntimeMode: "approval-required", + }), + ), + ); + expect(ceilingFailure).toBeInstanceOf(ScheduledTaskManualRunRuntimeCeilingError); + + const otherTaskId = ScheduledTaskId.make("scheduled-task:other-project"); + yield* scheduler.upsert( + taskInput({ id: otherTaskId, threadId: null, project: otherProjectId }), + ); + const scopeFailure = yield* Effect.flip( + scheduler.runNowIdempotent(manualRunInput({ taskId: otherTaskId, key: "other-project" })), + ); + assert.equal(scopeFailure._tag, "ScheduledTaskManualRunScopeError"); + }), + ); + }); +}); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index cb59c96afe78..b1f2a8379323 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -1,6 +1,11 @@ import { CommandId, MessageId, + type OrchestrationV2PolicyCeiling, + type OrchestrationV2RunStatus, + type OrchestrationV2ThreadProjection, + ProjectId, + type RunId, ScheduledTask, ScheduledTaskError, ScheduledTaskId, @@ -21,6 +26,7 @@ import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import * as PubSub from "effect/PubSub"; import * as Ref from "effect/Ref"; import * as Result from "effect/Result"; @@ -29,6 +35,8 @@ import * as Stream from "effect/Stream"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; +import type { CommandReceiptV2 } from "../orchestration-v2/CommandReceiptStore.ts"; +import { makeKeyedSerialExecutor } from "../orchestration-v2/KeyedSerialExecutor.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; import { isMissedFixedTimeRun, isSameSchedule, nextScheduledRunAt } from "./Schedule.ts"; @@ -66,6 +74,58 @@ interface ScheduledTaskRow { readonly run_count: number; } +export interface ScheduledTaskManualRunInput { + readonly id: ScheduledTaskId; + readonly commandId: CommandId; + readonly messageId: MessageId; + readonly unboundThreadId: ThreadId; + readonly projectId: ProjectId; + readonly policyCeiling: OrchestrationV2PolicyCeiling; +} + +export interface ScheduledTaskManualRunResult { + readonly task: ScheduledTask | null; + readonly threadId: ThreadId; + readonly messageId: MessageId; + readonly runId: RunId; + readonly status: OrchestrationV2RunStatus; + readonly replayed: boolean; + readonly receipt: CommandReceiptV2; +} + +export class ScheduledTaskManualRunScopeError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunScopeError", + { taskId: ScheduledTaskId, message: Schema.String }, +) {} + +export class ScheduledTaskManualRunRuntimeCeilingError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunRuntimeCeilingError", + { taskId: ScheduledTaskId, message: Schema.String }, +) {} + +export class ScheduledTaskManualRunInteractionCeilingError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunInteractionCeilingError", + { taskId: ScheduledTaskId, message: Schema.String }, +) {} + +export class ScheduledTaskManualRunConflictError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunConflictError", + { taskId: ScheduledTaskId, commandId: CommandId, message: Schema.String }, +) {} + +export class ScheduledTaskManualRunNotFoundError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunNotFoundError", + { taskId: ScheduledTaskId, message: Schema.String }, +) {} + +export type ScheduledTaskManualRunError = + | ScheduledTaskError + | ScheduledTaskManualRunScopeError + | ScheduledTaskManualRunRuntimeCeilingError + | ScheduledTaskManualRunInteractionCeilingError + | ScheduledTaskManualRunConflictError + | ScheduledTaskManualRunNotFoundError; + export class ScheduledTaskService extends Context.Service< ScheduledTaskService, { @@ -85,6 +145,9 @@ export class ScheduledTaskService extends Context.Service< readonly runNow: ( input: ScheduledTaskRunNowInput, ) => Effect.Effect; + readonly runNowIdempotent: ( + input: ScheduledTaskManualRunInput, + ) => Effect.Effect; } >()("t3/scheduledTasks/ScheduledTaskService") {} @@ -121,6 +184,23 @@ function errorMessage(error: unknown): string { return String(error); } +function runtimeModeRank(mode: ScheduledTask["runtimeMode"]): number { + switch (mode) { + case "approval-required": + return 0; + case "auto-accept-edits": + return 1; + case "auto": + return 2; + case "full-access": + return 3; + } +} + +function interactionModeRank(mode: ScheduledTask["interactionMode"]): number { + return mode === "plan" ? 0 : 1; +} + const decodeRow = (row: ScheduledTaskRow) => Effect.gen(function* () { const id = ScheduledTaskId.make(row.task_id); @@ -166,6 +246,7 @@ export const layer = Layer.effect( const threadLaunch = yield* ThreadLaunchService.ThreadLaunchService; const threadManagement = yield* ThreadManagementService.ThreadManagementService; const activeRuns = yield* Ref.make>(new Set()); + const taskMutations = yield* makeKeyedSerialExecutor(); // Sliding(1) coalesces the dirty-signal: every notification triggers a // full list() re-emit anyway, so a slow subscriber only ever needs the // latest signal β€” an unbounded backlog would just grow memory. @@ -269,6 +350,191 @@ export const layer = Layer.effect( return task; }); + const readCommandReceipt = (taskId: ScheduledTaskId, commandId: CommandId) => + threadManagement + .getCommandReceipt(commandId) + .pipe( + Effect.mapError((cause) => + taskError("Could not read the scheduled task run receipt.", { taskId, cause }), + ), + ); + + const authorizeFreshManualRun = Effect.fn("ScheduledTaskService.authorizeFreshManualRun")( + function* (task: ScheduledTask, input: ScheduledTaskManualRunInput) { + const loadScopedThread = (threadId: ThreadId, role: "caller" | "target") => + threadManagement.getProjectThread({ projectId: input.projectId, threadId }).pipe( + Effect.mapError((cause) => + cause._tag === "ThreadManagementThreadNotFoundError" + ? new ScheduledTaskManualRunScopeError({ + taskId: input.id, + message: `The scheduled task run ${role} is not active in the calling project.`, + }) + : taskError(`Could not authorize the scheduled task run ${role}.`, { + taskId: input.id, + cause, + }), + ), + ); + const caller = yield* loadScopedThread(input.policyCeiling.callerThreadId, "caller"); + if (caller.thread.archivedAt !== null) { + return yield* new ScheduledTaskManualRunScopeError({ + taskId: input.id, + message: "The scheduled task run caller is archived.", + }); + } + const target = + task.threadId === null + ? null + : yield* loadScopedThread(ThreadId.make(task.threadId), "target"); + if (target !== null && target.thread.archivedAt !== null) { + return yield* new ScheduledTaskManualRunScopeError({ + taskId: input.id, + message: "The scheduled task target is archived.", + }); + } + const runtimeMode = target?.thread.runtimeMode ?? task.runtimeMode; + const interactionMode = target?.thread.interactionMode ?? task.interactionMode; + if ( + runtimeModeRank(runtimeMode) > runtimeModeRank(input.policyCeiling.runtimeMode) || + runtimeModeRank(runtimeMode) > runtimeModeRank(caller.thread.runtimeMode) + ) { + return yield* new ScheduledTaskManualRunRuntimeCeilingError({ + taskId: input.id, + message: `Target runtime mode ${runtimeMode} exceeds the caller ceiling.`, + }); + } + if ( + interactionModeRank(interactionMode) > + interactionModeRank(input.policyCeiling.interactionMode) || + interactionModeRank(interactionMode) > interactionModeRank(caller.thread.interactionMode) + ) { + return yield* new ScheduledTaskManualRunInteractionCeilingError({ + taskId: input.id, + message: `Target interaction mode ${interactionMode} exceeds the caller ceiling.`, + }); + } + }, + ); + + const findAcceptedManualRun = Effect.fn("ScheduledTaskService.findAcceptedManualRun")( + function* (input: ScheduledTaskManualRunInput) { + const initialMessageCommandId = CommandId.make(`${input.commandId}:initial-message`); + const initialMessageReceipt = yield* readCommandReceipt(input.id, initialMessageCommandId); + let receipt = Option.getOrUndefined(initialMessageReceipt); + if (receipt === undefined) { + const primaryReceipt = yield* readCommandReceipt(input.id, input.commandId); + if (Option.isNone(primaryReceipt)) { + return Option.none(); + } + if (primaryReceipt.value.status === "rejected") { + return yield* taskError("This manual run request was previously rejected.", { + taskId: input.id, + }); + } + if (primaryReceipt.value.commandType === "thread.create") { + if (primaryReceipt.value.threadId !== input.unboundThreadId) { + return yield* new ScheduledTaskManualRunConflictError({ + taskId: input.id, + commandId: input.commandId, + message: "The idempotency key belongs to a different scheduled task run.", + }); + } + // Thread creation committed, but the scheduled prompt did not. Let + // ThreadLaunchService replay the create receipt and finish the same + // task-specific initial message. + return Option.none(); + } + receipt = primaryReceipt.value; + } + if (receipt.status === "rejected") { + return yield* taskError("This manual run request was previously rejected.", { + taskId: input.id, + }); + } + if (receipt.commandType !== "message.dispatch") { + return yield* new ScheduledTaskManualRunConflictError({ + taskId: input.id, + commandId: input.commandId, + message: "The idempotency key belongs to a different command.", + }); + } + + const caller = yield* threadManagement + .getThreadProjection(input.policyCeiling.callerThreadId) + .pipe( + Effect.mapError((cause) => + taskError("Could not authorize the scheduled task run caller.", { + taskId: input.id, + cause, + }), + ), + ); + const target = yield* threadManagement.getThreadProjection(receipt.threadId).pipe( + Effect.mapError((cause) => + taskError("Could not load the accepted scheduled task run.", { + taskId: input.id, + cause, + }), + ), + ); + if ( + caller.thread.deletedAt !== null || + caller.thread.archivedAt !== null || + caller.thread.projectId !== input.projectId || + target.thread.deletedAt !== null || + target.thread.projectId !== input.projectId + ) { + return yield* new ScheduledTaskManualRunScopeError({ + taskId: input.id, + message: "The accepted run is no longer authorized in the calling project.", + }); + } + if ( + runtimeModeRank(target.thread.runtimeMode) > + runtimeModeRank(input.policyCeiling.runtimeMode) || + runtimeModeRank(target.thread.runtimeMode) > runtimeModeRank(caller.thread.runtimeMode) + ) { + return yield* new ScheduledTaskManualRunRuntimeCeilingError({ + taskId: input.id, + message: `Target runtime mode ${target.thread.runtimeMode} exceeds the caller ceiling.`, + }); + } + if ( + interactionModeRank(target.thread.interactionMode) > + interactionModeRank(input.policyCeiling.interactionMode) || + interactionModeRank(target.thread.interactionMode) > + interactionModeRank(caller.thread.interactionMode) + ) { + return yield* new ScheduledTaskManualRunInteractionCeilingError({ + taskId: input.id, + message: `Target interaction mode ${target.thread.interactionMode} exceeds the caller ceiling.`, + }); + } + + const message = target.messages.find((candidate) => candidate.id === input.messageId); + const run = + message?.runId === null || message?.runId === undefined + ? undefined + : target.runs.find((candidate) => candidate.id === message.runId); + if (message === undefined || run === undefined) { + return yield* new ScheduledTaskManualRunConflictError({ + taskId: input.id, + commandId: input.commandId, + message: "The idempotency key was already accepted for a different scheduled task run.", + }); + } + return Option.some({ + task: null, + threadId: target.thread.id, + messageId: message.id, + runId: run.id, + status: run.status, + replayed: true, + receipt, + } satisfies ScheduledTaskManualRunResult); + }, + ); + // Run-state columns (last_run_*, run_count) are intentionally absent from // the conflict clause: they are owned by the run transitions below, and a // concurrent settings save must not overwrite an in-flight increment. @@ -420,9 +686,45 @@ export const layer = Layer.effect( ), ); + const acceptedManualRun = Effect.fn("ScheduledTaskService.acceptedManualRun")(function* ( + input: ScheduledTaskManualRunInput, + projection: OrchestrationV2ThreadProjection, + receiptCommandId: CommandId, + ) { + const receiptOption = yield* readCommandReceipt(input.id, receiptCommandId); + if (Option.isNone(receiptOption) || receiptOption.value.status !== "accepted") { + return yield* taskError( + "The scheduled task dispatch did not produce an accepted receipt.", + { + taskId: input.id, + }, + ); + } + const message = projection.messages.find((candidate) => candidate.id === input.messageId); + const run = + message?.runId === null || message?.runId === undefined + ? undefined + : projection.runs.find((candidate) => candidate.id === message.runId); + if (message === undefined || run === undefined) { + return yield* taskError("The scheduled task dispatch is missing its durable run.", { + taskId: input.id, + }); + } + return { + task: null, + threadId: projection.thread.id, + messageId: message.id, + runId: run.id, + status: run.status, + replayed: false, + receipt: receiptOption.value, + } satisfies ScheduledTaskManualRunResult; + }); + const runTask = Effect.fn("ScheduledTaskService.runTask")(function* ( task: ScheduledTask, trigger: "scheduled" | "manual", + manualRun?: ScheduledTaskManualRunInput, ) { const reserved = yield* Ref.modify(activeRuns, (active) => { if (active.has(task.id)) return [false, active] as const; @@ -432,127 +734,174 @@ export const layer = Layer.effect( }); if (!reserved) { if (trigger === "manual") { + if (manualRun !== undefined) { + return yield* new ScheduledTaskManualRunConflictError({ + taskId: task.id, + commandId: manualRun.commandId, + message: "Schedule task is already running.", + }); + } return yield* taskError("Schedule task is already running.", { taskId: task.id }); } - return task; + return { task, manualRun: null, dispatchError: null }; } - return yield* Effect.gen(function* () { - const startedAt = yield* localNow; - const startedAtIso = iso(startedAt); - - // The in-memory snapshot may be stale: re-read before touching run - // state. The task may have been deleted, paused, or postponed since - // the poll loaded it β€” none of those may fire. - const active = yield* findTask(task.id); - if (active === null) { - // A manual run on a just-deleted task must fail loudly, not report - // a successful run that never dispatched. - if (trigger === "manual") { - return yield* taskError("Schedule task not found.", { taskId: task.id }); - } - return task; - } - if ( - trigger === "scheduled" && - (!active.enabled || - active.nextRunAt === null || - DateTime.toEpochMillis(DateTime.makeUnsafe(active.nextRunAt)) > - DateTime.toEpochMillis(startedAt)) - ) { - return active; - } + return yield* taskMutations + .withLock( + task.id, + Effect.gen(function* () { + const startedAt = yield* localNow; + const startedAtIso = iso(startedAt); - yield* markRunning(active.id, startedAtIso); - yield* notifyChanged; + // The in-memory snapshot may be stale: re-read after admission and + // before touching run state. CRUD uses this same task lock. + const active = yield* findTask(task.id); + if (active === null) { + if (trigger === "manual") { + if (manualRun !== undefined) { + return yield* new ScheduledTaskManualRunNotFoundError({ + taskId: task.id, + message: "Schedule task not found.", + }); + } + return yield* taskError("Schedule task not found.", { taskId: task.id }); + } + return { task, manualRun: null, dispatchError: null }; + } + if (manualRun !== undefined && active.projectId !== manualRun.projectId) { + return yield* new ScheduledTaskManualRunScopeError({ + taskId: active.id, + message: "The scheduled task does not belong to the calling project.", + }); + } + if (manualRun !== undefined) { + yield* authorizeFreshManualRun(active, manualRun); + } + if ( + trigger === "scheduled" && + (!active.enabled || + active.nextRunAt === null || + DateTime.toEpochMillis(DateTime.makeUnsafe(active.nextRunAt)) > + DateTime.toEpochMillis(startedAt)) + ) { + return { task: active, manualRun: null, dispatchError: null }; + } - const fireKey = `${active.id}:${DateTime.toEpochMillis(startedAt)}:${trigger}`; - const commandId = CommandId.make(`scheduled-task:${fireKey}`); - const messageId = MessageId.make(`scheduled-task-message:${fireKey}`); - // Dispatch from the fresh row so prompt/model/binding edits made - // after the poll read are honoured. - const prompt = automationPrompt(active); - - // Effect.exit (not Effect.result) so defects and interruptions in the - // dispatch are also captured and recorded as a failed run instead of - // aborting before markCompleted. - const result = - active.threadId === null - ? yield* Effect.exit( - threadLaunch.launch({ - commandId, - projectId: active.projectId, - title: active.title, - modelSelection: active.modelSelection, - runtimeMode: active.runtimeMode, - interactionMode: active.interactionMode, - workspaceStrategy: active.workspaceStrategy, - initialMessage: { - messageId, - text: prompt, - attachments: [], - }, - createdBy: active.createdBy, - creationSource: active.creationSource, - }), - ) - : yield* Effect.exit( - threadManagement.sendToThread({ - projectId: active.projectId, - commandId, - threadId: ThreadId.make(active.threadId), - messageId, - text: prompt, - attachments: [], - modelSelection: active.modelSelection, - mode: "auto", - createdBy: active.createdBy, - creationSource: active.creationSource, - }), - ); - - const completedAt = yield* localNow; - const runSucceeded = result._tag === "Success"; - const lastRunStatus = runSucceeded ? ("succeeded" as const) : ("failed" as const); - const lastRunError = runSucceeded ? null : errorMessage(result.cause); - // Re-read the task so the next run is computed from the schedule as it - // is *now* (the user may have edited or deleted it while we ran). - const current = yield* findTask(task.id); - const scheduleSource = current ?? task; - const completed: ScheduledTask = { - ...scheduleSource, - updatedAt: iso(completedAt), - lastRunAt: startedAtIso, - nextRunAt: nextRunAt(scheduleSource, completedAt), - lastRunStatus, - lastRunError, - runCount: scheduleSource.runCount + 1, - }; - if (current !== null) { - // startedAtIso in the guard ensures this writes only to the row this - // run marked as running β€” a task deleted mid-run and recreated with - // the same id (idempotent commandId replay) must not be stamped. - yield* markCompleted({ - id: task.id, - completedAtIso: completed.updatedAt, - nextRunAtIso: completed.nextRunAt, - status: lastRunStatus, - error: lastRunError, - startedAtIso, - }); - yield* notifyChanged; - } - return completed; - }).pipe( - Effect.onError((cause) => releaseStuckRun(task, errorMessage(cause))), - Effect.ensuring( - Ref.update(activeRuns, (active) => { - const next = new Set(active); - next.delete(task.id); - return next; - }), - ), - ); + yield* markRunning(active.id, startedAtIso); + yield* notifyChanged; + + const fireKey = `${active.id}:${DateTime.toEpochMillis(startedAt)}:${trigger}`; + const commandId = manualRun?.commandId ?? CommandId.make(`scheduled-task:${fireKey}`); + const messageId = + manualRun?.messageId ?? MessageId.make(`scheduled-task-message:${fireKey}`); + const prompt = automationPrompt(active); + + // Effect.exit captures defects and interruption so bookkeeping is + // terminal even when dispatch is rejected before a turn starts. + const result = + active.threadId === null + ? yield* Effect.exit( + threadLaunch + .launch({ + commandId, + ...(manualRun === undefined ? {} : { threadId: manualRun.unboundThreadId }), + projectId: active.projectId, + title: active.title, + modelSelection: active.modelSelection, + runtimeMode: active.runtimeMode, + interactionMode: active.interactionMode, + ...(manualRun === undefined + ? {} + : { policyCeiling: manualRun.policyCeiling }), + workspaceStrategy: active.workspaceStrategy, + initialMessage: { + messageId, + text: prompt, + attachments: [], + }, + createdBy: active.createdBy, + creationSource: active.creationSource, + }) + .pipe( + Effect.flatMap((launched) => + manualRun === undefined + ? Effect.succeed(null) + : acceptedManualRun( + manualRun, + launched.projection, + CommandId.make(`${commandId}:initial-message`), + ), + ), + ), + ) + : yield* Effect.exit( + threadManagement + .sendToThread({ + projectId: active.projectId, + commandId, + threadId: ThreadId.make(active.threadId), + messageId, + text: prompt, + attachments: [], + modelSelection: active.modelSelection, + ...(manualRun === undefined + ? {} + : { policyCeiling: manualRun.policyCeiling }), + mode: "auto", + createdBy: active.createdBy, + creationSource: active.creationSource, + }) + .pipe( + Effect.flatMap((sent) => + manualRun === undefined + ? Effect.succeed(null) + : acceptedManualRun(manualRun, sent.projection, commandId), + ), + ), + ); + + const completedAt = yield* localNow; + const runSucceeded = result._tag === "Success"; + const lastRunStatus = runSucceeded ? ("succeeded" as const) : ("failed" as const); + const lastRunError = runSucceeded ? null : errorMessage(result.cause); + const current = yield* findTask(task.id); + const scheduleSource = current ?? task; + const completed: ScheduledTask = { + ...scheduleSource, + updatedAt: iso(completedAt), + lastRunAt: startedAtIso, + nextRunAt: nextRunAt(scheduleSource, completedAt), + lastRunStatus, + lastRunError, + runCount: scheduleSource.runCount + 1, + }; + if (current !== null) { + yield* markCompleted({ + id: task.id, + completedAtIso: completed.updatedAt, + nextRunAtIso: completed.nextRunAt, + status: lastRunStatus, + error: lastRunError, + startedAtIso, + }); + yield* notifyChanged; + } + return { + task: completed, + manualRun: runSucceeded ? result.value : null, + dispatchError: runSucceeded ? null : result.cause, + }; + }).pipe(Effect.onError((cause) => releaseStuckRun(task, errorMessage(cause)))), + ) + .pipe( + Effect.ensuring( + Ref.update(activeRuns, (active) => { + const next = new Set(active); + next.delete(task.id); + return next; + }), + ), + ); }); // A due fixed-time run that is long past its slot (server was off or @@ -600,7 +949,8 @@ export const layer = Layer.effect( ? rescheduleMissedRun(task, now) : runTask(task, "scheduled") ).pipe( - Effect.catch((cause) => + Effect.asVoid, + Effect.catchCause((cause) => Effect.logWarning("Scheduled task run failed", { taskId: task.id, cause }), ), ), @@ -690,7 +1040,6 @@ export const layer = Layer.effect( const upsert: ScheduledTaskService["Service"]["upsert"] = (input) => Effect.gen(function* () { - const now = yield* localNow; const uuid = input.commandId === undefined ? yield* crypto.randomUUIDv4.pipe( @@ -704,77 +1053,87 @@ export const layer = Layer.effect( ScheduledTaskId.make( input.commandId ? `scheduled-task:${input.commandId}` : `scheduled-task:${uuid}`, ); - // Look up by the *resolved* id so idempotent creates (commandId replays) - // keep their run history, and so real load failures propagate instead - // of silently resetting an existing row. - const existingTask = yield* findTask(id); - // Keep the existing next_run_at when the schedule itself is untouched: - // editing a title or prompt must not postpone (or resurrect) a due - // run β€” only schedule/enabled changes restart the clock. - const scheduleUnchanged = - existingTask !== null && - existingTask.enabled === input.enabled && - isSameSchedule(existingTask.schedule, input.schedule); - const task: ScheduledTask = { + return yield* taskMutations.withLock( id, - title: input.title, - prompt: input.prompt, - enabled: input.enabled, - schedule: input.schedule, - projectId: input.projectId, - threadId: input.threadId ?? null, - workspaceStrategy: input.workspaceStrategy, - modelSelection: input.modelSelection, - runtimeMode: input.runtimeMode, - interactionMode: input.interactionMode, - createdBy: existingTask?.createdBy ?? input.createdBy ?? "user", - creationSource: input.creationSource ?? "web", - createdAt: existingTask?.createdAt ?? iso(now), - updatedAt: iso(now), - nextRunAt: scheduleUnchanged - ? existingTask.nextRunAt - : nextRunAt({ enabled: input.enabled, schedule: input.schedule }, now), - lastRunAt: existingTask?.lastRunAt ?? null, - lastRunStatus: existingTask?.lastRunStatus ?? "never", - lastRunError: existingTask?.lastRunError ?? null, - runCount: existingTask?.runCount ?? 0, - }; - yield* saveTask(task); - yield* notifyChanged; - return { task }; + Effect.gen(function* () { + const now = yield* localNow; + // Look up by the *resolved* id so idempotent creates (commandId replays) + // keep their run history, and so real load failures propagate instead + // of silently resetting an existing row. + const existingTask = yield* findTask(id); + // Keep the existing next_run_at when the schedule itself is untouched: + // editing a title or prompt must not postpone (or resurrect) a due + // run β€” only schedule/enabled changes restart the clock. + const scheduleUnchanged = + existingTask !== null && + existingTask.enabled === input.enabled && + isSameSchedule(existingTask.schedule, input.schedule); + const task: ScheduledTask = { + id, + title: input.title, + prompt: input.prompt, + enabled: input.enabled, + schedule: input.schedule, + projectId: input.projectId, + threadId: input.threadId ?? null, + workspaceStrategy: input.workspaceStrategy, + modelSelection: input.modelSelection, + runtimeMode: input.runtimeMode, + interactionMode: input.interactionMode, + createdBy: existingTask?.createdBy ?? input.createdBy ?? "user", + creationSource: input.creationSource ?? "web", + createdAt: existingTask?.createdAt ?? iso(now), + updatedAt: iso(now), + nextRunAt: scheduleUnchanged + ? existingTask.nextRunAt + : nextRunAt({ enabled: input.enabled, schedule: input.schedule }, now), + lastRunAt: existingTask?.lastRunAt ?? null, + lastRunStatus: existingTask?.lastRunStatus ?? "never", + lastRunError: existingTask?.lastRunError ?? null, + runCount: existingTask?.runCount ?? 0, + }; + yield* saveTask(task); + yield* notifyChanged; + return { task }; + }), + ); }); const setEnabled: ScheduledTaskService["Service"]["setEnabled"] = (input) => - Effect.gen(function* () { - const existing = yield* loadTask(input.id); - if (existing.enabled === input.enabled) return { task: existing }; - const now = yield* localNow; - const next = nextRunAt({ enabled: input.enabled, schedule: existing.schedule }, now); - // RETURNING so a task deleted between the load and this UPDATE is a - // visible not-found error, not a false success. - const updated = yield* sql<{ task_id: string }>` - UPDATE scheduled_tasks - SET enabled = ${input.enabled ? 1 : 0}, - next_run_at = ${next}, - updated_at = ${iso(now)} - WHERE task_id = ${input.id} - RETURNING task_id - `.pipe( - Effect.mapError((cause) => - taskError("Could not update schedule task.", { taskId: input.id, cause }), - ), - ); - if (updated.length === 0) { - return yield* taskError("Schedule task not found.", { taskId: input.id }); - } - yield* notifyChanged; - return { - task: { ...existing, enabled: input.enabled, nextRunAt: next, updatedAt: iso(now) }, - }; - }); + taskMutations.withLock( + input.id, + Effect.gen(function* () { + const existing = yield* loadTask(input.id); + if (existing.enabled === input.enabled) return { task: existing }; + const now = yield* localNow; + const next = nextRunAt({ enabled: input.enabled, schedule: existing.schedule }, now); + const updated = yield* sql<{ task_id: string }>` + UPDATE scheduled_tasks + SET enabled = ${input.enabled ? 1 : 0}, + next_run_at = ${next}, + updated_at = ${iso(now)} + WHERE task_id = ${input.id} + RETURNING task_id + `.pipe( + Effect.mapError((cause) => + taskError("Could not update schedule task.", { taskId: input.id, cause }), + ), + ); + if (updated.length === 0) { + return yield* taskError("Schedule task not found.", { taskId: input.id }); + } + yield* notifyChanged; + return { + task: { ...existing, enabled: input.enabled, nextRunAt: next, updatedAt: iso(now) }, + }; + }), + ); const deleteTask: ScheduledTaskService["Service"]["delete"] = (input) => - deleteRow(input.id).pipe(Effect.andThen(notifyChanged), Effect.as({ id: input.id })); + taskMutations.withLock( + input.id, + deleteRow(input.id).pipe(Effect.andThen(notifyChanged), Effect.as({ id: input.id })), + ); const runNow: ScheduledTaskService["Service"]["runNow"] = (input: ScheduledTaskRunNowInput) => Effect.gen(function* () { @@ -784,7 +1143,33 @@ export const layer = Layer.effect( taskError("Could not run schedule task.", { taskId: input.id, cause }), ), ); - return { task: next }; + return { task: next.task }; + }); + + const runNowIdempotent: ScheduledTaskService["Service"]["runNowIdempotent"] = (input) => + Effect.gen(function* () { + const replay = yield* findAcceptedManualRun(input); + if (Option.isSome(replay)) return replay.value; + const task = yield* findTask(input.id); + if (task === null) { + return yield* new ScheduledTaskManualRunNotFoundError({ + taskId: input.id, + message: "Schedule task not found.", + }); + } + const outcome = yield* runTask(task, "manual", input); + if (outcome.dispatchError !== null) { + return yield* taskError("Could not dispatch schedule task run.", { + taskId: input.id, + cause: outcome.dispatchError, + }); + } + if (outcome.manualRun === null) { + return yield* taskError("The schedule task run was not dispatched.", { + taskId: input.id, + }); + } + return { ...outcome.manualRun, task: outcome.task }; }); return ScheduledTaskService.of({ @@ -794,6 +1179,7 @@ export const layer = Layer.effect( setEnabled, delete: deleteTask, runNow, + runNowIdempotent, }); }), ); diff --git a/docs/orchestration-v2/orchestrator-mcp-server.md b/docs/orchestration-v2/orchestrator-mcp-server.md index 215a350a2f04..a494fc739ea1 100644 --- a/docs/orchestration-v2/orchestrator-mcp-server.md +++ b/docs/orchestration-v2/orchestrator-mcp-server.md @@ -141,7 +141,7 @@ selection model-visible without allowing a request that cannot run. ## Tool Surface -The server exposes eleven orchestration tools. +The server exposes the following orchestration tools. ### `orchestrator_capabilities` @@ -229,6 +229,25 @@ Interrupts the original delegated run through the normal V2 `run.interrupt` command. It is idempotent for terminal tasks and accepts an optional cancellation reason. Use `t3_thread_interrupt` to interrupt a later follow-up run. +### `run_scheduled_task_now` + +Dispatches one existing scheduled task through `ScheduledTaskService` without changing +the task definition or enabled state. The request requires a stable `clientRequestId`. +An exact retry returns the original thread, message, run, and accepted command receipt; +reusing that key for another task is rejected. + +Fresh runs re-read the task and any bound target while holding the scheduler's task +admission lock. The caller and target also participate in the V2 thread dispatch locks, +so a task edit, caller downgrade, or target upgrade cannot create a broader run between +preflight and command acceptance. Unbound tasks launch through `ThreadLaunchService`; +bound tasks use `ThreadManagementService`. + +Success means the scheduled prompt was durably accepted for dispatch, not that the +provider turn completed. A manual attempt records run status and count and computes the +next occurrence from the current schedule. On an accepted replay after the task was +deleted, bookkeeping fields are returned as `null` because only the original dispatch +receipt remains authoritative. + ### `create_threads` Creates between one and twenty ordinary top-level T3 threads: diff --git a/docs/user/scheduled-tasks.md b/docs/user/scheduled-tasks.md new file mode 100644 index 000000000000..f81c330696fb --- /dev/null +++ b/docs/user/scheduled-tasks.md @@ -0,0 +1,12 @@ +# Scheduled tasks + +Scheduled tasks send a saved prompt on a recurring interval or at selected local times. A task can +post back into an existing thread or create a new thread for each run. + +An agent with access to T3 Code orchestration can run an existing task immediately. This manual run +does not enable, disable, or otherwise edit the schedule. It does count as a run and recalculates the +next occurrence from the task's current schedule. + +An accepted manual-run result means T3 Code durably queued or started the prompt. The agent turn may +still be running. Retrying with the same request key returns the original target run instead of +starting a duplicate. diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 541e65767b0b..bbf7a5c71736 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -73,6 +73,13 @@ const OrchestrationV2CreationFields = { creationSource: OrchestrationV2CreationSource, } as const; +export const OrchestrationV2PolicyCeiling = Schema.Struct({ + callerThreadId: ThreadId, + runtimeMode: RuntimeMode, + interactionMode: ProviderInteractionMode, +}); +export type OrchestrationV2PolicyCeiling = typeof OrchestrationV2PolicyCeiling.Type; + export const OrchestrationV2NativeRefStrength = Schema.Literals(["strong", "weak", "none"]); export type OrchestrationV2NativeRefStrength = typeof OrchestrationV2NativeRefStrength.Type; @@ -2035,6 +2042,7 @@ export const OrchestrationV2Command = Schema.Union([ modelSelection: ModelSelection, runtimeMode: RuntimeMode, interactionMode: ProviderInteractionMode, + policyCeiling: Schema.optional(OrchestrationV2PolicyCeiling), branch: Schema.NullOr(TrimmedNonEmptyString), worktreePath: Schema.NullOr(TrimmedNonEmptyString), }), @@ -2186,6 +2194,7 @@ export const OrchestrationV2Command = Schema.Union([ /** Seed the temporary title and generate a durable replacement for the first message. */ titleSeed: Schema.optional(TrimmedNonEmptyString), modelSelection: Schema.optional(ModelSelection), + policyCeiling: Schema.optional(OrchestrationV2PolicyCeiling), sourcePlanRef: Schema.optional(Schema.Struct({ threadId: ThreadId, planId: PlanId })), restartContinuationOfRunId: Schema.optional(RunId), delegatedCompletion: Schema.optional( diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index 98a891af273d..8e5d90fe3640 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -3,6 +3,7 @@ import * as Schema from "effect/Schema"; import * as SchemaTransformation from "effect/SchemaTransformation"; import { + CommandId, ContextTransferId, IsoDateTime, MessageId, @@ -565,6 +566,34 @@ export const OrchestratorMcpDeleteScheduledTaskResult = Schema.Struct({ export type OrchestratorMcpDeleteScheduledTaskResult = typeof OrchestratorMcpDeleteScheduledTaskResult.Type; +export const OrchestratorMcpRunScheduledTaskNowInput = Schema.Struct({ + scheduledTaskId: ScheduledTaskId, + clientRequestId: OrchestratorMcpClientRequestId.annotate({ + description: + "Required stable idempotency key. Reuse it only when retrying this exact manual run.", + }), +}); +export type OrchestratorMcpRunScheduledTaskNowInput = + typeof OrchestratorMcpRunScheduledTaskNowInput.Type; + +export const OrchestratorMcpRunScheduledTaskNowResult = Schema.Struct({ + scheduledTaskId: ScheduledTaskId, + threadId: ThreadId, + messageId: MessageId, + runId: RunId, + status: OrchestrationV2RunStatus, + replayed: Schema.Boolean, + receipt: Schema.Struct({ + commandId: CommandId, + acceptedAt: IsoDateTime, + resultSequence: NonNegativeInt, + }), + nextRunAt: Schema.NullOr(IsoDateTime), + runCount: Schema.NullOr(NonNegativeInt), +}); +export type OrchestratorMcpRunScheduledTaskNowResult = + typeof OrchestratorMcpRunScheduledTaskNowResult.Type; + export class OrchestratorMcpFailure extends Schema.TaggedErrorClass()( "OrchestratorMcpFailure", { diff --git a/packages/shared/src/t3McpToolPresentation.ts b/packages/shared/src/t3McpToolPresentation.ts index a9dcd685cc0c..4af296082bc7 100644 --- a/packages/shared/src/t3McpToolPresentation.ts +++ b/packages/shared/src/t3McpToolPresentation.ts @@ -44,6 +44,7 @@ const T3_MCP_TOOLS: Record< displayName: "Delete a scheduled task", summaryAction: "schedule-delete", }, + run_scheduled_task_now: { displayName: "Run a scheduled task now" }, create_threads: { displayName: "Create T3 threads", summaryAction: "thread-create" }, t3_thread_start: { displayName: "Start a T3 thread", summaryAction: "thread-create" }, t3_thread_list: { displayName: "List T3 threads", summaryAction: "thread-list" }, From c6ff539a2a793fc3e1c04378ca131a5649264f11 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 29 Aug 2026 19:49:45 -0700 Subject: [PATCH 2/9] fix(mcp): serialize scheduled run admission --- ...OrchestratorMcpToolkit.integration.test.ts | 16 ++ .../src/orchestration-v2/Orchestrator.ts | 7 + .../ScheduledTaskService.test.ts | 260 +++++++++++++++-- .../scheduledTasks/ScheduledTaskService.ts | 265 ++++++++++++------ .../contracts/src/orchestratorMcp.test.ts | 13 + packages/contracts/src/orchestratorMcp.ts | 9 +- 6 files changed, 458 insertions(+), 112 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index ba5ae8c23cd1..24d6d0977564 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -42,6 +42,7 @@ import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as PubSub from "effect/PubSub"; import * as Ref from "effect/Ref"; +import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import { McpSchema, McpServer } from "effect/unstable/ai"; @@ -1315,6 +1316,10 @@ describe("orchestrator MCP toolkit", () => { }, required: expect.arrayContaining(["scheduledTaskId", "clientRequestId"]), }); + const runNowInputSchema = runScheduledTaskNowTool?.tool.inputSchema as + | { readonly properties?: Record } + | undefined; + expect(runNowInputSchema?.properties?.clientRequestId).toEqual({ type: "string" }); expect(runScheduledTaskNowTool?.tool.annotations).toMatchObject({ destructiveHint: true, idempotentHint: true, @@ -1385,6 +1390,17 @@ describe("orchestrator MCP toolkit", () => { runCount: 1, }); + const malformedScheduledRunNowCall = yield* Effect.result( + invoke("run_scheduled_task_now", { + scheduledTaskId, + clientRequestId: "run-scheduled-\ud800", + }), + ); + expect(Result.isFailure(malformedScheduledRunNowCall)).toBe(true); + if (Result.isFailure(malformedScheduledRunNowCall)) { + expect(String(malformedScheduledRunNowCall.failure)).toContain("well-formed Unicode"); + } + // delete_scheduled_task removes it entirely. const scheduledDeleteCall = yield* invoke("delete_scheduled_task", { scheduledTaskId }); expect(scheduledDeleteCall.isError).toBe(false); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 585a23b20844..b72b0b65fb79 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -2967,6 +2967,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } } + if (projection.thread.deletedAt !== null || projection.thread.archivedAt !== null) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Thread ${command.threadId} is ${projection.thread.deletedAt !== null ? "deleted" : "archived"}.`, + }); + } yield* enforcePolicyCeiling({ command, projectId: projection.thread.projectId, diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index 56450bc90572..0fb2c5ae3de7 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -18,6 +18,7 @@ import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; @@ -35,12 +36,7 @@ import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterReg import * as ThreadLaunch from "../orchestration-v2/ThreadLaunchService.ts"; import * as ThreadManagement from "../orchestration-v2/ThreadManagementService.ts"; import { makeOrchestratorV2ReplayLayerWithRegistry } from "../orchestration-v2/testkit/ProviderReplayHarness.ts"; -import { - ScheduledTaskManualRunConflictError, - ScheduledTaskManualRunRuntimeCeilingError, - ScheduledTaskService, - layer as scheduledTaskLayer, -} from "./ScheduledTaskService.ts"; +import * as ScheduledTasks from "./ScheduledTaskService.ts"; const projectId = ProjectId.make("project:scheduled-run-now"); const otherProjectId = ProjectId.make("project:scheduled-run-now-other"); @@ -75,6 +71,9 @@ const adapter = { interface HarnessOptions { readonly getProject?: ProjectService.ProjectService["Service"]["getById"]; readonly providerAdapter?: ProviderAdapterV2Shape | null; + readonly mapThreadManagement?: ( + service: ThreadManagement.ThreadManagementService["Service"], + ) => ThreadManagement.ThreadManagementService["Service"]; } function makeHarness(options: HarnessOptions = {}) { @@ -88,6 +87,18 @@ function makeHarness(options: HarnessOptions = {}) { { databaseLayer: database, runEffectWorker: false }, ); const threads = ThreadManagement.layer.pipe(Layer.provide(orchestrator)); + const schedulerThreads = + options.mapThreadManagement === undefined + ? threads + : Layer.effect( + ThreadManagement.ThreadManagementService, + Effect.gen(function* () { + const service = yield* ThreadManagement.ThreadManagementService; + return ThreadManagement.ThreadManagementService.of( + options.mapThreadManagement!(service), + ); + }), + ).pipe(Layer.provide(threads)); const receipts = CommandReceiptStore.layer.pipe(Layer.provide(database)); const externalServices = Layer.mergeAll( Layer.succeed( @@ -124,8 +135,8 @@ function makeHarness(options: HarnessOptions = {}) { const launches = ThreadLaunch.layer.pipe( Layer.provide(Layer.mergeAll(externalServices, threads, receipts, IdAllocator.layer)), ); - const scheduler = scheduledTaskLayer.pipe( - Layer.provide(Layer.mergeAll(database, launches, threads)), + const scheduler = ScheduledTasks.layer.pipe( + Layer.provide(Layer.mergeAll(database, launches, schedulerThreads)), ); return Layer.mergeAll(database, orchestrator, threads, launches, scheduler).pipe( Layer.provideMerge(NodeServices.layer), @@ -212,7 +223,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { (it) => { it.effect("dispatches through ThreadManagement and ThreadLaunch", () => Effect.gen(function* () { - const scheduler = yield* ScheduledTaskService; + const scheduler = yield* ScheduledTasks.ScheduledTaskService; const threads = yield* ThreadManagement.ThreadManagementService; yield* createThread({ commandId: "command:caller:create", threadId: callerThreadId }); yield* createThread({ commandId: "command:bound:create", threadId: boundThreadId }); @@ -272,8 +283,9 @@ describe("ScheduledTaskService.runNowIdempotent", () => { it.layer(makeHarness(), { timeout: "30 seconds" })("accepted retry behavior", (it) => { it.effect("replays after task deletion and rejects the same key for another task", () => Effect.gen(function* () { - const scheduler = yield* ScheduledTaskService; + const scheduler = yield* ScheduledTasks.ScheduledTaskService; const threads = yield* ThreadManagement.ThreadManagementService; + const orchestrator = yield* OrchestratorV2; yield* createThread({ commandId: "command:caller:retry", threadId: callerThreadId }); yield* createThread({ commandId: "command:bound:retry", threadId: boundThreadId }); @@ -281,6 +293,12 @@ describe("ScheduledTaskService.runNowIdempotent", () => { const firstInput = manualRunInput({ taskId: firstTaskId, key: "accepted-replay" }); yield* scheduler.upsert(taskInput({ id: firstTaskId, threadId: boundThreadId })); const accepted = yield* scheduler.runNowIdempotent(firstInput); + yield* orchestrator.dispatch({ + type: "thread.runtime-mode.set", + commandId: CommandId.make("command:caller:accepted-replay:downgrade"), + threadId: callerThreadId, + runtimeMode: "approval-required", + }); yield* scheduler.delete({ id: firstTaskId }); const replayed = yield* scheduler.runNowIdempotent(firstInput); expect(replayed).toEqual({ ...accepted, task: null, replayed: true }); @@ -294,18 +312,171 @@ describe("ScheduledTaskService.runNowIdempotent", () => { commandId: firstInput.commandId, }), ); - expect(conflict).toBeInstanceOf(ScheduledTaskManualRunConflictError); + expect(conflict).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); }), ); }); + it.effect("serializes a manual command from receipt lookup through bookkeeping", () => + Effect.gen(function* () { + const receiptLookupCompleted = yield* Deferred.make(); + const releaseReceiptLookup = yield* Deferred.make(); + const taskId = ScheduledTaskId.make("scheduled-task:concurrent-command"); + const input = manualRunInput({ taskId, key: "concurrent-command" }); + let gated = false; + const harness = makeHarness({ + mapThreadManagement: (service) => ({ + ...service, + getCommandReceipt: (commandId) => + service.getCommandReceipt(commandId).pipe( + Effect.flatMap((receipt) => { + if (gated || commandId !== input.commandId || Option.isSome(receipt)) { + return Effect.succeed(receipt); + } + gated = true; + return Deferred.succeed(receiptLookupCompleted, undefined).pipe( + Effect.andThen(Deferred.await(releaseReceiptLookup)), + Effect.as(receipt), + ); + }), + ), + }), + }); + + yield* Effect.gen(function* () { + const scheduler = yield* ScheduledTasks.ScheduledTaskService; + const threads = yield* ThreadManagement.ThreadManagementService; + yield* createThread({ commandId: "command:caller:concurrent", threadId: callerThreadId }); + yield* createThread({ commandId: "command:bound:concurrent", threadId: boundThreadId }); + yield* scheduler.upsert(taskInput({ id: taskId, threadId: boundThreadId })); + + const first = yield* scheduler.runNowIdempotent(input).pipe(Effect.forkChild); + yield* Deferred.await(receiptLookupCompleted); + const retry = yield* scheduler.runNowIdempotent(input).pipe(Effect.forkChild); + yield* Deferred.succeed(releaseReceiptLookup, undefined); + + const accepted = yield* Fiber.join(first); + const replayed = yield* Fiber.join(retry); + expect(accepted.replayed).toBe(false); + expect(replayed).toEqual({ ...accepted, task: null, replayed: true }); + expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); + expect((yield* scheduler.list()).tasks.find((task) => task.id === taskId)).toMatchObject({ + runCount: 1, + nextRunAt: accepted.task?.nextRunAt, + }); + }).pipe(Effect.provide(harness), Effect.scoped); + }), + ); + + it.effect("rejects concurrent cross-task reuse before mutating the losing task", () => + Effect.gen(function* () { + const receiptLookupCompleted = yield* Deferred.make(); + const releaseReceiptLookup = yield* Deferred.make(); + const firstTaskId = ScheduledTaskId.make("scheduled-task:concurrent-key-first"); + const secondTaskId = ScheduledTaskId.make("scheduled-task:concurrent-key-second"); + const firstInput = manualRunInput({ taskId: firstTaskId, key: "concurrent-cross-task" }); + const secondInput = { + ...manualRunInput({ taskId: secondTaskId, key: "concurrent-cross-task" }), + commandId: firstInput.commandId, + }; + let gated = false; + const harness = makeHarness({ + mapThreadManagement: (service) => ({ + ...service, + getCommandReceipt: (commandId) => + service.getCommandReceipt(commandId).pipe( + Effect.flatMap((receipt) => { + if (gated || commandId !== firstInput.commandId || Option.isSome(receipt)) { + return Effect.succeed(receipt); + } + gated = true; + return Deferred.succeed(receiptLookupCompleted, undefined).pipe( + Effect.andThen(Deferred.await(releaseReceiptLookup)), + Effect.as(receipt), + ); + }), + ), + }), + }); + + yield* Effect.gen(function* () { + const scheduler = yield* ScheduledTasks.ScheduledTaskService; + const threads = yield* ThreadManagement.ThreadManagementService; + yield* createThread({ commandId: "command:caller:cross-task", threadId: callerThreadId }); + yield* createThread({ commandId: "command:bound:cross-task", threadId: boundThreadId }); + yield* scheduler.upsert(taskInput({ id: firstTaskId, threadId: boundThreadId })); + yield* scheduler.upsert(taskInput({ id: secondTaskId, threadId: boundThreadId })); + + const first = yield* scheduler.runNowIdempotent(firstInput).pipe(Effect.forkChild); + yield* Deferred.await(receiptLookupCompleted); + const second = yield* scheduler.runNowIdempotent(secondInput).pipe(Effect.forkChild); + yield* Deferred.succeed(releaseReceiptLookup, undefined); + + yield* Fiber.join(first); + const conflict = yield* Fiber.join(second).pipe(Effect.flip); + expect(conflict).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); + expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); + const tasks = (yield* scheduler.list()).tasks; + expect(tasks.find((task) => task.id === firstTaskId)?.runCount).toBe(1); + expect(tasks.find((task) => task.id === secondTaskId)?.runCount).toBe(0); + }).pipe(Effect.provide(harness), Effect.scoped); + }), + ); + + it.layer(makeHarness(), { timeout: "30 seconds" })("fresh target lifecycle acceptance", (it) => { + it.effect("rejects when archive wins after the caller preflight", () => + Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const threads = yield* ThreadManagement.ThreadManagementService; + yield* createThread({ commandId: "command:caller:archive-race", threadId: callerThreadId }); + yield* createThread({ commandId: "command:target:archive-race", threadId: boundThreadId }); + const preflightCompleted = yield* Deferred.make(); + const allowAcceptance = yield* Deferred.make(); + + const send = yield* Effect.gen(function* () { + const target = yield* threads.getProjectThread({ projectId, threadId: boundThreadId }); + expect(target.thread.archivedAt).toBeNull(); + yield* Deferred.succeed(preflightCompleted, undefined); + yield* Deferred.await(allowAcceptance); + return yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("command:target:archive-race:message"), + threadId: boundThreadId, + messageId: MessageId.make("message:target:archive-race"), + text: "Run after the preflight.", + attachments: [], + policyCeiling: { + callerThreadId, + runtimeMode: "full-access", + interactionMode: "default", + }, + dispatchMode: { type: "start_immediately" }, + createdBy: "agent", + creationSource: "mcp", + }); + }).pipe(Effect.forkChild); + + yield* Deferred.await(preflightCompleted); + yield* orchestrator.dispatch({ + type: "thread.archive", + commandId: CommandId.make("command:target:archive-race:archive"), + threadId: boundThreadId, + }); + yield* Deferred.succeed(allowAcceptance, undefined); + const failure = yield* Fiber.join(send).pipe(Effect.flip); + expect(failure._tag).toBe("OrchestratorDispatchError"); + expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(0); + }), + ); + }); + it.layer(makeHarness({ providerAdapter: null }), { timeout: "30 seconds" })( "partial launch failure behavior", (it) => { it.effect("does not duplicate an unbound thread after provider dispatch rejection", () => Effect.gen(function* () { - const scheduler = yield* ScheduledTaskService; + const scheduler = yield* ScheduledTasks.ScheduledTaskService; const orchestrator = yield* OrchestratorV2; yield* createThread({ commandId: "command:caller:partial", threadId: callerThreadId }); const taskId = ScheduledTaskId.make("scheduled-task:partial-launch"); @@ -365,7 +536,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { }); yield* Effect.gen(function* () { - const scheduler = yield* ScheduledTaskService; + const scheduler = yield* ScheduledTasks.ScheduledTaskService; const orchestrator = yield* OrchestratorV2; yield* createThread({ commandId: "command:caller:race", threadId: callerThreadId }); const taskId = ScheduledTaskId.make("scheduled-task:admission-race"); @@ -382,7 +553,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { const overlap = yield* Effect.flip( scheduler.runNowIdempotent(manualRunInput({ taskId, key: "admission-race-overlap" })), ); - expect(overlap).toBeInstanceOf(ScheduledTaskManualRunConflictError); + expect(overlap).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); yield* orchestrator.dispatch({ type: "thread.runtime-mode.set", @@ -418,7 +589,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { }); yield* Effect.gen(function* () { - const scheduler = yield* ScheduledTaskService; + const scheduler = yield* ScheduledTasks.ScheduledTaskService; yield* createThread({ commandId: "command:caller:recurring", threadId: callerThreadId }); const taskId = ScheduledTaskId.make("scheduled-task:recurring-overlap"); yield* scheduler.upsert(taskInput({ id: taskId, threadId: null })); @@ -438,7 +609,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { const overlap = yield* Effect.flip( scheduler.runNowIdempotent(manualRunInput({ taskId, key: "recurring-overlap-manual" })), ); - expect(overlap).toBeInstanceOf(ScheduledTaskManualRunConflictError); + expect(overlap).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); yield* Deferred.succeed(allowLookup, undefined); const snapshot = yield* Fiber.join(completed); expect(Option.getOrThrow(snapshot)?.tasks.find((task) => task.id === taskId)).toMatchObject( @@ -451,10 +622,61 @@ describe("ScheduledTaskService.runNowIdempotent", () => { }), ); + it.effect("does not apply a stale missed-run reschedule after the task is disabled", () => + Effect.gen(function* () { + const firstDispatchEntered = yield* Deferred.make(); + const releaseFirstDispatch = yield* Deferred.make(); + let lookupCount = 0; + const harness = makeHarness({ + getProject: (id) => { + lookupCount += 1; + return lookupCount === 1 + ? Deferred.succeed(firstDispatchEntered, undefined).pipe( + Effect.andThen(Deferred.await(releaseFirstDispatch)), + Effect.as(id === projectId ? Option.some(project) : Option.none()), + ) + : Effect.succeed(id === projectId ? Option.some(project) : Option.none()); + }, + }); + + yield* Effect.gen(function* () { + const scheduler = yield* ScheduledTasks.ScheduledTaskService; + const sql = yield* SqlClient.SqlClient; + const blockerTaskId = ScheduledTaskId.make("scheduled-task:a-poll-blocker"); + const staleTaskId = ScheduledTaskId.make("scheduled-task:z-stale-fixed-time"); + yield* scheduler.upsert({ + ...taskInput({ id: staleTaskId, threadId: null }), + schedule: { type: "fixed_time", timeOfDay: "09:00" }, + }); + yield* scheduler.upsert(taskInput({ id: blockerTaskId, threadId: null })); + const staleDueAt = "1969-12-30T00:00:00.000Z"; + yield* sql` + UPDATE scheduled_tasks + SET next_run_at = ${staleDueAt} + WHERE task_id IN (${blockerTaskId}, ${staleTaskId}) + `; + + const poll = yield* TestClock.adjust(Duration.seconds(5)).pipe(Effect.forkChild); + yield* Deferred.await(firstDispatchEntered); + yield* scheduler.setEnabled({ id: staleTaskId, enabled: false }); + yield* Deferred.succeed(releaseFirstDispatch, undefined); + yield* Fiber.join(poll); + + expect( + (yield* scheduler.list()).tasks.find((task) => task.id === staleTaskId), + ).toMatchObject({ + enabled: false, + nextRunAt: null, + runCount: 0, + }); + }).pipe(Effect.provide(harness), Effect.scoped); + }), + ); + it.layer(makeHarness(), { timeout: "30 seconds" })("fresh scope and ceiling checks", (it) => { it.effect("rejects a higher-mode bound target and a task from another project", () => Effect.gen(function* () { - const scheduler = yield* ScheduledTaskService; + const scheduler = yield* ScheduledTasks.ScheduledTaskService; yield* createThread({ commandId: "command:caller:limited", threadId: callerThreadId, @@ -473,7 +695,9 @@ describe("ScheduledTaskService.runNowIdempotent", () => { }), ), ); - expect(ceilingFailure).toBeInstanceOf(ScheduledTaskManualRunRuntimeCeilingError); + expect(ceilingFailure).toBeInstanceOf( + ScheduledTasks.ScheduledTaskManualRunRuntimeCeilingError, + ); const otherTaskId = ScheduledTaskId.make("scheduled-task:other-project"); yield* scheduler.upsert( diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index b1f2a8379323..7f6fd1747f46 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -4,7 +4,9 @@ import { type OrchestrationV2PolicyCeiling, type OrchestrationV2RunStatus, type OrchestrationV2ThreadProjection, + ProviderInteractionMode, ProjectId, + RuntimeMode, type RunId, ScheduledTask, ScheduledTaskError, @@ -95,28 +97,93 @@ export interface ScheduledTaskManualRunResult { export class ScheduledTaskManualRunScopeError extends Schema.TaggedErrorClass()( "ScheduledTaskManualRunScopeError", - { taskId: ScheduledTaskId, message: Schema.String }, -) {} + { + taskId: ScheduledTaskId, + reason: Schema.Literals([ + "caller-not-in-project", + "target-not-in-project", + "caller-archived", + "target-archived", + "task-not-in-project", + "accepted-run-unauthorized", + ]), + }, +) { + override get message(): string { + switch (this.reason) { + case "caller-not-in-project": + return "The scheduled task run caller is not active in the calling project."; + case "target-not-in-project": + return "The scheduled task target is not active in the calling project."; + case "caller-archived": + return "The scheduled task run caller is archived."; + case "target-archived": + return "The scheduled task target is archived."; + case "task-not-in-project": + return "The scheduled task does not belong to the calling project."; + case "accepted-run-unauthorized": + return "The accepted run is no longer authorized in the calling project."; + } + } +} export class ScheduledTaskManualRunRuntimeCeilingError extends Schema.TaggedErrorClass()( "ScheduledTaskManualRunRuntimeCeilingError", - { taskId: ScheduledTaskId, message: Schema.String }, -) {} + { taskId: ScheduledTaskId, targetMode: RuntimeMode, ceilingMode: RuntimeMode }, +) { + override get message(): string { + return `Target runtime mode ${this.targetMode} exceeds the caller ceiling ${this.ceilingMode}.`; + } +} export class ScheduledTaskManualRunInteractionCeilingError extends Schema.TaggedErrorClass()( "ScheduledTaskManualRunInteractionCeilingError", - { taskId: ScheduledTaskId, message: Schema.String }, -) {} + { + taskId: ScheduledTaskId, + targetMode: ProviderInteractionMode, + ceilingMode: ProviderInteractionMode, + }, +) { + override get message(): string { + return `Target interaction mode ${this.targetMode} exceeds the caller ceiling ${this.ceilingMode}.`; + } +} export class ScheduledTaskManualRunConflictError extends Schema.TaggedErrorClass()( "ScheduledTaskManualRunConflictError", - { taskId: ScheduledTaskId, commandId: CommandId, message: Schema.String }, -) {} + { + taskId: ScheduledTaskId, + commandId: CommandId, + reason: Schema.Literals([ + "different-scheduled-run", + "different-command", + "different-scheduled-task-run", + "task-already-running", + ]), + }, +) { + override get message(): string { + switch (this.reason) { + case "different-scheduled-run": + return "The idempotency key belongs to a different scheduled task run."; + case "different-command": + return "The idempotency key belongs to a different command."; + case "different-scheduled-task-run": + return "The idempotency key was already accepted for a different scheduled task run."; + case "task-already-running": + return "Schedule task is already running."; + } + } +} export class ScheduledTaskManualRunNotFoundError extends Schema.TaggedErrorClass()( "ScheduledTaskManualRunNotFoundError", - { taskId: ScheduledTaskId, message: Schema.String }, -) {} + { taskId: ScheduledTaskId }, +) { + override get message(): string { + return "Schedule task not found."; + } +} export type ScheduledTaskManualRunError = | ScheduledTaskError @@ -247,6 +314,7 @@ export const layer = Layer.effect( const threadManagement = yield* ThreadManagementService.ThreadManagementService; const activeRuns = yield* Ref.make>(new Set()); const taskMutations = yield* makeKeyedSerialExecutor(); + const manualRunMutations = yield* makeKeyedSerialExecutor(); // Sliding(1) coalesces the dirty-signal: every notification triggers a // full list() re-emit anyway, so a slow subscriber only ever needs the // latest signal β€” an unbounded backlog would just grow memory. @@ -367,7 +435,7 @@ export const layer = Layer.effect( cause._tag === "ThreadManagementThreadNotFoundError" ? new ScheduledTaskManualRunScopeError({ taskId: input.id, - message: `The scheduled task run ${role} is not active in the calling project.`, + reason: role === "caller" ? "caller-not-in-project" : "target-not-in-project", }) : taskError(`Could not authorize the scheduled task run ${role}.`, { taskId: input.id, @@ -379,7 +447,7 @@ export const layer = Layer.effect( if (caller.thread.archivedAt !== null) { return yield* new ScheduledTaskManualRunScopeError({ taskId: input.id, - message: "The scheduled task run caller is archived.", + reason: "caller-archived", }); } const target = @@ -389,28 +457,33 @@ export const layer = Layer.effect( if (target !== null && target.thread.archivedAt !== null) { return yield* new ScheduledTaskManualRunScopeError({ taskId: input.id, - message: "The scheduled task target is archived.", + reason: "target-archived", }); } const runtimeMode = target?.thread.runtimeMode ?? task.runtimeMode; const interactionMode = target?.thread.interactionMode ?? task.interactionMode; - if ( - runtimeModeRank(runtimeMode) > runtimeModeRank(input.policyCeiling.runtimeMode) || - runtimeModeRank(runtimeMode) > runtimeModeRank(caller.thread.runtimeMode) - ) { + const runtimeCeiling = + runtimeModeRank(input.policyCeiling.runtimeMode) <= + runtimeModeRank(caller.thread.runtimeMode) + ? input.policyCeiling.runtimeMode + : caller.thread.runtimeMode; + if (runtimeModeRank(runtimeMode) > runtimeModeRank(runtimeCeiling)) { return yield* new ScheduledTaskManualRunRuntimeCeilingError({ taskId: input.id, - message: `Target runtime mode ${runtimeMode} exceeds the caller ceiling.`, + targetMode: runtimeMode, + ceilingMode: runtimeCeiling, }); } - if ( - interactionModeRank(interactionMode) > - interactionModeRank(input.policyCeiling.interactionMode) || - interactionModeRank(interactionMode) > interactionModeRank(caller.thread.interactionMode) - ) { + const interactionCeiling = + interactionModeRank(input.policyCeiling.interactionMode) <= + interactionModeRank(caller.thread.interactionMode) + ? input.policyCeiling.interactionMode + : caller.thread.interactionMode; + if (interactionModeRank(interactionMode) > interactionModeRank(interactionCeiling)) { return yield* new ScheduledTaskManualRunInteractionCeilingError({ taskId: input.id, - message: `Target interaction mode ${interactionMode} exceeds the caller ceiling.`, + targetMode: interactionMode, + ceilingMode: interactionCeiling, }); } }, @@ -436,7 +509,7 @@ export const layer = Layer.effect( return yield* new ScheduledTaskManualRunConflictError({ taskId: input.id, commandId: input.commandId, - message: "The idempotency key belongs to a different scheduled task run.", + reason: "different-scheduled-run", }); } // Thread creation committed, but the scheduled prompt did not. Let @@ -455,7 +528,7 @@ export const layer = Layer.effect( return yield* new ScheduledTaskManualRunConflictError({ taskId: input.id, commandId: input.commandId, - message: "The idempotency key belongs to a different command.", + reason: "different-command", }); } @@ -486,28 +559,7 @@ export const layer = Layer.effect( ) { return yield* new ScheduledTaskManualRunScopeError({ taskId: input.id, - message: "The accepted run is no longer authorized in the calling project.", - }); - } - if ( - runtimeModeRank(target.thread.runtimeMode) > - runtimeModeRank(input.policyCeiling.runtimeMode) || - runtimeModeRank(target.thread.runtimeMode) > runtimeModeRank(caller.thread.runtimeMode) - ) { - return yield* new ScheduledTaskManualRunRuntimeCeilingError({ - taskId: input.id, - message: `Target runtime mode ${target.thread.runtimeMode} exceeds the caller ceiling.`, - }); - } - if ( - interactionModeRank(target.thread.interactionMode) > - interactionModeRank(input.policyCeiling.interactionMode) || - interactionModeRank(target.thread.interactionMode) > - interactionModeRank(caller.thread.interactionMode) - ) { - return yield* new ScheduledTaskManualRunInteractionCeilingError({ - taskId: input.id, - message: `Target interaction mode ${target.thread.interactionMode} exceeds the caller ceiling.`, + reason: "accepted-run-unauthorized", }); } @@ -520,7 +572,7 @@ export const layer = Layer.effect( return yield* new ScheduledTaskManualRunConflictError({ taskId: input.id, commandId: input.commandId, - message: "The idempotency key was already accepted for a different scheduled task run.", + reason: "different-scheduled-task-run", }); } return Option.some({ @@ -738,7 +790,7 @@ export const layer = Layer.effect( return yield* new ScheduledTaskManualRunConflictError({ taskId: task.id, commandId: manualRun.commandId, - message: "Schedule task is already running.", + reason: "task-already-running", }); } return yield* taskError("Schedule task is already running.", { taskId: task.id }); @@ -761,7 +813,6 @@ export const layer = Layer.effect( if (manualRun !== undefined) { return yield* new ScheduledTaskManualRunNotFoundError({ taskId: task.id, - message: "Schedule task not found.", }); } return yield* taskError("Schedule task not found.", { taskId: task.id }); @@ -771,7 +822,7 @@ export const layer = Layer.effect( if (manualRun !== undefined && active.projectId !== manualRun.projectId) { return yield* new ScheduledTaskManualRunScopeError({ taskId: active.id, - message: "The scheduled task does not belong to the calling project.", + reason: "task-not-in-project", }); } if (manualRun !== undefined) { @@ -910,23 +961,47 @@ export const layer = Layer.effect( task: ScheduledTask, now: DateTime.DateTime, ) { - const next = nextRunAt(task, now); - yield* Effect.logInfo("Skipping missed schedule task run", { - taskId: task.id, - missedRunAt: task.nextRunAt, - rescheduledTo: next, - }); - yield* sql` - UPDATE scheduled_tasks - SET next_run_at = ${next}, - updated_at = ${iso(now)} - WHERE task_id = ${task.id} - `.pipe( - Effect.mapError((cause) => - taskError("Could not reschedule missed schedule task run.", { taskId: task.id, cause }), - ), + yield* taskMutations.withLock( + task.id, + Effect.gen(function* () { + const current = yield* findTask(task.id); + if ( + current === null || + !current.enabled || + current.nextRunAt === null || + current.lastRunStatus === "running" + ) { + return; + } + const dueAt = DateTime.makeUnsafe(current.nextRunAt); + if ( + DateTime.toEpochMillis(dueAt) > DateTime.toEpochMillis(now) || + !isMissedFixedTimeRun(current.schedule, dueAt, now) + ) { + return; + } + const next = nextRunAt(current, now); + yield* Effect.logInfo("Skipping missed schedule task run", { + taskId: current.id, + missedRunAt: current.nextRunAt, + rescheduledTo: next, + }); + yield* sql` + UPDATE scheduled_tasks + SET next_run_at = ${next}, + updated_at = ${iso(now)} + WHERE task_id = ${current.id} + `.pipe( + Effect.mapError((cause) => + taskError("Could not reschedule missed schedule task run.", { + taskId: current.id, + cause, + }), + ), + ); + yield* notifyChanged; + }), ); - yield* notifyChanged; }); const runDueTasks = Effect.fn("ScheduledTaskService.runDueTasks")(function* () { @@ -951,7 +1026,9 @@ export const layer = Layer.effect( ).pipe( Effect.asVoid, Effect.catchCause((cause) => - Effect.logWarning("Scheduled task run failed", { taskId: task.id, cause }), + Cause.hasInterrupts(cause) + ? Effect.failCause(cause) + : Effect.logWarning("Scheduled task run failed", { taskId: task.id, cause }), ), ), { concurrency: 1, discard: true }, @@ -1147,30 +1224,32 @@ export const layer = Layer.effect( }); const runNowIdempotent: ScheduledTaskService["Service"]["runNowIdempotent"] = (input) => - Effect.gen(function* () { - const replay = yield* findAcceptedManualRun(input); - if (Option.isSome(replay)) return replay.value; - const task = yield* findTask(input.id); - if (task === null) { - return yield* new ScheduledTaskManualRunNotFoundError({ - taskId: input.id, - message: "Schedule task not found.", - }); - } - const outcome = yield* runTask(task, "manual", input); - if (outcome.dispatchError !== null) { - return yield* taskError("Could not dispatch schedule task run.", { - taskId: input.id, - cause: outcome.dispatchError, - }); - } - if (outcome.manualRun === null) { - return yield* taskError("The schedule task run was not dispatched.", { - taskId: input.id, - }); - } - return { ...outcome.manualRun, task: outcome.task }; - }); + manualRunMutations.withLock( + input.commandId, + Effect.gen(function* () { + const replay = yield* findAcceptedManualRun(input); + if (Option.isSome(replay)) return replay.value; + const task = yield* findTask(input.id); + if (task === null) { + return yield* new ScheduledTaskManualRunNotFoundError({ + taskId: input.id, + }); + } + const outcome = yield* runTask(task, "manual", input); + if (outcome.dispatchError !== null) { + return yield* taskError("Could not dispatch schedule task run.", { + taskId: input.id, + cause: outcome.dispatchError, + }); + } + if (outcome.manualRun === null) { + return yield* taskError("The schedule task run was not dispatched.", { + taskId: input.id, + }); + } + return { ...outcome.manualRun, task: outcome.task }; + }), + ); return ScheduledTaskService.of({ list, diff --git a/packages/contracts/src/orchestratorMcp.test.ts b/packages/contracts/src/orchestratorMcp.test.ts index 3e698f882c78..3198ab8b9dfb 100644 --- a/packages/contracts/src/orchestratorMcp.test.ts +++ b/packages/contracts/src/orchestratorMcp.test.ts @@ -5,6 +5,7 @@ import { OrchestratorMcpCreateThreadsInput, OrchestratorMcpDelegateTaskInput, OrchestratorMcpDelegateTaskResult, + OrchestratorMcpRunScheduledTaskNowInput, OrchestratorMcpThreadInterruptInput, OrchestratorMcpThreadListInput, OrchestratorMcpThreadReadInput, @@ -16,6 +17,9 @@ import { const decodeCreateThreadsInput = Schema.decodeUnknownSync(OrchestratorMcpCreateThreadsInput); const decodeDelegateTaskInput = Schema.decodeUnknownSync(OrchestratorMcpDelegateTaskInput); const decodeDelegateTaskResult = Schema.decodeUnknownSync(OrchestratorMcpDelegateTaskResult); +const decodeRunScheduledTaskNowInput = Schema.decodeUnknownSync( + OrchestratorMcpRunScheduledTaskNowInput, +); const decodeThreadInterruptInput = Schema.decodeUnknownSync(OrchestratorMcpThreadInterruptInput); const decodeThreadListInput = Schema.decodeUnknownSync(OrchestratorMcpThreadListInput); const decodeThreadReadInput = Schema.decodeUnknownSync(OrchestratorMcpThreadReadInput); @@ -24,6 +28,15 @@ const decodeThreadStartInput = Schema.decodeUnknownSync(OrchestratorMcpThreadSta const decodeThreadWaitInput = Schema.decodeUnknownSync(OrchestratorMcpThreadWaitInput); describe("orchestrator MCP contracts", () => { + it("rejects malformed Unicode in manual-run idempotency keys", () => { + expect(() => + decodeRunScheduledTaskNowInput({ + scheduledTaskId: "scheduled-task:unicode-key", + clientRequestId: "manual-run-\ud800", + }), + ).toThrow(/well-formed Unicode/); + }); + it("decodes cross-provider delegated task requests and durable results", () => { const request = decodeDelegateTaskInput({ task: "Inspect the workspace and report the result.", diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index 8e5d90fe3640..81758e906f21 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -46,6 +46,13 @@ const OrchestratorMcpTitle = TrimmedNonEmptyString.check(Schema.isMaxLength(512) const OrchestratorMcpClientRequestId = TrimmedNonEmptyString.check( Schema.isMaxLength(256), ).annotate({ description: "Stable idempotency key to reuse when retrying this mutation." }); +const OrchestratorMcpWellFormedClientRequestId = TrimmedNonEmptyString.check( + Schema.isMaxLength(256), +).check( + Schema.makeFilter( + (value) => value.isWellFormed() || "Idempotency key must contain well-formed Unicode.", + ), +); /** * OpenCode 1.15 has been observed serializing nested MCP union objects as JSON @@ -568,7 +575,7 @@ export type OrchestratorMcpDeleteScheduledTaskResult = export const OrchestratorMcpRunScheduledTaskNowInput = Schema.Struct({ scheduledTaskId: ScheduledTaskId, - clientRequestId: OrchestratorMcpClientRequestId.annotate({ + clientRequestId: OrchestratorMcpWellFormedClientRequestId.annotate({ description: "Required stable idempotency key. Reuse it only when retrying this exact manual run.", }), From db655de410b29f2bacab382aaf2294bc43148267 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 29 Aug 2026 19:53:36 -0700 Subject: [PATCH 3/9] fix(contracts): validate Unicode keys portably --- packages/contracts/src/orchestratorMcp.ts | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index 81758e906f21..813a351524fc 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -46,11 +46,27 @@ const OrchestratorMcpTitle = TrimmedNonEmptyString.check(Schema.isMaxLength(512) const OrchestratorMcpClientRequestId = TrimmedNonEmptyString.check( Schema.isMaxLength(256), ).annotate({ description: "Stable idempotency key to reuse when retrying this mutation." }); + +function isWellFormedUnicode(value: string): boolean { + for (let index = 0; index < value.length; index += 1) { + const codeUnit = value.charCodeAt(index); + if (codeUnit >= 0xd800 && codeUnit <= 0xdbff) { + if (index + 1 >= value.length) return false; + const nextCodeUnit = value.charCodeAt(index + 1); + if (nextCodeUnit < 0xdc00 || nextCodeUnit > 0xdfff) return false; + index += 1; + } else if (codeUnit >= 0xdc00 && codeUnit <= 0xdfff) { + return false; + } + } + return true; +} + const OrchestratorMcpWellFormedClientRequestId = TrimmedNonEmptyString.check( Schema.isMaxLength(256), ).check( Schema.makeFilter( - (value) => value.isWellFormed() || "Idempotency key must contain well-formed Unicode.", + (value) => isWellFormedUnicode(value) || "Idempotency key must contain well-formed Unicode.", ), ); From a87358672aa630d0941ebe12875e0936189fa730 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 29 Aug 2026 20:01:46 -0700 Subject: [PATCH 4/9] refactor(server): distinguish manual run failures --- apps/server/src/mcp/OrchestratorMcpService.ts | 12 +- .../ScheduledTaskService.test.ts | 10 +- .../scheduledTasks/ScheduledTaskService.ts | 187 +++++++++++------- 3 files changed, 134 insertions(+), 75 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 8e5ed5dd2009..68b34dd22b4e 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -193,13 +193,21 @@ function errorMessage(error: unknown): string { function scheduledTaskManualRunFailure(error: ScheduledTaskManualRunError): OrchestratorMcpFailure { switch (error._tag) { case "ScheduledTaskManualRunNotFoundError": - case "ScheduledTaskManualRunScopeError": + case "ScheduledTaskManualRunCallerScopeError": + case "ScheduledTaskManualRunTargetScopeError": + case "ScheduledTaskManualRunCallerArchivedError": + case "ScheduledTaskManualRunTargetArchivedError": + case "ScheduledTaskManualRunTaskScopeError": + case "ScheduledTaskManualRunAcceptedRunScopeError": return failure("task_not_found", error.message); case "ScheduledTaskManualRunRuntimeCeilingError": return failure("runtime_mode_escalation_denied", error.message); case "ScheduledTaskManualRunInteractionCeilingError": return failure("interaction_mode_escalation_denied", error.message); - case "ScheduledTaskManualRunConflictError": + case "ScheduledTaskManualRunReceiptThreadConflictError": + case "ScheduledTaskManualRunCommandConflictError": + case "ScheduledTaskManualRunMessageConflictError": + case "ScheduledTaskManualRunAlreadyRunningError": return failure("invalid_request", error.message); case "ScheduledTaskError": return failure("orchestration_error", error.message); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index 0fb2c5ae3de7..460b8b80054e 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -312,7 +312,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { commandId: firstInput.commandId, }), ); - expect(conflict).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); + expect(conflict).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunMessageConflictError); expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); }), ); @@ -415,7 +415,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { yield* Fiber.join(first); const conflict = yield* Fiber.join(second).pipe(Effect.flip); - expect(conflict).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); + expect(conflict).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunMessageConflictError); expect((yield* threads.getThreadProjection(boundThreadId)).messages).toHaveLength(1); const tasks = (yield* scheduler.list()).tasks; expect(tasks.find((task) => task.id === firstTaskId)?.runCount).toBe(1); @@ -553,7 +553,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { const overlap = yield* Effect.flip( scheduler.runNowIdempotent(manualRunInput({ taskId, key: "admission-race-overlap" })), ); - expect(overlap).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); + expect(overlap).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunAlreadyRunningError); yield* orchestrator.dispatch({ type: "thread.runtime-mode.set", @@ -609,7 +609,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { const overlap = yield* Effect.flip( scheduler.runNowIdempotent(manualRunInput({ taskId, key: "recurring-overlap-manual" })), ); - expect(overlap).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunConflictError); + expect(overlap).toBeInstanceOf(ScheduledTasks.ScheduledTaskManualRunAlreadyRunningError); yield* Deferred.succeed(allowLookup, undefined); const snapshot = yield* Fiber.join(completed); expect(Option.getOrThrow(snapshot)?.tasks.find((task) => task.id === taskId)).toMatchObject( @@ -706,7 +706,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { const scopeFailure = yield* Effect.flip( scheduler.runNowIdempotent(manualRunInput({ taskId: otherTaskId, key: "other-project" })), ); - assert.equal(scopeFailure._tag, "ScheduledTaskManualRunScopeError"); + assert.equal(scopeFailure._tag, "ScheduledTaskManualRunTaskScopeError"); }), ); }); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index 7f6fd1747f46..e657b93bc385 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -95,35 +95,62 @@ export interface ScheduledTaskManualRunResult { readonly receipt: CommandReceiptV2; } -export class ScheduledTaskManualRunScopeError extends Schema.TaggedErrorClass()( - "ScheduledTaskManualRunScopeError", +export class ScheduledTaskManualRunCallerScopeError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunCallerScopeError", + { taskId: ScheduledTaskId, callerThreadId: ThreadId, projectId: ProjectId }, +) { + override get message(): string { + return "The scheduled task run caller is not active in the calling project."; + } +} + +export class ScheduledTaskManualRunTargetScopeError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunTargetScopeError", + { taskId: ScheduledTaskId, targetThreadId: ThreadId, projectId: ProjectId }, +) { + override get message(): string { + return "The scheduled task target is not active in the calling project."; + } +} + +export class ScheduledTaskManualRunCallerArchivedError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunCallerArchivedError", + { taskId: ScheduledTaskId, callerThreadId: ThreadId }, +) { + override get message(): string { + return "The scheduled task run caller is archived."; + } +} + +export class ScheduledTaskManualRunTargetArchivedError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunTargetArchivedError", + { taskId: ScheduledTaskId, targetThreadId: ThreadId }, +) { + override get message(): string { + return "The scheduled task target is archived."; + } +} + +export class ScheduledTaskManualRunTaskScopeError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunTaskScopeError", + { taskId: ScheduledTaskId, projectId: ProjectId }, +) { + override get message(): string { + return "The scheduled task does not belong to the calling project."; + } +} + +export class ScheduledTaskManualRunAcceptedRunScopeError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunAcceptedRunScopeError", { taskId: ScheduledTaskId, - reason: Schema.Literals([ - "caller-not-in-project", - "target-not-in-project", - "caller-archived", - "target-archived", - "task-not-in-project", - "accepted-run-unauthorized", - ]), + callerThreadId: ThreadId, + targetThreadId: ThreadId, + projectId: ProjectId, }, ) { override get message(): string { - switch (this.reason) { - case "caller-not-in-project": - return "The scheduled task run caller is not active in the calling project."; - case "target-not-in-project": - return "The scheduled task target is not active in the calling project."; - case "caller-archived": - return "The scheduled task run caller is archived."; - case "target-archived": - return "The scheduled task target is archived."; - case "task-not-in-project": - return "The scheduled task does not belong to the calling project."; - case "accepted-run-unauthorized": - return "The accepted run is no longer authorized in the calling project."; - } + return "The accepted run is no longer authorized in the calling project."; } } @@ -149,30 +176,39 @@ export class ScheduledTaskManualRunInteractionCeilingError extends Schema.Tagged } } -export class ScheduledTaskManualRunConflictError extends Schema.TaggedErrorClass()( - "ScheduledTaskManualRunConflictError", - { - taskId: ScheduledTaskId, - commandId: CommandId, - reason: Schema.Literals([ - "different-scheduled-run", - "different-command", - "different-scheduled-task-run", - "task-already-running", - ]), - }, +export class ScheduledTaskManualRunReceiptThreadConflictError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunReceiptThreadConflictError", + { taskId: ScheduledTaskId, commandId: CommandId, receiptThreadId: ThreadId }, +) { + override get message(): string { + return "The idempotency key belongs to a different scheduled task run."; + } +} + +export class ScheduledTaskManualRunCommandConflictError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunCommandConflictError", + { taskId: ScheduledTaskId, commandId: CommandId }, ) { override get message(): string { - switch (this.reason) { - case "different-scheduled-run": - return "The idempotency key belongs to a different scheduled task run."; - case "different-command": - return "The idempotency key belongs to a different command."; - case "different-scheduled-task-run": - return "The idempotency key was already accepted for a different scheduled task run."; - case "task-already-running": - return "Schedule task is already running."; - } + return "The idempotency key belongs to a different command."; + } +} + +export class ScheduledTaskManualRunMessageConflictError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunMessageConflictError", + { taskId: ScheduledTaskId, commandId: CommandId, messageId: MessageId }, +) { + override get message(): string { + return "The idempotency key was already accepted for a different scheduled task run."; + } +} + +export class ScheduledTaskManualRunAlreadyRunningError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunAlreadyRunningError", + { taskId: ScheduledTaskId, commandId: CommandId }, +) { + override get message(): string { + return "Schedule task is already running."; } } @@ -187,10 +223,18 @@ export class ScheduledTaskManualRunNotFoundError extends Schema.TaggedErrorClass export type ScheduledTaskManualRunError = | ScheduledTaskError - | ScheduledTaskManualRunScopeError + | ScheduledTaskManualRunCallerScopeError + | ScheduledTaskManualRunTargetScopeError + | ScheduledTaskManualRunCallerArchivedError + | ScheduledTaskManualRunTargetArchivedError + | ScheduledTaskManualRunTaskScopeError + | ScheduledTaskManualRunAcceptedRunScopeError | ScheduledTaskManualRunRuntimeCeilingError | ScheduledTaskManualRunInteractionCeilingError - | ScheduledTaskManualRunConflictError + | ScheduledTaskManualRunReceiptThreadConflictError + | ScheduledTaskManualRunCommandConflictError + | ScheduledTaskManualRunMessageConflictError + | ScheduledTaskManualRunAlreadyRunningError | ScheduledTaskManualRunNotFoundError; export class ScheduledTaskService extends Context.Service< @@ -433,10 +477,17 @@ export const layer = Layer.effect( threadManagement.getProjectThread({ projectId: input.projectId, threadId }).pipe( Effect.mapError((cause) => cause._tag === "ThreadManagementThreadNotFoundError" - ? new ScheduledTaskManualRunScopeError({ - taskId: input.id, - reason: role === "caller" ? "caller-not-in-project" : "target-not-in-project", - }) + ? role === "caller" + ? new ScheduledTaskManualRunCallerScopeError({ + taskId: input.id, + callerThreadId: threadId, + projectId: input.projectId, + }) + : new ScheduledTaskManualRunTargetScopeError({ + taskId: input.id, + targetThreadId: threadId, + projectId: input.projectId, + }) : taskError(`Could not authorize the scheduled task run ${role}.`, { taskId: input.id, cause, @@ -445,9 +496,9 @@ export const layer = Layer.effect( ); const caller = yield* loadScopedThread(input.policyCeiling.callerThreadId, "caller"); if (caller.thread.archivedAt !== null) { - return yield* new ScheduledTaskManualRunScopeError({ + return yield* new ScheduledTaskManualRunCallerArchivedError({ taskId: input.id, - reason: "caller-archived", + callerThreadId: input.policyCeiling.callerThreadId, }); } const target = @@ -455,9 +506,9 @@ export const layer = Layer.effect( ? null : yield* loadScopedThread(ThreadId.make(task.threadId), "target"); if (target !== null && target.thread.archivedAt !== null) { - return yield* new ScheduledTaskManualRunScopeError({ + return yield* new ScheduledTaskManualRunTargetArchivedError({ taskId: input.id, - reason: "target-archived", + targetThreadId: target.thread.id, }); } const runtimeMode = target?.thread.runtimeMode ?? task.runtimeMode; @@ -506,10 +557,10 @@ export const layer = Layer.effect( } if (primaryReceipt.value.commandType === "thread.create") { if (primaryReceipt.value.threadId !== input.unboundThreadId) { - return yield* new ScheduledTaskManualRunConflictError({ + return yield* new ScheduledTaskManualRunReceiptThreadConflictError({ taskId: input.id, commandId: input.commandId, - reason: "different-scheduled-run", + receiptThreadId: primaryReceipt.value.threadId, }); } // Thread creation committed, but the scheduled prompt did not. Let @@ -525,10 +576,9 @@ export const layer = Layer.effect( }); } if (receipt.commandType !== "message.dispatch") { - return yield* new ScheduledTaskManualRunConflictError({ + return yield* new ScheduledTaskManualRunCommandConflictError({ taskId: input.id, commandId: input.commandId, - reason: "different-command", }); } @@ -557,9 +607,11 @@ export const layer = Layer.effect( target.thread.deletedAt !== null || target.thread.projectId !== input.projectId ) { - return yield* new ScheduledTaskManualRunScopeError({ + return yield* new ScheduledTaskManualRunAcceptedRunScopeError({ taskId: input.id, - reason: "accepted-run-unauthorized", + callerThreadId: input.policyCeiling.callerThreadId, + targetThreadId: receipt.threadId, + projectId: input.projectId, }); } @@ -569,10 +621,10 @@ export const layer = Layer.effect( ? undefined : target.runs.find((candidate) => candidate.id === message.runId); if (message === undefined || run === undefined) { - return yield* new ScheduledTaskManualRunConflictError({ + return yield* new ScheduledTaskManualRunMessageConflictError({ taskId: input.id, commandId: input.commandId, - reason: "different-scheduled-task-run", + messageId: input.messageId, }); } return Option.some({ @@ -787,10 +839,9 @@ export const layer = Layer.effect( if (!reserved) { if (trigger === "manual") { if (manualRun !== undefined) { - return yield* new ScheduledTaskManualRunConflictError({ + return yield* new ScheduledTaskManualRunAlreadyRunningError({ taskId: task.id, commandId: manualRun.commandId, - reason: "task-already-running", }); } return yield* taskError("Schedule task is already running.", { taskId: task.id }); @@ -820,9 +871,9 @@ export const layer = Layer.effect( return { task, manualRun: null, dispatchError: null }; } if (manualRun !== undefined && active.projectId !== manualRun.projectId) { - return yield* new ScheduledTaskManualRunScopeError({ + return yield* new ScheduledTaskManualRunTaskScopeError({ taskId: active.id, - reason: "task-not-in-project", + projectId: manualRun.projectId, }); } if (manualRun !== undefined) { From 6eef5305e7c5e189f284d5cdc9d7236cc32c74e6 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 30 Aug 2026 10:35:22 -0700 Subject: [PATCH 5/9] fix(server): align receipt service typing --- apps/server/src/orchestration-v2/Orchestrator.ts | 4 ++-- docs/orchestration-v2/orchestrator-mcp-server.md | 6 ++++++ 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index b72b0b65fb79..5abb960f5f0c 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -41,7 +41,7 @@ import * as Stream from "effect/Stream"; import { CheckpointServiceV2 } from "./CheckpointService.ts"; import { CommandPolicyV2 } from "./CommandPolicy.ts"; -import { CommandReceiptStoreV2, type CommandReceiptStoreV2Shape } from "./CommandReceiptStore.ts"; +import { CommandReceiptStoreV2 } from "./CommandReceiptStore.ts"; import { ContextHandoffServiceV2 } from "./ContextHandoffService.ts"; import { EventSinkV2 } from "./EventSink.ts"; import type { OrchestrationEffectRequestV2, PendingOrchestrationEffectV2 } from "./EffectOutbox.ts"; @@ -181,7 +181,7 @@ export interface OrchestratorV2Shape { readonly dispatch: ( command: OrchestrationV2Command, ) => Effect.Effect; - readonly getCommandReceipt: CommandReceiptStoreV2Shape["getByCommandId"]; + readonly getCommandReceipt: CommandReceiptStoreV2["Service"]["getByCommandId"]; readonly getThreadProjection: ( threadId: ThreadId, ) => Effect.Effect; diff --git a/docs/orchestration-v2/orchestrator-mcp-server.md b/docs/orchestration-v2/orchestrator-mcp-server.md index a494fc739ea1..759f17b725fc 100644 --- a/docs/orchestration-v2/orchestrator-mcp-server.md +++ b/docs/orchestration-v2/orchestrator-mcp-server.md @@ -248,6 +248,12 @@ next occurrence from the current schedule. On an accepted replay after the task deleted, bookkeeping fields are returned as `null` because only the original dispatch receipt remains authoritative. +For an unbound task, thread creation and initial-message dispatch are separate +durable commands. If a fresh policy check or provider admission rejects the +message after creation, the empty thread is retained because concurrent work +may already own it. The failed key remains a rejected replay and does not +repeat schedule bookkeeping; use a new key after correcting the rejection. + ### `create_threads` Creates between one and twenty ordinary top-level T3 threads: From 35f1f65e8be68b1879e6e61af5b140e46e34a312 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 30 Aug 2026 13:38:47 -0700 Subject: [PATCH 6/9] fix(server): resume partial scheduled launches once --- .../ScheduledTaskService.test.ts | 80 ++++++++++++++++++- .../scheduledTasks/ScheduledTaskService.ts | 79 +++++++++++------- 2 files changed, 130 insertions(+), 29 deletions(-) diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index 460b8b80054e..d59eb95a87de 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -74,6 +74,9 @@ interface HarnessOptions { readonly mapThreadManagement?: ( service: ThreadManagement.ThreadManagementService["Service"], ) => ThreadManagement.ThreadManagementService["Service"]; + readonly mapThreadLaunch?: ( + service: ThreadLaunch.ThreadLaunchService["Service"], + ) => ThreadLaunch.ThreadLaunchService["Service"]; } function makeHarness(options: HarnessOptions = {}) { @@ -135,10 +138,20 @@ function makeHarness(options: HarnessOptions = {}) { const launches = ThreadLaunch.layer.pipe( Layer.provide(Layer.mergeAll(externalServices, threads, receipts, IdAllocator.layer)), ); + const schedulerLaunches = + options.mapThreadLaunch === undefined + ? launches + : Layer.effect( + ThreadLaunch.ThreadLaunchService, + Effect.gen(function* () { + const service = yield* ThreadLaunch.ThreadLaunchService; + return ThreadLaunch.ThreadLaunchService.of(options.mapThreadLaunch!(service)); + }), + ).pipe(Layer.provide(launches)); const scheduler = ScheduledTasks.layer.pipe( - Layer.provide(Layer.mergeAll(database, launches, schedulerThreads)), + Layer.provide(Layer.mergeAll(database, schedulerLaunches, schedulerThreads)), ); - return Layer.mergeAll(database, orchestrator, threads, launches, scheduler).pipe( + return Layer.mergeAll(database, orchestrator, threads, schedulerLaunches, scheduler).pipe( Layer.provideMerge(NodeServices.layer), ); } @@ -518,6 +531,69 @@ describe("ScheduledTaskService.runNowIdempotent", () => { }, ); + it.effect("resumes an accepted unbound create without counting another schedule run", () => + Effect.gen(function* () { + let failAfterCreate = true; + const harness = makeHarness({ + mapThreadLaunch: (service) => + ThreadLaunch.ThreadLaunchService.of({ + launch: (input) => + failAfterCreate + ? Effect.gen(function* () { + failAfterCreate = false; + const { initialMessage: _initialMessage, ...createOnly } = input; + yield* service.launch(createOnly); + return yield* new ThreadLaunch.ThreadLaunchError({ + operation: "dispatch-message", + commandId: input.commandId, + projectId: input.projectId, + ...(input.threadId === undefined ? {} : { threadId: input.threadId }), + cause: "Simulated failure after durable thread creation.", + }); + }) + : service.launch(input), + }), + }); + + yield* Effect.gen(function* () { + const scheduler = yield* ScheduledTasks.ScheduledTaskService; + const threads = yield* ThreadManagement.ThreadManagementService; + yield* createThread({ + commandId: "command:caller:partial-create-resume", + threadId: callerThreadId, + }); + const taskId = ScheduledTaskId.make("scheduled-task:partial-create-resume"); + const targetThreadId = ThreadId.make("thread:scheduled-run-now:partial-create-resume"); + const input = manualRunInput({ + taskId, + key: "partial-create-resume", + unboundThreadId: targetThreadId, + }); + yield* scheduler.upsert(taskInput({ id: taskId, threadId: null })); + + const firstFailure = yield* scheduler.runNowIdempotent(input).pipe(Effect.flip); + expect(firstFailure._tag).toBe("ScheduledTaskError"); + const afterFailure = (yield* scheduler.list()).tasks.find((task) => task.id === taskId); + expect(afterFailure).toMatchObject({ runCount: 1, lastRunStatus: "failed" }); + expect((yield* threads.getThreadProjection(targetThreadId)).messages).toHaveLength(0); + + const resumed = yield* scheduler.runNowIdempotent(input); + expect(resumed).toMatchObject({ + threadId: targetThreadId, + messageId: input.messageId, + replayed: false, + receipt: { commandType: "message.dispatch", status: "accepted" }, + }); + expect(resumed.task).toMatchObject({ + runCount: 1, + lastRunStatus: "failed", + nextRunAt: afterFailure?.nextRunAt, + }); + expect((yield* threads.getThreadProjection(targetThreadId)).messages).toHaveLength(1); + }).pipe(Effect.provide(harness), Effect.scoped); + }), + ); + it.effect("serializes active admission and rejects a caller downgrade before V2 acceptance", () => Effect.gen(function* () { const lookupEntered = yield* Deferred.make(); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index e657b93bc385..f7a988059f5d 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -95,6 +95,11 @@ export interface ScheduledTaskManualRunResult { readonly receipt: CommandReceiptV2; } +type ScheduledTaskManualRunReceiptLookup = + | { readonly type: "none" } + | { readonly type: "accepted"; readonly result: ScheduledTaskManualRunResult } + | { readonly type: "accepted_unbound_create" }; + export class ScheduledTaskManualRunCallerScopeError extends Schema.TaggedErrorClass()( "ScheduledTaskManualRunCallerScopeError", { taskId: ScheduledTaskId, callerThreadId: ThreadId, projectId: ProjectId }, @@ -548,7 +553,7 @@ export const layer = Layer.effect( if (receipt === undefined) { const primaryReceipt = yield* readCommandReceipt(input.id, input.commandId); if (Option.isNone(primaryReceipt)) { - return Option.none(); + return { type: "none" } satisfies ScheduledTaskManualRunReceiptLookup; } if (primaryReceipt.value.status === "rejected") { return yield* taskError("This manual run request was previously rejected.", { @@ -565,8 +570,11 @@ export const layer = Layer.effect( } // Thread creation committed, but the scheduled prompt did not. Let // ThreadLaunchService replay the create receipt and finish the same - // task-specific initial message. - return Option.none(); + // task-specific initial message without recording another schedule + // attempt. + return { + type: "accepted_unbound_create", + } satisfies ScheduledTaskManualRunReceiptLookup; } receipt = primaryReceipt.value; } @@ -627,15 +635,18 @@ export const layer = Layer.effect( messageId: input.messageId, }); } - return Option.some({ - task: null, - threadId: target.thread.id, - messageId: message.id, - runId: run.id, - status: run.status, - replayed: true, - receipt, - } satisfies ScheduledTaskManualRunResult); + return { + type: "accepted", + result: { + task: null, + threadId: target.thread.id, + messageId: message.id, + runId: run.id, + status: run.status, + replayed: true, + receipt, + }, + } satisfies ScheduledTaskManualRunReceiptLookup; }, ); @@ -829,6 +840,7 @@ export const layer = Layer.effect( task: ScheduledTask, trigger: "scheduled" | "manual", manualRun?: ScheduledTaskManualRunInput, + resumeAcceptedManualCreate = false, ) { const reserved = yield* Ref.modify(activeRuns, (active) => { if (active.has(task.id)) return [false, active] as const; @@ -889,8 +901,10 @@ export const layer = Layer.effect( return { task: active, manualRun: null, dispatchError: null }; } - yield* markRunning(active.id, startedAtIso); - yield* notifyChanged; + if (!resumeAcceptedManualCreate) { + yield* markRunning(active.id, startedAtIso); + yield* notifyChanged; + } const fireKey = `${active.id}:${DateTime.toEpochMillis(startedAt)}:${trigger}`; const commandId = manualRun?.commandId ?? CommandId.make(`scheduled-task:${fireKey}`); @@ -968,16 +982,18 @@ export const layer = Layer.effect( const lastRunError = runSucceeded ? null : errorMessage(result.cause); const current = yield* findTask(task.id); const scheduleSource = current ?? task; - const completed: ScheduledTask = { - ...scheduleSource, - updatedAt: iso(completedAt), - lastRunAt: startedAtIso, - nextRunAt: nextRunAt(scheduleSource, completedAt), - lastRunStatus, - lastRunError, - runCount: scheduleSource.runCount + 1, - }; - if (current !== null) { + const completed: ScheduledTask = resumeAcceptedManualCreate + ? scheduleSource + : { + ...scheduleSource, + updatedAt: iso(completedAt), + lastRunAt: startedAtIso, + nextRunAt: nextRunAt(scheduleSource, completedAt), + lastRunStatus, + lastRunError, + runCount: scheduleSource.runCount + 1, + }; + if (current !== null && !resumeAcceptedManualCreate) { yield* markCompleted({ id: task.id, completedAtIso: completed.updatedAt, @@ -993,7 +1009,11 @@ export const layer = Layer.effect( manualRun: runSucceeded ? result.value : null, dispatchError: runSucceeded ? null : result.cause, }; - }).pipe(Effect.onError((cause) => releaseStuckRun(task, errorMessage(cause)))), + }).pipe( + Effect.onError((cause) => + resumeAcceptedManualCreate ? Effect.void : releaseStuckRun(task, errorMessage(cause)), + ), + ), ) .pipe( Effect.ensuring( @@ -1279,14 +1299,19 @@ export const layer = Layer.effect( input.commandId, Effect.gen(function* () { const replay = yield* findAcceptedManualRun(input); - if (Option.isSome(replay)) return replay.value; + if (replay.type === "accepted") return replay.result; const task = yield* findTask(input.id); if (task === null) { return yield* new ScheduledTaskManualRunNotFoundError({ taskId: input.id, }); } - const outcome = yield* runTask(task, "manual", input); + const outcome = yield* runTask( + task, + "manual", + input, + replay.type === "accepted_unbound_create", + ); if (outcome.dispatchError !== null) { return yield* taskError("Could not dispatch schedule task run.", { taskId: input.id, From ef91f4a9fbfeb88dc505ad26be1f0ccf41b2ce01 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 30 Aug 2026 13:41:32 -0700 Subject: [PATCH 7/9] test(server): preserve newer schedule bookkeeping --- .../ScheduledTaskService.test.ts | 26 ++++++++++++++++--- .../orchestrator-mcp-server.md | 4 +++ docs/user/scheduled-tasks.md | 5 ++++ 3 files changed, 31 insertions(+), 4 deletions(-) diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index d59eb95a87de..3892fc6e07d8 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -531,7 +531,7 @@ describe("ScheduledTaskService.runNowIdempotent", () => { }, ); - it.effect("resumes an accepted unbound create without counting another schedule run", () => + it.effect("resumes an accepted unbound create without rewriting newer task bookkeeping", () => Effect.gen(function* () { let failAfterCreate = true; const harness = makeHarness({ @@ -577,6 +577,23 @@ describe("ScheduledTaskService.runNowIdempotent", () => { expect(afterFailure).toMatchObject({ runCount: 1, lastRunStatus: "failed" }); expect((yield* threads.getThreadProjection(targetThreadId)).messages).toHaveLength(0); + const newerInput = manualRunInput({ + taskId, + key: "partial-create-newer-run", + unboundThreadId: ThreadId.make("thread:scheduled-run-now:partial-create-newer-run"), + }); + const newerRun = yield* scheduler.runNowIdempotent(newerInput); + const afterNewerRun = (yield* scheduler.list()).tasks.find((task) => task.id === taskId); + expect(newerRun.receipt).toMatchObject({ + commandType: "message.dispatch", + status: "accepted", + }); + expect(afterNewerRun).toMatchObject({ + runCount: 2, + lastRunStatus: "succeeded", + nextRunAt: expect.any(String), + }); + const resumed = yield* scheduler.runNowIdempotent(input); expect(resumed).toMatchObject({ threadId: targetThreadId, @@ -585,10 +602,11 @@ describe("ScheduledTaskService.runNowIdempotent", () => { receipt: { commandType: "message.dispatch", status: "accepted" }, }); expect(resumed.task).toMatchObject({ - runCount: 1, - lastRunStatus: "failed", - nextRunAt: afterFailure?.nextRunAt, + runCount: 2, + lastRunStatus: "succeeded", + nextRunAt: afterNewerRun?.nextRunAt, }); + expect(resumed.task).toEqual(afterNewerRun); expect((yield* threads.getThreadProjection(targetThreadId)).messages).toHaveLength(1); }).pipe(Effect.provide(harness), Effect.scoped); }), diff --git a/docs/orchestration-v2/orchestrator-mcp-server.md b/docs/orchestration-v2/orchestrator-mcp-server.md index 759f17b725fc..d2eca91c2170 100644 --- a/docs/orchestration-v2/orchestrator-mcp-server.md +++ b/docs/orchestration-v2/orchestrator-mcp-server.md @@ -253,6 +253,10 @@ durable commands. If a fresh policy check or provider admission rejects the message after creation, the empty thread is retained because concurrent work may already own it. The failed key remains a rejected replay and does not repeat schedule bookkeeping; use a new key after correcting the rejection. +If thread creation was accepted but no initial-message receipt exists, an exact +retry may finish that missing message through the ordinary launch path. This +resume does not rewrite `runCount`, `nextRunAt`, or the latest run status. Those +fields may describe the earlier failed attempt or a newer logical run. ### `create_threads` diff --git a/docs/user/scheduled-tasks.md b/docs/user/scheduled-tasks.md index f81c330696fb..7bd1549439c2 100644 --- a/docs/user/scheduled-tasks.md +++ b/docs/user/scheduled-tasks.md @@ -10,3 +10,8 @@ next occurrence from the task's current schedule. An accepted manual-run result means T3 Code durably queued or started the prompt. The agent turn may still be running. Retrying with the same request key returns the original target run instead of starting a duplicate. + +For tasks that create a new thread, thread creation and prompt dispatch commit separately. If thread +creation succeeds but prompt dispatch fails, retrying the same request can finish that prompt. This +resume does not increment the run count or replace the latest run status or next occurrence, so a +newer task run remains the authoritative schedule summary. From 6566a43c524747f0c2998b094631022841e81e20 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 4 Sep 2026 23:08:23 -0700 Subject: [PATCH 8/9] test(server): update scheduler service fixtures --- apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts | 2 ++ apps/server/src/relay/AgentAwarenessRelay.test.ts | 1 + 2 files changed, 3 insertions(+) diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index 24d6d0977564..4160f782970e 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -467,6 +467,8 @@ const unusedScheduledTaskStubLayer = Layer.succeed( setEnabled: () => Effect.die("ScheduledTaskService.setEnabled is unused in this test"), delete: () => Effect.die("ScheduledTaskService.delete is unused in this test"), runNow: () => Effect.die("ScheduledTaskService.runNow is unused in this test"), + runNowIdempotent: () => + Effect.die("ScheduledTaskService.runNowIdempotent is unused in this test"), }), ); diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index d19d819e595c..200babd981a2 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -126,6 +126,7 @@ const makeTestRelay = Effect.fnUntraced(function* ( getShellSnapshot: unused, ensureLegacyTranscript: unused, dispatch: unused, + getCommandReceipt: unused, getThreadProjection: unused, getCheckpointContext: unused, getThreadSnapshot: unused, From c5a6c9755adea95bd684df1efd84ec13dfe0e69b Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 4 Sep 2026 23:46:04 -0700 Subject: [PATCH 9/9] fix(mcp): validate scheduled task identifiers --- ...OrchestratorMcpToolkit.integration.test.ts | 44 +++++++++++++++++++ .../contracts/src/orchestratorMcp.test.ts | 21 ++++++++- packages/contracts/src/orchestratorMcp.ts | 8 +++- 3 files changed, 71 insertions(+), 2 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index 4160f782970e..258d325dec55 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -1402,6 +1402,50 @@ describe("orchestrator MCP toolkit", () => { if (Result.isFailure(malformedScheduledRunNowCall)) { expect(String(malformedScheduledRunNowCall.failure)).toContain("well-formed Unicode"); } + expect(yield* Ref.get(scheduledManualRuns)).toHaveLength(1); + + for (const malformedScheduledTaskId of [ + "scheduled-task:\ud800", + "scheduled-task:\udc00", + ]) { + const malformedTaskRunNowCall = yield* Effect.result( + invoke("run_scheduled_task_now", { + scheduledTaskId: malformedScheduledTaskId, + clientRequestId: "run-scheduled-valid-unicode-key", + }), + ); + expect(Result.isFailure(malformedTaskRunNowCall)).toBe(true); + if (Result.isFailure(malformedTaskRunNowCall)) { + expect(String(malformedTaskRunNowCall.failure)).toContain("well-formed Unicode"); + } + } + expect(yield* Ref.get(scheduledManualRuns)).toHaveLength(1); + + const validNonBmpScheduledTaskId = ScheduledTaskId.make("scheduled-task:πŸš€"); + const scheduledTaskTemplate = storedAfterCreate[0]; + if (scheduledTaskTemplate === undefined) { + return yield* Effect.die("Scheduled task fixture is missing"); + } + const validNonBmpTask = { + ...scheduledTaskTemplate, + id: validNonBmpScheduledTaskId, + title: "run a Unicode task", + prompt: "run a Unicode task", + }; + yield* Ref.update(scheduledStore, (tasks) => [...tasks, validNonBmpTask]); + const validNonBmpRunNowCall = yield* invoke("run_scheduled_task_now", { + scheduledTaskId: validNonBmpScheduledTaskId, + clientRequestId: "run-scheduled-πŸš€", + }); + expect(validNonBmpRunNowCall.isError).toBe(false); + expect(validNonBmpRunNowCall.structuredContent).toMatchObject({ + scheduledTaskId: validNonBmpScheduledTaskId, + threadId: parentThreadId, + }); + expect(yield* Ref.get(scheduledManualRuns)).toHaveLength(2); + yield* Ref.update(scheduledStore, (tasks) => + tasks.filter((task) => task.id !== validNonBmpScheduledTaskId), + ); // delete_scheduled_task removes it entirely. const scheduledDeleteCall = yield* invoke("delete_scheduled_task", { scheduledTaskId }); diff --git a/packages/contracts/src/orchestratorMcp.test.ts b/packages/contracts/src/orchestratorMcp.test.ts index 3198ab8b9dfb..71bc03d9d700 100644 --- a/packages/contracts/src/orchestratorMcp.test.ts +++ b/packages/contracts/src/orchestratorMcp.test.ts @@ -28,13 +28,32 @@ const decodeThreadStartInput = Schema.decodeUnknownSync(OrchestratorMcpThreadSta const decodeThreadWaitInput = Schema.decodeUnknownSync(OrchestratorMcpThreadWaitInput); describe("orchestrator MCP contracts", () => { - it("rejects malformed Unicode in manual-run idempotency keys", () => { + it("accepts well-formed Unicode and rejects malformed manual-run identifiers", () => { + expect( + decodeRunScheduledTaskNowInput({ + scheduledTaskId: "scheduled-task:πŸš€", + clientRequestId: "manual-run:πŸš€", + }), + ).toEqual({ + scheduledTaskId: "scheduled-task:πŸš€", + clientRequestId: "manual-run:πŸš€", + }); + expect(() => decodeRunScheduledTaskNowInput({ scheduledTaskId: "scheduled-task:unicode-key", clientRequestId: "manual-run-\ud800", }), ).toThrow(/well-formed Unicode/); + + for (const malformedScheduledTaskId of ["scheduled-task:\ud800", "scheduled-task:\udc00"]) { + expect(() => + decodeRunScheduledTaskNowInput({ + scheduledTaskId: malformedScheduledTaskId, + clientRequestId: "manual-run-valid-key", + }), + ).toThrow(/well-formed Unicode/); + } }); it("decodes cross-provider delegated task requests and durable results", () => { diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index 813a351524fc..233bb0477bc7 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -70,6 +70,12 @@ const OrchestratorMcpWellFormedClientRequestId = TrimmedNonEmptyString.check( ), ); +const OrchestratorMcpWellFormedScheduledTaskId = ScheduledTaskId.check( + Schema.makeFilter( + (value) => isWellFormedUnicode(value) || "Scheduled task id must contain well-formed Unicode.", + ), +); + /** * OpenCode 1.15 has been observed serializing nested MCP union objects as JSON * strings. Keep the structured object as the documented form while accepting @@ -590,7 +596,7 @@ export type OrchestratorMcpDeleteScheduledTaskResult = typeof OrchestratorMcpDeleteScheduledTaskResult.Type; export const OrchestratorMcpRunScheduledTaskNowInput = Schema.Struct({ - scheduledTaskId: ScheduledTaskId, + scheduledTaskId: OrchestratorMcpWellFormedScheduledTaskId, clientRequestId: OrchestratorMcpWellFormedClientRequestId.annotate({ description: "Required stable idempotency key. Reuse it only when retrying this exact manual run.",