> ## Documentation Index
> Fetch the complete documentation index at: https://docs.akter.dev/llms.txt
> Use this file to discover all available pages before exploring further.

# Server API

> Actor declarations and server composition in @rikalabs/akter.

# Server API

Due-work data placement is internal: relay claims and connection-holder liveness use the runtime's data shard map, whose ordinary database default is one range spanning all buckets. This does not add an `Actors.layer` or `Database.postgres` option, change retry identity, or change receiver deduplication. Optional internal Neki shard targeting has no supported topology configuration yet, and does not establish Neki support ([ADR 0067](../decisions/0067-due-work-shard-ranges.md)).

**Responsibility:** define actor declarations and server composition.\
**Authority:** API design.\
**Owner role:** API/Effect.
**Change policy:** a change requires compatibility review against docs/api/versioning.md.

The Effect-native server API is one package, `@rikalabs/akter`. Its root entry exports `Actor`, `Actors`, `Intent`, errors, identity (`ActorRef`, `Caller`, `User`, `System`, `Anonymous`, `Principal`, `CurrentCaller`, `Tenant`), and the `Policy`, `Handle`, `Intents`, and context types. The root entry holds browser-safe declarations only. Runtime construction and serving are imported separately from `@rikalabs/akter/runtime`: `Actors.layer`, `Actors.serve` for HTTP commands and queries (M3.2), and the `Auth` providers; `Fleet` is a target API and is not exported yet. The same entry exports `Inspector.serve({ auth, basePath?, runner?, region? })`, read-only routes over the single-table `durable.*_v2` inspection views ([ADR 0095](../decisions/0095-single-table-inspection-views.md)) for the authenticated principal's tenant, with each actor's placement joined from `durable.placements_v2` by the server and so unchanged in the responses and, inside a runtime, over that runner's own memory (live rates, latency, resident activations, connections and a redacted command stream, which name `runner` and `region` as given), which `akter dev` serves ([inspection views](../operations/inspection-views.md#the-local-inspector)), `Telemetry.serve({ basePath? })`, `GET /metrics` in Prometheus text ([observability](../operations/03-observability.md)), and `Operators.serve({ auth, basePath? })`, the operator routes for inspection, actor export, receipt reads, defects, dead-letter repair, and the audit log, authenticated only by an `OperatorAuth.make` or `OperatorAuth.tokens` provider whose grant lists action- and resource-scoped capabilities ([ADR 0050](../decisions/0050-operator-authority-and-audited-repair.md)). `GET /operator/actors/:type/:id/export?tenant=` needs the `export` action, which no other action implies, and answers one actor's seed: `format`, `actor`, `created`, the stored `state` and its `stateVersion`, the pending `intents` and `jobs` with `dueInMs` measured from the export, and `omitted` counts of the receipts, events, workflows, dead letters, owned-table rows, and blob entries it does not carry. It runs bound to the tenant, as the runtime's tenant role when row-level security is on. It carries no tenant, caller, or credential. It reads in one read-only snapshot, is audited before it answers, answers 404 for an actor the tenant does not have, and answers 409 `ExportRefused` for state that does not decode or more than 10,000 pending intents or jobs. The runtime entry exports the `Seed` schema and `SeedJson` codec. `akter export Room/r1 --url <runner> --tenant <t> --output r1.seed` writes the file for its owner alone and never replaces one.

## Implemented foundation subset

Actor, query, and job layers may register concurrently beyond the PostgreSQL pool's checkout bound. Startup waits and retries pre-statement checkout refusals with backoff; placement, schema, workflow, and payload incompatibilities still refuse the deployment. This does not increase the pool bound or retry external command overload refusals in place.

Served commands with a supplied id resolve retained receipts in the fenced turn, not a pre-delivery read. External authorization still runs before delivery and before returning the outcome; expiry is checked at fenced admission against a clock read after acquiring the generation lock and rechecked using a fresh database clock read after the transaction ends. An ordinary warm command without handler SQL uses two database flights and seven statements on Postgres; a warm replay uses two flights and five statements without running the handler ([ADR 0072](../decisions/0072-served-command-in-two-round-trips.md)). PGlite performs the same checks as sequential engine calls.

The runtime and served protocol support Bun 1.4.2+ and Node 24+. Provide the matching Effect platform layers: `BunCrypto` and `BunHttpServer` from `@effect/platform-bun`, or `NodeCrypto` and `NodeHttpServer` from `@effect/platform-node`. Actor definitions, `Actors.layer`, `Actors.serve`, `Database`, and `ActorTest` do not change between runtimes. Node HTTP serving uses `NodeHttpServer.layer(createServer, options)`, with `createServer` from `node:http`; file-backed PGlite keeps its Linux/macOS local-filesystem limits.

