Skip to content

feat(api): scale-out phase C1, coordination across API instances - #668

Merged
jonwiggins merged 13 commits into
mainfrom
feat/scale-out-coordination
Oct 10, 2026
Merged

jonwiggins merged 13 commits into
mainfrom
feat/scale-out-coordination

Conversation

@jonwiggins

@jonwiggins jonwiggins commented Oct 10, 2026 •

Copy link
Copy Markdown
Owner

Phase C1 of docs/plans/scale-out.md: the coordination that has to hold once several API pods share one database. It builds on the phase A schema (leases, ws_upgrade_tokens, inbound_webhook_deliveries, ticket_sync_claims, the unique partial index on active PR reviews). The chart still pins one replica; that limit is lifted in phase D.

What moved where

Before (in process) Now
claimLockChain promise mutex in the task and workflow workers services/claim-lock.ts withClaimLock(keys, fn): one READ COMMITTED transaction, pg_advisory_xact_lock(hashtext(key)) on the global key (claim:tasks, claim:workflows; per-repo and per-Job limits are counted under it), lock_timeout 10 s (a timeout resolves null and the worker re-queues). The count and the CAS claim run on the transaction (taskService.claimTransitionIn(claim, …), claimWorkflowRunIn(claim, …)), so a lock holder never needs a second pooled connection; the announcement (events, webhooks, reconciler) is registered with claim.afterCommit and runs after the commit. transitionTask / transitionWorkflowRunCas are split into a write step and an announce step to support this; their behaviour is unchanged.
queue.add(…, { repeat }) + cleanRepeatJobs wiping every queue at boot services/repeat-jobs.ts: upsertJobScheduler(<queue>.<tick>) with stable ids. A new scheduler first ticks one interval out, as the old repeatables did. cleanRepeatJobs and its boot calls are gone. Boot now runs removeLegacyRepeatables, which removes only the hash-keyed entries an older build registered.
Pollers overlapping across instances services/poller-lease.ts underPollerLease('poller:<name>') for ticket sync, schedule checker, external PR review, repo cleanup, PR watcher, reconcile resync, skill sync, token validation. The config directory apply uses its own config-sync:apply lease. Each sweep is also safe when two do overlap: ticket sync claims ticket_sync_claims before creating a task; the schedule checker advances next_fire_at by CAS before firing (advanceScheduleCas); an external PR review launch names the review the sweep saw (seenVersion) and lands only while that is still the review; the unique-index conflict counts as someone else's launch.
Slack/Linear in-memory set, Redis SET NX delivery dedupe services/inbound-delivery-service.ts: INSERT … ON CONFLICT DO NOTHING RETURNING into inbound_webhook_deliveries, used by every receiver (GitHub, GitLab, Slack, Linear, Jira, PagerDuty, Sentry).
(none) workers/sweep-worker.ts: a housekeeping tick every 30 s under a lease, with a registerSweep(name, fn) registry. It deletes deliveries older than 24 h, expired ws_upgrade_tokens and leases expired over a day ago.
In-memory WS upgrade token map ws_upgrade_tokens: the SHA-256 is stored, the token is consumed by DELETE … RETURNING, 30 s TTL.
Optio-chat activeConnections map lease chat:<userId> per socket, renewed while open
GitHub user-token refresh lock (per process) lease github-refresh:<userId>. A waiter reads the token the holder stored. The installation-token cache stays per instance (documented: GitHub mints as many tokens as are asked for).

Docs: docs/production-eks.md covers what is shared through the database and what stays per pod. The per-IP WS cap is per replica. CHANGELOG.md has an entry under [Unreleased].

Behaviour notes:

  • Ticket claims and deleted tasks: a claim whose task was deleted can be taken over, so a deleted ticket task is re-created as it was before claims existed. A claim with no task (its sweep died between the claim and createTask) can be taken over after 10 minutes. Both takeovers happen in the same single statement, so two racing sweeps still get one claim.
  • Legacy repeatable cleanup at boot: cleanLegacyRepeatJobs removes any repeat entry on a known queue whose key is not a scheduler id this build knows. That keeps an upgraded deployment from ticking twice. Phase B's attach-sweep scheduler must be added to REPEAT_SCHEDULERS when it lands, or every boot will remove it before the worker upserts it again.

