From 4f6f1a57286d1884ce9757f838ae7c9aee6ce990 Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 11:45:48 +0530 Subject: [PATCH 1/8] feat(watcher): per-source workers + global Claude session ceiling 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. --- claudear.example.toml | 19 +- crates/claudear-config/src/config.rs | 52 ++++- crates/claudear-engine/src/api/auth.rs | 1 + crates/claudear-engine/src/api/routes.rs | 1 + crates/claudear-engine/src/watcher.rs | 183 +++++++++++++----- .../src/webhook/configurator.rs | 1 + src/webhook/server.rs | 1 + 7 files changed, 206 insertions(+), 52 deletions(-) diff --git a/claudear.example.toml b/claudear.example.toml index 7733653f..315ebb16 100644 --- a/claudear.example.toml +++ b/claudear.example.toml @@ -23,9 +23,18 @@ poll_interval_ms = 300000 # Max issues to process per poll cycle (default: 5) max_issues_per_cycle = 5 -# Max concurrent issue processing (default: 1) +# Max concurrent issue processing PER SOURCE for sources without an explicit +# override (default: 1). Each source runs its own worker that polls and fixes +# independently; this caps how many fixes a single source runs at once. max_concurrent = 1 +# Global ceiling on concurrent Claude sessions across EVERY source and path +# (fix, QA, retry, deploy-QA), regardless of how high per-source max_concurrent +# limits sum (default: 12). This is the machine-wide safety bound. Clamped to +# >= 1. The built-in Sentry and GitHub sources default to a per-source +# max_concurrent of 6; all other sources default to the global max_concurrent. +max_concurrent_sessions = 12 + # Delay between processing issues in ms (default: 5000) processing_delay_ms = 5000 @@ -272,6 +281,10 @@ poll_interval_ms = 60000 # Auto-resolve issues on Linear/Sentry when PRs merge (default: false) auto_resolve_on_merge = false +# Max concurrent fixes for the GitHub source worker (default: 6). GitHub issues +# frequently retry, so it ships with a wider worker than the global default. +# max_concurrent = 6 + # Optional: Webhook secret for verifying GitHub webhook signatures # Set via GITHUB_WEBHOOK_SECRET env var for security webhook_secret = "" @@ -446,7 +459,9 @@ escalation_threshold_percent = 50 # Set via SENTRY_CLIENT_SECRET env var for security client_secret = "" -# Per-source rate limiting (overrides global values if set) +# Per-source rate limiting (overrides global values if set). +# The Sentry worker defaults to max_concurrent = 6 when omitted (wider than the +# global default since Sentry is high-volume and retry-heavy). max_issues_per_cycle = 2 max_concurrent = 4 diff --git a/crates/claudear-config/src/config.rs b/crates/claudear-config/src/config.rs index ee3d9cec..85eb0a7b 100644 --- a/crates/claudear-config/src/config.rs +++ b/crates/claudear-config/src/config.rs @@ -428,8 +428,14 @@ pub struct Config { pub db_path: PathBuf, /// Maximum issues to process per poll cycle. pub max_issues_per_cycle: usize, - /// Maximum concurrent issue processing. + /// Maximum concurrent issue processing (per source, default for sources + /// without an explicit override). pub max_concurrent: usize, + /// Global ceiling on concurrent Claude sessions across every source and + /// path (fix, QA, retry, deploy-QA). Per-source `max_concurrent` bounds how + /// many a single source runs; this bounds the machine-wide total regardless + /// of how high the per-source limits are summed. Clamped to >= 1. + pub max_concurrent_sessions: usize, /// Delay between processing issues (ms). pub processing_delay_ms: u64, /// Maximum number of activity entries to keep in the IPC server (default: 10,000). @@ -774,6 +780,7 @@ impl Default for Config { db_path: PathBuf::from("claudear.db"), max_issues_per_cycle: 5, max_concurrent: 1, + max_concurrent_sessions: 12, processing_delay_ms: 5000, max_activity_entries: 10_000, ipc_timeout_secs: 30, @@ -1721,6 +1728,10 @@ pub struct GitHubConfig { /// GitHub App configuration (nested under [scm.github.app]). #[serde(default)] pub app: GitHubAppConfig, + /// Maximum concurrent issue processing for the GitHub source (overrides + /// global `max_concurrent`). Defaults to a wider worker since GitHub issues + /// frequently retry. + pub max_concurrent: Option, } impl Default for GitHubConfig { @@ -1737,6 +1748,7 @@ impl Default for GitHubConfig { trigger_labels: Vec::new(), trigger_states: Vec::new(), app: GitHubAppConfig::default(), + max_concurrent: Some(6), } } } @@ -1756,6 +1768,7 @@ impl GitHubConfig { trigger_labels: vec!["auto-implement".to_string(), "claude".to_string()], trigger_states: vec!["open".to_string()], app: GitHubAppConfig::default(), + max_concurrent: Some(6), } } } @@ -2034,7 +2047,9 @@ impl Default for SentryConfig { escalation_threshold_percent: 50, client_secret: None, max_issues_per_cycle: None, - max_concurrent: None, + // Sentry issues are high-volume and often need retries, so run a + // wider worker by default. Override in [issues.sentry] if needed. + max_concurrent: Some(6), poll_interval_ms: None, } } @@ -3730,6 +3745,13 @@ impl Config { .as_ref() .and_then(|c| c.max_concurrent) .unwrap_or(self.max_concurrent), + "github_issues" => self.scm.github.max_concurrent.unwrap_or(self.max_concurrent), + "helpscout" => self + .issues + .helpscout + .as_ref() + .and_then(|c| c.max_concurrent) + .unwrap_or(self.max_concurrent), _ => self.max_concurrent, } } @@ -4774,6 +4796,8 @@ api_key = "key" sentry: Some(SentryConfig { auth_token: "tok".into(), org_slug: "org".into(), + // Opt out of the Sentry default (6) to exercise fallback. + max_concurrent: None, ..Default::default() }), ..Default::default() @@ -4785,6 +4809,30 @@ api_key = "key" assert_eq!(config.max_concurrent_for("unknown"), 4); } + #[test] + fn test_default_worker_concurrency_for_heavy_sources() { + // Sentry and GitHub ship with a wider default worker; everything else + // falls back to the global max_concurrent (1 by default). + let config = Config { + issues: IssuesConfig { + sentry: Some(SentryConfig { + auth_token: "tok".into(), + org_slug: "org".into(), + ..Default::default() + }), + ..Default::default() + }, + ..Default::default() + }; + assert_eq!(config.max_concurrent, 1); + assert_eq!(config.max_concurrent_sessions, 12); + assert_eq!(config.max_concurrent_for("sentry"), 6); + // GitHub config is always present (not optional) and defaults to 6. + assert_eq!(config.max_concurrent_for("github_issues"), 6); + assert_eq!(config.max_concurrent_for("linear"), 1); + assert_eq!(config.max_concurrent_for("unknown"), 1); + } + #[test] fn test_per_source_max_concurrent_overrides_global() { let config = Config { diff --git a/crates/claudear-engine/src/api/auth.rs b/crates/claudear-engine/src/api/auth.rs index 8ac7ccd8..66dced56 100644 --- a/crates/claudear-engine/src/api/auth.rs +++ b/crates/claudear-engine/src/api/auth.rs @@ -922,6 +922,7 @@ mod tests { db_path: ":memory:".into(), max_issues_per_cycle: 5, max_concurrent: 1, + max_concurrent_sessions: 12, processing_delay_ms: 5000, max_activity_entries: 100, ipc_timeout_secs: 30, diff --git a/crates/claudear-engine/src/api/routes.rs b/crates/claudear-engine/src/api/routes.rs index dd4927c3..6caf6c35 100644 --- a/crates/claudear-engine/src/api/routes.rs +++ b/crates/claudear-engine/src/api/routes.rs @@ -3053,6 +3053,7 @@ mod tests { db_path: ":memory:".into(), max_issues_per_cycle: 5, max_concurrent: 1, + max_concurrent_sessions: 12, processing_delay_ms: 5000, max_activity_entries: 100, ipc_timeout_secs: 30, diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 3e7c4913..7cbf93df 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -36,7 +36,7 @@ use std::pin::Pin; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; use tokio::sync::futures::Notified; -use tokio::sync::{Notify, RwLock}; +use tokio::sync::{Notify, RwLock, Semaphore}; use tokio::time::{interval, Duration}; /// A candidate issue ready for dispatch: the issue, its match result, and the @@ -303,6 +303,12 @@ pub struct Watcher { rate_limit_pause_until: RwLock>>, /// Notifies waiters when a processing slot becomes available. slot_available: Notify, + /// Global ceiling on concurrent Claude sessions across every source and + /// path. Each `process_issue` call holds one permit for its whole lifetime, + /// so the total number of issues being worked on at once never exceeds + /// `config.max_concurrent_sessions`, no matter how high per-source + /// `max_concurrent` limits sum. Per-source limits still apply on top. + session_limiter: Semaphore, /// Optional LLM analyzer for enhanced analysis across the pipeline. llm_analyzer: Option>, /// Intent classifier for QA-vs-fix routing. Backend selected by `qa.use_llm`: @@ -377,6 +383,8 @@ impl Watcher { )) }; + let session_limit = options.config.max_concurrent_sessions.max(1); + Self { agent: options.agent, qa_agent: options.qa_agent, @@ -403,6 +411,7 @@ impl Watcher { last_seen_releases: RwLock::new(HashMap::new()), rate_limit_pause_until: RwLock::new(HashMap::new()), slot_available: Notify::new(), + session_limiter: Semaphore::new(session_limit), llm_analyzer, intent_classifier, spawn_handles: tokio::sync::Mutex::new(Vec::new()), @@ -1075,6 +1084,10 @@ impl Watcher { self.config.max_issues_per_cycle ); tracing::info!(" Max concurrent: {} (global)", self.config.max_concurrent); + tracing::info!( + " Max concurrent sessions: {} (global ceiling)", + self.config.max_concurrent_sessions + ); for source in &self.sources { let src_max_issues = self.config.max_issues_per_cycle_for(source.name()); let src_max_concurrent = self.config.max_concurrent_for(source.name()); @@ -1146,65 +1159,68 @@ impl Watcher { /// Run the source polling loop. /// - /// Polls each source at its configured interval. Housekeeping is handled - /// separately by [`HousekeepingWorker`]. + /// Each source gets its own long-lived worker task that polls that source + /// at its configured interval, independently of the others. A slow source + /// can no longer delay a fast one, and each source's own `max_concurrent` + /// budget (plus the global session ceiling) bounds how many fixes it runs. + /// Housekeeping is handled separately by [`HousekeepingWorker`]. async fn run_source_poll_loop(self: &Arc, poll_interval: u64) -> Result<()> { - // Build per-source timer state: (source index, interval_ms, last_poll) - let now = std::time::Instant::now(); - let mut source_timers: Vec<(usize, u64, std::time::Instant)> = self - .sources - .iter() - .enumerate() - .map(|(i, source)| { - let src_interval = self.config.poll_interval_ms_for(source.name()).max(1); - tracing::info!( - source = source.name(), - interval_ms = src_interval, - "Per-source poll interval" - ); - (i, src_interval, now) - }) - .collect(); + let mut workers = Vec::with_capacity(self.sources.len()); + for (idx, source) in self.sources.iter().enumerate() { + let src_interval = self.config.poll_interval_ms_for(source.name()).max(1000); + tracing::info!( + source = source.name(), + interval_ms = src_interval, + max_concurrent = self.config.max_concurrent_for(source.name()), + "Starting source worker" + ); + let watcher = Arc::clone(self); + workers.push(tokio::spawn(async move { + watcher.run_source_worker(idx, src_interval).await; + })); + } - // Determine the base tick: minimum source interval or global, whichever is smallest. - // Cap at 1s to avoid busy-looping when all intervals are large. - let min_source_interval = source_timers - .iter() - .map(|(_, ms, _)| *ms) - .min() - .unwrap_or(poll_interval); - let base_tick_ms = min_source_interval.min(poll_interval).max(1000); - let mut base_timer = interval(Duration::from_millis(base_tick_ms)); - base_timer.tick().await; // Skip immediate first tick + // Nothing configured: idle until stopped so the caller's select! still + // has a future to hold. + if workers.is_empty() { + let _ = poll_interval; + while self.is_running.load(Ordering::SeqCst) { + tokio::time::sleep(Duration::from_millis(1000)).await; + } + return Ok(()); + } + + for worker in workers { + let _ = worker.await; + } + Ok(()) + } + + /// Worker loop for a single source: poll it every `interval_ms` until the + /// watcher stops. The first poll fires one interval after start (the + /// initial fan-out poll already ran in [`Self::start`]). + async fn run_source_worker(self: Arc, source_idx: usize, interval_ms: u64) { + let mut timer = interval(Duration::from_millis(interval_ms)); + timer.tick().await; // consume the immediate first tick while self.is_running.load(Ordering::SeqCst) { - base_timer.tick().await; + timer.tick().await; if !self.is_running.load(Ordering::SeqCst) { break; } if self.is_rate_limit_paused().await { continue; } - - // Poll each source whose interval has elapsed - for (src_idx, src_interval_ms, last_poll) in &mut source_timers { - let src_interval = Duration::from_millis(*src_interval_ms); - if last_poll.elapsed() >= src_interval { - let source = &self.sources[*src_idx]; - if let Err(e) = self.poll_source(source).await { - tracing::error!( - component = "watcher", - source = source.name(), - error = %e, - "Error polling source" - ); - } - *last_poll = std::time::Instant::now(); - } + let source = &self.sources[source_idx]; + if let Err(e) = self.poll_source(source).await { + tracing::error!( + component = "watcher", + source = source.name(), + error = %e, + "Error polling source" + ); } } - - Ok(()) } /// Stop the watcher. @@ -4230,6 +4246,26 @@ Create a PR with your changes.{custom_instructions}"#, } } + // Global session ceiling: hold one permit for the whole processing run. + // The per-source `_claim` above bounds how many of THIS source run at + // once; this bounds the machine-wide total across every source and path. + // The per-source claim is already held while we wait here, so no other + // item for this source can slip past its own budget meanwhile. + if self.session_limiter.available_permits() == 0 { + tracing::info!( + short_id = %issue.short_id, + limit = self.config.max_concurrent_sessions, + "Global session ceiling reached, waiting for a free Claude session slot" + ); + } + let _session_permit = match self.session_limiter.acquire().await { + Ok(permit) => permit, + Err(e) => { + tracing::error!(short_id = %issue.short_id, error = %e, "Session limiter closed, skipping issue"); + return false; + } + }; + let intent_label = match intent { Some(Intent::Question) => "question", Some(Intent::Bug) => "bug", @@ -5583,6 +5619,7 @@ mod tests { db_path: std::path::PathBuf::from(":memory:"), max_issues_per_cycle: 5, max_concurrent: 2, + max_concurrent_sessions: 12, processing_delay_ms: 1000, max_activity_entries: 100, ipc_timeout_secs: 30, @@ -5651,6 +5688,56 @@ mod tests { })) } + fn create_test_watcher_with_config(config: Config) -> Arc { + let notifier = Arc::new(MockNotifier::new(true)); + let tracker = Arc::new(SqliteTracker::in_memory().unwrap()); + let agent: Arc = + Arc::new(claudear_integrations::runner::ClaudeAgentRunner::new( + claudear_integrations::runner::ClaudeRunnerConfig::default(), + tracker.clone(), + )); + Arc::new(Watcher::new(WatcherOptions { + config, + sources: vec![], + notifier, + tracker: tracker.clone(), + inferrer: None, + embedding_client: None, + review_watcher: None, + issue_embedding_service: None, + code_search_service: None, + discord_search_service: None, + discord_index_orchestrator: None, + relationships: None, + github_client: None, + scm_provider: None, + user_registry: UserRegistry::new(std::collections::HashMap::new()), + agent, + classification_agent: None, + repo_classification_agent: None, + qa_agent: None, + dry_run: false, + llm_engine: None, + })) + } + + #[test] + fn test_session_limiter_sized_from_config() { + let mut config = test_config(); + config.max_concurrent_sessions = 3; + let watcher = create_test_watcher_with_config(config); + assert_eq!(watcher.session_limiter.available_permits(), 3); + } + + #[test] + fn test_session_limiter_clamps_zero_to_one() { + let mut config = test_config(); + config.max_concurrent_sessions = 0; + let watcher = create_test_watcher_with_config(config); + // A zero ceiling would deadlock every issue, so it clamps to 1. + assert_eq!(watcher.session_limiter.available_permits(), 1); + } + #[test] fn test_watcher_new() { let notifier = Arc::new(MockNotifier::new(true)); diff --git a/crates/claudear-integrations/src/webhook/configurator.rs b/crates/claudear-integrations/src/webhook/configurator.rs index 443bfaf8..2310340a 100644 --- a/crates/claudear-integrations/src/webhook/configurator.rs +++ b/crates/claudear-integrations/src/webhook/configurator.rs @@ -1688,6 +1688,7 @@ mod tests { db_path: "/tmp/test.db".into(), max_issues_per_cycle: 5, max_concurrent: 1, + max_concurrent_sessions: 12, processing_delay_ms: 5000, max_activity_entries: 100, ipc_timeout_secs: 30, diff --git a/src/webhook/server.rs b/src/webhook/server.rs index 63dce8ba..ae650419 100644 --- a/src/webhook/server.rs +++ b/src/webhook/server.rs @@ -1313,6 +1313,7 @@ mod tests { db_path: std::path::PathBuf::from(":memory:"), max_issues_per_cycle: 5, max_concurrent: 2, + max_concurrent_sessions: 12, processing_delay_ms: 1000, max_activity_entries: 100, ipc_timeout_secs: 30, From fdcf624c61846094741951a5054e0a73ed14de69 Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 11:55:29 +0530 Subject: [PATCH 2/8] fix(watcher): address Hansi review on session ceiling - 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. --- crates/claudear-engine/src/watcher.rs | 124 ++++++++++++++++++++++++-- 1 file changed, 115 insertions(+), 9 deletions(-) diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 7cbf93df..20e8bb8b 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -3438,6 +3438,13 @@ Create a PR with your changes.{custom_instructions}"#, .intent_classifier .clone() .expect("intent_classifier present (checked by qa_split_enabled)"); + // The agent-based classifier launches a Claude session per call, so + // it must count against the global ceiling too; otherwise a worker + // could spawn classification sessions while every permit is held by + // fixes. The local-LLM backend runs no Claude session, so it is not + // gated. This happens before dispatch, so it never double-acquires + // with the permit `process_issue` holds for an actual run. + let gate_classification = !self.config.qa.use_llm; let mut intents: Vec = Vec::with_capacity(ordered.len()); for (issue, _) in &ordered { // Ground the classification in the Discord reply thread (if any) so a @@ -3453,12 +3460,17 @@ Create a PR with your changes.{custom_instructions}"#, crate::processing::TranscriptTrust::ClaudearOnly, ) .await; - intents.push( + let intent = if gate_classification { + let _classify_permit = self.session_limiter.acquire().await.ok(); classifier .classify_intent(issue, conversation.as_deref()) .await - .unwrap_or(Intent::Fix), - ); + } else { + classifier + .classify_intent(issue, conversation.as_deref()) + .await + }; + intents.push(intent.unwrap_or(Intent::Fix)); } // Partition preserving prioritisation order within each bucket. Only @@ -4240,12 +4252,6 @@ Create a PR with your changes.{custom_instructions}"#, return false; }; - if let Some(ref tip) = deploy_qa_tip { - if !self.claim_deploy_qa_tip(tip, &issue.short_id) { - return false; - } - } - // Global session ceiling: hold one permit for the whole processing run. // The per-source `_claim` above bounds how many of THIS source run at // once; this bounds the machine-wide total across every source and path. @@ -4266,6 +4272,24 @@ Create a PR with your changes.{custom_instructions}"#, } }; + // Re-check the rate-limit pause after waiting for a permit: a pause can + // begin while queued here, and a saturated queue must not keep launching + // sessions through it. Bail before claiming the deploy_qa tip so a paused + // run never leaves a tip stuck in `Running`. + if self.is_rate_limit_paused().await { + tracing::info!( + short_id = %issue.short_id, + "Skipping issue: watcher paused for Claude rate limit while waiting for a session slot" + ); + return false; + } + + if let Some(ref tip) = deploy_qa_tip { + if !self.claim_deploy_qa_tip(tip, &issue.short_id) { + return false; + } + } + let intent_label = match intent { Some(Intent::Question) => "question", Some(Intent::Bug) => "bug", @@ -6720,6 +6744,88 @@ mod tests { watcher.active_processing.fetch_sub(1, Ordering::SeqCst); } + #[tokio::test] + async fn test_session_ceiling_blocks_processing_and_frees_on_completion() { + let notifier = Arc::new(MockNotifier::new(true)); + let tracker = Arc::new(SqliteTracker::in_memory().unwrap()); + let issue = Issue::new("1", "CEIL-1", "Ceiling issue", "http://example.com/1", "mock"); + let source = Arc::new(MockSource::with_issues("mock", vec![issue.clone()])) + as Arc; + + let mut config = test_config(); + config.max_concurrent_sessions = 1; + config.processing_delay_ms = 0; + + let watcher = Arc::new(Watcher::new(WatcherOptions { + config, + sources: vec![source.clone()], + notifier, + tracker: tracker.clone(), + inferrer: None, + embedding_client: None, + review_watcher: None, + issue_embedding_service: None, + code_search_service: None, + discord_search_service: None, + discord_index_orchestrator: None, + relationships: None, + github_client: None, + scm_provider: None, + user_registry: UserRegistry::new(std::collections::HashMap::new()), + agent: Arc::new(claudear_integrations::runner::ClaudeAgentRunner::new( + claudear_integrations::runner::ClaudeRunnerConfig::default(), + tracker.clone(), + )), + classification_agent: None, + repo_classification_agent: None, + qa_agent: None, + dry_run: false, + llm_engine: None, + })); + watcher.is_running.store(true, Ordering::SeqCst); + + // Saturate the global ceiling by holding its only session permit. + let held = watcher + .session_limiter + .try_acquire() + .expect("the single session permit should be free"); + assert_eq!(watcher.session_limiter.available_permits(), 0); + + // Start a real processing run. It must park on the session permit. + let w = Arc::clone(&watcher); + let handle = tokio::spawn(async move { + w.process_issue( + source, + issue, + MatchResult::matched("Test", MatchPriority::Normal), + None, + None, + None, + ) + .await + }); + + // While the ceiling is saturated the run cannot proceed. + tokio::time::sleep(Duration::from_millis(100)).await; + assert!( + !handle.is_finished(), + "processing ran even though the session ceiling was saturated" + ); + + // Freeing the slot lets the parked run proceed to completion. + drop(held); + let _ = tokio::time::timeout(Duration::from_secs(10), handle) + .await + .expect("processing did not complete after a session slot was freed"); + + // Completion returns the permit to the ceiling. + assert_eq!( + watcher.session_limiter.available_permits(), + 1, + "completing a run must free its session slot" + ); + } + #[tokio::test] async fn test_dispatch_lane_wakes_for_slot_freed_during_rate_limit_check() { let issue = Issue::new("1", "T-1", "Test Issue", "http://example.com/1", "mock"); From 9c7351451ceb7910fb2ac6172021ddb9aabfb079 Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 12:04:04 +0530 Subject: [PATCH 3/8] feat(qa): default QA concurrency to 6 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. --- claudear.example.toml | 5 +++-- crates/claudear-config/src/config.rs | 2 +- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/claudear.example.toml b/claudear.example.toml index 315ebb16..c728be87 100644 --- a/claudear.example.toml +++ b/claudear.example.toml @@ -1031,8 +1031,9 @@ use_llm = false # Max questions answered concurrently per source, independent of the fix-lane # concurrency budget (max_concurrent / [sources.].max_concurrent). Increase # this to answer more questions in parallel without affecting fix throughput. -# Clamped to 1 if set to 0 to avoid deadlock. (default: 1) -max_concurrent = 1 +# Still bounded by the global max_concurrent_sessions ceiling. +# Clamped to 1 if set to 0 to avoid deadlock. (default: 6) +max_concurrent = 6 # ============================================ # Reply Action diff --git a/crates/claudear-config/src/config.rs b/crates/claudear-config/src/config.rs index 85eb0a7b..a12612b6 100644 --- a/crates/claudear-config/src/config.rs +++ b/crates/claudear-config/src/config.rs @@ -888,7 +888,7 @@ impl Default for QaConfig { answer_timeout_secs: 600, max_qa_per_cycle: 20, use_llm: false, - max_concurrent: 1, + max_concurrent: 6, } } } From 8310d41c9b6f77007a9c54770c316038bca2de8e Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 12:11:02 +0530 Subject: [PATCH 4/8] fix(watcher): space source-worker polls with MissedTickBehavior::Delay 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. --- crates/claudear-engine/src/watcher.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 20e8bb8b..38295de6 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -37,7 +37,7 @@ use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; use tokio::sync::futures::Notified; use tokio::sync::{Notify, RwLock, Semaphore}; -use tokio::time::{interval, Duration}; +use tokio::time::{interval, Duration, MissedTickBehavior}; /// A candidate issue ready for dispatch: the issue, its match result, and the /// decided routing `Intent` (`None` for non-QA-eligible / QA-disabled sources). @@ -1201,6 +1201,11 @@ impl Watcher { /// initial fan-out poll already ran in [`Self::start`]). async fn run_source_worker(self: Arc, source_idx: usize, interval_ms: u64) { let mut timer = interval(Duration::from_millis(interval_ms)); + // A poll can run longer than the interval (slow API, waiting for dispatch + // slots). Delay, not the default Burst, so missed ticks are not replayed + // back-to-back afterwards; the next poll is spaced a full interval after + // the previous one finishes, matching the old loop's post-completion reset. + timer.set_missed_tick_behavior(MissedTickBehavior::Delay); timer.tick().await; // consume the immediate first tick while self.is_running.load(Ordering::SeqCst) { From f127a7527563eda1f1541adcd5aee20594b9d7de Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 12:36:58 +0530 Subject: [PATCH 5/8] test(watcher): prove the session ceiling caps concurrent processing 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. --- crates/claudear-engine/src/watcher.rs | 92 +++++++++++++++++---------- 1 file changed, 59 insertions(+), 33 deletions(-) diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 595b76c5..4da97cf5 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -7359,15 +7359,22 @@ mod tests { } #[tokio::test] - async fn test_session_ceiling_blocks_processing_and_frees_on_completion() { + async fn test_session_ceiling_caps_concurrent_processing_and_frees_slots() { let notifier = Arc::new(MockNotifier::new(true)); let tracker = Arc::new(SqliteTracker::in_memory().unwrap()); - let issue = Issue::new("1", "CEIL-1", "Ceiling issue", "http://example.com/1", "mock"); - let source = Arc::new(MockSource::with_issues("mock", vec![issue.clone()])) + let issues = vec![ + Issue::new("1", "CEIL-1", "Ceiling issue 1", "http://example.com/1", "mock"), + Issue::new("2", "CEIL-2", "Ceiling issue 2", "http://example.com/2", "mock"), + Issue::new("3", "CEIL-3", "Ceiling issue 3", "http://example.com/3", "mock"), + ]; + let source = Arc::new(MockSource::with_issues("mock", issues.clone())) as Arc; let mut config = test_config(); - config.max_concurrent_sessions = 1; + // Global ceiling of 2; a generous per-source budget so the per-source + // lane is never the binding limit (the ceiling is what we test). + config.max_concurrent_sessions = 2; + config.max_concurrent = 10; config.processing_delay_ms = 0; let watcher = Arc::new(Watcher::new(WatcherOptions { @@ -7398,46 +7405,65 @@ mod tests { })); watcher.is_running.store(true, Ordering::SeqCst); - // Saturate the global ceiling by holding its only session permit. - let held = watcher + // Saturate the global ceiling by holding both session permits, standing + // in for two runs already in flight. + let held_a = watcher .session_limiter .try_acquire() - .expect("the single session permit should be free"); + .expect("the first session permit should be free"); + let held_b = watcher + .session_limiter + .try_acquire() + .expect("the second session permit should be free"); assert_eq!(watcher.session_limiter.available_permits(), 0); - // Start a real processing run. It must park on the session permit. - let w = Arc::clone(&watcher); - let handle = tokio::spawn(async move { - w.process_issue( - source, - issue, - MatchResult::matched("Test", MatchPriority::Normal), - None, - None, - None, - None, - ) - .await - }); + // Start three more real processing runs. With the ceiling saturated, + // none may begin work: concurrent processing cannot exceed the limit. + let handles: Vec<_> = issues + .into_iter() + .map(|issue| { + let w = Arc::clone(&watcher); + let s = Arc::clone(&source); + tokio::spawn(async move { + w.process_issue( + s, + issue, + MatchResult::matched("Test", MatchPriority::Normal), + None, + None, + None, + None, + ) + .await + }) + }) + .collect(); - // While the ceiling is saturated the run cannot proceed. tokio::time::sleep(Duration::from_millis(100)).await; assert!( - !handle.is_finished(), - "processing ran even though the session ceiling was saturated" + handles.iter().all(|h| !h.is_finished()), + "a run started even though the session ceiling was saturated" ); + assert_eq!(watcher.session_limiter.available_permits(), 0); - // Freeing the slot lets the parked run proceed to completion. - drop(held); - let _ = tokio::time::timeout(Duration::from_secs(10), handle) - .await - .expect("processing did not complete after a session slot was freed"); - - // Completion returns the permit to the ceiling. + // Freeing the ceiling lets the parked runs proceed; each must release + // its slot on completion, so all permits return. + drop(held_a); + drop(held_b); + let drained = tokio::time::timeout(Duration::from_secs(10), async { + for handle in handles { + let _ = handle.await; + } + }) + .await; + assert!( + drained.is_ok(), + "processing did not complete after the session slots were freed" + ); assert_eq!( watcher.session_limiter.available_permits(), - 1, - "completing a run must free its session slot" + 2, + "completing every run must free its session slot" ); } From 8be4ddcc98e47d40153bb22a6783ad571c75d9eb Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 12:50:21 +0530 Subject: [PATCH 6/8] fix(watcher): recheck shutdown after acquiring a session permit; test 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. --- crates/claudear-engine/src/watcher.rs | 135 ++++++++++++++++++++++++++ 1 file changed, 135 insertions(+) diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 4da97cf5..388507b9 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -4802,6 +4802,19 @@ Create a PR with your changes.{custom_instructions}"#, } }; + // Re-check shutdown after waiting for a permit: the watcher can begin + // stopping while this run is queued on a saturated ceiling, and a drain + // must not have new runs start behind it once a slot frees. Bail before + // claiming the deploy_qa tip so a stopped run never leaves a tip in + // `Running`. + if self.is_stopped() { + tracing::info!( + short_id = %issue.short_id, + "Not starting issue processing: watcher began stopping while waiting for a session slot" + ); + return IssueRun::Stopping; + } + // Re-check the rate-limit pause after waiting for a permit: a pause can // begin while queued here, and a saturated queue must not keep launching // sessions through it. Bail before claiming the deploy_qa tip so a paused @@ -7467,6 +7480,128 @@ mod tests { ); } + /// Source whose `fetch_issues` blocks at a gate, standing in for a slow API, + /// and counts how many times it was polled. + struct BlockingFetchSource { + name: &'static str, + gate: Arc, + fetch_calls: AtomicUsize, + } + + #[async_trait] + impl IssueSource for BlockingFetchSource { + fn name(&self) -> &str { + self.name + } + fn display_name(&self) -> &str { + self.name + } + async fn fetch_issues(&self) -> Result> { + self.fetch_calls.fetch_add(1, AtomicOrdering::SeqCst); + self.gate.pass().await; + Ok(vec![]) + } + fn matches_criteria(&self, _issue: &Issue) -> MatchResult { + MatchResult::matched("blocking match", MatchPriority::Normal) + } + async fn build_issue_context(&self, issue: &Issue) -> Result { + Ok(format!("Context for {}", issue.short_id)) + } + async fn get_issue(&self, id: &str) -> Result { + Err(claudear_core::error::Error::source( + self.name, + format!("Issue {id} not found"), + )) + } + } + + #[tokio::test] + async fn test_source_workers_poll_independently() { + // A source stuck in a slow `fetch_issues` must not delay another + // source's polling: each runs in its own worker. + let gate = Arc::new(Gate::default()); + let slow = Arc::new(BlockingFetchSource { + name: "slow", + gate: Arc::clone(&gate), + fetch_calls: AtomicUsize::new(0), + }); + let fast = Arc::new(MockSource::new("fast")); + let sources: Vec> = vec![ + Arc::clone(&slow) as Arc, + Arc::clone(&fast) as Arc, + ]; + + let mut config = test_config(); + config.poll_interval_ms = 1000; // one poll per second per source + + let tracker: Arc = + Arc::new(SqliteTracker::in_memory().unwrap()); + let watcher = Arc::new(Watcher::new(WatcherOptions { + config, + sources, + notifier: Arc::new(MockNotifier::new(true)), + tracker: Arc::clone(&tracker), + inferrer: None, + embedding_client: None, + review_watcher: None, + issue_embedding_service: None, + code_search_service: None, + discord_search_service: None, + discord_index_orchestrator: None, + relationships: None, + github_client: None, + scm_provider: None, + user_registry: UserRegistry::new(std::collections::HashMap::new()), + agent: Arc::new(claudear_integrations::runner::ClaudeAgentRunner::new( + claudear_integrations::runner::ClaudeRunnerConfig::default(), + tracker, + )), + classification_agent: None, + repo_classification_agent: None, + qa_agent: None, + dry_run: true, + llm_engine: None, + })); + watcher.is_running.store(true, Ordering::SeqCst); + + let poller = { + let w = Arc::clone(&watcher); + tokio::spawn(async move { w.run_source_poll_loop(1000).await }) + }; + + // Wait (robustly, not on a fixed sleep that flakes under CI load) until + // the fast source has polled at least twice while the slow source stays + // stuck in its first fetch. + let deadline = std::time::Instant::now() + Duration::from_secs(20); + while fast.fetch_call_count() < 2 && std::time::Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(100)).await; + } + + let fast_polls = fast.fetch_call_count(); + let slow_polls = slow.fetch_calls.load(AtomicOrdering::SeqCst); + assert!( + fast_polls >= 2, + "the fast source should have polled repeatedly while the slow source was blocked, got {fast_polls}" + ); + // The slow source is parked in its first fetch, so it cannot poll again: + // its worker makes no progress while the fast one keeps going. + assert_eq!( + slow_polls, 1, + "the slow source should still be parked in its first fetch, got {slow_polls}" + ); + assert!( + !gate.was_cancelled(), + "the slow source's fetch should still be waiting at the gate" + ); + + // Release the slow source and stop; the loop must wind down cleanly. + gate.open(); + watcher.stop(); + let _ = tokio::time::timeout(Duration::from_secs(10), poller) + .await + .expect("the poll loop did not stop after the watcher was stopped"); + } + #[tokio::test] async fn test_dispatch_lane_wakes_for_slot_freed_during_rate_limit_check() { let issue = Issue::new("1", "T-1", "Test Issue", "http://example.com/1", "mock"); From 766180dc7056f3650faf6aa7581c515caf3cfe6b Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 13:00:44 +0530 Subject: [PATCH 7/8] fix(watcher): bound detached retrieval judges by the session ceiling 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) 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. --- crates/claudear-engine/src/processing.rs | 17 +++++++++++++++++ crates/claudear-engine/src/watcher.rs | 8 ++++++-- src/webhook/server.rs | 1 + 3 files changed, 24 insertions(+), 2 deletions(-) diff --git a/crates/claudear-engine/src/processing.rs b/crates/claudear-engine/src/processing.rs index 2f50729f..97a53c31 100644 --- a/crates/claudear-engine/src/processing.rs +++ b/crates/claudear-engine/src/processing.rs @@ -726,6 +726,11 @@ pub struct IssueProcessor { /// provider, local-LLM-based when `agent.use_llm` is set. `None` falls back /// to the label/source heuristic. pub intent_classifier: Option>, + /// Shared global Claude-session ceiling. When set, detached agent work the + /// processor spawns (the retrieval-quality judge) acquires a permit per + /// session so it counts against the same limit as the main run. `None` + /// leaves that work ungated (callers without a watcher-level ceiling). + pub session_limiter: Option>, } /// Everything the caller provides to `IssueProcessor::run()`. @@ -3643,6 +3648,9 @@ impl IssueProcessor { let tracker = self.tracker.clone(); let analyzer = self.llm_analyzer.clone(); let agent = self.qa_agent.clone().unwrap_or_else(|| self.agent.clone()); + // Each agent-backed score is a Claude session, so bound it by the same + // global ceiling as the main run (no-op for the local-LLM backend). + let session_limiter = self.session_limiter.clone(); // Issue identity captured for the detached task's timeline events. let issue_id = issue.id.clone(); let short_id = issue.short_id.clone(); @@ -3712,8 +3720,15 @@ impl IssueProcessor { let tracker = &tracker; let agent = agent.as_ref(); let scored = &scored; + let session_limiter = session_limiter.as_ref(); futures::stream::iter(items.iter()) .for_each_concurrent(RETRIEVAL_JUDGE_CONCURRENCY, |item| async move { + // Hold a global session permit for the scoring call so + // these detached sessions never exceed the ceiling. + let _session_permit = match session_limiter { + Some(limiter) => limiter.acquire().await.ok(), + None => None, + }; if let Some(score) = crate::agent_classifier::score_chunk_relevance_via_agent( agent, @@ -5242,6 +5257,7 @@ mod tests { github_client: None, llm_analyzer: None, intent_classifier: None, + session_limiter: None, }; let input = ProcessingInput { @@ -6356,6 +6372,7 @@ mod tests { github_client: None, llm_analyzer: None, intent_classifier: None, + session_limiter: None, } } diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 388507b9..03dbb634 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -418,7 +418,9 @@ pub struct Watcher { /// so the total number of issues being worked on at once never exceeds /// `config.max_concurrent_sessions`, no matter how high per-source /// `max_concurrent` limits sum. Per-source limits still apply on top. - session_limiter: Semaphore, + /// Shared (`Arc`) so detached session-spawning work — e.g. the retrieval + /// judge in [`IssueProcessor`] — counts against the same ceiling. + session_limiter: Arc, /// Optional LLM analyzer for enhanced analysis across the pipeline. llm_analyzer: Option>, /// Intent classifier for QA-vs-fix routing. Backend selected by `qa.use_llm`: @@ -554,7 +556,7 @@ impl Watcher { last_seen_releases: RwLock::new(HashMap::new()), rate_limit_pause_until: RwLock::new(HashMap::new()), slot_available: Notify::new(), - session_limiter: Semaphore::new(session_limit), + session_limiter: Arc::new(Semaphore::new(session_limit)), llm_analyzer, intent_classifier, spawn_handles: tokio::sync::Mutex::new(Vec::new()), @@ -5046,6 +5048,7 @@ Create a PR with your changes.{custom_instructions}"#, github_client: self.github_client.clone(), llm_analyzer: self.llm_analyzer.clone(), intent_classifier: self.intent_classifier.clone(), + session_limiter: Some(Arc::clone(&self.session_limiter)), }; let input = ProcessingInput { @@ -5827,6 +5830,7 @@ Create a PR with your changes.{custom_instructions}"#, github_client: self.github_client.clone(), llm_analyzer: self.llm_analyzer.clone(), intent_classifier: self.intent_classifier.clone(), + session_limiter: Some(Arc::clone(&self.session_limiter)), }; let input = ProcessingInput { diff --git a/src/webhook/server.rs b/src/webhook/server.rs index 812897f7..347ea27c 100644 --- a/src/webhook/server.rs +++ b/src/webhook/server.rs @@ -1183,6 +1183,7 @@ async fn process_issue( github_client: None, llm_analyzer: None, intent_classifier: None, + session_limiter: None, }; let input = ProcessingInput { From 8c90bd635a6c6e1ed1240928d610b6761fe05816 Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 5 Oct 2026 13:21:50 +0530 Subject: [PATCH 8/8] style: cargo fmt --- crates/claudear-config/src/config.rs | 6 +++++- crates/claudear-engine/src/watcher.rs | 28 ++++++++++++++++++++++----- 2 files changed, 28 insertions(+), 6 deletions(-) diff --git a/crates/claudear-config/src/config.rs b/crates/claudear-config/src/config.rs index ac7416cf..c2e9333a 100644 --- a/crates/claudear-config/src/config.rs +++ b/crates/claudear-config/src/config.rs @@ -3823,7 +3823,11 @@ impl Config { .as_ref() .and_then(|c| c.max_concurrent) .unwrap_or(self.max_concurrent), - "github_issues" => self.scm.github.max_concurrent.unwrap_or(self.max_concurrent), + "github_issues" => self + .scm + .github + .max_concurrent + .unwrap_or(self.max_concurrent), "helpscout" => self .issues .helpscout diff --git a/crates/claudear-engine/src/watcher.rs b/crates/claudear-engine/src/watcher.rs index 03dbb634..f0376d96 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -7380,12 +7380,30 @@ mod tests { let notifier = Arc::new(MockNotifier::new(true)); let tracker = Arc::new(SqliteTracker::in_memory().unwrap()); let issues = vec![ - Issue::new("1", "CEIL-1", "Ceiling issue 1", "http://example.com/1", "mock"), - Issue::new("2", "CEIL-2", "Ceiling issue 2", "http://example.com/2", "mock"), - Issue::new("3", "CEIL-3", "Ceiling issue 3", "http://example.com/3", "mock"), + Issue::new( + "1", + "CEIL-1", + "Ceiling issue 1", + "http://example.com/1", + "mock", + ), + Issue::new( + "2", + "CEIL-2", + "Ceiling issue 2", + "http://example.com/2", + "mock", + ), + Issue::new( + "3", + "CEIL-3", + "Ceiling issue 3", + "http://example.com/3", + "mock", + ), ]; - let source = Arc::new(MockSource::with_issues("mock", issues.clone())) - as Arc; + let source = + Arc::new(MockSource::with_issues("mock", issues.clone())) as Arc; let mut config = test_config(); // Global ceiling of 2; a generous per-source budget so the per-source