Change streams
A change stream is a long-lived HTTP connection that receives one server-sent event per row a
committed write changed. data.svc serves it.
Two schema blocks declare one, and they are halves of the same thing. events: says what is
published and on which topic; subscriptions: says who may listen and what they see. A
subscription whose source no events: block publishes to is refused at boot, because a stream that
stays silent forever is worse than a refusal: it looks like it is working.
What a subscriber gets
Section titled “What a subscriber gets”Read this first. It is the part that surprises people.
A frame is a notification. It names the entity, the operation and the row, so a client knows what
to re-read. It carries no row data unless the subscription declares rls: true, and that setting
is what makes row data safe to send: the server re-reads the row through the engine, as the
subscriber, before putting it on the wire. Every access rule, row rule and mask that guards a
read guards the stream, because it is the same read. A row the subscriber may not see produces no
frame at all.
There is no shape in which rows leave a stream unguarded. An entity that has row rules cannot
declare a subscription without rls: true. That combination is refused at boot.
| Subscription | What a frame carries |
|---|---|
rls: false (the default) | Entity, operation, row id, time, tenant. No values |
rls: true | The same, plus the row as this subscriber may read it, shaped by the publisher’s include: / exclude: |
rls: true, row not visible to this subscriber | Nothing. No frame is written |
rls: true, operation is delete | The notification. A removed row has no values left to read |
The route
Section titled “The route”GET /subscribe/{entity}| Parameter | Default | Meaning |
|---|---|---|
name | the only subscription | Which declared subscription to open. Required when the entity declares more than one |
ops | what the subscription declares | Narrow to create, update, delete. Never widens |
heartbeat | 25s | Keep-alive interval, 1s to 5m |
lifetime | 30m | How long this connection stays open, up to 4h |
Connecting is a read: the caller must be able to read the entity, and the server proves it by performing one. A caller the REST surface refuses cannot open a stream on the same entity.
Every refusal names what is wrong:
| Status | Code | Meaning |
|---|---|---|
| 404 | entity_not_found | No such entity in this tenant’s schema |
| 404 | no_subscription_declared | The entity declares no subscriptions:, so it has no stream |
| 404 | subscription_not_found | No subscription by that name; the response lists the ones there are |
| 400 | subscription_ambiguous | The entity declares several and the request named none |
| 400 | bad_ops / bad_heartbeat / bad_lifetime | A parameter outside what it accepts. Out of range is an error, never a silent clamp |
| 403 | — | The caller cannot read the entity |
Response headers
Section titled “Response headers”| Header | Value |
|---|---|
Content-Type | text/event-stream |
Cache-Control | no-cache |
Connection | keep-alive |
X-Accel-Buffering | no |
X-Accel-Buffering: no is set because nginx and CloudFront buffer text/event-stream by default,
which would hold frames until a buffer fills.
Frames
Section titled “Frames”retry: 2000: subscribed entity=order subscription=mine closing-in=30m0s
id: 01M2P1R18V26CMVBZ1DC0AS1Z1event: entity.createdata: {"id":"01M2P1R18V26CMVBZ1DC0AS1Z1","type":"entity.create","entity":"order","operation":"create","subject":"01M2P1R18HT8575N8GVM2H5WPN","time":"2026-09-17T02:25:12Z","tenant":"acme:test:shop:t1","data":{"id":"01M2P1R18HT8575N8GVM2H5WPN","title":"first"}}
: ping
: closing, reconnectretry: tells an EventSource how soon to reconnect. id: and event: let a browser dispatch by
type and deduplicate. Payloads are JSON on one line, because the event-stream format splits on
newlines.
A connection ends on client disconnect, on a failed write, or when it reaches its lifetime. The
last case writes : closing, reconnect first, and an EventSource reconnects on its own. Every
exit path unsubscribes and releases the buffer.
| Field | Type | Content |
|---|---|---|
id | string | Event id, a ULID. One per changed row |
type | string | entity.create, entity.update or entity.delete |
entity | string | Entity name |
operation | string | create, update or delete |
subject | string | The row’s primary key |
time | RFC 3339 | Publication time, UTC |
tenant | string | Tenant of the write |
data | object | Present only on an rls: true stream. See What a subscriber gets |
One frame is written per row a write changed, not one per statement. An update touching forty rows produces forty frames, each naming its own row, because a notification that does not say which row changed says nothing a client can act on.
Declaring a stream
Section titled “Declaring a stream”entities: - name: order events: - name: changed operations: [create, update, delete] exclude: [internal_notes, card_last4] condition: language: expr expression: payload.status != "draft" subscriptions: - name: mine rls: true eventtypes: [create, update]events:, the publishing side
Section titled “events:, the publishing side”| Key | Type | Default | Meaning |
|---|---|---|---|
name | string | — | Label; also the publish-error metadata key |
operations | list | every write | create, update, delete. Anything else is refused at boot |
condition | script ref | — | Guard; a falsy result skips this event. An evaluation error aborts the mutation |
topic | string | entity.<name> | Publish topic. Two entries on one entity cannot share a topic |
format | string | json | Only json is served; anything else is refused at boot |
include | list | — | Allow-list of fields kept in data. Must name fields of the entity |
exclude | list | — | Deny-list. Ignored when include is set. Must name fields of the entity |
include: and exclude: are the field selection for the stream, and they hold whether a frame’s
data came from the write or from a re-read.
A publish failure is recorded in the request metadata and does not fail the mutation. A guard evaluation failure does.
subscriptions:, the listening side
Section titled “subscriptions:, the listening side”| Key | Values | Meaning |
|---|---|---|
name | string | How a client picks this stream. Unique per entity |
source | topic name | Which topic to read. Defaults to entity.<name>; must be a topic some events: block publishes to |
eventtypes | create, update, delete | Which operations reach a subscriber |
scope | tenant, user | user narrows to the writes this caller made. Every stream is tenant-scoped regardless |
rls | bool | Whether row data travels, by re-reading each row as the subscriber |
transports | sse | The only transport served. Anything else is refused at boot |
filter | script ref | Extra predicate per subscriber, truthy to deliver |
Two keys are refused rather than accepted and ignored:
fieldgroup:named field groups, a concept the schema language does not have. It selected nothing. Choose fields on the publishing side instead.- A
subscriptions:list at the top of a schema file rather than on an entity. Nothing ever read it.
Delivery contract
Section titled “Delivery contract”Each connection gets a channel of 64 events. Publication is non-blocking: the publisher walks the matching subscribers and enqueues without waiting for any of them to drain.
| Condition | Result |
|---|---|
| Buffer has room | Event enqueued; the subscriber’s delivered counter increments |
| Buffer full | Event dropped for that subscriber only; its dropped counter increments; others are unaffected |
| Subscriber already unsubscribed | Skipped before the send |
| Filter rejects the event | Skipped; counted as neither delivered nor dropped |
A slow consumer loses events. It never slows the write that produced them and never blocks another subscriber.
Two paths produce no events at all:
- Bulk create. The bulk insert path runs only the validation phases, never the post-execution phase where publication happens. Rows inserted in bulk are silent.
- Entities with no
events:block. Declaring a subscription is not enough; declare the event that feeds it. A subscription with no matching publisher is refused at boot, so this shows up as a boot failure rather than a quiet stream.
Scope of the fan-out
Section titled “Scope of the fan-out”The event router is per tenant and per datastore, the same unit as the engine beside it. A tenant’s writes cannot reach another tenant’s stream even before the per-subscriber tenant filter applies.
A stream does not cross process boundaries. Two instances behind a load balancer each publish their own writes to their own subscribers, so a client connected to instance A does not see a write that landed on instance B. Pin a stream to one instance, treat a reconnect as a cue to re-read, or use durable delivery where missing an event is not acceptable.
Durable delivery
Section titled “Durable delivery”A stream is a live view: a client that is not connected misses what happened while it was away, and
no history is kept. Where that is not good enough, declare a delivery target. The same events:
block then also goes to a queue that outlives the request, the process and the reconnect.
# boot documentevents: targets: billing: url: https://receiver.example/hooks/orders secret: a-shared-secret timeout: 10s topics: [entity.sales_order, entity.invoice]| Key | Default | Meaning |
|---|---|---|
events.targets.<name>.url | — | Where a delivery is posted. Must be http or https |
events.targets.<name>.secret | — | Signs the body. Required: an unsigned delivery is one a receiver cannot trust |
events.targets.<name>.timeout | 10s | Bounds one attempt. A slow receiver is retried, not waited on |
events.targets.<name>.topics | — | Which topics this target wants. Required |
A target that cannot work refuses boot rather than being skipped: no URL, a URL that is not http, no secret, no topics, a misspelt key. A declared target that silently receives nothing is the failure the queue exists to remove.
The schema does not change. A URL and a signing secret belong to whoever runs the deployment, so nothing about delivery appears in a tenant’s model.
What the guarantee is
Section titled “What the guarantee is”The queue row is written inside the write’s own transaction. A write that rolls back leaves nothing to deliver, and a write that commits cannot commit without its event recorded.
Delivery is at-least-once. A receiver can be sent the same event more than once, after a timeout it
did not report or a crash between accepting and recording. Every attempt carries the same
X-Kisai-Delivery, so a receiver that remembers the ids it has handled turns at-least-once into
exactly-once.
What a receiver gets
Section titled “What a receiver gets”POST /hooks/ordersContent-Type: application/jsonX-Kisai-Delivery: U01M2P1R18V26CMVBZ1DC0AS1Z1X-Kisai-Topic: entity.sales_orderX-Kisai-Event: entity.createX-Kisai-Timestamp: 1789245912X-Kisai-Signature: 6f1c…The body is the event, the same shape a stream frame carries. The signature is the HMAC-SHA256 of
"<timestamp>.<body>" under the target’s secret. Verify it before trusting anything else: a URL is
not a secret, and anything that can reach it can post to it. The timestamp is inside the signed
string, so a captured delivery cannot be replayed under a later one.
| Your answer | What happens |
|---|---|
| Any 2xx | Accepted. The delivery is done |
| 408 or 429 | Retried later with backoff |
| Any other 4xx | You have said the request itself is wrong. The delivery stops and waits for an operator |
| 5xx, a timeout, a refused connection | Retried later with backoff |
Backoff doubles from two seconds up to ten minutes, for up to twelve attempts, after which the delivery waits rather than being retried for ever.
Running the queue
Section titled “Running the queue”scheduler: outbox: every: 10sThe drainer inside each instance already ticks on its own, and is nudged the moment a write records an event. The scheduler slot makes the cadence a deployment decision and takes a per-tenant lease, so a fleet does not all wake at once to contend for the same rows. It is not scheduled at all where no target is declared.
| Route | Method | What it does |
|---|---|---|
/admin/data/events/pending | GET | Deliveries waiting or stuck, failed first |
/admin/data/events/stats | GET | Queue depth by status, this instance’s counters, the declared targets |
/admin/data/events/drain | POST | Attempt the due deliveries now |
/admin/data/events/requeue | POST | Put named failed deliveries back, with {"ids": ["U…"]} |
Requeue moves a failed delivery only. One that is pending would lose its backoff, and one that is claimed is in flight somewhere.
Counters
Section titled “Counters”| Scope | Counters |
|---|---|
| Per subscriber | delivered, dropped |
| Per manager | active subscribers, active topics, published, delivered, dropped |
Both are read through the Go API. No HTTP endpoint publishes them yet.
Related
Section titled “Related”| For | See |
|---|---|
The events:, triggers: and pointcuts: keys in context | Entities |
| Who may read the entity at all | Access rules |
| The rules that filter a re-read row | Row-level security |
| The boot keys for targets and the drain schedule | Configuration |
| Every served route, parameter and status code | Data API reference |