diff --git a/docs/evidence/spec-X/README.md b/docs/evidence/spec-X/README.md new file mode 100644 index 000000000..c5e53a456 --- /dev/null +++ b/docs/evidence/spec-X/README.md @@ -0,0 +1,35 @@ +# Slice X verification + +Implemented mapping-driven Slack/GitHub declarations, provider envelope ingress, +and lowering to the existing webhook inbox executor. The generator also accepts +an adapter checkout. No kernel production code changed. + +Commands and captured outputs: + +- [Surface typecheck](surface-typecheck.txt): exit 0. +- [Surface regression typecheck and generated helper check](surface-regressions.txt): exit 0. +- [SDK source and type contracts](sdk-typecheck.txt): exit 0. +- [SDK test typecheck](sdk-test-types.txt): exit 0. +- [Surface suite](surface-tests.txt): 28 passed. +- [Provider executor and codegen tests](provider-tests.txt): 7 passed. These + exercise the compiled subscriptions through the real kernel CLI, including + provider/type/payload nonmatches, distinct event IDs, and durable deduplication. +- [Existing inbox watcher tests](inbox-watcher.txt): 3 passed. +- [Full SDK suite attempt](sdk.txt): 185 failed, 1109 passed, 10 skipped, + 14 errors. This is **not a green full-suite result**. The provider HTTP tests + fail at socket creation with `listen EPERM: operation not permitted 127.0.0.1`. + The full transcript also contains Unix socket permission failures, a filesystem + watch `EMFILE`, and a wrapper identification timeout. The SDK transcript excerpt + links the complete local log. The additional provider kernel CLI tests were + completed separately after this full-suite attempt. + +Local dependencies were installed from the npm cache for the SDK. The surface's +pre-existing npm lockfile omits its relay-helpers peer dependency; its installed +dependencies were copied from the local helpers worktree, with `ai-hist` copied +from the SDK installation. The SDK used this worktree's built surface through a +local node_modules link. These setup changes do not modify tracked lockfiles. + +One surface regression attempt hit a transient `ENOSPC` while creating a temp +directory; the captured rerun passed. Socket restrictions remain unresolved: +this session cannot request execution outside the sandbox. The HTTP/daemon tests +are retained and must be rerun in an environment that permits local sockets. diff --git a/docs/evidence/spec-X/inbox-watcher.txt b/docs/evidence/spec-X/inbox-watcher.txt new file mode 100644 index 000000000..2edfa024f --- /dev/null +++ b/docs/evidence/spec-X/inbox-watcher.txt @@ -0,0 +1,23 @@ +Command: PATH=/Users/khaliqgant/.cargo/bin:$PATH sh ../ops/cargo.sh test --offline -p relayflowd --test trigger_watcher +Working directory: /Users/khaliqgant/fl-slice-X/kernel + + Compiling bitflags v2.13.1 + Compiling getrandom v0.4.3 + Compiling rustix v1.1.4 + Compiling errno v0.3.14 + Compiling fastrand v2.5.0 + Compiling rusqlite v0.37.0 + Compiling relayflowd-journal v0.1.0 (/Users/khaliqgant/fl-slice-X/kernel/relayflowd-journal) + Compiling tempfile v3.27.0 + Compiling relayflowd v0.1.0 (/Users/khaliqgant/fl-slice-X/kernel/relayflowd) + Finished `test` profile [unoptimized + debuginfo] target(s) in 13.38s + Running tests/trigger_watcher.rs (/Users/khaliqgant/.relayflows-toolchain/target/1463204285/debug/deps/trigger_watcher-dcffffe2d25ba6f1) + +running 3 tests +test retains_bad_and_unregistered_events_while_consuming_filter_nonmatches ... ok +test failed_archive_retries_the_same_durable_run ... ok +test journals_payload_and_filename_key_then_archives_and_dedupes_replay ... ok + +test result: ok. 3 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 0.08s + +Exit code: 0 diff --git a/docs/evidence/spec-X/provider-tests.txt b/docs/evidence/spec-X/provider-tests.txt new file mode 100644 index 000000000..c3d171344 --- /dev/null +++ b/docs/evidence/spec-X/provider-tests.txt @@ -0,0 +1,15 @@ +Command: RELAYFLOWD_BIN=/Users/khaliqgant/.relayflows-toolchain/target/1463204285/debug/relayflowd ./node_modules/.bin/vitest run tests/provider-trigger-executor.test.ts tests/generate-triggers.test.ts +Working directory: /Users/khaliqgant/fl-slice-X/packages/sdk + + RUN v2.1.9 /Users/khaliqgant/fl-slice-X/packages/sdk + + ✓ tests/provider-trigger-executor.test.ts (4 tests) 138ms + ✓ tests/generate-triggers.test.ts (3 tests) 645ms + ✓ discovers new adapters, preserves exact event names, and prefers adapter-local mappings 369ms + + Test Files 2 passed (2) + Tests 7 passed (7) + Start at 23:55:10 + Duration 881ms (transform 82ms, setup 0ms, collect 128ms, tests 784ms, environment 0ms, prepare 69ms) + +Exit code: 0 diff --git a/docs/evidence/spec-X/sdk-test-types.txt b/docs/evidence/spec-X/sdk-test-types.txt new file mode 100644 index 000000000..38c3ef79c --- /dev/null +++ b/docs/evidence/spec-X/sdk-test-types.txt @@ -0,0 +1,9 @@ +Command: npm run typecheck:tests +Working directory: /Users/khaliqgant/fl-slice-X/packages/sdk + + +> @relayflows/sdk@2.0.8 typecheck:tests +> tsc -p tsconfig.tests.json + + +Exit code: 0 diff --git a/docs/evidence/spec-X/sdk-typecheck.txt b/docs/evidence/spec-X/sdk-typecheck.txt new file mode 100644 index 000000000..37703fa78 --- /dev/null +++ b/docs/evidence/spec-X/sdk-typecheck.txt @@ -0,0 +1,9 @@ +Command: npm run typecheck +Working directory: /Users/khaliqgant/fl-slice-X/packages/sdk + + +> @relayflows/sdk@2.0.8 typecheck +> tsc --noEmit && tsc -p tsconfig.type-tests.json + + +Exit code: 0 diff --git a/docs/evidence/spec-X/sdk.txt b/docs/evidence/spec-X/sdk.txt new file mode 100644 index 000000000..dd0d6a579 --- /dev/null +++ b/docs/evidence/spec-X/sdk.txt @@ -0,0 +1,99 @@ +Command: PATH=/Users/khaliqgant/.cargo/bin:/Users/khaliqgant/.bun/bin:$PATH RELAYFLOWD_BIN=/Users/khaliqgant/.relayflows-toolchain/target/1463204285/debug/relayflowd npm test +Working directory: /Users/khaliqgant/fl-slice-X/packages/sdk + +Full captured transcript: /private/tmp/fl-slice-X-sdk-tests-yk4hiyhm.log + +Captured webhook failure excerpt: + FAIL tests/webhook-live.test.ts > executes and deduplicates 'app_mention' only for its provider and matching payload + FAIL tests/webhook-live.test.ts > executes and deduplicates 'pull_request' only for its provider and matching payload +Error: webhook integration timed out: FAILED [webhook_server] listen EPERM: operation not permitted 127.0.0.1 + + ❯ until tests/webhook-live.test.ts:39:9 + 37| const deadline = Date.now() + 10_000; + 38| while (Date.now() < deadline) { if (await predicate()) return; await… + 39| throw new Error(`webhook integration timed out: ${detail()}`); + | ^ + 40| } + 41| async function daemon(dir: string): Promise { + ❯ setup tests/webhook-live.test.ts:60:3 + ❯ tests/webhook-live.test.ts:84:25 + +⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯[107/187]⎯ + + FAIL tests/webhook-live.test.ts > executes and deduplicates 'reaction_added' only for its provider and matching payload +Error: webhook integration timed out: FAILED [webhook_server] listen EPERM: operation not permitted 127.0.0.1 + + ❯ until tests/webhook-live.test.ts:39:9 + 37| const deadline = Date.now() + 10_000; + 38| while (Date.now() < deadline) { if (await predicate()) return; await… + 39| throw new Error(`webhook integration timed out: ${detail()}`); + | ^ + 40| } + 41| async function daemon(dir: string): Promise { + ❯ setup tests/webhook-live.test.ts:60:3 + ❯ tests/webhook-live.test.ts:84:25 + +⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯[108/187]⎯ + + FAIL tests/webhook-live.test.ts > flows serve-webhook writes JSON before the daemon starts, then journals and archives exactly once +Error: webhook integration timed out: FAILED [webhook_server] listen EPERM: operation not permitted 127.0.0.1 + + ❯ until tests/webhook-live.test.ts:39:9 + 37| const deadline = Date.now() + 10_000; + 38| while (Date.now() < deadline) { if (await predicate()) return; await… + 39| throw new Error(`webhook integration timed out: ${detail()}`); + | ^ + 40| } + 41| async function daemon(dir: string): Promise { + ❯ setup tests/webhook-live.test.ts:60:3 + ❯ tests/webhook-live.test.ts:116:25 + +⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯[109/187]⎯ + + FAIL tests/webhook-live.test.ts > replays a dropped file after SIGKILL before spawn +Error: webhook integration timed out: FAILED [webhook_server] listen EPERM: operation not permitted 127.0.0.1 + + ❯ until tests/webhook-live.test.ts:39:9 + 37| const deadline = Date.now() + 10_000; + 38| while (Date.now() < deadline) { if (await predicate()) return; await… + 39| throw new Error(`webhook integration timed out: ${detail()}`); + | ^ + 40| } + 41| async function daemon(dir: string): Promise { + ❯ setup tests/webhook-live.test.ts:60:3 + ❯ tests/webhook-live.test.ts:136:25 + +⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯[110/187]⎯ + + FAIL tests/webhook-live.test.ts > resumes the same journal after SIGKILL after spawn and before acknowledgement +Error: webhook integration timed out: FAILED [webhook_server] listen EPERM: operation not permitted 127.0.0.1 + + ❯ until tests/webhook-live.test.ts:39:9 + 37| const deadline = Date.now() + 10_000; + 38| while (Date.now() < deadline) { if (await predicate()) return; await… + 39| throw new Error(`webhook integration timed out: ${detail()}`); + | ^ + 40| } + 41| async function daemon(dir: string): Promise { + ❯ setup tests/webhook-live.test.ts:60:3 + ❯ tests/webhook-live.test.ts:149:25 + +⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯[111/187]⎯ + + FAIL tests/webhook.test.ts > webhook ingress > routes provider envelopes to isolated inboxes and rejects spoofed or unsupported events + FAIL tests/webhook.test.ts > webhook ingress > accepts JSON through atomic files without creating a daemon + FAIL tests/webhook.test.ts > webhook ingress > rejects malformed, oversized, traversal, and non-POST requests + FAIL tests/webhook.test.ts > webhook ingress > fails closed on a symlink inbox target +Error: listen EPERM: operation not permitted 127.0.0.1 +⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯⎯[112/187]⎯ + + +Captured suite summary: + Test Files 25 failed | 55 passed | 1 skipped (81) + Tests 185 failed | 1109 passed | 10 skipped (1304) + Errors 14 errors + Start at 23:52:24 + Duration 68.17s (transform 1.54s, setup 0ms, collect 15.17s, tests 323.65s, environment 9ms, prepare 2.95s) + + +Exit code: 1 diff --git a/docs/evidence/spec-X/surface-regressions.txt b/docs/evidence/spec-X/surface-regressions.txt new file mode 100644 index 000000000..f78af0e89 --- /dev/null +++ b/docs/evidence/spec-X/surface-regressions.txt @@ -0,0 +1,10 @@ +Command: npm run typecheck:regressions +Working directory: /Users/khaliqgant/fl-slice-X/packages/surface + + +> @relayflows/surface@2.0.8 typecheck:regressions +> tsc -p ../../regressions/tsconfig.json && tsc -p tsconfig.test.json && node scripts/check-generated-helpers.mjs + +HELPERS_GENERATED_OK index.ts, slack.ts + +Exit code: 0 diff --git a/docs/evidence/spec-X/surface-tests.txt b/docs/evidence/spec-X/surface-tests.txt new file mode 100644 index 000000000..415eba168 --- /dev/null +++ b/docs/evidence/spec-X/surface-tests.txt @@ -0,0 +1,22 @@ +Command: PATH=/Users/khaliqgant/.bun/bin:$PATH npm test +Working directory: /Users/khaliqgant/fl-slice-X/packages/surface + +> @relayflows/surface@2.0.8 test +> bun run build && tsc -p tsconfig.test.json && vitest run + +$ tsc + + RUN v2.1.9 /Users/khaliqgant/fl-slice-X/packages/surface + + ✓ tests/triggers.test.ts (4 tests) 4ms + ✓ tests/provider-triggers.test.ts (3 tests) 5ms + ✓ tests/flow.test.ts (20 tests) 7ms + ✓ tests/helpers.snapshot.test.ts (1 test) 661ms + ✓ regenerates helpers byte-identically from the pinned adapter 660ms + + Test Files 4 passed (4) + Tests 28 passed (28) + Start at 23:51:18 + Duration 1.55s (transform 64ms, setup 0ms, collect 913ms, tests 677ms, environment 0ms, prepare 387ms) + +Exit code: 0 diff --git a/docs/evidence/spec-X/surface-typecheck.txt b/docs/evidence/spec-X/surface-typecheck.txt new file mode 100644 index 000000000..aec0e9839 --- /dev/null +++ b/docs/evidence/spec-X/surface-typecheck.txt @@ -0,0 +1,9 @@ +Command: npm run typecheck +Working directory: /Users/khaliqgant/fl-slice-X/packages/surface + + +> @relayflows/surface@2.0.8 typecheck +> tsc --noEmit + + +Exit code: 0 diff --git a/packages/sdk/src/cli/serve-webhook.ts b/packages/sdk/src/cli/serve-webhook.ts index 1c216de5c..94141ad28 100644 --- a/packages/sdk/src/cli/serve-webhook.ts +++ b/packages/sdk/src/cli/serve-webhook.ts @@ -4,6 +4,7 @@ import { createServer, type Server, type ServerResponse } from 'node:http'; import { join, resolve } from 'node:path'; import { TextDecoder } from 'node:util'; import type { CliIo } from '../cli.js'; +import { providerInboxEvent } from '../trigger-executor.js'; const MAX_BODY_BYTES = 1024 * 1024; const NAME = /^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$/; @@ -32,7 +33,7 @@ function reply(response: ServerResponse, status: number, body: object): void { response.end(JSON.stringify(body)); } -/** POST / accepts JSON, including scalar values. No daemon connection. */ +/** POST / accepts JSON; /providers/ accepts typed event envelopes. */ export async function startWebhookServer(dataDir: string, port: number): Promise { const inbox = join(resolve(dataDir), 'inbox'); await directory(inbox); @@ -45,7 +46,8 @@ export async function startWebhookServer(dataDir: string, port: number): Promise reply(response, 405, { error: 'method_not_allowed' }); return; } - const name = request.url?.slice(1); + const providerRoute = request.url?.startsWith('/providers/') ?? false; + const name = request.url?.slice(providerRoute ? '/providers/'.length : 1); if (!name || !NAME.test(name)) { request.resume(); reply(response, 404, { error: 'invalid_webhook_name' }); @@ -74,6 +76,15 @@ export async function startWebhookServer(dataDir: string, port: number): Promise reply(response, 400, { error: 'invalid_json' }); return; } + if (providerRoute) { + try { + payload = providerInboxEvent(name, payload); + } catch (error) { + reply(response, 400, { error: 'invalid_provider_event', + message: error instanceof Error ? error.message : 'invalid provider event' }); + return; + } + } const target = join(inbox, name); await directory(target); const id = randomUUID(); diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index 962df481d..80caf6a60 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -216,3 +216,4 @@ export { export { createFlow, type CreateFlowOptions, type CreatedFlow } from './create-flow.js'; export { renderProgress, type ProgressEvent } from './progress.js'; +export { webhookTriggerSpec } from './trigger-executor.js'; diff --git a/packages/sdk/src/trigger-executor.ts b/packages/sdk/src/trigger-executor.ts new file mode 100644 index 000000000..49ac797bc --- /dev/null +++ b/packages/sdk/src/trigger-executor.ts @@ -0,0 +1,34 @@ +import { providerEventTypes, webhook, type TriggerSource, type WebhookFilter } from '@relayflows/surface'; +import type { TriggerSpec } from './spec.js'; + +/** Lower a surface subscription to the existing inbox executor contract. */ +export function webhookTriggerSpec(id: string, source: TriggerSource): TriggerSpec { + if (source.kind !== 'webhook') throw new TypeError('unsupported trigger kind'); + const trigger = webhook(source.name, source.filter); + return { + id, + executor: trigger.name, + eventType: trigger.name, + ...(trigger.filter === undefined ? {} : { pattern: trigger.filter }), + // The inbox watcher supplies the durable file ID as the event key. + dedupeKeyTemplate: '{{event.type}}', + }; +} + +/** Provider ingress uses the path as authority; it never trusts a body to reroute. */ +export function providerInboxEvent(provider: string, value: unknown): WebhookFilter { + if (!Object.hasOwn(providerEventTypes, provider)) throw new TypeError(`unknown provider: ${provider}`); + if (value === null || typeof value !== 'object' || Array.isArray(value)) { + throw new TypeError('provider event must be an object'); + } + const event = value as Record; + const types: readonly string[] = providerEventTypes[provider as keyof typeof providerEventTypes]; + if (typeof event.type !== 'string' || !types.includes(event.type)) { + throw new TypeError(`unknown event type for ${provider}`); + } + if (event.provider !== undefined && event.provider !== provider) throw new TypeError('provider does not match inbox'); + if (event.payload === null || typeof event.payload !== 'object' || Array.isArray(event.payload)) { + throw new TypeError('provider event payload must be an object'); + } + return webhook(provider, { ...event, provider } as WebhookFilter).filter!; +} diff --git a/packages/sdk/tests/generate-triggers.test.ts b/packages/sdk/tests/generate-triggers.test.ts new file mode 100644 index 000000000..4fb98e4a9 --- /dev/null +++ b/packages/sdk/tests/generate-triggers.test.ts @@ -0,0 +1,62 @@ +import { execFileSync } from 'node:child_process'; +import { mkdtempSync, mkdirSync, readFileSync, readdirSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { afterEach, expect, it } from 'vitest'; + +const generator = fileURLToPath(new URL('../../../scripts/generate-triggers.mjs', import.meta.url)); +const dirs: string[] = []; +afterEach(() => { for (const dir of dirs.splice(0)) rmSync(dir, { recursive: true, force: true }); }); +function temporary(): string { + const dir = mkdtempSync(join(tmpdir(), 'provider-codegen-')); + dirs.push(dir); + return dir; +} +function generate(...args: string[]): string { + return execFileSync(process.execPath, [generator, ...args], { encoding: 'utf8', stdio: 'pipe' }); +} +function mapping(root: string, path: string, value: unknown): void { + const target = join(root, 'packages', path); + mkdirSync(target, { recursive: true }); + writeFileSync(join(target, 'test.mapping.yaml'), JSON.stringify(value)); +} + +it('reproduces all checked-in modules from the pinned adapter mappings', () => { + expect(generate('--check')).toContain('Checked 2 provider trigger modules'); +}); + +it('discovers new adapters, preserves exact event names, and prefers adapter-local mappings', () => { + const root = temporary(); + const out = join(root, 'generated'); + mapping(root, 'core/mappings', { adapter: { name: 'github' }, webhooks: { stale: {} } }); + mapping(root, 'github', { adapter: { name: 'github' }, webhooks: { pull_request: { extract: ['action'] } } }); + mapping(root, 'new-provider', { provider: 'new-provider', webhooks: { 'file.created': {}, 'file.deleted': {} } }); + mapping(root, 'no-events', { provider: 'no-events', webhooks: {} }); + generate('--adapters-dir', root, '--out-dir', out); + expect(readdirSync(out).sort()).toEqual(['github.ts', 'index.ts', 'new-provider.ts']); + expect(readFileSync(join(out, 'github.ts'), 'utf8')).toContain('pull_request(action?: string)'); + expect(readFileSync(join(out, 'github.ts'), 'utf8')).not.toContain('stale'); + const provider = readFileSync(join(out, 'new-provider.ts'), 'utf8'); + expect(provider).toContain('export const new_provider'); + expect(provider).toContain('file_created(filter?: WebhookFilter)'); + expect(provider).toContain('providerTrigger("new-provider", "file.created", filter)'); + generate('--adapters-dir', root, '--out-dir', out, '--check'); + mapping(root, 'no-events', { provider: 'no-events', webhooks: { added: {} } }); + generate('--adapters-dir', root, '--out-dir', out); + expect(readdirSync(out)).toContain('no-events.ts'); + mapping(root, 'no-events', { provider: 'no-events', webhooks: {} }); + generate('--adapters-dir', root, '--out-dir', out); + expect(readdirSync(out)).not.toContain('no-events.ts'); + writeFileSync(join(out, 'new-provider.ts'), 'stale'); + expect(() => generate('--adapters-dir', root, '--out-dir', out, '--check')).toThrow(/drifted/); +}); + +it('fails closed on malformed mappings and colliding method names before writing output', () => { + for (const webhooks of [[], { broken: null }, { 'file.created': {}, file_created: {} }]) { + const root = temporary(); + mapping(root, 'example', { provider: 'example', webhooks }); + expect(() => generate('--adapters-dir', root, '--out-dir', join(root, 'out'))).toThrow(); + expect(readdirSync(root)).toEqual(['packages']); + } +}); diff --git a/packages/sdk/tests/provider-trigger-executor.test.ts b/packages/sdk/tests/provider-trigger-executor.test.ts new file mode 100644 index 000000000..7a76078e5 --- /dev/null +++ b/packages/sdk/tests/provider-trigger-executor.test.ts @@ -0,0 +1,60 @@ +import { execFileSync } from 'node:child_process'; +import { mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { afterEach, expect, it } from 'vitest'; +import { github, slack } from '@relayflows/surface'; +import { compileSpec, toKernelSpec } from '../src/compile.js'; +import { providerInboxEvent, webhookTriggerSpec } from '../src/trigger-executor.js'; + +const binary = process.env['RELAYFLOWD_BIN'] ?? resolve('../../kernel/target/debug/relayflowd'); +const dirs: string[] = []; +afterEach(() => { for (const dir of dirs.splice(0)) rmSync(dir, { recursive: true, force: true }); }); + +it('validates and snapshots provider envelopes before they reach an inbox', () => { + const input = { type: 'pull_request', payload: { action: 'opened' } }; + const event = providerInboxEvent('github', input); + input.payload.action = 'closed'; + expect(event).toEqual({ provider: 'github', type: 'pull_request', payload: { action: 'opened' } }); + expect(Object.isFrozen(event.payload)).toBe(true); + expect(() => providerInboxEvent('unknown', input)).toThrow(/unknown provider/); + expect(() => providerInboxEvent('toString', input)).toThrow(/unknown provider/); + expect(() => providerInboxEvent('slack', input)).toThrow(/unknown event type/); + expect(() => providerInboxEvent('github', { ...input, provider: 'slack' })).toThrow(/does not match/); + for (const payload of [undefined, null, 1, [], { invalid: Infinity }]) { + expect(() => providerInboxEvent('github', { type: 'pull_request', payload })).toThrow(); + } +}); + +it.each([ + { source: slack.mention('C123'), type: 'app_mention', payload: { channel: 'C123' }, rejected: { channel: 'C456' } }, + { source: slack.reaction('eyes'), type: 'reaction_added', payload: { reaction: 'eyes' }, rejected: { reaction: 'heart' } }, + { source: github.pull_request('opened'), type: 'pull_request', payload: { action: 'opened' }, rejected: { action: 'closed' } }, +])('the kernel executes compiled $type subscriptions with provider isolation and durable dedupe', ({ source, type, payload, rejected }) => { + const dir = mkdtempSync(join(tmpdir(), 'provider-executor-')); + dirs.push(dir); + const effect = join(dir, 'effect.txt'); + const spec = join(dir, 'binding.json'); + writeFileSync(spec, JSON.stringify(toKernelSpec(compileSpec({ + version: '0.1.0', name: 'provider-test', + triggers: [webhookTriggerSpec('event', source)], + steps: [{ id: 'effect', type: 'deterministic', command: `printf accepted >> '${effect}'` }], + })))); + const submit = (envelope: unknown, key: string, executor = source.name) => JSON.parse(execFileSync(binary, [ + '--data-dir', dir, 'run', spec, '--event', JSON.stringify({ type: executor, payload: envelope, key }), + ], { encoding: 'utf8', stdio: 'pipe' })) as { matched: boolean; deduped: boolean; run: { run_id: string } | null }; + const event = providerInboxEvent(source.name, { type, payload }); + for (const envelope of [ + { ...event, provider: 'other' }, { ...event, type: 'other' }, { ...event, payload: rejected }, + ]) { + expect(submit(envelope, 'rejected')).toMatchObject({ matched: false, run: null }); + } + expect(submit(event, 'wrong-inbox', 'other')).toMatchObject({ matched: false, run: null }); + const first = submit(event, 'event-1.json'); + expect(first).toMatchObject({ matched: true, deduped: false }); + expect(readFileSync(effect, 'utf8')).toBe('accepted'); + expect(submit(event, 'event-1.json')).toMatchObject({ matched: true, deduped: true, run: null }); + expect(readFileSync(effect, 'utf8')).toBe('accepted'); + expect(submit(event, 'event-2.json')).toMatchObject({ matched: true, deduped: false }); + expect(readFileSync(effect, 'utf8')).toBe('acceptedaccepted'); +}); diff --git a/packages/sdk/tests/webhook-live.test.ts b/packages/sdk/tests/webhook-live.test.ts index 3108a8e0c..cc2cf8133 100644 --- a/packages/sdk/tests/webhook-live.test.ts +++ b/packages/sdk/tests/webhook-live.test.ts @@ -7,6 +7,9 @@ import { join, resolve } from 'node:path'; import { setTimeout as delay } from 'node:timers/promises'; import { socketPathFor } from '../src/daemon-connection.js'; import { JournalClient } from '../src/journal-client.js'; +import { github, slack, webhook, type TriggerSource } from '@relayflows/surface'; +import { webhookTriggerSpec } from '../src/trigger-executor.js'; +import { compileSpec, toKernelSpec } from '../src/compile.js'; const binary = process.env['RELAYFLOWD_BIN'] ?? resolve('../../kernel/target/debug/relayflowd'); const children: ChildProcess[] = []; @@ -44,15 +47,15 @@ async function daemon(dir: string): Promise { }, process.output); return process.child; } -async function setup(command = 'printf accepted'): Promise<{ dir: string; base: string }> { +async function setup(command = 'printf accepted', source: TriggerSource = webhook('release')): Promise<{ dir: string; base: string }> { const dir = await mkdtemp(join(tmpdir(), 'flows-inbox-live-')); directories.push(dir); await mkdir(join(dir, 'triggers')); - await writeFile(join(dir, 'triggers', 'release.json'), JSON.stringify({ - version: '0.1.0', name: 'release', - triggers: [{ id: 'release', executor: 'release', event_type: 'release', dedupe_key_template: '{{event.type}}' }], + await writeFile(join(dir, 'triggers', `${source.name}.json`), JSON.stringify(toKernelSpec(compileSpec({ + version: '0.1.0', name: source.name, + triggers: [webhookTriggerSpec('event', source)], steps: [{ id: 'log', type: 'deterministic', command }], - })); + })))); const receiver = start(process.execPath, [resolve('dist/cli.js'), 'serve-webhook', '--data-dir', dir, '--port', '0']); await until(() => /WEBHOOK http:\/\/127.0.0.1:\d+/.test(receiver.output()), receiver.output); return { dir, base: receiver.output().match(/http:\/\/127.0.0.1:\d+/)![0] }; @@ -73,6 +76,42 @@ async function readJournal(dir: string, runFile: string) { return client.journalRead(runFile.slice(0, -8), 1); } +it.each([ + { source: slack.mention('C123'), type: 'app_mention', payload: { channel: 'C123' }, rejected: { channel: 'C456' } }, + { source: slack.reaction('eyes'), type: 'reaction_added', payload: { reaction: 'eyes' }, rejected: { reaction: 'heart' } }, + { source: github.pull_request('opened'), type: 'pull_request', payload: { action: 'opened' }, rejected: { action: 'closed' } }, +])('executes and deduplicates $type only for its provider and matching payload', async ({ source, type, payload, rejected }) => { + const { dir, base } = await setup('printf provider-accepted', source); + const ids: string[] = []; + for (const body of [{ type, payload }, { type, payload: rejected }]) { + const response = await fetch(`${base}/providers/${source.name}`, { method: 'POST', body: JSON.stringify(body) }); + expect(response.status).toBe(202); + ids.push(`${(await response.json() as { id: string }).id}.json`); + } + // Even a generic webhook writer cannot bypass the provider/type filter. + for (const body of [ + { provider: 'other', type, payload }, + { provider: source.name, type: 'other', payload }, + ]) { + const response = await fetch(`${base}/${source.name}`, { method: 'POST', body: JSON.stringify(body) }); + expect(response.status).toBe(202); + ids.push(`${(await response.json() as { id: string }).id}.json`); + } + await daemon(dir); + const processed = join(dir, 'inbox-processed', source.name); + await until(() => ids.every(id => existsSync(join(processed, id)))); + const files = await journals(dir); + expect(files).toHaveLength(1); + const journal = await readJournal(dir, files[0]!); + expect(journal.entries.find(entry => entry.entry_type === 'run.spawned')?.payload['event']) + .toEqual({ provider: source.name, type, payload }); + expect(journal.entries.filter(entry => entry.entry_type === 'step.completed')).toHaveLength(1); + await rename(join(processed, ids[0]!), join(dir, 'inbox', source.name, ids[0]!)); + await until(() => existsSync(join(processed, ids[0]!))); + expect(await journals(dir)).toEqual(files); + expect((await readJournal(dir, files[0]!)).entries.filter(entry => entry.entry_type === 'step.completed')).toHaveLength(1); +}, 20_000); + it('flows serve-webhook writes JSON before the daemon starts, then journals and archives exactly once', async () => { const { dir, base } = await setup(); const filename = await post(base); diff --git a/packages/sdk/tests/webhook.test.ts b/packages/sdk/tests/webhook.test.ts index 593394646..be0073189 100644 --- a/packages/sdk/tests/webhook.test.ts +++ b/packages/sdk/tests/webhook.test.ts @@ -5,7 +5,8 @@ import { tmpdir } from 'node:os'; import type { Server } from 'node:http'; import { startWebhookServer, parseWebhookArgs } from '../src/cli/serve-webhook.js'; import { preflightWebhookTriggers } from '../src/preflight.js'; -import { webhook } from '@relayflows/surface'; +import { webhook, slack, github } from '@relayflows/surface'; +import { webhookTriggerSpec } from '../src/trigger-executor.js'; import { runCli } from '../src/cli.js'; import { runDirectFlow } from '../src/cli/direct-run.js'; @@ -30,6 +31,41 @@ async function receiver(dir: string): Promise { } describe('webhook ingress', () => { + it('lowers provider subscriptions to provider inbox executors with payload filters', () => { + expect(webhookTriggerSpec('mention', slack.mention('C123'))).toEqual({ + id: 'mention', executor: 'slack', eventType: 'slack', dedupeKeyTemplate: '{{event.type}}', + pattern: { provider: 'slack', type: 'app_mention', payload: { channel: 'C123' } }, + }); + expect(preflightWebhookTriggers([slack.mention('C123'), slack.reaction('eyes'), github.pull_request()], ['slack'])) + .toEqual([expect.objectContaining({ kind: 'no_executor', executor: 'github' })]); + expect(webhookTriggerSpec('plain', webhook('release')).pattern).toBeUndefined(); + }); + it('routes provider envelopes to isolated inboxes and rejects spoofed or unsupported events', async () => { + const dir = await temporary(); + const base = await receiver(dir); + const post = (provider: string, body: unknown) => fetch(`${base}/providers/${provider}`, { + method: 'POST', body: JSON.stringify(body), + }); + for (const [provider, type, payload] of [ + ['slack', 'app_mention', { channel: 'C123', text: 'hello' }], + ['github', 'pull_request', { action: 'opened', number: 42 }], + ] as const) { + const response = await post(provider, { type, payload }); + expect(response.status).toBe(202); + const { id } = await response.json() as { id: string }; + expect(JSON.parse(await readFile(join(dir, 'inbox', provider, `${id}.json`), 'utf8'))) + .toEqual({ provider, type, payload }); + } + for (const body of [null, {}, { type: 'pull_request', payload: {} }, + { provider: 'github', type: 'app_mention', payload: {} }, + { type: 'app_mention', payload: [] }, { type: 'app_mention', payload: null }]) { + expect((await post('slack', body)).status).toBe(400); + } + expect((await post('unknown', { type: 'push', payload: {} })).status).toBe(400); + expect((await post('slack%2Fgithub', { type: 'push', payload: {} })).status).toBe(404); + expect(await readdir(join(dir, 'inbox', 'slack'))).toHaveLength(1); + expect(await readdir(join(dir, 'inbox'))).toEqual(['github', 'slack']); + }); it('requires registered executor names with the exact refusal', () => { expect(preflightWebhookTriggers([webhook('unregistered')], [])).toEqual([{ severity: 'refusal', kind: 'no_executor', executor: 'unregistered', diff --git a/packages/surface/README.md b/packages/surface/README.md index 8121887b8..28cc1c616 100644 --- a/packages/surface/README.md +++ b/packages/surface/README.md @@ -26,8 +26,9 @@ lifecycle contract are in `docs/SURFACE.md`, "The authored operation lifecycle". This is an in-repository foundation, not a registry-published package. The repository's `flows run` command can execute a directly authored `.flow.ts` with required JSON input. Durable authored-root resume remains tracked in issue -#132. Resident trigger handlers (`flow.on(...)`) are gate-2 work and are not yet -part of this package. +#132. `flow.on(...)` records webhook and generated provider subscriptions. +See [provider event triggers](src/triggers/README.md) for declarations, inbox +envelopes, and executor bindings; deploying handler bodies is tracked in #301. The repository pins Bun through `surface/bun.lock`. From a fresh checkout: diff --git a/packages/surface/package.json b/packages/surface/package.json index 5c174c414..760a4c9c7 100644 --- a/packages/surface/package.json +++ b/packages/surface/package.json @@ -15,6 +15,16 @@ "types": "./dist/runtime.d.ts", "import": "./dist/runtime.js", "require": "./dist/runtime.js" + }, + "./triggers": { + "types": "./dist/triggers/index.d.ts", + "import": "./dist/triggers/index.js", + "require": "./dist/triggers/index.js" + }, + "./triggers/*": { + "types": "./dist/triggers/*.d.ts", + "import": "./dist/triggers/*.js", + "require": "./dist/triggers/*.js" } }, "files": [ @@ -23,6 +33,7 @@ ], "scripts": { "gen": "node ../../scripts/generate-helpers.mjs", + "gen:triggers": "node ../../scripts/generate-triggers.mjs", "build": "tsc", "prepare": "bun run build", "typecheck": "tsc --noEmit", diff --git a/packages/surface/src/index.ts b/packages/surface/src/index.ts index dac12590c..27dc4fd26 100644 --- a/packages/surface/src/index.ts +++ b/packages/surface/src/index.ts @@ -25,4 +25,6 @@ export { flowRunWritebackIdempotency, type SlackHelper, type SlackReceipt } from export type { Helpers } from "./helpers/index.js"; export type { MemoryHelper, MemoryFinding, MemoryRecallOptions, HistoryEntry, TrajectoryEntry } from "./memory.js"; export { webhook, type TriggerSource, type WebhookFilter, type WebhookValue } from "./triggers.js"; +export * from "./triggers/index.js"; +export type { ProviderTriggerSource } from "./provider-trigger.js"; export type { PluginMethod, PluginPrimitive } from './plugin-contract.js'; diff --git a/packages/surface/src/provider-trigger.ts b/packages/surface/src/provider-trigger.ts new file mode 100644 index 000000000..98a163f36 --- /dev/null +++ b/packages/surface/src/provider-trigger.ts @@ -0,0 +1,26 @@ +import { webhook, type TriggerSource, type WebhookFilter } from "./triggers.js"; + +/** A provider inbox subscription over an envelope's type and payload. */ +export interface ProviderTriggerSource extends TriggerSource { + readonly name: Provider; + readonly filter: WebhookFilter & { readonly provider: Provider; readonly type: Event }; +} + +/** Shared implementation for generated provider declarations. */ +export function providerTrigger

