Skip to content

feat: scale-out phase D: Optio Local relay across API instances, replica rule lifted - #670

Merged
jonwiggins merged 12 commits into
mainfrom
feat/scale-out-local
Oct 11, 2026
Merged

jonwiggins merged 12 commits into
mainfrom
feat/scale-out-local

Conversation

@jonwiggins

@jonwiggins jonwiggins commented Oct 11, 2026 •

Copy link
Copy Markdown
Owner

Phase D of docs/plans/scale-out.md: the Optio Local relay works when several API instances share one Postgres and one Redis, and the chart no longer pins api.replicas to 1.

Merge order: after phase B (re-attachable runs). This PR touches nothing B touches, but lifting the replica rule assumes runs are no longer tied to one API process.

What moved where

Was (one process) Now
daemonsByHost: which process has the daemon Lease local-host:<hostId> in leases (30 s TTL, renewed every 10 s), held per daemon connection (<instance>#<connection>). Seized on hello ("newest daemon wins") after the instance subscribes to the host's channel; publishes host-taken. The old owner checks the lease really moved, then closes its socket (4000) without marking the host offline; it also finds out at its next renew. A closing connection releases only its own lease, and markHostOffline is skipped while any connection holds one.
isHostOnline / capability flags from the socket map "This instance has the socket, or the lease is live". The hello's capabilities go on local_hosts.capabilities (new migration 1792020000_local_host_capabilities) before the lease is taken, read only while it is live. hostsWithClaudeCredentials joins hosts and live leases.
"Newest wins" closed the older socket in the same process (and left its viewers frozen) Also across instances; a replaced connection's viewers are closed with 4503 so they reattach to the new one (fixes the same-process case too).
Viewer sockets, pending attaches, grid arbiter per terminal Stay with the owner. A viewer on another instance is a RemoteViewer proxy (socket-shaped) in the owner's maps; its frames go out on optio:local:term:<id>:out addressed to it, so the attach handshake, snapshot-then-live ordering and the arbiter (fed by view reports over the bus) are unchanged. The viewer's instance (HomeViewer) relays only the frames of the owner it attached through.
ws/local-terminal-stream.ts: daemon elsewhere = "Host is offline"; input/resize/view went nowhere attachViewer returns a handle that routes input, resize, view reports and pongs locally or over optio:local:host:<id>:in.
Spawn / kill / input / dirs / limits / credentials / backfill on the wrong instance (spawn parked, kill marked exited while the PTY ran) sendToHost is async: local, or a bus request naming the holder it read, acknowledged by the owner (no ack in 5 s = offline, as before). Daemon answers (dirs-result, limits-refresh-result, credentials-result, transcript-backfill) are relayed back to the asking instance's reply channel, where the existing request maps resolve them.
Remote viewers' liveness Every 10 s each instance reports its remote viewers to their owners (owners drop unreported ones and ask unknown ones to reconnect) and checks the owner still holds the lease (viewers of a daemon that moved or whose owner died are closed with 4503).
local-sweep-worker.ts: setInterval on every instance BullMQ job scheduler local-sweep.sweep (in REPEAT_SCHEDULERS) under poller:local-sweep: stale hosts, hosts whose lease lapsed (owner died), terminals parked host-offline whose host is back, unacknowledged launches.
lastAutoAttempt (Claude token auto-refresh) Lease claude-auto-refresh:<hostId> per window, shared by every instance.
interactionWrittenAt Still a per-instance first filter, plus a conditional update so it's one write a minute across instances.
Backfill bookkeeping Per instance (a cache: each instance asks at most once per window; frames are idempotent).

The bus is services/local-bus.ts: one subscriber per instance (messageBuffer), started after listen and retried with backoff like session-bus.ts; binary envelopes (JSON header + raw payload). Channels: optio:local:host:<hostId>:in (owner only), optio:local:term:<terminalId>:out (instances with a viewer), optio:local:reply:<instanceId> (replies by request id + relayed daemon answers). Output reaches Redis only while a viewer is on another instance. Without Redis, each instance still serves its own daemons and viewers.

Chart: the api.replicas=1 / autoscaling guard is gone from _helpers.tpl; with api.autoscaling.enabled the Deployment omits replicas (the HPA owns it); rollouts stay Recreate. values.yaml, values.production.yaml (still 1 replica by default) and docs/production-eks.md say what is per pod: the per-address WebSocket cap and caches. scripts/verify-production-chart.ts now asserts 3 replicas and the HPA render instead of failing. Docs: docs/optio-local.md "Several API instances", the plan's §4 status, CLAUDE.md.

E2E matrix (apps/api/e2e/scale-out-local.e2e.test.ts, three real API servers, fake daemon/viewers)

Four runs (three standalone, one inside the full tier): 6/6 passed every time, ~16 s per run.

Case What it checks Result
(a) daemon on A, viewer on B B gets the replay + snapshot, live output in order, size frames (yours after its resize); B's input and resize reach the daemon; a kill through C reaches it and the exit status reaches the viewer on B pass ×4
(b) REST through B/C, daemon on A spawn (through B) acknowledged and running; dirs add (B) gets the daemon's dirs-result and the new list; limits refresh (C) gets its result; REST input (B) reaches the PTY; kill (B) reaches the daemon and the row stays running until the real exit pass ×4
(c) daemon reconnects to B A closes the old daemon socket with 4000; the viewer on A gets 4503, reattaches (as the clients do) and streams the new connection on B; input from A reaches B's daemon; still online everywhere, terminal never exited pass ×4
(d) A SIGKILLed while owning the host daemon reconnects to B and reports its terminal running; C's viewer (streaming through A) is closed 4503 within a heartbeat and reattaches through B; a viewer on B streams directly; nothing marked exited; survivors say live pass ×4
(e) local Job run started through B, daemon on A three runs, each spawned on the daemon, started → usage → exit 0 → run completed with its cost, whichever server's worker dispatched it (spawn through a non-owner is pinned deterministically in (b)) pass ×4
(f) online state before the hello all three say offline; after it all three say online (host list capabilities) and live (session recovery); after the daemon closes all three say offline / reconnecting, host row offline, terminal still running pass ×4

Tiers

Tier Result
pnpm turbo test pass (13/13 tasks; API 189 files). New: local-bus.test.ts (10), local-relay-cluster.test.ts (13)
pnpm --filter @optio/api test:integration pass (56 files, 457 tests). New: local-bus.int.test.ts (7, real Redis + Postgres: two bus instances route a request and its reply; lease seized → old owner closes on host-taken or at renew, host stays online; same-instance reconnect keeps the new lease)
pnpm --filter @optio/api test:e2e pass (26 files, 118 tests), incl. the existing Local files (local-terminal, local-run, local-daemon-auth, now on the shared fakes)
pnpm --filter @optio/web test:e2e pass (107 tests)
pnpm format:check, pnpm turbo typecheck pass
helm lint helm/optio --set encryption.key=test pass
helm template with api.replicas=3 (default and production values), and with api.autoscaling.enabled=true render cleanly: replicas: 3, Recreate; with the HPA the Deployment has no replicas. scripts/verify-production-chart.ts passes
gen:swift / gen:kotlin not needed: no shared type changed
Live (real cluster) not run: this phase was to stay off the local cluster and never touch the real daemon. The plan's phase-E manual pass (isolated daemon against a 3-replica release in its own namespace) covers it

Notes

  • (c) and (d): a viewer streaming a daemon connection that went away is closed with 4503 and reattaches, rather than being switched over silently. The daemon's new connection starts with no subscriptions, so a new snapshot is needed anyway, and the web, iOS and Android clients already reset and replay on 4503.
  • No new shared types. One migration (local_hosts.capabilities, jsonb).

Generated with Claude Code

Review fixes

All 15 findings from the high-effort review are fixed. Commits 3488a2b5..542a318c.

# Finding Fix Test
1 Ordering breaks in REDIS_MODE=cluster (PUBLISH lands on any master) One choice in redis-config.ts: publishOrdered / createOrderedSubscriber / subscribeOrdered / onOrderedMessage. Cluster mode uses sharded pub/sub (SPUBLISH / SSUBSCRIBE, a cluster client with shardedSubscribers); standalone keeps PUBLISH / SUBSCRIBE. Both local-bus.ts and C2's session-bus.ts use it. redis-cluster.int.test.ts "keeps publish order…": 1000 frames on one channel arrive in order through the Local bus and through the session bus on the three-master test cluster (and request/reply works there)
2 An ack timeout re-sent the spawn deliverToHost returns sent / offline (nobody listening) / unacknowledged. Unacknowledged leaves the row launching and is never re-sent; handleStarted accepts launching or pending; the stuck-launching sweep fails it otherwise. sendToHost counts only offline as not sent. local-cluster.int.test.ts: owner acks after the window → one spawn, still launching after a hello flush and the sweep, running on started; started settles a pending row; unit test of the three outcomes
3 host-taken race host-taken only triggers a re-read; a connection supersedes itself only when claimed and the lease's holder (read after the message) is not its own. Releases are always WHERE holder = <own holder id> (releaseLease, also from claimHost). e2e (g): two hellos at once on A and B, 20 iterations, fast renew; exactly one survives (the other gets 4000), every server says online. Unit: an early host-taken is re-read after the claim and the connection stays
4 onHostIn dropped frames before the claim Matches h === holder regardless of claimed; frames before the claim completes are queued on the connection and handled right after it Unit: a send naming the connection mid-claim reaches the daemon (and is acked) right after the claim, not before
5 Relayed replies out of order Per-host promise chain on the receiving instance Unit: a slow first backfill frame still commits before the done frame
6 Socket closed during claimHost daemonDisconnected waits for the claim in flight, which releases exactly what it seized, then the offline mark runs local-cluster.int.test.ts: close at 0/1/3/8/15 ms into the claim → no lease, host offline every time
7 Heartbeat reported attaching viewers HomeViewer.attached is set only after the attach publish was heard; the heartbeat skips the rest Unit: no viewers report while the attach is held; one after it goes out
8 Credential over Redis The refresh runs on the daemon's instance as an owner task (runOnOwner / onOwnerTask): ask, validate, store there; only { ok, hostId } / { ok:false, error } travels back. The owner refuses a forwarded credentials frame and never relays a credentials-result. local-auth-refresh.int.test.ts: another instance asks; the secret is stored; a PSUBSCRIBE tap on optio:local:* sees no sk-ant-oat01 in any message; a smuggled credentials send is refused. Unit: no relay of credentials-result
9 request() waited out the timeout with nobody listening ask() returns unheard at once when PUBLISH/SPUBLISH reports 0 receivers; request() returns null local-bus.test.ts: null / unheard well under a second; timeout when heard but unanswered; also in the cluster test
10 Heartbeat overlap, a query per host In-flight guard; one leaseHolders(keys) query per tick Unit: two concurrent heartbeats → one lease read covering both hosts
11 GET /api/local/hosts queried per host listHosts reads the lease with the rows (one query); capabilitiesOf(row) reads capabilities off the loaded row local-cluster.int.test.ts "reads capabilities off the rows it loaded"
12 Two sources of truth for online getHost / listHosts / pickOnlineHost derive state from the lease (in a cluster); the column stays last-known for alerts, and the sweep still sends offline notifications local-cluster.int.test.ts "online is the lease"
13 Interaction throttle set from "now" Conditional UPDATE returning the stored stamp (or a read of it when nothing was written); the throttle starts from that local-cluster.int.test.ts: a stamp written 59.7 s ago by another instance → no write now, a write 0.6 s later (a throttle from "now" would hold it a minute)
14 Duplicated request lifecycles / types / grid checks HostRequests (local-host-requests.ts) used by credentials, dirs and limits; one exported DaemonReply; one saneGrid (local-grid.ts) for the daemon's replay grid and bus grids existing dirs / limits / auth-refresh tests on the helper
15 Unused capability helpers hostHasClaudeCredentials, hostCanRefreshLimits, hostCanManageDirs removed; tests use hostCapabilities —

Results after the fixes

Tier Result
pnpm turbo test pass (13/13; API 189 files)
pnpm --filter @optio/api test:integration pass (57 files, 465 tests), incl. the cluster-mode ordering test
pnpm --filter @optio/api test:e2e pass (26 files, 119 tests)
scale-out-local.e2e.test.ts 7/7 passed in the full tier and in three more standalone runs; (g) took about 38 s each time
pnpm --filter @optio/web test:e2e pass (107)
pnpm format:check, pnpm turbo typecheck, helm lint, helm template api.replicas=3 pass

PR #669 (phase B) was still open when this was pushed; origin/main hasn't moved, so no merge was needed.

jonwiggins and others added 12 commits October 10, 2026 18:11
One Redis subscriber per instance (started after listen, retried with
backoff, like the session bus) carrying binary envelopes on three kinds of
channel: a host's inbound channel (only its owner listens), a terminal's
outbound channel (instances with a viewer of it listen) and an instance's
reply channel (request replies correlated by id with a timeout, and direct
messages). An instance never hears its own messages.

Co-Authored-By: Claude <noreply@anthropic.com>
The instance a daemon connects to owns the host by the lease
local-host:<id>, held per daemon connection and seized on hello (newest
daemon wins); the previous owner closes its socket on host-taken or when
its renew fails, and a host is marked offline only while no connection
holds its lease. isHostOnline reads the lease; the hello's capabilities
are stored on local_hosts.capabilities (new migration) so every instance
answers the same.

Viewers, spawn, kill, input, dirs, limits, credentials and transcript
backfill on any instance reach the owner over the Local bus, each message
naming the lease holder it is for: frames to a daemon elsewhere are
acknowledged by its owner, daemon answers are relayed back to the asker,
and a remote viewer is a socket-shaped proxy in the owner's maps, so the
attach handshake and the grid arbiter work unchanged. A heartbeat reports
remote viewers to their owners and reconnects (4503) the ones whose daemon
moved. A same-instance reconnect now closes the old connection's viewers
too (they were left frozen). The token auto-refresh window is a lease and
the interaction throttle a conditional update.

Co-Authored-By: Claude <noreply@anthropic.com>
The sweep was a setInterval on every instance. It is now the BullMQ job
scheduler local-sweep.sweep, run under poller:local-sweep: stale and
unowned hosts go offline, terminals parked as host-offline whose host is
back are started, and launches an online daemon never acknowledged fail.

Co-Authored-By: Claude <noreply@anthropic.com>
scale-out-local.e2e.test.ts on startApiCluster: a viewer on B streams a
daemon on A (snapshot, output, input, resize, a status change through C);
spawn, input, kill, dirs and limits through B reach the daemon on A; the
daemon reconnects to B and A's viewer follows it; A is SIGKILLed while it
owns the host and the terminal lives on; a local Job run started through B
completes; every server agrees on online state. The scripted daemon and
viewer are shared by the Local e2e files (test-utils/e2e/local-fakes.ts).

Co-Authored-By: Claude <noreply@anthropic.com>
api.replicas may be any number, or api.autoscaling may be enabled (the HPA
then owns the count and the Deployment leaves replicas out). Rollouts stay
Recreate: migrations run at boot. The values and docs/production-eks.md say
what is per pod now: the per-address WebSocket cap and caches.

Co-Authored-By: Claude <noreply@anthropic.com>
Co-Authored-By: Claude <noreply@anthropic.com>
…mode

A broadcast PUBLISH reaches subscribers through whichever node took it, so
two frames on one channel could arrive swapped in REDIS_MODE=cluster. The
session bus and the Local bus now publish and subscribe through one choice
in redis-config.ts (publishOrdered / createOrderedSubscriber): sharded
pub/sub (SPUBLISH / SSUBSCRIBE) in cluster mode, plain pub/sub standalone.
The Local bus's request settles at once when nobody listens (ask() tells
unheard from unanswered). redis-cluster.int.test.ts sends 1000 frames on
one channel through each bus on the three-master test cluster.

