From db95d619c2f78f7bb558b12bd3f7b8b72d133cf2 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sat, 12 Sep 2026 20:54:56 +0200 Subject: [PATCH] feat(sdk): webhook receiver loaded-flow admission (#303) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds an explicit --allow allowlist to `flows serve-webhook`. When set, the receiver refuses any name whose trigger wasn't advertised at boot with `webhook_flow_unknown` (404) before reading the body — so a rogue POST can't accumulate events in an inbox the daemon will never drain. This is the admission half of #303. Durable authored-handler execution against the journal is deferred: per #303's own scoping note, that half depends on kernel authored-handler registration + resume protocol which sits outside the SDK. - parseWebhookArgs learns `--allow name1,name2`. Empty list and invalid names refuse at parse time so misconfiguration surfaces as a CLI error rather than a mute receiver. - startWebhookServer takes `options.admittedNames?: ReadonlySet`. When set, non-admitted names return `{ error: 'webhook_flow_unknown', name }` before any body read, bounded by the same limits as invalid names. - runServeWebhook echoes `ADMITTED ` after `WEBHOOK http://...` so operators can confirm the loaded set from the receiver's stdout. - Two new tests: parse-strict on --allow, and end-to-end admission (release/push admitted, rogue refused, unadmitted names never write). Co-Authored-By: Claude Opus 4.7 (1M context) Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 --- packages/sdk/src/cli.ts | 4 +-- packages/sdk/src/cli/serve-webhook.ts | 52 ++++++++++++++++++++++----- packages/sdk/tests/webhook.test.ts | 23 +++++++++++- 3 files changed, 68 insertions(+), 11 deletions(-) diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index f5337e970..51e37ed42 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -49,7 +49,7 @@ type ParsedArgs = | ReplayArgs | BuildArgs | DeployArgs - | { command: 'serve-webhook'; dataDir: string; port: number } + | { command: 'serve-webhook'; dataDir: string; port: number; admitted?: readonly string[] } | { command: 'cloud-run'; value: string; json: boolean; wait: boolean } | { command: 'check'; json: boolean; watch: boolean; value: string } | { command: 'run'; bucket: string | undefined; reuseFromRunId: string | undefined; localAgent: boolean; dataDir: string; input: string | undefined; json: boolean; spawn: boolean; noObserverLink: boolean; allowHumanInfluenced: boolean; value: string } @@ -69,7 +69,7 @@ const USAGE = [ 'flows deploy @sha256: --to ', 'flows run @sha256: [--bucket ] [--data-dir ] [--json]', 'flows check [--watch] [--json] ', - 'flows serve-webhook --data-dir --port

', + 'flows serve-webhook --data-dir

--port

[--allow [,]]', 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir

] [--local-agent] [--reuse-from ] ', 'flows run --cloud [--json] [--wait] ', 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir ] [--local-agent] --input ', diff --git a/packages/sdk/src/cli/serve-webhook.ts b/packages/sdk/src/cli/serve-webhook.ts index 94141ad28..5d24c945e 100644 --- a/packages/sdk/src/cli/serve-webhook.ts +++ b/packages/sdk/src/cli/serve-webhook.ts @@ -10,13 +10,13 @@ const MAX_BODY_BYTES = 1024 * 1024; const NAME = /^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$/; export function parseWebhookArgs(args: readonly string[]): { - command: 'serve-webhook'; dataDir: string; port: number; + command: 'serve-webhook'; dataDir: string; port: number; admitted?: readonly string[]; } | undefined { const values = new Map(); for (let i = 0; i < args.length; i += 2) { const flag = args[i]!; const value = args[i + 1]; - if (!['--data-dir', '--port'].includes(flag) || values.has(flag) + if (!['--data-dir', '--port', '--allow'].includes(flag) || values.has(flag) || !value || value.startsWith('-')) return undefined; values.set(flag, value); } @@ -25,7 +25,17 @@ export function parseWebhookArgs(args: readonly string[]): { if (!dataDir || !portText || !/^\d+$/.test(portText)) return undefined; const port = Number(portText); if (!Number.isInteger(port) || port < 0 || port > 65535) return undefined; - return { command: 'serve-webhook', dataDir, port }; + const allow = values.get('--allow'); + // --allow is comma-separated. Empty list is not admitted (would refuse everything); + // treat as parse failure so the author sees the typo rather than a mute receiver. + const admitted = allow === undefined ? undefined + : allow.split(',').map(entry => entry.trim()); + if (admitted !== undefined && (admitted.length === 0 || admitted.some(entry => !NAME.test(entry)))) { + return undefined; + } + return admitted === undefined + ? { command: 'serve-webhook', dataDir, port } + : { command: 'serve-webhook', dataDir, port, admitted }; } function reply(response: ServerResponse, status: number, body: object): void { @@ -33,10 +43,22 @@ function reply(response: ServerResponse, status: number, body: object): void { response.end(JSON.stringify(body)); } -/** POST / accepts JSON; /providers/ accepts typed event envelopes. */ -export async function startWebhookServer(dataDir: string, port: number): Promise { +/** + * POST / accepts JSON; /providers/ accepts typed event envelopes. + * + * When `admittedNames` is set, unknown names (or providers) refuse with + * `webhook_flow_unknown`. This closes the "any-name" ingress opened by slice E + * (#333) so the receiver only accepts inbox writes for triggers whose flows + * have been explicitly loaded — the loaded-flow admission half of #303. + * (Durable handler execution against the journal is deferred; per #303 that + * half depends on kernel authored-handler registration + resume protocol.) + */ +export async function startWebhookServer( + dataDir: string, port: number, options: { admittedNames?: ReadonlySet } = {}, +): Promise { const inbox = join(resolve(dataDir), 'inbox'); await directory(inbox); + const admitted = options.admittedNames; // TODO https://github.com/AgentWorkforce/flows/issues/301: provider signatures, // public ingress and Cloud mount provisioning belong to the deployment slice. const server = createServer(async (request, response) => { @@ -53,6 +75,15 @@ export async function startWebhookServer(dataDir: string, port: number): Promise reply(response, 404, { error: 'invalid_webhook_name' }); return; } + if (admitted !== undefined && !admitted.has(name)) { + // Loaded-flow admission (#303): the receiver refuses any name whose + // trigger wasn't advertised at boot, so a rogue POST can't accumulate + // events in an inbox the daemon will never drain. The check runs before + // any body read so an unadmitted caller can't waste MAX_BODY_BYTES. + request.resume(); + reply(response, 404, { error: 'webhook_flow_unknown', name }); + return; + } let temporary: string | undefined; try { let bytes = 0; @@ -125,12 +156,17 @@ async function directory(path: string): Promise { } export async function runServeWebhook( - options: { dataDir: string; port: number }, io: CliIo, + options: { dataDir: string; port: number; admitted?: readonly string[] }, io: CliIo, ): Promise<0 | 1> { try { - const server = await startWebhookServer(options.dataDir, options.port); + const admittedNames = options.admitted === undefined ? undefined : new Set(options.admitted); + const server = await startWebhookServer(options.dataDir, options.port, { admittedNames }); const address = server.address(); - io.stdout(`WEBHOOK http://127.0.0.1:${typeof address === 'object' && address ? address.port : options.port}`); + const port = typeof address === 'object' && address ? address.port : options.port; + io.stdout(`WEBHOOK http://127.0.0.1:${port}`); + if (admittedNames !== undefined) { + io.stdout(`ADMITTED ${[...admittedNames].sort().join(',')}`); + } await new Promise((accept, reject) => { const stop = (): void => { server.close(error => error ? reject(error) : accept()); diff --git a/packages/sdk/tests/webhook.test.ts b/packages/sdk/tests/webhook.test.ts index be0073189..9cedd1f68 100644 --- a/packages/sdk/tests/webhook.test.ts +++ b/packages/sdk/tests/webhook.test.ts @@ -124,9 +124,30 @@ describe('webhook ingress', () => { }); it('parses CLI options strictly', () => { expect(parseWebhookArgs(['--data-dir', 'data', '--port', '0'])).toEqual({ command: 'serve-webhook', dataDir: 'data', port: 0 }); + expect(parseWebhookArgs(['--data-dir', 'data', '--port', '0', '--allow', 'release,push'])) + .toEqual({ command: 'serve-webhook', dataDir: 'data', port: 0, admitted: ['release', 'push'] }); for (const args of [[], ['--port', '80'], ['--data-dir', 'x', '--port', '65536'], - ['--data-dir', 'x', '--port', '1', '--port', '2'], ['--data-dir', 'x', '--port', '3.1']]) { + ['--data-dir', 'x', '--port', '1', '--port', '2'], ['--data-dir', 'x', '--port', '3.1'], + // #303 admission: empty --allow list and invalid names refuse at parse time + ['--data-dir', 'x', '--port', '0', '--allow', ''], + ['--data-dir', 'x', '--port', '0', '--allow', '../oops'], + ['--data-dir', 'x', '--port', '0', '--allow', 'ok,../oops']]) { expect(parseWebhookArgs(args)).toBeUndefined(); } }); + it('admits only loaded flow trigger names when the allowlist is set (#303)', async () => { + const dir = await temporary(); + const server = await startWebhookServer(dir, 0, { admittedNames: new Set(['release', 'push']) }); + servers.push(server); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('missing HTTP address'); + const base = `http://127.0.0.1:${address.port}`; + expect((await fetch(`${base}/release`, { method: 'POST', body: '{"ok":true}' })).status).toBe(202); + expect((await fetch(`${base}/push`, { method: 'POST', body: '{"ok":true}' })).status).toBe(202); + const unknown = await fetch(`${base}/rogue`, { method: 'POST', body: '{"ok":true}' }); + expect(unknown.status).toBe(404); + expect(await unknown.json()).toEqual({ error: 'webhook_flow_unknown', name: 'rogue' }); + // Unadmitted names never write to disk — the daemon can't drain what wasn't accumulated. + expect(await readdir(join(dir, 'inbox'))).toEqual(['push', 'release']); + }); });