Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions docs/contributing/architecture/primitives.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
14 changes: 14 additions & 0 deletions docs/contributing/setup-manifest.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down
124 changes: 124 additions & 0 deletions docs/guides/package-subscriptions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>
}
```

- 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
Expand Down
128 changes: 128 additions & 0 deletions packages/shared/src/json-schema-subset.node.test.ts
Original file line number Diff line number Diff line change
@@ -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.'])
})
Loading
Loading