From 6831bb4d7b2bf7bf841c07de098ec023ddb26a5c Mon Sep 17 00:00:00 2001 From: Henry Park Date: Wed, 18 Mar 2026 12:29:58 -0700 Subject: [PATCH] fix: full_job routine concurrency tracks linked job lifetime (#1372) * fix: add FullJobWatcher to track full_job lifecycle for concurrency (#1318) full_job routines previously bypassed max_concurrent and global concurrency limits because execute_full_job() returned RunStatus::Ok immediately after dispatch. This meant running_count was decremented and the routine_run row was finalized before the actual job completed. Introduce FullJobWatcher struct that polls store.get_job() every 5s until the linked job reaches a non-active state, then maps the final JobState to RunStatus. execute_full_job now creates and awaits the watcher, keeping both the DB-level running row and the in-memory running_count elevated for the full job duration. Co-Authored-By: Claude Opus 4.6 (1M context) * test: full_job concurrency regression tests (issue #1318) Add two integration tests verifying full_job routine concurrency: 1. full_job_max_concurrent_blocks_second_fire_while_first_active: Inserts a Running routine_run (simulating an in-flight full_job) and verifies fire_manual returns MaxConcurrent error for max_concurrent=1. 2. global_concurrency_counts_live_full_job_runs: Elevates running_count to simulate a live full_job holding the global slot, verifies check_cron_triggers skips due routines, then releases the slot and verifies the routine fires. Also makes running_count_for_test() unconditionally public so integration tests (separate crate) can access it. Co-Authored-By: Claude Opus 4.6 (1M context) * style: fmt and clippy fixes for full_job concurrency tests Co-Authored-By: Claude Opus 4.6 (1M context) * fix: address PR review feedback on FullJobWatcher - Add #[doc(hidden)] to running_count_for_test() to hide from public API - Derive MAX_POLLS from POLL_INTERVAL to keep constants coupled - Check job state before first sleep to finalize promptly for fast jobs - Update execute_full_job doc comment to reflect blocking behavior Co-Authored-By: Claude Opus 4.6 (1M context) --------- Co-authored-by: Claude Opus 4.6 (1M context) --- src/agent/routine_engine.rs | 106 +++++++++++++++-- tests/e2e_routine_heartbeat.rs | 206 +++++++++++++++++++++++++++++++++ 2 files changed, 305 insertions(+), 7 deletions(-) diff --git a/src/agent/routine_engine.rs b/src/agent/routine_engine.rs index ec8ab851..14360d85 100644 --- a/src/agent/routine_engine.rs +++ b/src/agent/routine_engine.rs @@ -88,6 +88,12 @@ impl RoutineEngine { } } + /// Expose the running count for integration tests. + #[doc(hidden)] + pub fn running_count_for_test(&self) -> &Arc { + &self.running_count + } + /// Refresh the in-memory event trigger cache from DB. pub async fn refresh_event_cache(&self) { match self.store.list_event_routines().await { @@ -508,6 +514,88 @@ impl RoutineEngine { } } +/// Watches a dispatched full_job until the linked scheduler job completes. +/// +/// Polls `store.get_job(job_id)` at a fixed interval until the job leaves +/// an active state (Pending/InProgress/Stuck). Maps the final `JobState` to +/// a `RunStatus` for the routine run. +struct FullJobWatcher { + store: Arc, + job_id: Uuid, + routine_name: String, +} + +impl FullJobWatcher { + /// Poll interval between DB checks. + const POLL_INTERVAL: Duration = Duration::from_secs(5); + /// Safety ceiling: 24 hours, derived from POLL_INTERVAL. + const MAX_POLLS: u32 = (24 * 60 * 60) / Self::POLL_INTERVAL.as_secs() as u32; + + fn new(store: Arc, job_id: Uuid, routine_name: String) -> Self { + Self { + store, + job_id, + routine_name, + } + } + + /// Block until the linked job finishes and return the mapped status + summary. + async fn wait_for_completion(&self) -> (RunStatus, Option) { + let mut polls = 0u32; + + let final_status = loop { + // Check job state before sleeping so we finalize promptly + // if the job is already done (e.g. fast-failing jobs). + match self.store.get_job(self.job_id).await { + Ok(Some(job_ctx)) => { + if !job_ctx.state.is_active() { + break Self::map_job_state(&job_ctx.state); + } + } + Ok(None) => { + tracing::warn!( + routine = %self.routine_name, + job_id = %self.job_id, + "full_job disappeared from DB while polling" + ); + break RunStatus::Failed; + } + Err(e) => { + tracing::error!( + routine = %self.routine_name, + job_id = %self.job_id, + "Error polling full_job state: {}", e + ); + break RunStatus::Failed; + } + } + + polls += 1; + if polls >= Self::MAX_POLLS { + tracing::error!( + routine = %self.routine_name, + job_id = %self.job_id, + "full_job timed out after 24 hours, treating as failed" + ); + break RunStatus::Failed; + } + + tokio::time::sleep(Self::POLL_INTERVAL).await; + }; + + let summary = format!("Job {} finished ({})", self.job_id, final_status); + (final_status, Some(summary)) + } + + fn map_job_state(state: &crate::context::JobState) -> RunStatus { + use crate::context::JobState; + match state { + JobState::Failed | JobState::Cancelled => RunStatus::Failed, + _ => RunStatus::Ok, // Completed / Submitted / Accepted + } + } +} + /// Shared context passed to the execution function. struct EngineContext { config: RoutineConfig, @@ -682,8 +770,10 @@ fn sanitize_routine_name(name: &str) -> String { /// /// Fire-and-forget: creates a job via `Scheduler::dispatch_job` (which handles /// creation, metadata, persistence, and scheduling), links the routine run to -/// the job, and returns immediately. The job runs independently via the -/// existing Worker/Scheduler with full tool access. +/// the job, then watches it via `FullJobWatcher` until it reaches a +/// non-active state (not Pending/InProgress/Stuck). Returns the final +/// `RunStatus` mapped from the job outcome. This keeps the routine run +/// active for the full job lifetime so concurrency guardrails apply. async fn execute_full_job( ctx: &EngineContext, routine: &Routine, @@ -738,13 +828,15 @@ async fn execute_full_job( routine = %routine.name, job_id = %job_id, max_iterations = max_iterations, - "Dispatched full job for routine" + "Dispatched full job for routine, watching for completion" ); - let summary = format!( - "Dispatched job {job_id} for full execution with tool access (max_iterations: {max_iterations})" - ); - Ok((RunStatus::Ok, Some(summary), None)) + // Watch the job until it finishes — keeps the routine run active + // so concurrency guardrails (running_count, routine_runs status) + // remain enforced for the full job lifetime. + let watcher = FullJobWatcher::new(ctx.store.clone(), job_id, routine.name.clone()); + let (status, summary) = watcher.wait_for_completion().await; + Ok((status, summary, None)) } /// Execute a lightweight routine with optional tool support. diff --git a/tests/e2e_routine_heartbeat.rs b/tests/e2e_routine_heartbeat.rs index 3388feb8..25432f3d 100644 --- a/tests/e2e_routine_heartbeat.rs +++ b/tests/e2e_routine_heartbeat.rs @@ -823,4 +823,210 @@ mod tests { "Deleted routine must not fire after cache refresh" ); } + + // ----------------------------------------------------------------------- + // Test: full_job per-routine concurrency blocks second fire (issue #1318) + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn full_job_max_concurrent_blocks_second_fire_while_first_active() { + use ironclaw::agent::routine::{ + NotifyConfig, Routine, RoutineAction, RoutineGuardrails, RoutineRun, RunStatus, Trigger, + }; + use ironclaw::error::RoutineError; + + let (db, _tmp) = create_test_db().await; + let ws = create_workspace(&db); + + // Stub LLM — fire_manual will be rejected before any LLM call + let trace = LlmTrace::single_turn( + "stub", + "stub", + vec![TraceStep { + request_hint: None, + response: TraceResponse::Text { + content: "ROUTINE_OK".to_string(), + input_tokens: 10, + output_tokens: 5, + }, + expected_tool_results: vec![], + }], + ); + let llm = Arc::new(TraceLlm::from_trace(trace)); + let (notify_tx, _notify_rx) = tokio::sync::mpsc::channel(4); + let tools = Arc::new(ToolRegistry::new()); + let safety = Arc::new(SafetyLayer::new(&SafetyConfig { + max_output_length: 100_000, + injection_check_enabled: false, + })); + + let engine = Arc::new(RoutineEngine::new( + RoutineConfig::default(), + db.clone(), + llm, + ws, + notify_tx, + None, // no scheduler — rejected before dispatch + tools, + safety, + )); + + // Create a full_job routine with max_concurrent = 1 + let routine = Routine { + id: Uuid::new_v4(), + name: "concurrent-guard".to_string(), + description: "test max_concurrent for full_job".to_string(), + user_id: "default".to_string(), + enabled: true, + trigger: Trigger::Manual, + action: RoutineAction::FullJob { + title: "t".to_string(), + description: "d".to_string(), + max_iterations: 3, + tool_permissions: vec![], + }, + guardrails: RoutineGuardrails { + cooldown: Duration::from_secs(0), + max_concurrent: 1, + dedup_window: None, + }, + notify: NotifyConfig::default(), + last_run_at: None, + next_fire_at: None, + run_count: 0, + consecutive_failures: 0, + state: serde_json::json!({}), + created_at: Utc::now(), + updated_at: Utc::now(), + }; + db.create_routine(&routine).await.expect("create_routine"); + + // Simulate first full_job run still active: the fix keeps the + // routine_run in Running state while the linked job executes. + let active_run = RoutineRun { + id: Uuid::new_v4(), + routine_id: routine.id, + trigger_type: "cron".to_string(), + trigger_detail: None, + started_at: Utc::now(), + completed_at: None, + status: RunStatus::Running, + result_summary: None, + tokens_used: None, + job_id: None, + created_at: Utc::now(), + }; + db.create_routine_run(&active_run) + .await + .expect("create_routine_run"); + + // Attempt to fire the same routine again — must be rejected + let result = engine.fire_manual(routine.id, None).await; + assert!( + matches!(result, Err(RoutineError::MaxConcurrent { .. })), + "second fire while first full_job active must be rejected by max_concurrent=1, got: {:?}", + result + ); + } + + // ----------------------------------------------------------------------- + // Test: global running_count tracks live full_job runs (issue #1318) + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn global_concurrency_counts_live_full_job_runs() { + use std::sync::atomic::Ordering; + + let (db, _tmp) = create_test_db().await; + let ws = create_workspace(&db); + + let trace = LlmTrace::single_turn( + "test-global-limit", + "check", + vec![TraceStep { + request_hint: None, + response: TraceResponse::Text { + content: "ROUTINE_OK".to_string(), + input_tokens: 50, + output_tokens: 5, + }, + expected_tool_results: vec![], + }], + ); + let llm = Arc::new(TraceLlm::from_trace(trace)); + let (notify_tx, _notify_rx) = tokio::sync::mpsc::channel(16); + let tools = Arc::new(ToolRegistry::new()); + let safety = Arc::new(SafetyLayer::new(&SafetyConfig { + max_output_length: 100_000, + injection_check_enabled: true, + })); + + // Configure global limit of 1 + let config = RoutineConfig { + max_concurrent_routines: 1, + ..RoutineConfig::default() + }; + + let engine = Arc::new(RoutineEngine::new( + config, + db.clone(), + llm, + ws, + notify_tx, + None, + tools, + safety, + )); + + // Insert a due cron routine + let mut routine = make_routine( + "global-limit-test", + Trigger::Cron { + schedule: "* * * * *".to_string(), + timezone: None, + }, + "Check status.", + ); + routine.next_fire_at = Some(Utc::now() - chrono::Duration::minutes(1)); + db.create_routine(&routine).await.expect("create_routine"); + + // Simulate one full_job from another routine holding the global slot. + // With the fix, running_count stays elevated for the full job duration. + engine + .running_count_for_test() + .fetch_add(1, Ordering::Relaxed); + + // check_cron_triggers should see global limit hit and skip + engine.check_cron_triggers().await; + tokio::time::sleep(Duration::from_millis(100)).await; + + let runs = db + .list_routine_runs(routine.id, 10) + .await + .expect("list_routine_runs"); + assert!( + runs.is_empty(), + "cron routine must not fire when global limit is reached by live full_job" + ); + + // Release the global slot + engine + .running_count_for_test() + .fetch_sub(1, Ordering::Relaxed); + + // Now the routine should fire + engine.check_cron_triggers().await; + tokio::time::sleep(Duration::from_millis(200)).await; + + // Because the first check skipped it, next_fire_at is unchanged — + // the second check should see it as still due and fire it. + let runs_after = db + .list_routine_runs(routine.id, 10) + .await + .expect("list_routine_runs"); + assert!( + !runs_after.is_empty(), + "cron routine should fire after global slot is released" + ); + } }