From 4a3ac85114e3467c69dc5fc6626b5193240676bd Mon Sep 17 00:00:00 2001 From: drewstone Date: Fri, 11 Sep 2026 19:34:06 -0700 Subject: [PATCH 1/4] fix(workflows): revalidate current trigger grants before dispatch --- src/workflow.ts | 36 +++++++++++++++++++++++++++++------- 1 file changed, 29 insertions(+), 7 deletions(-) diff --git a/src/workflow.ts b/src/workflow.ts index dff155f6..17cf3313 100644 --- a/src/workflow.ts +++ b/src/workflow.ts @@ -131,12 +131,34 @@ export class IntegrationWorkflowRuntime { event: IntegrationTriggerEvent, handler: (input: { event: IntegrationTriggerEvent; workflows: InstalledIntegrationWorkflow[] }) => Promise | void, ): Promise<{ matched: InstalledIntegrationWorkflow[] }> { - const workflows = (await this.store.list()) - .filter((workflow) => - workflow.status === 'active' - && workflow.subscription.connectionId === event.connectionId - && workflow.subscription.trigger === event.trigger - ) + // A workflow outlives its installation-time grant. Recheck the current + // grant before routing a new event; an active subscription is not authority. + // This assumes verified ingress. Deferred effects must recheck authority + // again at execution time rather than treating this snapshot as a lease. + const workflows: InstalledIntegrationWorkflow[] = [] + for (const workflow of await this.store.list()) { + if ( + workflow.status !== 'active' + || workflow.subscription.status !== 'active' + || workflow.subscription.connectionId !== event.connectionId + || workflow.subscription.trigger !== event.trigger + ) continue + + const grant = await this.grants.get(workflow.triggerGrantId) + if ( + !grant + || grant.id !== workflow.triggerGrantId + || grant.status !== 'active' + || grant.manifestId !== workflow.manifestId + || !sameActor(grant.owner, workflow.owner) + || !sameActor(grant.grantee, workflow.grantee) + || grant.connectionId !== event.connectionId + || grant.connectorId !== event.connectorId + || !grant.allowedTriggers.includes(event.trigger) + ) continue + + workflows.push(workflow) + } await handler({ event, workflows }) return { matched: workflows } } @@ -156,4 +178,4 @@ function findTriggerGrant(grants: IntegrationGrant[], requirementId: string, tri function sameActor(a: IntegrationActor, b: IntegrationActor): boolean { return a.type === b.type && a.id === b.id -} +} \ No newline at end of file From ef061d4e37fff27cb89dcabe937779b08cc8c30e Mon Sep 17 00:00:00 2001 From: drewstone Date: Fri, 11 Sep 2026 19:35:15 -0700 Subject: [PATCH 2/4] test(workflows): cover revoked and mismatched trigger grants --- tests/workflow-authorization.test.ts | 155 +++++++++++++++++++++++++++ 1 file changed, 155 insertions(+) create mode 100644 tests/workflow-authorization.test.ts diff --git a/tests/workflow-authorization.test.ts b/tests/workflow-authorization.test.ts new file mode 100644 index 00000000..a9abe47d --- /dev/null +++ b/tests/workflow-authorization.test.ts @@ -0,0 +1,155 @@ +import { describe, expect, it, vi } from 'vitest' +import { + InMemoryConnectionStore, + InMemoryIntegrationGrantStore, + InMemoryIntegrationWorkflowStore, + IntegrationHub, + createIntegrationRuntime, + createIntegrationWorkflowRuntime, + createMockIntegrationProvider, + type IntegrationGrant, + type IntegrationManifest, + type IntegrationTriggerEvent, +} from '../src/index' + +const owner = { type: 'user' as const, id: 'owner' } +const grantee = { type: 'agent' as const, id: 'agent' } +const manifest: IntegrationManifest = { + id: 'inbound-agent', + requirements: [{ + id: 'gmail-trigger', connectorId: 'gmail', mode: 'trigger', + reason: 'Wake the authorized agent when mail arrives.', + requiredTriggers: ['message.received'], + }], +} + +async function setup() { + const grants = new InMemoryIntegrationGrantStore() + const store = new InMemoryIntegrationWorkflowStore() + const hub = new IntegrationHub({ + providers: [createMockIntegrationProvider()], + store: new InMemoryConnectionStore(), + capabilitySecret: 'workflow-test-only', + }) + await hub.upsertConnection({ + id: 'connection', owner, providerId: 'mock', connectorId: 'gmail', + status: 'active', grantedScopes: ['email.read'], + createdAt: new Date(0).toISOString(), updatedAt: new Date(0).toISOString(), + }) + const runtime = createIntegrationWorkflowRuntime({ + runtime: createIntegrationRuntime({ hub, grants }), hub, grants, store, + }) + const installed = await runtime.install({ + workflow: { + id: 'incoming-mail', manifest, + trigger: { requirementId: 'gmail-trigger', triggerId: 'message.received' }, + }, + owner, grantee, + }) + const grant = grants.get(installed.triggerGrantId)! + const event: IntegrationTriggerEvent = { + id: 'event', providerId: 'mock', connectorId: 'gmail', + connectionId: 'connection', trigger: 'message.received', + occurredAt: new Date(0).toISOString(), payload: { subject: 'hello' }, + } + return { runtime, grants, store, installed, grant, event } +} + +const changedGrants: Array<[string, (grant: IntegrationGrant) => IntegrationGrant]> = [ + ['revoked grant', (g) => ({ ...g, status: 'revoked' })], + ['wrong manifest', (g) => ({ ...g, manifestId: 'other' })], + ['changed owner id', (g) => ({ ...g, owner: { ...g.owner, id: 'other' } })], + ['changed owner type', (g) => ({ ...g, owner: { ...g.owner, type: 'team' } })], + ['changed grantee id', (g) => ({ ...g, grantee: { ...g.grantee, id: 'other' } })], + ['changed grantee type', (g) => ({ ...g, grantee: { ...g.grantee, type: 'app' } })], + ['different connection', (g) => ({ ...g, connectionId: 'other' })], + ['different connector', (g) => ({ ...g, connectorId: 'other' })], + ['removed trigger', (g) => ({ ...g, allowedTriggers: [] })], +] + +describe('workflow dispatch authorization', () => { + it('dispatches an active binding and awaits the consumer', async () => { + const f = await setup() + let completed = false + const result = await f.runtime.dispatchEvent(f.event, async ({ workflows }) => { + expect(workflows.map((w) => w.id)).toEqual([f.installed.id]) + await Promise.resolve() + completed = true + }) + expect(result.matched).toEqual([f.installed]) + expect(completed).toBe(true) + }) + + it('rechecks the grant after a previously successful dispatch', async () => { + const f = await setup() + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toHaveLength(1) + f.grants.put({ ...f.grant, status: 'revoked' }) + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) + }) + + it.each(changedGrants)('excludes a %s', async (_name, change) => { + const f = await setup() + f.grants.put(change(f.grant)) + const handler = vi.fn() + expect((await f.runtime.dispatchEvent(f.event, handler)).matched).toEqual([]) + // Preserve the existing empty-match callback contract. + expect(handler).toHaveBeenCalledExactlyOnceWith({ event: f.event, workflows: [] }) + }) + + it('excludes a deleted grant', async () => { + const f = await setup() + f.grants.delete(f.grant.id) + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) + }) + + it('rejects a store result for a different grant id', async () => { + const f = await setup() + vi.spyOn(f.grants, 'get').mockReturnValue({ ...f.grant, id: 'other' }) + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) + }) + + it.each(['paused', 'error'] as const)('excludes a %s subscription', async (status) => { + const f = await setup() + f.store.put({ ...f.installed, subscription: { ...f.installed.subscription, status } }) + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) + }) + + it.each(['paused', 'error'] as const)('excludes a %s workflow', async (status) => { + const f = await setup() + f.store.put({ ...f.installed, status }) + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) + }) + + it.each(['connectionId', 'trigger'] as const)('does not read grants for an unrelated %s', async (field) => { + const f = await setup() + const get = vi.spyOn(f.grants, 'get') + expect((await f.runtime.dispatchEvent({ ...f.event, [field]: 'other' }, () => {})).matched).toEqual([]) + expect(get).not.toHaveBeenCalled() + }) + + it('does not route a different connector sharing a connection id', async () => { + const f = await setup() + expect((await f.runtime.dispatchEvent({ ...f.event, connectorId: 'other' }, () => {})).matched).toEqual([]) + }) + + it('propagates grant-store failures without delivering partial matches', async () => { + const f = await setup() + vi.spyOn(f.grants, 'get').mockImplementation(() => { throw new Error('store unavailable') }) + const handler = vi.fn() + await expect(f.runtime.dispatchEvent(f.event, handler)).rejects.toThrow('store unavailable') + expect(handler).not.toHaveBeenCalled() + }) + + it('propagates the consumer failure for retry', async () => { + const f = await setup() + await expect(f.runtime.dispatchEvent(f.event, () => { throw new Error('enqueue failed') })).rejects.toThrow('enqueue failed') + }) + + it('does not cache grant state across dispatches', async () => { + const f = await setup() + const get = vi.spyOn(f.grants, 'get') + await f.runtime.dispatchEvent(f.event, () => {}) + await f.runtime.dispatchEvent(f.event, () => {}) + expect(get).toHaveBeenCalledTimes(2) + }) +}) From 1ab5db508ef1d9e1160ffa5fd156fe7b636431a3 Mon Sep 17 00:00:00 2001 From: drewstone Date: Sat, 12 Sep 2026 15:13:56 -0700 Subject: [PATCH 3/4] ci: verify pull requests without publishing or deployment permissions --- .github/workflows/verify-pr.yml | 41 +++++++++++++++++++++++++++++++++ 1 file changed, 41 insertions(+) create mode 100644 .github/workflows/verify-pr.yml diff --git a/.github/workflows/verify-pr.yml b/.github/workflows/verify-pr.yml new file mode 100644 index 00000000..661c253c --- /dev/null +++ b/.github/workflows/verify-pr.yml @@ -0,0 +1,41 @@ +name: Verify PR + +on: + pull_request: + types: [opened, synchronize, reopened, ready_for_review] + +permissions: + contents: read + +concurrency: + group: verify-pr-${{ github.event.pull_request.number }} + cancel-in-progress: true + +jobs: + verify: + runs-on: ubuntu-latest + timeout-minutes: 25 + steps: + - uses: actions/checkout@v4 + with: + persist-credentials: false + - uses: pnpm/action-setup@v4 + - uses: actions/setup-node@v4 + with: + node-version: 22 + cache: pnpm + - run: pnpm install --frozen-lockfile + - run: pnpm run typecheck + - run: pnpm run test + - run: pnpm run release + - name: Record tested source + run: | + mkdir -p verification + git rev-parse HEAD > verification/commit.txt + git archive --format=tar.gz -o verification/source.tar.gz HEAD + sha256sum verification/source.tar.gz > verification/SHA256SUMS + - uses: actions/upload-artifact@v4 + with: + name: verified-source-${{ github.event.pull_request.number }} + path: verification/ + retention-days: 7 From caed0f3fee6b78c389bb6e3b2b5ce8d6e7d6ff5e Mon Sep 17 00:00:00 2001 From: Drew Stone Date: Fri, 18 Sep 2026 00:18:20 -0600 Subject: [PATCH 4/4] perf(workflows): batch trigger grant reads before dispatch Read every candidate workflow's current grant in one listByIds call, falling back to concurrent point reads for stores without it, and skip the read when no workflow matches the event. --- src/workflow.ts | 56 ++++++++++++---------- tests/workflow-authorization.test.ts | 72 ++++++++++++++++++++++++---- 2 files changed, 95 insertions(+), 33 deletions(-) diff --git a/src/workflow.ts b/src/workflow.ts index 17cf3313..eb35ed7d 100644 --- a/src/workflow.ts +++ b/src/workflow.ts @@ -135,30 +135,36 @@ export class IntegrationWorkflowRuntime { // grant before routing a new event; an active subscription is not authority. // This assumes verified ingress. Deferred effects must recheck authority // again at execution time rather than treating this snapshot as a lease. - const workflows: InstalledIntegrationWorkflow[] = [] - for (const workflow of await this.store.list()) { - if ( - workflow.status !== 'active' - || workflow.subscription.status !== 'active' - || workflow.subscription.connectionId !== event.connectionId - || workflow.subscription.trigger !== event.trigger - ) continue - - const grant = await this.grants.get(workflow.triggerGrantId) - if ( - !grant - || grant.id !== workflow.triggerGrantId - || grant.status !== 'active' - || grant.manifestId !== workflow.manifestId - || !sameActor(grant.owner, workflow.owner) - || !sameActor(grant.grantee, workflow.grantee) - || grant.connectionId !== event.connectionId - || grant.connectorId !== event.connectorId - || !grant.allowedTriggers.includes(event.trigger) - ) continue - - workflows.push(workflow) - } + const candidates = (await this.store.list()).filter((workflow) => + workflow.status === 'active' + && workflow.subscription.status === 'active' + && workflow.subscription.connectionId === event.connectionId + && workflow.subscription.trigger === event.trigger + ) + // One batched read keeps webhook latency flat as a trigger fans out to + // many workflows backed by a remote grant store. + const grantIds = [...new Set(candidates.map((workflow) => workflow.triggerGrantId))] + const grants = grantIds.length === 0 + ? [] + : this.grants.listByIds + ? await this.grants.listByIds(grantIds) + : await Promise.all(grantIds.map((grantId) => this.grants.get(grantId))) + const grantsById = new Map( + grants + .filter((grant): grant is IntegrationGrant => Boolean(grant)) + .map((grant) => [grant.id, grant]), + ) + const workflows = candidates.filter((workflow) => { + const grant = grantsById.get(workflow.triggerGrantId) + if (!grant) return false + return grant.status === 'active' + && grant.manifestId === workflow.manifestId + && sameActor(grant.owner, workflow.owner) + && sameActor(grant.grantee, workflow.grantee) + && grant.connectionId === event.connectionId + && grant.connectorId === event.connectorId + && grant.allowedTriggers.includes(event.trigger) + }) await handler({ event, workflows }) return { matched: workflows } } @@ -178,4 +184,4 @@ function findTriggerGrant(grants: IntegrationGrant[], requirementId: string, tri function sameActor(a: IntegrationActor, b: IntegrationActor): boolean { return a.type === b.type && a.id === b.id -} \ No newline at end of file +} diff --git a/tests/workflow-authorization.test.ts b/tests/workflow-authorization.test.ts index a9abe47d..ae4e24ea 100644 --- a/tests/workflow-authorization.test.ts +++ b/tests/workflow-authorization.test.ts @@ -8,6 +8,7 @@ import { createIntegrationWorkflowRuntime, createMockIntegrationProvider, type IntegrationGrant, + type IntegrationGrantStore, type IntegrationManifest, type IntegrationTriggerEvent, } from '../src/index' @@ -23,8 +24,20 @@ const manifest: IntegrationManifest = { }], } -async function setup() { +// A store without the optional batch read, as a minimal persistent adapter may be. +function pointReadStore(inner: InMemoryIntegrationGrantStore): IntegrationGrantStore { + return { + get: (grantId) => inner.get(grantId), + put: (grant) => inner.put(grant), + listByManifest: (manifestId, grantee) => inner.listByManifest(manifestId, grantee), + listByGrantee: (grantee) => inner.listByGrantee(grantee), + delete: (grantId) => inner.delete(grantId), + } +} + +async function setup(options: { pointReads?: boolean } = {}) { const grants = new InMemoryIntegrationGrantStore() + const runtimeGrants = options.pointReads ? pointReadStore(grants) : grants const store = new InMemoryIntegrationWorkflowStore() const hub = new IntegrationHub({ providers: [createMockIntegrationProvider()], @@ -37,7 +50,7 @@ async function setup() { createdAt: new Date(0).toISOString(), updatedAt: new Date(0).toISOString(), }) const runtime = createIntegrationWorkflowRuntime({ - runtime: createIntegrationRuntime({ hub, grants }), hub, grants, store, + runtime: createIntegrationRuntime({ hub, grants: runtimeGrants }), hub, grants: runtimeGrants, store, }) const installed = await runtime.install({ workflow: { @@ -52,7 +65,15 @@ async function setup() { connectionId: 'connection', trigger: 'message.received', occurredAt: new Date(0).toISOString(), payload: { subject: 'hello' }, } - return { runtime, grants, store, installed, grant, event } + // A distinct manifest yields a distinct trigger grant for the same connection. + const install = (id: string) => runtime.install({ + workflow: { + id, manifest: { ...manifest, id: `${manifest.id}-${id}` }, + trigger: { requirementId: 'gmail-trigger', triggerId: 'message.received' }, + }, + owner, grantee, + }) + return { runtime, grants, store, installed, grant, event, install } } const changedGrants: Array<[string, (grant: IntegrationGrant) => IntegrationGrant]> = [ @@ -104,6 +125,12 @@ describe('workflow dispatch authorization', () => { it('rejects a store result for a different grant id', async () => { const f = await setup() + vi.spyOn(f.grants, 'listByIds').mockReturnValue([{ ...f.grant, id: 'other' }]) + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) + }) + + it('rejects a point-read result for a different grant id', async () => { + const f = await setup({ pointReads: true }) vi.spyOn(f.grants, 'get').mockReturnValue({ ...f.grant, id: 'other' }) expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([]) }) @@ -123,8 +150,34 @@ describe('workflow dispatch authorization', () => { it.each(['connectionId', 'trigger'] as const)('does not read grants for an unrelated %s', async (field) => { const f = await setup() const get = vi.spyOn(f.grants, 'get') + const listByIds = vi.spyOn(f.grants, 'listByIds') expect((await f.runtime.dispatchEvent({ ...f.event, [field]: 'other' }, () => {})).matched).toEqual([]) expect(get).not.toHaveBeenCalled() + expect(listByIds).not.toHaveBeenCalled() + }) + + it('reads every candidate grant in one batched call', async () => { + const f = await setup() + const second = await f.install('second-mail') + const third = await f.install('third-mail') + const get = vi.spyOn(f.grants, 'get') + const listByIds = vi.spyOn(f.grants, 'listByIds') + const { matched } = await f.runtime.dispatchEvent(f.event, () => {}) + expect(matched.map((w) => w.id).sort()).toEqual([f.installed.id, second.id, third.id].sort()) + expect(listByIds).toHaveBeenCalledOnce() + expect(listByIds.mock.calls[0]![0].sort()).toEqual( + [f.installed.triggerGrantId, second.triggerGrantId, third.triggerGrantId].sort(), + ) + expect(get).not.toHaveBeenCalled() + }) + + it('falls back to concurrent point reads without a batch read', async () => { + const f = await setup({ pointReads: true }) + const second = await f.install('second-mail') + f.grants.put({ ...f.grants.get(second.triggerGrantId)!, status: 'revoked' }) + const get = vi.spyOn(f.grants, 'get') + expect((await f.runtime.dispatchEvent(f.event, () => {})).matched).toEqual([f.installed]) + expect(get).toHaveBeenCalledTimes(2) }) it('does not route a different connector sharing a connection id', async () => { @@ -132,9 +185,12 @@ describe('workflow dispatch authorization', () => { expect((await f.runtime.dispatchEvent({ ...f.event, connectorId: 'other' }, () => {})).matched).toEqual([]) }) - it('propagates grant-store failures without delivering partial matches', async () => { - const f = await setup() - vi.spyOn(f.grants, 'get').mockImplementation(() => { throw new Error('store unavailable') }) + it.each([ + ['batch', false, 'listByIds'], + ['point', true, 'get'], + ] as const)('propagates %s grant-store failures without delivering partial matches', async (_name, pointReads, method) => { + const f = await setup({ pointReads }) + vi.spyOn(f.grants, method).mockImplementation(() => { throw new Error('store unavailable') }) const handler = vi.fn() await expect(f.runtime.dispatchEvent(f.event, handler)).rejects.toThrow('store unavailable') expect(handler).not.toHaveBeenCalled() @@ -147,9 +203,9 @@ describe('workflow dispatch authorization', () => { it('does not cache grant state across dispatches', async () => { const f = await setup() - const get = vi.spyOn(f.grants, 'get') + const listByIds = vi.spyOn(f.grants, 'listByIds') await f.runtime.dispatchEvent(f.event, () => {}) await f.runtime.dispatchEvent(f.event, () => {}) - expect(get).toHaveBeenCalledTimes(2) + expect(listByIds).toHaveBeenCalledTimes(2) }) })