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: 2 additions & 2 deletions packages/sdk/src/cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand All @@ -69,7 +69,7 @@ const USAGE = [
'flows deploy <flow>@sha256:<digest> --to <file-bucket-uri>',
'flows run <flow>@sha256:<digest> [--bucket <file-bucket-uri>] [--data-dir <dir>] [--json]',
'flows check [--watch] [--json] <flow.ts|flow.yaml|spec.json>',
'flows serve-webhook --data-dir <dir> --port <p>',
'flows serve-webhook --data-dir <dir> --port <p> [--allow <name>[,<name>]]',
'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir <dir>] [--local-agent] [--reuse-from <run-id>] <flow.yaml|spec.json>',
'flows run --cloud [--json] [--wait] <flow.yaml|spec.json>',
'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir <dir>] [--local-agent] <flow.ts> --input <inline-json-or-file>',
Expand Down
52 changes: 44 additions & 8 deletions packages/sdk/src/cli/serve-webhook.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, string>();
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);
}
Expand All @@ -25,18 +25,40 @@ 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 {
response.writeHead(status, { 'content-type': 'application/json' });
response.end(JSON.stringify(body));
}

/** POST /<name> accepts JSON; /providers/<provider> accepts typed event envelopes. */
export async function startWebhookServer(dataDir: string, port: number): Promise<Server> {
/**
* POST /<name> accepts JSON; /providers/<provider> 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<string> } = {},
): Promise<Server> {
const inbox = join(resolve(dataDir), 'inbox');
await directory(inbox);
const admitted = options.admittedNames;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: When the caller mutates the Set passed to startWebhookServer after startup, the receiver admits newly added names or rejects removed ones. Snapshot options.admittedNames at startup so the loaded-flow allowlist remains fixed for the server lifetime.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At packages/sdk/src/cli/serve-webhook.ts, line 61:

<comment>When the caller mutates the `Set` passed to `startWebhookServer` after startup, the receiver admits newly added names or rejects removed ones. Snapshot `options.admittedNames` at startup so the loaded-flow allowlist remains fixed for the server lifetime.</comment>

<file context>
@@ -25,18 +25,40 @@ export function parseWebhookArgs(args: readonly string[]): {
+): Promise<Server> {
   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.
</file context>
Suggested change
const admitted = options.admittedNames;
const admitted = options.admittedNames === undefined ? undefined : new Set(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) => {
Expand All @@ -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;
Expand Down Expand Up @@ -125,12 +156,17 @@ async function directory(path: string): Promise<void> {
}

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<void>((accept, reject) => {
const stop = (): void => {
server.close(error => error ? reject(error) : accept());
Expand Down
23 changes: 22 additions & 1 deletion packages/sdk/tests/webhook.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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']);
});
});
Loading