`Database.postgres` holds a turn pool (`maxConnections`, default 50), an off-turn pool (`offTurnConnections`, default 10), and a query pool (`queryConnections`, default 10), each handing out connections first come, first served. Queries of types with owned tables or blobs stay on the off-turn pool so their bound client shares the transaction. Other queries use a caught-up replica when configured, otherwise the query pool. It defaults server TCP keepalive settings on its turn, off-turn, query, and optional replica and coordination pools to 5 seconds idle, 2 seconds between probes, and 3 probes. Explicit `startupParameters`, `startupOptions`, and the connection URL's `options` retain their normal driver precedence over the defaults; replica and coordination settings are independent. Set the three keepalive values to `0` to use the operating system's defaults. This helps detect a vanished TCP peer, not guarantee total recovery latency or interrupt every running SQL statement; pooler and provider compatibility remain unverified ([deployment](../operations/01-deployment.md#postgres-connections-across-runners)).

`Database.postgres({ coordination: { url: Redacted.make(coordinationUrl), ... } })` adds an independent unsharded primary pool for retention/workflow resource rows, Cluster runner registration and shard locks, singleton table-lease checks, and the fleet maintainer lock. Every runner must configure the same authority; omission uses its ordinary off-turn pool for a single-database deployment. The optional pool accepts the ordinary PostgreSQL pool settings and the same keepalive defaults; include its connections in the deployment budget. Its coordination-only migration runs automatically and needs DDL privileges; first-boot coordination and Cluster-table preparation are serialized on that pool. With `neki: true`, its fixed bootstrap DDL uses autocommit and propagation barriers too. Actor data, workflow manifests, and the fleet logical slot stay on the data connection. Neki routing and multi-shard failover remain unverified ([ADR 0066](../decisions/0066-authoritative-coordination.md)).

The executable runtime runs embedded or served on Postgres or PGlite, with one runner or, on Postgres, several runners sharing one database (M2). It implements the [ADR 0010](../decisions/0010-one-way-effect-native-api.md) definition and handler shape and [ADR 0011](../decisions/0011-direct-commands-outbox-and-performance.md) direct commands for commands, server reducers, queries, and keyed state ([ADR 0013](../decisions/0013-m0-reconciliation.md)). Owned tables ([Drizzle integration](04-drizzle.md)), events, intents, jobs, workflows with their deploy compatibility check (M2.7, M2.8), commutative merging, streams, and connections are implemented. `Fleet` is a target API ([ADR 0056](../decisions/0056-fleet-views.md)) and is not implemented.

* `Actor.command(tag, { payload?, success?, error? })` and `Actor.query(tag, { payload?, success?, error?, watch? })` accept service-free schemas; `payload` is a schema or struct fields, which stand for `Schema.Struct(fields)`. Omitted payload/success is `Schema.Void` and omitted `error` is `Schema.Never`; JSON codecs preserve it through persistence. Declared errors must be yieldable tagged errors.
* `Actor.table(pgTable(...))` declares an owned Drizzle table: it adds `routing_key`, `tenant_id`, and `actor_id` and prefixes the primary key, unique constraints, and indexes with them. `Actor.make(name, { tables })` lists the actor type's tables; a table belongs to one actor type. `turn.rows(table)` offers scoped `one`, `all`, `count`, `insert`, `update(...).where`, `delete().where`, and `upsert` in the turn transaction; `read.rows(table)` offers only the reads; `turn.group` and `read.group` run read-only selects and inner/left joins across the placement group. Migration `0005_tables` records each table's owner in `actor_tables`, and startup fails when a table is missing, keyed without ownership, or claimed by another actor type. `ActorTest.inspect` reports `rows` per owned table. See [Drizzle integration](04-drizzle.md) for the operation matrix.
* `Actor.blob(name)` declares binary storage; `name` is 1–80 letters, digits, `-`, or `_`, starting with a letter. `Actor.make(name, { blobs })` lists the actor type's blobs; a name appears once per actor type, and two actor types may declare the same name without sharing entries. `turn.blob(B)` offers `get(entry)` (an `Option<Uint8Array>` of every chunk in order), `set(entry, bytes)`, `append(entry, bytes)`, `compact(entry)`, and `delete(entry)` (removes the entry, so `get` returns none; a missing entry is no error) on the turn transaction; `read.blob(B)` offers only `get` and reads committed entries. Entry names are well-formed strings of 1–512 UTF-8 bytes without NUL. An entry holds at most 8 MiB (`MAX_ENTRY_BYTES`, 8,388,608 bytes) across all its chunks, because `get` returns it as one row and the Postgres driver closes a connection on any message over 16 MiB; a `set` or `append` past the cap is a deterministic defect. `policy.maxBlobBytes` (default 67,108,864 bytes, 64 MiB) caps the bytes of all of one actor's entries together; a `set` counts only the entry's new bytes, and a write past the cap is a deterministic defect that leaves the entry as it was, even if the handler catches the defect. `policy.maxBlobEntries` (default 10,000) caps the entries of all of one actor's blobs together: a `set` or `append` that would create an entry past it is a deterministic defect that creates nothing, even if the handler catches the defect, while writes to existing entries are never refused by the count; `delete` frees an entry's bytes and its slot ([ADR 0044](../decisions/0044-blob-entry-quota-and-deletion.md), proposed). Bytes are copied when a write is called, never count toward `maxStateBytes`, and roll back with a declared failure or a failed commit. An undeclared blob, a malformed entry name or byte value, a write from a query, or a capability used after its turn or query ends or from a fiber forked inside it is a deterministic defect. Migration `0009_blobs` adds `actor_blobs`. `ActorTest.inspect` reports `blobs`, the entry count per declared blob. There is no entry listing yet, and no external object storage.
* `Actor.make(name, { key?, placement?, state?, events?, feeds?, tables?, blobs?, api?, internal?, createdBy?, schedules?, jobs?, subscriptions?, policy?, access? })`. The whole definition is validated before anything is published, so a rejected definition leaves no table ownership or other registration behind. `placement` is `"tenant"` (default), `"actor"`, `"authority"`, or `{ parent: P }` ([ADR 0033](../decisions/0033-parent-actor-placement.md)): rows live on the shard of the actor-placed root of `P`'s chain. `"authority"` colocates a tenant's authority-placed actors like `"tenant"`, with every key in bucket -128, which a Neki topology keeps on the authoritative shard beside the tables it does not route, for an actor whose turns read or write those tables ([ADR 0097](../decisions/0097-authority-placed-control-plane-actors.md)). `P` must be placed by `"actor"` or `{ parent }` (a tenant- or authority-placed parent is a type error and throws), chains are at most four levels below the root, and a parent-placed actor needs a `key` or `createdBy` and cannot be a singleton. Its id is `c1.<byte length of parent id>.<parent id>.<local id>`: `X.idOf(parentId, local)` builds a named child's id, whose `key` validates the local part, and `X.get(id)` takes the full id; a minted one has no `X.create()` (`never`, and it dies at runtime) and is minted only by `turn.mint(X)` in a turn of its parent type, which returns `c1.<len>.<parent id>.<uuidv8>`. Served routes take the full id, and a malformed id or parent part answers `400 InvalidInput { code: "decode" }`. `api` and `internal` are records whose keys must equal their command tags, checked in types and at runtime. `key` is an id schema (named, `X.get(id)`), `Actor.singleton` (`X.get()`, registered through `Sharding.registerSingleton` on the single embedded runner), or omitted (minted: `X.create()` mints a UUIDv7 and is `never` otherwise; `turn.mint(X)` inside a parent's command turn mints a deterministic UUIDv8 for an actor that declares `createdBy` (every unkeyed actor accepts UUIDv8 ids, so its minted children stay reachable if the policy is later removed), [ADR 0025](../decisions/0025-turn-mint.md); see [context](02-context.md)). `internal` commands never appear on public handles or `X.api`; `ActorTest.actor` reaches them through a package-internal registry.
* `policy` accepts `executionTimeout`, `lockWait`, `deliveryTimeout`, `maxStateBytes`, `hibernateAfter`, `mailboxCapacity`, `keepReceipts`, `keepEvents`, `maxBlobBytes`, and `maxBlobEntries`, with defaults 30 s, 2 s, 30 s, 65,536 bytes of the complete encoded state object, 60 s, none declared (the runtime still refuses past 1,024 queued commands; see `admission` below), 7 days, 30 days, 64 MiB, and 10,000 entries, plus `keepWorkflows` (7 days), `maxScheduleLag` (1 day), `holdEventsForSubscribers` (7 days), `reauthorizeEvery` (60 s, from 1 s to 1 h), `connections` (`"park"`), `watch`, and `allowedSubscriberTypes`, described below. `createdBy`, `schedules`, and `jobs` sit on the definition itself, because they declare what the actor does rather than limit it. Values are positive integers to 2^31 − 1, except the five horizons (`keepReceipts`, `keepEvents`, `keepWorkflows`, `maxScheduleLag`, and `holdEventsForSubscribers`), which may reach about ten years; a `createdBy` command from another actor is rejected in types and at `Actor.make`. `executionTimeout` also bounds each query: past it the query fails `ActorError` `Timeout` and its running read is cancelled on the server. Cancellation needs an unmultiplexed Postgres pool (the default); with `multiplex: true`, or on PGlite, the query still fails at the deadline but its statement runs on.
* Retention ([retention](../operations/retention.md)). Every runtime sweeps once a minute. A receipt is deleted once `keepReceipts` has passed since its command id was issued (a timer's receipt counts from its due time), never sooner than `deliveryTimeout` after the id expires, and never while an outbox row with its id is pending, so a redelivered intent or job route still finds it. Events are deleted once `keepEvents` has passed since they were emitted, oldest first per actor, and `event_sequence` is never reset, so replay after a pruned cursor fails `RetentionGap`. Dead letters are not pruned. Sweeps of one actor type take turns on an authoritative transaction-owned row lock, with a local data fence when the coordination pool is separate ([ADR 0066](../decisions/0066-authoritative-coordination.md)). `ActorTest` runtimes never sweep on their own; `ActorTest.cleanup` runs one sweep now, and `cleanup` from `@rikalabs/akter/testing` does the same in a production-layer runtime. Migration `0010_retention` adds the indexes cleanup reads.
* `X.toLayer(handlers)` takes one handler per `api` and `internal` command, or an Effect that builds them from services; each handler's requirements are inferred, and a builder's typed failure types the layer (on a singleton it fails the activation instead). Handlers take only their payload and read the turn with `yield* X.Turn`, which supplies `id`, `ref`, `caller`, `principal` (an `Option`), `commandId`, and schema-decoded state with `state.set(patch)`. `X.Turn` is a distinct service per actor, so using another actor's turn leaves it unprovided in the layer's type. State changes become visible only after commit. Captured request/reply handles cannot run in turns; escaped state setters die.
* Storage follows [ADR 0006](../decisions/0006-scale-rules-placement-and-query-tiers.md) and [ADR 0011](../decisions/0011-direct-commands-outbox-and-performance.md): every per-actor and per-tenant framework row leads its primary key with `routing_key` (XXH3-64 of a versioned placement encoding), and state values are `bytea` in codec version 1, one plain zstd frame with no dictionary ([contract 06](../contracts/06-storage-ownership.md)). A warm activation reuses the committed state it last read or wrote and reads `actor_state` only after acquiring a new generation. Migration `0003_routing_state` rebuilds the foundation tables with composite foreign keys to `actor_generations` and refuses to run on a database that already holds actors. The first registration of an actor type records its placement and encoding in `actor_placements`; a later deployment with a different placement fails at startup instead of forking actors under a second routing key, except that a tenant-placed type declared `"authority"` has its rows moved to the authority keys in one transaction at that startup, which a database that routes tables refuses.
* Commands are direct: Cluster carries a volatile message and keeps no message storage, and the receipt committed in the turn is the only durable admission record. The command deadline covers the whole turn transaction plus transaction-local `statement_timeout`/`lock_timeout`; a timeout, retryable SQL failure, or stale generation restarts the activation and fails the attempt `ActorUnavailable` after its backoff, and the handle retries with the same command id until `deliveryTimeout`. `deliveryTimeout` only stops the caller waiting. `createdBy` fails non-creating commands `NotCreated` without a receipt until the creating command commits; adding the policy to actors with pre-existing rows requires an application migration.
* `state` is `Actor.state(fields, { migrations? })`; missing keys decode from field defaults, and `set` and `$version` are reserved keys. `migrations: [Actor.migration(From, To, upcast), ...]` upcasts older stored state. `Actor.make` rejects a chain unless each `to` is the next `from` and the last `to` is `fields`. Stored state carries a `$version` row (the number of migrations applied). A turn or query decodes older rows through the remaining steps; a successful turn writes every key at the current version and deletes keys the current shape no longer has, and a declared failure discards those writes. A stored version newer than the chain, or an upcast that throws or produces an invalid shape, is a deterministic defect that leaves the stored rows unchanged. `ActorTest.seed(ref, state, version)` writes old-shape rows under the actor type's recorded placement.
* `Actor.reducer(tag, { state, payload?, error?, reduce, batch? })` declares a pure transition. `state` must be the actor's own `Actor.state` value, checked structurally in types (the reducer's fields must equal the actor's) and by identity at `Actor.make`, so a separately built but equal `Actor.state` compiles and then throws; a reducer is an `api` member only. `reduce(state, payload)` returns a `Result` of the complete next state or a declared error. The handle method has a command's shape and framework reasons and replies with the committed state. `X.toLayer` has no entry for a reducer (declaring one fails to compile); the actor still registers through `X.toLayer`, with `X.toLayer({})` when it has no commands. On the server a reducer call is an ordinary fenced, receipted turn: the runtime decodes state (upcasting if needed), applies `reduce`, validates the result against the state schema, and writes only changed keys (`reduce` receives its own copy, so mutating it in place is still committed), or every key after an upcast. A failure is a declared failure with the usual rollback and terminal receipt; a throwing `reduce` or an invalid returned state is a deterministic defect. With `batch: { combine }` the reducer replies `void`, declares no errors, and its `reduce` must not fail, all checked in types (declaring an `error` also throws). Calls of one batched reducer already waiting, consecutively, in the actor's mailbox merge: their payloads are combined in order with `combine` and reduced once, in one turn of up to 1,024 calls with one receipt per original command id; a lone call is never delayed, and callers on any runner merge on the owner. Merging is only correct when `reduce(reduce(s, a), b)` equals `reduce(s, combine(a, b))`; `/testing` exports `checkBatchLaw({ reducer, state?, input?, maxInputs?, runs? })`, a seeded property check to run for every batched reducer. The Promise client runs reducers optimistically over HTTP ([implemented subset](03-typescript-sdk.md#implemented-subset-m34)).
* `X.toQueryLayer(build)` implements every query in `api`; handlers read `yield* X.Read` (`id`, `ref`, `caller`, `principal`, and committed `state`). A query applies the caller authorization, reads committed state on the caller's node, and never activates the actor, takes the generation fence, writes a receipt, or uses a command id; its framework errors are `ActorUnavailable`, `Unauthorized`, and `Timeout`. Authorization is checked again before the result is returned, so a caller revoked while the handler ran gets `Unauthorized`; a query read past `executionTimeout` gets `Timeout` and is cancelled on Postgres. A query handler making a request/reply call is a deterministic defect. `internal` holds commands only. With `Database.postgres({ replica })`, queries read that streaming replica once it has replayed the caller's commit version, and read the primary when it is behind or failing, or when the actor type declares owned tables or blobs ([ADR 0052](../decisions/0052-read-your-writes-commit-versions.md)); in-process queries wait for the highest version any command sent through their runtime returned.
* Intents and timers ([contract 05](../contracts/05-messaging.md)). `X.intents(id)` (`X.intents()` for a singleton) returns one method per `api` and `internal` command (reducers are not intent targets), except an `internal` command named as a subscription `handler`, which only subscription deliveries reach; each method returns `Effect<void, never, Actor.InTurn>`; `id` is a plain string checked against the key schema at runtime, and the target shares the sending turn's tenant. `Actor.InTurn` is provided only inside command turns and removed by `X.toLayer` (providing a hand-built `Actor.InTurn` compiles, but any intent staged under it dies), so `X.intents` and `Intent.cancel` outside a command turn, including in `X.toQueryLayer` handlers, leave an unsatisfiable requirement. A command handler that acquires a handle with `X.get` makes `X.toLayer` fail to compile with `Request/reply inside a turn: use X.intents(id)`; a captured handle still dies at runtime. `Intent.after(duration)`, `Intent.at(dateTime)`, and `Intent.key(key)` pipe onto an intent (or any Effect that stages intents) the way `Actor.commandId` does; the outermost setting wins. `Intent.key` names the intent within its sending actor, and staging the same key again, in the same or a later turn, replaces the pending row. `Intent.cancel(key)` deletes the sender's pending row with that key when the turn commits. An intent Effect run outside the turn that staged it dies with `Intent capability escaped its turn`.
* Staged intents are written to `actor_outbox` (migration `0004_outbox`) on the sender's `routing_key` in the commit statement; a declared failure or rollback writes none. Due times use the database clock. After commit, a runner's relay claims due rows from every `(bucket, kind, due_at_ms)` index range (the bucket is the top eight bits of `routing_key`) and delivers each claimed row as a direct command with the intent id as its command id and `System({ source: "actor" | "timer", ref: <sender>, onBehalfOf: <sender turn principal> })` as its caller, and deletes the row after the receiver's receipt (success or declared failure) commits. Delivery may repeat; the receiver's receipt deduplicates it. An intent id is a v1 command id whose expiry is its due time plus the deployment retry window, which becomes the receiver receipt's `expires_at_ms`. Delivery is trusted internal recovery: it does not call `access` or `authorize` or check command-id expiry, so revoking a caller does not stop its committed intents. Any delivery without a committed receipt (a receiver defect, `NotCreated`, an unregistered or unavailable receiver, a delivery timeout, or an unreadable row) leaves the row and retries it after `min(2^(attempts − 1), relay.maxBackoff)` seconds without limit; intents have no dead letter yet, and the operator signal is the `Outbox delivery failed; retrying with backoff` warning (annotated with the reason) plus `actor_outbox.attempts`. Replacing or cancelling a key deletes the pending row; a timer whose delivery has already started is still delivered once ([contract 05](../contracts/05-messaging.md)). Every runner claims rows with `FOR UPDATE SKIP LOCKED` (migration `0011_relay`), only as many as it has free delivery slots (`relay.deliveryConcurrency`, default 16, at most `relay.passLimit` per claim), and each claim moves the row's `due_at_ms` past `max(relay.claimLease, backoff(attempts))`, so no runner scans it again while its delivery runs and a claimed row never waits locally. Every settling write names the claim it holds, so a runner whose lease has passed to another changes nothing. A relay that dies holding a claim (a crash, or a SQL error on its delete or reschedule) leaves the row until the lease ends; any runner then redelivers it, and the receipt deduplicates. A row whose settle keeps dying is therefore tried at most once per `relay.maxBackoff` and never holds back newer rows. `actor_outbox.attempts` counts claims, and `scheduled_at_ms` keeps the time a row first became due. Each runner polls every `relay.poll` (1 second, ±10% jitter), which is the correctness path; a commit that writes a row due now wakes only its own runner's relay, and a freed delivery slot claims again at once while more rows are due. Rows another transaction holds locked do not hide later due rows: a claim that skipped candidates widens the next claim's probe. No runner sends another a wake. There is no ordering guarantee between intents.
* Jobs ([contract 08](../contracts/08-background-work.md)). `const J = Actor.job(tag, { payload?, success?, progress? })` declares a job class: `payload` is a record of struct fields (the instance's fields, as in `Schema.TaggedClass`), `success` is the schema of the executor's return value, `Schema.Void` when omitted, and `progress` is the schema of the transient frames its executor may report. The value is itself the class: `J.make(...)`, `new J(...)`, `instanceof`, and schema decoding all work, and one job may be bound by several actors with different routes. An actor binds the jobs its turns may enqueue in `jobs: { [Tag]: { job: J, ...settings } }`, keyed by each job's tag. `turn.enqueue(J.make(...))` exists only on `X.Turn` and accepts only bound jobs (an unbound instance is a type error and dies at runtime); it records the job in the turn and returns immediately. An `enqueue` captured and run after its turn dies with `Job capability escaped its turn`. A binding takes `retry: { times }` (retries after the first attempt, an integer from 0 to 100, default 3), `progressEvery` (the least time between one attempt's progress sends, 50 ms to 1 minute, default 250 ms; checked against that range, then rounded down to whole milliseconds, like the other job timings), `timeout`, `concurrency`, `onSuccess`, `onDeadLetter`, and `onCancelled`. Routes must name a command in this actor's `api` or `internal` (checked in types and at `Actor.make`); `onSuccess`'s payload must accept `J`'s `success` type and `onDeadLetter`'s payload must accept `Actor.DeadLetter(J)`, or the definition does not compile. A job's executor failure and ambiguous outcomes follow the settlement contract below, so a job declares no `error`.
* `X.toJobLayer(executors)` implements one executor per bound job as `(job) => Effect<Success, unknown, R>`, directly or from an Effect that builds them; executors read `yield* X.Executor` (`jobId`, `attempt`, `principal`, the enqueueing actor's `ref`, and `progress(J, frame)` for a job that declares `progress`; see [context](02-context.md)). Progress is best-effort and never durable: a frame that does not encode, exceeds 4 KiB, or names another job is dropped with a warning, a call after its attempt ends does nothing. It reaches only the enqueueing actor's own connection members and streams that list the job under `progress` ([ADR 0030](../decisions/0030-executor-progress-frames.md)). If the executors' or the build Effect's requirements include `SqlClient`, `PgClient`, or `PgliteClient`, `X.toJobLayer` fails to compile with `Executors have no database capability`, and at runtime an executor's context is the layer's build context with those clients removed and `Tenant` set to the enqueueing actor's tenant. A client obtained outside Effect's context, such as a driver pool created in the build, is not detected. A runner without an actor's job layer never claims that actor's jobs; they stay due for a runner that has it.
* Enqueued jobs are written to `actor_outbox` as rows of kind `job` (migrations `0008_effects` and `0026_jobs`) in the turn's commit statement, targeting the enqueueing actor, with a v1 command id as the job id; a declared failure, defect, or rollback writes none. After commit, the relay runs each due job's executor. Each runner claims due jobs it has executors for, up to its pool's free permits (`executors.concurrency`, default 64), and runs them outside the relay pass. The claim records the attempt number and a lease (`executors.lease`, default 60 seconds) in the row, so a crash during the call is counted; the pool renews the lease every third of its length while the attempt runs, interrupts an attempt whose renewal finds the row claimed by another attempt, and interrupts one whose renewals have not reached the database for one lease, measured on the runner; an attempt whose claim is already one lease old when its executor would start is abandoned without calling it. An attempt times out after `jobs[Tag].timeout` (default 30 seconds). If the result does not encode under the `onSuccess` command's payload schema (for example a finite number routed to an integer payload), the provider has already applied the call, so the job is dead-lettered at once as `ambiguous` instead of executed again. When the executor succeeds, the same statement that records the result turns the row into an intent to `onSuccess` with the encoded return value (or deletes the row when there is no `onSuccess`), and the relay delivers it like any intent: a direct command whose command id is the job id and whose caller is `System({ source: "job", ref: <actor>, onBehalfOf: <enqueueing turn's principal> })`. The receipt keeps delivery to one committed `onSuccess` turn per job id; an executor that ran twice because a crash lost its first result routes only the recorded result. A failed attempt `n` retries after `min(base × 2^(n − 1), max)` from `jobs[Tag].retry.backoff` (default 1 second and 256 seconds). The first success of any attempt routes; a success that arrives after the job was dead-lettered routes nothing and marks the stored dead letter `ambiguous`, with the warning `Job succeeded after it was dead-lettered`. Two attempts of one job can overlap after a lease expires while the first still runs, so executors must be idempotent under `jobId`. When retries run out, the relay records the job in `actor_dead_letters` and, in the same transaction, turns the row into an intent to `onDeadLetter` with `Actor.DeadLetter(J)` payload `{ jobId, job, attempts, cause, ambiguous }`, delivered once in a new turn with the job id as its command id; without `onDeadLetter`, or when the stored job no longer decodes under its class, the row is deleted and the dead letter stays for operators. `ambiguous` is `false` only when the last attempt failed with a typed error; a defect, timeout, interruption, or crash after the claim leaves the provider outcome unknown, so it is `true`. Executors must treat a typed failure as "the provider did not apply the call" and use `jobId` as the provider's idempotency key, because a later attempt may follow an attempt whose outcome is unknown. Warnings `Job attempt failed; retrying with backoff` and `Job dead-lettered after its last attempt` are the operator signals. Job cancellation and per-actor concurrency caps are implemented by M2.13 ([ADR 0024](../decisions/0024-effect-cancellation-and-per-actor-concurrency.md)); see [job cancellation and caps](#job-cancellation-and-caps). A slow executor holds a pool permit, not a delivery slot, so it does not delay intents or timers. See [multi-runner relay and cron](#multi-runner-relay-and-cron).
* A deterministic defect rolls back, writes no receipt, returns `Die`, logs `Deterministic actor defect` with actor, id, tenant, command, and command id, and runs inside the span `akter.<Actor>/<Command>`. No user hook runs.
* Constructing a command Effect creates one operation. Its first execution mints an ID using the database clock; rerunning that same Effect reuses it. To retry across processes, save an ID from `(yield* Actors).mintCommandId` and apply `call.pipe(Actor.commandId(id))` before its first execution. Minting reads the database clock, so while the database is unreachable `mintCommandId`, and a call that mints its own ID, fail `ActorError` `ActorUnavailable`, which a caller retries like any other delivery failure. Every external retry rechecks access and expiry.
* `/runtime` exports `Actors.layer({ authorize?, retryWindowMs?, maxResidentActors?, admission?, rowLevelSecurity?, payloadWriterWindow?, observability? })`, where `observability.defects` (default 1,000) bounds the defect spans a runner keeps for `akter defects list` and `observability.sampleEvery` (default 15 seconds) is how often one runner samples the database gauges ([ADR 0049](../decisions/0049-observability-names-metrics-and-defect-spans.md)), `Database.postgres({ url: Redacted.make(url), maxConnections?, ... })`, and `Database.pglite(config?)`. `maxResidentActors` (default 10,000) is how many activations one runner keeps in memory. A command that needs a new activation past it fails `RunnerAtCapacity`, and the handle retries it until an idle actor hibernates or `deliveryTimeout` passes. Raise it only with the memory you give the process. `admission: { concurrency?, wait?, requests? }` sheds load a runner cannot serve promptly ([ADR 0077](../decisions/0077-admission-control.md)): at most `concurrency` (default 64) external commands run at once, from admission to reply; up to as many more wait for a slot in arrival order for at most `wait` (default 100 ms); any other command is refused at once with `ActorUnavailable` and a `retryAfter`, before anything of it runs, so it writes no receipt and its retry under the same command id runs it once. That refusal is not retried inside the runtime, so an in-process caller sees it at once. A caller disconnect or delivery timeout stops its reply wait but keeps the executing attempt's slot until its runner RPC settles. `Actors.serve` also refuses HTTP actor command routes before authenticating or reading them when `admission.requests` (default 64) command request handlers are already processing authentication, bodies, or execution, or when runtime admission is full; queries do not use the command gate, but their pool checkouts may be refused with `ActorUnavailable`. MCP tool calls use runtime admission after envelope authentication and parsing. An actor whose `policy` declares no `mailboxCapacity` still holds at most 1,024 commands, waiting and in its current batch; the next is refused with `ActorUnavailable` rather than `MailboxFull`, which stays reserved for a declared capacity. `fleet` registers [fleet views](04-drizzle.md#fleet-views) the runtime maintains, and `Actors.serve({ fleet })` serves them at `GET /fleet/{View}`; it needs Postgres with `wal_level=logical` and `akter fleet setup`, and PGlite refuses it. `rowLevelSecurity: { role }` opts in to row-level security: command turns and queries run as `role` with their actor's tenant in `durable.tenant`, and startup fails unless the database is set up as [deployment](../operations/01-deployment.md#row-level-security) describes ([ADR 0051](../decisions/0051-row-level-security.md)). `maxConnections` (default 50) is the turn pool: each command turn leases one of its sessions for the whole turn and pipelines its statements on it. `offTurnConnections` (default 10) is a second pool for queries, command-id minting, the relay, migrations, and Cluster runner storage. Both pools open connections only as load needs them, so one runner holds up to `maxConnections + offTurnConnections` (60 by default); keep the total across runners below the server's `max_connections`. Supply a platform `Crypto` layer, such as `BunCrypto.layer`. The optional `authorize` callback is the global hook: it receives the captured `caller`, `ref`, `command`, and `kind` at admission and before returning an outcome, beside the actor's own `access` policy ([ADR 0059](../decisions/0059-caller-and-tenant-defaults.md)). `User.make({ subject })` and `System({ source, ref?, onBehalfOf? })` are application-trusted attribution, not authentication. Code in the application's own process, outside a turn and not through `Actors.serve`, runs as `System({ source: "process" })` in tenant `"default"`; `Actor.as(caller)` and `Actor.tenant(tenant)` are opt-in for trusted code acting on behalf of a user or tenant, and scope the caller and tenant around acquiring handles; `get` takes no options.
* `bun create @akter [dir] [--template counter|chat]` runs `@akter/create`, which copies a standalone app pinned to the core's exact versions. The app's `src/database.ts` provides `Database.postgres({ url })` when `DATABASE_URL` is set and otherwise `Database.pglite({ dataDir })`, with `DATA_DIR` defaulting to `./.data`. File-backed PGlite serves one process per data directory, for development and for one-process production within the limits of ADR 0035; see the [quickstart](../quickstart.md).
* `Database.pglite({ dataDir })` with a filesystem `dataDir` (M4.14, [ADR 0035](../decisions/0035-pglite-embedded-production-backend.md)) takes an exclusive, non-blocking `flock` on `<dataDir>/.akter.lock` before PGlite opens and holds it until the layer's scope closes. Layer build fails with `DataDirLocked { dataDir }` while another process, or another open layer in this one, holds it, and with `DataDirVersion { dataDir, found, expected }` when `<dataDir>/PG_VERSION` names another Postgres major than the pinned PGlite embeds (18); both are exported from `/runtime`. `relaxedDurability: true` with a `dataDir` is a defect. In-memory PGlite (no `dataDir`, or `memory://`) takes no lock. Linux and macOS only; the lock is advisory and unreliable on network filesystems.
* `Database.pglite` creates a fresh owned instance per layer build; a supplied `liveClient` config keeps its own methods and lifetime. Owned instances drain tracked queries before closing — the pinned driver deadlocks when closed mid-exchange. Cluster runner bookkeeping uses in-memory storage on PGlite because `SqlRunnerStorage` would reserve the sole connection; this confers no independent-connection or concurrent-ownership guarantee. `Database.postgres` includes a scoped `regclass` codec for the pinned rc.116 driver restart bug; it does not modify global driver configuration.
* `/runtime` also exports `Runner.socket({ address, listenAddress?, transport, shardsPerGroup?, shardLockExpiration?, shardLockRefreshInterval?, refreshAssignmentsInterval?, entityTerminationTimeout? })`, provided to `Actors.layer`. `transport` is a platform TCP server/client layer, such as `Layer.merge(layerSocketServer, layerClientProtocol)` from `@effect/platform-bun/BunClusterSocket`; Akter supplies NDJSON serialization. Defaults are 256 shards, table leases, 35-second expiration, 10-second lock refresh (capped at a third of expiration), one-second assignment refresh, and 15-second entity shutdown. Advertisement must be unique and directly reachable by every peer; binding can differ. Table leases retain the singleton lease check. Postgres is required. Count, mode, and expiration are persisted and incompatible joins are refused. Readiness reports `routing` until registered and holding every currently assigned shard. `Runner.mtls({ deployment, credentials, refreshEvery?, refreshTimeout?, unhealthyAfter? })` is the authenticated transport: TLS 1.3 without session resumption, where both sides require the peer's chain to reach an authority in `ca`, the certificate to be valid now, and its subject alternative names to be exactly the URI `Runner.identity(deployment)` (`spiffe://akter/deployment/<deployment>`). `credentials` is an Effect of PEM `{ ca, certificate, key: Redacted }`, run at startup and every `refreshEvery` (default one minute) with a `refreshTimeout` (default 10 seconds); invalid or late credentials refuse startup, or on refresh are logged and ignored. Readiness reports `peering` once the current certificate has expired or loads have failed for `unhealthyAfter` (default 5 minutes). `RunnerAuthority.make()`, `.from({ certificate, key })` and `authority.issue({ deployment, validFor?, notBefore? })` create development credentials. The plaintext platform layers trust every peer, so use them only on a network isolated to one deployment's runners. See [deployment](../guides/deploy.md#several-runners) and [ADR 0086](../decisions/0086-runner-mutual-tls.md).
* The retry window defaults to 86,400,000 ms, accepts 1–2,592,000,000 ms, and is recorded in the database. A different setting fails startup. State, receipt, and creation-marker migrations run at startup, so use a disposable database for the example. Protocol/ID v1 is unchanged; no arbitrary rolling mixed-version support is claimed.
* `Database.postgres({ neki: true, ... })` opts into the prepared Neki startup migration protocol and single-shard turn sessions. It also provides `Database.Neki`, a context reference (default `false`); a process that opens its own Postgres client and runs no actors provides `Layer.succeed(Database.Neki, true)` itself. `Database.schemaChange(effect, lock)` (or `effect.pipe(Database.schemaChange(lock))`) runs an application's own startup schema statements on the ambient `SqlClient` under advisory lock `lock`: on Postgres in one transaction; on Neki on one reserved session, each statement autocommitted and each DDL statement followed by a propagation wait, so every statement must be safe to rerun after a crash ([ADR 0091](../decisions/0091-neki-control-plane-database.md)). Startup DDL runs outside transactions with propagation barriers and durable statement progress; ordinary Postgres and PGlite keep transactional migrations. Neki remains unverified, including advisory-lock placement, and requires migration metadata in an authoritative unsharded group. See [Neki migration operations](../operations/02-migrations.md#neki-mode-prepared-unverified).
* `/testing` exports `ActorTest.layer({ database?, as?, authorize?, retryWindowMs?, maxResidentActors?, admission?, relay?, executors?, observability?, rowLevelSecurity? })`, with a fresh tenant per build, committed `inspect(ref)` (`generation`, `state`, `receipts`, and pending `outbox` rows the actor sent), `receiptsFor(ref, command)`, `crashNext(point, { commandId? })`, `pauseNext(point, { commandId? })` (with `commandId`, only that command's turn takes the fault, so an intent delivery or another command passing the same point does not), `invalidate(ref)`, `advance(duration)`, and `now`. `advance` moves the leases of this runner's running executor attempts forward by the same duration, as their renewals would, then moves the outbox clock forward and waits for in-flight deliveries and executor attempts and claims again until a claim finds nothing and nothing is running, including intents those deliveries stage; a row whose settle died stays claimed until the clock moves past its lease; `now` is that clock as a `DateTime.Utc`, for `Intent.at`. `database` defaults to a fresh in-memory PGlite; a `dataDir` config persists across builds and a `Redacted` URL uses Postgres. `test.actor(X, id?, { seed? })` returns `{ system, inspect }` where `system` reaches every command including internal ones with a `System` caller that inherits the configured principal. `seed` is the path of a seed file `akter export` wrote, read through Effect's `FileSystem`, so a test that names one provides a file system. In one transaction it writes the actor's row, state, and pending work, staged as the caller the test runs as and due relative to the test's clock, so `advance` delivers it; the actor starts at generation 0 with no receipts or events. An existing actor, a seed of another actor type, a malformed or other-format file, and a job the actor does not bind each die and write nothing. `checkBatchLaw` checks a batched reducer's merge law over generated states and input lists and dies with a reproducible seed on failure. Fault points are `beforeDelivery`, `queued` (a command is in its actor's mailbox), `beforeHandler`, `beforeCommit`, `afterCommit`, `beforeFlush` (reached only by actors with connection members, after the turn commits and before its broadcasts flush), and the relay's `afterClaim`, `beforeOutboxDelete`, `beforeExecute`, `afterExecute`, and `beforeRenew`. `TurnHooks` is available only through the testing entry for process-level faults; a `TurnHooks` service provided around `ActorTest.layer` still sees every point no queued fault takes. `test.connect(ref, Member, params)` opens a connection held by this runner's in-process transport and returns `{ connectionId, cursor, send, frames, messages, resyncDone, close }`: `cursor` is the flushed-through event cursor at open, `frames` is a `Stream` of decoded server frames and `messages` the same with the holder's `Resync` and `ResyncReplayed` control messages and each frame's cursor stamp, both failing with the session's `ActorError` when it ends; `resyncDone` answers a `Resync`. A declared `open` failure fails `connect` with that error. `test.progress` lists every progress message this runner's executor pool sent, in order, each marked `dropped` when `test.dropProgress(predicate)` dropped it between the pool and the owner; nothing delivers progress to an owner yet, so this test sink is the only receiver. `test.hibernate(ref)` ends the actor's activation on this runner as `hibernateAfter` would, sealing every holder so its connections stay open without a resync. The in-process transport is not exported from `/runtime` yet; `Actors.serve` holds WebSocket sessions through the same holder. `/testing` also exports the framework-neutral `conformance` cases and `describeConformance`; the same named cases run on PGlite and Postgres, and cases needing a second connection are reported skipped on PGlite.
* `/testing` also exports `ActorTest.simulateCluster({ seed, faults, faultRate?, primary?, commands?, attempt?, within?, settle? }, program)` and `clusterSimulationSeeds`, the multi-runner counterpart of `ActorTest.simulate`, run inside an `ActorTest.cluster` layer. `simulation.command(label, effect)` takes an effect that gets its own handles, since it runs on a seeded runner, and sends its minted command id again through the runners in turn until a receipt answers. Faults are `crashBeforeCommit`, `crashAfterCommit`, `dropReply`, `runnerKill`, `heartbeatLoss`, `connectionLoss`, and `primaryFailover`; see deterministic simulation.
* `/testing` also exports `ActorTest.cluster({ database, runners, shardLockExpiration, actors, runnerActors?, holdersOnly?, as?, authorize?, retryWindowMs?, maxResidentActors?, relay?, executors? })`, a layer providing `ActorCluster`. It runs `runners` actor runtimes in one process against one Postgres database; each is a distinct Cluster runner with its own address, connection pool, and SQL shard locks that expire `shardLockExpiration` after its last heartbeat (whole seconds, rounded up). `actors` is the application layer every runner builds (`X.toLayer`, `X.toQueryLayer`, ...), and `runnerActors(index)` adds layers only some runners build, such as a job layer one runner lacks; `holdersOnly` lists runners that hold connections but are never assigned actor shards, so killing an actor's owner never kills the holder of its connections; all runners share one fresh tenant. PGlite is refused with an error, because its single connection cannot host several runners. The database must belong to this cluster alone, since Cluster's runner and lock tables are shared by everything that connects to it. Runners call each other over an in-process transport that encodes every message as a socket would. The layer waits until every runner that hosts actors holds the shards the hash ring assigns it. `ActorCluster` offers `on(runner)(effect)`, which runs `effect` with that runner's `Actors` and `ActorTest`, so handles dispatch through it and `crashNext`, `pauseNext`, `inspect`, and `advance` belong to it; `kill(runner)`, which cuts the runner's database and runner connections at once, so Postgres rolls back its open turns and its locks and heartbeat are left to expire; `restart(runner)`, which starts it again as a new process under a new address; `pauseHeartbeat(runner)`, which returns `{ resume }` and stops the runner writing heartbeats and lock refreshes while it keeps serving the shards it holds, so other runners take them after expiry and only the generation fence separates the two; `owner(ref)`, the index of the runner holding an unexpired lock on the actor's shard, or `undefined` (while a runner's heartbeat is paused it may still serve the actor too); `holdReplies`, which delays every database reply to every runner without dropping it, so a COMMIT sent meanwhile is made by the database and unknown to its runner; `failover`, which cuts every runner's database connections at once, discards the held replies, and lets the runners reconnect to the same database; `tenant` and `sql`, the cluster's tenant and a pool of its own for inspection; and `ready`, which waits for the shards to settle again. An effect already running through a runner when it is killed is not interrupted; its calls fail or time out, so drive retries through a survivor. A command caught by a kill waits about `shardLockExpiration` for a new owner, so keep the expiration well below the actors' `Delivery.timeout` (30 s by default). For quick detection the harness polls runner changes every 250 ms (Cluster's default is 3 s) and uses 32 shards.

See the [protocol](../decisions/0007-foundation-command-protocol.md) and [reconciliation](../decisions/0013-m0-reconciliation.md) decisions. Support remains scoped by the [support matrix](../operations/support-matrix.md); public socket runners have three-process loopback evidence on Postgres, not provider or separate-host certification.

## Definitions

`Actor.make(name, definition)` is the only way to make an actor, and the definition is its only shape: there is no piping or later configuration. Every section is data; code lives in layers. See [ADR 0010](../decisions/0010-one-way-effect-native-api.md).

```ts title="counter.ts" theme={"theme":"css-variables"}
import { Effect, Result, Schema } from "effect"
import { Actor } from "@rikalabs/akter"

export const CounterId = Schema.String.pipe(Schema.brand("CounterId"))
export class Overflow extends Schema.TaggedError<Overflow>()("Overflow", { max: Schema.Int }) {}
export const CountChanged = Actor.event("CountChanged", { count: Schema.Int })

export const CounterState = Actor.state({
  count: Schema.Int.pipe(Schema.withDecodingDefault(Effect.succeed(0))),
})

export const Increment = Actor.reducer("Increment", {
  state: CounterState,
  payload: Schema.Int,
  error: Overflow,
  reduce: (state, amount) =>
    state.count + amount > 1_000
      ? Result.fail(new Overflow({ max: 1_000 }))
      : Result.succeed({ count: state.count + amount }),
})
export const Reset = Actor.command("Reset")
export const GetCount = Actor.query("GetCount", { success: Schema.Int })

export const Counter = Actor.make("Counter", {
  key: CounterId,
  state: CounterState,
  events: [CountChanged],
  api: { Increment, Reset, GetCount },
  schedules: { "0 * * * *": Reset },
  policy: { hibernateAfter: "30 seconds" },
})
```

| Section | Content |
| - | - |
| `key` | id schema (named, `X.get(id)`), `Actor.singleton` (`X.get()`), or omitted (minted, `X.create()` or `turn.mint(X)`) |
| `placement` | `"tenant"` (default), `"actor"`, or `{ parent: P }` ([ADR 0033](../decisions/0033-parent-actor-placement.md)) |
| `state` | one `Actor.state(fields, { migrations })`; missing keys decode from defaults |
| `tables` | `Actor.table` Drizzle tables with framework ownership columns |
| `blobs` | `Actor.blob` database-backed `bytea` chunks |
| `events` | `Actor.event` values |
| `feeds` | events from `events` that `Actors.serve` serves as SSE event feeds; none are served unless listed |
| `createdBy` | the one command that creates the actor; other commands fail `NotCreated` until it commits |
| `schedules` | cron expressions or `@every` intervals mapped to commands whose payload is `Schema.Void` |
| `jobs` | `{ [Tag]: { job, timeout?, progressEvery?, retry?, concurrency?, onSuccess?, onDeadLetter?, onCancelled? } }`: the `Actor.job` values turns may enqueue, with this actor's settings and routes |
| `api` | public commands, reducers, queries, streams, connections, and workflows; each key equals its member's tag; omitted for an actor with none |
| `internal` | commands callable only by System callers: outbox intents, job routes, schedules, and subscription deliveries |
| `subscriptions` | `Actor.subscription` members that follow another actor type's events (M3.7, [ADR 0026](../decisions/0026-cross-actor-event-subscriptions.md)) |
| `policy` | `hibernateAfter`, `executionTimeout`, `lockWait`, `deliveryTimeout`, `maxStateBytes`, `mailboxCapacity`, `keepReceipts`, `keepEvents`, `maxBlobBytes`, `maxBlobEntries`, `connections`, `reauthorizeEvery`, `watch` (`maxPerActor`, `minInterval`, `reconcileEvery`), `maxScheduleLag`, `keepWorkflows`, `allowedSubscriberTypes`, `holdEventsForSubscribers` |

Members:

* `Actor.command(tag, { payload?, success?, error? })` runs an effectful server handler. `payload` is a schema or struct fields (standing for `Schema.Struct(fields)`) and gives the handle method's argument; an omitted payload is a zero-argument call. `error` is one tagged error class or a `Schema.Union` of them; framework reasons never masquerade as declared failures. Listing it under `internal` instead of `api` removes it from public handles and transports.

* `Actor.reducer(tag, { state, payload?, error?, reduce, batch? })` is a pure transition with no server handler; `state` is the actor's own `Actor.state`, and `reduce` returns a `Result` of the next state or a declared error. It runs optimistically in browser handles; with `batch: { combine }` queued calls fold in their order into one turn, it returns `void`, and declares no error. `combine` needs ordered fold equivalence, not commutativity.

* `Actor.query` and `Actor.workflow` declare reads and durable workflows.

* `Actor.query(tag, { payload?, success?, error?, watch? })` with `watch: true` (M6.2, [ADR 0055](../decisions/0055-query-observation.md)) declares a query that can also be observed. Its handler in `X.toQueryLayer` may require only `X.Read`: a handler that needs another service does not compile, and at runtime a watched handler is provided no service but `X.Read`, so one that got around the types fails with a defect (a watch ends with `SessionEnded` cause `Defect`). The in-process handle gains `handle.Query.watch(input)`, a `Stream` of the query's outputs that fails with its declared errors or `ActorError` reason `ActorUnavailable`, `Unauthorized`, `RunnerAtCapacity`, `SessionEnded`, `NotCreated` (the actor has no generation row; a watch never creates one), or `Timeout`. The stream starts with the current result and then sends the newest result after each commit that wrote state, an emitted event class, a table, or a blob that the last run read; it skips intermediate states, sends nothing for a result equal to the last, and is not an event history (follow events with `read.follow` for that). The runtime records the reads per rerun: `state` (all keyed state), `read.events(E)` per class, `read.rows(T)` per table, and `read.blob(B)` per blob. A rerun that reads `read.group` is a defect of the handler in process (the watch dies with `InvalidInput { code: "not_watchable" }` as its cause), because `InvalidInput` is boundary-only; a served watch ends with it. `read.id`, `read.ref`, `read.caller`, and `read.principal` are recorded but do not narrow a rerun. A handler that reads the clock or a random source gets a result as of its rerun and no invalidation for it; the reconcile rerun bounds the staleness. A watch is a parked connection: it keeps the activation resident no more than an idle socket does, it counts apart from streams and feeds, and it survives its actor's hibernation. Retention pruning of a read event class reaches a watch only through the reconcile rerun. Singleton actors cannot declare a watchable query yet. It is served at `POST /actors/{Actor}/{id}/{Query}/watch` over SSE (see [protocol](../contracts/protocol.md)), exposed as `handle.Query.watch(input, { signal? })` on the Promise client, and followed by `useWatch` in `@akter/react`. OpenAPI lists it as the operation `<Actor>.<Query>.watch` (a session, so the MCP endpoint lists no tool for it, as for streams) with `x-durable-transport: sse` and `x-durable-element` referencing `<Actor>.<Query>.element`; a served watch of a query without `watch: true` is `400 InvalidInput { code: "not_watchable" }`, and a rerun that reads `read.group` ends it with an `end` message carrying the same error.

* `policy.watch` sets the limits of an actor's watches: `maxPerActor` (default 1,000 open watches; the next is `RunnerAtCapacity`), `minInterval` (default 100 ms between reruns of one watch), and `reconcileEvery` (default 30 s, from 5 s to 1 h). A runner runs at most 64 watch reruns at once; a rerun beyond that waits and starts from the newest version seen.

* `Actor.stream(tag, { payload?, success, error? })` declares a live, non-persisted feed. Its handler in `X.toLayer` takes the payload and returns a `Stream` of `success` values that runs on the activation for one subscriber, with `X.Read` plus `read.follow(Event, { after })`; it fails only with its declared `error`. The handle method returns a `Stream` of the same elements that fails with those errors or `ActorError` reason `ActorUnavailable`, `Unauthorized`, `RunnerAtCapacity` (past 256 subscriptions on one actor), or `SessionEnded`. An open subscription keeps the activation resident whatever `policy.connections` says; the activation's end (a move, a runner's death, eviction, or shutdown) reaches the subscriber as `SessionEnded { cause: "ActivationEnded", resync: false }`, never a silent completion, and a handler whose stream ends by itself completes the subscriber's stream. Elements pass through a 256-element window on the owner; one that stays full for 30 seconds ends the subscription with `SessionEnded { cause: "SlowConsumer", resync: true }`. The subscriber's runner and then the owner call `authorize` with `kind: "stream"`, and the owner reauthorizes every `reauthorizeEvery`; a denial ends it with `Unauthorized`, discarding elements not yet delivered. `Actors.serve` serves streams over SSE (M3.3, below). Singleton actors cannot declare them yet.

* `Actor.connection(tag, { payload?, server, client?, session?, error?, stampCursor?, progress? })` declares a typed session: `payload` is the opening params, `server` the frames the actor sends, `client` the frames the client sends, `session` the per-connection fields (at most 16 KiB encoded), and `error` the declared failures of opening. Its `X.toLayer` entry is `{ open, frame, close?, resync? }`: short handlers that run with `X.Connection`, not one long-lived stream, so the activation can hibernate between frames. Frames reach the activation through a framework connection entity beside the command entity, never through the command mailbox.

* `progress: { jobs, to? }` on a connection member, and `progress: { jobs }` on a stream member, opt into executor progress of jobs the actor binds in `jobs` that declare a `progress` schema (anything else is rejected at `Actor.make`). A connection receives `Progress { job, jobId, attempt, seq, frame }` envelopes, apart from member frames and never replayed: by default (`to: "principal"`) only connections whose caller has the enqueueing turn's principal, with `to: "all"` every open connection of the member. A stream handler reads `read.progress(J, { jobId? })`. Progress is coalesced (newest per job), may be lost (a gap in `seq`), never causes `SlowConsumer`, and stops once the job's route, cancellation, or settle commits. An actor type with no such member costs nothing: executor pools send it no progress. Over the served WebSocket, progress arrives as the server message `t: "progress"`, never inside a `frame` message.

* `policy.connections` is `"park"` (default: open connections do not keep the activation resident) or `"keepAwake"` (they count as activity; residency is still not guaranteed). `policy.reauthorizeEvery` is the revocation bound for connections and streams: default 60 s, from 1 s to 1 h.

* The actor's `access` policy and the runtime's `authorize` hook are also called on every open and stream subscription and then every `reauthorizeEvery`, with the member tag as `command` and a new `kind` field (`"command" | "query" | "open" | "stream" | "feed" | "watch" | "reauthorize" | "content"`, the last for a content grant or read with `command` set to `<blob>.grant` or `<blob>.get`). A query watch is authorized with `kind: "watch"` and the query tag as `command`, before anything is read, and reauthorized with `kind: "reauthorize"` and `of: "watch"` every `reauthorizeEvery`; reruns use the caller stored at open and call neither. An event feed is authorized with `kind: "feed"` and each event tag it reads as `command`. A reauthorization carries `of` (`"open"`, `"stream"`, `"feed"`, or `"watch"`), so a feed's event tag can't be mistaken for a connection member with the same name. A policy or hook that allows unknown names or kinds allows sessions, so both should deny kinds they don't know.

* Connections, streams, and these two policies are implemented (M2.10, [ADR 0023](../decisions/0023-connections-parking-and-streams.md)): connections with the in-process holder transport described under `/testing`, and streams as a streaming request from the subscriber's runner to the owner. Connection members are served over WebSocket and streams over SSE (M3.3, below).

* `const E = Actor.event(tag, fields)` declares a durable event. The value is a `Schema.TaggedClass` whose identifier is its tag (`E.make`, `new E`, and `instanceof` work; the event type is `typeof E.Type`), and the tag is stored with each committed event, so tags must be unique within an actor's `events`. Listing it in `events` lets that actor's command turns `turn.emit` it and its queries replay it with `read.events` ([context](02-context.md#events)). Changing an event's schema needs the same care as changing a command's payload: stored events are decoded with the current class.

* `Actor.job(tag, { payload?, success?, progress? })` declares a class of external I/O that turns stage with `turn.enqueue` and `X.toJobLayer` executors run after commit. `Actor.DeadLetter(J)` is the schema of `{ jobId, job, attempts, cause, ambiguous }`, the payload of an `onDeadLetter` command. `Actor.Cancelled(J)` is the payload of an `onCancelled` command ([ADR 0024](../decisions/0024-effect-cancellation-and-per-actor-concurrency.md); implemented in M2.13).

* `Actor.subscription(tag, { delivery, retired?, handler, route? })` (M3.7, [ADR 0026](../decisions/0026-cross-actor-event-subscriptions.md)) declares a subscription to the event classes its `delivery` names. `delivery` is `Actor.Delivery({ source, events })`, declared once and also used as the handler's payload schema; it records its `source` and `events`. `handler` is an `internal` command, reachable only by subscription deliveries, whose payload accepts that delivery; several subscriptions may share one handler across sources or event subsets. A delivery is a union of `Event` (`source`, `cursor`, `event`, `commandId`, `timestamp`), `RetentionGap` (`after`, `resumeAfter`), and `Rejected` (a dynamic subscription whose `from` cursor the source refused). With `route: (event, source) => id` (or `Actor.singleton`) every matching event from any source in the tenant reaches the routed subscriber; without it, the subscriber follows the source ids its turns pass to `turn.subscribe`. Deliveries are System command turns through the relay, applied once per source cursor and in cursor order per source; they wake a sleeping subscriber. `policy.allowedSubscriberTypes` on a source restricts subscriber types by name, and `policy.holdEventsForSubscribers` (default 7 days) bounds how long subscribers hold its events from pruning. Implemented: declarations, routed and dynamic delivery, epochs, receipts and the applied cursor, `Rejected`, `RetentionGap` (an id-routed row counts its gap in `gaps` instead), `policy.allowedSubscriberTypes` (checked at `Actor.make`), `policy.holdEventsForSubscribers`, startup widening of dynamic rows, and deletion, during the retention sweep, of rows whose declaration a registered subscriber type no longer has once they have been due for a day. A subscription whose source type no runner has registered is inert rather than a startup failure, because layers of one runtime register concurrently. A routed subscription's source-side row is created by the publishing turn, from the routed declarations its runner registers, so every runner that serves the source type must also register the subscriber type's layer (its commands layer, which carries the declaration). The runtime records each routed declaration in `actor_routed_subscriptions`, and registering a source type fails, after waiting up to 5 seconds for concurrently built layers, when a subscriber type recorded as routing from it isn't registered on the same runner. A runner that dropped a routed declaration removes its record when it registers. During a rolling deploy that adds a routed subscription, runners already serving the source keep running; they are checked only when they start. Dynamic subscriptions have no such requirement.

Type checks replace lists that must agree: an `api`, `internal`, or `jobs` key must equal its member's or job's tag; `schedules`, `createdBy`, and job routes must name a command in `api` or `internal`; schedule targets take no payload; a binding's `onSuccess` command payload must accept its executor's return type; and its `onDeadLetter` command payload must accept `Actor.DeadLetter(J)`. The tag, `api` key, handler key, and handle method are the same PascalCase name.

### Multi-runner relay and cron

From [ADR 0021](../decisions/0021-multi-runner-relay-singleton-and-cron.md). The relay, executor, and per-job timing settings are implemented (M2.4), and so are `schedules` and `maxScheduleLag` (M2.5), with the time zones and fixed intervals of [ADR 0042](../decisions/0042-cron-time-zones-intervals-and-daylight-saving.md) ([below](#cron-time-zones-and-intervals)). Every default equals the M1 behaviour.

```ts theme={"theme":"css-variables"}
import { Actor } from "@rikalabs/akter"
import { Actors } from "@rikalabs/akter/runtime"
import { Schema } from "effect"

const Moderate = Actor.job("Moderate", {
  payload: { body: Schema.String },
  success: Schema.Boolean,
})
const Moderated = Actor.command("Moderated", { payload: Schema.Boolean })
const Send = Actor.command("Send")

export const Newsletter = Actor.make("Newsletter", {
  key: Actor.singleton,
  internal: { Moderated, Send },
  schedules: { "0 8 * * *": Send },
  jobs: {
    Moderate: {
      job: Moderate,
      timeout: "30 seconds",
      retry: { times: 3, backoff: { base: "1 second", max: "256 seconds" } },
      onSuccess: Moderated,
    },
  },
  policy: { maxScheduleLag: "1 day" },
})

export const runtime = Actors.layer({
  relay: {
    poll: "1 second",
    passLimit: 256,
    deliveryConcurrency: 16,
    subscriptionConcurrency: 16,
    subscriptionBatch: 16,
    claimLease: "37 seconds",
    maxBackoff: "256 seconds",
  },
  executors: { concurrency: 64, lease: "60 seconds" },
})
```

* `timeout` bounds one executor attempt; attempt `n` that fails waits `min(base × 2^(n − 1), max)`. Durations run from 1 ms to 2^31 − 1 ms and `max ≥ base`; `Actor.make` rejects anything else.
* `relay.subscriptionConcurrency` (subscription delivery slots per runner, separate from intents; feed expansions and control registrations each get as many) and `relay.subscriptionBatch` (events read per subscription claim) are implemented (M3.7, [ADR 0026](../decisions/0026-cross-actor-event-subscriptions.md)). A runner that registers a subscriber type claims feed, control, and subscription rows in the same statement as its intents and jobs, so a relay pass stays one round trip. A feed expansion that makes rows due leases the ones this runner can deliver, up to its free subscription-delivery slots, in the same statement, and starts their deliveries at once, so a delivery right after a commit needs no separate claim pass. A row another runner has claimed meanwhile is left alone.
* `relay.claimLease` defaults to the largest registered `executionTimeout + lockWait` plus 5 seconds. `executors.lease` is renewed every third of its length while an attempt runs and must be at least 3 seconds; `Actors.layer` throws otherwise. Counts run from 1 to 1,000,000 and durations from 1 ms to 2^31 − 1 ms. On graceful shutdown the relay stops claiming, releases every intent it claimed but has not settled (an interrupted delivery whose receiver already committed replays the receipt on redelivery), and interrupts running attempts, which keep `ambiguous = true` and are taken over after their lease.
* **Changed from M1:** a runner without a job's executor never claims its rows, so the 5-second retry and its warning go away; a slow executor no longer delays intents; two attempts of one job can overlap after a lease expires, so executors must be idempotent under `jobId`; a relay that dies holding a claim delays that row's redelivery until the lease ends, so `ActorTest.advance` must move past the lease to redeliver it; `advance` waits for in-flight deliveries and executor attempts; `actor_outbox.attempts` counts claims, including expired leases, rather than failed deliveries; a success that arrives after its job was dead-lettered marks the stored dead letter `ambiguous`.
* Cron expressions are parsed with `Cron.parse` (five or six fields) and evaluated in UTC against the database clock unless they name a zone ([below](#cron-time-zones-and-intervals)); the timer key names the zone and spells the parsed schedule canonically (sorted value lists, `*` for a full field, seconds only when not `0`), so `"0  8 * * 1-5"` and `"0 8 * * 1,2,3,4,5"` are one key, `$cron:UTC 0 8 * * 1,2,3,4,5`, and a deployment that respells a schedule keeps its row. `Actor.make` rejects an unparsable expression, two equivalent schedules (such as `1-5` and `1,2,3,4,5`), a target that is not a command of the actor or takes input, and an invalid `maxScheduleLag`. A cron tick's caller is `System({ source: "cron", ref })`, so a target may be an `internal` command that handles, HTTP, and the Promise client cannot call. On a singleton, cron runs in the default tenant only.
* **Changed from M1:** `Intent.key` values starting with `$`, `$cron:` among them, are reserved: `Intent.key` throws and `Intent.cancel` dies with `Intent.key values starting with $ are reserved`.

### Runtime control: readiness and drain

Implemented in M4.2 ([ADR 0003](../decisions/0003-failure-scoping-drain-and-hosted-trust.md)). `Actors.layer` also provides `RuntimeControl`:

```ts theme={"theme":"css-variables"}
import type { Effect } from "effect"

import { RuntimeControl } from "@rikalabs/akter/runtime"

const ready = RuntimeControl.use((control) => control.readiness)
// { ready: true } | { ready: false, reason: "draining" | "drained" | "storage" | "routing" | "peering" | "unregistered" }

const report = RuntimeControl.use((control) => control.drain({ deadline: "30 seconds" }))
// { outcome: "clean" | "deadline-expired", interruptedTurns, interruptedJobs }
```

* **Readiness** is true once the database answers a `SELECT 1` within 2 seconds, the layer has built (so migrations and the deployment check passed), at least one actor, query, or job layer is registered, sharding is up, a public runner has healthy registration and every currently assigned shard, and the runner is not draining. It never waits for actors to wake or workflows to finish. Incomplete public-runner acquisition reports `routing`; registration snapshots are checked on each probe.
* **`drain({ deadline })`** makes the runner unready and refuses new external commands through its handles with `ActorUnavailable`. Its actors refuse new turns the same way, so callers on other runners retry until the actors move. The relay stops claiming intents, timers, jobs, and subscription deliveries, and releases deliveries it had claimed, which are due again at once. In-flight turns and job attempts then get until `deadline`, which counts from the call and covers stopping the relay and background sweeps too: releasing a claimed delivery writes to the database, and when a slow or contended database holds that write past the deadline, the stop carries on in the background while the drain returns `deadline-expired`. Past the deadline, turns are interrupted: each rolls back, or, if its `COMMIT` was already sent, its receipt answers the caller's retry. Job attempts are interrupted too; each keeps its claim and `ambiguous = true` and is taken over after its lease. `deadline-expired` reports the counts and logs a warning. Live workflow runs are not turns and are not waited for: they end unrecorded when the layer closes and replay on another runner. There is no default deadline. A second `drain` waits for the first and returns its report.
* `Actors.serve` answers readiness at `GET {basePath}/ready` without credentials: `200 { ready: true }`, or `503 { ready: false, reason }`, with `cache-control: no-store` ([ADR 0053](../decisions/0053-served-readiness-route.md)). On Postgres the storage check behind it runs at most once a second, whatever the probe rate; on PGlite readiness does not probe the database, because its one connection is held by any in-flight turn. The OpenAPI document lists it as `durable.ready`.
* The runner keeps its shard locks until its layer closes. Close it once `drain` returns, and a graceful exit releases them at once.

### Cron time zones and intervals

From [ADR 0042](../decisions/0042-cron-time-zones-intervals-and-daylight-saving.md), implemented (M2.5). The four entries below get the timer keys `$cron:UTC 0 8 * * *`, `$cron:America/New_York 0 8 * * 1,2,3,4,5`, `$cron:Europe/London 0 8 * * 1,2,3,4,5` (the same expression in another zone is another entry), and `$cron:@every 5400000ms`.

```ts theme={"theme":"css-variables"}
import { Actor } from "@rikalabs/akter"

const Digest = Actor.command("Digest")
const OpenDesk = Actor.command("OpenDesk")
const OpenLondonDesk = Actor.command("OpenLondonDesk")
const Reconcile = Actor.command("Reconcile")

export const Desk = Actor.make("Desk", {
  key: Actor.singleton,
  internal: { Digest, OpenDesk, OpenLondonDesk, Reconcile },
  schedules: {
    "0 8 * * *": Digest,
    "CRON_TZ=America/New_York 0 8 * * 1-5": OpenDesk,
    "CRON_TZ=Europe/London 0 8 * * 1-5": OpenLondonDesk,
    "@every 90 minutes": Reconcile,
  },
})
```

* `CRON_TZ=<zone>` takes an IANA zone name the runtime knows; the key keeps the name as declared, so aliases such as `US/Eastern` and `America/New_York` are separate entries. Without the prefix the zone is `UTC`. `Actor.make` rejects an unknown zone or a fixed offset.
* `@every <duration>` takes a `Duration.Input` string of whole milliseconds, at least 1 second, and fires at every multiple of it since the Unix epoch, so `@every 1 day` fires at 00:00 UTC. `Actor.make` rejects a zone prefix on an interval.
* Daylight saving: a wall-clock time a spring-forward gap skips fires once at the first instant after the gap (`CRON_TZ=America/New_York 30 2 * * *` fires at 03:00 EDT on 2027-03-14), and a wall-clock time a fall-back transition repeats fires once, at its first occurrence (`30 1 * * *` in that zone fires at 01:30 EDT on 2026-11-01, not again at 01:30 EST).
* After downtime a pending zoned or interval tick fires once inside `maxScheduleLag` and is skipped outside it, and the next tick is the first scheduled time after now.

### Job cancellation and caps

Implemented in M2.13 (migration `0015_effect_control`) with the accepted defaults of [ADR 0024](../decisions/0024-effect-cancellation-and-per-actor-concurrency.md). `perActor` is validated at `Actor.make`, and `executors.cancelCheck` at `Actors.layer` (at least 1 second; above `lease / 3` it is lowered to `lease / 3`).

```ts theme={"theme":"css-variables"}
import { Actor } from "@rikalabs/akter"
import { Actors } from "@rikalabs/akter/runtime"
import { Effect, Schema } from "effect"

const SendReminder = Actor.job("SendReminder", { payload: { userId: Schema.String } })
const CapturePayment = Actor.job("CapturePayment", { payload: { amount: Schema.Int } })
const Remind = Actor.command("Remind", { payload: Schema.String })
const Forget = Actor.command("Forget")
const ReminderSent = Actor.command("ReminderSent")
const ReminderCancelled = Actor.command("ReminderCancelled", {
  payload: Actor.Cancelled(SendReminder),
})
const Captured = Actor.command("Captured")
const CaptureCancelled = Actor.command("CaptureCancelled", {
  payload: Actor.Cancelled(CapturePayment),
})

export const Billing = Actor.make("Billing", {
  key: Schema.String,
  api: { Remind, Forget },
  internal: { ReminderSent, ReminderCancelled, Captured, CaptureCancelled },
  jobs: {
    SendReminder: { job: SendReminder, onSuccess: ReminderSent, onCancelled: ReminderCancelled },
    CapturePayment: {
      job: CapturePayment,
      concurrency: { perActor: 1 },
      onSuccess: Captured,
      onCancelled: CaptureCancelled,
    },
  },
})

export const BillingLive = Billing.toLayer({
  Remind: Effect.fn(function* (userId) {
    const turn = yield* Billing.Turn
    yield* turn.enqueue(SendReminder.make({ userId }), { key: "reminder", after: "24 hours" })
  }),
  Forget: Effect.fn(function* () {
    yield* (yield* Billing.Turn).cancelJob("reminder")
  }),
  ReminderSent: () => Effect.void,
  ReminderCancelled: () => Effect.void,
  Captured: () => Effect.void,
  CaptureCancelled: () => Effect.void,
})

export const runtime = Actors.layer({
  executors: { concurrency: 64, lease: "60 seconds", cancelCheck: "20 seconds" },
})
```

* `turn.enqueue(job, { key?, after?, at? })` names the job within its actor and may delay its first attempt; enqueueing again with a live key cancels the earlier job, then records the new one under a new job id. Keys are stored with the reserved prefix `$job:`, which `Intent.key` and `Intent.cancel` reject.
* `turn.cancelJob(key)` applies in the commit statement. A never-attempted job is deleted. An attempted one is never attempted again and is reported once through `onCancelled` with `Actor.Cancelled(J)` payload `{ jobId, job, attempts, outcome, ambiguous }`, where `outcome` is `Succeeded { value }`, `Failed { cause }`, or `Unknown { cause }` and `ambiguous` is `true` exactly for `Unknown`. A running attempt is interrupted at its next lease renewal (every `executors.cancelCheck`, default `lease / 3`) or at once on the same runner. Without `onCancelled`, `Succeeded` routes to `onSuccess`, `Unknown` is dead-lettered as ambiguous, and `Failed` is deleted with the info log `Cancelled job dropped after a failed attempt`. A completed job is not affected.
* A binding's `concurrency: { perActor: n }` (1–64) bounds that job's running attempts for one actor across all runners; extra due rows wait in order and never delay other actors. There is no deployment-wide limit per job type.

### Event and job payload evolution

Implemented in M4.7 (migration `0021_payload_versions`) with the accepted defaults of [ADR 0032](../decisions/0032-event-and-effect-payload-evolution.md).

```ts theme={"theme":"css-variables"}
import { Actor } from "@rikalabs/akter"
import { Schema } from "effect"

const OrderId = Schema.String
const Money = Schema.Struct({ amount: Schema.Number, currency: Schema.String })
const Receipt = Schema.Struct({ id: Schema.String })

const V1 = { orderId: OrderId, amount: Schema.Number }
const V2 = { orderId: OrderId, total: Money }

const toV2 = Actor.migration(V1, V2, (v1) => ({
  orderId: v1.orderId,
  total: { amount: v1.amount, currency: "USD" },
}))

export const OrderPlaced = Actor.event("OrderPlaced", V2, { migrations: [toV2] })

export const ChargeCard = Actor.job("ChargeCard", {
  payload: V2,
  success: Receipt,
  migrations: [toV2],
})
```

* `Actor.event(tag, fields, { migrations?, writeVersion? })` and `Actor.job(tag, { payload?, success?, progress?, migrations?, writeVersion? })` take a chain of `Actor.migration(from, to, upcast, { downcast? })` steps. Each `to` has the next `from`'s fields and the last `to` has the class's fields, compared key by key; `_tag` is never part of a step. An invalid chain throws when the class is built.
* A value's payload version is the number of steps applied to reach its shape; it is stored beside the value in `payload_version`. Every read decodes through the chain from the stored version: `read.events`, `read.follow`, connection contexts, served feeds, subscription deliveries, workflow waits, job attempts, and the `onDeadLetter` and `onCancelled` routes. A stored version newer than the chain, or an upcast that throws or produces an invalid value, is a defect for that read, never a skip. A subscription delivery that cannot decode backs its row off at that event. A job attempt whose payload does not decode never calls the executor, so the row keeps the ambiguity earlier attempts left.
* `migrations: { from: k, steps }` declares a chain whose first `k` steps were dropped; later version numbers do not move.
* `writeVersion: n − 1` keeps writing the previous version during the first phase of a rolling deploy that adds step `n`; every step above it needs a `downcast`.
* `Actors.layer` refuses to start when a recorded version is above the chain's current version (a rollback past a schema change), when a chain starts above an event version not yet marked cleared, when a chain starts above the version of a pending job or dead letter, and when an actor type drops an event class a subscription still has undelivered events of. Each runtime records the version it writes before it takes a shard and heartbeats a writer row per written version every half `payloadWriterWindow` (default 2 minutes, at least 1 second); a runtime that has not refreshed within the window refuses new turns with `ActorUnavailable` until it does.
* `akter payloads check --entry <module> --database-url <url> [--json]` runs the startup check read-only and exits 1 when a deployment would be refused. `akter payloads clear` marks a superseded event version cleared once its `superseded_at_ms + keepEvents` has passed, no writer row for it was refreshed within its window plus the longest `executionTimeout`, and no event of that version remains; it exits 1 while a version stays uncleared. `checkPayloads` and `clearPayloads` in `@rikalabs/akter/runtime` are the same operations.

### Existing-schema adoption

`Actor.table(existing, { owner: { tenant, actor }, access? })` adopts a table other code writes ([Drizzle integration](04-drizzle.md#adopting-an-existing-table)). The CLI, with `--entry <module>` for the actor definitions and `--database-url <url>` for a login that owns or can alter the table:

* `akter adopt plan [--table <name>] [--json]` reads the catalog and prints what adoption would meet and the SQL each step would run; it changes nothing and exits 1 while a table has a problem.
* `akter adopt observe <table>` adds `routing_key bigint`, records the adoption in `actor_adoptions`, and installs the statement triggers that record each write in `actor_adoption_writes`. `akter adopt observe <table> --report [--since 7d] [--clear] [--json]` prints the recorded writers by login, `application_name`, and operation.
* `akter adopt backfill <table> [--batch 1000]` fills `routing_key` in resumable batches; it refuses rows whose mapped columns are `NULL` or empty.
* `akter adopt enforce <table> --writer-role <role> [--allow <role>]... [--quiet 7d]` refuses, listing every reason, unless the table is ready, then revokes write privileges from every role but the writer and the allowed ones and installs the guard trigger ([migrations guide](../operations/02-migrations.md#adopting-an-existing-schema)). `akter adopt release <table> --to observe` undoes it.
* `akter adopt status --database-url <url> [--json]` lists each adopted table with its mode, unbackfilled rows, last legacy write, and writer and allowed roles.

`Actors.layer({ adoption: { role } })` is the writer role: a turn of an actor type that owns an enforced adopted table runs as `role` with its tenant set, and the runtime refuses to start an enforced table without it. With `rowLevelSecurity` on, both must name the same role.

`planAdoption`, `observeAdoption`, `adoptionWriters`, `backfillAdoption`, `enforceAdoption`, `releaseAdoption`, and `adoptionStatus` in `@rikalabs/akter/runtime` are the same operations. A refused command exits 1 with the reason (`AdoptionRefused`).

## Content

Built by M4.13 ([ADR 0034](../decisions/0034-tenant-scoped-content-addressed-blobs.md)). Shared content is immutable bytes stored once per tenant; actors hold named references to it.

* `Actor.content(name)` declares a content blob and sits in `blobs` beside `Actor.blob`; names share one namespace per actor type.
* `Content.upload(bytes)` stores up to 64 MiB for the ambient tenant, outside turns and queries, and returns `ContentRef { hash, size, grant }`; past 64 MiB it fails `ContentTooLarge` and stores nothing. The server computes the SHA-256; uploading bytes the tenant already holds stores nothing new and returns a fresh grant.
* `Content.grant(X, id, C, name)` returns a fresh grant for the entry `name` that actor `X/id` references under content blob `C`, or `Option.none()` when it holds none. The actor's `access` and the global `authorize` see the ambient caller with kind `content` and command `<blob>.grant`. It names the blob because an actor type may declare several; the served route carries it too.
* A grant is `g1.<key id>.<expires ms>.<mac>`, an HMAC-SHA-256 over the deployment, tenant, hash, size, and expiry, and lasts one hour. `ContentRef` is a schema, so a command takes it as input; `attach` fails with the typed `InvalidContentRef` (`malformed`, `invalid`, or `expired`).
* `Actors.layer({ content: { keys, grace, skew } })` configures grants: `keys` lists `{ id, secret }` with a secret of at least 32 bytes; the first signs and every listed key verifies. `grace` (default 24 hours) and `skew` (default 60 seconds) set the sweep horizon and the attach margin. A runtime without `content` refuses to register an actor type that declares content. `ActorTest.layer` uses a fixed test key unless given `content`.
* `Actors.serve`'s `limits.contentBytes` (default and maximum 64 MiB) bounds `POST /content`.
* `ContentStore` is the runtime service these need. Job executors never receive it: a job layer that requires it does not type-check, and an attempt runs without it.

## Accepted M4 targets

Designs from the M4 ADRs, accepted 2026-09-28. Hosted ingress and embedded PGlite are built (see `Auth.assertion` below and `Database.pglite` above); the cold tier is not.

* **Hosted ingress (M4.8, [ADR 0031](../decisions/0031-hosted-ingress-tenant-directory-and-regions.md)).** `Auth.assertion({ issuer, audience, region, keys })`, where `keys` is a key-set URL or static keys, is the only provider a hosted runner uses. Hosted deployments configure hosted API keys and JWT settings whose tenant is `{ claim: "org_id" }` or `{ fixed: "default" }`; custom `Auth.make` code does not run at the edge. `akter tenants create <tenant> --region <region>` accepts only the primary region until L.1 adds `move`.
* **Embedded PGlite (M4.14, [ADR 0035](../decisions/0035-pglite-embedded-production-backend.md)).** `Database.pglite({ dataDir })` takes an exclusive `flock` on `<dataDir>/.akter.lock` and fails with `DataDirLocked` while another process holds it, or `DataDirVersion` for a `dataDir` from another Postgres major; it refuses `relaxedDurability` with a `dataDir`.
* **Cold tier (L.2, [ADR 0036](../decisions/0036-cold-tier.md)).** `policy.coldAfter` (default 30 days, or `"never"`) and `Actors.layer({ coldStorage })` with an S3-compatible adapter; nothing goes cold without `coldStorage`.

## Layers

```ts theme={"theme":"css-variables"}
import { Effect } from "effect"
import { CountChanged, Counter } from "./counter.ts"

export const CounterLive = Counter.toLayer({
  Reset: Effect.fn(function* () {
    const turn = yield* Counter.Turn
    yield* turn.state.set({ count: 0 })
    yield* turn.emit(CountChanged.make({ count: 0 }))
  }),
})

export const CounterReads = Counter.toQueryLayer({
  GetCount: Effect.fn(function* () {
    const read = yield* Counter.Read
    return read.state.count
  }),
})
```

* `X.toLayer(handlers)` implements commands (public and internal), streams, connections, and workflows. Reducers have no entry. Given an Effect that builds the handlers, the builder runs once per layer build, and its typed failure is the layer's error; for a singleton it runs once per activation on its owner, in the activation's scope, so a singleton's failing builder fails that activation's commands rather than startup, and its layer has no error type: it replaces wake hooks; `Effect.addFinalizer` replaces sleep hooks; `Effect.forkScoped` replaces `run` on singletons; a `Ref` replaces `vars`. There is no defect hook: deterministic defects are recorded in the turn span ([ADR 0012](../decisions/0012-workflows-internals-effects-defects-merging-regions.md)).
* `X.toQueryLayer(handlers)` implements queries against committed data.
* `X.toJobLayer(executors)` implements job executors and may run on separate processes. An executor returns the value routed to its binding's `onSuccess` command (or `void`); the framework delivers it through the outbox with the job id as the command id, and delivers `onDeadLetter` with an `Actor.DeadLetter(J)` payload when retries are exhausted.

Each takes the handler map directly, or an Effect that builds it; there is no options object. Each handler's payload type comes from its member, so `Effect.fn(function* (payload) { ... })` needs no annotation, and each handler's requirements are inferred into the layer's. Handlers take only their payload; the exceptions are a connection's `{ open, frame, close?, resync? }` entry and a stream handler, which returns a `Stream` ([ADR 0023](../decisions/0023-connections-parking-and-streams.md)). Context is a typed service per phase (`X.Turn`, `X.Read`, `X.Connection`, `X.Workflow`, `X.Executor`); see [context](02-context.md). A small actor may keep its contract and layers in one module; split `<actor>/contract.ts`, `layer.ts`, `queries.ts`, and `jobs.ts`, with a `workflows/` folder, when responsibilities justify it. Workflow bodies use Effect's `Activity` and `DurableClock` on the framework's own `WorkflowEngine`, which stores steps on the owner's shard ([ADR 0012](../decisions/0012-workflows-internals-effects-defects-merging-regions.md)).

## Calling actors

```ts theme={"theme":"css-variables"}
import { Intent } from "@rikalabs/akter"
import { Effect } from "effect"
import { Counter, CounterId } from "./counter.ts"

const id = CounterId.make("c1")

const outside = Effect.gen(function* () {
  const counter = yield* Counter.get(id)
  const state = yield* counter.Increment(5) // request/reply
})

const insideATurn = Effect.gen(function* () {
  const later = yield* Counter.intents(id)
  yield* later.Reset() // commits with the turn, delivered after
  yield* later.Reset().pipe(Intent.after("1 hour"), Intent.key("idle"))
  yield* Intent.cancel("idle")
})
```

Outside a turn, every call is request/reply and direct: the command runs in its owner's turn, and the committed receipt is its only durable admission record ([ADR 0011](../decisions/0011-direct-commands-outbox-and-performance.md)). The handle retries retryable failures with the same command id. Work that must survive a caller crash is an intent written by a turn, or a workflow.

Inside a turn, `X.intents(id)` returns the same method shape as durable intents. They are committed with the turn, delivered after commit, and deduplicated by the receiver's receipt. Self-intents use `X.intents(turn.id)`. Calling `X.get` inside a turn is a type error, and a captured handle dies at runtime. Workflows start as `later.Ship(payload)` inside a turn and `counter.Ship(payload)` outside; outside calls return a `WorkflowRun`.

Workflow members, runs, and steps follow [ADR 0022](../decisions/0022-workflow-engine-storage-and-version-markers.md):

```ts theme={"theme":"css-variables"}
import { Actor } from "@rikalabs/akter"
import { Effect, Schema } from "effect"

const OrderId = Schema.String
const Address = Schema.String
const Label = Schema.String
const Reservation = Schema.String
const Quote = Schema.Number
const Order = Schema.Struct({ id: Schema.String })
class ShippingFailed extends Schema.TaggedError<ShippingFailed>()("ShippingFailed", {}) {}
class OutOfStock extends Schema.TaggedError<OutOfStock>()("OutOfStock", {}) {}
const Paid = Actor.event("Paid", { orderId: Schema.String })
declare const inventory: { readonly reserve: (order: typeof Order.Type) => Effect.Effect<string> }
const address = "1 Main St"

export const Ship = Actor.workflow("Ship", {
  payload: { orderId: OrderId, address: Address },
  success: Label,
  error: Schema.Union([ShippingFailed, OutOfStock]),
  key: ({ orderId }) => orderId, // optional; defaults to the start's command id
  versions: { "fraud-check": { current: 2, min: 1 } }, // fraud-check 1 still runs "fraud"
})
// every step is a typed, module-level constructor with an explicit, static name
export const Reserve = Ship.step("reserve", {
  payload: Order,
  success: Reservation,
  error: OutOfStock,
})
export const CoolOff = Ship.sleep("cool-off")
export const AwaitPaid = Ship.wait("paid", Paid)
export const FirstQuote = Ship.race("first-quote", { success: Quote })

export const PlaceOrder = Actor.command("PlaceOrder", {
  payload: { orderId: OrderId, address: Address },
  success: Schema.String,
})

export const Orders = Actor.make("Orders", {
  key: OrderId,
  events: [Paid],
  api: { Ship, PlaceOrder },
})

export const OrdersLive = Orders.toLayer({
  Ship: Effect.fn(function* ({ orderId }) {
    const order = { id: orderId }
    const reservation = yield* Reserve.run(order, (o) => inventory.reserve(o))
    yield* CoolOff("1 hour")
    const paid = yield* AwaitPaid({ where: (e) => e.orderId === order.id, timeout: "1 day" }) // Option<Paid>
    return reservation
  }),
  PlaceOrder: Effect.fn(function* (payload) {
    const turn = yield* Orders.Turn
    const later = yield* Orders.intents(turn.id)
    return yield* later.Ship(payload) // in a turn: the execution id
  }),
})

const outside = Effect.gen(function* () {
  const order = yield* Orders.get("o1")
  const run = yield* order.Ship({ orderId: "o1", address }) // WorkflowRun<Label, ShippingFailed | OutOfStock>
  yield* run.poll // Option<Workflow.Result>
  yield* run.result // waits for completion
  yield* run.interrupt // idempotent, receipted
  const same = yield* Orders.run(Ship, run.executionId) // reattach from a stored id
})
```

A workflow intent returns the execution id (`Effect<string, never, Actor.InTurn>`) and mints its id when staged, unlike other intents, which return `void`. The framework drives Effect's `WorkflowEngine` directly with these ids, so applications don't call Effect's `Workflow.execute` or `Workflow.poll`, and a child `Workflow.execute` inside a body is unsupported. `poll` and `result` are admitted like queries; `interrupt` is a receipted public command authorized like the workflow member. `policy.keepWorkflows` (default `"7 days"`) keeps finished results for `poll`.

Workflow bodies use only these constructors: `Ship.step` (an activity; actor calls happen only inside its `execute`), `Ship.sleep` (a durable clock), `Ship.wait` (an owner-event wait, `Option.none()` after its timeout), and `Ship.race` (the first of several effects). They compile to Effect's `Activity`, `DurableClock`, and `DurableDeferred`, but using those primitives directly in a body, or a constructor created inside a body, dies with `Unregistered workflow step`. Two constructors with one name on a member throw. Step names are static: there are no dynamic or keyed names, and calling a step twice in one execution returns its first recorded result. The registered constructors and `versions` form the workflow's manifest, which `akter workflows check` and startup compare with open executions: a removed or renamed step, or a changed result schema for a step an open execution has recorded, refuses the deploy until those executions finish.

Caller and tenant are ambient: the edge sets them per request, `ActorTest.layer` per test, and `Actor.as(caller)` and `Actor.tenant(tenant)` around an Effect. A route the application serves itself, outside `Actors.serve`, is in-process code and calls actors as the trusted `System({ source: "process" })`: it must authenticate the request and acquire handles inside `Actor.as` and `Actor.tenant` built from the verified credential ([ADR 0059](../decisions/0059-caller-and-tenant-defaults.md)). `get` takes no options. `Actor.commandId(id)` supplies an explicit command id. Acquiring a handle writes nothing; the first turn creates durable rows.

## Composition

Applications run actors in three forms:

* **embedded:** provide the layers and `Actors.layer`, and call handles as Effects;
* **served:** add `Actors.serve` for HTTP, WebSocket, SSE, and OpenAPI;
* **hosted:** run the same layers on managed runners with Neki.

`Actors.serve` requires authentication; `Auth.none` is the explicit public opt-out. Authentication sets `CurrentCaller` at the edge. Authorization is the actor's: `Actor.make(name, { access })` takes `({ caller, ref, command, kind, of }) => boolean | Effect<boolean>` with the `kind` values of `authorize`, and when both exist both must allow. With neither, `System` callers are allowed and every `User` and `Anonymous` caller is denied with `Unauthorized` `access_denied`, so a served actor answers `403` until it declares a policy. `Actor.access.public` is a ready-made policy that allows every caller and kind:

```ts theme={"theme":"css-variables"}
import { Actor } from "@rikalabs/akter"
import { Schema } from "effect"

const Increment = Actor.command("Increment", { payload: Schema.Int, success: Schema.Int })

export const Counter = Actor.make("Counter", {
  key: Schema.String,
  api: { Increment },
  access: Actor.access.public,
})
```

> **Warning:** `Actor.access.public` opens the actor to anyone who can reach the server, `Anonymous` callers under `Auth.none` included. Use it only for demos and deliberately public actors; write an `access` policy for everything else. The API uses runtime schemas at every transport and persistence boundary and preserves the command identity across retries.

Target API from [ADR 0027](../decisions/0027-served-protocol.md): `Actors.serve({ actors, auth, basePath?, openapi?, origins?, limits? })` takes actor definitions (their layers are provided as usual) and returns a layer of `HttpRouter` routes. `auth` is `Auth.jwt({ issuer, audience, jwks, tenant, algorithms?, subject?, clockTolerance? })`, `Auth.make(authenticate)`, or `Auth.none`; a provider returns a `User` or `Anonymous` caller, the tenant, and the credential's expiry, never a `System` caller. `Actor.make` gains `feeds`, the events an actor type serves as SSE feeds. `openapi: { path }` serves an OpenAPI 3.1 document derived from the same `HttpApi` as the routes.

Committed HTTP command outcomes also expose `durable-replayed`: `false` for a newly committed command and `true` for a currently authorized receipt replay, including declared errors. The runtime supplies the marker from its receipt path, so concurrent retries do not require a preflight inspector lookup. CORS exposes this header with the other durable response metadata; the command's JSON result is unchanged.

**Implemented (M3.2, HTTP only):** `Actors.serve({ actors, auth, basePath?, openapi?, mcp?, origins?, limits? })` is a layer that adds routes to the ambient `HttpRouter`; serve it with `HttpRouter.serve(layer)` and a platform server layer such as `BunHttpServer.layer`, and provide the runtime (`Actors.layer`) and each actor's `X.toLayer`/`X.toQueryLayer`. Every served actor must be registered in that runtime, or the layer fails to build. Routes: `GET /protocol` (unauthenticated: `{ protocol, retryWindowMs, now }`), `POST /command-ids` (authenticated like a command; `{ commandId }` minted from the database clock, nothing written), `POST /actors/{Actor}/{id}/{Member}` for every public command, reducer, and query (singletons omit `{id}`; a minted actor's `{id}` is a UUIDv7 from `X.create()` or a lowercase UUIDv8 from `turn.mint`, as [ADR 0025](../decisions/0025-turn-mint.md) amends ADR 0027, and a request to a turn-minted child's `createdBy` command carries no mint proof, so it fails `Unauthorized`), `OPTIONS` preflight, and `openapi.path` when set. Commands require `Idempotency-Key` (bare or quoted) and answer `x-request-id` with it; while the runner is overloaded a command route answers `503 ActorUnavailable` with `retry-after` before authenticating or reading the request (see `admission` above), and every committed or replayed command answers `durable-version`; queries take neither, and read a configured replica only once it has replayed their `durable-min-version` (a malformed one is `400 InvalidInput { code: "decode" }`). Bodies are UTF-8 JSON (`application/json`; any other content type is `415`, invalid UTF-8 is `400 InvalidInput { code: "decode" }`), at most `limits.requestBytes` (default 1 MiB); `authorization` and `cookie` together are at most `limits.credentialBytes` (default 8 KiB). `origins` lists browser origins allowed besides the server's own (same origin means the request URL's scheme and host; forwarding headers such as `x-forwarded-proto` are not trusted, so behind a TLS-terminating proxy list the public origin); allowed origins get CORS headers exposing `x-request-id`, `retry-after`, `durable-now`, and `durable-version`. `Auth.make(authenticate)` wraps a custom provider that reads `authorization: Bearer`; `Auth.make({ authenticate, cookies: { name }, bearer? })` wraps one whose credential is the cookie `name` (an HTTP token, or `Auth.make` throws). Only a provider that names a cookie receives the request's cookies, and `bearer: true` says it also reads `authorization`. The OpenAPI document lists the provider's credentials as alternative security requirements on every authenticated operation: `bearer` (`http`, `bearerFormat: JWT` for `Auth.jwt`) and `cookie` (`apiKey` `in: cookie` with the cookie's name), and `Actors.serve` fails at startup on a hand-built provider with two credentials of one scheme ([ADR 0045](../decisions/0045-cookie-security-schemes.md)); `Auth.jwt` verifies RS256/384/512, PS256/384/512, ES256/384, or EdDSA tokens from static keys or a JWKS URL (which needs an `HttpClient` layer), requires `exp`, and checks `iss`, `aud`, and `nbf` within `clockTolerance` (default 30 s). A declared error's `httpApiStatus` must be an unreserved 4xx (not 400, 401, 403, 404, 409, 410, 413, 415, or 429) and its tag cannot be `ActorError` or `Defect`; `Actor.make` throws otherwise. `mcp: { path, name?, version? }` (M6.6, [ADR 0060](../decisions/0060-generated-protocols-mcp-and-python-client.md)) serves a stateless MCP 2026-07-28 endpoint at `path`: `POST` only (`GET` and `DELETE` answer `405`), authenticated per request like every route, with `server/discover`, `tools/list`, and `tools/call`. Each public command, reducer, and query is a tool named by its operation id (`Room.Post`), plus `durable.commandIds`; a tool takes `id`, `commandId`, and `input` as its OpenAPI operation takes the path segment, `Idempotency-Key`, and body, and answers the route's own JSON as text (`isError` for an error, with the body the route would send). `mcp.path` may not take a protocol route, an `/actors` path, or `openapi.path`. Workflow starts are not served yet. A member returning a non-void output answers `200` with its JSON, `null` for `undefined`; a void output answers `204`. Trailing slashes on `basePath` are ignored, so `basePath: "/"` is the same as no base path and `"/api/"` is `"/api"`. `Actors.serve` fails at startup when a served member's operation id would equal a protocol route's (`durable.protocol`, `durable.commandIds`), or when `openapi.path` is `/protocol`, `/command-ids`, or under `/actors`. [Generating clients](05-generated-clients.md) covers generating a client from the OpenAPI document; a generated client must keep the same `Idempotency-Key` across its retries. `Actors.serve` fails at startup when the runtime's retry window is below 60 seconds, because served clients mint ids up to a second or a round trip behind the database clock. Deployments must serve `Actors.serve` behind TLS: credentials and `Idempotency-Key` travel in headers, and the framework cannot tell whether a proxy in front of it terminates TLS, so it does not reject plain HTTP itself.

**Implemented (M4.8, hosted runners):** `Auth.assertion({ issuer, audience, region, keys, refreshEvery? })` is the only provider a hosted runner uses ([ADR 0031](../decisions/0031-hosted-ingress-tenant-directory-and-regions.md)). It reads the edge's signed assertion from `durable-assertion`, or a WebSocket frame's `authorization: Bearer`, and ignores `authorization` and cookies, so the caller and tenant come only from the assertion. `issuer` is the edge's, `audience` the deployment id, and `region` the runner's own region. `keys` is the control plane's key-set URL (`{ keys: [{ kid, kty: "OKP", crv: "Ed25519", x, nbf?, exp? }] }`, which needs an `HttpClient` layer) or static keys of the same shape. A key-set URL is reread every `refreshEvery` (default 5 minutes), which bounds how long a removed key is still accepted, and on an unknown `kid` at most once a minute. A runner that can't reread a stale key set answers `503 ActorUnavailable` rather than keep trusting it. With a key-set URL, `Actors.serve` also serves `POST <basePath>/assertion-keys/refresh`, where the edge pushes an immediate reread after revoking a key (the [protocol contract](../contracts/protocol.md#hosted-assertions-adr-0031) has its format), so revocation doesn't wait for `refreshEvery`. The server then rebuilds the canonical request binding from the request it received, the body included, and refuses a mismatch with `401 invalid_credentials` before any turn. The wire format is in the [protocol contract](../contracts/protocol.md#hosted-assertions-adr-0031). The OpenAPI document lists the credential as an `apiKey` scheme `in: header` named `durable-assertion`. `@rikalabs/akter` also exports `requestDigest`, `reauthenticationDigest`, `canonicalRequest`, `AssertionClaims`, and `AssertionKeySet`, which the edge signs with.

**Implemented (M3.3, WebSocket):** each connection member is served at `GET /actors/{Actor}/{id}/{Connection}` as a WebSocket upgrade that must offer the subprotocol `akter.v1` first (the server selects it); any other `GET` there is `400 InvalidInput { code: "unsupported_protocol" }`. The origin check, `401` for a bad upgrade credential, and `503 RunnerAtCapacity` once 1,000 sockets on the runtime await `hello` are answered before the upgrade. Every message is one JSON text frame with a `t` discriminator (`serve/frames.ts` holds both directions' schemas): the client's first message is `hello { authorization?, params }` within 10 seconds; `authorization` has the header's form, `Bearer <token>`, and a non-browser client may instead authenticate the upgrade request, in which case both must name the same caller and tenant. After the holder's open the server sends `open { connectionId, baseline?, reauthenticateBy? }`, then `frame { frame, cursor?, event? }` member frames, `resync { after?, reason, deadline }` (`deadline` in milliseconds) and `resyncReplayed` after an ungraceful owner death, `reauthenticate { by }` a minute before the credential's expiry or at half its remaining life, `reauthenticated { by? }`, and last `end { error? }` before the close code in ADR 0027's table. The client sends `frame { frame }`, `resyncDone { through? }`, and `reauthenticate { authorization }`; any other `t`, a second `hello`, or a frame its member's schema rejects ends the session with `InvalidInput` (close `4400`), a message over 64 KiB with `SessionEnded` `Defect` (close `1009`), and a binary message with `InvalidInput` (close `1003`). A credential's `expiresAt` caps the session: the holder ends it with `Unauthorized` `expired` at that time unless a renewal for the same caller arrived, and `access` and `authorize` run again with `kind: "reauthorize"` on each renewal. Under `stampCursor: false`, `open` has no `baseline` and `resync` no `after`. Commands never travel over the socket. Ping intervals and idle timeouts are the HTTP server's (`BunHttpServer.layer({ websocket: { sendPings, idleTimeout } })`). OpenAPI lists each connection as a `GET` operation `<Actor>.<Connection>` with `x-durable-transport: websocket`, `x-durable-subprotocol`, and `x-durable-frames` referencing `<Actor>.<Connection>.params`, `.server`, and `.client` under `components.schemas`.

**Implemented (M3.3, SSE event feeds):** an actor type lists the events it serves in `feeds` (each must be in `events`; singletons can't declare feeds yet), and `Actors.serve` answers `GET /actors/{Actor}/{id}/events?event={Event}&event=…&after={cursor}` with `text/event-stream`. At least one `event` is required, and every one must be in `feeds`: a missing or undeclared event is `404 InvalidInput { code: "unknown_event" }`, exactly like one that doesn't exist, and more than 16 distinct events is `400 too_many_filters`. `after` is exclusive, and `Last-Event-ID` overrides it. The server authenticates the request, answers `404 NotCreated` for an actor with no generation row without creating one, and then opens the feed at its holder, which calls `access` and `authorize` with `kind: "feed"` once per event tag (`403 access_denied` on a denial). It then answers `404` with the `UnknownCursor` body or `410` with the `RetentionGap` body before streaming. Each event is one message: `id` is its cursor, `event` its tag, and `data` is `{ event, commandId, timestamp }` with the event schema-encoded. A `: keepalive` comment line is sent every 15 seconds. Committed events after the cursor are read from `actor_events` in pages of 256 without waking the actor. Live events come from the owner, which broadcasts every committed feed event after commit to the actor's feed rows. They are deduplicated by cursor, so the race between the read and live delivery neither loses nor repeats an event. An owner loss, a broadcast gap, or a feed more than 1,024 frames behind rereads from the last cursor sent instead of ending, so the client sees one gap-free stream. Revocation (`reauthorize` with `of: "feed"` every `reauthorizeEvery`), the credential's expiry (`Unauthorized expired`), and a pruned reread (`RetentionGap`) end the feed with an `event: end` message whose `data` is the error. A client then reconnects with its last cursor. An idle feed parks with its actor like any connection, and an actor type with feeds loads its feed rows in a cold activation's first turn. At most 10,000 feeds per actor are open at once; the next is `503 RunnerAtCapacity`. OpenAPI lists the feed as `<Actor>.events` with `x-durable-transport: sse`. HTTP/1.1 browsers open at most six connections per origin, so a page that follows many feeds should be served over HTTP/2 or HTTP/3.

**Implemented (M3.3, SSE streams):** each `Actor.stream` member is served at `POST /actors/{Actor}/{id}/{Stream}` with its payload as the JSON body. After authentication and payload decoding, the response is `200 text/event-stream`. Each element is one `event: element` message whose `data` is the encoded success value, and the stream's last message is `event: end`. Its `data` is `null` when the handler's stream ended by itself; otherwise it is the `ActorError` envelope (for example `SessionEnded` `ActivationEnded` or `SlowConsumer`, or `Unauthorized` from the actor's `access` or the runtime's `authorize` with `kind: "stream"`, which runs once the subscription starts), or the member's declared error. Messages carry no `id`, because streams have no cursor; a stream that needs continuity takes its cursor in its payload and uses `read.follow`. A `: keepalive` comment is sent every 15 seconds. OpenAPI lists each stream as a `POST` operation with `x-durable-transport: sse` and `x-durable-element` referencing `<Actor>.<Stream>.element`.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.