Repository navigation
feat(api): scale-out phase C1, coordination across API instances - #668
Merged
Merged
Conversation
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.
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 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
claimLockChainpromise mutex in the task and workflow workersservices/claim-lock.tswithClaimLock(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_timeout10 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 withclaim.afterCommitand runs after the commit.transitionTask/transitionWorkflowRunCasare split into a write step and an announce step to support this; their behaviour is unchanged.queue.add(…, { repeat })+cleanRepeatJobswiping every queue at bootservices/repeat-jobs.ts:upsertJobScheduler(<queue>.<tick>)with stable ids. A new scheduler first ticks one interval out, as the old repeatables did.cleanRepeatJobsand its boot calls are gone. Boot now runsremoveLegacyRepeatables, which removes only the hash-keyed entries an older build registered.services/poller-lease.tsunderPollerLease('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 ownconfig-sync:applylease. Each sweep is also safe when two do overlap: ticket sync claimsticket_sync_claimsbefore creating a task; the schedule checker advancesnext_fire_atby 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.SET NXdelivery dedupeservices/inbound-delivery-service.ts:INSERT … ON CONFLICT DO NOTHING RETURNINGintoinbound_webhook_deliveries, used by every receiver (GitHub, GitLab, Slack, Linear, Jira, PagerDuty, Sentry).workers/sweep-worker.ts: a housekeeping tick every 30 s under a lease, with aregisterSweep(name, fn)registry. It deletes deliveries older than 24 h, expiredws_upgrade_tokensand leases expired over a day ago.ws_upgrade_tokens: the SHA-256 is stored, the token is consumed byDELETE … RETURNING, 30 s TTL.activeConnectionsmapchat:<userId>per socket, renewed while opengithub-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.mdcovers what is shared through the database and what stays per pod. The per-IP WS cap is per replica.CHANGELOG.mdhas an entry under[Unreleased].Behaviour notes:
createTask) can be taken over after 10 minutes. Both takeovers happen in the same single statement, so two racing sweeps still get one claim.cleanLegacyRepeatJobsremoves 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 toREPEAT_SCHEDULERSwhen it lands, or every boot will remove it before the worker upserts it again.Tests
services/scale-out-coordination.int.test.ts):syncAllTickets()create one task.checkDueScheduleTriggers()fire once.sweepExternalPrs()launch one review and one run, and relaunch on new commits once.e2e/scale-out-coordination.e2e.test.ts, two servers onstartApiCluster, auth on):Review fixes
From the medium-effort review (15 items), plus one follow-up:
github-refresh:<userId>.withLeasegainedwait: { timeoutMs, pollMs }. It is the one lease-wait implementation: the refresh and the config apply use it, and the housekeeping tick now usesunderPollerLease.AbortSignal.anywith the lease's signal).bad_refresh_tokenthe tokens are compare-and-deleted in SQL on the refresh-token row's IV (secret-service.deleteSecretsIfUnchanged).getTokenForUserawaits the refresh, so a rejection falls back to the PAT. Unit test with a rejecting refresh.claimKeyhas no length cap. Tested with a 400-character repo URL (the stress test) and a 400-character key.SkillSyncBusyError. The queue'sdefaultJobOptionsgive every job 8 attempts with exponential backoff, which works with C2's stablesync-one-<id>job id.withClaimLocksetsisolationLevel: "read committed"explicitly, with a comment saying why.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.ConfigApplyBusyError(route test).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.recordTicketClaimTaskis logged; the task is still queued and the provider loop goes on (unit test).now() + 30 sand compared withnow()in the consumingDELETE … RETURNING(integration test).acquireLeaseerror closes the socket with 1011 and a readable message.SET LOCAL statement_timeout = '1s'. On timeout the delivery is accepted with a warning (unit test).afterCommit(step).claimTransitionIn/claimWorkflowRunInregister their announcement there,withClaimLockruns the steps after the commit and drops them on rollback, and no caller handles an announcement. CLAUDE.md namesclaimTransitionInas the in-transaction form oftransitionTask.recordSyncResulttakes 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 clearsresolved_shain 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/mainwith C2 (#667) is merged in.skill-sync-worker.tskeeps C2's database-backed sync and itssync-one-<id>job id, with the lease and retry on top.sync-duerepeat job stays a scheduler (scheduleRepeat), with C2's boot pass.index.tsstarts both the housekeeping and glance sweep workers. The glance sweep keeps its ownupsertJobScheduler; its queue is not inREPEAT_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.