Skip to content

feat: scale-out phase C2 — sessions, glance and skills on every API instance - #667

Merged
jonwiggins merged 9 commits into
mainfrom
feat/scale-out-sessions
Oct 10, 2026
Merged

jonwiggins merged 9 commits into
mainfrom
feat/scale-out-sessions

Conversation

@jonwiggins

@jonwiggins jonwiggins commented Oct 10, 2026 •

Copy link
Copy Markdown
Owner

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)

  • A chat turn still runs on the instance that received the prompt (the old exec path; recoverSessionTurns stays). 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.
  • Interrupt is durable: interruptSessionChat stops this instance's exec at once, records the interrupt on the running turn's row (finished_at set while state = '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-session and share revocation publish end / share_revoked control 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.
  • Deviation from the plan's wording, on purpose: one instance-wide channel instead of 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 (watchSessionAccess covers 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)

  • Every in-memory map and timer is a glance_state row: 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 (queue.upsertJobScheduler("glance-sweep", …)), each tick run under the sweep:glance lease (withLease). A timer fires once, up to one interval late (OPTIO_GLANCE_SWEEP_INTERVAL).
  • The APNs / FCM one-second Watch coalescers stay per instance (documented in docs/ios-push.md, docs/android-push.md and 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)

  • The sync worker clones into scratch under /tmp (OPTIO_SKILLS_SCRATCH_DIR) and writes every file to installed_skill_files with 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.
  • Spawn-time reads come from the database only (readInstalledSkillFiles(skillId, sha), which refuses files from another commit).
  • Upgrade path: a skill synced to the old cache volume (a resolved_sha, no files in the database) is due again on the first periodic sync (dueSkillIds).
  • Helm: the installedSkillsCache PVC template, volume and mount are removed (the value is ignored); docs/production-eks.md says helm upgrade deletes 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.

Tier Command Result
Unit pnpm turbo test 13/13 tasks green (api 182 files / 3030 tests; web 739; shared 748; cli 227; …)
Integration pnpm --filter @optio/api test:integration 54 files / 425 tests passed
Pipeline e2e pnpm --filter @optio/api test:e2e 24 files / 107 tests passed
Format / types pnpm format:check, pnpm turbo typecheck clean / 13/13
Helm helm lint helm/optio --set encryption.key=test 0 charts failed

New tests:

  • apps/api/e2e/scale-out-sessions.e2e.test.ts on startApiCluster (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 (row interrupted, both sockets see idle/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

  1. Boot no longer interrupts live turns. The instance running a chat turn holds the lease 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). recoverSessionTurns settles only running turns 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.
  2. The 10 s recheck also reads the session's state (auth disabled too); a socket on an ended pod session closes within one recheck (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.
  3. Interrupts name their turn: requestTurnInterrupt returns the flagged ids, the bus message carries turnId, the receiver acts only if its running turn matches, and nothing is published when nothing was flagged.
  4. startSessionBus runs after app.listen, in the background with backoff (1 s doubling to 30 s); REST serves while Redis is down.
  5. A failed SUBSCRIBE disconnects its client and resets state; isSessionBusStarted() is true only after a successful SUBSCRIBE.
  6. TURN_CONTROL_POLL_MS via parseIntEnv.

Glance
2. A terminal's previous attention state is read and replaced in one step (swapGlance: a transaction under pg_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 upsert WHERE value IS DISTINCT FROM excluded.value; the previous counts are one read (getGlanceMany), passed into syncWatch; a swap to the same value writes nothing.

Skills
3. sync-one jobs carry the id sync-one-<skillId> — a dash, not the colon the review asked for, because BullMQ throws Custom Id cannot contain : for an id with a single colon — and recordSyncResult upserts on (skill_id, path) inside its transaction.
6. cloneIntoScratch removes its scratch root on any throw.
10. Every successful sync writes a .optio-synced sentinel (content: the resolved sha); dueSkillIds / filesSynced look for it, readInstalledSkillFiles skips 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-due at start.

New and changed tests

  • e2e (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 a live idle status.
  • Integration: recovery settles a lapsed-lease turn and a lease-less legacy turn but not a leased one; an interrupt for another turn id is ignored; interruptSessionChat stops 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's xmin unchanged; 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.ts now lapses the turn's lease before recovery, as a restart would.
  • Unit: 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 in recordSyncResult.

Tier results after the fixes (final tree)

Tier Result
pnpm turbo test 13/13 tasks (api 183 files / 3034 tests)
pnpm --filter @optio/api test:integration 54 files / 432 tests passed
pnpm --filter @optio/api test:e2e 24 files / 108 tests passed
pnpm format:check, pnpm turbo typecheck, helm lint clean

…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.
@jonwiggins
jonwiggins merged commit 25d1619 into main Oct 10, 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