mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
* feat: structured fallback deliverables for failed/stuck jobs (#221) When a job fails or gets stuck, build a FallbackDeliverable that captures partial results, action statistics, cost, timing, and repair attempts. This replaces opaque error strings with structured data users can act on. - Add FallbackDeliverable, LastAction, ActionStats types in context/fallback.rs - Store fallback in JobContext.metadata["fallback_deliverable"] on failure - Surface fallback in job_status tool output and SSE job_result events - Update mark_failed() and mark_stuck() in worker to build fallback - 8 unit tests covering zero/mixed actions, truncation, timing, serialization Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address review comments on fallback deliverables - Fix doc comment: "200 chars" -> "200 bytes (UTF-8 safe)" since truncate_str operates on byte length, not character count. - Add code comment documenting that SSE fallback_deliverable is currently always None (forward-compatible infrastructure). Co-Authored-By: Claude Opus 4.6 <[email protected]> * refactor: take Option<&FallbackDeliverable> instead of &Option<…> Addresses Gemini review feedback: idiomatic Rust prefers Option<&T> over &Option<T> for borrowed optional values. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: guard against non-object metadata and add fallback test - store_fallback_in_metadata now resets metadata to {} when it's any non-object type (string, array, number), not just null. Prevents panic on index assignment. - Add test_job_status_includes_fallback_deliverable to verify the fallback field is surfaced in job_status tool output. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: use sanitized output in fallback preview + add integration tests Security fix: FallbackDeliverable::build() now uses output_sanitized instead of output_raw, preventing secrets/PII from leaking through the job_status tool and SSE job_result events. Also adds: - test_fallback_uses_sanitized_output: proves raw secrets don't leak - test_store_fallback_in_metadata_roundtrip: full serialize/deserialize - test_store_fallback_handles_non_object_metadata: edge case coverage - test_store_fallback_none_is_noop: None input is safe Addresses serrrfirat review feedback on PR #236. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: harden fallback deliverables against review findings - Truncate failure_reason to 1000 bytes to prevent metadata bloat - Add tracing::warn on fallback serialization failure (was silently discarded) - Fix module/struct docs to cover stuck jobs, remove stale SSE claim - Fix job.rs test to use real FallbackDeliverable field names - Add tests for failure_reason truncation and completed_at=None elapsed time - Fix pre-existing clippy warning in settings.rs (field_reassign_with_default) Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address Copilot review findings on fallback deliverables - Fix output_raw/output_sanitized field swap in ActionRecord::succeed() so sanitized data actually goes into the sanitized field (security) - Return None instead of empty Memory when get_memory fails in build_fallback, with tracing::warn for observability - Replace manual elapsed calculation with ctx.elapsed() which already clamps negative durations Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: resolve rebase conflicts and update tests for parameter swap - Add fallback field to SseEvent::JobResult in job_monitor - Fix type annotation in fallback deliverable test - Update test_action_record_succeed_sets_fields for new parameter order - Use create_job_for_user in test (API changed on main) Co-Authored-By: Claude Opus 4.6 <[email protected]> * chore: trigger CI re-check after rebase * fix: fall back to error message for failed action output_preview When the last action is a failed tool call, output_sanitized is None, leaving output_preview empty. Now falls back to the action's error message so users see what went wrong. [skip-regression-check] * ci: add safety comments to test code for no-panics check The CI no-panics grep check cannot distinguish test code inside src/ files from production code. Add // safety: test annotations to .unwrap(), .expect(), and assert!() calls in #[cfg(test)] modules. * fix: clarify succeed() doc and avoid clone in output_preview - Fix doc comment: output_raw is stored as pretty-printed JSON string, not a raw JSON value - Borrow string slice directly in fallback preview to avoid cloning potentially large sanitized outputs before truncation * refactor: reuse floor_char_boundary in truncate_str Replace hand-rolled UTF-8 boundary logic with existing crate::util::floor_char_boundary to reduce duplication. * fix: rename SSE fallback field to fallback_deliverable for consistency The SSE JobResult field was named `fallback` while everywhere else (metadata key, job_status tool) uses `fallback_deliverable`. Align the SSE wire format to avoid forcing clients to handle two names. --------- Co-authored-by: Claude Opus 4.6 <[email protected]>
298 lines
11 KiB
Rust
298 lines
11 KiB
Rust
//! Background job monitor that forwards Claude Code output to the main agent loop.
|
|
//!
|
|
//! When the main agent kicks off a sandbox job (especially Claude Code), this
|
|
//! monitor subscribes to the broadcast event channel and injects relevant
|
|
//! assistant messages back into the channel manager's stream. This lets the
|
|
//! main agent see what the sub-agent is producing and surface it to the user.
|
|
//!
|
|
//! ```text
|
|
//! Container ──NDJSON──► Orchestrator ──broadcast──► JobMonitor
|
|
//! │
|
|
//! inject_tx (mpsc)
|
|
//! │
|
|
//! ▼
|
|
//! Agent Loop
|
|
//! ```
|
|
|
|
use tokio::sync::{broadcast, mpsc};
|
|
use tokio::task::JoinHandle;
|
|
use uuid::Uuid;
|
|
|
|
use crate::channels::IncomingMessage;
|
|
use crate::channels::web::types::SseEvent;
|
|
|
|
/// Route context for forwarding job monitor events back to the user's channel.
|
|
#[derive(Debug, Clone)]
|
|
pub struct JobMonitorRoute {
|
|
pub channel: String,
|
|
pub user_id: String,
|
|
pub thread_id: Option<String>,
|
|
}
|
|
|
|
/// Spawn a background task that watches for events from a specific job and
|
|
/// injects assistant messages into the agent loop.
|
|
///
|
|
/// The monitor forwards:
|
|
/// - `SseEvent::JobMessage` (assistant role): injected as incoming messages so
|
|
/// the main agent can read and relay to the user.
|
|
/// - `SseEvent::JobResult`: injected as a completion notice, then the task exits.
|
|
///
|
|
/// Tool use/result and status events are intentionally skipped (too noisy for
|
|
/// the main agent's context window).
|
|
pub fn spawn_job_monitor(
|
|
job_id: Uuid,
|
|
mut event_rx: broadcast::Receiver<(Uuid, SseEvent)>,
|
|
inject_tx: mpsc::Sender<IncomingMessage>,
|
|
route: JobMonitorRoute,
|
|
) -> JoinHandle<()> {
|
|
let short_id = job_id.to_string()[..8].to_string();
|
|
|
|
tokio::spawn(async move {
|
|
tracing::info!(job_id = %short_id, "Job monitor started successfully");
|
|
|
|
loop {
|
|
match event_rx.recv().await {
|
|
Ok((ev_job_id, event)) => {
|
|
if ev_job_id != job_id {
|
|
continue;
|
|
}
|
|
|
|
match event {
|
|
SseEvent::JobMessage { role, content, .. } if role == "assistant" => {
|
|
let mut msg = IncomingMessage::new(
|
|
route.channel.clone(),
|
|
route.user_id.clone(),
|
|
format!("[Job {}] Claude Code: {}", short_id, content),
|
|
)
|
|
.into_internal();
|
|
if let Some(ref thread_id) = route.thread_id {
|
|
msg = msg.with_thread(thread_id.clone());
|
|
}
|
|
if inject_tx.send(msg).await.is_err() {
|
|
tracing::debug!(
|
|
job_id = %short_id,
|
|
"Inject channel closed, stopping monitor"
|
|
);
|
|
break;
|
|
}
|
|
}
|
|
SseEvent::JobResult { status, .. } => {
|
|
let mut msg = IncomingMessage::new(
|
|
route.channel.clone(),
|
|
route.user_id.clone(),
|
|
format!(
|
|
"[Job {}] Container finished (status: {})",
|
|
short_id, status
|
|
),
|
|
)
|
|
.into_internal();
|
|
if let Some(ref thread_id) = route.thread_id {
|
|
msg = msg.with_thread(thread_id.clone());
|
|
}
|
|
let _ = inject_tx.send(msg).await;
|
|
tracing::debug!(
|
|
job_id = %short_id,
|
|
status = %status,
|
|
"Job monitor exiting (job finished)"
|
|
);
|
|
break;
|
|
}
|
|
_ => {
|
|
// Skip tool_use, tool_result, status events
|
|
}
|
|
}
|
|
}
|
|
Err(broadcast::error::RecvError::Lagged(n)) => {
|
|
tracing::warn!(
|
|
job_id = %short_id,
|
|
skipped = n,
|
|
"Job monitor lagged, some events were dropped"
|
|
);
|
|
}
|
|
Err(broadcast::error::RecvError::Closed) => {
|
|
tracing::debug!(
|
|
job_id = %short_id,
|
|
"Broadcast channel closed, stopping monitor"
|
|
);
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
fn test_route() -> JobMonitorRoute {
|
|
JobMonitorRoute {
|
|
channel: "cli".to_string(),
|
|
user_id: "user-1".to_string(),
|
|
thread_id: Some("thread-1".to_string()),
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_monitor_forwards_assistant_messages() {
|
|
let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16);
|
|
let (inject_tx, mut inject_rx) = mpsc::channel::<IncomingMessage>(16);
|
|
|
|
let job_id = Uuid::new_v4();
|
|
let _handle = spawn_job_monitor(job_id, event_tx.subscribe(), inject_tx, test_route());
|
|
|
|
// Send an assistant message
|
|
event_tx
|
|
.send((
|
|
job_id,
|
|
SseEvent::JobMessage {
|
|
job_id: job_id.to_string(),
|
|
role: "assistant".to_string(),
|
|
content: "I found a bug".to_string(),
|
|
},
|
|
))
|
|
.unwrap();
|
|
|
|
let msg = tokio::time::timeout(std::time::Duration::from_secs(1), inject_rx.recv())
|
|
.await
|
|
.unwrap()
|
|
.unwrap();
|
|
|
|
assert_eq!(msg.channel, "cli");
|
|
assert_eq!(msg.user_id, "user-1");
|
|
assert_eq!(msg.thread_id, Some("thread-1".to_string()));
|
|
assert!(msg.content.contains("I found a bug"));
|
|
assert!(msg.is_internal, "monitor messages must be marked internal");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_monitor_ignores_other_jobs() {
|
|
let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16);
|
|
let (inject_tx, mut inject_rx) = mpsc::channel::<IncomingMessage>(16);
|
|
|
|
let job_id = Uuid::new_v4();
|
|
let other_job_id = Uuid::new_v4();
|
|
let _handle = spawn_job_monitor(job_id, event_tx.subscribe(), inject_tx, test_route());
|
|
|
|
// Send a message for a different job
|
|
event_tx
|
|
.send((
|
|
other_job_id,
|
|
SseEvent::JobMessage {
|
|
job_id: other_job_id.to_string(),
|
|
role: "assistant".to_string(),
|
|
content: "wrong job".to_string(),
|
|
},
|
|
))
|
|
.unwrap();
|
|
|
|
// Should not receive anything
|
|
let result =
|
|
tokio::time::timeout(std::time::Duration::from_millis(100), inject_rx.recv()).await;
|
|
assert!(
|
|
result.is_err(),
|
|
"should have timed out, no message expected"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_monitor_exits_on_job_result() {
|
|
let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16);
|
|
let (inject_tx, mut inject_rx) = mpsc::channel::<IncomingMessage>(16);
|
|
|
|
let job_id = Uuid::new_v4();
|
|
let handle = spawn_job_monitor(job_id, event_tx.subscribe(), inject_tx, test_route());
|
|
|
|
// Send a completion event
|
|
event_tx
|
|
.send((
|
|
job_id,
|
|
SseEvent::JobResult {
|
|
job_id: job_id.to_string(),
|
|
status: "completed".to_string(),
|
|
session_id: None,
|
|
fallback_deliverable: None,
|
|
},
|
|
))
|
|
.unwrap();
|
|
|
|
// Should receive the completion message
|
|
let msg = tokio::time::timeout(std::time::Duration::from_secs(1), inject_rx.recv())
|
|
.await
|
|
.unwrap()
|
|
.unwrap();
|
|
assert!(msg.content.contains("finished"));
|
|
|
|
// The monitor task should exit
|
|
tokio::time::timeout(std::time::Duration::from_secs(1), handle)
|
|
.await
|
|
.expect("monitor should have exited")
|
|
.expect("monitor task should not panic");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_monitor_skips_tool_events() {
|
|
let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16);
|
|
let (inject_tx, mut inject_rx) = mpsc::channel::<IncomingMessage>(16);
|
|
|
|
let job_id = Uuid::new_v4();
|
|
let _handle = spawn_job_monitor(job_id, event_tx.subscribe(), inject_tx, test_route());
|
|
|
|
// Send tool use event (should be skipped)
|
|
event_tx
|
|
.send((
|
|
job_id,
|
|
SseEvent::JobToolUse {
|
|
job_id: job_id.to_string(),
|
|
tool_name: "shell".to_string(),
|
|
input: serde_json::json!({"command": "ls"}),
|
|
},
|
|
))
|
|
.unwrap();
|
|
|
|
// Send user message (should be skipped)
|
|
event_tx
|
|
.send((
|
|
job_id,
|
|
SseEvent::JobMessage {
|
|
job_id: job_id.to_string(),
|
|
role: "user".to_string(),
|
|
content: "user prompt".to_string(),
|
|
},
|
|
))
|
|
.unwrap();
|
|
|
|
// Should not receive anything for tool events or user messages
|
|
let result =
|
|
tokio::time::timeout(std::time::Duration::from_millis(100), inject_rx.recv()).await;
|
|
assert!(
|
|
result.is_err(),
|
|
"should have timed out, no message expected"
|
|
);
|
|
}
|
|
|
|
/// Regression test: external channels must not be able to spoof the
|
|
/// `is_internal` flag via metadata keys. A message created through
|
|
/// the normal `IncomingMessage::new` + `with_metadata` path must
|
|
/// always have `is_internal == false`, regardless of metadata content.
|
|
#[test]
|
|
fn test_external_metadata_cannot_spoof_internal_flag() {
|
|
let msg = IncomingMessage::new("wasm_channel", "attacker", "pwned").with_metadata(
|
|
serde_json::json!({
|
|
"__internal_job_monitor": true,
|
|
"is_internal": true,
|
|
}),
|
|
);
|
|
assert!(
|
|
!msg.is_internal,
|
|
"with_metadata must not set is_internal — only into_internal() can"
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn test_into_internal_sets_flag() {
|
|
let msg = IncomingMessage::new("monitor", "system", "test").into_internal();
|
|
assert!(msg.is_internal);
|
|
}
|
|
}
|