mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-27 08:00:17 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d531adaf18 | ||
|
|
302fa8a38d | ||
|
|
1b9a8ad1b3 | ||
|
|
8c1553e2c9 |
+209
-2
@@ -25,6 +25,7 @@ use crate::agent::routine::{
|
|||||||
};
|
};
|
||||||
use crate::channels::{IncomingMessage, OutgoingResponse};
|
use crate::channels::{IncomingMessage, OutgoingResponse};
|
||||||
use crate::config::RoutineConfig;
|
use crate::config::RoutineConfig;
|
||||||
|
use crate::context::JobState;
|
||||||
use crate::db::Database;
|
use crate::db::Database;
|
||||||
use crate::error::RoutineError;
|
use crate::error::RoutineError;
|
||||||
use crate::llm::{ChatMessage, CompletionRequest, FinishReason, LlmProvider};
|
use crate::llm::{ChatMessage, CompletionRequest, FinishReason, LlmProvider};
|
||||||
@@ -180,6 +181,130 @@ impl RoutineEngine {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Sync dispatched routine runs with their linked background job status.
|
||||||
|
///
|
||||||
|
/// Full-job routines are fire-and-forget: the routine run is created with
|
||||||
|
/// `Running` status when the job is dispatched, but the run record is never
|
||||||
|
/// updated when the background job completes or fails. This method checks
|
||||||
|
/// all `Running` routine runs that have a linked job, queries the job's
|
||||||
|
/// current state, and updates the routine run accordingly. It also sends
|
||||||
|
/// failure/success notifications that would otherwise be lost.
|
||||||
|
pub async fn sync_dispatched_runs(&self) {
|
||||||
|
let runs = match self.store.list_dispatched_routine_runs().await {
|
||||||
|
Ok(r) => r,
|
||||||
|
Err(e) => {
|
||||||
|
tracing::debug!("Failed to list dispatched routine runs: {}", e);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
for run in runs {
|
||||||
|
let Some(job_id) = run.job_id else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
|
||||||
|
// Check the linked job's current state
|
||||||
|
let job = match self.store.get_job(job_id).await {
|
||||||
|
Ok(Some(j)) => j,
|
||||||
|
Ok(None) => {
|
||||||
|
// Job was deleted — mark the routine run as failed
|
||||||
|
tracing::warn!(
|
||||||
|
run_id = %run.id,
|
||||||
|
job_id = %job_id,
|
||||||
|
"Linked job not found, marking routine run as failed"
|
||||||
|
);
|
||||||
|
self.complete_dispatched_run(
|
||||||
|
&run,
|
||||||
|
RunStatus::Failed,
|
||||||
|
"Linked job not found (may have been deleted)",
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
tracing::debug!(
|
||||||
|
run_id = %run.id,
|
||||||
|
job_id = %job_id,
|
||||||
|
"Failed to query linked job: {}", e
|
||||||
|
);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
// Extract the reason from the most recent state transition
|
||||||
|
let last_reason = job.transitions.last().and_then(|t| t.reason.clone());
|
||||||
|
|
||||||
|
// Map job state to routine run status
|
||||||
|
let (new_status, summary) = match job.state {
|
||||||
|
JobState::Completed | JobState::Submitted | JobState::Accepted => {
|
||||||
|
let summary =
|
||||||
|
last_reason.unwrap_or_else(|| "Job completed successfully".to_string());
|
||||||
|
(RunStatus::Ok, summary)
|
||||||
|
}
|
||||||
|
JobState::Failed => {
|
||||||
|
let summary = last_reason
|
||||||
|
.unwrap_or_else(|| "Job failed (no error message recorded)".to_string());
|
||||||
|
(RunStatus::Failed, summary)
|
||||||
|
}
|
||||||
|
JobState::Cancelled => (RunStatus::Failed, "Job was cancelled".to_string()),
|
||||||
|
// Still in progress — skip
|
||||||
|
JobState::Pending | JobState::InProgress | JobState::Stuck => continue,
|
||||||
|
};
|
||||||
|
|
||||||
|
tracing::info!(
|
||||||
|
run_id = %run.id,
|
||||||
|
job_id = %job_id,
|
||||||
|
status = %new_status,
|
||||||
|
"Syncing dispatched routine run with completed job"
|
||||||
|
);
|
||||||
|
|
||||||
|
self.complete_dispatched_run(&run, new_status, &summary)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Complete a dispatched routine run and send the appropriate notification.
|
||||||
|
async fn complete_dispatched_run(&self, run: &RoutineRun, status: RunStatus, summary: &str) {
|
||||||
|
if let Err(e) = self
|
||||||
|
.store
|
||||||
|
.complete_routine_run(run.id, status, Some(summary), None)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
tracing::error!(
|
||||||
|
run_id = %run.id,
|
||||||
|
"Failed to update dispatched routine run: {}", e
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Look up the routine to get its notify config and name
|
||||||
|
match self.store.get_routine(run.routine_id).await {
|
||||||
|
Ok(Some(routine)) => {
|
||||||
|
send_notification(
|
||||||
|
&self.notify_tx,
|
||||||
|
&routine.notify,
|
||||||
|
&routine.name,
|
||||||
|
status,
|
||||||
|
Some(summary),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
Ok(None) => {
|
||||||
|
tracing::debug!(
|
||||||
|
routine_id = %run.routine_id,
|
||||||
|
"Routine not found for notification (may have been deleted)"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
tracing::debug!(
|
||||||
|
routine_id = %run.routine_id,
|
||||||
|
"Failed to look up routine for notification: {}", e
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Fire a routine manually (from tool call or CLI).
|
/// Fire a routine manually (from tool call or CLI).
|
||||||
///
|
///
|
||||||
/// Bypasses cooldown checks (those only apply to cron/event triggers).
|
/// Bypasses cooldown checks (those only apply to cron/event triggers).
|
||||||
@@ -534,9 +659,10 @@ async fn execute_full_job(
|
|||||||
);
|
);
|
||||||
|
|
||||||
let summary = format!(
|
let summary = format!(
|
||||||
"Dispatched job {job_id} for full execution with tool access (max_iterations: {max_iterations})"
|
"Dispatched job {job_id} for full execution with tool access (max_iterations: {max_iterations}). \
|
||||||
|
Status will be updated when the job completes."
|
||||||
);
|
);
|
||||||
Ok((RunStatus::Ok, Some(summary), None))
|
Ok((RunStatus::Running, Some(summary), None))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Execute a lightweight routine (single LLM call).
|
/// Execute a lightweight routine (single LLM call).
|
||||||
@@ -712,6 +838,7 @@ pub fn spawn_cron_ticker(
|
|||||||
loop {
|
loop {
|
||||||
ticker.tick().await;
|
ticker.tick().await;
|
||||||
engine.check_cron_triggers().await;
|
engine.check_cron_triggers().await;
|
||||||
|
engine.sync_dispatched_runs().await;
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -756,4 +883,84 @@ mod tests {
|
|||||||
let _ = status.to_string();
|
let _ = status.to_string();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_running_status_does_not_notify() {
|
||||||
|
// Running status should not trigger notifications (job still in progress)
|
||||||
|
let config = NotifyConfig {
|
||||||
|
on_success: true,
|
||||||
|
on_failure: true,
|
||||||
|
on_attention: true,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
// RunStatus::Running maps to false in send_notification's match
|
||||||
|
let should_notify = match RunStatus::Running {
|
||||||
|
RunStatus::Ok => config.on_success,
|
||||||
|
RunStatus::Attention => config.on_attention,
|
||||||
|
RunStatus::Failed => config.on_failure,
|
||||||
|
RunStatus::Running => false,
|
||||||
|
};
|
||||||
|
assert!(!should_notify);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_full_job_dispatch_returns_running_status() {
|
||||||
|
// Verify the status text for Running is "running"
|
||||||
|
assert_eq!(RunStatus::Running.to_string(), "running");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Regression test for #697: full_job routines were immediately marked Ok
|
||||||
|
/// on dispatch, so failures/completions were never synced back. The fix
|
||||||
|
/// changed dispatch to return Running and added sync_dispatched_runs which
|
||||||
|
/// maps terminal job states to routine run statuses.
|
||||||
|
#[test]
|
||||||
|
fn test_job_state_to_run_status_mapping() {
|
||||||
|
use crate::context::JobState;
|
||||||
|
|
||||||
|
// Helper that replicates the mapping logic from sync_dispatched_runs
|
||||||
|
let map_state = |state: JobState, reason: Option<&str>| -> Option<(RunStatus, String)> {
|
||||||
|
let last_reason = reason.map(|s| s.to_string());
|
||||||
|
match state {
|
||||||
|
JobState::Completed | JobState::Submitted | JobState::Accepted => {
|
||||||
|
let summary =
|
||||||
|
last_reason.unwrap_or_else(|| "Job completed successfully".to_string());
|
||||||
|
Some((RunStatus::Ok, summary))
|
||||||
|
}
|
||||||
|
JobState::Failed => {
|
||||||
|
let summary = last_reason
|
||||||
|
.unwrap_or_else(|| "Job failed (no error message recorded)".to_string());
|
||||||
|
Some((RunStatus::Failed, summary))
|
||||||
|
}
|
||||||
|
JobState::Cancelled => Some((RunStatus::Failed, "Job was cancelled".to_string())),
|
||||||
|
JobState::Pending | JobState::InProgress | JobState::Stuck => None,
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
// Terminal states produce a status update
|
||||||
|
let (status, _) = map_state(JobState::Completed, None).unwrap();
|
||||||
|
assert_eq!(status, RunStatus::Ok);
|
||||||
|
|
||||||
|
let (status, _) = map_state(JobState::Submitted, None).unwrap();
|
||||||
|
assert_eq!(status, RunStatus::Ok);
|
||||||
|
|
||||||
|
let (status, _) = map_state(JobState::Accepted, None).unwrap();
|
||||||
|
assert_eq!(status, RunStatus::Ok);
|
||||||
|
|
||||||
|
let (status, summary) = map_state(JobState::Failed, Some("OOM killed")).unwrap();
|
||||||
|
assert_eq!(status, RunStatus::Failed);
|
||||||
|
assert_eq!(summary, "OOM killed");
|
||||||
|
|
||||||
|
let (status, summary) = map_state(JobState::Failed, None).unwrap();
|
||||||
|
assert_eq!(status, RunStatus::Failed);
|
||||||
|
assert!(summary.contains("no error message"));
|
||||||
|
|
||||||
|
let (status, _) = map_state(JobState::Cancelled, None).unwrap();
|
||||||
|
assert_eq!(status, RunStatus::Failed);
|
||||||
|
|
||||||
|
// In-progress states should NOT produce a status update (skip)
|
||||||
|
assert!(map_state(JobState::Pending, None).is_none());
|
||||||
|
assert!(map_state(JobState::InProgress, None).is_none());
|
||||||
|
assert!(map_state(JobState::Stuck, None).is_none());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -423,4 +423,29 @@ impl RoutineStore for LibSqlBackend {
|
|||||||
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn list_dispatched_routine_runs(&self) -> Result<Vec<RoutineRun>, DatabaseError> {
|
||||||
|
let conn = self.connect().await?;
|
||||||
|
let mut rows = conn
|
||||||
|
.query(
|
||||||
|
&format!(
|
||||||
|
"SELECT {} FROM routine_runs \
|
||||||
|
WHERE status = 'running' AND job_id IS NOT NULL",
|
||||||
|
ROUTINE_RUN_COLUMNS
|
||||||
|
),
|
||||||
|
params![],
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
||||||
|
|
||||||
|
let mut runs = Vec::new();
|
||||||
|
while let Some(row) = rows
|
||||||
|
.next()
|
||||||
|
.await
|
||||||
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
||||||
|
{
|
||||||
|
runs.push(row_to_routine_run_libsql(&row)?);
|
||||||
|
}
|
||||||
|
Ok(runs)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -303,6 +303,10 @@ pub trait RoutineStore: Send + Sync {
|
|||||||
run_id: Uuid,
|
run_id: Uuid,
|
||||||
job_id: Uuid,
|
job_id: Uuid,
|
||||||
) -> Result<(), DatabaseError>;
|
) -> Result<(), DatabaseError>;
|
||||||
|
/// List routine runs that were dispatched as full_job (status = 'running'
|
||||||
|
/// with a linked job_id). Used by the routine engine to sync completion
|
||||||
|
/// status from the background job.
|
||||||
|
async fn list_dispatched_routine_runs(&self) -> Result<Vec<RoutineRun>, DatabaseError>;
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
|
|||||||
@@ -494,6 +494,10 @@ impl RoutineStore for PgBackend {
|
|||||||
) -> Result<(), DatabaseError> {
|
) -> Result<(), DatabaseError> {
|
||||||
self.store.link_routine_run_to_job(run_id, job_id).await
|
self.store.link_routine_run_to_job(run_id, job_id).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn list_dispatched_routine_runs(&self) -> Result<Vec<RoutineRun>, DatabaseError> {
|
||||||
|
self.store.list_dispatched_routine_runs().await
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ==================== ToolFailureStore ====================
|
// ==================== ToolFailureStore ====================
|
||||||
|
|||||||
@@ -1295,6 +1295,17 @@ impl Store {
|
|||||||
.await?;
|
.await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn list_dispatched_routine_runs(&self) -> Result<Vec<RoutineRun>, DatabaseError> {
|
||||||
|
let conn = self.conn().await?;
|
||||||
|
let rows = conn
|
||||||
|
.query(
|
||||||
|
"SELECT * FROM routine_runs WHERE status = 'running' AND job_id IS NOT NULL",
|
||||||
|
&[],
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
rows.iter().map(row_to_routine_run).collect()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "postgres")]
|
#[cfg(feature = "postgres")]
|
||||||
|
|||||||
Reference in New Issue
Block a user