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 anActors.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).
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) 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), Telemetry.serve({ basePath? }), GET /metrics in Prometheus text (observability), 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). 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). 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).
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).
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 definition and handler shape and ADR 0011 direct commands for commands, server reducers, queries, and keyed state (ADR 0013). Owned tables (Drizzle integration), 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) and is not implemented.
Actor.command(tag, { payload?, success?, error? })andActor.query(tag, { payload?, success?, error?, watch? })accept service-free schemas;payloadis a schema or struct fields, which stand forSchema.Struct(fields). Omitted payload/success isSchema.Voidand omittederrorisSchema.Never; JSON codecs preserve it through persistence. Declared errors must be yieldable tagged errors.Actor.table(pgTable(...))declares an owned Drizzle table: it addsrouting_key,tenant_id, andactor_idand 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 scopedone,all,count,insert,update(...).where,delete().where, andupsertin the turn transaction;read.rows(table)offers only the reads;turn.groupandread.grouprun read-only selects and inner/left joins across the placement group. Migration0005_tablesrecords each table’s owner inactor_tables, and startup fails when a table is missing, keyed without ownership, or claimed by another actor type.ActorTest.inspectreportsrowsper owned table. See Drizzle integration for the operation matrix.Actor.blob(name)declares binary storage;nameis 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)offersget(entry)(anOption<Uint8Array>of every chunk in order),set(entry, bytes),append(entry, bytes),compact(entry), anddelete(entry)(removes the entry, sogetreturns none; a missing entry is no error) on the turn transaction;read.blob(B)offers onlygetand 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, becausegetreturns it as one row and the Postgres driver closes a connection on any message over 16 MiB; asetorappendpast 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; asetcounts 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: asetorappendthat 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;deletefrees an entry’s bytes and its slot (ADR 0044, proposed). Bytes are copied when a write is called, never count towardmaxStateBytes, 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. Migration0009_blobsaddsactor_blobs.ActorTest.inspectreportsblobs, 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.placementis"tenant"(default),"actor", or{ parent: P }(ADR 0033): rows live on the shard of the actor-placed root ofP’s chain.Pmust be placed by"actor"or{ parent }(a tenant-placed parent is a type error and throws), chains are at most four levels below the root, and a parent-placed actor needs akeyorcreatedByand cannot be a singleton. Its id isc1.<byte length of parent id>.<parent id>.<local id>:X.idOf(parentId, local)builds a named child’s id, whosekeyvalidates the local part, andX.get(id)takes the full id; a minted one has noX.create()(never, and it dies at runtime) and is minted only byturn.mint(X)in a turn of its parent type, which returnsc1.<len>.<parent id>.<uuidv8>. Served routes take the full id, and a malformed id or parent part answers400 InvalidInput { code: "decode" }.apiandinternalare records whose keys must equal their command tags, checked in types and at runtime.keyis an id schema (named,X.get(id)),Actor.singleton(X.get(), registered throughSharding.registerSingletonon the single embedded runner), or omitted (minted:X.create()mints a UUIDv7 and isneverotherwise;turn.mint(X)inside a parent’s command turn mints a deterministic UUIDv8 for an actor that declarescreatedBy(every unkeyed actor accepts UUIDv8 ids, so its minted children stay reachable if the policy is later removed), ADR 0025; see context).internalcommands never appear on public handles orX.api;ActorTest.actorreaches them through a package-internal registry.policyacceptsexecutionTimeout,lockWait,deliveryTimeout,maxStateBytes,hibernateAfter,mailboxCapacity,keepReceipts,keepEvents,maxBlobBytes, andmaxBlobEntries, 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; seeadmissionbelow), 7 days, 30 days, 64 MiB, and 10,000 entries, pluskeepWorkflows(7 days),maxScheduleLag(1 day),holdEventsForSubscribers(7 days),reauthorizeEvery(60 s, from 1 s to 1 h),connections("park"),watch, andallowedSubscriberTypes, described below.createdBy,schedules, andjobssit 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, andholdEventsForSubscribers), which may reach about ten years; acreatedBycommand from another actor is rejected in types and atActor.make.executionTimeoutalso bounds each query: past it the query failsActorErrorTimeoutand its running read is cancelled on the server. Cancellation needs an unmultiplexed Postgres pool (the default); withmultiplex: true, or on PGlite, the query still fails at the deadline but its statement runs on.- Retention (retention). Every runtime sweeps once a minute. A receipt is deleted once
keepReceiptshas passed since its command id was issued (a timer’s receipt counts from its due time), never sooner thandeliveryTimeoutafter 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 oncekeepEventshas passed since they were emitted, oldest first per actor, andevent_sequenceis never reset, so replay after a pruned cursor failsRetentionGap. 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).ActorTestruntimes never sweep on their own;ActorTest.cleanupruns one sweep now, andcleanupfrom@rikalabs/akter/testingdoes the same in a production-layer runtime. Migration0010_retentionadds the indexes cleanup reads. X.toLayer(handlers)takes one handler perapiandinternalcommand, 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 withyield* X.Turn, which suppliesid,ref,caller,principal(anOption),commandId, and schema-decoded state withstate.set(patch).X.Turnis 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 and ADR 0011: 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 arebyteain codec version 1, one plain zstd frame with no dictionary (contract 06). A warm activation reuses the committed state it last read or wrote and readsactor_stateonly after acquiring a new generation. Migration0003_routing_staterebuilds the foundation tables with composite foreign keys toactor_generationsand refuses to run on a database that already holds actors. The first registration of an actor type records its placement and encoding inactor_placements; a later deployment with a different placement fails at startup instead of forking actors under a second routing key. - 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 attemptActorUnavailableafter its backoff, and the handle retries with the same command id untildeliveryTimeout.deliveryTimeoutonly stops the caller waiting.createdByfails non-creating commandsNotCreatedwithout a receipt until the creating command commits; adding the policy to actors with pre-existing rows requires an application migration. stateisActor.state(fields, { migrations? }); missing keys decode from field defaults, andsetand$versionare reserved keys.migrations: [Actor.migration(From, To, upcast), ...]upcasts older stored state.Actor.makerejects a chain unless eachtois the nextfromand the lasttoisfields. Stored state carries a$versionrow (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.statemust be the actor’s ownActor.statevalue, checked structurally in types (the reducer’s fields must equal the actor’s) and by identity atActor.make, so a separately built but equalActor.statecompiles and then throws; a reducer is anapimember only.reduce(state, payload)returns aResultof 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.toLayerhas no entry for a reducer (declaring one fails to compile); the actor still registers throughX.toLayer, withX.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), appliesreduce, validates the result against the state schema, and writes only changed keys (reducereceives 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 throwingreduceor an invalid returned state is a deterministic defect. Withbatch: { combine }the reducer repliesvoid, declares no errors, and itsreducemust not fail, all checked in types (declaring anerroralso throws). Calls of one batched reducer already waiting, consecutively, in the actor’s mailbox merge: their payloads are combined in order withcombineand 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 whenreduce(reduce(s, a), b)equalsreduce(s, combine(a, b));/testingexportscheckBatchLaw({ 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).X.toQueryLayer(build)implements every query inapi; handlers readyield* X.Read(id,ref,caller,principal, and committedstate). 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 areActorUnavailable,Unauthorized, andTimeout. Authorization is checked again before the result is returned, so a caller revoked while the handler ran getsUnauthorized; a query read pastexecutionTimeoutgetsTimeoutand is cancelled on Postgres. A query handler making a request/reply call is a deterministic defect.internalholds commands only. WithDatabase.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); in-process queries wait for the highest version any command sent through their runtime returned.- Intents and timers (contract 05).
X.intents(id)(X.intents()for a singleton) returns one method perapiandinternalcommand (reducers are not intent targets), except aninternalcommand named as a subscriptionhandler, which only subscription deliveries reach; each method returnsEffect<void, never, Actor.InTurn>;idis a plain string checked against the key schema at runtime, and the target shares the sending turn’s tenant.Actor.InTurnis provided only inside command turns and removed byX.toLayer(providing a hand-builtActor.InTurncompiles, but any intent staged under it dies), soX.intentsandIntent.canceloutside a command turn, including inX.toQueryLayerhandlers, leave an unsatisfiable requirement. A command handler that acquires a handle withX.getmakesX.toLayerfail to compile withRequest/reply inside a turn: use X.intents(id); a captured handle still dies at runtime.Intent.after(duration),Intent.at(dateTime), andIntent.key(key)pipe onto an intent (or any Effect that stages intents) the wayActor.commandIddoes; the outermost setting wins.Intent.keynames 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 withIntent capability escaped its turn. - Staged intents are written to
actor_outbox(migration0004_outbox) on the sender’srouting_keyin 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 ofrouting_key) and delivers each claimed row as a direct command with the intent id as its command id andSystem({ 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’sexpires_at_ms. Delivery is trusted internal recovery: it does not callaccessorauthorizeor 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 aftermin(2^(attempts − 1), relay.maxBackoff)seconds without limit; intents have no dead letter yet, and the operator signal is theOutbox delivery failed; retrying with backoffwarning (annotated with the reason) plusactor_outbox.attempts. Replacing or cancelling a key deletes the pending row; a timer whose delivery has already started is still delivered once (contract 05). Every runner claims rows withFOR UPDATE SKIP LOCKED(migration0011_relay), only as many as it has free delivery slots (relay.deliveryConcurrency, default 16, at mostrelay.passLimitper claim), and each claim moves the row’sdue_at_mspastmax(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 perrelay.maxBackoffand never holds back newer rows.actor_outbox.attemptscounts claims, andscheduled_at_mskeeps the time a row first became due. Each runner polls everyrelay.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).
const J = Actor.job(tag, { payload?, success?, progress? })declares a job class:payloadis a record of struct fields (the instance’s fields, as inSchema.TaggedClass),successis the schema of the executor’s return value,Schema.Voidwhen omitted, andprogressis 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 injobs: { [Tag]: { job: J, ...settings } }, keyed by each job’s tag.turn.enqueue(J.make(...))exists only onX.Turnand 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. Anenqueuecaptured and run after its turn dies withJob capability escaped its turn. A binding takesretry: { 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, andonCancelled. Routes must name a command in this actor’sapiorinternal(checked in types and atActor.make);onSuccess’s payload must acceptJ’ssuccesstype andonDeadLetter’s payload must acceptActor.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 noerror. X.toJobLayer(executors)implements one executor per bound job as(job) => Effect<Success, unknown, R>, directly or from an Effect that builds them; executors readyield* X.Executor(jobId,attempt,principal, the enqueueing actor’sref, andprogress(J, frame)for a job that declaresprogress; see context). 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 underprogress(ADR 0030). If the executors’ or the build Effect’s requirements includeSqlClient,PgClient, orPgliteClient,X.toJobLayerfails to compile withExecutors have no database capability, and at runtime an executor’s context is the layer’s build context with those clients removed andTenantset 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_outboxas rows of kindjob(migrations0008_effectsand0026_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 afterjobs[Tag].timeout(default 30 seconds). If the result does not encode under theonSuccesscommand’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 asambiguousinstead of executed again. When the executor succeeds, the same statement that records the result turns the row into an intent toonSuccesswith the encoded return value (or deletes the row when there is noonSuccess), and the relay delivers it like any intent: a direct command whose command id is the job id and whose caller isSystem({ source: "job", ref: <actor>, onBehalfOf: <enqueueing turn's principal> }). The receipt keeps delivery to one committedonSuccessturn per job id; an executor that ran twice because a crash lost its first result routes only the recorded result. A failed attemptnretries aftermin(base × 2^(n − 1), max)fromjobs[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 letterambiguous, with the warningJob 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 underjobId. When retries run out, the relay records the job inactor_dead_lettersand, in the same transaction, turns the row into an intent toonDeadLetterwithActor.DeadLetter(J)payload{ jobId, job, attempts, cause, ambiguous }, delivered once in a new turn with the job id as its command id; withoutonDeadLetter, or when the stored job no longer decodes under its class, the row is deleted and the dead letter stays for operators.ambiguousisfalseonly 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 istrue. Executors must treat a typed failure as “the provider did not apply the call” and usejobIdas the provider’s idempotency key, because a later attempt may follow an attempt whose outcome is unknown. WarningsJob attempt failed; retrying with backoffandJob dead-lettered after its last attemptare the operator signals. Job cancellation and per-actor concurrency caps are implemented by M2.13 (ADR 0024); see 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. - A deterministic defect rolls back, writes no receipt, returns
Die, logsDeterministic actor defectwith actor, id, tenant, command, and command id, and runs inside the spanakter.<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).mintCommandIdand applycall.pipe(Actor.commandId(id))before its first execution. Minting reads the database clock, so while the database is unreachablemintCommandId, and a call that mints its own ID, failActorErrorActorUnavailable, which a caller retries like any other delivery failure. Every external retry rechecks access and expiry. /runtimeexportsActors.layer({ authorize?, retryWindowMs?, maxResidentActors?, admission?, rowLevelSecurity?, payloadWriterWindow?, observability? }), whereobservability.defects(default 1,000) bounds the defect spans a runner keeps forakter defects listandobservability.sampleEvery(default 15 seconds) is how often one runner samples the database gauges (ADR 0049),Database.postgres({ url: Redacted.make(url), maxConnections?, ... }), andDatabase.pglite(config?).maxResidentActors(default 10,000) is how many activations one runner keeps in memory. A command that needs a new activation past it failsRunnerAtCapacity, and the handle retries it until an idle actor hibernates ordeliveryTimeoutpasses. Raise it only with the memory you give the process.admission: { concurrency?, wait?, requests? }sheds load a runner cannot serve promptly (ADR 0077): at mostconcurrency(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 mostwait(default 100 ms); any other command is refused at once withActorUnavailableand aretryAfter, 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.servealso refuses HTTP actor command routes before authenticating or reading them whenadmission.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 withActorUnavailable. MCP tool calls use runtime admission after envelope authentication and parsing. An actor whosepolicydeclares nomailboxCapacitystill holds at most 1,024 commands, waiting and in its current batch; the next is refused withActorUnavailablerather thanMailboxFull, which stays reserved for a declared capacity.fleetregisters fleet views the runtime maintains, andActors.serve({ fleet })serves them atGET /fleet/{View}; it needs Postgres withwal_level=logicalandakter fleet setup, and PGlite refuses it.rowLevelSecurity: { role }opts in to row-level security: command turns and queries run asrolewith their actor’s tenant indurable.tenant, and startup fails unless the database is set up as deployment describes (ADR 0051).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 tomaxConnections + offTurnConnections(60 by default); keep the total across runners below the server’smax_connections. Supply a platformCryptolayer, such asBunCrypto.layer. The optionalauthorizecallback is the global hook: it receives the capturedcaller,ref,command, andkindat admission and before returning an outcome, beside the actor’s ownaccesspolicy (ADR 0059).User.make({ subject })andSystem({ source, ref?, onBehalfOf? })are application-trusted attribution, not authentication. Code in the application’s own process, outside a turn and not throughActors.serve, runs asSystem({ source: "process" })in tenant"default";Actor.as(caller)andActor.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;gettakes 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’ssrc/database.tsprovidesDatabase.postgres({ url })whenDATABASE_URLis set and otherwiseDatabase.pglite({ dataDir }), withDATA_DIRdefaulting 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.Database.pglite({ dataDir })with a filesystemdataDir(M4.14, ADR 0035) takes an exclusive, non-blockingflockon<dataDir>/.akter.lockbefore PGlite opens and holds it until the layer’s scope closes. Layer build fails withDataDirLocked { dataDir }while another process, or another open layer in this one, holds it, and withDataDirVersion { dataDir, found, expected }when<dataDir>/PG_VERSIONnames another Postgres major than the pinned PGlite embeds (18); both are exported from/runtime.relaxedDurability: truewith adataDiris a defect. In-memory PGlite (nodataDir, ormemory://) takes no lock. Linux and macOS only; the lock is advisory and unreliable on network filesystems.Database.pglitecreates a fresh owned instance per layer build; a suppliedliveClientconfig 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 becauseSqlRunnerStoragewould reserve the sole connection; this confers no independent-connection or concurrent-ownership guarantee.Database.postgresincludes a scopedregclasscodec for the pinned rc.116 driver restart bug; it does not modify global driver configuration./runtimealso exportsRunner.socket({ address, listenAddress?, transport, shardsPerGroup?, shardLockExpiration?, shardLockRefreshInterval?, refreshAssignmentsInterval?, entityTerminationTimeout? }), provided toActors.layer.transportis a platform TCP server/client layer, such asLayer.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 reportsroutinguntil 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 inca, the certificate to be valid now, and its subject alternative names to be exactly the URIRunner.identity(deployment)(spiffe://akter/deployment/<deployment>).credentialsis an Effect of PEM{ ca, certificate, key: Redacted }, run at startup and everyrefreshEvery(default one minute) with arefreshTimeout(default 10 seconds); invalid or late credentials refuse startup, or on refresh are logged and ignored. Readiness reportspeeringonce the current certificate has expired or loads have failed forunhealthyAfter(default 5 minutes).RunnerAuthority.make(),.from({ certificate, key })andauthority.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 and ADR 0086.- 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 providesDatabase.Neki, a context reference (defaultfalse); a process that opens its own Postgres client and runs no actors providesLayer.succeed(Database.Neki, true)itself.Database.schemaChange(effect, lock)(oreffect.pipe(Database.schemaChange(lock))) runs an application’s own startup schema statements on the ambientSqlClientunder advisory locklock: 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). 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./testingexportsActorTest.layer({ database?, as?, authorize?, retryWindowMs?, maxResidentActors?, admission?, relay?, executors?, observability?, rowLevelSecurity? }), with a fresh tenant per build, committedinspect(ref)(generation,state,receipts, and pendingoutboxrows the actor sent),receiptsFor(ref, command),crashNext(point, { commandId? }),pauseNext(point, { commandId? })(withcommandId, 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), andnow.advancemoves 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;nowis that clock as aDateTime.Utc, forIntent.at.databasedefaults to a fresh in-memory PGlite; adataDirconfig persists across builds and aRedactedURL uses Postgres.test.actor(X, id?, { seed? })returns{ system, inspect }wheresystemreaches every command including internal ones with aSystemcaller that inherits the configured principal.seedis the path of a seed fileakter exportwrote, read through Effect’sFileSystem, 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, soadvancedelivers 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.checkBatchLawchecks a batched reducer’s merge law over generated states and input lists and dies with a reproducible seed on failure. Fault points arebeforeDelivery,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’safterClaim,beforeOutboxDelete,beforeExecute,afterExecute, andbeforeRenew.TurnHooksis available only through the testing entry for process-level faults; aTurnHooksservice provided aroundActorTest.layerstill 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 }:cursoris the flushed-through event cursor at open,framesis aStreamof decoded server frames andmessagesthe same with the holder’sResyncandResyncReplayedcontrol messages and each frame’s cursor stamp, both failing with the session’sActorErrorwhen it ends;resyncDoneanswers aResync. A declaredopenfailure failsconnectwith that error.test.progresslists every progress message this runner’s executor pool sent, in order, each markeddroppedwhentest.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 ashibernateAfterwould, sealing every holder so its connections stay open without a resync. The in-process transport is not exported from/runtimeyet;Actors.serveholds WebSocket sessions through the same holder./testingalso exports the framework-neutralconformancecases anddescribeConformance; the same named cases run on PGlite and Postgres, and cases needing a second connection are reported skipped on PGlite./testingalso exportsActorTest.simulateCluster({ seed, faults, faultRate?, primary?, commands?, attempt?, within?, settle? }, program)andclusterSimulationSeeds, the multi-runner counterpart ofActorTest.simulate, run inside anActorTest.clusterlayer.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 arecrashBeforeCommit,crashAfterCommit,dropReply,runnerKill,heartbeatLoss,connectionLoss, andprimaryFailover; see deterministic simulation./testingalso exportsActorTest.cluster({ database, runners, shardLockExpiration, actors, runnerActors?, holdersOnly?, as?, authorize?, retryWindowMs?, maxResidentActors?, relay?, executors? }), a layer providingActorCluster. It runsrunnersactor 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 expireshardLockExpirationafter its last heartbeat (whole seconds, rounded up).actorsis the application layer every runner builds (X.toLayer,X.toQueryLayer, …), andrunnerActors(index)adds layers only some runners build, such as a job layer one runner lacks;holdersOnlylists 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.ActorClusterofferson(runner)(effect), which runseffectwith that runner’sActorsandActorTest, so handles dispatch through it andcrashNext,pauseNext,inspect, andadvancebelong 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, orundefined(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;tenantandsql, the cluster’s tenant and a pool of its own for inspection; andready, 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 aboutshardLockExpirationfor 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.
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.
counter.ts
Members:
-
Actor.command(tag, { payload?, success?, error? })runs an effectful server handler.payloadis a schema or struct fields (standing forSchema.Struct(fields)) and gives the handle method’s argument; an omitted payload is a zero-argument call.erroris one tagged error class or aSchema.Unionof them; framework reasons never masquerade as declared failures. Listing it underinternalinstead ofapiremoves it from public handles and transports. -
Actor.reducer(tag, { state, payload?, error?, reduce, batch? })is a pure transition with no server handler;stateis the actor’s ownActor.state, andreducereturns aResultof the next state or a declared error. It runs optimistically in browser handles; withbatch: { combine }queued calls fold in their order into one turn, it returnsvoid, and declares no error.combineneeds ordered fold equivalence, not commutativity. -
Actor.queryandActor.workflowdeclare reads and durable workflows. -
Actor.query(tag, { payload?, success?, error?, watch? })withwatch: true(M6.2, ADR 0055) declares a query that can also be observed. Its handler inX.toQueryLayermay require onlyX.Read: a handler that needs another service does not compile, and at runtime a watched handler is provided no service butX.Read, so one that got around the types fails with a defect (a watch ends withSessionEndedcauseDefect). The in-process handle gainshandle.Query.watch(input), aStreamof the query’s outputs that fails with its declared errors orActorErrorreasonActorUnavailable,Unauthorized,RunnerAtCapacity,SessionEnded,NotCreated(the actor has no generation row; a watch never creates one), orTimeout. 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 withread.followfor that). The runtime records the reads per rerun:state(all keyed state),read.events(E)per class,read.rows(T)per table, andread.blob(B)per blob. A rerun that readsread.groupis a defect of the handler in process (the watch dies withInvalidInput { code: "not_watchable" }as its cause), becauseInvalidInputis boundary-only; a served watch ends with it.read.id,read.ref,read.caller, andread.principalare 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 atPOST /actors/{Actor}/{id}/{Query}/watchover SSE (see protocol), exposed ashandle.Query.watch(input, { signal? })on the Promise client, and followed byuseWatchin@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) withx-durable-transport: sseandx-durable-elementreferencing<Actor>.<Query>.element; a served watch of a query withoutwatch: trueis400 InvalidInput { code: "not_watchable" }, and a rerun that readsread.groupends it with anendmessage carrying the same error. -
policy.watchsets the limits of an actor’s watches:maxPerActor(default 1,000 open watches; the next isRunnerAtCapacity),minInterval(default 100 ms between reruns of one watch), andreconcileEvery(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 inX.toLayertakes the payload and returns aStreamofsuccessvalues that runs on the activation for one subscriber, withX.Readplusread.follow(Event, { after }); it fails only with its declarederror. The handle method returns aStreamof the same elements that fails with those errors orActorErrorreasonActorUnavailable,Unauthorized,RunnerAtCapacity(past 256 subscriptions on one actor), orSessionEnded. An open subscription keeps the activation resident whateverpolicy.connectionssays; the activation’s end (a move, a runner’s death, eviction, or shutdown) reaches the subscriber asSessionEnded { 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 withSessionEnded { cause: "SlowConsumer", resync: true }. The subscriber’s runner and then the owner callauthorizewithkind: "stream", and the owner reauthorizes everyreauthorizeEvery; a denial ends it withUnauthorized, discarding elements not yet delivered.Actors.serveserves 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:payloadis the opening params,serverthe frames the actor sends,clientthe frames the client sends,sessionthe per-connection fields (at most 16 KiB encoded), anderrorthe declared failures of opening. ItsX.toLayerentry is{ open, frame, close?, resync? }: short handlers that run withX.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, andprogress: { jobs }on a stream member, opt into executor progress of jobs the actor binds injobsthat declare aprogressschema (anything else is rejected atActor.make). A connection receivesProgress { 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, withto: "all"every open connection of the member. A stream handler readsread.progress(J, { jobId? }). Progress is coalesced (newest per job), may be lost (a gap inseq), never causesSlowConsumer, 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 messaget: "progress", never inside aframemessage. -
policy.connectionsis"park"(default: open connections do not keep the activation resident) or"keepAwake"(they count as activity; residency is still not guaranteed).policy.reauthorizeEveryis the revocation bound for connections and streams: default 60 s, from 1 s to 1 h. -
The actor’s
accesspolicy and the runtime’sauthorizehook are also called on every open and stream subscription and then everyreauthorizeEvery, with the member tag ascommandand a newkindfield ("command" | "query" | "open" | "stream" | "feed" | "watch" | "reauthorize" | "content", the last for a content grant or read withcommandset to<blob>.grantor<blob>.get). A query watch is authorized withkind: "watch"and the query tag ascommand, before anything is read, and reauthorized withkind: "reauthorize"andof: "watch"everyreauthorizeEvery; reruns use the caller stored at open and call neither. An event feed is authorized withkind: "feed"and each event tag it reads ascommand. A reauthorization carriesof("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): 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 aSchema.TaggedClasswhose identifier is its tag (E.make,new E, andinstanceofwork; the event type istypeof E.Type), and the tag is stored with each committed event, so tags must be unique within an actor’sevents. Listing it ineventslets that actor’s command turnsturn.emitit and its queries replay it withread.events(context). 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 withturn.enqueueandX.toJobLayerexecutors run after commit.Actor.DeadLetter(J)is the schema of{ jobId, job, attempts, cause, ambiguous }, the payload of anonDeadLettercommand.Actor.Cancelled(J)is the payload of anonCancelledcommand (ADR 0024; implemented in M2.13). -
Actor.subscription(tag, { delivery, retired?, handler, route? })(M3.7, ADR 0026) declares a subscription to the event classes itsdeliverynames.deliveryisActor.Delivery({ source, events }), declared once and also used as the handler’s payload schema; it records itssourceandevents.handleris aninternalcommand, 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 ofEvent(source,cursor,event,commandId,timestamp),RetentionGap(after,resumeAfter), andRejected(a dynamic subscription whosefromcursor the source refused). Withroute: (event, source) => id(orActor.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 toturn.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.allowedSubscriberTypeson a source restricts subscriber types by name, andpolicy.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 ingapsinstead),policy.allowedSubscriberTypes(checked atActor.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 inactor_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.
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. The relay, executor, and per-job timing settings are implemented (M2.4), and so areschedules and maxScheduleLag (M2.5), with the time zones and fixed intervals of ADR 0042 (below). Every default equals the M1 behaviour.
timeoutbounds one executor attempt; attemptnthat fails waitsmin(base × 2^(n − 1), max). Durations run from 1 ms to 2^31 − 1 ms andmax ≥ base;Actor.makerejects anything else.relay.subscriptionConcurrency(subscription delivery slots per runner, separate from intents; feed expansions and control registrations each get as many) andrelay.subscriptionBatch(events read per subscription claim) are implemented (M3.7, ADR 0026). 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.claimLeasedefaults to the largest registeredexecutionTimeout + lockWaitplus 5 seconds.executors.leaseis renewed every third of its length while an attempt runs and must be at least 3 seconds;Actors.layerthrows 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 keepambiguous = trueand 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, soActorTest.advancemust move past the lease to redeliver it;advancewaits for in-flight deliveries and executor attempts;actor_outbox.attemptscounts claims, including expired leases, rather than failed deliveries; a success that arrives after its job was dead-lettered marks the stored dead letterambiguous. - 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); the timer key names the zone and spells the parsed schedule canonically (sorted value lists,*for a full field, seconds only when not0), 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.makerejects an unparsable expression, two equivalent schedules (such as1-5and1,2,3,4,5), a target that is not a command of the actor or takes input, and an invalidmaxScheduleLag. A cron tick’s caller isSystem({ source: "cron", ref }), so a target may be aninternalcommand that handles, HTTP, and the Promise client cannot call. On a singleton, cron runs in the default tenant only. - Changed from M1:
Intent.keyvalues starting with$,$cron:among them, are reserved:Intent.keythrows andIntent.canceldies withIntent.key values starting with $ are reserved.
Runtime control: readiness and drain
Implemented in M4.2 (ADR 0003).Actors.layer also provides RuntimeControl:
- Readiness is true once the database answers a
SELECT 1within 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 reportsrouting; registration snapshots are checked on each probe. drain({ deadline })makes the runner unready and refuses new external commands through its handles withActorUnavailable. 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 untildeadline, 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 returnsdeadline-expired. Past the deadline, turns are interrupted: each rolls back, or, if itsCOMMITwas already sent, its receipt answers the caller’s retry. Job attempts are interrupted too; each keeps its claim andambiguous = trueand is taken over after its lease.deadline-expiredreports 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 seconddrainwaits for the first and returns its report.Actors.serveanswers readiness atGET {basePath}/readywithout credentials:200 { ready: true }, or503 { ready: false, reason }, withcache-control: no-store(ADR 0053). 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 asdurable.ready.- The runner keeps its shard locks until its layer closes. Close it once
drainreturns, and a graceful exit releases them at once.
Cron time zones and intervals
From ADR 0042, 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.
CRON_TZ=<zone>takes an IANA zone name the runtime knows; the key keeps the name as declared, so aliases such asUS/EasternandAmerica/New_Yorkare separate entries. Without the prefix the zone isUTC.Actor.makerejects an unknown zone or a fixed offset.@every <duration>takes aDuration.Inputstring of whole milliseconds, at least 1 second, and fires at every multiple of it since the Unix epoch, so@every 1 dayfires at 00:00 UTC.Actor.makerejects 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
maxScheduleLagand is skipped outside it, and the next tick is the first scheduled time after now.
Job cancellation and caps
Implemented in M2.13 (migration0015_effect_control) with the accepted defaults of ADR 0024. 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).
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:, whichIntent.keyandIntent.cancelreject.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 throughonCancelledwithActor.Cancelled(J)payload{ jobId, job, attempts, outcome, ambiguous }, whereoutcomeisSucceeded { value },Failed { cause }, orUnknown { cause }andambiguousistrueexactly forUnknown. A running attempt is interrupted at its next lease renewal (everyexecutors.cancelCheck, defaultlease / 3) or at once on the same runner. WithoutonCancelled,Succeededroutes toonSuccess,Unknownis dead-lettered as ambiguous, andFailedis deleted with the info logCancelled 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 (migration0021_payload_versions) with the accepted defaults of ADR 0032.
Actor.event(tag, fields, { migrations?, writeVersion? })andActor.job(tag, { payload?, success?, progress?, migrations?, writeVersion? })take a chain ofActor.migration(from, to, upcast, { downcast? })steps. Eachtohas the nextfrom’s fields and the lasttohas the class’s fields, compared key by key;_tagis 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 theonDeadLetterandonCancelledroutes. 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 firstksteps were dropped; later version numbers do not move.writeVersion: n − 1keeps writing the previous version during the first phase of a rolling deploy that adds stepn; every step above it needs adowncast.Actors.layerrefuses 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 halfpayloadWriterWindow(default 2 minutes, at least 1 second); a runtime that has not refreshed within the window refuses new turns withActorUnavailableuntil 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 clearmarks a superseded event version cleared once itssuperseded_at_ms + keepEventshas passed, no writer row for it was refreshed within its window plus the longestexecutionTimeout, and no event of that version remains; it exits 1 while a version stays uncleared.checkPayloadsandclearPayloadsin@rikalabs/akter/runtimeare the same operations.
Existing-schema adoption
Actor.table(existing, { owner: { tenant, actor }, access? }) adopts a table other code writes (Drizzle integration). 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>addsrouting_key bigint, records the adoption inactor_adoptions, and installs the statement triggers that record each write inactor_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]fillsrouting_keyin resumable batches; it refuses rows whose mapped columns areNULLor 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).akter adopt release <table> --to observeundoes 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). Shared content is immutable bytes stored once per tenant; actors hold named references to it.Actor.content(name)declares a content blob and sits inblobsbesideActor.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 returnsContentRef { hash, size, grant }; past 64 MiB it failsContentTooLargeand 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 entrynamethat actorX/idreferences under content blobC, orOption.none()when it holds none. The actor’saccessand the globalauthorizesee the ambient caller with kindcontentand 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.ContentRefis a schema, so a command takes it as input;attachfails with the typedInvalidContentRef(malformed,invalid, orexpired). Actors.layer({ content: { keys, grace, skew } })configures grants:keyslists{ id, secret }with a secret of at least 32 bytes; the first signs and every listed key verifies.grace(default 24 hours) andskew(default 60 seconds) set the sweep horizon and the attach margin. A runtime withoutcontentrefuses to register an actor type that declares content.ActorTest.layeruses a fixed test key unless givencontent.Actors.serve’slimits.contentBytes(default and maximum 64 MiB) boundsPOST /content.ContentStoreis 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 (seeAuth.assertion below and Database.pglite above); the cold tier is not.
- Hosted ingress (M4.8, ADR 0031).
Auth.assertion({ issuer, audience, region, keys }), wherekeysis 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" }; customAuth.makecode does not run at the edge.akter tenants create <tenant> --region <region>accepts only the primary region until L.1 addsmove. - Embedded PGlite (M4.14, ADR 0035).
Database.pglite({ dataDir })takes an exclusiveflockon<dataDir>/.akter.lockand fails withDataDirLockedwhile another process holds it, orDataDirVersionfor adataDirfrom another Postgres major; it refusesrelaxedDurabilitywith adataDir. - Cold tier (L.2, ADR 0036).
policy.coldAfter(default 30 days, or"never") andActors.layer({ coldStorage })with an S3-compatible adapter; nothing goes cold withoutcoldStorage.
Layers
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.addFinalizerreplaces sleep hooks;Effect.forkScopedreplacesrunon singletons; aRefreplacesvars. There is no defect hook: deterministic defects are recorded in the turn span (ADR 0012).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’sonSuccesscommand (orvoid); the framework delivers it through the outbox with the job id as the command id, and deliversonDeadLetterwith anActor.DeadLetter(J)payload when retries are exhausted.
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). Context is a typed service per phase (X.Turn, X.Read, X.Connection, X.Workflow, X.Executor); see context. 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).
Calling actors
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:
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). 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.servefor 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:
Warning:Target API from ADR 0027:Actor.access.publicopens the actor to anyone who can reach the server,Anonymouscallers underAuth.noneincluded. Use it only for demos and deliberately public actors; write anaccesspolicy for everything else. The API uses runtime schemas at every transport and persistence boundary and preserves the command identity across retries.
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 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); 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) 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 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). 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 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. 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.