( + provider: P, event: E, payload?: WebhookFilter, +): ProviderTriggerSource { + if (typeof event !== "string" || event.length === 0) throw new TypeError("provider event type must not be empty"); + // Validate separately so an invalid payload cannot become a scalar filter leaf. + const filter = webhook(provider, payload).filter; + return webhook(provider, { + provider, type: event, ...(filter === undefined ? {} : { payload: filter }), + }) as ProviderTriggerSource; +} + +export function triggerArgument(value: string, name: string): string { + if (typeof value !== "string" || value.trim().length === 0) { + throw new TypeError(`${name} must be a nonempty string`); + } + return value; +} diff --git a/packages/surface/src/triggers/README.md b/packages/surface/src/triggers/README.md new file mode 100644 index 000000000..8a6c5640f --- /dev/null +++ b/packages/surface/src/triggers/README.md @@ -0,0 +1,75 @@ +# Provider event triggers + +```ts +import { flow, slack, github } from '@relayflows/surface'; + +export default flow('triage') + .on(slack.mention('C123'), async (f, event) => { f.done('success'); }) + .on(slack.reaction('eyes'), async (f, event) => { f.done('success'); }) + .on(github.pull_request('opened'), async (f, event) => { f.done('success'); }); +``` + +Provider namespaces also export from `@relayflows/surface/triggers` and +`@relayflows/surface/triggers/slack` (or `/github`). Declarations are immutable +webhook sources: their executor/inbox name is the provider, and their filter +matches the provider, event type, and requested payload fields. Register those +provider names in `flows.json`'s `executors` array for preflight. + +`slack.mention(channel)` subscribes to Slack's `app_mention` event; +`slack.reaction(emoji)` subscribes to `reaction_added`. Arguments match provider +values exactly: use the channel ID and reaction name from the incoming event. +`github.pull_request(action)` filters `payload.action`; omitting the action +accepts every pull request event. Upstream mappings do not declare action enums, +so the action parameter is a string. Other generated methods accept an optional +recursive payload filter, for example `github.push({ ref: 'refs/heads/main' })`. + +The receiver accepts `POST /providers/slack` with this JSON: + +```json +{"type":"app_mention","payload":{"channel":"C123","text":"Please review"}} +``` + +It writes `{ "provider": "slack", "type": "app_mention", "payload": {...} }` +atomically to `inbox/slack/.json`. The path supplies the provider; a conflicting +provider in the body or unknown event type is rejected. `payload` is the provider +event body (Slack's inner `event` object). The event names are the exact mapping +keys; dots and hyphens become underscores only in TypeScript method names. +Ingress expects this envelope, not a raw vendor HTTP delivery. Signature +verification and public ingress remain part of the deployment work tracked in +#301. Generic `POST /` retains its existing arbitrary-JSON contract. + +The SDK's `webhookTriggerSpec(id, source)` lowers a declaration to `TriggerSpec` +for `compileSpec`/`toKernelSpec`. Provision the resulting RunSpec in +`triggers/.json`, with executable steps, as in the existing webhook +executor. The watcher applies the generated filters, journals matching events, +deduplicates by inbox file ID, and archives consumed events. For example: + +```ts +import { compileSpec, toKernelSpec, webhookTriggerSpec } from '@relayflows/sdk'; +import { slack } from '@relayflows/surface'; + +const binding = toKernelSpec(compileSpec({ + version: '0.1.0', name: 'mentions', + triggers: [webhookTriggerSpec('mention', slack.mention('C123'))], + steps: [{ id: 'ack', type: 'deterministic', command: 'printf accepted' }], +})); // Write binding as JSON to triggers/slack.json before ingress starts. +``` + +Each provider inbox still binds one RunSpec. This extends E's declaration and +inbox contract; deploying TypeScript handler bodies remains the separate #301 +binding step. + +Generation requires the SDK's development dependencies. The default input is +the mapping YAML shipped in the pinned `@relayfile/adapter-core` dependency +(currently Slack and GitHub). A checkout supplies additional providers: + +```sh +node scripts/generate-triggers.mjs +node scripts/generate-triggers.mjs --check +node scripts/generate-triggers.mjs --adapters-dir /path/to/relayfile-adapters +``` + +The generator reads `packages/core/mappings/*.mapping.yaml`, then each adapter's +`packages//*.mapping.yaml`; adapter-local mappings take precedence. +Only adapters with a nonempty `webhooks:` section produce modules. It never +imports or executes adapter code. Rebuild the surface and SDK after regeneration. diff --git a/packages/surface/src/triggers/github.ts b/packages/surface/src/triggers/github.ts new file mode 100644 index 000000000..43c6959b5 --- /dev/null +++ b/packages/surface/src/triggers/github.ts @@ -0,0 +1,19 @@ +// GENERATED by scripts/generate-triggers.mjs — do not edit. + +import { providerTrigger, triggerArgument } from "../provider-trigger.js"; +import type { WebhookFilter } from "../triggers.js"; + +export const github = Object.freeze({ + issues(filter?: WebhookFilter) { + return providerTrigger("github", "issues", filter); + }, + pull_request(action?: string) { + return providerTrigger("github", "pull_request", action === undefined ? undefined : { action: triggerArgument(action, "action") }); + }, + pull_request_review(filter?: WebhookFilter) { + return providerTrigger("github", "pull_request_review", filter); + }, + push(filter?: WebhookFilter) { + return providerTrigger("github", "push", filter); + }, +}); diff --git a/packages/surface/src/triggers/index.ts b/packages/surface/src/triggers/index.ts new file mode 100644 index 000000000..0ae53b47e --- /dev/null +++ b/packages/surface/src/triggers/index.ts @@ -0,0 +1,10 @@ +// GENERATED by scripts/generate-triggers.mjs — do not edit. + +export { github } from "./github.js"; +export { slack } from "./slack.js"; + +/** Exact upstream event names, plus the generated Slack mention shorthand. */ +export const providerEventTypes = Object.freeze({ + "github": Object.freeze(["issues","pull_request","pull_request_review","push"] as const), + "slack": Object.freeze(["app_mention","message","reaction_added"] as const), +}); diff --git a/packages/surface/src/triggers/slack.ts b/packages/surface/src/triggers/slack.ts new file mode 100644 index 000000000..17bc258ad --- /dev/null +++ b/packages/surface/src/triggers/slack.ts @@ -0,0 +1,19 @@ +// GENERATED by scripts/generate-triggers.mjs — do not edit. + +import { providerTrigger, triggerArgument } from "../provider-trigger.js"; +import type { WebhookFilter } from "../triggers.js"; + +export const slack = Object.freeze({ + message(filter?: WebhookFilter) { + return providerTrigger("slack", "message", filter); + }, + reaction_added(filter?: WebhookFilter) { + return providerTrigger("slack", "reaction_added", filter); + }, + mention(channel: string) { + return providerTrigger("slack", "app_mention", { channel: triggerArgument(channel, "channel") }); + }, + reaction(emoji: string) { + return providerTrigger("slack", "reaction_added", { reaction: triggerArgument(emoji, "emoji") }); + }, +}); diff --git a/packages/surface/tests/provider-triggers.test.ts b/packages/surface/tests/provider-triggers.test.ts new file mode 100644 index 000000000..c32fb7564 --- /dev/null +++ b/packages/surface/tests/provider-triggers.test.ts @@ -0,0 +1,59 @@ +import { describe, expect, it } from 'vitest'; +import { flow, github, slack, type ProviderTriggerSource } from '@relayflows/surface'; +import { slack as slackSubpath } from '@relayflows/surface/triggers/slack'; +import { github as githubSubpath } from '@relayflows/surface/triggers'; +import { getFlowDefinition } from '@relayflows/surface/runtime'; + +describe('generated provider declarations', () => { + it('exposes typed, immutable provider subscriptions through the package exports', () => { + const mention: ProviderTriggerSource<'slack', 'app_mention'> = slack.mention('C123'); + expect(mention).toEqual({ kind: 'webhook', name: 'slack', filter: { + provider: 'slack', type: 'app_mention', payload: { channel: 'C123' }, + } }); + expect(slack.reaction('eyes').filter).toEqual({ + provider: 'slack', type: 'reaction_added', payload: { reaction: 'eyes' }, + }); + expect(github.pull_request('opened').filter).toEqual({ + provider: 'github', type: 'pull_request', payload: { action: 'opened' }, + }); + expect(github.pull_request().filter).toEqual({ provider: 'github', type: 'pull_request' }); + expect(slackSubpath).toBe(slack); + expect(githubSubpath).toBe(github); + expect(Object.isFrozen(slack)).toBe(true); + expect(Object.isFrozen(mention.filter.payload)).toBe(true); + }); + + it('snapshots filters and preserves routing through flow.on without running handlers', () => { + const filter = { repository: { name: 'flows' } }; + const source = github.push(filter); + const declared = flow('providers').on(source, async () => { throw new Error('handler ran'); }); + filter.repository.name = 'changed'; + expect(getFlowDefinition(declared).handlers[0]!.trigger).toEqual(source); + expect(source.filter.payload).toEqual({ repository: { name: 'flows' } }); + expect(Object.isFrozen(source.filter.payload)).toBe(true); + }); + + it('refuses invalid typed arguments and non-object filters at runtime', () => { + for (const argument of ['', ' ', null, 1, {}]) { + expect(() => slack.mention(argument as string)).toThrow(/channel/); + expect(() => slack.reaction(argument as string)).toThrow(/emoji/); + expect(() => github.pull_request(argument as string)).toThrow(/action/); + } + expect(() => github.push([] as never)).toThrow(/filter/); + expect(() => github.push({ value: NaN })).toThrow(/filter/); + }); +}); + +// Compile-time contracts are checked by the surface test tsconfig. +function invalidDeclarations(): void { + // @ts-expect-error channel is a string + slack.mention(123); + // @ts-expect-error emoji is required + slack.reaction(); + // @ts-expect-error action is a string, not an arbitrary object + github.pull_request({ action: 'opened' }); + // @ts-expect-error provider and event literals cannot be interchanged + const source: ProviderTriggerSource<'github', 'push'> = slack.mention('C123'); + void source; +} +void invalidDeclarations; diff --git a/scripts/generate-triggers.mjs b/scripts/generate-triggers.mjs new file mode 100644 index 000000000..2d60d9444 --- /dev/null +++ b/scripts/generate-triggers.mjs @@ -0,0 +1,116 @@ +import assert from 'node:assert/strict'; +import { existsSync, mkdirSync, readFileSync, readdirSync, rmSync, writeFileSync } from 'node:fs'; +import { createRequire } from 'node:module'; +import { dirname, join, resolve } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { parseArgs } from 'node:util'; + +const root = fileURLToPath(new URL('../', import.meta.url)); +const require = createRequire(join(root, 'packages/sdk/package.json')); +const { parse } = require('yaml'); +const { values } = parseArgs({ options: { + 'adapters-dir': { type: 'string' }, + 'out-dir': { type: 'string', default: join(root, 'packages/surface/src/triggers') }, + check: { type: 'boolean', default: false }, +} }); +const mappings = new Map(); +const object = value => value !== null && typeof value === 'object' && !Array.isArray(value); +const identifier = value => value.replaceAll(/[^A-Za-z0-9_$]/g, '_'); + +function readMappings(directory) { + for (const name of readdirSync(directory).sort()) { + if (!name.endsWith('.mapping.yaml')) continue; + const mapping = parse(readFileSync(join(directory, name), 'utf8'), { uniqueKeys: true }); + assert(object(mapping), `Invalid mapping: ${name}`); + const provider = mapping.adapter?.name ?? mapping.provider ?? name.replace('.mapping.yaml', ''); + assert(typeof provider === 'string' && /^[a-z][a-z0-9-]*$/.test(provider), `Invalid provider: ${provider}`); + if (mapping.webhooks === undefined) continue; + assert(object(mapping.webhooks), `Invalid webhooks: ${name}`); + for (const [event, definition] of Object.entries(mapping.webhooks)) { + assert(event.length > 0 && object(definition), `Invalid webhook: ${name}/${event}`); + assert(definition.extract === undefined || (Array.isArray(definition.extract) + && definition.extract.every(field => typeof field === 'string')), `Invalid extract: ${name}/${event}`); + } + // Adapter-local mappings supersede the core's fallback mappings as a whole. + mappings.set(provider, mapping.webhooks); + } +} + +if (values['adapters-dir']) { + const packages = join(resolve(values['adapters-dir']), 'packages'); + const core = join(packages, 'core/mappings'); + if (existsSync(core)) readMappings(core); + for (const entry of readdirSync(packages, { withFileTypes: true }).sort((a, b) => a.name.localeCompare(b.name))) { + if (entry.isDirectory() && entry.name !== 'core') readMappings(join(packages, entry.name)); + } +} else { + // Pinned SDK dependency includes the upstream mapping YAML in its npm tarball. + readMappings(join(dirname(require.resolve('@relayfile/adapter-core/package.json')), 'mappings')); +} +assert([...mappings.values()].some(events => Object.keys(events).length), 'No webhook mappings found'); + +const header = '// GENERATED by scripts/generate-triggers.mjs — do not edit.\n'; +const files = new Map(); +const exports = []; +const registry = {}; +const namespaces = new Set(); +for (const [provider, events] of [...mappings].sort(([a], [b]) => a.localeCompare(b))) { + if (!Object.keys(events).length) continue; + const namespace = identifier(provider); + assert(!namespaces.has(namespace) && !['index', 'providerEventTypes'].includes(namespace), `Provider name collision: ${provider}`); + namespaces.add(namespace); + const methods = []; + const names = new Set(); + const eventNames = Object.keys(events).sort(); + const add = (name, code) => { + assert(!names.has(name), `Trigger name collision: ${provider}.${name}`); + names.add(name); + methods.push(code); + }; + for (const event of eventNames) { + const name = identifier(event); + assert(/^[A-Za-z_$]/.test(name), `Invalid event identifier: ${event}`); + const prefix = `providerTrigger(${JSON.stringify(provider)}, ${JSON.stringify(event)}`; + if (events[event].extract?.includes('action')) { + add(name, ` ${name}(action?: string) {\n return ${prefix}, action === undefined ? undefined : { action: triggerArgument(action, "action") });\n },`); + } else { + add(name, ` ${name}(filter?: WebhookFilter) {\n return ${prefix}, filter);\n },`); + } + } + // Author-facing shorthand for Slack's raw Events API payloads. These are + // emitted only when the corresponding mapping capability exists. + if (provider === 'slack' && events.message) { + add('mention', ' mention(channel: string) {\n return providerTrigger("slack", "app_mention", { channel: triggerArgument(channel, "channel") });\n },'); + if (!eventNames.includes('app_mention')) eventNames.push('app_mention'); + } + if (provider === 'slack' && events.reaction_added) { + add('reaction', ' reaction(emoji: string) {\n return providerTrigger("slack", "reaction_added", { reaction: triggerArgument(emoji, "emoji") });\n },'); + } + const needsArgument = methods.some(method => method.includes('triggerArgument(')); + const needsFilter = methods.some(method => method.includes('WebhookFilter')); + files.set(`${provider}.ts`, `${header}\nimport { providerTrigger${needsArgument ? ', triggerArgument' : ''} } from "../provider-trigger.js";\n` + + (needsFilter ? 'import type { WebhookFilter } from "../triggers.js";\n' : '') + + `\nexport const ${namespace} = Object.freeze({\n${methods.join('\n')}\n});\n`); + exports.push(`export { ${namespace} } from "./${provider}.js";`); + registry[provider] = eventNames.sort(); +} +files.set('index.ts', `${header}\n${exports.join('\n')}\n\n` + + '/** Exact upstream event names, plus the generated Slack mention shorthand. */\n' + + `export const providerEventTypes = Object.freeze({\n${Object.entries(registry).map(([provider, events]) => + ` ${JSON.stringify(provider)}: Object.freeze(${JSON.stringify(events)} as const),`).join('\n')}\n});\n`); + +const destination = resolve(values['out-dir']); +if (values.check) { + const actual = existsSync(destination) ? readdirSync(destination).filter(name => name.endsWith('.ts')) : []; + const stale = [...new Set([...files.keys(), ...actual])].filter(name => !files.has(name) + || !existsSync(join(destination, name)) || readFileSync(join(destination, name), 'utf8') !== files.get(name)); + assert.equal(stale.length, 0, `Generated triggers drifted: ${stale.join(', ')}`); +} else { + mkdirSync(destination, { recursive: true }); + for (const name of readdirSync(destination)) { + if (name.endsWith('.ts') && !files.has(name) + && readFileSync(join(destination, name), 'utf8').startsWith(header)) rmSync(join(destination, name)); + } + for (const [name, content] of files) writeFileSync(join(destination, name), content); +} +console.log(`${values.check ? 'Checked' : 'Generated'} ${files.size - 1} provider trigger modules`);