diff --git a/docs/contributing/architecture/primitives.yaml b/docs/contributing/architecture/primitives.yaml index dbfdc2bde7..0e9ca704ed 100644 --- a/docs/contributing/architecture/primitives.yaml +++ b/docs/contributing/architecture/primitives.yaml @@ -601,6 +601,19 @@ primitives: - docs/contributing/architecture/authorization.md - docs/guides/package-subscriptions.md + - id: package-events-dispatch-queue + group: storage + name: Package events dispatch queue + summary: + Durable same-user delivery of package-emitted events (kody.emits / + events.dispatch) to subscribing packages. + code: + - packages/worker/src/package-events/ + - packages/worker/src/package-invocations/subscription-dispatch.ts + docs: + - docs/contributing/setup-manifest.md + - docs/guides/package-subscriptions.md + - id: community-assets-r2 group: storage name: Community asset storage (R2) diff --git a/docs/contributing/setup-manifest.md b/docs/contributing/setup-manifest.md index b8e6639b43..cadc5bc313 100644 --- a/docs/contributing/setup-manifest.md +++ b/docs/contributing/setup-manifest.md @@ -63,6 +63,20 @@ This project uses the following resources: reloads the metadata-only activity projection, acknowledges invalid or deleted activity, and retries transient lookup, subscription-discovery, or package-invocation infrastructure failures. +- Cloudflare Queue for durable package-emitted event dispatch + - Producer binding: `PACKAGE_EVENTS_DISPATCH_QUEUE` + - Queue: `kody-package-events-dispatch` + - Dead-letter queue: `kody-package-events-dispatch-dlq` + - The production consumer uses the same batch, retry, and DLQ settings as + platform-feedback dispatch. Production CI ensures both resources. + - Queue messages carry the full event (emitting user, source package, topic, + idempotency key, payload, and invocation depth). The consumer resolves the + emitting user's subscribed packages at delivery time, invokes handlers with + exactly-once idempotency, acknowledges terminal handler failures, and + retries pre-execution package-invocation infrastructure failures. + - Preview and local runtimes without this production-only queue binding + deliver inline through the same consumer code path so package events remain + testable. - Cloudflare Queue for isolated scheduled maintenance - Producer binding: `SCHEDULED_DISPATCH_QUEUE` - Queue: `kody-scheduled-dispatch` diff --git a/docs/guides/package-subscriptions.md b/docs/guides/package-subscriptions.md index eaffcc2530..66919c23eb 100644 --- a/docs/guides/package-subscriptions.md +++ b/docs/guides/package-subscriptions.md @@ -56,6 +56,130 @@ The result lists package id, `kody.id`, package name, topic, handler, description, and filters. Use this before debugging event dispatch, building fan-out, or deciding whether a package already subscribes to a topic. +## Package-emitted topics (`@scope/...`) + +Packages can define their own event topics and emit to them; every other package +saved by the same user that declares the topic in `kody.subscriptions` receives +the event. There is no cross-user delivery. + +### Declaring emitted topics + +Declare topics in `package.json#kody.emits`. Topics must use the scoped form +`@{username}/topic.name` with a lower-dot-case body, and the scope must match +the emitting package's npm scope: + +```json +{ + "name": "@kentcdodds/discord-gateway", + "kody": { + "id": "discord-gateway", + "description": "Discord gateway.", + "emits": { + "@kentcdodds/discord.message.created": { + "description": "A Discord message was created.", + "payloadSchema": { + "type": "object", + "properties": { + "messageId": { "type": "string", "minLength": 1 }, + "channelId": { "type": "string" } + }, + "required": ["messageId", "channelId"], + "additionalProperties": false + } + } + } + } +} +``` + +`payloadSchema` is optional. When present it must be a JSON Schema subset with +root `"type": "object"`; supported keywords are `type`, `description`, +`properties`, `required`, `additionalProperties` (boolean), `items`, `enum`, +`const`, `minLength`, `maxLength`, `minimum`, `maximum`, `minItems`, and +`maxItems`. Unsupported keywords fail package checks at publish time so authors +never rely on silently ignored constraints. Declared schemas appear in package +search/detail projections so subscribers can discover payload shapes. + +### Emitting + +Emit from any package runtime context (exports, subscription handlers, +package-owned jobs, services, apps, retrievers) with the `events` helper: + +```ts +import { events } from 'kody:runtime' + +await events.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: `discord:message-create:${message.id}`, + payload: { messageId: message.id, channelId: message.channelId }, +}) +``` + +Rules: + +- The topic must be declared in the emitting package's `kody.emits`. +- `idempotencyKey` is required; payloads must be JSON objects and are validated + against `payloadSchema` when declared. +- Payloads are capped at 64 KiB (canonical JSON). Store large data with + `packageStorage()` and emit a reference instead. +- `events.dispatch` is unavailable in ad hoc `execute` runs — topics belong to + packages, so emit from package code (or `packages.invoke` into one). + +### Delivery semantics + +Dispatch is asynchronous and durable: `events.dispatch` validates the event, +enqueues it on the `kody-package-events-dispatch` Queue (with DLQ), and returns +`{ topic, source, idempotencyKey, status: "enqueued" }` immediately. Emitters +never observe subscriber results or latency; check each subscriber's run records +for handler outcomes. + +The Queue consumer resolves the emitting user's subscribed packages at delivery +time and invokes each `subscription:@scope/topic` handler with: + +```ts +type PackageEventEnvelope = { + event: string + source: { type: 'package'; package_id: string; kody_id: string } + idempotency_key: string + payload: Record +} +``` + +- Per-subscriber invocations are exactly-once keyed on + `(source package, subscriber package, topic, idempotencyKey)`, so Queue + redelivery replays stored results instead of re-running handlers. +- Infrastructure failures before handler code runs retry via the Queue (3 + attempts, then the `kody-package-events-dispatch-dlq` dead-letter queue). + Terminal handler failures do not retry — a stored failed invocation replays + rather than re-running — and stay visible in run records. +- Event-driven chains carry the same nested invocation depth budget as + `packages.invoke` (max 8 hops), so emit cycles between packages terminate. +- In environments without the Queue binding (local dev, preview) — or when an + enqueue fails — dispatch falls back to inline delivery with the same consumer + code path and reports `status: "delivered_inline"` instead of `"enqueued"`. + +### Filters on package-emitted topics + +A subscription to a package-emitted topic may declare `filters`; every filter +key must be present in the event payload with an equal JSON value or the +subscriber is skipped: + +```json +{ + "kody": { + "subscriptions": { + "@kentcdodds/discord.message.created": { + "handler": "./src/on-general-chat-message.ts", + "filters": { "channelId": "1470913684598423592" } + } + } + } +} +``` + +Platform-owned topics (below) keep their existing behavior: their dispatchers +define whether and how `filters` apply. + ## `email.message.received` Accepted stored inbound email dispatches `email.message.received` after Kody diff --git a/packages/shared/src/json-schema-subset.node.test.ts b/packages/shared/src/json-schema-subset.node.test.ts new file mode 100644 index 0000000000..1a1cd07bf8 --- /dev/null +++ b/packages/shared/src/json-schema-subset.node.test.ts @@ -0,0 +1,128 @@ +import { expect, test } from 'vitest' +import { + listJsonSchemaSubsetProblems, + listJsonSchemaSubsetValueErrors, +} from './json-schema-subset.ts' + +test('listJsonSchemaSubsetProblems accepts the supported subset', () => { + expect( + listJsonSchemaSubsetProblems({ + type: 'object', + description: 'A Discord message event payload.', + properties: { + messageId: { type: 'string', minLength: 1 }, + channelId: { type: 'string' }, + kind: { enum: ['created', 'updated'] }, + attachments: { + type: 'array', + maxItems: 10, + items: { type: 'object', properties: { url: { type: 'string' } } }, + }, + priority: { type: ['integer', 'null'], minimum: 0, maximum: 5 }, + }, + required: ['messageId', 'channelId'], + additionalProperties: false, + }), + ).toEqual([]) +}) + +test('listJsonSchemaSubsetProblems rejects unsupported keywords and shapes', () => { + expect(listJsonSchemaSubsetProblems('nope')).toEqual([ + '# must be a JSON object schema.', + ]) + expect( + listJsonSchemaSubsetProblems({ type: 'object', $ref: '#/defs/x' }), + ).toEqual([expect.stringContaining('unsupported keyword "$ref"')]) + expect(listJsonSchemaSubsetProblems({ type: 'uuid' })).toEqual([ + expect.stringContaining('"uuid" is not supported'), + ]) + expect( + listJsonSchemaSubsetProblems({ + type: 'object', + properties: { nested: { pattern: '^a' } }, + }), + ).toEqual([expect.stringContaining('#.properties["nested"]')]) + expect( + listJsonSchemaSubsetProblems({ type: 'object', additionalProperties: {} }), + ).toEqual([expect.stringContaining('additionalProperties must be a boolean')]) + expect(listJsonSchemaSubsetProblems({ enum: [] })).toEqual([ + expect.stringContaining('enum must be a non-empty array'), + ]) + expect(listJsonSchemaSubsetProblems({ minLength: -1 })).toEqual([ + expect.stringContaining('minLength must be a non-negative integer'), + ]) +}) + +test('listJsonSchemaSubsetValueErrors validates values against the subset', () => { + const schema = { + type: 'object', + properties: { + messageId: { type: 'string', minLength: 1 }, + kind: { enum: ['created', 'updated'] }, + count: { type: 'integer', minimum: 0, maximum: 10 }, + tags: { type: 'array', maxItems: 2, items: { type: 'string' } }, + }, + required: ['messageId'], + additionalProperties: false, + } + + expect( + listJsonSchemaSubsetValueErrors(schema, { + messageId: 'm-1', + kind: 'created', + count: 3, + tags: ['a', 'b'], + }), + ).toEqual([]) + + expect(listJsonSchemaSubsetValueErrors(schema, {})).toEqual([ + 'payload is missing required property "messageId".', + ]) + expect( + listJsonSchemaSubsetValueErrors(schema, { messageId: '', kind: 'nope' }), + ).toEqual([ + 'payload.messageId must have at least 1 characters.', + 'payload.kind must be one of "created", "updated".', + ]) + expect( + listJsonSchemaSubsetValueErrors(schema, { messageId: 'm', count: 3.5 }), + ).toEqual(['payload.count must have type integer.']) + expect( + listJsonSchemaSubsetValueErrors(schema, { + messageId: 'm', + tags: ['a', 'b', 'c'], + }), + ).toEqual(['payload.tags must have at most 2 items.']) + expect( + listJsonSchemaSubsetValueErrors(schema, { messageId: 'm', tags: [1] }), + ).toEqual(['payload.tags[0] must have type string.']) + expect( + listJsonSchemaSubsetValueErrors(schema, { messageId: 'm', extra: true }), + ).toEqual(['payload has unexpected property "extra".']) + // Missing `properties` means an empty declared set, so + // additionalProperties: false rejects every key. + expect( + listJsonSchemaSubsetValueErrors( + { type: 'object', additionalProperties: false }, + { anything: 1 }, + ), + ).toEqual(['payload has unexpected property "anything".']) +}) + +test('listJsonSchemaSubsetValueErrors handles const and union types', () => { + expect( + listJsonSchemaSubsetValueErrors( + { const: { a: 1, b: [2] } }, + { b: [2], a: 1 }, + ), + ).toEqual([]) + expect(listJsonSchemaSubsetValueErrors({ const: 'x' }, 'y')).toEqual([ + 'payload must equal the const value "x".', + ]) + expect( + listJsonSchemaSubsetValueErrors({ type: ['string', 'null'] }, null), + ).toEqual([]) + expect( + listJsonSchemaSubsetValueErrors({ type: ['string', 'null'] }, 5), + ).toEqual(['payload must have type string | null.']) +}) diff --git a/packages/shared/src/json-schema-subset.ts b/packages/shared/src/json-schema-subset.ts new file mode 100644 index 0000000000..92ee22b259 --- /dev/null +++ b/packages/shared/src/json-schema-subset.ts @@ -0,0 +1,287 @@ +import { canonicalJsonStringify } from './canonical-json.ts' + +/** + * Deliberately small JSON Schema subset for package-event payload contracts. + * Only the keywords below are supported; schemas using anything else are + * rejected at publish time so authors never depend on silently ignored + * constraints. Validation is structural and side-effect free (no `$ref`, + * no `pattern`, no format registry), which keeps it safe to run against + * untrusted payloads in the dispatch path. + */ +export type JsonSchemaSubset = Record + +const supportedKeywords = new Set([ + 'type', + 'description', + 'properties', + 'required', + 'additionalProperties', + 'items', + 'enum', + 'const', + 'minLength', + 'maxLength', + 'minimum', + 'maximum', + 'minItems', + 'maxItems', +]) + +const supportedTypes = new Set([ + 'object', + 'array', + 'string', + 'number', + 'integer', + 'boolean', + 'null', +]) + +function isPlainObject(value: unknown): value is Record { + return Boolean(value) && typeof value === 'object' && !Array.isArray(value) +} + +function isNonNegativeInteger(value: unknown): value is number { + return typeof value === 'number' && Number.isInteger(value) && value >= 0 +} + +function listSchemaTypes(schema: Record): Array { + const type = schema['type'] + if (typeof type === 'string') return [type] + if (Array.isArray(type)) { + return type.filter((entry): entry is string => typeof entry === 'string') + } + return [] +} + +/** + * Validate a schema definition against the supported subset. Returns an + * empty array when the schema is usable with + * {@link listJsonSchemaSubsetValueErrors}. + */ +export function listJsonSchemaSubsetProblems( + schema: unknown, + path = '#', +): Array { + if (!isPlainObject(schema)) { + return [`${path} must be a JSON object schema.`] + } + const problems: Array = [] + for (const keyword of Object.keys(schema)) { + if (!supportedKeywords.has(keyword)) { + problems.push( + `${path} uses unsupported keyword "${keyword}". Supported keywords: ${[...supportedKeywords].join(', ')}.`, + ) + } + } + if ('type' in schema) { + const type = schema['type'] + const types = + typeof type === 'string' ? [type] : Array.isArray(type) ? type : null + if (!types || types.length === 0) { + problems.push(`${path}.type must be a type name or array of type names.`) + } else { + for (const entry of types) { + if (typeof entry !== 'string' || !supportedTypes.has(entry)) { + problems.push( + `${path}.type "${String(entry)}" is not supported. Supported types: ${[...supportedTypes].join(', ')}.`, + ) + } + } + } + } + if ('description' in schema && typeof schema['description'] !== 'string') { + problems.push(`${path}.description must be a string.`) + } + if ('properties' in schema) { + const properties = schema['properties'] + if (!isPlainObject(properties)) { + problems.push(`${path}.properties must be an object of schemas.`) + } else { + for (const [key, propertySchema] of Object.entries(properties)) { + problems.push( + ...listJsonSchemaSubsetProblems( + propertySchema, + `${path}.properties["${key}"]`, + ), + ) + } + } + } + if ('required' in schema) { + const required = schema['required'] + if ( + !Array.isArray(required) || + required.some((entry) => typeof entry !== 'string') + ) { + problems.push(`${path}.required must be an array of property names.`) + } + } + if ( + 'additionalProperties' in schema && + typeof schema['additionalProperties'] !== 'boolean' + ) { + problems.push( + `${path}.additionalProperties must be a boolean in the supported subset.`, + ) + } + if ('items' in schema) { + problems.push( + ...listJsonSchemaSubsetProblems(schema['items'], `${path}.items`), + ) + } + if ('enum' in schema) { + const enumValues = schema['enum'] + if (!Array.isArray(enumValues) || enumValues.length === 0) { + problems.push(`${path}.enum must be a non-empty array of values.`) + } + } + for (const keyword of ['minLength', 'maxLength', 'minItems', 'maxItems']) { + if (keyword in schema && !isNonNegativeInteger(schema[keyword])) { + problems.push(`${path}.${keyword} must be a non-negative integer.`) + } + } + for (const keyword of ['minimum', 'maximum']) { + if ( + keyword in schema && + (typeof schema[keyword] !== 'number' || !Number.isFinite(schema[keyword])) + ) { + problems.push(`${path}.${keyword} must be a finite number.`) + } + } + return problems +} + +function matchesType(type: string, value: unknown): boolean { + switch (type) { + case 'object': + return isPlainObject(value) + case 'array': + return Array.isArray(value) + case 'string': + return typeof value === 'string' + case 'number': + return typeof value === 'number' && Number.isFinite(value) + case 'integer': + return typeof value === 'number' && Number.isInteger(value) + case 'boolean': + return typeof value === 'boolean' + case 'null': + return value === null + default: + return false + } +} + +/** + * Validate a JSON value against a subset schema. Callers must have checked + * the schema with {@link listJsonSchemaSubsetProblems} first; unsupported + * keywords are ignored here. + */ +export function listJsonSchemaSubsetValueErrors( + schema: JsonSchemaSubset, + value: unknown, + path = 'payload', +): Array { + const errors: Array = [] + const types = listSchemaTypes(schema) + if (types.length > 0 && !types.some((type) => matchesType(type, value))) { + errors.push(`${path} must have type ${types.join(' | ')}.`) + return errors + } + if ('const' in schema) { + if ( + canonicalJsonStringify(value) !== canonicalJsonStringify(schema['const']) + ) { + errors.push( + `${path} must equal the const value ${canonicalJsonStringify(schema['const'])}.`, + ) + } + } + if (Array.isArray(schema['enum'])) { + const candidate = canonicalJsonStringify(value) + if ( + !schema['enum'].some( + (entry) => canonicalJsonStringify(entry) === candidate, + ) + ) { + errors.push( + `${path} must be one of ${schema['enum'] + .map((entry) => canonicalJsonStringify(entry)) + .join(', ')}.`, + ) + } + } + if (typeof value === 'string') { + const minLength = schema['minLength'] + if (isNonNegativeInteger(minLength) && value.length < minLength) { + errors.push(`${path} must have at least ${minLength} characters.`) + } + const maxLength = schema['maxLength'] + if (isNonNegativeInteger(maxLength) && value.length > maxLength) { + errors.push(`${path} must have at most ${maxLength} characters.`) + } + } + if (typeof value === 'number') { + const minimum = schema['minimum'] + if (typeof minimum === 'number' && value < minimum) { + errors.push(`${path} must be >= ${minimum}.`) + } + const maximum = schema['maximum'] + if (typeof maximum === 'number' && value > maximum) { + errors.push(`${path} must be <= ${maximum}.`) + } + } + if (Array.isArray(value)) { + const minItems = schema['minItems'] + if (isNonNegativeInteger(minItems) && value.length < minItems) { + errors.push(`${path} must have at least ${minItems} items.`) + } + const maxItems = schema['maxItems'] + if (isNonNegativeInteger(maxItems) && value.length > maxItems) { + errors.push(`${path} must have at most ${maxItems} items.`) + } + const items = schema['items'] + if (isPlainObject(items)) { + for (const [index, entry] of value.entries()) { + errors.push( + ...listJsonSchemaSubsetValueErrors(items, entry, `${path}[${index}]`), + ) + } + } + } + if (isPlainObject(value)) { + const required = schema['required'] + if (Array.isArray(required)) { + for (const key of required) { + if (typeof key === 'string' && !(key in value)) { + errors.push(`${path} is missing required property "${key}".`) + } + } + } + // Missing `properties` means an empty declared set, so + // `additionalProperties: false` rejects every key (standard JSON + // Schema behavior). + const properties = isPlainObject(schema['properties']) + ? schema['properties'] + : {} + for (const [key, propertySchema] of Object.entries(properties)) { + if (!(key in value) || !isPlainObject(propertySchema)) continue + errors.push( + ...listJsonSchemaSubsetValueErrors( + propertySchema, + value[key], + `${path}.${key}`, + ), + ) + } + if (schema['additionalProperties'] === false) { + for (const key of Object.keys(value)) { + if (!(key in properties)) { + errors.push(`${path} has unexpected property "${key}".`) + } + } + } + } + return errors +} diff --git a/packages/worker/src/package-events/dispatch-queue-names.ts b/packages/worker/src/package-events/dispatch-queue-names.ts new file mode 100644 index 0000000000..a9c88764bd --- /dev/null +++ b/packages/worker/src/package-events/dispatch-queue-names.ts @@ -0,0 +1,4 @@ +export const packageEventsDispatchQueueBinding = 'PACKAGE_EVENTS_DISPATCH_QUEUE' +export const packageEventsDispatchQueueName = 'kody-package-events-dispatch' +export const packageEventsDispatchDeadLetterQueueName = + 'kody-package-events-dispatch-dlq' diff --git a/packages/worker/src/package-events/dispatch-queue-producer.ts b/packages/worker/src/package-events/dispatch-queue-producer.ts new file mode 100644 index 0000000000..fe5277905b --- /dev/null +++ b/packages/worker/src/package-events/dispatch-queue-producer.ts @@ -0,0 +1,67 @@ +export type PackageEventsDispatchQueueMessage = { + userId: string + topic: string + idempotencyKey: string + payload: Record + source: { + packageId: string + kodyId: string + } + /** + * Runtime invocation depth carried across the queue boundary so + * event-driven package chains (A emits, B's handler emits, ...) keep the + * same cycle protection as synchronous packages.invoke chains. + */ + invokeDepth: number +} + +export function parsePackageEventsDispatchQueueMessage( + body: unknown, +): PackageEventsDispatchQueueMessage | null { + if (!body || typeof body !== 'object' || Array.isArray(body)) return null + const record = body as Record + const userId = record['userId'] + const topic = record['topic'] + const idempotencyKey = record['idempotencyKey'] + const payload = record['payload'] + const source = record['source'] + const invokeDepth = record['invokeDepth'] + if ( + typeof userId !== 'string' || + !userId.trim() || + typeof topic !== 'string' || + !topic.trim() || + typeof idempotencyKey !== 'string' || + !idempotencyKey.trim() || + !payload || + typeof payload !== 'object' || + Array.isArray(payload) || + !source || + typeof source !== 'object' || + Array.isArray(source) || + typeof invokeDepth !== 'number' || + !Number.isInteger(invokeDepth) || + invokeDepth < 0 + ) { + return null + } + const sourceRecord = source as Record + const packageId = sourceRecord['packageId'] + const kodyId = sourceRecord['kodyId'] + if ( + typeof packageId !== 'string' || + !packageId.trim() || + typeof kodyId !== 'string' || + !kodyId.trim() + ) { + return null + } + return { + userId: userId.trim(), + topic: topic.trim(), + idempotencyKey: idempotencyKey.trim(), + payload: payload as Record, + source: { packageId: packageId.trim(), kodyId: kodyId.trim() }, + invokeDepth, + } +} diff --git a/packages/worker/src/package-events/dispatch-queue.node.test.ts b/packages/worker/src/package-events/dispatch-queue.node.test.ts new file mode 100644 index 0000000000..0674d07ae0 --- /dev/null +++ b/packages/worker/src/package-events/dispatch-queue.node.test.ts @@ -0,0 +1,170 @@ +import { expect, test, vi } from 'vitest' +import { consoleError } from '#worker/test-support/console-spies.ts' +import { parsePackageEventsDispatchQueueMessage } from './dispatch-queue-producer.ts' + +const mocks = vi.hoisted(() => ({ + deliverPackageEvent: + vi.fn< + (input: Record) => Promise> + >(), +})) + +vi.mock('#worker/package-invocations/service.ts', () => ({ + deliverPackageEvent: mocks.deliverPackageEvent, +})) + +vi.mock('#worker/app-base-url.ts', () => ({ + getAppBaseUrl: () => 'https://kody.dev', +})) + +const { handlePackageEventsDispatchQueue } = await import('./dispatch-queue.ts') + +function createMessageBody(overrides: Record = {}) { + return { + userId: 'user-123', + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:123', + payload: { messageId: '123' }, + source: { packageId: 'pkg-gateway', kodyId: 'discord-gateway' }, + invokeDepth: 1, + ...overrides, + } +} + +function createQueueMessage(id: string, body: unknown) { + return { + id, + timestamp: new Date('2026-08-04T00:01:00.000Z'), + body, + attempts: 1, + ack: vi.fn<() => void>(), + retry: vi.fn<(options?: { delaySeconds?: number }) => void>(), + } +} + +function createBatch(messages: Array>) { + return { + queue: 'kody-package-events-dispatch', + messages, + ackAll: vi.fn<() => void>(), + retryAll: vi.fn<() => void>(), + } as unknown as MessageBatch +} + +test('package events queue delivers valid messages and acks invalid ones', async () => { + consoleError.mockImplementation(() => {}) + const valid = createQueueMessage('queue-valid', createMessageBody()) + const invalid = createQueueMessage('queue-invalid', { topic: ' ' }) + const handlerFailure = createQueueMessage( + 'queue-handler-failure', + createMessageBody({ idempotencyKey: 'discord:handler-failure' }), + ) + const infrastructureFailure = createQueueMessage( + 'queue-infra-failure', + createMessageBody({ idempotencyKey: 'discord:infra-failure' }), + ) + mocks.deliverPackageEvent + .mockResolvedValueOnce({ delivered: 1, failed: 0, subscribers: [] }) + // Terminal handler failures are final (the ledger replays them on + // retry), so the message still acks. + .mockResolvedValueOnce({ delivered: 0, failed: 1, subscribers: [] }) + .mockRejectedValueOnce(new Error('Package event dispatch was incomplete.')) + + await handlePackageEventsDispatchQueue( + createBatch([valid, invalid, handlerFailure, infrastructureFailure]), + { APP_DB: {} } as Env, + { + waitUntil: vi.fn<(promise: Promise) => void>(), + } as unknown as ExecutionContext, + ) + + expect(mocks.deliverPackageEvent).toHaveBeenCalledTimes(3) + expect(mocks.deliverPackageEvent).toHaveBeenNthCalledWith(1, { + env: expect.anything(), + baseUrl: 'https://kody.dev', + message: { + userId: 'user-123', + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:123', + payload: { messageId: '123' }, + source: { packageId: 'pkg-gateway', kodyId: 'discord-gateway' }, + invokeDepth: 1, + }, + waitUntil: expect.any(Function), + }) + for (const message of [valid, invalid, handlerFailure]) { + expect(message.ack).toHaveBeenCalledTimes(1) + expect(message.retry).not.toHaveBeenCalled() + } + expect(infrastructureFailure.ack).not.toHaveBeenCalled() + expect(infrastructureFailure.retry).toHaveBeenCalledWith({ + delaySeconds: 30, + }) + expect(consoleError).toHaveBeenCalledWith( + 'package-events-dispatch-message-invalid', + expect.objectContaining({ queueMessageId: 'queue-invalid' }), + ) + expect(consoleError).toHaveBeenCalledWith( + 'package-events-dispatch-subscribers-failed', + expect.objectContaining({ + queueMessageId: 'queue-handler-failure', + failed: 1, + }), + ) + expect(consoleError).toHaveBeenCalledWith( + 'package-events-dispatch-queue-processing-failed', + expect.objectContaining({ + queueMessageId: 'queue-infra-failure', + error: expect.any(Error), + }), + ) +}) + +test('package events queue message parsing rejects malformed bodies', () => { + expect(parsePackageEventsDispatchQueueMessage(createMessageBody())).toEqual( + createMessageBody(), + ) + expect(parsePackageEventsDispatchQueueMessage(null)).toBeNull() + expect(parsePackageEventsDispatchQueueMessage('nope')).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage(createMessageBody({ topic: ' ' })), + ).toBeNull() + // userId selects whose packages receive the event, so it enforces + // per-user isolation across the queue boundary. + expect( + parsePackageEventsDispatchQueueMessage(createMessageBody({ userId: ' ' })), + ).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage(createMessageBody({ userId: 42 })), + ).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage( + createMessageBody({ userId: ' user-123 ', topic: ' topic.a ' }), + ), + ).toMatchObject({ userId: 'user-123', topic: 'topic.a' }) + expect( + parsePackageEventsDispatchQueueMessage( + createMessageBody({ idempotencyKey: '' }), + ), + ).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage( + createMessageBody({ payload: ['nope'] }), + ), + ).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage( + createMessageBody({ source: { packageId: 'pkg-1' } }), + ), + ).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage( + createMessageBody({ invokeDepth: -1 }), + ), + ).toBeNull() + expect( + parsePackageEventsDispatchQueueMessage( + createMessageBody({ invokeDepth: 1.5 }), + ), + ).toBeNull() +}) diff --git a/packages/worker/src/package-events/dispatch-queue.ts b/packages/worker/src/package-events/dispatch-queue.ts new file mode 100644 index 0000000000..dbe6fb4720 --- /dev/null +++ b/packages/worker/src/package-events/dispatch-queue.ts @@ -0,0 +1,55 @@ +import { getAppBaseUrl } from '#worker/app-base-url.ts' +import { deliverPackageEvent } from '#worker/package-invocations/service.ts' +import { parsePackageEventsDispatchQueueMessage } from './dispatch-queue-producer.ts' + +const packageEventsDispatchRetryDelaySeconds = 30 + +export async function handlePackageEventsDispatchQueue( + batch: MessageBatch, + env: Env, + ctx: ExecutionContext, +) { + const baseUrl = getAppBaseUrl({ env }) + for (const queueMessage of batch.messages) { + const message = parsePackageEventsDispatchQueueMessage(queueMessage.body) + if (!message) { + console.error('package-events-dispatch-message-invalid', { + queueMessageId: queueMessage.id, + }) + queueMessage.ack() + continue + } + try { + const result = await deliverPackageEvent({ + env, + baseUrl, + message, + waitUntil: (promise) => ctx.waitUntil(promise), + }) + // Terminal handler failures are final for this delivery: the + // idempotency ledger would replay the stored failure on retry, so + // redelivery cannot help. They stay visible in each subscriber's + // run records. + if (result.failed > 0) { + console.error('package-events-dispatch-subscribers-failed', { + queueMessageId: queueMessage.id, + topic: message.topic, + sourcePackageId: message.source.packageId, + failed: result.failed, + delivered: result.delivered, + }) + } + queueMessage.ack() + } catch (error) { + console.error('package-events-dispatch-queue-processing-failed', { + queueMessageId: queueMessage.id, + topic: message.topic, + sourcePackageId: message.source.packageId, + error, + }) + queueMessage.retry({ + delaySeconds: packageEventsDispatchRetryDelaySeconds, + }) + } + } +} diff --git a/packages/worker/src/package-invocations/admin-package-subscriptions.ts b/packages/worker/src/package-invocations/admin-package-subscriptions.ts index 34a54125ce..6dda3f88cd 100644 --- a/packages/worker/src/package-invocations/admin-package-subscriptions.ts +++ b/packages/worker/src/package-invocations/admin-package-subscriptions.ts @@ -5,63 +5,23 @@ import { listPackageSubscriptions } from '#worker/package-registry/manifest.ts' import { listSavedPackagesByUserId } from '#worker/package-registry/repo.ts' import { loadPackageManifestBySourceId } from '#worker/package-registry/source.ts' import { type SavedPackageRecord } from '#worker/package-registry/types.ts' +import { + readPreExecutionPackageInvocationInfrastructureCode, + readRetryablePackageInvocationInfrastructureCode, +} from './infrastructure-codes.ts' import { invokePackageSubscription } from './service.ts' +export { + readPreExecutionPackageInvocationInfrastructureCode, + readRetryablePackageInvocationInfrastructureCode, +} + type LoadedAdminPackageSubscription = { savedPackage: SavedPackageRecord subscription: ReturnType[number] } const adminPackageSubscriptionConcurrency = 5 -const retryablePackageInvocationInfrastructureCodes = new Set([ - 'idempotency_lookup_failed', - 'idempotency_persistence_failed', - 'idempotency_conflict_unresolved', - 'invocation_in_progress', - 'invocation_failed', - 'idempotency_response_unavailable', -]) -const preExecutionPackageInvocationInfrastructureCodes = new Set([ - 'idempotency_lookup_failed', - 'idempotency_persistence_failed', - 'idempotency_conflict_unresolved', - 'invocation_in_progress', - 'artifact_preparation_failed', -]) - -function readPackageInvocationInfrastructureCode(input: { - response: { - status: number - body: Record - } - codes: ReadonlySet -}) { - if (input.response.status >= 200 && input.response.status < 300) return null - const error = input.response.body['error'] - if (!error || typeof error !== 'object' || Array.isArray(error)) return null - const code = (error as Record)['code'] - return typeof code === 'string' && input.codes.has(code) ? code : null -} - -export function readRetryablePackageInvocationInfrastructureCode(response: { - status: number - body: Record -}) { - return readPackageInvocationInfrastructureCode({ - response, - codes: retryablePackageInvocationInfrastructureCodes, - }) -} - -export function readPreExecutionPackageInvocationInfrastructureCode(response: { - status: number - body: Record -}) { - return readPackageInvocationInfrastructureCode({ - response, - codes: preExecutionPackageInvocationInfrastructureCodes, - }) -} async function mapSettledInChunks( items: ReadonlyArray, diff --git a/packages/worker/src/package-invocations/infrastructure-codes.ts b/packages/worker/src/package-invocations/infrastructure-codes.ts new file mode 100644 index 0000000000..a2099d4837 --- /dev/null +++ b/packages/worker/src/package-invocations/infrastructure-codes.ts @@ -0,0 +1,57 @@ +/** + * Package-invocation error codes that indicate infrastructure trouble + * rather than a user-code failure. Neutral module (no service imports) so + * dispatchers like subscription-dispatch can classify responses without + * creating an import cycle through admin-package-subscriptions → service. + */ +const retryablePackageInvocationInfrastructureCodes = new Set([ + 'idempotency_lookup_failed', + 'idempotency_persistence_failed', + 'idempotency_conflict_unresolved', + 'invocation_in_progress', + 'invocation_failed', + 'idempotency_response_unavailable', +]) + +/** Codes raised before user code ran, so Queue redelivery re-executes. */ +const preExecutionPackageInvocationInfrastructureCodes = new Set([ + 'idempotency_lookup_failed', + 'idempotency_persistence_failed', + 'idempotency_conflict_unresolved', + 'invocation_in_progress', + 'artifact_preparation_failed', +]) + +function readPackageInvocationInfrastructureCode(input: { + response: { + status: number + body: Record + } + codes: ReadonlySet +}) { + if (input.response.status >= 200 && input.response.status < 300) return null + const error = input.response.body['error'] + if (!error || typeof error !== 'object' || Array.isArray(error)) return null + const code = (error as Record)['code'] + return typeof code === 'string' && input.codes.has(code) ? code : null +} + +export function readRetryablePackageInvocationInfrastructureCode(response: { + status: number + body: Record +}) { + return readPackageInvocationInfrastructureCode({ + response, + codes: retryablePackageInvocationInfrastructureCodes, + }) +} + +export function readPreExecutionPackageInvocationInfrastructureCode(response: { + status: number + body: Record +}) { + return readPackageInvocationInfrastructureCode({ + response, + codes: preExecutionPackageInvocationInfrastructureCodes, + }) +} diff --git a/packages/worker/src/package-invocations/service.node.test.ts b/packages/worker/src/package-invocations/service.node.test.ts index 040e6c7c37..300abffdf1 100644 --- a/packages/worker/src/package-invocations/service.node.test.ts +++ b/packages/worker/src/package-invocations/service.node.test.ts @@ -5,6 +5,7 @@ import { createExecutePackageInvokeTools, createPackageRuntimeInvokeTools, createPackageEventTools, + deliverPackageEvent, invokePackageExport, invokePackageSubscription, } from './service.ts' @@ -691,8 +692,18 @@ function createManifest(input: { kodyId: string exportName: string entryPoint: string - emits?: Record - subscriptions?: Record + emits?: Record< + string, + { description: string; payloadSchema?: Record } + > + subscriptions?: Record< + string, + { + handler: string + description?: string + filters?: Record + } + > }) { return { name: input.name, @@ -932,9 +943,18 @@ function createRuntimeDispatchTools(db: D1Database) { }) } -function createRuntimeEventTools(db: D1Database) { +function createRuntimeEventTools( + db: D1Database, + options: { + envOverrides?: Record + packageInvokeDepth?: number + } = {}, +) { return createPackageEventTools({ - env: createEnv(db), + env: { + ...(createEnv(db) as unknown as Record), + ...options.envOverrides, + } as unknown as Env, baseUrl: 'https://kody.dev', callerContext: createMcpCallerContext({ baseUrl: 'https://kody.dev', @@ -957,7 +977,7 @@ function createRuntimeEventTools(db: D1Database) { name: './dispatch-message-created', idempotencyKey: 'message-1', }, - packageInvokeDepth: 0, + packageInvokeDepth: options.packageInvokeDepth ?? 0, }) } @@ -1220,7 +1240,150 @@ test('runtime invoke tools expose only the supported invoke helper', () => { expect(tools.invokeChecked).toBeUndefined() }) -test('package runtime dispatches declared events to same-user package subscriptions', async () => { +test('package runtime dispatch enqueues declared events to the package events queue', async () => { + const db = createDatabase() + seedRuntimeDispatchPackages() + repoMockModule.runBundledModuleWithRegistry.mockClear() + const send = vi.fn<(message: unknown) => Promise>( + async () => undefined, + ) + const tools = createRuntimeEventTools(db, { + envOverrides: { PACKAGE_EVENTS_DISPATCH_QUEUE: { send } }, + }) + + const result = await tools.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:123', + payload: { + messageId: '123', + channelId: '456', + }, + }) + + expect(send).toHaveBeenCalledTimes(1) + expect(send).toHaveBeenCalledWith({ + userId: 'user-123', + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:123', + payload: { + messageId: '123', + channelId: '456', + }, + source: { + packageId: 'pkg-gateway', + kodyId: 'discord-gateway', + }, + invokeDepth: 1, + }) + // Queued delivery never invokes subscribers inside the emitting request. + expect(repoMockModule.runBundledModuleWithRegistry).not.toHaveBeenCalled() + expect(result).toEqual({ + topic: '@kentcdodds/discord.message.created', + source: { + type: 'package', + packageId: 'pkg-gateway', + kodyId: 'discord-gateway', + }, + idempotencyKey: 'discord:message-create:123', + status: 'enqueued', + }) + + await expect( + tools.dispatch({ + topic: '@kentcdodds/discord.reaction.created', + idempotencyKey: 'discord:reaction-create:123', + payload: { + reactionId: '123', + }, + }), + ).rejects.toThrow( + /does not declare emitted event "@kentcdodds\/discord.reaction.created"/, + ) + expect(send).toHaveBeenCalledTimes(1) + + await expect( + tools.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:huge', + payload: { blob: 'x'.repeat(65 * 1024) }, + }), + ).rejects.toThrow(/payload is \d+ bytes; the maximum is 65536 bytes/) + expect(send).toHaveBeenCalledTimes(1) + + const depthCappedTools = createRuntimeEventTools(db, { + envOverrides: { PACKAGE_EVENTS_DISPATCH_QUEUE: { send } }, + packageInvokeDepth: 8, + }) + await expect( + depthCappedTools.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:deep', + payload: {}, + }), + ).rejects.toThrow(/exceeded the maximum nested invocation depth/) + expect(send).toHaveBeenCalledTimes(1) +}) + +test('package runtime dispatch validates payloads against the declared payloadSchema', async () => { + const db = createDatabase() + const { manifests, sourceFiles } = seedRuntimeDispatchPackages() + const gatewayManifest = manifests.get('source-gateway') as { + kody: { + emits?: Record< + string, + { description: string; payloadSchema?: Record } + > + } + } + gatewayManifest.kody.emits = { + '@kentcdodds/discord.message.created': { + description: 'A Discord message was created.', + payloadSchema: { + type: 'object', + properties: { + messageId: { type: 'string', minLength: 1 }, + channelId: { type: 'string' }, + }, + required: ['messageId'], + additionalProperties: false, + }, + }, + } + // Keep the serialized source snapshot in sync with the mutated manifest + // so every manifest loader observes the same emits declaration. + const gatewayFiles = sourceFiles.get('source-gateway') + if (gatewayFiles) { + gatewayFiles['package.json'] = JSON.stringify(gatewayManifest) + } + const send = vi.fn<(message: unknown) => Promise>( + async () => undefined, + ) + const tools = createRuntimeEventTools(db, { + envOverrides: { PACKAGE_EVENTS_DISPATCH_QUEUE: { send } }, + }) + + await expect( + tools.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:bad', + payload: { channelId: '456', extra: true }, + }), + ).rejects.toThrow( + /payload does not match the declared payloadSchema[\s\S]*missing required property "messageId"[\s\S]*unexpected property "extra"/, + ) + expect(send).not.toHaveBeenCalled() + + await expect( + tools.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:ok', + payload: { messageId: '123', channelId: '456' }, + }), + ).resolves.toMatchObject({ status: 'enqueued' }) + expect(send).toHaveBeenCalledTimes(1) +}) + +test('package events deliver to same-user package subscriptions with idempotent replay', async () => { const db = createDatabase() seedRuntimeDispatchPackages() repoMockModule.runBundledModuleWithRegistry.mockClear() @@ -1240,23 +1403,30 @@ test('package runtime dispatches declared events to same-user package subscripti logs: [], }), ) - const tools = createRuntimeEventTools(db) - - const first = await tools.dispatch({ + const message = { + userId: 'user-123', topic: '@kentcdodds/discord.message.created', idempotencyKey: 'discord:message-create:123', payload: { messageId: '123', channelId: '456', }, - }) - const second = await tools.dispatch({ - topic: '@kentcdodds/discord.message.created', - idempotencyKey: 'discord:message-create:123', - payload: { - messageId: '123', - channelId: '456', + source: { + packageId: 'pkg-gateway', + kodyId: 'discord-gateway', }, + invokeDepth: 1, + } + + const first = await deliverPackageEvent({ + env: createEnv(db), + baseUrl: 'https://kody.dev', + message, + }) + const second = await deliverPackageEvent({ + env: createEnv(db), + baseUrl: 'https://kody.dev', + message, }) expect(repoMockModule.runBundledModuleWithRegistry).toHaveBeenCalledTimes(1) @@ -1306,25 +1476,118 @@ test('package runtime dispatches declared events to same-user package subscripti }) const { manifests, sources } = seedRuntimeDispatchPackages() + repoMockModule.loadPackageManifestBySourceId.mockImplementation( + async (input: { sourceId: string }) => { + if (input.sourceId === 'source-subscriber') { + throw new Error('manifest unavailable') + } + return { + source: sources.get(input.sourceId), + manifest: manifests.get(input.sourceId), + } + }, + ) + await expect( + deliverPackageEvent({ + env: createEnv(db), + baseUrl: 'https://kody.dev', + message: { ...message, idempotencyKey: 'discord:manifest-error' }, + }), + ).rejects.toThrow( + /Failed to load package manifest for package event dispatch/, + ) +}) + +test('package event subscription filters gate delivery on payload values', async () => { + const db = createDatabase() + const { manifests, sourceFiles } = seedRuntimeDispatchPackages() + const subscriberManifest = manifests.get('source-subscriber') as { + kody: { + subscriptions?: Record< + string, + { handler: string; filters?: Record } + > + } + } + subscriberManifest.kody.subscriptions = { + '@kentcdodds/discord.message.created': { + handler: './src/handle-discord-message-created.ts', + filters: { channelId: '456' }, + }, + } + const subscriberFiles = sourceFiles.get('source-subscriber') + if (subscriberFiles) { + subscriberFiles['package.json'] = JSON.stringify(subscriberManifest) + } repoMockModule.runBundledModuleWithRegistry.mockClear() repoMockModule.runBundledModuleWithRegistry.mockResolvedValue({ - result: { ok: true }, + result: { handled: true }, logs: [], }) + const baseMessage = { + userId: 'user-123', + topic: '@kentcdodds/discord.message.created', + payload: {}, + source: { + packageId: 'pkg-gateway', + kodyId: 'discord-gateway', + }, + invokeDepth: 1, + } - await expect( - tools.dispatch({ - topic: '@kentcdodds/discord.reaction.created', - idempotencyKey: 'discord:reaction-create:123', - payload: { - reactionId: '123', - }, - }), - ).rejects.toThrow( - /does not declare emitted event "@kentcdodds\/discord.reaction.created"/, - ) + const filteredOut = await deliverPackageEvent({ + env: createEnv(db), + baseUrl: 'https://kody.dev', + message: { + ...baseMessage, + idempotencyKey: 'discord:other-channel', + payload: { messageId: '1', channelId: '999' }, + }, + }) + expect(filteredOut).toMatchObject({ delivered: 0, failed: 0 }) + expect(filteredOut.subscribers).toEqual([]) expect(repoMockModule.runBundledModuleWithRegistry).not.toHaveBeenCalled() + const matching = await deliverPackageEvent({ + env: createEnv(db), + baseUrl: 'https://kody.dev', + message: { + ...baseMessage, + idempotencyKey: 'discord:matching-channel', + payload: { messageId: '2', channelId: '456' }, + }, + }) + expect(matching).toMatchObject({ delivered: 1, failed: 0 }) + expect(repoMockModule.runBundledModuleWithRegistry).toHaveBeenCalledTimes(1) +}) + +test('package events fall back to inline delivery without a queue binding', async () => { + const db = createDatabase() + seedRuntimeDispatchPackages() + repoMockModule.runBundledModuleWithRegistry.mockClear() + repoMockModule.runBundledModuleWithRegistry.mockResolvedValue({ + result: { handled: true }, + logs: [], + }) + const tools = createRuntimeEventTools(db) + + const result = await tools.dispatch({ + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:inline', + payload: { messageId: 'inline-1' }, + }) + + expect(result).toMatchObject({ status: 'delivered_inline' }) + expect(repoMockModule.runBundledModuleWithRegistry).toHaveBeenCalledTimes(1) + expect( + repoMockModule.runBundledModuleWithRegistry.mock.calls[0]?.[3], + ).toMatchObject({ + event: '@kentcdodds/discord.message.created', + payload: { messageId: 'inline-1' }, + }) + + // Inline delivery failures are logged, never surfaced to the emitter. + const { manifests, sources } = seedRuntimeDispatchPackages() repoMockModule.loadPackageManifestBySourceId.mockImplementation( async (input: { sourceId: string }) => { if (input.sourceId === 'source-subscriber') { @@ -1336,18 +1599,45 @@ test('package runtime dispatches declared events to same-user package subscripti } }, ) - + consoleError.mockImplementation(() => {}) await expect( tools.dispatch({ topic: '@kentcdodds/discord.message.created', - idempotencyKey: 'discord:message-create:manifest-error', - payload: { - messageId: '123', - }, + idempotencyKey: 'discord:message-create:inline-error', + payload: { messageId: 'inline-2' }, + }), + ).resolves.toMatchObject({ status: 'delivered_inline' }) + expect(consoleError).toHaveBeenCalledWith( + 'package-events-inline-delivery-failed', + expect.objectContaining({ + topic: '@kentcdodds/discord.message.created', }), - ).rejects.toThrow( - /Failed to load package manifest for package event dispatch/, ) +}) + +test('deliverPackageEvent surfaces retryable infrastructure failures', async () => { + const db = createDatabase({ failClaim: true }) + seedRuntimeDispatchPackages() + repoMockModule.runBundledModuleWithRegistry.mockClear() + consoleError.mockImplementation(() => {}) + + await expect( + deliverPackageEvent({ + env: createEnv(db), + baseUrl: 'https://kody.dev', + message: { + userId: 'user-123', + topic: '@kentcdodds/discord.message.created', + idempotencyKey: 'discord:message-create:claim-failure', + payload: { messageId: '1' }, + source: { + packageId: 'pkg-gateway', + kodyId: 'discord-gateway', + }, + invokeDepth: 1, + }, + }), + ).rejects.toThrow('Package event dispatch was incomplete.') expect(repoMockModule.runBundledModuleWithRegistry).not.toHaveBeenCalled() }) diff --git a/packages/worker/src/package-invocations/service.ts b/packages/worker/src/package-invocations/service.ts index 5d5a623f87..807086e38b 100644 --- a/packages/worker/src/package-invocations/service.ts +++ b/packages/worker/src/package-invocations/service.ts @@ -19,8 +19,10 @@ import { createExecutePackageInvokeToolsWithToolFactories, createPackageRuntimeInvokeToolsWithToolFactories, } from './runtime-tool-factories.ts' +import { type PackageEventsDispatchQueueMessage } from '#worker/package-events/dispatch-queue-producer.ts' import { createPackageEventToolsWithToolFactories, + deliverPackageEventWithToolFactories, invokePackageSubscriptionWithToolFactories, } from './subscription-dispatch.ts' @@ -115,6 +117,18 @@ export async function invokePackageExport(input: { }) } +export async function deliverPackageEvent(input: { + env: Env + baseUrl: string + message: PackageEventsDispatchQueueMessage + waitUntil?: (promise: Promise) => void +}) { + return await deliverPackageEventWithToolFactories({ + ...input, + toolFactories: packageRuntimeToolFactories, + }) +} + export async function invokePackageSubscription(input: { env: Env baseUrl: string diff --git a/packages/worker/src/package-invocations/subscription-dispatch.ts b/packages/worker/src/package-invocations/subscription-dispatch.ts index bd67e728b8..f96cc4ac8c 100644 --- a/packages/worker/src/package-invocations/subscription-dispatch.ts +++ b/packages/worker/src/package-invocations/subscription-dispatch.ts @@ -1,5 +1,7 @@ import { canonicalJsonStringify } from '@kody-internal/shared/canonical-json.ts' +import { listJsonSchemaSubsetValueErrors } from '@kody-internal/shared/json-schema-subset.ts' import { toHex } from '@kody-internal/shared/hex.ts' +import { runWithDynamicWorkerEvaluationBudget } from '#mcp/executor.ts' import { type createMcpCallerContext } from '#mcp/context.ts' import { type PackageEventDispatchInput, @@ -13,10 +15,12 @@ import { listPackageEmittedEvents, listPackageSubscriptions, } from '#worker/package-registry/manifest.ts' +import { type PackageEventsDispatchQueueMessage } from '#worker/package-events/dispatch-queue-producer.ts' import { buildPackageSubscriptionArtifactName, normalizePackageSubscriptionTopic, } from '#worker/package-runtime/subscription-artifacts.ts' +import { readPreExecutionPackageInvocationInfrastructureCode } from './infrastructure-codes.ts' import { internalEmailSubscriptionTokenId, internalPackageEventSubscriptionTokenId, @@ -30,6 +34,12 @@ import { import { invokeSavedPackageModule } from './idempotent-module-invocation.ts' import { buildJsonErrorResponse } from './responses.ts' +/** + * Package event payloads ride inside a Queue message (128 KiB limit), so + * the payload itself is capped well below that to leave envelope headroom. + */ +export const maxPackageEventPayloadBytes = 64 * 1024 + function parsePackageEventDispatchInput(rawInput: PackageEventDispatchInput) { const input = rawInput && typeof rawInput === 'object' && !Array.isArray(rawInput) @@ -91,11 +101,30 @@ function readInvocationError(response: PackageInvocationResponse) { } } +/** + * Filter semantics for package-emitted topics: every declared filter key + * must be present in the event payload with a canonically-equal JSON value. + * Subscriptions without filters receive every event on the topic. + */ +export function packageEventFiltersMatchPayload(input: { + filters: Record | null + payload: Record +}) { + if (!input.filters) return true + return Object.entries(input.filters).every( + ([key, expected]) => + key in input.payload && + canonicalJsonStringify(input.payload[key]) === + canonicalJsonStringify(expected), + ) +} + async function loadMatchingPackageEventSubscriptions(input: { env: Env baseUrl: string userId: string topic: string + payload: Record }) { const savedPackages = await listSavedPackagesByUserId(input.env.APP_DB, { userId: input.userId, @@ -114,7 +143,12 @@ async function loadMatchingPackageEventSubscriptions(input: { ) }) const subscription = listPackageSubscriptions(loaded.manifest).find( - (candidate) => candidate.topic === input.topic, + (candidate) => + candidate.topic === input.topic && + packageEventFiltersMatchPayload({ + filters: candidate.filters, + payload: input.payload, + }), ) if (!subscription) return null return { @@ -128,6 +162,128 @@ async function loadMatchingPackageEventSubscriptions(input: { ) } +export type PackageEventDeliveryResult = { + topic: string + source: { type: 'package'; packageId: string; kodyId: string } + idempotencyKey: string + subscribers: Array<{ + packageId: string + kodyId: string + handler: string + status: 'completed' | 'replayed' | 'failed' + error?: { code: string; message: string } + }> + delivered: number + failed: number +} + +/** + * Queue-consumer side of package event dispatch: resolve the emitting + * user's matching subscriptions and invoke each handler with exactly-once + * idempotency. Throws when subscriber discovery fails or an invocation + * fails before user code ran (both retryable via Queue redelivery); user + * handler failures are terminal and reported in the returned summary (the + * idempotency ledger would replay a stored failure on retry anyway). + */ +export async function deliverPackageEventWithToolFactories(input: { + env: Env + baseUrl: string + message: PackageEventsDispatchQueueMessage + toolFactories: PackageRuntimeToolFactories + waitUntil?: (promise: Promise) => void +}): Promise { + const message = input.message + const subscriptions = await loadMatchingPackageEventSubscriptions({ + env: input.env, + baseUrl: input.baseUrl, + userId: message.userId, + topic: message.topic, + payload: message.payload, + }) + const envelope = { + event: message.topic, + source: { + type: 'package', + package_id: message.source.packageId, + kody_id: message.source.kodyId, + }, + idempotency_key: message.idempotencyKey, + payload: message.payload, + } + const subscribers: PackageEventDeliveryResult['subscribers'] = [] + const retryableInfrastructureErrors: Array = [] + await runWithDynamicWorkerEvaluationBudget(async () => { + for (const { savedPackage, subscription } of subscriptions) { + const response = await invokePackageSubscriptionWithToolFactories({ + env: input.env, + baseUrl: input.baseUrl, + savedPackage, + topic: message.topic, + params: envelope, + idempotencyKey: await buildPackageEventSubscriptionIdempotencyKey({ + sourcePackageId: message.source.packageId, + subscriberPackageId: savedPackage.id, + topic: message.topic, + idempotencyKey: message.idempotencyKey, + }), + source: `package:${message.source.kodyId}`, + actorTokenId: `${internalPackageEventSubscriptionTokenId}:${message.source.packageId}`, + actorDisplayName: `package:${message.source.kodyId}`, + runtimeInvokeDepth: message.invokeDepth, + toolFactories: input.toolFactories, + waitUntil: input.waitUntil, + }) + const retryableCode = + readPreExecutionPackageInvocationInfrastructureCode(response) + if (retryableCode) { + retryableInfrastructureErrors.push( + new Error( + `Retryable package invocation infrastructure response: ${retryableCode}.`, + ), + ) + } + const replayed = + (response.body['idempotency'] as { replayed?: unknown } | undefined) + ?.replayed === true + const status = + response.status >= 200 && response.status < 400 + ? replayed + ? 'replayed' + : 'completed' + : 'failed' + subscribers.push({ + packageId: savedPackage.id, + kodyId: savedPackage.kodyId, + handler: subscription.handler, + status, + ...(status === 'failed' + ? { error: readInvocationError(response) } + : {}), + }) + } + }) + if (retryableInfrastructureErrors.length > 0) { + throw new Error('Package event dispatch was incomplete.', { + cause: retryableInfrastructureErrors[0], + }) + } + const failed = subscribers.filter( + (subscriber) => subscriber.status === 'failed', + ).length + return { + topic: message.topic, + source: { + type: 'package', + packageId: message.source.packageId, + kodyId: message.source.kodyId, + }, + idempotencyKey: message.idempotencyKey, + subscribers, + delivered: subscribers.length - failed, + failed, + } +} + export function createPackageEventToolsWithToolFactories(input: { env: Env baseUrl: string @@ -173,65 +329,74 @@ export function createPackageEventToolsWithToolFactories(input: { `Package "${packageContext.kodyId}" does not declare emitted event "${request.topic}" in package.json#kody.emits.`, ) } - const subscriptions = await loadMatchingPackageEventSubscriptions({ - env: input.env, - baseUrl: input.baseUrl, + if (declaredEvent.payloadSchema) { + const errors = listJsonSchemaSubsetValueErrors( + declaredEvent.payloadSchema, + request.payload, + ) + if (errors.length > 0) { + throw new Error( + `events.dispatch payload does not match the declared payloadSchema for "${request.topic}":\n${errors.join('\n')}`, + ) + } + } + const payloadBytes = new TextEncoder().encode( + canonicalJsonStringify(request.payload), + ).byteLength + if (payloadBytes > maxPackageEventPayloadBytes) { + throw new Error( + `events.dispatch payload is ${payloadBytes} bytes; the maximum is ${maxPackageEventPayloadBytes} bytes. Store large data in packageStorage() and emit a reference instead.`, + ) + } + const message: PackageEventsDispatchQueueMessage = { userId: user.userId, topic: request.topic, - }) - const envelope = { - event: request.topic, + idempotencyKey: request.idempotencyKey, + payload: request.payload, source: { - type: 'package', - package_id: packageContext.packageId, - kody_id: packageContext.kodyId, + packageId: packageContext.packageId, + kodyId: packageContext.kodyId, }, - idempotency_key: request.idempotencyKey, - payload: request.payload, + invokeDepth: packageInvokeDepth + 1, } - const subscribers = [] - for (const { savedPackage, subscription } of subscriptions) { - const response = await invokePackageSubscriptionWithToolFactories({ + const queue = (input.env as Partial).PACKAGE_EVENTS_DISPATCH_QUEUE + let enqueued = false + if (queue) { + try { + await queue.send(message) + enqueued = true + } catch (error) { + console.error('package-events-dispatch-enqueue-failed', { + topic: message.topic, + sourcePackageId: message.source.packageId, + error, + }) + } + } + if (!enqueued) { + // No Queue binding (local dev / preview) or enqueue failure: + // deliver inline with the same consumer code path so events + // still reach subscribers. Delivery failures are logged, never + // surfaced to the emitter — matching queued semantics. + const inlineDelivery = deliverPackageEventWithToolFactories({ env: input.env, baseUrl: input.baseUrl, - savedPackage, - topic: request.topic, - params: envelope, - idempotencyKey: await buildPackageEventSubscriptionIdempotencyKey({ - sourcePackageId: packageContext.packageId, - subscriberPackageId: savedPackage.id, - topic: request.topic, - idempotencyKey: request.idempotencyKey, - }), - source: `package:${packageContext.kodyId}`, - actorTokenId: `${internalPackageEventSubscriptionTokenId}:${packageContext.packageId}`, - actorDisplayName: `package:${packageContext.kodyId}`, - runtimeInvokeDepth: packageInvokeDepth + 1, + message, toolFactories: input.toolFactories, waitUntil: input.waitUntil, + }).catch((error) => { + console.error('package-events-inline-delivery-failed', { + topic: message.topic, + sourcePackageId: message.source.packageId, + error, + }) }) - const replayed = - (response.body['idempotency'] as { replayed?: unknown } | undefined) - ?.replayed === true - const status = - response.status >= 200 && response.status < 400 - ? replayed - ? 'replayed' - : 'completed' - : 'failed' - subscribers.push({ - packageId: savedPackage.id, - kodyId: savedPackage.kodyId, - handler: subscription.handler, - status, - ...(status === 'failed' - ? { error: readInvocationError(response) } - : {}), - }) + if (input.waitUntil) { + input.waitUntil(inlineDelivery) + } else { + await inlineDelivery + } } - const failed = subscribers.filter( - (subscriber) => subscriber.status === 'failed', - ).length return { topic: request.topic, source: { @@ -240,9 +405,7 @@ export function createPackageEventToolsWithToolFactories(input: { kodyId: packageContext.kodyId, }, idempotencyKey: request.idempotencyKey, - subscribers, - delivered: subscribers.length - failed, - failed, + status: enqueued ? 'enqueued' : 'delivered_inline', } }, } diff --git a/packages/worker/src/package-registry/manifest.node.test.ts b/packages/worker/src/package-registry/manifest.node.test.ts index 4b44a4abb5..29709a6bd2 100644 --- a/packages/worker/src/package-registry/manifest.node.test.ts +++ b/packages/worker/src/package-registry/manifest.node.test.ts @@ -1,6 +1,7 @@ import { expect, test } from 'vitest' import { buildPackageSearchProjection, + listPackageEmittedEvents, parseAuthoredPackageJson, } from './manifest.ts' import { @@ -180,6 +181,54 @@ test('parseAuthoredPackageJson accepts services, subscriptions, emits, retriever description: 'A Discord message was created.', }, }) + expect(listPackageEmittedEvents(manifest)).toEqual([ + { + topic: '@kentcdodds/discord.message.created', + description: 'A Discord message was created.', + payloadSchema: null, + }, + ]) + + const withPayloadSchema = parseAuthoredPackageJson({ + content: JSON.stringify({ + name: '@kentcdodds/discord-gateway', + exports: { + '.': './index.ts', + }, + kody: { + id: 'discord-gateway', + description: 'Discord gateway package', + emits: { + '@kentcdodds/discord.message.created': { + description: 'A Discord message was created.', + payloadSchema: { + type: 'object', + properties: { + messageId: { type: 'string', minLength: 1 }, + }, + required: ['messageId'], + additionalProperties: false, + }, + }, + }, + }, + }), + manifestPath: 'package.json', + }) + expect(listPackageEmittedEvents(withPayloadSchema)).toEqual([ + { + topic: '@kentcdodds/discord.message.created', + description: 'A Discord message was created.', + payloadSchema: { + type: 'object', + properties: { + messageId: { type: 'string', minLength: 1 }, + }, + required: ['messageId'], + additionalProperties: false, + }, + }, + ]) expect(buildPackageSearchProjection(manifest).retrievers).toEqual([ { key: 'notes-search', @@ -390,6 +439,55 @@ test('parseAuthoredPackageJson rejects unsupported or invalid kody manifest exte }), ).toThrow(/must use the package scope "@kentcdodds"/) + expect(() => + parseAuthoredPackageJson({ + content: JSON.stringify({ + name: '@kentcdodds/discord-gateway', + exports: { + '.': './index.ts', + }, + kody: { + id: 'discord-gateway', + description: 'Discord gateway package', + emits: { + '@kentcdodds/discord.message.created': { + description: 'A Discord message was created.', + payloadSchema: { type: 'string' }, + }, + }, + }, + }), + manifestPath: 'package.json', + }), + ).toThrow(/payloadSchema must declare "type": "object"/) + + expect(() => + parseAuthoredPackageJson({ + content: JSON.stringify({ + name: '@kentcdodds/discord-gateway', + exports: { + '.': './index.ts', + }, + kody: { + id: 'discord-gateway', + description: 'Discord gateway package', + emits: { + '@kentcdodds/discord.message.created': { + description: 'A Discord message was created.', + payloadSchema: { + type: 'object', + properties: { + messageId: { type: 'string', pattern: '^[0-9]+$' }, + }, + }, + }, + }, + }, + }), + manifestPath: 'package.json', + }), + ).toThrow(/payloadSchema is not a supported JSON Schema subset/) + expect(() => parseAuthoredPackageJson({ content: JSON.stringify({ diff --git a/packages/worker/src/package-registry/manifest.ts b/packages/worker/src/package-registry/manifest.ts index 5f0b9c872c..8d578c09a2 100644 --- a/packages/worker/src/package-registry/manifest.ts +++ b/packages/worker/src/package-registry/manifest.ts @@ -1,4 +1,5 @@ import { getErrorMessage } from '@kody-internal/shared/error-message.ts' +import { listJsonSchemaSubsetProblems } from '@kody-internal/shared/json-schema-subset.ts' import { z } from 'zod' import { parseModuleSource, type ModuleAstNode } from '#worker/module-source.ts' import { @@ -40,7 +41,9 @@ function assertPackageEmittedEventTopics(input: { manifestPath?: string }) { const packageScope = getPackageNameScope(input.manifest.name) - for (const topic of Object.keys(input.manifest.kody.emits ?? {})) { + for (const [topic, emittedEvent] of Object.entries( + input.manifest.kody.emits ?? {}, + )) { if (!customPackageEventTopicPattern.test(topic)) { throw new Error( `Invalid ${input.manifestPath ?? packageManifestPath}:\nkody.emits topic "${topic}" must use the scoped form "@scope/topic.name" with a lower-dot-case topic body.`, @@ -52,6 +55,20 @@ function assertPackageEmittedEventTopics(input: { `Invalid ${input.manifestPath ?? packageManifestPath}:\nkody.emits topic "${topic}" must use the package scope "@${packageScope}".`, ) } + if (emittedEvent.payloadSchema !== undefined) { + const rootType = emittedEvent.payloadSchema['type'] + if (rootType !== 'object') { + throw new Error( + `Invalid ${input.manifestPath ?? packageManifestPath}:\nkody.emits topic "${topic}" payloadSchema must declare "type": "object" (event payloads are JSON objects).`, + ) + } + const problems = listJsonSchemaSubsetProblems(emittedEvent.payloadSchema) + if (problems.length > 0) { + throw new Error( + `Invalid ${input.manifestPath ?? packageManifestPath}:\nkody.emits topic "${topic}" payloadSchema is not a supported JSON Schema subset:\n${problems.join('\n')}`, + ) + } + } } } @@ -287,6 +304,7 @@ export function listPackageEmittedEvents(manifest: AuthoredPackageJson) { .map(([topic, emittedEvent]) => ({ topic, description: emittedEvent.description.trim(), + payloadSchema: emittedEvent.payloadSchema ?? null, })) .sort((left, right) => left.topic.localeCompare(right.topic)) } @@ -376,6 +394,7 @@ export type PackageSearchProjection = { emits?: Array<{ topic: string description: string + payloadSchema?: Record | null }> retrievers: Array webhooks: Array diff --git a/packages/worker/src/package-registry/types.ts b/packages/worker/src/package-registry/types.ts index 5aa0b0944a..3541c4aa38 100644 --- a/packages/worker/src/package-registry/types.ts +++ b/packages/worker/src/package-registry/types.ts @@ -113,6 +113,10 @@ const packageWebhooksSchema = z export const packageEmittedEventDefinitionSchema = z.object({ description: z.string().min(1), + // JSON Schema subset for dispatch-time payload validation; the supported + // keyword set is enforced in parseAuthoredPackageJson (manifest.ts) so + // schema problems surface as publish-time manifest errors. + payloadSchema: z.record(z.string(), z.unknown()).optional(), }) export type PackageEmittedEventDefinition = z.infer< diff --git a/packages/worker/src/queue-handler.node.test.ts b/packages/worker/src/queue-handler.node.test.ts index b899b88937..a91901cac3 100644 --- a/packages/worker/src/queue-handler.node.test.ts +++ b/packages/worker/src/queue-handler.node.test.ts @@ -1,4 +1,5 @@ import { expect, test, vi } from 'vitest' +import { packageEventsDispatchQueueName } from '#worker/package-events/dispatch-queue-names.ts' import { scheduledDispatchQueueName } from '#worker/scheduled/scheduled-dispatch-queue-names.ts' import { consoleError } from '#worker/test-support/console-spies.ts' @@ -6,6 +7,7 @@ const mocks = vi.hoisted(() => ({ handleCommunityActivityDispatchQueue: vi.fn(), handleEmailDeliveryQueue: vi.fn(), handleArtifactsRepoEventsQueue: vi.fn(), + handlePackageEventsDispatchQueue: vi.fn(), handlePlatformFeedbackDispatchQueue: vi.fn(), handleScheduledDispatchQueue: vi.fn(), })) @@ -29,6 +31,10 @@ vi.mock('#worker/repo/artifacts-event-queue.ts', () => ({ handleArtifactsRepoEventsQueue: mocks.handleArtifactsRepoEventsQueue, })) +vi.mock('#worker/package-events/dispatch-queue.ts', () => ({ + handlePackageEventsDispatchQueue: mocks.handlePackageEventsDispatchQueue, +})) + vi.mock('#worker/platform-feedback/dispatch-queue.ts', () => ({ handlePlatformFeedbackDispatchQueue: mocks.handlePlatformFeedbackDispatchQueue, @@ -57,6 +63,7 @@ test('worker queue routing isolates known queues and retries unknown queues', as const artifactsBatch = createBatch('kody-artifacts-repo-events') const feedbackBatch = createBatch('kody-platform-feedback-dispatch') const communityActivityBatch = createBatch('kody-community-activity-dispatch') + const packageEventsBatch = createBatch(packageEventsDispatchQueueName) const scheduledBatch = createBatch(scheduledDispatchQueueName) const unknownBatch = createBatch('unexpected-queue') @@ -64,6 +71,7 @@ test('worker queue routing isolates known queues and retries unknown queues', as await handleQueueBatch(artifactsBatch, env, ctx) await handleQueueBatch(feedbackBatch, env, ctx) await handleQueueBatch(communityActivityBatch, env, ctx) + await handleQueueBatch(packageEventsBatch, env, ctx) await handleQueueBatch(scheduledBatch, env, ctx) await handleQueueBatch(unknownBatch, env, ctx) @@ -90,6 +98,12 @@ test('worker queue routing isolates known queues and retries unknown queues', as env, ctx, ) + expect(mocks.handlePackageEventsDispatchQueue).toHaveBeenCalledTimes(1) + expect(mocks.handlePackageEventsDispatchQueue).toHaveBeenCalledWith( + packageEventsBatch, + env, + ctx, + ) expect(mocks.handleScheduledDispatchQueue).toHaveBeenCalledWith( scheduledBatch, env, diff --git a/packages/worker/src/queue-handler.ts b/packages/worker/src/queue-handler.ts index 02f9d294e3..ecbf2b4c6d 100644 --- a/packages/worker/src/queue-handler.ts +++ b/packages/worker/src/queue-handler.ts @@ -6,6 +6,8 @@ import { handleCommunityActivityDispatchQueue } from '#worker/community/activity import { communityActivityDispatchQueueName } from '#worker/community/activity-dispatch-queue-names.ts' import { handlePlatformFeedbackDispatchQueue } from '#worker/platform-feedback/dispatch-queue.ts' import { platformFeedbackDispatchQueueName } from '#worker/platform-feedback/dispatch-queue-names.ts' +import { handlePackageEventsDispatchQueue } from '#worker/package-events/dispatch-queue.ts' +import { packageEventsDispatchQueueName } from '#worker/package-events/dispatch-queue-names.ts' import { artifactsRepoEventsQueueName, handleArtifactsRepoEventsQueue, @@ -33,6 +35,9 @@ export async function handleQueueBatch( case communityActivityDispatchQueueName: await handleCommunityActivityDispatchQueue(batch, env, ctx) return + case packageEventsDispatchQueueName: + await handlePackageEventsDispatchQueue(batch, env, ctx) + return case scheduledDispatchQueueName: await handleScheduledDispatchQueue(batch, env, ctx) return diff --git a/packages/worker/worker-configuration.d.ts b/packages/worker/worker-configuration.d.ts index f4dbd20388..f435ad2417 100644 --- a/packages/worker/worker-configuration.d.ts +++ b/packages/worker/worker-configuration.d.ts @@ -1,5 +1,5 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types --config=packages/worker/wrangler.jsonc --env=production ./packages/worker/worker-configuration.d.ts` (hash: 9d2951fc705f84c5a9c8be6c57c33e36) +// Generated by Wrangler by running `wrangler types --config=packages/worker/wrangler.jsonc --env=production ./packages/worker/worker-configuration.d.ts` (hash: ff721fcd9d867202d94fb7f9ffbd1993) // Runtime types generated with workerd@1.20260730.1 2026-04-16 global_fetch_strictly_public,nodejs_compat interface __BaseEnv_Env { OAUTH_KV: KVNamespace; @@ -16,6 +16,7 @@ interface __BaseEnv_Env { PLATFORM_FEEDBACK_DISPATCH_QUEUE: Queue; COMMUNITY_ACTIVITY_DISPATCH_QUEUE: Queue; SCHEDULED_DISPATCH_QUEUE: Queue; + PACKAGE_EVENTS_DISPATCH_QUEUE: Queue; AUTH_RATE_LIMITER: RateLimit; LOADER: WorkerLoader; APP_LOADER: WorkerLoader; diff --git a/packages/worker/wrangler.jsonc b/packages/worker/wrangler.jsonc index dc9c6ef11a..59fd6cac7a 100644 --- a/packages/worker/wrangler.jsonc +++ b/packages/worker/wrangler.jsonc @@ -273,6 +273,10 @@ "binding": "SCHEDULED_DISPATCH_QUEUE", "queue": "kody-scheduled-dispatch", }, + { + "binding": "PACKAGE_EVENTS_DISPATCH_QUEUE", + "queue": "kody-package-events-dispatch", + }, ], "consumers": [ { @@ -311,6 +315,17 @@ "max_concurrency": 16, "dead_letter_queue": "kody-scheduled-dispatch-dlq", }, + { + "queue": "kody-package-events-dispatch", + "max_batch_size": 10, + "max_batch_timeout": 5, + "max_retries": 3, + // Each event fans out to dynamic Worker isolate loads per + // subscriber; cap consumer concurrency so an emit burst + // cannot multiply into unbounded isolate loads. + "max_concurrency": 16, + "dead_letter_queue": "kody-package-events-dispatch-dlq", + }, ], }, "kv_namespaces": [ diff --git a/tools/ci/production-queue-resources.node.test.ts b/tools/ci/production-queue-resources.node.test.ts index 98acc5e166..8e48843de8 100644 --- a/tools/ci/production-queue-resources.node.test.ts +++ b/tools/ci/production-queue-resources.node.test.ts @@ -19,6 +19,10 @@ function createProductionEnv() { binding: 'SCHEDULED_DISPATCH_QUEUE', queue: 'kody-scheduled-dispatch', }, + { + binding: 'PACKAGE_EVENTS_DISPATCH_QUEUE', + queue: 'kody-package-events-dispatch', + }, ], consumers: [ { @@ -57,6 +61,14 @@ function createProductionEnv() { max_concurrency: 16, dead_letter_queue: 'kody-scheduled-dispatch-dlq', }, + { + queue: 'kody-package-events-dispatch', + max_batch_size: 10, + max_batch_timeout: 5, + max_retries: 3, + max_concurrency: 16, + dead_letter_queue: 'kody-package-events-dispatch-dlq', + }, ], }, } @@ -89,6 +101,9 @@ test('production queue config requires all consumers and consistent producers', 'kody-community-activity-dispatch-dlq', scheduledDispatchQueueName: 'kody-scheduled-dispatch', scheduledDispatchDeadLetterQueueName: 'kody-scheduled-dispatch-dlq', + packageEventsDispatchQueueName: 'kody-package-events-dispatch', + packageEventsDispatchDeadLetterQueueName: + 'kody-package-events-dispatch-dlq', }) const missingFeedbackConsumer = createProductionEnv() @@ -98,7 +113,7 @@ test('production queue config requires all consumers and consistent producers', productionEnv: missingFeedbackConsumer, configPath: 'wrangler.jsonc', }), - ).toThrow('exactly 5 production Queue consumers') + ).toThrow('exactly 6 production Queue consumers') const mismatchedProducer = createProductionEnv() mismatchedProducer.queues.producers[0] = { @@ -140,6 +155,20 @@ test('production queue config requires all consumers and consistent producers', }), ).toThrow('must bind "SCHEDULED_DISPATCH_QUEUE" to "kody-scheduled-dispatch"') + const missingPackageEventsProducer = createProductionEnv() + missingPackageEventsProducer.queues.producers = + missingPackageEventsProducer.queues.producers.filter( + (producer) => producer.binding !== 'PACKAGE_EVENTS_DISPATCH_QUEUE', + ) + expect(() => + parseProductionQueueResources({ + productionEnv: missingPackageEventsProducer, + configPath: 'wrangler.jsonc', + }), + ).toThrow( + 'must bind "PACKAGE_EVENTS_DISPATCH_QUEUE" to "kody-package-events-dispatch"', + ) + for (const invalidSettings of [ { max_batch_size: 2 }, { max_batch_timeout: 1 }, diff --git a/tools/ci/production-queue-resources.ts b/tools/ci/production-queue-resources.ts index 16289fb4a5..cb9e3e0477 100644 --- a/tools/ci/production-queue-resources.ts +++ b/tools/ci/production-queue-resources.ts @@ -3,6 +3,11 @@ import { communityActivityDispatchQueueBinding, communityActivityDispatchQueueName, } from '../../packages/worker/src/community/activity-dispatch-queue-names.ts' +import { + packageEventsDispatchDeadLetterQueueName, + packageEventsDispatchQueueBinding, + packageEventsDispatchQueueName, +} from '../../packages/worker/src/package-events/dispatch-queue-names.ts' import { platformFeedbackDispatchDeadLetterQueueName, platformFeedbackDispatchQueueBinding, @@ -23,7 +28,7 @@ const emailDeliveryDeadLetterQueueName = 'kody-email-delivery-dlq' const expectedMaxBatchSize = 10 const expectedMaxBatchTimeout = 5 const expectedMaxRetries = 3 -const expectedConsumerCount = 5 +const expectedConsumerCount = 6 function readQueueConsumer(input: { consumers: Array @@ -137,6 +142,13 @@ export function parseProductionQueueResources(input: { maxBatchTimeout: 0, maxConcurrency: 16, }) + const packageEventsDispatch = readQueueConsumer({ + consumers, + queueName: packageEventsDispatchQueueName, + deadLetterQueueName: packageEventsDispatchDeadLetterQueueName, + configPath: input.configPath, + maxConcurrency: 16, + }) const producers = queueConfig.producers if (!Array.isArray(producers)) { throw new Error( @@ -161,6 +173,12 @@ export function parseProductionQueueResources(input: { queueName: scheduledDispatchQueueName, configPath: input.configPath, }) + readQueueProducer({ + producers, + binding: packageEventsDispatchQueueBinding, + queueName: packageEventsDispatchQueueName, + configPath: input.configPath, + }) return { emailDeliveryQueueName: emailDelivery.queue, emailDeliveryDeadLetterQueueName: emailDelivery.deadLetterQueue, @@ -174,5 +192,8 @@ export function parseProductionQueueResources(input: { communityActivityDispatch.deadLetterQueue, scheduledDispatchQueueName: scheduledDispatch.queue, scheduledDispatchDeadLetterQueueName: scheduledDispatch.deadLetterQueue, + packageEventsDispatchQueueName: packageEventsDispatch.queue, + packageEventsDispatchDeadLetterQueueName: + packageEventsDispatch.deadLetterQueue, } } diff --git a/tools/ci/production-resources.ts b/tools/ci/production-resources.ts index ae55683f80..879653738e 100644 --- a/tools/ci/production-resources.ts +++ b/tools/ci/production-resources.ts @@ -46,6 +46,8 @@ type ResolvedProductionBindings = { communityActivityDispatchDeadLetterQueueName: string scheduledDispatchQueueName: string scheduledDispatchDeadLetterQueueName: string + packageEventsDispatchQueueName: string + packageEventsDispatchDeadLetterQueueName: string } function parseArgs(argv: Array): { @@ -443,7 +445,7 @@ async function ensureProductionResources(options: CliOptions) { kvTitleOverride: options.kvTitleOverride, }) console.error( - `Ensuring production resources for worker: ${bindings.workerName} (D1: ${bindings.d1DatabaseName}, OAuth KV: ${bindings.oauthKvTitle}, Bundle KV: ${bindings.bundleArtifactsKvTitle}, Community R2: ${bindings.communityAssetsBucketName}, Email R2: ${bindings.emailBlobsBucketName}, Email Queue: ${bindings.emailDeliveryQueueName}, Email DLQ: ${bindings.emailDeliveryDeadLetterQueueName}, Artifacts Repo Events Queue: ${bindings.artifactsRepoEventsQueueName}, Artifacts Repo Events DLQ: ${bindings.artifactsRepoEventsDeadLetterQueueName}, Platform Feedback Queue: ${bindings.platformFeedbackDispatchQueueName}, Platform Feedback DLQ: ${bindings.platformFeedbackDispatchDeadLetterQueueName}, Community Activity Queue: ${bindings.communityActivityDispatchQueueName}, Community Activity DLQ: ${bindings.communityActivityDispatchDeadLetterQueueName}, Scheduled Dispatch Queue: ${bindings.scheduledDispatchQueueName}, Scheduled Dispatch DLQ: ${bindings.scheduledDispatchDeadLetterQueueName})`, + `Ensuring production resources for worker: ${bindings.workerName} (D1: ${bindings.d1DatabaseName}, OAuth KV: ${bindings.oauthKvTitle}, Bundle KV: ${bindings.bundleArtifactsKvTitle}, Community R2: ${bindings.communityAssetsBucketName}, Email R2: ${bindings.emailBlobsBucketName}, Email Queue: ${bindings.emailDeliveryQueueName}, Email DLQ: ${bindings.emailDeliveryDeadLetterQueueName}, Artifacts Repo Events Queue: ${bindings.artifactsRepoEventsQueueName}, Artifacts Repo Events DLQ: ${bindings.artifactsRepoEventsDeadLetterQueueName}, Platform Feedback Queue: ${bindings.platformFeedbackDispatchQueueName}, Platform Feedback DLQ: ${bindings.platformFeedbackDispatchDeadLetterQueueName}, Community Activity Queue: ${bindings.communityActivityDispatchQueueName}, Community Activity DLQ: ${bindings.communityActivityDispatchDeadLetterQueueName}, Scheduled Dispatch Queue: ${bindings.scheduledDispatchQueueName}, Scheduled Dispatch DLQ: ${bindings.scheduledDispatchDeadLetterQueueName}, Package Events Queue: ${bindings.packageEventsDispatchQueueName}, Package Events DLQ: ${bindings.packageEventsDispatchDeadLetterQueueName})`, ) const d1 = ensureD1Database({ @@ -544,6 +546,18 @@ async function ensureProductionResources(options: CliOptions) { name: bindings.scheduledDispatchDeadLetterQueueName, dryRun: options.dryRun, }) + const packageEventsDispatchQueue = await ensureCloudflareQueue({ + accountId: accountId ?? 'dry-run-account', + apiToken: apiToken ?? 'dry-run-token', + name: bindings.packageEventsDispatchQueueName, + dryRun: options.dryRun, + }) + const packageEventsDispatchDeadLetterQueue = await ensureCloudflareQueue({ + accountId: accountId ?? 'dry-run-account', + apiToken: apiToken ?? 'dry-run-token', + name: bindings.packageEventsDispatchDeadLetterQueueName, + dryRun: options.dryRun, + }) const emailSendingDomain = resolveEmailSendingDomain(options.dryRun) const emailEventSubscription = await ensureEmailSendingEventSubscription({ accountId: accountId ?? 'dry-run-account', @@ -625,6 +639,12 @@ async function ensureProductionResources(options: CliOptions) { console.log( `scheduled_dispatch_dead_letter_queue_name=${scheduledDispatchDeadLetterQueue.name}`, ) + console.log( + `package_events_dispatch_queue_name=${packageEventsDispatchQueue.name}`, + ) + console.log( + `package_events_dispatch_dead_letter_queue_name=${packageEventsDispatchDeadLetterQueue.name}`, + ) console.log(`email_event_subscription_id=${emailEventSubscription.id}`) console.log( `artifacts_event_subscription_id=${artifactsEventSubscription.id}`,