Repository navigation
feat(watcher): per-source workers + global Claude session ceiling - #170
ArnabChatterjee20k wants to merge 9 commits into
Conversation
Give each source its own polling worker instead of one shared sequential loop, so a slow source cannot delay others. Add a global max_concurrent_sessions ceiling that every path (fix, QA, retry, deploy-QA) acquires from in process_issue, bounding total Claude sessions regardless of how high per-source max_concurrent limits sum. Default per-source workers: sentry=6, github=6, others fall back to the global max_concurrent (1). Wire github_issues and helpscout into max_concurrent_for. All configurable via config.
🟢 Tier S · Ready to merge
The pull request introduces independent polling workers per source and a shared semaphore limiting Claude sessions across processing, classification, and detached retrieval scoring. It adds configuration defaults and concurrency tests and updates configuration construction sites. The latest commits only reformat configuration lookup and test setup without changing behavior. Latest changes: The newest commits apply formatting-only changes to the GitHub concurrency lookup and the session-ceiling test setup.
📂 Walkthrough · 7
Reviewed the commits since |
- Gate poll-time agent intent classification with the session limiter (agent backend only), so classification sessions count against the global ceiling instead of bypassing it. - Recheck the rate-limit pause after acquiring a session permit, and claim the deploy_qa tip only after that recheck, so a run paused while queued never leaves a tip stuck in Running. - Add a behavioral test proving a saturated ceiling parks processing and that completion frees the slot.
|
Thanks @hansi-codes — addressed all three findings in fdcf624: 1. Include polling-time agent calls in the session ceiling — The agent-based intent classifier in 2. Recheck the rate-limit pause after waiting for a permit — 3. Test the ceiling through processing behavior — Added @hansi-codes review |
Raise qa.max_concurrent default from 1 to 6 so read-only answers run wider by default, matching the heavy fix workers. Still bounded by the global max_concurrent_sessions ceiling.
A poll can exceed its interval (slow API, waiting for dispatch slots). With the default Burst behavior the worker would replay every missed tick back-to-back afterwards, bursting redundant API calls. Delay keeps a full interval between the end of one poll and the start of the next, matching the old loop's post-completion spacing.
|
Addressed the remaining inline suggestion in 8310d41: Avoid replaying missed polls after a slow source cycle — the per-source worker now sets @hansi-codes review |
# Conflicts: # crates/claudear-engine/src/watcher.rs
Replace the single-run ceiling test with one that saturates a ceiling of 2 and shows three concurrent process_issue runs all park (none begin work) while the limit is held, then that every run frees its session slot on completion. Exercises the real acquire/release path, not just the initial permit count.
|
@hansi-codes review |
… worker independence - process_issue now rechecks is_stopped() after waiting on the session permit: the watcher can begin stopping while a run is queued on a saturated ceiling, and a drain must not have new runs start behind it once a slot frees. Bails before claiming the deploy_qa tip. - Add test_source_workers_poll_independently: a source stuck in a slow fetch_issues must not delay another source's polling.
|
@hansi-codes review |
The retrieval-quality judge spawns up to RETRIEVAL_JUDGE_CONCURRENCY agent sessions in a detached task, which bypassed the global ceiling (a ceiling of 1 could run the main session plus scoring sessions). Share the limiter (now Arc<Semaphore>) into IssueProcessor and acquire a permit per agent-backed scoring call so these sessions count against the same ceiling. The local-LLM backend runs no Claude session and stays ungated; callers without a watcher ceiling pass None.
|
@hansi-codes review |
What
Reshape concurrency around a per-source worker model with a global Claude-session ceiling. Each source now runs its own polling worker with its own concurrency budget, and one shared semaphore bounds the total number of Claude sessions across the whole machine.
Why
Two problems with the old design:
run_source_poll_loopcomputed onebase_tick = min(all source intervals)and walked sources in order,poll_source(...).awaitone at a time. A slow source could delay every other source's poll.AtomicUsize + Notify. With N sources each allowed their own budget — plus independent gates in the retry and deploy-QA paths — nothing bounded the total number of Claude CLI processes running at once. N sources at 6 each = 18+ concurrent sessions with no ceiling.Architecture: before / after
Before
After
Two independent levels:
max_concurrent_for(source)): how many fixes a single source runs at once. Unchanged mechanism, now owned by an independent worker.max_concurrent_sessions): a singletokio::sync::Semaphore.process_issue— the one funnel for fix, QA, retry, and deploy-QA — acquires a permit at the top and holds it for the entire run. This unifies the four previously-uncoordinated gates into one machine-wide bound.Config changes
max_concurrent_sessions[config]12>= 1.max_concurrent1max_concurrent[scm.github]6max_concurrent[issues.sentry]6None → 6(high-volume).Also wired
github_issuesandhelpscoutintomax_concurrent_for()(previously only linear/sentry/jira/gitlab were resolved;helpscouthad the field but it was dead, and GitHub had none).Default worker settings (as requested)
sentry→ 6github_issues→ 6max_concurrent)Example
With this config: Sentry and GitHub each run up to 6 fixes, Linear runs 1, but the total Claude sessions in flight never exceeds 12 — if both heavy sources saturate, further issues (including QA, retries, and deploy-QA runs) queue on the global semaphore until a slot frees. Lower
max_concurrent_sessionsto tighten the machine bound below the sum of per-source limits.Behavioral notes
max_concurrent_sessions = 0clamps to 1 (a zero ceiling would deadlock every issue).is_runningafter every tick and exits, same as before.Tests
test_default_worker_concurrency_for_heavy_sources— sentry/github default to 6, others to 1, global field defaults to 12.test_session_limiter_sized_from_config/test_session_limiter_clamps_zero_to_one.test_per_source_max_concurrent_falls_back_to_globalfor the new Sentry default.cargo build,cargo test --no-run, and clippy clean.🤖 Generated with Claude Code