Skip to content
83 changes: 82 additions & 1 deletion apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ import {
type OrchestratorMcpDeleteScheduledTaskResult,
type OrchestratorMcpListScheduledTasksResult,
type OrchestratorMcpRuntimeMode,
type OrchestratorMcpRunScheduledTaskNowInput,
type OrchestratorMcpRunScheduledTaskNowResult,
type OrchestratorMcpScheduledTask,
type OrchestratorMcpScheduleTaskInput,
type OrchestratorMcpScheduleTaskResult,
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -126,6 +131,10 @@ export interface OrchestratorMcpServiceShape {
scope: McpInvocationScope,
input: OrchestratorMcpDeleteScheduledTaskInput,
) => Effect.Effect<OrchestratorMcpDeleteScheduledTaskResult, OrchestratorMcpFailure>;
readonly runScheduledTaskNow: (
scope: McpInvocationScope,
input: OrchestratorMcpRunScheduledTaskNowInput,
) => Effect.Effect<OrchestratorMcpRunScheduledTaskNowResult, OrchestratorMcpFailure>;
readonly listThreads: (
scope: McpInvocationScope,
input: OrchestratorMcpThreadListInput,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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);
Expand Down
157 changes: 156 additions & 1 deletion apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import {
type ProviderOptionDescriptor,
ProviderThreadId,
ProviderTurnId,
RunId,
type ScheduledTask,
ScheduledTaskId,
type ScheduledTaskUpsertInput,
Expand All @@ -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";
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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"),
}),
);

Expand Down Expand Up @@ -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<ReadonlyArray<ScheduledTask>>([]);
const scheduledManualRuns = yield* Ref.make<ReadonlyArray<ScheduledTaskManualRunResult>>(
[],
);
const scheduledTaskStubLayer = Layer.succeed(
ScheduledTaskService,
ScheduledTaskService.of({
Expand All @@ -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(
Expand Down Expand Up @@ -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<string, unknown> }
| 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 },
Expand Down Expand Up @@ -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);
Expand All @@ -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.
Expand Down
6 changes: 6 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
17 changes: 17 additions & 0 deletions apps/server/src/mcp/toolkits/orchestrator/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import {
OrchestratorMcpListScheduledTasksResult,
OrchestratorMcpScheduleTaskInput,
OrchestratorMcpScheduleTaskResult,
OrchestratorMcpRunScheduledTaskNowInput,
OrchestratorMcpRunScheduledTaskNowResult,
OrchestratorMcpTaskCancelInput,
OrchestratorMcpTaskCancelResult,
OrchestratorMcpUpdateScheduledTaskInput,
Expand Down Expand Up @@ -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.",
Expand Down Expand Up @@ -258,6 +274,7 @@ export const OrchestratorToolkit = Toolkit.make(
ListScheduledTasksTool,
UpdateScheduledTaskTool,
DeleteScheduledTaskTool,
RunScheduledTaskNowTool,
CreateThreadsTool,
ThreadStartTool,
ThreadListTool,
Expand Down
Loading
Loading