diff --git a/claudear.example.toml b/claudear.example.toml index 7f107b85..ff2c2814 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 @@ -301,6 +310,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 = "" @@ -475,7 +488,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 @@ -1065,8 +1080,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 9c4332cd..c2e9333a 100644 --- a/crates/claudear-config/src/config.rs +++ b/crates/claudear-config/src/config.rs @@ -447,8 +447,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). @@ -793,6 +799,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, @@ -900,7 +907,7 @@ impl Default for QaConfig { answer_timeout_secs: 600, max_qa_per_cycle: 20, use_llm: false, - max_concurrent: 1, + max_concurrent: 6, } } } @@ -1740,6 +1747,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 { @@ -1756,6 +1767,7 @@ impl Default for GitHubConfig { trigger_labels: Vec::new(), trigger_states: Vec::new(), app: GitHubAppConfig::default(), + max_concurrent: Some(6), } } } @@ -1775,6 +1787,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), } } } @@ -2053,7 +2066,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, } } @@ -3808,6 +3823,17 @@ 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, } } @@ -4852,6 +4878,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() @@ -4863,6 +4891,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 2a1a860b..d2193746 100644 --- a/crates/claudear-engine/src/api/routes.rs +++ b/crates/claudear-engine/src/api/routes.rs @@ -3055,6 +3055,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/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 f6fcbaec..f0376d96 100644 --- a/crates/claudear-engine/src/watcher.rs +++ b/crates/claudear-engine/src/watcher.rs @@ -38,8 +38,8 @@ 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::time::{interval_at, Duration}; +use tokio::sync::{Notify, RwLock, Semaphore}; +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). @@ -413,6 +413,14 @@ 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. + /// 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`: @@ -518,6 +526,8 @@ impl Watcher { )) }; + let session_limit = options.config.max_concurrent_sessions.max(1); + Self { agent: options.agent, qa_agent: options.qa_agent, @@ -546,6 +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: Arc::new(Semaphore::new(session_limit)), llm_analyzer, intent_classifier, spawn_handles: tokio::sync::Mutex::new(Vec::new()), @@ -1260,6 +1271,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()); @@ -1323,34 +1338,62 @@ impl Watcher { /// Run the source polling loop until the watcher stops. /// - /// Polls each source at its configured interval, logging failed polls. + /// 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`](crate::housekeeping::HousekeepingWorker). async fn run_source_poll_loop(self: &Arc, poll_interval: u64) { - let now = std::time::Instant::now(); - let mut schedule: Vec<_> = self - .sources - .iter() - .map(|source| { - let interval_ms = self.config.poll_interval_ms_for(source.name()).max(1); - tracing::info!( - source = source.name(), - interval_ms, - "Per-source poll interval" - ); - (source, Duration::from_millis(interval_ms), now) - }) - .collect(); + let min_period_ms = MIN_SOURCE_POLL_PERIOD.as_millis() as u64; + 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(min_period_ms); + 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; + })); + } - let poll_interval = Duration::from_millis(poll_interval); - let period = schedule - .iter() - .map(|(_, interval, _)| *interval) - .min() - .unwrap_or(poll_interval) - .min(poll_interval) - .max(MIN_SOURCE_POLL_PERIOD); - let mut timer = interval_at(tokio::time::Instant::now() + period, period); + // Nothing configured: idle until stopped so the caller's join! still + // has a future to hold. + if workers.is_empty() { + let _ = poll_interval; + while self.is_running.load(Ordering::SeqCst) { + tokio::time::sleep(MIN_SOURCE_POLL_PERIOD).await; + } + return; + } + + for worker in workers { + let _ = worker.await; + } + } + + /// 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`]). + /// + /// The timer ticks on a fixed [`MIN_SOURCE_POLL_PERIOD`] granularity so a + /// stop is noticed promptly even when `interval_ms` is large; the source is + /// actually polled only once its interval has elapsed, and the clock is + /// reset *after* the poll completes so a slow poll (slow API, waiting for + /// dispatch slots) never triggers a burst of catch-up polls. + async fn run_source_worker(self: Arc, source_idx: usize, interval_ms: u64) { + let source_interval = Duration::from_millis(interval_ms); + let mut timer = interval(MIN_SOURCE_POLL_PERIOD); + timer.set_missed_tick_behavior(MissedTickBehavior::Delay); + timer.tick().await; // consume the immediate first tick + let mut last_poll = std::time::Instant::now(); while self.is_running.load(Ordering::SeqCst) { timer.tick().await; @@ -1360,20 +1403,19 @@ impl Watcher { if self.is_rate_limit_paused().await { continue; } - - for (source, interval, last_poll) in &mut schedule { - if last_poll.elapsed() >= *interval { - if let Err(error) = self.poll_source(source).await { - tracing::error!( - component = "watcher", - source = source.name(), - error = %error, - "Error polling source" - ); - } - *last_poll = std::time::Instant::now(); - } + if last_poll.elapsed() < source_interval { + continue; + } + 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" + ); } + last_poll = std::time::Instant::now(); } } @@ -3872,6 +3914,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 @@ -3887,12 +3936,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 @@ -4730,6 +4784,51 @@ Create a PR with your changes.{custom_instructions}"#, return IssueRun::Skipped; } + // 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 IssueRun::Skipped; + } + }; + + // 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 + // 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 IssueRun::Paused; + } + if let Some(ref tip) = deploy_qa_tip { if !self.claim_deploy_qa_tip(tip, &issue.short_id) { return IssueRun::Skipped; @@ -4949,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 { @@ -5730,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 { @@ -6133,6 +6234,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, @@ -6241,6 +6343,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)); @@ -7223,6 +7375,255 @@ mod tests { watcher.active_processing.fetch_sub(1, Ordering::SeqCst); } + #[tokio::test] + 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 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(); + // 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 { + 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 both session permits, standing + // in for two runs already in flight. + let held_a = watcher + .session_limiter + .try_acquire() + .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 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(); + + tokio::time::sleep(Duration::from_millis(100)).await; + assert!( + 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 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(), + 2, + "completing every run must free its session slot" + ); + } + + /// 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"); 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 19a0ceed..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 { @@ -1374,6 +1375,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,