Repository navigation
feat: scale-out phase C2 — sessions, glance and skills on every API instance - #667
Merged
Merged
Conversation
…stances A pod session's chat turn still runs on the instance that received the prompt, but its live frames reach viewers connected to any instance through the session bus (services/session-bus.ts): one Redis channel every instance subscribes to once at boot; the publisher delivers to its own sockets first and skips its own message when it comes back. An interrupt is durable: it stops this instance's exec at once, records finished_at on the still-running turn row, and nudges the bus; the instance running the turn acts on the message and polls the row every 2 s (watchTurnControl), so a lost message or a session ended elsewhere still stops it. Ending a session and revoking a share publish control messages that close every instance's sockets for the session; the 10 s access recheck stays as the fallback.
…hedule Every in-memory map and timer of the glance service is a glance_state row shared by every API instance (services/glance-state-service.ts): attention:<terminal>, needs-you:<user>, running:<user>, host-offline:<host>, and the timers watch-end:<user> (providers unioned in one statement) and snooze:<terminal> as rows with due_at. Due rows are claimed atomically (DELETE ... RETURNING) by a 30 s sweep, a BullMQ job scheduler with a stable id (upsertJobScheduler "glance-sweep") whose ticks run under the sweep:glance lease, so a timer fires once, up to one interval late. The APNs / FCM one-second Watch coalescers stay per instance; the push docs say so.
…emoved The skill sync worker clones into scratch under /tmp and writes every file of a synced skill to installed_skill_files with the manifest in one transaction. Spawn-time reads come from the database only (readInstalledSkillFiles(skillId, sha), which refuses files from another commit), so a run started on any API instance gets the skill. A re-sync at the same commit is a no-op only when its files are in the database, and a skill synced to the old cache volume before the upgrade is due again on the first periodic sync (dueSkillIds). The chart's installedSkillsCache PVC, volume and mount are gone (the value is ignored); docs/production-eks.md covers the upgrade.
On startApiCluster (two real servers, one database, one Redis): a viewer on server B sees the frames of a turn running on A; an interrupt sent through B stops the turn on A; a share revoked and a session ended through B close the sockets on A within the bus, not the 10 s recheck; and a skill synced by whichever server took the job is injected into a run started through the other.
…es overlap Review fixes for the marketplace skills in the database: - Every successful sync writes a `.optio-synced` sentinel row (content: the resolved commit); "synced at this commit" (dueSkillIds, filesSynced) is that row, so a skill with no files is not due forever, and readInstalledSkillFiles skips it. - A sync-one job carries the id `sync-one-<skillId>` (BullMQ refuses a custom id with a single colon), so duplicates collapse while queued, and recordSyncResult upserts on (skill_id, path) inside its transaction, so two overlapping syncs both succeed. - One directory walk gives the manifest, the totals, the cap check (before a file that would cross it is read) and the stored bytes. - cloneIntoScratch removes its scratch root itself when the clone throws. - The worker enqueues one sync-due when it starts, so skills synced before the upgrade are back within seconds.
…rites nothing Review fixes for the glance bookkeeping: - A terminal's previous attention state is read and replaced in one step (swapGlance: a transaction under pg_advisory_xact_lock(hashtext(key))), so N concurrent working -> needs_you changes alert exactly once. - A host's outage is claimed with INSERT ... ON CONFLICT DO NOTHING RETURNING (claimGlance); only the instance that created the row alerts. - remember() is one multi-row upsert written only where the value is distinct, the previous counts are one read, and a swap to the same value writes nothing: a terminal change that changes nothing writes no row.
…pts them Review fixes for session chat across instances: - The instance running a chat turn holds the lease chat-turn:<turnId> (60 s, renewed every 20 s, taken before the row is inserted, released when the turn settles). recoverSessionTurns settles only running turns whose lease is missing or expired, at boot and on every housekeeping sweep (the glance sweep job), and tells viewers the session is idle. Before, any replica booting marked every running turn interrupted and the 2 s control poll then killed live turns on healthy instances. - An interrupt names its turn: requestTurnInterrupt returns the flagged turn ids, the bus message carries turnId, the receiving instance acts only if its running turn matches, and nothing is published when nothing was flagged. - The 10 s access recheck also reads the session's state (auth disabled too): a socket on an ended pod session closes within one recheck. - The session bus subscribes after the server listens, in the background with backoff (1 s doubling to 30 s); a failed SUBSCRIBE disconnects its client, and isSessionBusStarted() is true only after one succeeded. - TURN_CONTROL_POLL_MS is read with parseIntEnv.
This was referenced Oct 10, 2026
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 C2 of
docs/plans/scale-out.md(§3 "State moves"): pod-session chat, glance / push bookkeeping and marketplace skills no longer depend on one API process.What moved where
Session chat, interrupt, end-session, share revocation (
services/session-bus.ts,session-turn-service.ts,session-sharing-service.ts,ws/session-chat.ts)recoverSessionTurnsstays). Its live frames reach viewers on any instance through the session bus: one Redis channel,optio:session-bus, that every instance subscribes to once at boot (startSessionBus, before the first socket can connect). The publishing instance delivers to its own sockets first and skips its own message when it comes back, so a local viewer never waits on Redis.interruptSessionChatstops this instance's exec at once, records the interrupt on the running turn's row (finished_atset whilestate = 'running', a state no other write produces — no new column, no migration), and nudges the bus. The instance running the turn acts on the bus message at once, and polls the row every 2 s (watchTurnControl/turnShouldStop), so a lost message, a session ended from any instance, or a boot-time recovery elsewhere still stops the exec.end/share_revokedcontrol messages; every instance closes the sockets it holds for that session (all sockets on end, collaborators on revoke). The 10 s DB access recheck stays as the fallback.optio:session:{id}with a subscriber per chat socket. Control messages must reach the instance running the turn (which may hold no chat socket for the session) and the instances holding terminal sockets (watchSessionAccesscovers pod terminals and chat), neither of which a per-chat-socket subscription reaches. One subscriber per instance also avoids a Redis connection per chat socket. The cost is that every instance receives every session's frames, which is small at the replica counts the chart allows. The plan's §3 table records this.Glance / push bookkeeping (
services/glance-state-service.ts,glance-service.ts,workers/glance-sweep-worker.ts)glance_staterow:attention:<terminal>,needs-you:<user>,running:<user>,host-offline:<host>, and the timerswatch-end:<user>(providers unioned in one statement) andsnooze:<terminal>as rows withdue_at.DELETE … RETURNING) by a 30 s sweep: a BullMQ job scheduler with a stable id (queue.upsertJobScheduler("glance-sweep", …)), each tick run under thesweep:glancelease (withLease). A timer fires once, up to one interval late (OPTIO_GLANCE_SWEEP_INTERVAL).docs/ios-push.md,docs/android-push.mdand the plan): worst case one duplicate push, folded by the collapse id.Marketplace skills in the database (
workers/skill-sync-worker.ts,services/installed-skill-service.ts,agent-environment-service.ts)/tmp(OPTIO_SKILLS_SCRATCH_DIR) and writes every file toinstalled_skill_fileswith the manifest in one transaction; a re-sync at the same commit is a no-op only when the files at that commit are in the database.readInstalledSkillFiles(skillId, sha), which refuses files from another commit).resolved_sha, no files in the database) is due again on the first periodic sync (dueSkillIds).installedSkillsCachePVC template, volume and mount are removed (the value is ignored);docs/production-eks.mdsayshelm upgradedeletes the claim and that nothing in it needs a backup.Docs: CHANGELOG under
[Unreleased], the plan's §3 rows marked shipped,docs/ios-push.md/docs/android-push.md(bookkeeping and the per-instance coalescer),docs/production-eks.md(PVC removal).Tests
All run on this branch's final tree.
pnpm turbo testpnpm --filter @optio/api test:integrationpnpm --filter @optio/api test:e2epnpm format:check,pnpm turbo typecheckhelm lint helm/optio --set encryption.key=testNew tests:
apps/api/e2e/scale-out-sessions.e2e.test.tsonstartApiCluster(two real servers, one DB, one Redis, fake runtime, auth on): a chat socket on B receives the frames of a turn running on A (same frames, same order as A's socket); an interrupt sent through B stops the[[mock:hang]]turn on A (rowinterrupted, both sockets seeidle/resumable); a share revoked through B closes the collaborator's socket on A with 4403 in under 5 s (the bus, not the 10 s recheck); a session ended through B closes the owner's socket on A and interrupts its running turn; a skill synced by whichever server took the job is injected into a Job run started through the other. A pod-session chat turn runs fine under the fake runtime, so no integration-only fallback was needed.session-bus.int.test.ts: the bus through real Redis — fan-out to other instances, own-message skip, interrupt from the bus, from the durable flag and from the session ending.glance-state.int.test.ts: row semantics, timer union, exactly-once claims across concurrent sweeps, the Watch's quiet-period end and a snooze expiry fired by a sweep on "another instance", one offline alert per outage, and the BullMQ scheduler (two workers started → one scheduler, due rows fired).skill-sync-worker.int.test.ts: files and manifest written from a real git clone, read with no shared filesystem, same-commit no-op, new commit replaces the set, failed sync keeps nothing, cascade on delete, and a pre-upgrade row (a commit, no files) found due and repaired.Not done here: the live tier (
scripts/smoke-e2e.sh) — this branch was not deployed to the local cluster. Chat turns still use the old exec path; moving them to the run protocol is a later phase, as the plan says.Review fixes
A medium-effort review found 15 issues; 14 are fixed here (the 15th, the live tier, runs after merge). Commits
5c7e40b4,8c316166,0ad63ba6,6e5524a0.Sessions
chat-turn:<turnId>(60 s TTL, renewed every 20 s, taken before the row is inserted so no recovery pass ever sees an unleased running turn, released when the turn settles).recoverSessionTurnssettles onlyrunningturns whose lease is missing or expired — at boot and on every 30 s housekeeping tick (the glance sweep job) — and broadcasts an idle status to the settled session's viewers. A crashed instance's turn is settled within about a TTL plus one sweep.OPTIO_SESSION_ACCESS_RECHECK_MS). Liveness is read on the timer and cached for the per-input check, so no query per keystroke is added. Local terminal streams are exempt: they replay an exited terminal's final screen on purpose.requestTurnInterruptreturns the flagged ids, the bus message carriesturnId, the receiver acts only if its running turn matches, and nothing is published when nothing was flagged.startSessionBusruns afterapp.listen, in the background with backoff (1 s doubling to 30 s); REST serves while Redis is down.isSessionBusStarted()is true only after a successful SUBSCRIBE.TURN_CONTROL_POLL_MSviaparseIntEnv.Glance
2. A terminal's previous attention state is read and replaced in one step (
swapGlance: a transaction underpg_advisory_xact_lock(hashtext(key))).7. Host-offline is claimed with
INSERT … ON CONFLICT DO NOTHING RETURNING(claimGlance); only the creator alerts.11.
remember()is one multi-row upsertWHERE value IS DISTINCT FROM excluded.value; the previous counts are one read (getGlanceMany), passed intosyncWatch; a swap to the same value writes nothing.Skills
3.
sync-onejobs carry the idsync-one-<skillId>— a dash, not the colon the review asked for, because BullMQ throwsCustom Id cannot contain :for an id with a single colon — andrecordSyncResultupserts on(skill_id, path)inside its transaction.6.
cloneIntoScratchremoves its scratch root on any throw.10. Every successful sync writes a
.optio-syncedsentinel (content: the resolved sha);dueSkillIds/filesSyncedlook for it,readInstalledSkillFilesskips it, and a root-level file of that name in a skill is not synced.12. One directory walk (
readSkillDir) yields the manifest, totals, cap check (before a file that would cross it is read) and the stored bytes.14. The worker enqueues one
sync-dueat start.New and changed tests
scale-out-sessions.e2e.test.ts): a turn running on A ([[mock:sleep:20000]]) is still running after a third server boots (cluster.add()), then completes normally with aliveidle status.interruptSessionChatstops a local turn without publishing, names a remote turn on the bus, and publishes nothing when nothing runs; a socket on a session ended without a bus message closes on the recheck; 8 concurrent working → needs_you changes alert exactly once; a no-op terminal change leaves every row'sxminunchanged; 5 concurrent host-offline events alert once; an empty skill is synced, read as[], and not due again; three overlapping syncs of one skill succeed; a failed clone leaves no scratch.session-security.int.test.tsnow lapses the turn's lease before recovery, as a restart would.session-bus.test.ts(failed SUBSCRIBE cleanup, background retry with backoff, stop cancels a retry);readSkillDir(one walk, cap, sentinel name skipped, empty skill); sentinel row inrecordSyncResult.Tier results after the fixes (final tree)
pnpm turbo testpnpm --filter @optio/api test:integrationpnpm --filter @optio/api test:e2epnpm format:check,pnpm turbo typecheck,helm lint