Skip to content

feat(watcher): per-source workers + global Claude session ceiling - #170

Open
ArnabChatterjee20k wants to merge 9 commits into
mainfrom
feat-source-workers
Open

ArnabChatterjee20k wants to merge 9 commits into
mainfrom
feat-source-workers

Conversation

@ArnabChatterjee20k

Copy link
Copy Markdown
Member

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:

  1. Polling was a single sequential loop. run_source_poll_loop computed one base_tick = min(all source intervals) and walked sources in order, poll_source(...).await one at a time. A slow source could delay every other source's poll.
  2. There was no global cap on Claude sessions. Concurrency was gated per-source and per-lane (fix vs QA) via an ad-hoc 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

          ┌──────────────────────────────────────────┐
          │  run_source_poll_loop (ONE loop)          │
          │  base_tick = min(all intervals)           │
          │  for each source: poll_source().await     │  <- sequential
          └──────────────────────────────────────────┘
 per-source gate (fix lane)   per-source gate (QA lane)   retry gate   deploy-qa gate
        AtomicUsize+Notify, counted per source, NO shared global cap
                      → total Claude sessions = unbounded

After

   ┌─ worker(sentry)  interval=Xs ─┐   each source = its own task,
   ├─ worker(github)  interval=Ys ─┤   polls independently, never
   ├─ worker(linear)  interval=Zs ─┤   blocked by a slow sibling
   └─ worker(...)                  ─┘
                │ poll_source → dispatch (per-source max_concurrent still applies)
                ▼
        process_issue()  ── acquires 1 permit from ──▶  Semaphore(max_concurrent_sessions)
        (fix / QA / retry / deploy-QA ALL funnel here)   GLOBAL CEILING, held for the run

Two independent levels:

  • Per-source worker budget (max_concurrent_for(source)): how many fixes a single source runs at once. Unchanged mechanism, now owned by an independent worker.
  • Global session ceiling (max_concurrent_sessions): a single tokio::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

Key Location Default Meaning
max_concurrent_sessions top-level [config] 12 New. Global ceiling on concurrent Claude sessions across every source and path. Clamped to >= 1.
max_concurrent top-level 1 Now documented as the per-source default for sources without an override.
max_concurrent [scm.github] 6 New field. GitHub worker default (retry-heavy).
max_concurrent [issues.sentry] 6 Default changed None → 6 (high-volume).

Also wired github_issues and helpscout into max_concurrent_for() (previously only linear/sentry/jira/gitlab were resolved; helpscout had the field but it was dead, and GitHub had none).

Default worker settings (as requested)

  • sentry → 6
  • github_issues → 6
  • everything else → 1 (falls back to global max_concurrent)

Example

# Global
max_concurrent = 1              # per-source default for unlisted sources
max_concurrent_sessions = 12    # hard machine-wide ceiling on Claude sessions

[scm.github]
# max_concurrent = 6            # (default) GitHub runs up to 6 fixes at once

[issues.sentry]
# max_concurrent = 6            # (default) Sentry runs up to 6 fixes at once

[issues.linear]
# unset → falls back to global max_concurrent = 1

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_sessions to tighten the machine bound below the sum of per-source limits.

Behavioral notes

  • QA, retry, and deploy-QA now also count against the global ceiling (previously only loosely bounded). This is intentional — the ceiling bounds all Claude sessions.
  • max_concurrent_sessions = 0 clamps to 1 (a zero ceiling would deadlock every issue).
  • Shutdown: each worker checks is_running after 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.
  • Updated test_per_source_max_concurrent_falls_back_to_global for the new Sentry default.
  • All 377 watcher + 339 config unit tests pass; cargo build, cargo test --no-run, and clippy clean.

🤖 Generated with Claude Code

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.
@hansi-codes

hansi-codes Bot commented Oct 5, 2026 •

Copy link
Copy Markdown

🟢 Tier S · Ready to merge

