feat(package-events): queue-backed delivery, payload schemas, and subscription filters - #1229
Conversation
…scription filters Package-emitted events (kody.emits / events.dispatch) now enqueue on a new kody-package-events-dispatch Queue (with DLQ) instead of synchronously fanning out inside the emitting request. The consumer resolves same-user subscribers at delivery time, applies subscription filters (exact-match on top-level payload values), invokes handlers with the existing exactly-once idempotency keys, retries pre-execution infrastructure failures, and acks terminal handler failures. Environments without the queue binding (local dev, preview) deliver inline through the same consumer code path. kody.emits entries gain an optional payloadSchema (documented JSON Schema subset, validated at publish and dispatch time), payloads are capped at 64 KiB, and event chains carry the nested invocation depth across the queue boundary so emit cycles terminate. Production CI provisions the new queue pair like the existing dispatch queues. Co-authored-by: Kent C. Dodds <me+github@kentcdodds.com>
📝 WalkthroughWalkthroughPackage-emitted events now support validated JSON Schema payloads, filters, durable queue delivery, idempotent subscriber handling, retries, dead-letter queues, and inline fallback delivery. Worker routing and production resource provisioning support the new package-events queue. ChangesPackage event delivery
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant PackageRuntime
participant SubscriptionDispatch
participant PackageEventsQueue
participant QueueHandler
participant SubscriberInvocation
PackageRuntime->>SubscriptionDispatch: events.dispatch(topic, payload, idempotencyKey)
SubscriptionDispatch->>PackageEventsQueue: enqueue validated event message
PackageEventsQueue->>QueueHandler: deliver queue batch
QueueHandler->>SubscriberInvocation: deliverPackageEvent(message)
SubscriberInvocation-->>QueueHandler: delivery result or retryable error
Possibly related PRs
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches 💡 1📝 Generate docstrings 💡
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
🔎 Preview deployed: https://kody-pr-1229.kody-a99.workers.dev Worker: Mocks:
|
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
packages/worker/src/package-invocations/subscription-dispatch.ts (1)
376-408: 🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick winThe result reports
status: 'enqueued'after inline delivery.The inline branch runs when no queue binding exists or when
queue.sendthrows. The return value still reports'enqueued'. Callers and package authors then cannot tell a durable enqueue from a best-effort inline fan-out, and an enqueue failure is invisible to them. Report the actual outcome instead.♻️ Proposed change
if (!enqueued) { @@ if (input.waitUntil) { input.waitUntil(inlineDelivery) } else { await inlineDelivery } } return { topic: request.topic, source: { type: 'package', packageId: packageContext.packageId, kodyId: packageContext.kodyId, }, idempotencyKey: request.idempotencyKey, - status: 'enqueued', + status: enqueued ? 'enqueued' : 'delivered_inline', }Update the
PackageEventTools['dispatch']return type to include the new status value.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-invocations/subscription-dispatch.ts` around lines 376 - 408, Update the dispatch flow around deliverPackageEventWithToolFactories so inline delivery returns the new non-enqueued status instead of 'enqueued', while preserving 'enqueued' for successful queue sends. Extend PackageEventTools['dispatch'] and any related result types to include the new status value, and ensure queue-send failures report the inline outcome rather than masking the failure.
🧹 Nitpick comments (6)
packages/shared/src/json-schema-subset.ts (1)
277-283: 🎯 Functional Correctness | 🔵 Trivial | 💤 Low value
additionalProperties: falseis ignored whenpropertiesis absent.The check requires
propertiesto be present. A schema like{ type: 'object', additionalProperties: false }therefore accepts any keys. Standard JSON Schema rejects every property in that case. Treat a missingpropertiesas an empty property set.♻️ Proposed change
- const properties = isPlainObject(schema['properties']) - ? schema['properties'] - : null - if (properties) { + 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 && properties) { + if (schema['additionalProperties'] === false) { for (const key of Object.keys(value)) { if (!(key in properties)) { errors.push(`${path} has unexpected property "${key}".`) } } }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/shared/src/json-schema-subset.ts` around lines 277 - 283, Update the additionalProperties validation in the object-schema checking logic so additionalProperties: false treats a missing properties definition as an empty allowed-property set. Ensure every key in the value is reported as unexpected when properties is absent, while preserving the existing behavior when properties is provided.packages/worker/src/package-events/dispatch-queue-producer.ts (1)
69-74: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueDecide whether
enqueuePackageEventDispatchshould be the producer API.
enqueuePackageEventDispatchis only exported fromdispatch-queue-producer.ts; the dispatch path constructs aPackageEventsDispatchQueueMessageand callsPACKAGE_EVENTS_DISPATCH_QUEUE.send(message)directly. Call this helper from the dispatch path or remove it if the direct send path is intentional.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-events/dispatch-queue-producer.ts` around lines 69 - 74, Resolve the unused producer API around enqueuePackageEventDispatch: either update the dispatch path to call this helper when sending the PackageEventsDispatchQueueMessage, or remove the helper and its export if direct PACKAGE_EVENTS_DISPATCH_QUEUE.send usage is intentional. Keep a single consistent producer path.packages/worker/src/queue-handler.node.test.ts (1)
65-65: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueUse the exported queue-name constant for consistency.
Line 66 uses
scheduledDispatchQueueName, but Line 65 hardcodes'kody-package-events-dispatch'. ImportpackageEventsDispatchQueueNamefrom#worker/package-events/dispatch-queue-names.tsto match the surrounding cases.♻️ Proposed change
- const packageEventsBatch = createBatch('kody-package-events-dispatch') + const packageEventsBatch = createBatch(packageEventsDispatchQueueName)🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/queue-handler.node.test.ts` at line 65, Update the packageEventsBatch setup to use the exported packageEventsDispatchQueueName constant instead of the hardcoded queue-name string, importing it from `#worker/package-events/dispatch-queue-names.ts` and preserving the existing createBatch flow.packages/worker/src/package-invocations/service.node.test.ts (1)
1327-1357: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winThe manifest mutation leaves the
sourceFilesfixture stale.
seedRuntimeDispatchPackagesserializes each manifest intosourceFilesat the time it builds the map ('package.json': JSON.stringify(manifests.get('source-gateway')), Line 824). Lines 1338-1351 mutategatewayManifest.kody.emitsafter that snapshot exists.loadPackageManifestBySourceIdreads the live object and seespayloadSchema, butloadPackageSourceBySourceIdreturns the stalepackage.jsontext without it.The test passes today only because dispatch validation reads the manifest through
loadPackageManifestBySourceId. If dispatch later resolves the manifest from source files, this test stops validatingpayloadSchemaand still passes.Pass the emitted-event declaration into the seeding helper, or re-serialize
sourceFilesafter mutating the manifest, so both loaders agree.♻️ Proposed fix: re-serialize the source file after mutation
test('package runtime dispatch validates payloads against the declared payloadSchema', async () => { const db = createDatabase() - const { manifests } = seedRuntimeDispatchPackages() + const { manifests, sourceFiles } = seedRuntimeDispatchPackages() const gatewayManifest = manifests.get('source-gateway') as { kody: { emits?: Record< string, { description: string; payloadSchema?: Record<string, unknown> } > } } 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, }, }, } + const gatewayFiles = sourceFiles.get('source-gateway') + if (gatewayFiles) { + gatewayFiles['package.json'] = JSON.stringify(gatewayManifest) + }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-invocations/service.node.test.ts` around lines 1327 - 1357, After mutating gatewayManifest.kody.emits in the payload-schema dispatch test, re-serialize the updated source-gateway manifest into its corresponding sourceFiles package.json entry before creating the runtime event tools. Ensure loadPackageManifestBySourceId and loadPackageSourceBySourceId both observe the same emitted-event declaration.packages/worker/wrangler.jsonc (1)
320-322: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick winConsider a
max_concurrencycap on this consumer.Each package event can invoke several subscriber handlers, and every handler load runs a dynamic Worker isolate. Without
max_concurrency, Cloudflare scales this consumer freely, so a burst of emitted events can multiply into a large number of concurrent isolate loads and subrequests.The
kody-scheduled-dispatchconsumer at Line 315 already caps concurrency at 16 for comparable per-message work. The sibling dispatch queues omit the cap, so this is a deliberate choice to confirm rather than a defect.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/wrangler.jsonc` around lines 320 - 322, Confirm the intended concurrency policy for the consumer configuration containing max_batch_size, max_batch_timeout, and max_retries; if it should match kody-scheduled-dispatch, add a max_concurrency setting capped at 16, otherwise leave the current configuration unchanged.packages/worker/src/package-events/dispatch-queue.node.test.ts (1)
123-157: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winAdd a rejection case for a blank
userId.The parser rejects a non-string or blank
userId.userIdselects which user's subscribed packages receive the event, so it is the field that enforces per-user isolation across the queue boundary. The current cases covertopic,idempotencyKey,payload,source, andinvokeDepth, but notuserId.Also consider asserting that the parser trims values, because the equality check at Line 124 uses an already-trimmed fixture.
As per coding guidelines: "Every signed-in user must have a fully isolated personal assistant".
💚 Proposed additional cases
expect( parsePackageEventsDispatchQueueMessage(createMessageBody({ topic: ' ' })), ).toBeNull() + 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(🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-events/dispatch-queue.node.test.ts` around lines 123 - 157, Add a malformed-body assertion in the test "package events queue message parsing rejects malformed bodies" that passes a blank userId to parsePackageEventsDispatchQueueMessage and expects null; also add coverage for surrounding whitespace to verify the parser trims accepted userId values, while preserving the existing trimmed-fixture equality assertion.Source: Coding guidelines
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Outside diff comments:
In `@packages/worker/src/package-invocations/subscription-dispatch.ts`:
- Around line 376-408: Update the dispatch flow around
deliverPackageEventWithToolFactories so inline delivery returns the new
non-enqueued status instead of 'enqueued', while preserving 'enqueued' for
successful queue sends. Extend PackageEventTools['dispatch'] and any related
result types to include the new status value, and ensure queue-send failures
report the inline outcome rather than masking the failure.
---
Nitpick comments:
In `@packages/shared/src/json-schema-subset.ts`:
- Around line 277-283: Update the additionalProperties validation in the
object-schema checking logic so additionalProperties: false treats a missing
properties definition as an empty allowed-property set. Ensure every key in the
value is reported as unexpected when properties is absent, while preserving the
existing behavior when properties is provided.
In `@packages/worker/src/package-events/dispatch-queue-producer.ts`:
- Around line 69-74: Resolve the unused producer API around
enqueuePackageEventDispatch: either update the dispatch path to call this helper
when sending the PackageEventsDispatchQueueMessage, or remove the helper and its
export if direct PACKAGE_EVENTS_DISPATCH_QUEUE.send usage is intentional. Keep a
single consistent producer path.
In `@packages/worker/src/package-events/dispatch-queue.node.test.ts`:
- Around line 123-157: Add a malformed-body assertion in the test "package
events queue message parsing rejects malformed bodies" that passes a blank
userId to parsePackageEventsDispatchQueueMessage and expects null; also add
coverage for surrounding whitespace to verify the parser trims accepted userId
values, while preserving the existing trimmed-fixture equality assertion.
In `@packages/worker/src/package-invocations/service.node.test.ts`:
- Around line 1327-1357: After mutating gatewayManifest.kody.emits in the
payload-schema dispatch test, re-serialize the updated source-gateway manifest
into its corresponding sourceFiles package.json entry before creating the
runtime event tools. Ensure loadPackageManifestBySourceId and
loadPackageSourceBySourceId both observe the same emitted-event declaration.
In `@packages/worker/src/queue-handler.node.test.ts`:
- Line 65: Update the packageEventsBatch setup to use the exported
packageEventsDispatchQueueName constant instead of the hardcoded queue-name
string, importing it from `#worker/package-events/dispatch-queue-names.ts` and
preserving the existing createBatch flow.
In `@packages/worker/wrangler.jsonc`:
- Around line 320-322: Confirm the intended concurrency policy for the consumer
configuration containing max_batch_size, max_batch_timeout, and max_retries; if
it should match kody-scheduled-dispatch, add a max_concurrency setting capped at
16, otherwise leave the current configuration unchanged.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: 8f49479d-dd0a-46a0-91f4-09460c703bd5
📒 Files selected for processing (24)
docs/contributing/architecture/primitives.yamldocs/contributing/setup-manifest.mddocs/guides/package-subscriptions.mdpackages/shared/src/json-schema-subset.node.test.tspackages/shared/src/json-schema-subset.tspackages/worker/src/package-events/dispatch-queue-names.tspackages/worker/src/package-events/dispatch-queue-producer.tspackages/worker/src/package-events/dispatch-queue.node.test.tspackages/worker/src/package-events/dispatch-queue.tspackages/worker/src/package-invocations/admin-package-subscriptions.tspackages/worker/src/package-invocations/infrastructure-codes.tspackages/worker/src/package-invocations/service.node.test.tspackages/worker/src/package-invocations/service.tspackages/worker/src/package-invocations/subscription-dispatch.tspackages/worker/src/package-registry/manifest.node.test.tspackages/worker/src/package-registry/manifest.tspackages/worker/src/package-registry/types.tspackages/worker/src/queue-handler.node.test.tspackages/worker/src/queue-handler.tspackages/worker/worker-configuration.d.tspackages/worker/wrangler.jsonctools/ci/production-queue-resources.node.test.tstools/ci/production-queue-resources.tstools/ci/production-resources.ts
Report delivered_inline when dispatch falls back past the queue, apply additionalProperties: false with an absent properties set, drop the unused enqueue helper, cap the queue consumer at 16 concurrent invocations, keep mutated manifest fixtures in sync with their serialized source snapshots, and cover userId parsing/trimming in the queue message parser. Co-authored-by: Kent C. Dodds <me+github@kentcdodds.com>
What
Package-emitted events (
kody.emits+events.dispatch) move from synchronous in-request fan-out to durable, queue-backed delivery, and gain payload contracts and working subscription filters.events.dispatchvalidates the event and enqueues it on a newkody-package-events-dispatchQueue (with-dlqdead-letter queue), returning{ topic, source, idempotencyKey, status: "enqueued" }immediately. The old synchronous fan-out (and its per-subscriber return payload) is removed. The Queue consumer resolves the emitting user's subscribers at delivery time and invokes handlers with the existing exactly-once idempotency keys, so Queue redelivery replays instead of re-running. Pre-execution infrastructure failures retry (3 attempts, then DLQ); terminal handler failures ack and stay visible in run records. Consumer concurrency is capped at 16.SCHEDULED_DISPATCH_QUEUEdegrades. The fallback reportsstatus: "delivered_inline"so callers can tell durable enqueue from best-effort inline delivery.payloadSchemaonkody.emitsentries. Optional, documented JSON Schema subset (packages/shared/src/json-schema-subset.ts); unsupported keywords fail package checks at publish time, payloads are validated at dispatch time, and schemas project into search/detail so subscribers can discover payload shapes. Payloads are capped at 64 KiB canonical JSON (Queue message headroom).filtersnow apply to package-emitted topics: every filter key must be present in the payload with a canonically-equal JSON value, otherwise the subscriber is skipped. Platform topics keep their existing dispatcher-defined behavior.packages.invokechains and cannot loop forever.tools/ci/production-resources.tsensures the new queue pair on deploy, exactly like the existing dispatch queues;production-queue-resources.tsenforces the wrangler config shape.Deliberately not included
A publish-time topic→subscriber D1 index was considered and deferred: manifest loads are commit-keyed and cached per isolate, every existing platform dispatcher uses the same delivery-time scan, and queued delivery already moves that scan off the emitter's request path. Adding a D1 migration right after the migration squash was not worth the marginal win.
Testing
userIdisolation and trimming), consumer ack/retry semantics, queue routing, manifestpayloadSchemavalidation and projection, dispatch enqueue contract (declared-topic gate, schema gate, size cap, depth cap), delivery fan-out with replay, filters gating, inline fallback status, and retryable-infrastructure classification.npm run validaterun locally (format, lint, typecheck, migrations/guardrail checks, MCP tests, e2e, 1837 unit tests).additionalPropertiessemantics, consumer concurrency cap, fixture staleness, parser coverage).System recap — adds a new primitive (high risk)
Mode: recap · Base:
main@bde1cc1e· Head:7a9d1898Classification: adds — introduces the
package-events-dispatch-queueprimitive (Cloudflare Queue + consumer for package-emitted events) and extends the saved-packages manifest contract withkody.emits[].payloadSchema.Primitives touched
package-events-dispatch-queueevents.dispatch, consumer fan-out insubscription-dispatch.tssaved-packageskody.emits[].payloadSchema(JSON Schema subset), filters applied for package-emitted topicsplatform-feedback-dispatch-queueemail-delivery-queuewebhookspackage-registryroot only, no webhook behavior changeSystem map
A package export calls
events.dispatch, which enqueues on the new package-events queue; the consumer fans out to same-user subscription handlers with exactly-once idempotency.Legend: green = composes (wiring only) · amber = extended by this PR · red = new primitive · gray = context (unchanged, included only when an edge crosses it).
Before / after
Invariants
per-user-isolation: delivery remains strictly same-user; the queue message carries the emittinguserIdand subscriber resolution is scoped to it.no-per-event-shared-writes: no new per-event D1 writes; durability comes from the Queue, and per-delivery history stays in the per-user RunLog DO ledger.Summary by CodeRabbit