Skip to content
Open
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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ agent-slack
│ ├── list <channel> # workflows bookmarked in a channel
│ ├── preview <trigger-id> # trigger metadata (no side effects)
│ ├── get <id> # workflow definition + form fields
│ └── run <trigger-id> # trip a workflow trigger
│ └── run <trigger-id> # trip a trigger (--field for form submission)
└── canvas
├── create # markdown file/blob → canvas
└── get <canvas-url-or-id> # canvas → markdown
Expand Down
2 changes: 2 additions & 0 deletions skills/agent-slack/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,8 @@ Ordinary `message send` and `message edit` calls auto-convert lists. `message se

Slack-native drafts (`message draft list|create|update|delete`) manage drafts that appear in the user's Slack client; `create` posts nothing. They use undocumented session endpoints and require browser-style auth (xoxc/xoxd).

`workflow run <trigger>` trips a trigger. Add `--field "Title=Value"` (repeatable) to submit a workflow's form when one is presented; field titles are validated against the workflow schema. Form submission requires browser-style auth (xoxc/xoxd).

## Conditional references

- Read [references/targets.md](references/targets.md) only when choosing between a message URL, channel, or user target, or when resolving multiple workspaces.
Expand Down
88 changes: 76 additions & 12 deletions src/cli/workflow-command.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,21 @@ import {
resolveShortcutUrl,
runWorkflow,
} from "../slack/workflows.ts";
import {
requireBrowserAuth,
submitWorkflow,
validateFieldInputs,
} from "../slack/workflow-submit.ts";

type WorkspaceOption = {
workspace?: string;
};

type RunOptions = WorkspaceOption & {
channel: string;
field?: string[];
};

export function registerWorkflowCommand(input: { program: Command; ctx: CliContext }): void {
const workflowCmd = input.program
.command("workflow")
Expand Down Expand Up @@ -105,27 +115,81 @@ export function registerWorkflowCommand(input: { program: Command; ctx: CliConte

workflowCmd
.command("run")
.description("Trip a workflow trigger")
.description("Trip a workflow trigger (with --field, submits form data)")
.argument("<trigger-id>", "Trigger ID (Ft...)")
.requiredOption("--channel <id-or-name>", "Channel where the workflow is bookmarked")
.option(
"--field <title=value>",
"Form field value (repeatable)",
(v, prev: string[]) => {
prev.push(v);
return prev;
},
[] as string[],
)
.option(
"--workspace <url>",
"Workspace selector (full URL or unique substring; required if you have multiple workspaces)",
)
.action(async (...args) => {
const [triggerId, options] = args as [string, WorkspaceOption & { channel: string }];
const [triggerId, options] = args as [string, RunOptions];
try {
const workspaceUrl = input.ctx.effectiveWorkspaceUrl(options.workspace);
const payload = await input.ctx.withAutoRefresh({
workspaceUrl,
work: async () => {
const { client } = await input.ctx.getClientForWorkspace(workspaceUrl);
const channelId = await resolveChannelId(client, options.channel);
const shortcutUrl = await resolveShortcutUrl(client, { channelId, triggerId });
return await runWorkflow(client, { shortcutUrl, channelId, triggerId });
},
});
console.log(JSON.stringify(pruneEmpty(payload), null, 2));
const fieldArgs = options.field ?? [];

if (fieldArgs.length === 0) {
// Trip-only (existing behavior, no WebSocket)
const payload = await input.ctx.withAutoRefresh({
workspaceUrl,
work: async () => {
const { client } = await input.ctx.getClientForWorkspace(workspaceUrl);
const channelId = await resolveChannelId(client, options.channel);
const shortcutUrl = await resolveShortcutUrl(client, { channelId, triggerId });
return await runWorkflow(client, { shortcutUrl, channelId, triggerId });
},
});
console.log(JSON.stringify(pruneEmpty(payload), null, 2));
} else {
// Parse --field args
const fields = new Map<string, string>();
for (const arg of fieldArgs) {
const eqIdx = arg.indexOf("=");
if (eqIdx < 1) {
throw new Error(`Invalid --field format: "${arg}". Expected Title=value`);
}
fields.set(arg.substring(0, eqIdx), arg.substring(eqIdx + 1));
}

const payload = await input.ctx.withAutoRefresh({
workspaceUrl,
work: async () => {
const { client, auth } = await input.ctx.getClientForWorkspace(workspaceUrl);
requireBrowserAuth(auth);

const channelId = await resolveChannelId(client, options.channel);

// Preview → schema → validate before opening WebSocket
const preview = await previewWorkflow(client, triggerId);
const schema = await getWorkflowSchema(client, preview.workflow.id);
const errors = validateFieldInputs(fields, schema);
if (errors.length > 0) {
throw new Error(errors.join("\n"));
}

const shortcutUrl = await resolveShortcutUrl(client, { channelId, triggerId });
return await submitWorkflow({
client,
auth,
shortcutUrl,
channelId,
triggerId,
fields,
schema,
});
},
});
console.log(JSON.stringify(pruneEmpty(payload), null, 2));
}
} catch (err: unknown) {
console.error(input.ctx.errorMessage(err));
process.exitCode = 1;
Expand Down
208 changes: 208 additions & 0 deletions src/lib/rtm-websocket.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,208 @@
/**
* Short-lived RTM WebSocket connection using node:https upgrade.
* Bun delivers HTTP 101 via the response callback (not the 'upgrade' event),
* so we read WebSocket frames directly from the response stream.
*/
import { request } from "node:https";

export type RtmConnection = {
waitForMessage: (
predicate: (msg: Record<string, unknown>) => boolean,
timeoutMs?: number,
) => Promise<Record<string, unknown>>;
close: () => void;
};

export function connectRtm(input: {
wsUrl: string;
cookie: string;
timeoutMs?: number;
}): Promise<RtmConnection> {
const { wsUrl, cookie, timeoutMs = 5000 } = input;
const url = new URL(wsUrl.replace("wss://", "https://"));
const keyBytes = new Uint8Array(16);
crypto.getRandomValues(keyBytes);
const wsKey = btoa(String.fromCharCode(...keyBytes));

const messages: Record<string, unknown>[] = [];
const listeners: ((msg: Record<string, unknown>) => void)[] = [];
let buffer: Buffer = Buffer.alloc(0);
let req: ReturnType<typeof request> | null = null;
let res: { destroy: () => void } | null = null;
let closed = false;

function onFrame(text: string): void {
try {
const parsed: unknown = JSON.parse(text);
if (typeof parsed === "object" && parsed !== null) {
const msg = parsed as Record<string, unknown>;
messages.push(msg);
for (const listener of listeners) {
listener(msg);
}
}
} catch {
// ignore non-JSON frames
}
}

function close(): void {
if (closed) {
return;
}
closed = true;
try {
res?.destroy();
} catch {
// ignore
}
try {
req?.destroy();
} catch {
// ignore
}
}

return new Promise<RtmConnection>((resolve, reject) => {
const timeout = setTimeout(() => {
close();
reject(new Error("RTM WebSocket connection timed out"));
}, timeoutMs);

req = request(
{
hostname: url.hostname,
port: 443,
path: url.pathname + url.search,
method: "GET",
headers: {
Upgrade: "websocket",
Connection: "Upgrade",
Cookie: cookie,
Origin: "https://app.slack.com",
"Sec-WebSocket-Key": wsKey,
"Sec-WebSocket-Version": "13",
},
},
(response) => {
res = response;
if (response.statusCode !== 101) {
clearTimeout(timeout);
close();
reject(new Error(`RTM WebSocket connection failed: HTTP ${response.statusCode}`));
return;
}

response.on("data", (chunk: Buffer) => {
buffer = Buffer.concat([buffer, chunk]);
const { texts, remaining } = extractTextFrames(buffer);
buffer = remaining;
for (const text of texts) {
onFrame(text);
}
});

// Wait briefly for the hello message before resolving
setTimeout(() => {
clearTimeout(timeout);
resolve({ waitForMessage, close });
}, 300);
},
);

req.on("error", (err: Error) => {
clearTimeout(timeout);
close();
reject(new Error(`RTM WebSocket connection failed: ${err.message}`));
});

req.end();
});

function waitForMessage(
predicate: (msg: Record<string, unknown>) => boolean,
waitTimeoutMs = 15000,
): Promise<Record<string, unknown>> {
// Check already-received messages first
for (const msg of messages) {
if (predicate(msg)) {
return Promise.resolve(msg);
}
}

return new Promise<Record<string, unknown>>((resolve, reject) => {
const timer = setTimeout(() => {
const idx = listeners.indexOf(listener);
if (idx >= 0) {
listeners.splice(idx, 1);
}
reject(new Error(`Timed out waiting for message (${waitTimeoutMs / 1000}s)`));
}, waitTimeoutMs);

function listener(msg: Record<string, unknown>): void {
if (predicate(msg)) {
clearTimeout(timer);
const idx = listeners.indexOf(listener);
if (idx >= 0) {
listeners.splice(idx, 1);
}
resolve(msg);
}
}

listeners.push(listener);
});
}
}

/**
* Parse WebSocket text frames from a raw byte buffer.
* Server→client frames are unmasked per RFC 6455.
*/
export function extractTextFrames(buf: Buffer): {
texts: string[];
remaining: Buffer;
} {
const texts: string[] = [];
let pos = 0;

while (pos < buf.length) {
if (buf.length - pos < 2) {
break;
}
const opcode = buf[pos]! & 0x0f;
const masked = (buf[pos + 1]! & 0x80) !== 0;
let payloadLen = buf[pos + 1]! & 0x7f;
let headerLen = 2;

if (payloadLen === 126) {
if (buf.length - pos < 4) {
break;
}
payloadLen = buf.readUInt16BE(pos + 2);
headerLen = 4;
} else if (payloadLen === 127) {
if (buf.length - pos < 10) {
break;
}
payloadLen = Number(buf.readBigUInt64BE(pos + 2));
headerLen = 10;
}
if (masked) {
headerLen += 4;
}

const totalLen = headerLen + payloadLen;
if (buf.length - pos < totalLen) {
break;
}

if (opcode === 1) {
const payload = buf.subarray(pos + headerLen, pos + totalLen);
texts.push(payload.toString("utf8"));
}
pos += totalLen;
}

return { texts, remaining: buf.subarray(pos) };
}
Loading
Loading