Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion docs/workflows.md
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,9 @@ because they share an owner.

The destination channel is accepted only when the owner is a member and the Bot is in that same
non-deleted channel. A missing, foreign or deleted channel is one refusal and does not reveal which
part of the check failed.
part of the check failed. The create tool may omit `channelId`: when the owner and Bot share exactly
one live channel the store resolves it server-side; when they share several, creation is refused with
the available channel choices so the Agent can ask the person instead of inventing an opaque id.

Deleting an Agent cascades its workflow rows and steps, so a pending plan cannot survive as stale work
for an identity that no longer exists.
Expand Down
2 changes: 2 additions & 0 deletions scripts/release-preflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1646,6 +1646,8 @@ function checkDurableWorkflowState(): void {
"attachmentVisibleToWorkflow(",
"eq(attachments.channelId, run.channelId)",
"eq(channelAgents.agentId, identity.agentId)",
"resolveChannel(input)",
"more than one channel with me",
"addAsset(",
"listAssets(",
"removeAsset(",
Expand Down
13 changes: 8 additions & 5 deletions server/src/plugins/builtin-workflows.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,11 @@ const TOOLS: readonly McpTool[] = Object.freeze([
inputSchema: {
type: "object",
properties: {
channelId: { type: "string" },
channelId: {
type: "string",
description:
"Optional destination channel. Omit when you and the person share only one channel; if there are several, the tool will tell you which choices exist.",
},
title: { type: "string" },
steps: {
type: "array",
Expand All @@ -64,7 +68,7 @@ const TOOLS: readonly McpTool[] = Object.freeze([
},
},
},
required: ["channelId", "title", "steps"],
required: ["title", "steps"],
},
},
{
Expand Down Expand Up @@ -377,8 +381,7 @@ export async function callTool(

try {
if (toolName === "create_workflow") {
const channelId = requiredString(args, "channelId");
if (isFailure(channelId)) return channelId;
const channelId = stringArg(args, "channelId");
const title = requiredString(args, "title");
if (isFailure(title)) return title;
const steps = stepsArg(args.steps);
Expand All @@ -387,7 +390,7 @@ export async function callTool(
compactPlan(
await tools.create({
...identity,
channelId,
...(channelId ? { channelId } : {}),
title,
steps,
}),
Expand Down
64 changes: 54 additions & 10 deletions server/src/workflows/store.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import {
and,
asc,
desc,
eq,
inArray,
isNotNull,
Expand Down Expand Up @@ -63,7 +64,7 @@ export type WorkflowStepInput = {
};

export type WorkflowInput = WorkflowIdentity & {
channelId: string;
channelId?: string;
title: string;
steps: WorkflowStepInput[];
};
Expand Down Expand Up @@ -383,9 +384,39 @@ export function createWorkflowStore(database: Database): WorkflowStore {
);
}

async function channelAllowed(input: WorkflowInput): Promise<void> {
const [row] = await database
.select({ id: channels.id })
async function resolveChannel(input: WorkflowInput): Promise<string> {
if (input.channelId !== undefined) {
const [row] = await database
.select({ id: channels.id })
.from(channels)
.innerJoin(
channelMemberships,
and(
eq(channelMemberships.channelId, channels.id),
eq(channelMemberships.userId, input.ownerUserId),
),
)
.innerJoin(
channelAgents,
and(
eq(channelAgents.channelId, channels.id),
eq(channelAgents.agentId, input.agentId),
),
)
.where(
and(eq(channels.id, input.channelId), isNull(channels.deletedAt)),
)
.limit(1);
if (!row) {
throw new WorkflowRefusedError(
"That channel is not shared by this person and this Bot.",
);
}
return row.id;
}

const candidates = await database
.select({ id: channels.id, name: channels.name })
.from(channels)
.innerJoin(
channelMemberships,
Expand All @@ -401,13 +432,26 @@ export function createWorkflowStore(database: Database): WorkflowStore {
eq(channelAgents.agentId, input.agentId),
),
)
.where(and(eq(channels.id, input.channelId), isNull(channels.deletedAt)))
.limit(1);
if (!row) {
.where(isNull(channels.deletedAt))
.orderBy(desc(channels.createdAt), desc(channels.id))
.limit(6);

const only = candidates[0];
if (!only) {
throw new WorkflowRefusedError(
"We have no channel for this workflow. Start one and ask again.",
);
}
if (candidates.length > 1) {
const labels = candidates
.slice(0, 5)
.map((candidate) => `${candidate.name} (${candidate.id})`);
if (candidates.length > 5) labels.push("and others");
throw new WorkflowRefusedError(
"That channel is not shared by this person and this Bot.",
`You are in more than one channel with me — ${labels.join(", ")}. Say which one.`,
);
}
return only.id;
}

function preparedSteps(input: WorkflowStepInput[]) {
Expand Down Expand Up @@ -616,15 +660,15 @@ export function createWorkflowStore(database: Database): WorkflowStore {
MAX_WORKFLOW_TITLE_CODE_POINTS,
);
const steps = preparedSteps(input.steps);
await channelAllowed(input);
const channelId = await resolveChannel(input);

const id = `workflow_${crypto.randomUUID()}`;
await database.transaction(async (transaction) => {
await transaction.insert(workflowRuns).values({
id,
ownerUserId: input.ownerUserId,
agentId: input.agentId,
channelId: input.channelId,
channelId,
title,
});
await transaction.insert(workflowSteps).values(
Expand Down
23 changes: 23 additions & 0 deletions server/tests/builtin-workflows.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,29 @@ describe("workflow identity boundary", () => {
});
});

test("lets create omit the opaque channel id so the store can resolve it", async () => {
let received: Parameters<WorkflowTools["create"]>[0] | undefined;
useWorkflowTools({
async create(input) {
received = input;
return plan;
},
} as WorkflowTools);

const result = await callTool(CONNECTION, "create_workflow", {
title: "Video pipeline",
steps: [{ key: "render", instruction: "Render scene one." }],
});

expect(result.isError).toBe(false);
expect(received).toMatchObject({
ownerUserId: "user_owner",
agentId: "bot_worker",
title: "Video pipeline",
});
expect(received).not.toHaveProperty("channelId");
});

test("refuses calls without connection identity", async () => {
useWorkflowTools({} as WorkflowTools);
const result = await callTool(
Expand Down
23 changes: 23 additions & 0 deletions server/tests/workflow-store.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,29 @@ describe("durable workflow creation", () => {
expect(plan.steps[1]?.dependsOn).toEqual(["script"]);
});

test("resolves the only shared channel when create omits channelId", async () => {
const { owner, agentId, channel } = await setUp();
const input = planInput(owner, agentId, channel.id);
const { channelId: _channelId, ...withoutChannel } = input;

const plan = await store.create(withoutChannel);

expect(plan.channelId).toBe(channel.id);
});

test("refuses an omitted channelId when more than one shared channel exists", async () => {
const { owner, agentId, channel } = await setUp();
const second = await createChannel(owner, [agentId]);
const input = planInput(owner, agentId, channel.id);
const { channelId: _channelId, ...withoutChannel } = input;

await expect(store.create(withoutChannel)).rejects.toThrow(
/more than one channel/,
);

expect(second.id).not.toBe(channel.id);
});

test("refuses forward, missing and duplicate dependencies", async () => {
const { owner, agentId, channel } = await setUp();

Expand Down
Loading