diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index ea0abc1e9785..68b34dd22b4e 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,30 @@ function errorMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } +function scheduledTaskManualRunFailure(error: ScheduledTaskManualRunError): OrchestratorMcpFailure { + switch (error._tag) { + case "ScheduledTaskManualRunNotFoundError": + 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 "ScheduledTaskManualRunReceiptThreadConflictError": + case "ScheduledTaskManualRunCommandConflictError": + case "ScheduledTaskManualRunMessageConflictError": + case "ScheduledTaskManualRunAlreadyRunningError": + 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 +1168,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..258d325dec55 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, @@ -41,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"; @@ -71,7 +73,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"; @@ -462,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"), }), ); @@ -590,6 +597,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 +621,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 +1307,26 @@ 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"]), + }); + 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, + 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 +1373,80 @@ 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, + }); + + 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"); + } + 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 }); expect(scheduledDeleteCall.isError).toBe(false); @@ -1317,6 +1456,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..5abb960f5f0c 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -181,6 +181,7 @@ export interface OrchestratorV2Shape { readonly dispatch: ( command: OrchestrationV2Command, ) => Effect.Effect; + readonly getCommandReceipt: CommandReceiptStoreV2["Service"]["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,19 @@ 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, + runtimeMode: projection.thread.runtimeMode, + interactionMode: projection.thread.interactionMode, + }); if (projection.thread.settledOverride !== null) { const now = yield* DateTime.now; const thread: OrchestrationV2AppThread = { @@ -7183,8 +7275,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 +7447,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 +7556,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/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, diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts new file mode 100644 index 000000000000..3892fc6e07d8 --- /dev/null +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -0,0 +1,807 @@ +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 SqlClient from "effect/unstable/sql/SqlClient"; + +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 * as ScheduledTasks 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; + readonly mapThreadManagement?: ( + service: ThreadManagement.ThreadManagementService["Service"], + ) => ThreadManagement.ThreadManagementService["Service"]; + readonly mapThreadLaunch?: ( + service: ThreadLaunch.ThreadLaunchService["Service"], + ) => ThreadLaunch.ThreadLaunchService["Service"]; +} + +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 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( + 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 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, schedulerLaunches, schedulerThreads)), + ); + return Layer.mergeAll(database, orchestrator, threads, schedulerLaunches, 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* ScheduledTasks.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* 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 }); + + 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* 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 }); + 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(ScheduledTasks.ScheduledTaskManualRunMessageConflictError); + 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.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); + 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* ScheduledTasks.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("resumes an accepted unbound create without rewriting newer task bookkeeping", () => + 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 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, + messageId: input.messageId, + replayed: false, + receipt: { commandType: "message.dispatch", status: "accepted" }, + }); + expect(resumed.task).toMatchObject({ + 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); + }), + ); + + 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* ScheduledTasks.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(ScheduledTasks.ScheduledTaskManualRunAlreadyRunningError); + + 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* 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 })); + 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(ScheduledTasks.ScheduledTaskManualRunAlreadyRunningError); + 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.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* ScheduledTasks.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( + ScheduledTasks.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, "ScheduledTaskManualRunTaskScopeError"); + }), + ); + }); +}); diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index cb59c96afe78..f7a988059f5d 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -1,6 +1,13 @@ import { CommandId, MessageId, + type OrchestrationV2PolicyCeiling, + type OrchestrationV2RunStatus, + type OrchestrationV2ThreadProjection, + ProviderInteractionMode, + ProjectId, + RuntimeMode, + type RunId, ScheduledTask, ScheduledTaskError, ScheduledTaskId, @@ -21,6 +28,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 +37,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 +76,172 @@ 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; +} + +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 }, +) { + 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, + callerThreadId: ThreadId, + targetThreadId: ThreadId, + projectId: ProjectId, + }, +) { + override get message(): string { + return "The accepted run is no longer authorized in the calling project."; + } +} + +export class ScheduledTaskManualRunRuntimeCeilingError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunRuntimeCeilingError", + { 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, + targetMode: ProviderInteractionMode, + ceilingMode: ProviderInteractionMode, + }, +) { + override get message(): string { + return `Target interaction mode ${this.targetMode} exceeds the caller ceiling ${this.ceilingMode}.`; + } +} + +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 { + 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."; + } +} + +export class ScheduledTaskManualRunNotFoundError extends Schema.TaggedErrorClass()( + "ScheduledTaskManualRunNotFoundError", + { taskId: ScheduledTaskId }, +) { + override get message(): string { + return "Schedule task not found."; + } +} + +export type ScheduledTaskManualRunError = + | ScheduledTaskError + | ScheduledTaskManualRunCallerScopeError + | ScheduledTaskManualRunTargetScopeError + | ScheduledTaskManualRunCallerArchivedError + | ScheduledTaskManualRunTargetArchivedError + | ScheduledTaskManualRunTaskScopeError + | ScheduledTaskManualRunAcceptedRunScopeError + | ScheduledTaskManualRunRuntimeCeilingError + | ScheduledTaskManualRunInteractionCeilingError + | ScheduledTaskManualRunReceiptThreadConflictError + | ScheduledTaskManualRunCommandConflictError + | ScheduledTaskManualRunMessageConflictError + | ScheduledTaskManualRunAlreadyRunningError + | ScheduledTaskManualRunNotFoundError; + export class ScheduledTaskService extends Context.Service< ScheduledTaskService, { @@ -85,6 +261,9 @@ export class ScheduledTaskService extends Context.Service< readonly runNow: ( input: ScheduledTaskRunNowInput, ) => Effect.Effect; + readonly runNowIdempotent: ( + input: ScheduledTaskManualRunInput, + ) => Effect.Effect; } >()("t3/scheduledTasks/ScheduledTaskService") {} @@ -121,6 +300,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 +362,8 @@ 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(); + 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. @@ -269,6 +467,189 @@ 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" + ? 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, + }), + ), + ); + const caller = yield* loadScopedThread(input.policyCeiling.callerThreadId, "caller"); + if (caller.thread.archivedAt !== null) { + return yield* new ScheduledTaskManualRunCallerArchivedError({ + taskId: input.id, + callerThreadId: input.policyCeiling.callerThreadId, + }); + } + const target = + task.threadId === null + ? null + : yield* loadScopedThread(ThreadId.make(task.threadId), "target"); + if (target !== null && target.thread.archivedAt !== null) { + return yield* new ScheduledTaskManualRunTargetArchivedError({ + taskId: input.id, + targetThreadId: target.thread.id, + }); + } + const runtimeMode = target?.thread.runtimeMode ?? task.runtimeMode; + const interactionMode = target?.thread.interactionMode ?? task.interactionMode; + 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, + targetMode: runtimeMode, + ceilingMode: runtimeCeiling, + }); + } + 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, + targetMode: interactionMode, + ceilingMode: interactionCeiling, + }); + } + }, + ); + + 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 { type: "none" } satisfies ScheduledTaskManualRunReceiptLookup; + } + 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 ScheduledTaskManualRunReceiptThreadConflictError({ + taskId: input.id, + commandId: input.commandId, + receiptThreadId: primaryReceipt.value.threadId, + }); + } + // Thread creation committed, but the scheduled prompt did not. Let + // ThreadLaunchService replay the create receipt and finish the same + // task-specific initial message without recording another schedule + // attempt. + return { + type: "accepted_unbound_create", + } satisfies ScheduledTaskManualRunReceiptLookup; + } + 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 ScheduledTaskManualRunCommandConflictError({ + taskId: input.id, + commandId: input.commandId, + }); + } + + 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 ScheduledTaskManualRunAcceptedRunScopeError({ + taskId: input.id, + callerThreadId: input.policyCeiling.callerThreadId, + targetThreadId: receipt.threadId, + projectId: input.projectId, + }); + } + + 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 ScheduledTaskManualRunMessageConflictError({ + taskId: input.id, + commandId: input.commandId, + messageId: input.messageId, + }); + } + return { + type: "accepted", + result: { + task: null, + threadId: target.thread.id, + messageId: message.id, + runId: run.id, + status: run.status, + replayed: true, + receipt, + }, + } satisfies ScheduledTaskManualRunReceiptLookup; + }, + ); + // 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 +801,46 @@ 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, + resumeAcceptedManualCreate = false, ) { const reserved = yield* Ref.modify(activeRuns, (active) => { if (active.has(task.id)) return [false, active] as const; @@ -432,127 +850,180 @@ export const layer = Layer.effect( }); if (!reserved) { if (trigger === "manual") { + if (manualRun !== undefined) { + return yield* new ScheduledTaskManualRunAlreadyRunningError({ + taskId: task.id, + commandId: manualRun.commandId, + }); + } 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, + }); + } + 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 ScheduledTaskManualRunTaskScopeError({ + taskId: active.id, + projectId: manualRun.projectId, + }); + } + 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; - }), - ), - ); + 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}`); + 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 = 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, + 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) => + resumeAcceptedManualCreate ? Effect.void : 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 @@ -561,23 +1032,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* () { @@ -600,8 +1095,11 @@ export const layer = Layer.effect( ? rescheduleMissedRun(task, now) : runTask(task, "scheduled") ).pipe( - Effect.catch((cause) => - Effect.logWarning("Scheduled task run failed", { taskId: task.id, cause }), + Effect.asVoid, + Effect.catchCause((cause) => + Cause.hasInterrupts(cause) + ? Effect.failCause(cause) + : Effect.logWarning("Scheduled task run failed", { taskId: task.id, cause }), ), ), { concurrency: 1, discard: true }, @@ -690,7 +1188,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 +1201,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,9 +1291,42 @@ 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) => + manualRunMutations.withLock( + input.commandId, + Effect.gen(function* () { + const replay = yield* findAcceptedManualRun(input); + 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, + replay.type === "accepted_unbound_create", + ); + 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, subscribeList, @@ -794,6 +1334,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..d2eca91c2170 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,35 @@ 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. + +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. +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` 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..7bd1549439c2 --- /dev/null +++ b/docs/user/scheduled-tasks.md @@ -0,0 +1,17 @@ +# 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. + +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. 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.test.ts b/packages/contracts/src/orchestratorMcp.test.ts index 3e698f882c78..71bc03d9d700 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,34 @@ const decodeThreadStartInput = Schema.decodeUnknownSync(OrchestratorMcpThreadSta const decodeThreadWaitInput = Schema.decodeUnknownSync(OrchestratorMcpThreadWaitInput); describe("orchestrator MCP contracts", () => { + 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", () => { 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 98a891af273d..233bb0477bc7 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, @@ -46,6 +47,35 @@ 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) => isWellFormedUnicode(value) || "Idempotency key must contain well-formed Unicode.", + ), +); + +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 @@ -565,6 +595,34 @@ export const OrchestratorMcpDeleteScheduledTaskResult = Schema.Struct({ export type OrchestratorMcpDeleteScheduledTaskResult = typeof OrchestratorMcpDeleteScheduledTaskResult.Type; +export const OrchestratorMcpRunScheduledTaskNowInput = Schema.Struct({ + scheduledTaskId: OrchestratorMcpWellFormedScheduledTaskId, + clientRequestId: OrchestratorMcpWellFormedClientRequestId.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" },