Tests

  • Integration (services/scale-out-coordination.int.test.ts):
    • two claimers under a limit of 3 over 50 iterations each, with runs finishing at random, never exceed the limit. A watcher polls the count throughout, and the test fails with the advisory lock removed.
    • the claim locks never deadlock.
    • two poller-lease holders: one runs.
    • ticket claims: one of two wins, release works, and both takeovers work.
    • two concurrent syncAllTickets() create one task.
    • two schedule CAS calls: one wins.
    • two concurrent checkDueScheduleTriggers() fire once.
    • two concurrent sweepExternalPrs() launch one review and one run, and relaunch on new commits once.
    • one inbound delivery: one of two receivers claims it.
    • one WS token: one of two upgrades consumes it.
    • schedulers: one per id, and legacy repeatables are removed.
  • Pipeline e2e (e2e/scale-out-coordination.e2e.test.ts, two servers on startApiCluster, auth on):
    • a WS upgrade token minted on A is accepted on B, and A then rejects it with 4401.
    • one signed GitLab delivery posted to A and B fires one Job run.
    • a schedule trigger due now fires once with both servers sweeping every second.
    • after both boots there is exactly one scheduler per repeat job and no stray repeatables.

Review fixes

From the medium-effort review (15 items), plus one follow-up:

  1. GitHub token refresh (with 15). A token is only ever refreshed under github-refresh:<userId>.
    • withLease gained wait: { timeoutMs, pollMs }. It is the one lease-wait implementation: the refresh and the config apply use it, and the housekeeping tick now uses underPollerLease.
    • Inside the lease the stored token is re-read first.
    • The GitHub call has a 15 s timeout (AbortSignal.any with the lease's signal).
    • A wait that times out returns the stored token if it still works, else null, so the caller falls back to the PAT. Nothing refreshes without the lease.
    • On bad_refresh_token the tokens are compare-and-deleted in SQL on the refresh-token row's IV (secret-service.deleteSecretsIfUnchanged).
    • Integration test: two instances refreshing one user make one GitHub call and keep the new tokens. Compare-and-delete is tested both ways.
  2. getTokenForUser awaits the refresh, so a rejection falls back to the PAT. Unit test with a rejecting refresh.
  3. claimKey has no length cap. Tested with a 400-character repo URL (the stress test) and a 400-character key.
  4. A sync-one that finds the skill's lease held throws SkillSyncBusyError. The queue's defaultJobOptions give every job 8 attempts with exponential backoff, which works with C2's stable sync-one-<id> job id.
  5. withClaimLock sets isolationLevel: "read committed" explicitly, with a comment saying why.
  6. SET LOCAL lock_timeout = '10s': a timeout resolves null, which the workers already treat as re-queue with delay. Integration test: a holder that never commits makes a claimer give up after about 10 s without running.
  7. Sync now answers 409 with a readable message on ConfigApplyBusyError (route test).
  8. A launch without seenVersion (a person's) relaunches the current review without a CAS. Integration test: a stale sweep view is refused, then a person's launch relaunches.
  9. A failed recordTicketClaimTask is logged; the task is still queued and the provider loop goes on (unit test).
  10. WS upgrade token expiry is set with now() + 30 s and compared with now() in the consuming DELETE … RETURNING (integration test).
  11. Optio chat lease:
    • TTL 30 s, renewed every 10 s.
    • The release on close is awaited with a 2 s cap.
    • An acquireLease error closes the socket with 1011 and a readable message.
  12. The inbound delivery claim runs with SET LOCAL statement_timeout = '1s'. On timeout the delivery is accepted with a warning (unit test).
  13. The claim context has afterCommit(step). claimTransitionIn / claimWorkflowRunIn register their announcement there, withClaimLock runs the steps after the commit and drops them on rollback, and no caller handles an announcement. CLAUDE.md names claimTransitionIn as the in-transaction form of transitionTask.
  14. The per-repo and per-Job advisory keys are gone; a comment says per-repo and per-Job limits are counted under the global lock.
  • Follow-up: skill ref CAS. recordSyncResult takes the ref and subpath the sync read. In its transaction it locks the skill row (FOR UPDATE) and writes nothing if they changed. The PATCH that changes the ref clears resolved_sha in the same statement, so the skill is due at once. Integration test: the ref moves while a sync of main is in flight; nothing of main's sync lands, the skill is due, and the next sync writes the new ref's files.

Merge: origin/main with C2 (#667) is merged in.

  • skill-sync-worker.ts keeps C2's database-backed sync and its sync-one-<id> job id, with the lease and retry on top.
  • The sync-due repeat job stays a scheduler (scheduleRepeat), with C2's boot pass.
  • index.ts starts both the housekeeping and glance sweep workers. The glance sweep keeps its own upsertJobScheduler; its queue is not in REPEAT_SCHEDULERS, so the legacy cleanup at boot never touches it.

Results (after the review fixes and the merge with main)

  • pnpm turbo test: 13/13 tasks pass (api 3056 tests).
  • pnpm --filter @optio/api test:integration: 55 files, 450 tests pass.
  • pnpm --filter @optio/api test:e2e: 25 files, 112 tests pass.
  • pnpm format:check: clean.
  • pnpm turbo typecheck: 13/13.

The count-then-claim that admits a Repo Task run or a Job run under its
concurrency limits takes pg_advisory_xact_lock on the global key and the
repo's (or the Job's) key, in sorted order, inside one transaction
(services/claim-lock.ts, withClaimLock). The count and the CAS claim run
on that transaction (claimTransitionIn, claimWorkflowRunIn), so a holder
never needs a second pooled connection, and the claim is announced
(events, webhooks, the reconciler) after the commit. The per-process
promise-chain mutexes are gone.
Every receiver (GitHub, GitLab, Slack, Linear, Jira, PagerDuty, Sentry)
claims the provider's delivery id with one INSERT ... ON CONFLICT DO
NOTHING RETURNING into inbound_webhook_deliveries
(services/inbound-delivery-service.ts), so a delivery fires once however
many API instances receive copies. The per-process set and the Redis
SET NX path are gone. Fails open when the database errors, as before.
A token minted by one API instance is accepted by whichever instance the
upgrade lands on: createWsToken stores its SHA-256 with a 30 s expiry,
and validateWsToken consumes it with one DELETE ... RETURNING, so exactly
one upgrade gets it. The in-memory map and its timer are gone.
…apply

- Optio chat: one conversation per user across instances; each socket
  holds the lease chat:<userId> (renewed while open, released on close).
- A user's GitHub token refresh runs under github-refresh:<userId>
  (GitHub rotates the refresh token); an instance that finds it held waits
  and reads what the holder stored. The installation-token cache stays per
  instance, documented.
- The config directory's apply holds config-sync:apply: a tick skips when
  another instance is applying, Sync now waits up to a minute.
- Every repeat job is a BullMQ job scheduler with a stable id
  (services/repeat-jobs.ts, upsertJobScheduler): instances booting in any
  order leave one schedule per tick, and a new one first ticks one
  interval out. cleanRepeatJobs, which wiped every instance's schedules
  at boot, is gone; boot only removes the hash-keyed repeatables an older
  build registered.
- Every sweep (ticket sync, schedule checker, external PR review, repo
  cleanup, PR watcher, reconcile resync, skill sync, token validation)
  runs under the lease poller:<name> (services/poller-lease.ts), and is
  safe to overlap anyway: ticket sync claims (source, ticket, repo) in
  ticket_sync_claims before creating a task (a claim whose task was
  deleted, or whose sweep died, is taken over); the schedule checker
  advances next_fire_at by CAS before firing; an external PR review
  launch lands only while the review is still the one the sweep saw, and
  the one-active-review-per-PR conflict is someone else's launch.
- A housekeeping tick (workers/sweep-worker.ts, 30 s, under a lease,
  registerSweep) deletes inbound deliveries older than a day, expired
  WebSocket upgrade tokens and long-expired leases.
- Integration tests race two callers through every claim: two claimers
  under a limit of N over 50 iterations, two ticket syncs, two schedule
  sweeps, two external-PR sweeps, one lease, one delivery, one WS token,
  and the schedulers.
On startApiCluster (two real API servers, auth on, one fake-runtime dir):
a WebSocket upgrade token minted on A is accepted by B once; one signed
GitLab delivery posted to A and B fires once; a schedule trigger due now
fires once with both servers sweeping; both boots leave exactly one
scheduler per repeat job.
docs/production-eks.md says what is shared through the database now and
what stays per pod (the per-address WebSocket cap, the installation-token
cache).
- withClaimLock takes only the global keys (claim:tasks, claim:workflows);
  a repo's or a Job's own limit is counted under the same lock.
- The transaction is READ COMMITTED explicitly, and waits for the lock at
  most 10 s (SET LOCAL lock_timeout): a timeout resolves null, which the
  workers already treat as 're-queue with delay'.
- The claim context carries afterCommit(step): claimTransitionIn and
  claimWorkflowRunIn register their announcement there, and withClaimLock
  runs the steps after the commit (dropped on rollback), so no caller
  handles an announcement by hand. CLAUDE.md names it as the
  in-transaction form of transitionTask.
- claimKey takes keys of any length (hashtext does).
- withLease takes wait: { timeoutMs, pollMs }, the one lease-wait
  implementation: the GitHub refresh and the config apply use it, and the
  housekeeping tick runs under underPollerLease.
- A user's GitHub token is only ever refreshed under github-refresh:<id>.
  Inside, the stored token is re-read first (another instance may just
  have refreshed); the call to GitHub has a 15 s timeout; a wait that
  times out returns the stored token while it still works, else null so
  the caller falls back to the PAT. getTokenForUser awaits the refresh so
  a rejection falls back too. On bad_refresh_token the stored tokens are
  deleted only while the refresh token is the one sent (a compare-and-
  delete on the row's IV, secret-service deleteSecretsIfUnchanged).
- Sync now answers 409 with a readable message while another instance is
  applying the configuration directory.
- A person's PR review launch (no seenVersion) relaunches the review as
  it is; only a sweep's launch is a CAS on the version it saw.
- Ticket sync: failing to record a task on its claim is logged and the
  task is still queued; the provider's sweep goes on.
- WebSocket upgrade tokens: expiry minted and compared with the
  database's now().
- Optio chat: 30 s lease renewed every 10 s; release on close awaited with
  a 2 s cap; a failed acquire closes with 1011 and a readable reason.
- Inbound delivery claim: one second (statement_timeout); on timeout the
  delivery is accepted and processed, with a warning.
- Skill sync: a sync-one that finds the skill's lease held throws a
  retryable error; the queue's jobs retry with exponential backoff.
- Integration tests: after-commit steps, the 10 s lock timeout, a
  400-character key, a person's relaunch, two instances refreshing one
  user's GitHub token (one call), compare-and-delete on a refused refresh,
  WS token expiry on the database clock, the lease wait.
…63a45a59af5

# Conflicts:
#	apps/api/src/index.ts
#	apps/api/src/workers/skill-sync-worker.test.ts
#	apps/api/src/workers/skill-sync-worker.ts
…it resolved

recordSyncResult takes the ref and subpath the sync read and, in the same
transaction, locks the skill row (FOR UPDATE) and writes nothing unless
they still match. A PATCH that moves the ref clears resolved_sha in the
same statement, so the skill is due at once; a sync of the old ref that
finishes after it is discarded (its sync-one, sharing the job id, may
have absorbed the PATCH's), and the next sync-due takes the new ref.
Integration test: the ref moves while a sync of main is in flight; nothing
of main's lands, the skill is due, the next sync writes next's files.
@jonwiggins
jonwiggins merged commit 60c8d10 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