diff --git a/docs/workflows.md b/docs/workflows.md index 0def65858..ebae3fb58 100644 --- a/docs/workflows.md +++ b/docs/workflows.md @@ -7,6 +7,22 @@ review chains, or plans that sleep and resume later. This layer deliberately reuses the deployment's existing scheduler and `work_items` queue. It does not create a second lease, timer, or coordinator. +## Agent tools + +The durable store is exposed through the first-party **Workflows** plugin, not as an ungoverned +internal shortcut. An administrator must enable the Workflows catalogue entry and grant its tools to +a Bot before that Bot can create or mutate workflow state. Calls derive the owner and Bot from the +active run connection; tool arguments cannot substitute another owner or Agent id. + +The tool surface can create/list/read workflows, pause/resume/cancel them, start/wait/complete/fail +or retry steps, and add/list/remove per-step asset metadata. Waiting a step persists an exact future +timestamp and lets the shared durable wake bridge resume it later instead of keeping a browser turn +alive or polling in a loop. + +This surface changes workflow state only. It does not grant Browser, Files, Computer, connector or +host access; execution of a step still uses the Bot's existing grants and policy at the time the +action is attempted. + ## State model A workflow run records its owner, Bot, destination channel, title and lifecycle: diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index fd89db7dc..8e2ce221a 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1479,6 +1479,9 @@ function checkDurableWorkflowState(): void { const wake = read("server/src/workflows/wake.ts"); const server = read("server/src/index.ts"); const assetMigration = read("server/drizzle/0044_workflow_assets.sql"); + const builtin = read("server/src/plugins/builtin-workflows.ts"); + const transport = read("server/src/plugins/transport.ts"); + const catalogue = read("server/src/plugins/catalogue.ts"); for (const evidence of [ 'workflowRunStatus = pgEnum("workflow_run_status"', @@ -1555,6 +1558,47 @@ function checkDurableWorkflowState(): void { } } + for (const evidence of ["useWorkflowTools(workflowStore)"]) { + if (!server.includes(evidence)) { + fail( + `workflows: builtin workflow tool store is not installed: ${evidence}`, + ); + } + } + for (const evidence of [ + '"builtin-workflows": builtinWorkflows', + '| "builtin-workflows"', + ]) { + if (!transport.includes(evidence)) { + fail(`workflows: builtin workflow transport is missing ${evidence}`); + } + } + for (const evidence of [ + 'key: "workflows"', + 'transport: "builtin-workflows"', + '"create_workflow"', + '"cancel_workflow"', + '"add_workflow_asset"', + ]) { + if (!catalogue.includes(evidence)) { + fail(`workflows: governed catalogue surface is missing ${evidence}`); + } + } + for (const evidence of [ + "connection.actorId", + "connection.botId", + '"create_workflow"', + '"wait_workflow_step"', + '"complete_workflow_step"', + '"retry_workflow_step"', + '"add_workflow_asset"', + '"list_workflow_assets"', + ]) { + if (!builtin.includes(evidence)) { + fail(`workflows: builtin tool boundary is missing ${evidence}`); + } + } + for (const evidence of [ 'CREATE TYPE "public"."workflow_run_status"', 'CREATE TYPE "public"."workflow_step_status"', diff --git a/server/src/index.ts b/server/src/index.ts index ac6970c95..ff4d6db9f 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -101,6 +101,7 @@ import { hostAccessTools } from "./host-access/tools"; import { createOnboardingStore } from "./people/onboarding"; import { createPeopleStore } from "./people/store"; import { useRoutineTools } from "./plugins/builtin-routines"; +import { useWorkflowTools } from "./plugins/builtin-workflows"; import { useComposioClient } from "./plugins/composio"; import { createComposioClient } from "./plugins/composio-adapter"; import { redirectUriFor } from "./plugins/oauth"; @@ -419,6 +420,7 @@ const pluginStore = createPluginStore({ const routineStore = createRoutineStore(database); useRoutineTools(routineStore); const workflowStore = createWorkflowStore(database); +useWorkflowTools(workflowStore); /** * Other Bots this person deliberately put in the same live channel as this Bot. diff --git a/server/src/plugins/builtin-workflows.ts b/server/src/plugins/builtin-workflows.ts new file mode 100644 index 000000000..d9b4c63d6 --- /dev/null +++ b/server/src/plugins/builtin-workflows.ts @@ -0,0 +1,503 @@ +import { cutAtCodeUnits } from "../channels/text"; +import { + type WorkflowAssetInput, + type WorkflowIdentity, + WorkflowNotFoundError, + type WorkflowPlan, + WorkflowRefusedError, + type WorkflowStepInput, + type WorkflowStore, +} from "../workflows/store"; +import { MAX_RESULT_CHARS, type McpCallResult, type McpTool } from "./mcp"; + +export type WorkflowTools = Pick< + WorkflowStore, + | "create" + | "get" + | "listFor" + | "pause" + | "resume" + | "cancel" + | "startStep" + | "waitStep" + | "completeStep" + | "failStep" + | "retryStep" + | "addAsset" + | "listAssets" + | "removeAsset" +>; + +let installed: WorkflowTools | null = null; + +export function useWorkflowTools(tools: WorkflowTools | null): void { + installed = tools; +} + +type Connection = { + url: string; + token?: string; + actorId?: string; + botId?: string; +}; + +const TOOLS: readonly McpTool[] = Object.freeze([ + { + name: "create_workflow", + description: + "Create a durable multi-step workflow for the current person and Bot. Dependencies may name only earlier step keys in this workflow.", + inputSchema: { + type: "object", + properties: { + channelId: { type: "string" }, + title: { type: "string" }, + steps: { + type: "array", + items: { + type: "object", + properties: { + key: { type: "string" }, + instruction: { type: "string" }, + dependsOn: { type: "array", items: { type: "string" } }, + }, + required: ["key", "instruction"], + }, + }, + }, + required: ["channelId", "title", "steps"], + }, + }, + { + name: "list_workflows", + description: + "List durable workflows owned by the current person and this Bot.", + inputSchema: { type: "object", properties: {} }, + }, + { + name: "get_workflow", + description: "Read one durable workflow and all of its current step state.", + inputSchema: { + type: "object", + properties: { id: { type: "string" } }, + required: ["id"], + }, + }, + { + name: "pause_workflow", + description: + "Pause an active workflow without deleting its durable step, wait, attempt or asset state.", + inputSchema: { + type: "object", + properties: { id: { type: "string" } }, + required: ["id"], + }, + }, + { + name: "resume_workflow", + description: + "Resume a paused workflow. Existing persisted waits and attempts are preserved.", + inputSchema: { + type: "object", + properties: { id: { type: "string" } }, + required: ["id"], + }, + }, + { + name: "cancel_workflow", + description: + "Cancel a workflow and all still-active steps. This is terminal for that workflow.", + inputSchema: { + type: "object", + properties: { id: { type: "string" } }, + required: ["id"], + }, + }, + { + name: "start_workflow_step", + description: + "Atomically move one ready workflow step to running and increment its durable attempt count.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + }, + required: ["id", "stepKey"], + }, + }, + { + name: "wait_workflow_step", + description: + "Put one running step to sleep until an exact future RFC3339 timestamp. Use this instead of busy-polling an external generation job.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + waitUntil: { + type: "string", + description: + "Absolute future RFC3339 timestamp with Z or an explicit numeric UTC offset.", + }, + provider: { + type: "string", + description: + "Optional provider label only. Never put credentials or secrets here.", + }, + }, + required: ["id", "stepKey", "waitUntil"], + }, + }, + { + name: "complete_workflow_step", + description: + "Mark one running step succeeded. Newly satisfied dependent steps become ready atomically.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + }, + required: ["id", "stepKey"], + }, + }, + { + name: "fail_workflow_step", + description: + "Mark one running step failed and persist a bounded failure reason for recovery or retry.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + reason: { type: "string" }, + }, + required: ["id", "stepKey", "reason"], + }, + }, + { + name: "retry_workflow_step", + description: + "Move a failed step back to ready when its dependencies still succeeded. The prior attempt count remains durable.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + }, + required: ["id", "stepKey"], + }, + }, + { + name: "add_workflow_asset", + description: + "Attach bounded asset metadata to one workflow step. References may only be attachment: or workspace:; this grants no file access.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + direction: { type: "string", enum: ["input", "output"] }, + mediaKind: { + type: "string", + enum: ["image", "video", "audio", "file", "text"], + }, + ref: { type: "string" }, + label: { type: "string" }, + }, + required: ["id", "stepKey", "direction", "mediaKind", "ref"], + }, + }, + { + name: "list_workflow_assets", + description: + "List asset metadata for one workflow, optionally limited to one step. Underlying file bytes remain behind existing file and attachment permissions.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + stepKey: { type: "string" }, + }, + required: ["id"], + }, + }, + { + name: "remove_workflow_asset", + description: + "Remove one workflow asset metadata entry. This does not delete the underlying file or attachment.", + inputSchema: { + type: "object", + properties: { + id: { type: "string" }, + assetId: { type: "string" }, + }, + required: ["id", "assetId"], + }, + }, +]); + +export async function listTools(): Promise { + return TOOLS.map((tool) => ({ ...tool })); +} + +export const listNeedsCredential = false; + +type FailedResult = McpCallResult & { isError: true }; + +const failure = (message: string): FailedResult => ({ + text: message, + isError: true, + truncated: false, +}); + +function result(value: unknown): McpCallResult { + const text = JSON.stringify(value, null, 2) ?? "null"; + if (text.length <= MAX_RESULT_CHARS) { + return { text, isError: false, truncated: false }; + } + return { + text: `${cutAtCodeUnits(text, MAX_RESULT_CHARS)}\n\n[truncated: the tool returned ${text.length} characters]`, + isError: false, + truncated: true, + }; +} + +function stringArg( + args: Record, + key: string, +): string | undefined { + const value = args[key]; + return typeof value === "string" && value.trim() !== "" ? value : undefined; +} + +function requiredString( + args: Record, + key: string, +): string | FailedResult { + const value = stringArg(args, key); + return value ?? failure(`A workflow call needs ${key}.`); +} + +function identityOf(connection: Connection): WorkflowIdentity | FailedResult { + const ownerUserId = connection.actorId?.trim(); + if (!ownerUserId) { + return failure( + "A workflow belongs to somebody, and this run is not attributed to anybody.", + ); + } + const agentId = connection.botId?.trim(); + if (!agentId) { + return failure( + "A workflow belongs to a Bot, and this run does not name one.", + ); + } + return { ownerUserId, agentId }; +} + +function isFailure(value: unknown): value is FailedResult { + return ( + typeof value === "object" && + value !== null && + "isError" in value && + (value as McpCallResult).isError === true + ); +} + +function stepsArg(value: unknown): WorkflowStepInput[] | FailedResult { + if (!Array.isArray(value)) return failure("A workflow needs a steps array."); + const steps: WorkflowStepInput[] = []; + for (const raw of value) { + if (!raw || typeof raw !== "object") { + return failure("Every workflow step must be an object."); + } + const item = raw as Record; + const key = typeof item.key === "string" ? item.key : undefined; + const instruction = + typeof item.instruction === "string" ? item.instruction : undefined; + if (!key || !instruction) { + return failure("Every workflow step needs key and instruction."); + } + if ( + item.dependsOn !== undefined && + (!Array.isArray(item.dependsOn) || + item.dependsOn.some((entry) => typeof entry !== "string")) + ) { + return failure( + "A workflow step dependsOn must be an array of step keys.", + ); + } + steps.push({ + key, + instruction, + ...(Array.isArray(item.dependsOn) + ? { dependsOn: item.dependsOn as string[] } + : {}), + }); + } + return steps; +} + +function absoluteFuture(value: string | undefined): Date | undefined { + if (!value || !/(?:Z|[+-]\d{2}:\d{2})$/i.test(value)) return undefined; + const parsed = new Date(value); + if (Number.isNaN(parsed.getTime()) || parsed.getTime() <= Date.now()) { + return undefined; + } + return parsed; +} + +function compactPlan(plan: WorkflowPlan | null) { + if (!plan) return null; + return { + id: plan.id, + title: plan.title, + status: plan.status, + channelId: plan.channelId, + createdAt: plan.createdAt, + updatedAt: plan.updatedAt, + steps: plan.steps.map((step) => ({ + key: step.key, + status: step.status, + attempts: step.attempts, + provider: step.provider, + waitUntil: step.waitUntil, + failureReason: step.failureReason, + dependsOn: step.dependsOn, + instruction: step.instruction, + })), + }; +} + +export async function callTool( + connection: Connection, + toolName: string, + args: Record, +): Promise { + const identity = identityOf(connection); + if (isFailure(identity)) return identity; + const tools = installed; + if (!tools) return failure("Workflows is not available in this deployment."); + + try { + if (toolName === "create_workflow") { + const channelId = requiredString(args, "channelId"); + if (isFailure(channelId)) return channelId; + const title = requiredString(args, "title"); + if (isFailure(title)) return title; + const steps = stepsArg(args.steps); + if (isFailure(steps)) return steps; + return result( + compactPlan( + await tools.create({ + ...identity, + channelId, + title, + steps, + }), + ), + ); + } + if (toolName === "list_workflows") { + return result(await tools.listFor(identity)); + } + + const id = requiredString(args, "id"); + if (isFailure(id)) return id; + + if (toolName === "get_workflow") { + const plan = await tools.get(identity, id); + if (!plan) return failure("That workflow does not exist."); + return result(compactPlan(plan)); + } + if (toolName === "pause_workflow") { + return result(compactPlan(await tools.pause(identity, id))); + } + if (toolName === "resume_workflow") { + return result(compactPlan(await tools.resume(identity, id))); + } + if (toolName === "cancel_workflow") { + return result(compactPlan(await tools.cancel(identity, id))); + } + if (toolName === "list_workflow_assets") { + return result( + await tools.listAssets(identity, id, stringArg(args, "stepKey")), + ); + } + if (toolName === "remove_workflow_asset") { + const assetId = requiredString(args, "assetId"); + if (isFailure(assetId)) return assetId; + await tools.removeAsset(identity, id, assetId); + return result({ removed: assetId }); + } + + const stepKey = requiredString(args, "stepKey"); + if (isFailure(stepKey)) return stepKey; + + if (toolName === "start_workflow_step") { + return result(await tools.startStep(identity, id, stepKey)); + } + if (toolName === "wait_workflow_step") { + const waitUntil = absoluteFuture(stringArg(args, "waitUntil")); + if (!waitUntil) { + return failure( + "waitUntil must be an absolute future RFC3339 timestamp with Z or a numeric UTC offset.", + ); + } + const provider = stringArg(args, "provider"); + return result( + await tools.waitStep(identity, id, stepKey, { + waitUntil, + ...(provider ? { provider } : {}), + }), + ); + } + if (toolName === "complete_workflow_step") { + return result( + compactPlan(await tools.completeStep(identity, id, stepKey)), + ); + } + if (toolName === "fail_workflow_step") { + const reason = requiredString(args, "reason"); + if (isFailure(reason)) return reason; + return result(await tools.failStep(identity, id, stepKey, reason)); + } + if (toolName === "retry_workflow_step") { + return result(await tools.retryStep(identity, id, stepKey)); + } + if (toolName === "add_workflow_asset") { + const direction = stringArg(args, "direction"); + const mediaKind = stringArg(args, "mediaKind"); + const ref = requiredString(args, "ref"); + if (isFailure(ref)) return ref; + if (direction !== "input" && direction !== "output") { + return failure("direction must be input or output."); + } + if ( + mediaKind !== "image" && + mediaKind !== "video" && + mediaKind !== "audio" && + mediaKind !== "file" && + mediaKind !== "text" + ) { + return failure("mediaKind must be image, video, audio, file or text."); + } + const label = stringArg(args, "label"); + const input: WorkflowAssetInput = { + direction, + mediaKind, + ref, + ...(label ? { label } : {}), + }; + return result(await tools.addAsset(identity, id, stepKey, input)); + } + return failure(`Unknown workflow tool: ${toolName}.`); + } catch (error) { + if ( + error instanceof WorkflowRefusedError || + error instanceof WorkflowNotFoundError + ) { + return failure(error.message); + } + return failure("The workflow operation could not be completed."); + } +} diff --git a/server/src/plugins/catalogue.ts b/server/src/plugins/catalogue.ts index afd265821..d575946ed 100644 --- a/server/src/plugins/catalogue.ts +++ b/server/src/plugins/catalogue.ts @@ -299,6 +299,32 @@ export const CATALOGUE: readonly CatalogueEntry[] = Object.freeze([ ]), docsUrl: "https://github.com/CopilotKit/OpenBot/blob/main/docs/routines.md", }, + { + key: "workflows", + title: "Workflows", + vendor: "OpenBot", + summary: + "Durable multi-step plans with dependency, wait, retry and per-step asset state.", + host: "builtin://workflows", + path: "/", + transport: "builtin-workflows", + auth: Object.freeze({ kind: "builtin" }), + writeTools: Object.freeze([ + "create_workflow", + "pause_workflow", + "resume_workflow", + "cancel_workflow", + "start_workflow_step", + "wait_workflow_step", + "complete_workflow_step", + "fail_workflow_step", + "retry_workflow_step", + "add_workflow_asset", + "remove_workflow_asset", + ]), + docsUrl: + "https://github.com/CopilotKit/OpenBot/blob/main/docs/workflows.md", + }, ]); const BY_KEY = new Map(CATALOGUE.map((entry) => [entry.key, entry])); diff --git a/server/src/plugins/transport.ts b/server/src/plugins/transport.ts index 60a18b333..29abb854b 100644 --- a/server/src/plugins/transport.ts +++ b/server/src/plugins/transport.ts @@ -1,4 +1,5 @@ import * as builtinRoutines from "./builtin-routines"; +import * as builtinWorkflows from "./builtin-workflows"; import * as composio from "./composio"; import * as driveRest from "./google-drive-rest"; import type { ListedTool, McpCallResult } from "./mcp"; @@ -119,6 +120,7 @@ export type TransportKind = | "mcp" | "google-drive-rest" | "builtin-routines" + | "builtin-workflows" | "composio"; /** @@ -144,6 +146,7 @@ const TRANSPORTS: Record = { mcp, "google-drive-rest": driveRest, "builtin-routines": builtinRoutines, + "builtin-workflows": builtinWorkflows, composio, }; diff --git a/server/tests/builtin-workflows.test.ts b/server/tests/builtin-workflows.test.ts new file mode 100644 index 000000000..3cf3dcec5 --- /dev/null +++ b/server/tests/builtin-workflows.test.ts @@ -0,0 +1,154 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { + callTool, + listTools, + type WorkflowTools, + useWorkflowTools, +} from "../src/plugins/builtin-workflows"; + +const CONNECTION = { + url: "builtin://workflows", + actorId: "user_owner", + botId: "bot_worker", +}; + +const now = new Date("2026-09-19T12:00:00.000Z"); +const step = { + id: "workflow_step_1", + workflowId: "workflow_1", + key: "render", + position: 0, + instruction: "Render scene one.", + dependsOn: [] as string[], + status: "ready" as const, + attempts: 0, + provider: null, + waitUntil: null, + failureReason: null, + startedAt: null, + finishedAt: null, + createdAt: now, + updatedAt: now, +}; +const plan = { + id: "workflow_1", + ownerUserId: "user_owner", + agentId: "bot_worker", + channelId: "channel_1", + title: "Video pipeline", + status: "active" as const, + finishedAt: null, + createdAt: now, + updatedAt: now, + steps: [step], +}; + +afterEach(() => useWorkflowTools(null)); + +describe("workflow tool catalogue", () => { + test("lists the durable workflow controls", async () => { + const names = (await listTools()).map((tool) => tool.name); + expect(names).toContain("create_workflow"); + expect(names).toContain("wait_workflow_step"); + expect(names).toContain("add_workflow_asset"); + expect(names).toContain("cancel_workflow"); + }); +}); + +describe("workflow identity boundary", () => { + test("create derives owner and Bot from the connection", async () => { + let received: Parameters[0] | undefined; + useWorkflowTools({ + async create(input) { + received = input; + return plan; + }, + async get() { + return plan; + }, + async listFor() { + return [plan]; + }, + async pause() { + return plan; + }, + async resume() { + return plan; + }, + async cancel() { + return plan; + }, + async startStep() { + return step; + }, + async waitStep() { + return { ...step, status: "waiting" as const }; + }, + async completeStep() { + return plan; + }, + async failStep() { + return { ...step, status: "failed" as const }; + }, + async retryStep() { + return step; + }, + async addAsset() { + return { + id: "asset_1", + workflowId: plan.id, + stepKey: "render", + direction: "input", + mediaKind: "image", + ref: "workspace:scene/reference.png", + label: null, + createdAt: now, + }; + }, + async listAssets() { + return []; + }, + async removeAsset() {}, + }); + + const result = await callTool(CONNECTION, "create_workflow", { + ownerUserId: "user_attacker", + agentId: "bot_attacker", + channelId: "channel_1", + title: "Video pipeline", + steps: [{ key: "render", instruction: "Render scene one." }], + }); + + expect(result.isError).toBe(false); + expect(received).toMatchObject({ + ownerUserId: "user_owner", + agentId: "bot_worker", + channelId: "channel_1", + title: "Video pipeline", + }); + }); + + test("refuses calls without connection identity", async () => { + useWorkflowTools({} as WorkflowTools); + const result = await callTool( + { url: "builtin://workflows" }, + "list_workflows", + {}, + ); + expect(result.isError).toBe(true); + expect(result.text).toContain("attributed"); + }); +}); + +describe("workflow wait input", () => { + test("refuses a local timestamp with no explicit offset", async () => { + useWorkflowTools({} as WorkflowTools); + const result = await callTool(CONNECTION, "wait_workflow_step", { + id: "workflow_1", + stepKey: "render", + waitUntil: "2026-09-20T14:30:00", + }); + expect(result.isError).toBe(true); + expect(result.text).toContain("RFC3339"); + }); +}); diff --git a/server/tests/plugin-catalogue.test.ts b/server/tests/plugin-catalogue.test.ts index 7aa8c49f7..4f5c1df34 100644 --- a/server/tests/plugin-catalogue.test.ts +++ b/server/tests/plugin-catalogue.test.ts @@ -105,7 +105,7 @@ describe("which servers this deployment will talk to", () => { // First-party and in-process: there is no host outside this process to reach, so the // https requirement below does not apply. Asserted positively instead, so this branch // cannot quietly become a loophole for a future entry that DOES dial a real host. - expect(entry.host).toBe("builtin://routines"); + expect(entry.host).toBe(`builtin://${entry.key}`); } else { expect(entry.host.startsWith("https://")).toBe(true); }