The incremental changes preserve behavior and introduce no new defects.

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.

Verdict New comments Fixed Still open
✅ Approved 0 0 0
📂 Walkthrough · 7
File Change
claudear.example.toml Documents per-source concurrency budgets and the global session ceiling.
crates/claudear-config/src/config.rs Adds concurrency configuration, source defaults, and tests; reformats the GitHub lookup.
crates/claudear-engine/src/watcher.rs Adds independent source workers, shared session gating, and concurrency tests; reformats test setup.
crates/claudear-engine/src/processing.rs Shares the session limiter with detached retrieval scoring.
crates/claudear-engine/src/api/auth.rs, crates/claudear-engine/src/api/routes.rs Updates test configuration literals for the new configuration field.
crates/claudear-integrations/src/webhook/configurator.rs Updates test configuration construction for the session ceiling.
src/webhook/server.rs Updates processor construction and test configuration for the shared limiter.

Reviewed the commits since 766180d · Details · Comment @hansi-codes review to re-run, or mention @hansi-codes with a question.

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Tier B · 1 blocking finding to address. Summary

Comment thread crates/claudear-engine/src/watcher.rs
Comment thread crates/claudear-engine/src/watcher.rs
Comment thread crates/claudear-engine/src/watcher.rs
- 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.
@ArnabChatterjee20k

Copy link
Copy Markdown
Member Author

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 poll_source now acquires a session permit around each classify_intent call (agent backend only; the local-LLM backend runs no Claude session so it's left ungated). This runs before dispatch, so it never double-acquires with the permit process_issue holds.

2. Recheck the rate-limit pause after waiting for a permit — process_issue now rechecks is_rate_limit_paused() immediately after acquiring the permit and bails if a pause began while queued. The deploy_qa tip is now claimed after that recheck, so a paused run can no longer leave a tip stuck in Running.

3. Test the ceiling through processing behavior — Added test_session_ceiling_blocks_processing_and_frees_on_completion: it saturates the ceiling by holding the only permit, asserts a real process_issue run parks (doesn't finish) while saturated, then releases the permit and asserts the run completes and the slot returns.

@hansi-codes review

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Tier A · Looks good to merge. Summary

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Tier A · Looks good to merge. Summary

Comment thread crates/claudear-engine/src/watcher.rs Outdated
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.
@ArnabChatterjee20k

Copy link
Copy Markdown
Member Author

Addressed the remaining inline suggestion in 8310d41:

Avoid replaying missed polls after a slow source cycle — the per-source worker now sets MissedTickBehavior::Delay on its interval timer, so a poll that runs longer than its interval no longer triggers a burst of back-to-back replayed polls. The next poll is spaced a full interval after the previous one finishes, matching the old loop's post-completion reset.

@hansi-codes review

@ArnabChatterjee20k
ArnabChatterjee20k changed the base branch from feat-mcp-list-filters to main October 5, 2026 06:43

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Tier A · Looks good to merge. Summary

Comment thread crates/claudear-engine/src/watcher.rs
# 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.
@ArnabChatterjee20k

Copy link
Copy Markdown
Member Author

@hansi-codes review

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Tier B · 1 blocking finding to address. Summary

Comment thread crates/claudear-engine/src/watcher.rs
Comment thread crates/claudear-engine/src/watcher.rs
… 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.
@ArnabChatterjee20k

Copy link
Copy Markdown
Member Author

@hansi-codes review

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 Tier S · Looks good to merge. Summary

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Tier B · 1 blocking finding to address. Summary

Comment thread crates/claudear-engine/src/watcher.rs
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.
@ArnabChatterjee20k

Copy link
Copy Markdown
Member Author

@hansi-codes review

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Tier A · Looks good to merge. Summary

Comment thread crates/claudear-engine/src/processing.rs

@hansi-codes hansi-codes Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Tier A · Looks good to merge. Summary

This branch has not been deployed

No deployments
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