mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
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) <[email protected]> * 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) <[email protected]> * style: fmt and clippy fixes for full_job concurrency tests Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * 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) <[email protected]> --------- Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
42ffefabe4
commit
6831bb4d7b
@@ -88,6 +88,12 @@ impl RoutineEngine {
|
||||
}
|
||||
}
|
||||
|
||||
/// Expose the running count for integration tests.
|
||||
#[doc(hidden)]
|
||||
pub fn running_count_for_test(&self) -> &Arc<AtomicUsize> {
|
||||
&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<dyn Database>,
|
||||
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<dyn Database>, 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<String>) {
|
||||
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.
|
||||
|
||||
@@ -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"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user