mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-28 16:39:31 +00:00
Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d531adaf18 | ||
|
|
302fa8a38d | ||
|
|
1b9a8ad1b3 | ||
|
|
8c1553e2c9 |
Generated
-29
@@ -864,16 +864,6 @@ dependencies = [
|
||||
"windows-link",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "chrono-tz"
|
||||
version = "0.10.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a6139a8597ed92cf816dfb33f5dd6cf0bb93a6adc938f11039f371bc5bcd26c3"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"phf 0.12.1",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "cipher"
|
||||
version = "0.4.4"
|
||||
@@ -2882,7 +2872,6 @@ dependencies = [
|
||||
"bollard",
|
||||
"bytes",
|
||||
"chrono",
|
||||
"chrono-tz",
|
||||
"clap",
|
||||
"clap_complete",
|
||||
"cron",
|
||||
@@ -3903,15 +3892,6 @@ dependencies = [
|
||||
"phf_shared 0.11.3",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "phf"
|
||||
version = "0.12.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "913273894cec178f401a31ec4b656318d95473527be05c0752cc41cdc32be8b7"
|
||||
dependencies = [
|
||||
"phf_shared 0.12.1",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "phf"
|
||||
version = "0.13.1"
|
||||
@@ -3986,15 +3966,6 @@ dependencies = [
|
||||
"uncased",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "phf_shared"
|
||||
version = "0.12.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "06005508882fb681fd97892ecff4b7fd0fee13ef1aa569f8695dae7ab9099981"
|
||||
dependencies = [
|
||||
"siphasher",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "phf_shared"
|
||||
version = "0.13.1"
|
||||
|
||||
@@ -73,7 +73,6 @@ toml = "0.8"
|
||||
# Core types
|
||||
uuid = { version = "1", features = ["v4", "v5", "serde"] }
|
||||
chrono = { version = "0.4", features = ["serde"] }
|
||||
chrono-tz = "0.10"
|
||||
rust_decimal = { version = "1", features = ["serde", "serde-with-str", "maths"] }
|
||||
rust_decimal_macros = "1"
|
||||
|
||||
|
||||
+209
-2
@@ -25,6 +25,7 @@ use crate::agent::routine::{
|
||||
};
|
||||
use crate::channels::{IncomingMessage, OutgoingResponse};
|
||||
use crate::config::RoutineConfig;
|
||||
use crate::context::JobState;
|
||||
use crate::db::Database;
|
||||
use crate::error::RoutineError;
|
||||
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).
|
||||
///
|
||||
/// Bypasses cooldown checks (those only apply to cron/event triggers).
|
||||
@@ -534,9 +659,10 @@ async fn execute_full_job(
|
||||
);
|
||||
|
||||
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).
|
||||
@@ -712,6 +838,7 @@ pub fn spawn_cron_ticker(
|
||||
loop {
|
||||
ticker.tick().await;
|
||||
engine.check_cron_triggers().await;
|
||||
engine.sync_dispatched_runs().await;
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -756,4 +883,84 @@ mod tests {
|
||||
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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,4 +26,3 @@ pub mod routines;
|
||||
pub mod settings;
|
||||
#[allow(dead_code)]
|
||||
pub mod static_files;
|
||||
pub mod webhooks;
|
||||
|
||||
@@ -1,210 +0,0 @@
|
||||
//! Public webhook trigger endpoint for routine webhook triggers.
|
||||
//!
|
||||
//! `POST /api/webhooks/{path}` — matches the path against routines with
|
||||
//! `Trigger::Webhook { path, secret }`, validates the secret via constant-time
|
||||
//! comparison, and fires the matching routine through the message pipeline.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::{
|
||||
Json,
|
||||
extract::{Path, State},
|
||||
http::{HeaderMap, StatusCode},
|
||||
};
|
||||
use subtle::ConstantTimeEq;
|
||||
|
||||
use crate::agent::routine::{RoutineAction, Trigger};
|
||||
use crate::channels::IncomingMessage;
|
||||
use crate::channels::web::server::GatewayState;
|
||||
|
||||
/// Handle incoming webhook POST to `/api/webhooks/{path}`.
|
||||
///
|
||||
/// This endpoint is **public** (no gateway auth token required) but protected
|
||||
/// by the per-routine webhook secret sent via the `X-Webhook-Secret` header.
|
||||
pub async fn webhook_trigger_handler(
|
||||
State(state): State<Arc<GatewayState>>,
|
||||
Path(path): Path<String>,
|
||||
headers: HeaderMap,
|
||||
) -> Result<Json<serde_json::Value>, (StatusCode, String)> {
|
||||
let store = state.store.as_ref().ok_or((
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
"Database not available".to_string(),
|
||||
))?;
|
||||
|
||||
// Load all routines and find one whose Trigger::Webhook path matches.
|
||||
let routines = store
|
||||
.list_all_routines()
|
||||
.await
|
||||
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
|
||||
|
||||
let matched = routines.into_iter().find(|r| {
|
||||
if !r.enabled {
|
||||
return false;
|
||||
}
|
||||
match &r.trigger {
|
||||
Trigger::Webhook { path: Some(wp), .. } => *wp == path,
|
||||
Trigger::Webhook { path: None, .. } => path == r.id.to_string(),
|
||||
_ => false,
|
||||
}
|
||||
});
|
||||
|
||||
let routine = matched.ok_or((
|
||||
StatusCode::NOT_FOUND,
|
||||
"No routine matches this webhook path".to_string(),
|
||||
))?;
|
||||
|
||||
// Validate the webhook secret if one is configured on the routine.
|
||||
if let Trigger::Webhook {
|
||||
secret: Some(expected_secret),
|
||||
..
|
||||
} = &routine.trigger
|
||||
{
|
||||
let provided_secret = headers
|
||||
.get("x-webhook-secret")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.unwrap_or("");
|
||||
|
||||
if !bool::from(provided_secret.as_bytes().ct_eq(expected_secret.as_bytes())) {
|
||||
return Err((
|
||||
StatusCode::UNAUTHORIZED,
|
||||
"Invalid webhook secret".to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
// Build the prompt from the routine action.
|
||||
let prompt = match &routine.action {
|
||||
RoutineAction::Lightweight { prompt, .. } => prompt.clone(),
|
||||
RoutineAction::FullJob {
|
||||
title, description, ..
|
||||
} => format!("{}: {}", title, description),
|
||||
};
|
||||
|
||||
let content = format!("[routine:{}] {}", routine.name, prompt);
|
||||
let thread_id = format!(
|
||||
"routine-{}-{}",
|
||||
routine.id,
|
||||
chrono::Utc::now().timestamp_millis()
|
||||
);
|
||||
let msg = IncomingMessage::new("gateway", &routine.user_id, content).with_thread(thread_id);
|
||||
|
||||
let tx_guard = state.msg_tx.read().await;
|
||||
let tx = tx_guard.as_ref().ok_or((
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
"Channel not started".to_string(),
|
||||
))?;
|
||||
|
||||
tx.send(msg).await.map_err(|_| {
|
||||
(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
"Channel closed".to_string(),
|
||||
)
|
||||
})?;
|
||||
|
||||
Ok(Json(serde_json::json!({
|
||||
"status": "triggered",
|
||||
"routine_id": routine.id,
|
||||
"routine_name": routine.name,
|
||||
})))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// Verify constant-time comparison logic for webhook secrets.
|
||||
#[test]
|
||||
fn test_webhook_secret_constant_time_comparison() {
|
||||
let expected = "my-secret-token";
|
||||
|
||||
// Matching secret
|
||||
let provided = "my-secret-token";
|
||||
assert!(bool::from(provided.as_bytes().ct_eq(expected.as_bytes())));
|
||||
|
||||
// Wrong secret
|
||||
let wrong = "wrong-secret";
|
||||
assert!(!bool::from(wrong.as_bytes().ct_eq(expected.as_bytes())));
|
||||
|
||||
// Empty secret
|
||||
let empty = "";
|
||||
assert!(!bool::from(empty.as_bytes().ct_eq(expected.as_bytes())));
|
||||
}
|
||||
|
||||
/// Verify that webhook path matching logic works for both explicit paths
|
||||
/// and fallback to routine ID.
|
||||
#[test]
|
||||
fn test_webhook_path_matching() {
|
||||
use chrono::Utc;
|
||||
use uuid::Uuid;
|
||||
|
||||
let routine_id = Uuid::parse_str("550e8400-e29b-41d4-a716-446655440000").unwrap();
|
||||
|
||||
let routine = crate::agent::routine::Routine {
|
||||
id: routine_id,
|
||||
name: "test-routine".to_string(),
|
||||
description: "A test routine".to_string(),
|
||||
user_id: "test-user".to_string(),
|
||||
enabled: true,
|
||||
trigger: Trigger::Webhook {
|
||||
path: Some("my-hook".to_string()),
|
||||
secret: None,
|
||||
},
|
||||
action: RoutineAction::Lightweight {
|
||||
prompt: "do stuff".to_string(),
|
||||
context_paths: vec![],
|
||||
max_tokens: 4096,
|
||||
},
|
||||
guardrails: crate::agent::routine::RoutineGuardrails::default(),
|
||||
notify: crate::agent::routine::NotifyConfig::default(),
|
||||
last_run_at: None,
|
||||
next_fire_at: None,
|
||||
run_count: 0,
|
||||
consecutive_failures: 0,
|
||||
state: serde_json::Value::Null,
|
||||
created_at: Utc::now(),
|
||||
updated_at: Utc::now(),
|
||||
};
|
||||
|
||||
// Explicit path match
|
||||
let matches_explicit = match &routine.trigger {
|
||||
Trigger::Webhook { path: Some(wp), .. } => *wp == "my-hook",
|
||||
_ => false,
|
||||
};
|
||||
assert!(matches_explicit);
|
||||
|
||||
// Should NOT match wrong path
|
||||
let matches_wrong = match &routine.trigger {
|
||||
Trigger::Webhook { path: Some(wp), .. } => *wp == "other-hook",
|
||||
_ => false,
|
||||
};
|
||||
assert!(!matches_wrong);
|
||||
|
||||
// Routine with no explicit path falls back to ID
|
||||
let routine_no_path = crate::agent::routine::Routine {
|
||||
trigger: Trigger::Webhook {
|
||||
path: None,
|
||||
secret: None,
|
||||
},
|
||||
..routine
|
||||
};
|
||||
let matches_id = match &routine_no_path.trigger {
|
||||
Trigger::Webhook { path: None, .. } => {
|
||||
routine_no_path.id.to_string() == "550e8400-e29b-41d4-a716-446655440000"
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
assert!(matches_id);
|
||||
|
||||
// Disabled routine should not match
|
||||
let disabled_routine = crate::agent::routine::Routine {
|
||||
enabled: false,
|
||||
trigger: Trigger::Webhook {
|
||||
path: Some("my-hook".to_string()),
|
||||
secret: None,
|
||||
},
|
||||
..routine_no_path
|
||||
};
|
||||
let should_skip = !disabled_routine.enabled;
|
||||
assert!(should_skip);
|
||||
}
|
||||
}
|
||||
@@ -37,7 +37,6 @@ use crate::channels::web::handlers::jobs::{
|
||||
use crate::channels::web::handlers::skills::{
|
||||
skills_install_handler, skills_list_handler, skills_remove_handler, skills_search_handler,
|
||||
};
|
||||
use crate::channels::web::handlers::webhooks::webhook_trigger_handler;
|
||||
use crate::channels::web::log_layer::LogBroadcaster;
|
||||
use crate::channels::web::sse::SseManager;
|
||||
use crate::channels::web::types::*;
|
||||
@@ -201,8 +200,7 @@ pub async fn start_server(
|
||||
// Public routes (no auth)
|
||||
let public = Router::new()
|
||||
.route("/api/health", get(health_handler))
|
||||
.route("/oauth/callback", get(oauth_callback_handler))
|
||||
.route("/api/webhooks/{path}", post(webhook_trigger_handler));
|
||||
.route("/oauth/callback", get(oauth_callback_handler));
|
||||
|
||||
// Protected routes (require auth)
|
||||
let auth_state = AuthState { token: auth_token };
|
||||
|
||||
@@ -423,4 +423,29 @@ impl RoutineStore for LibSqlBackend {
|
||||
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
||||
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,
|
||||
job_id: Uuid,
|
||||
) -> 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]
|
||||
|
||||
@@ -494,6 +494,10 @@ impl RoutineStore for PgBackend {
|
||||
) -> Result<(), DatabaseError> {
|
||||
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 ====================
|
||||
|
||||
@@ -1295,6 +1295,17 @@ impl Store {
|
||||
.await?;
|
||||
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")]
|
||||
|
||||
+23
-345
@@ -1,53 +1,11 @@
|
||||
//! Time utility tool.
|
||||
|
||||
use async_trait::async_trait;
|
||||
use chrono::{DateTime, FixedOffset, Utc};
|
||||
use chrono_tz::Tz;
|
||||
use chrono::{DateTime, Utc};
|
||||
|
||||
use crate::context::JobContext;
|
||||
use crate::tools::tool::{Tool, ToolError, ToolOutput, require_str};
|
||||
|
||||
/// Parse a timezone string into a `chrono_tz::Tz`, returning a clear error.
|
||||
fn parse_timezone(tz_str: &str) -> Result<Tz, ToolError> {
|
||||
tz_str.parse::<Tz>().map_err(|_| {
|
||||
ToolError::InvalidParameters(format!(
|
||||
"Unknown timezone '{}'. Use IANA names like 'America/New_York' or 'Europe/London'.",
|
||||
tz_str
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
/// Parse an input timestamp string. Accepts RFC 3339 with offset, or naive
|
||||
/// datetime in `YYYY-MM-DDTHH:MM:SS` / `YYYY-MM-DD HH:MM:SS` format
|
||||
/// (interpreted as UTC unless `default_tz` is provided).
|
||||
fn parse_input_timestamp(
|
||||
input: &str,
|
||||
default_tz: Option<Tz>,
|
||||
) -> Result<DateTime<FixedOffset>, ToolError> {
|
||||
// Try RFC 3339 first (has offset info)
|
||||
if let Ok(dt) = DateTime::parse_from_rfc3339(input) {
|
||||
return Ok(dt);
|
||||
}
|
||||
// Try common formats without offset — interpret in default_tz or UTC
|
||||
for fmt in &["%Y-%m-%dT%H:%M:%S", "%Y-%m-%d %H:%M:%S"] {
|
||||
if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(input, fmt) {
|
||||
let tz = default_tz.unwrap_or(Tz::UTC);
|
||||
let local = naive.and_local_timezone(tz).single().ok_or_else(|| {
|
||||
ToolError::InvalidParameters(format!(
|
||||
"Ambiguous or invalid datetime '{}' in timezone '{}'",
|
||||
input, tz
|
||||
))
|
||||
})?;
|
||||
return Ok(local.fixed_offset());
|
||||
}
|
||||
}
|
||||
Err(ToolError::InvalidParameters(format!(
|
||||
"Invalid timestamp '{}'. Use RFC 3339 (e.g. '2026-03-07T12:00:00Z') \
|
||||
or 'YYYY-MM-DD HH:MM:SS' format.",
|
||||
input
|
||||
)))
|
||||
}
|
||||
|
||||
/// Tool for getting current time and date operations.
|
||||
pub struct TimeTool;
|
||||
|
||||
@@ -58,7 +16,7 @@ impl Tool for TimeTool {
|
||||
}
|
||||
|
||||
fn description(&self) -> &str {
|
||||
"Get current time, convert timezones, format timestamps, or calculate time differences."
|
||||
"Get current time, convert timezones, or calculate time differences."
|
||||
}
|
||||
|
||||
fn parameters_schema(&self) -> serde_json::Value {
|
||||
@@ -67,28 +25,20 @@ impl Tool for TimeTool {
|
||||
"properties": {
|
||||
"operation": {
|
||||
"type": "string",
|
||||
"enum": ["now", "parse", "convert", "format", "diff"],
|
||||
"enum": ["now", "parse", "format", "diff"],
|
||||
"description": "The time operation to perform"
|
||||
},
|
||||
"timestamp": {
|
||||
"type": "string",
|
||||
"description": "ISO 8601 timestamp (for parse/convert/format/diff operations)"
|
||||
"description": "ISO 8601 timestamp (for parse/format/diff operations)"
|
||||
},
|
||||
"format": {
|
||||
"type": "string",
|
||||
"description": "Output format string (for format operation)"
|
||||
},
|
||||
"timestamp2": {
|
||||
"type": "string",
|
||||
"description": "Second timestamp (for diff operation)"
|
||||
},
|
||||
"timezone": {
|
||||
"type": "string",
|
||||
"description": "IANA timezone name, e.g. 'America/New_York' (for now/convert/format/parse)"
|
||||
},
|
||||
"to_timezone": {
|
||||
"type": "string",
|
||||
"description": "Target IANA timezone for convert operation"
|
||||
},
|
||||
"format_string": {
|
||||
"type": "string",
|
||||
"description": "strftime format string (for format operation), default: '%Y-%m-%d %H:%M:%S %Z'"
|
||||
}
|
||||
},
|
||||
"required": ["operation"]
|
||||
@@ -107,91 +57,36 @@ impl Tool for TimeTool {
|
||||
let result = match operation {
|
||||
"now" => {
|
||||
let now = Utc::now();
|
||||
let mut result = serde_json::json!({
|
||||
"utc_iso": now.to_rfc3339(),
|
||||
serde_json::json!({
|
||||
"iso": now.to_rfc3339(),
|
||||
"unix": now.timestamp(),
|
||||
"unix_millis": now.timestamp_millis()
|
||||
});
|
||||
if let Some(tz_str) = params.get("timezone").and_then(|v| v.as_str()) {
|
||||
let tz = parse_timezone(tz_str)?;
|
||||
let local = now.with_timezone(&tz);
|
||||
result["local_iso"] = serde_json::json!(local.to_rfc3339());
|
||||
result["timezone"] = serde_json::json!(tz_str);
|
||||
}
|
||||
result
|
||||
})
|
||||
}
|
||||
"parse" => {
|
||||
let timestamp = require_str(¶ms, "timestamp")?;
|
||||
let tz = params
|
||||
.get("timezone")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(parse_timezone)
|
||||
.transpose()?;
|
||||
|
||||
let dt = parse_input_timestamp(timestamp, tz)?;
|
||||
let utc = dt.with_timezone(&Utc);
|
||||
|
||||
let mut result = serde_json::json!({
|
||||
"iso": utc.to_rfc3339(),
|
||||
"unix": utc.timestamp(),
|
||||
"unix_millis": utc.timestamp_millis()
|
||||
});
|
||||
if let Some(tz) = tz {
|
||||
let local = dt.with_timezone(&tz);
|
||||
result["local_iso"] = serde_json::json!(local.to_rfc3339());
|
||||
result["timezone"] = serde_json::json!(tz.to_string());
|
||||
}
|
||||
result
|
||||
}
|
||||
"convert" => {
|
||||
let timestamp = require_str(¶ms, "timestamp")?;
|
||||
let to_tz_str = require_str(¶ms, "to_timezone")?;
|
||||
let to_tz = parse_timezone(to_tz_str)?;
|
||||
|
||||
let from_tz = params
|
||||
.get("timezone")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(parse_timezone)
|
||||
.transpose()?;
|
||||
|
||||
let dt = parse_input_timestamp(timestamp, from_tz)?;
|
||||
let converted = dt.with_timezone(&to_tz);
|
||||
let dt: DateTime<Utc> = timestamp.parse().map_err(|e| {
|
||||
ToolError::InvalidParameters(format!("invalid timestamp: {}", e))
|
||||
})?;
|
||||
|
||||
serde_json::json!({
|
||||
"input": timestamp,
|
||||
"output": converted.to_rfc3339(),
|
||||
"timezone": to_tz.to_string()
|
||||
"iso": dt.to_rfc3339(),
|
||||
"unix": dt.timestamp(),
|
||||
"unix_millis": dt.timestamp_millis()
|
||||
})
|
||||
}
|
||||
"format" => {
|
||||
let timestamp = require_str(¶ms, "timestamp")?;
|
||||
let fmt = params
|
||||
.get("format_string")
|
||||
.and_then(|v| v.as_str())
|
||||
.unwrap_or("%Y-%m-%d %H:%M:%S %Z");
|
||||
|
||||
let tz = params
|
||||
.get("timezone")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(parse_timezone)
|
||||
.transpose()?;
|
||||
|
||||
let dt = parse_input_timestamp(timestamp, None)?;
|
||||
let formatted = if let Some(tz) = tz {
|
||||
dt.with_timezone(&tz).format(fmt).to_string()
|
||||
} else {
|
||||
dt.format(fmt).to_string()
|
||||
};
|
||||
|
||||
serde_json::json!({ "formatted": formatted })
|
||||
}
|
||||
"diff" => {
|
||||
let ts1 = require_str(¶ms, "timestamp")?;
|
||||
|
||||
let ts2 = require_str(¶ms, "timestamp2")?;
|
||||
|
||||
let dt1 = parse_input_timestamp(ts1, None)?;
|
||||
let dt2 = parse_input_timestamp(ts2, None)?;
|
||||
let dt1: DateTime<Utc> = ts1.parse().map_err(|e| {
|
||||
ToolError::InvalidParameters(format!("invalid timestamp: {}", e))
|
||||
})?;
|
||||
let dt2: DateTime<Utc> = ts2.parse().map_err(|e| {
|
||||
ToolError::InvalidParameters(format!("invalid timestamp2: {}", e))
|
||||
})?;
|
||||
|
||||
let diff = dt2.signed_duration_since(dt1);
|
||||
|
||||
@@ -217,220 +112,3 @@ impl Tool for TimeTool {
|
||||
false // Internal tool, no external data
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::context::JobContext;
|
||||
use serde_json::json;
|
||||
|
||||
fn test_ctx() -> JobContext {
|
||||
JobContext::new("test-job", "test time tool")
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_now_utc() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(json!({"operation": "now"}), &test_ctx())
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
assert!(v["utc_iso"].as_str().is_some());
|
||||
assert!(v["iso"].as_str().is_some());
|
||||
assert!(v["unix"].as_i64().is_some());
|
||||
// No timezone requested — no local_iso
|
||||
assert!(v.get("local_iso").is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_now_with_timezone() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({"operation": "now", "timezone": "America/New_York"}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
assert!(v["local_iso"].as_str().is_some());
|
||||
assert_eq!(v["timezone"].as_str().unwrap(), "America/New_York");
|
||||
// local_iso should contain a non-UTC offset
|
||||
let local = v["local_iso"].as_str().unwrap();
|
||||
assert!(!local.ends_with('Z') || local.contains("-04:00") || local.contains("-05:00"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_now_invalid_timezone() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({"operation": "now", "timezone": "Not/A/Zone"}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
let err = result.unwrap_err();
|
||||
assert!(err.to_string().contains("Unknown timezone"));
|
||||
assert!(err.to_string().contains("Not/A/Zone"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_convert_timezone() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "convert",
|
||||
"timestamp": "2026-03-07T12:00:00Z",
|
||||
"to_timezone": "Asia/Tokyo"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
// UTC 12:00 -> JST 21:00 (UTC+9)
|
||||
let output = v["output"].as_str().unwrap();
|
||||
assert!(output.contains("21:00:00"));
|
||||
assert_eq!(v["timezone"].as_str().unwrap(), "Asia/Tokyo");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_convert_dst_boundary() {
|
||||
let tool = TimeTool;
|
||||
// US spring forward: 2026-03-08 2:00 AM EST -> 3:00 AM EDT
|
||||
// Before DST: EST = UTC-5, After: EDT = UTC-4
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "convert",
|
||||
"timestamp": "2026-03-08T06:30:00Z",
|
||||
"to_timezone": "America/New_York"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
// UTC 06:30 on Mar 8 -> after spring forward, EDT (UTC-4) = 02:30
|
||||
// But DST springs forward at 2 AM -> 3 AM, so 06:30 UTC = 01:30 EST or 02:30 EDT
|
||||
let output = v["output"].as_str().unwrap();
|
||||
assert!(output.contains("2026-03-08"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_format_with_timezone() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "format",
|
||||
"timestamp": "2026-03-07T12:00:00Z",
|
||||
"timezone": "Europe/London",
|
||||
"format_string": "%Y-%m-%d %H:%M %Z"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
let formatted = v["formatted"].as_str().unwrap();
|
||||
assert!(formatted.contains("2026-03-07"));
|
||||
assert!(formatted.contains("12:00")); // London = UTC in March (before DST)
|
||||
assert!(formatted.contains("GMT"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_format_default_format_string() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "format",
|
||||
"timestamp": "2026-06-15T18:30:00Z",
|
||||
"timezone": "America/Los_Angeles"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
let formatted = v["formatted"].as_str().unwrap();
|
||||
// UTC 18:30 -> PDT (UTC-7) = 11:30
|
||||
assert!(formatted.contains("11:30:00"));
|
||||
assert!(formatted.contains("PDT"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_parse_naive_with_timezone() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "parse",
|
||||
"timestamp": "2026-03-07 09:00:00",
|
||||
"timezone": "America/New_York"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
// 09:00 EST = 14:00 UTC (EST = UTC-5 in March before DST)
|
||||
let iso = v["iso"].as_str().unwrap();
|
||||
assert!(iso.contains("14:00:00"));
|
||||
assert_eq!(v["timezone"].as_str().unwrap(), "America/New_York");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_diff() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "diff",
|
||||
"timestamp": "2026-03-07T00:00:00Z",
|
||||
"timestamp2": "2026-03-07T02:30:00Z"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let v: serde_json::Value = result.result.clone();
|
||||
assert_eq!(v["hours"].as_i64().unwrap(), 2);
|
||||
assert_eq!(v["minutes"].as_i64().unwrap(), 150);
|
||||
assert_eq!(v["seconds"].as_i64().unwrap(), 9000);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_convert_missing_to_timezone() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(
|
||||
json!({
|
||||
"operation": "convert",
|
||||
"timestamp": "2026-03-07T12:00:00Z"
|
||||
}),
|
||||
&test_ctx(),
|
||||
)
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_unknown_operation() {
|
||||
let tool = TimeTool;
|
||||
let result = tool
|
||||
.execute(json!({"operation": "explode"}), &test_ctx())
|
||||
.await;
|
||||
assert!(result.is_err());
|
||||
assert!(
|
||||
result
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("unknown operation")
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user