Co-Authored-By: Claude <noreply@anthropic.com>
- host-taken only makes a connection re-read the lease; it gives the host
  up only once claimed and the lease names another connection. Frames
  naming a connection's holder that arrive before its claim completes are
  handled right after it. A connection closed mid-claim gives back exactly
  its own lease before the offline mark runs.
- An owner that heard a frame but didn't acknowledge it in time is
  "unacknowledged", never "not sent" (deliverToHost): a spawn stays
  launching and is never re-sent; started settles launching or pending.
- A Claude token never crosses Redis: the refresh runs on the daemon's
  instance as an owner task (runOnOwner) and only the outcome comes back;
  credentials requests are never forwarded nor their answers relayed.
- Relayed daemon answers are delivered in order per host; a remote viewer
  is reported only once its attach went out; the heartbeat runs one tick at
  a time with one lease query (leaseHolders).
- One pending-request helper (HostRequests) for dirs, limits and
  credentials; one DaemonReply type; one grid validator (saneGrid); the
  interaction throttle starts from the stamp as stored; the unused
  hostHasClaudeCredentials / hostCanRefreshLimits / hostCanManageDirs go.

Co-Authored-By: Claude <noreply@anthropic.com>
…ost list

getHost / listHosts / pickOnlineHost derive state from the host's lease in
the same query, so every instance says the same; the local_hosts.state
column stays the last-known state for alerts and the sweep. The host list
reads capabilities off the rows it loaded (capabilitiesOf), no query per
host. local-cluster.int.test.ts covers the late spawn acknowledgement, a
close mid-claim, lease-derived reads and the interaction throttle.

Co-Authored-By: Claude <noreply@anthropic.com>
Case (g) of scale-out-local.e2e.test.ts, 20 iterations with fast lease
renewal: one connection gives way with 4000, the other survives three
renews later, and every server says the host is online.

Co-Authored-By: Claude <noreply@anthropic.com>
…s, online)

Co-Authored-By: Claude <noreply@anthropic.com>
…8e8b99883b6

# Conflicts:
#	CHANGELOG.md
#	apps/api/src/services/repeat-jobs.ts
#	docs/plans/scale-out.md
@jonwiggins
jonwiggins merged commit 5bc58ab into main Oct 11, 2026
13 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant