Repository navigation
feat: scale-out phase D: Optio Local relay across API instances, replica rule lifted - #670
Merged
Merged
Conversation
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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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 pinsapi.replicasto 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
daemonsByHost: which process has the daemonlocal-host:<hostId>inleases(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; publisheshost-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, andmarkHostOfflineis skipped while any connection holds one.isHostOnline/ capability flags from the socket maplocal_hosts.capabilities(new migration1792020000_local_host_capabilities) before the lease is taken, read only while it is live.hostsWithClaudeCredentialsjoins hosts and live leases.RemoteViewerproxy (socket-shaped) in the owner's maps; its frames go out onoptio:local:term:<id>:outaddressed 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 nowhereattachViewerreturns a handle that routes input, resize, view reports and pongs locally or overoptio:local:host:<id>:in.sendToHostis 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.local-sweep-worker.ts:setIntervalon every instancelocal-sweep.sweep(inREPEAT_SCHEDULERS) underpoller: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)claude-auto-refresh:<hostId>per window, shared by every instance.interactionWrittenAtThe bus is
services/local-bus.ts: one subscriber per instance (messageBuffer), started afterlistenand retried with backoff likesession-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; withapi.autoscaling.enabledthe Deployment omitsreplicas(the HPA owns it); rollouts stayRecreate.values.yaml,values.production.yaml(still 1 replica by default) anddocs/production-eks.mdsay what is per pod: the per-address WebSocket cap and caches.scripts/verify-production-chart.tsnow 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.
replay+ snapshot, live output in order,sizeframes (yoursafter its resize); B's input and resize reach the daemon; a kill through C reaches it and theexitstatus reaches the viewer on Bdirs-resultand the new list; limits refresh (C) gets its result; REST input (B) reaches the PTY; kill (B) reaches the daemon and the row staysrunninguntil the realexitlivestarted→ usage →exit 0→ runcompletedwith its cost, whichever server's worker dispatched it (spawn through a non-owner is pinned deterministically in (b))live(session recovery); after the daemon closes all three say offline /reconnecting, host rowoffline, terminal stillrunningTiers
pnpm turbo testlocal-bus.test.ts(10),local-relay-cluster.test.ts(13)pnpm --filter @optio/api test:integrationlocal-bus.int.test.ts(7, real Redis + Postgres: two bus instances route a request and its reply; lease seized → old owner closes onhost-takenor at renew, host stays online; same-instance reconnect keeps the new lease)pnpm --filter @optio/api test:e2elocal-terminal,local-run,local-daemon-auth, now on the shared fakes)pnpm --filter @optio/web test:e2epnpm format:check,pnpm turbo typecheckhelm lint helm/optio --set encryption.key=testhelm templatewithapi.replicas=3(default and production values), and withapi.autoscaling.enabled=truereplicas: 3, Recreate; with the HPA the Deployment has noreplicas.scripts/verify-production-chart.tspassesgen:swift/gen:kotlinNotes
local_hosts.capabilities, jsonb).Generated with Claude Code
Review fixes
All 15 findings from the high-effort review are fixed. Commits
3488a2b5..542a318c.REDIS_MODE=cluster(PUBLISH lands on any master)redis-config.ts:publishOrdered/createOrderedSubscriber/subscribeOrdered/onOrderedMessage. Cluster mode uses sharded pub/sub (SPUBLISH/SSUBSCRIBE, a cluster client withshardedSubscribers); standalone keeps PUBLISH / SUBSCRIBE. Bothlocal-bus.tsand C2'ssession-bus.tsuse 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)deliverToHostreturnssent/offline(nobody listening) /unacknowledged. Unacknowledged leaves the rowlaunchingand is never re-sent;handleStartedacceptslaunchingorpending; the stuck-launching sweep fails it otherwise.sendToHostcounts onlyofflineas not sent.local-cluster.int.test.ts: owner acks after the window → one spawn, stilllaunchingafter a hello flush and the sweep,runningonstarted;startedsettles apendingrow; unit test of the three outcomeshost-takenonly 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 alwaysWHERE holder = <own holder id>(releaseLease, also fromclaimHost).onHostIndropped frames before the claimh === holderregardless ofclaimed; frames before the claim completes are queued on the connection and handled right after itsendnaming the connection mid-claim reaches the daemon (and is acked) right after the claim, not beforedoneframeclaimHostdaemonDisconnectedwaits for the claim in flight, which releases exactly what it seized, then the offline mark runslocal-cluster.int.test.ts: close at 0/1/3/8/15 ms into the claim → no lease, host offline every timeHomeViewer.attachedis set only after theattachpublish was heard; the heartbeat skips the restviewersreport while the attach is held; one after it goes outrunOnOwner/onOwnerTask): ask, validate, store there; only{ ok, hostId }/{ ok:false, error }travels back. The owner refuses a forwardedcredentialsframe and never relays acredentials-result.local-auth-refresh.int.test.ts: another instance asks; the secret is stored; a PSUBSCRIBE tap onoptio:local:*sees nosk-ant-oat01in any message; a smuggledcredentialssend is refused. Unit: no relay ofcredentials-resultrequest()waited out the timeout with nobody listeningask()returnsunheardat once when PUBLISH/SPUBLISH reports 0 receivers;request()returns nulllocal-bus.test.ts: null /unheardwell under a second;timeoutwhen heard but unanswered; also in the cluster testleaseHolders(keys)query per tickGET /api/local/hostsqueried per hostlistHostsreads the lease with the rows (one query);capabilitiesOf(row)reads capabilities off the loaded rowlocal-cluster.int.test.ts"reads capabilities off the rows it loaded"getHost/listHosts/pickOnlineHostderivestatefrom the lease (in a cluster); the column stays last-known for alerts, and the sweep still sends offline notificationslocal-cluster.int.test.ts"online is the lease"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)HostRequests(local-host-requests.ts) used by credentials, dirs and limits; one exportedDaemonReply; onesaneGrid(local-grid.ts) for the daemon's replay grid and bus gridshostHasClaudeCredentials,hostCanRefreshLimits,hostCanManageDirsremoved; tests usehostCapabilitiesResults after the fixes
pnpm turbo testpnpm --filter @optio/api test:integrationpnpm --filter @optio/api test:e2escale-out-local.e2e.test.tspnpm --filter @optio/web test:e2epnpm format:check,pnpm turbo typecheck,helm lint,helm template api.replicas=3PR #669 (phase B) was still open when this was pushed;
origin/mainhasn't moved, so no merge was needed.