Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions crates/ironclaw_common/src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,8 @@ pub enum AppEvent {
#[serde(skip_serializing_if = "Option::is_none")]
call_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
duration_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
thread_id: Option<String>,
},
#[serde(rename = "tool_result")]
Expand Down Expand Up @@ -420,6 +422,7 @@ mod tests {
error: None,
parameters: None,
call_id: None,
duration_ms: None,
thread_id: None,
},
AppEvent::ToolResult {
Expand Down
25 changes: 25 additions & 0 deletions crates/ironclaw_engine/src/executor/orchestrator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -779,6 +779,7 @@ async fn handle_execute_code_step(
// trace correlation.
call_id: format!("codeact-step-{}", exec_ctx.step_id.0),
error: error_msg,
duration_ms: code_start.elapsed().as_millis() as u64,
params_summary: None,
},
);
Expand Down Expand Up @@ -929,6 +930,7 @@ async fn handle_execute_action(
action_name: name.clone(),
call_id: call_id.clone(),
error,
duration_ms: 0,
params_summary: None,
},
&call_id,
Expand Down Expand Up @@ -962,6 +964,7 @@ async fn handle_execute_action(
action_name: name.clone(),
call_id: call_id.clone(),
error: reason,
duration_ms: 0,
params_summary: None,
},
&call_id,
Expand Down Expand Up @@ -1029,6 +1032,7 @@ async fn handle_execute_action(
action_name: name.clone(),
call_id: call_id.clone(),
error,
duration_ms: 0,
params_summary: None,
},
&call_id,
Expand All @@ -1045,6 +1049,7 @@ async fn handle_execute_action(

// 4. Execute
let ps = summarize_params(&name, &params);
let execution_start = std::time::Instant::now();
match effects
.execute_action(&name, params, &lease, &exec_ctx)
.await
Expand All @@ -1061,6 +1066,7 @@ async fn handle_execute_action(
.and_then(|v| v.as_str())
.map(String::from)
.unwrap_or_else(|| r.output.to_string());
let duration_ms = r.duration.as_millis() as u64;
emit_and_record(
thread,
event_tx,
Expand All @@ -1069,6 +1075,11 @@ async fn handle_execute_action(
action_name: name.clone(),
call_id: call_id.clone(),
error: error_msg,
duration_ms: if duration_ms > 0 {
duration_ms
} else {
execution_start.elapsed().as_millis() as u64
},
params_summary: ps.clone(),
},
&call_id,
Expand Down Expand Up @@ -1149,6 +1160,7 @@ async fn handle_execute_action(
action_name: name.clone(),
call_id: call_id.clone(),
error: e.to_string(),
duration_ms: execution_start.elapsed().as_millis() as u64,
params_summary: ps,
},
&call_id,
Expand Down Expand Up @@ -1261,6 +1273,7 @@ async fn handle_execute_actions_parallel(
action_name: pc.name.clone(),
call_id: pc.call_id.clone(),
error,
duration_ms: 0,
params_summary: None,
};
preflight.push(Some(PfOutcome::Error {
Expand Down Expand Up @@ -1292,6 +1305,7 @@ async fn handle_execute_actions_parallel(
action_name: pc.name.clone(),
call_id: pc.call_id.clone(),
error: reason,
duration_ms: 0,
params_summary: None,
};
preflight.push(Some(PfOutcome::Error {
Expand Down Expand Up @@ -1386,6 +1400,7 @@ async fn handle_execute_actions_parallel(
action_name: pc.name.clone(),
call_id: pc.call_id.clone(),
error,
duration_ms: 0,
params_summary: None,
};
preflight.push(Some(PfOutcome::Error {
Expand Down Expand Up @@ -1560,6 +1575,7 @@ async fn execute_single_action(
exec_ctx: &ThreadExecutionContext,
params_summary: Option<String>,
) -> (serde_json::Value, EventKind, serde_json::Value) {
let execution_start = std::time::Instant::now();
match effects.execute_action(name, params, lease, exec_ctx).await {
Ok(r) => {
// Surface wrapped errors as ActionFailed (see resolve_tool_future
Expand All @@ -1571,11 +1587,17 @@ async fn execute_single_action(
.and_then(|v| v.as_str())
.map(String::from)
.unwrap_or_else(|| r.output.to_string());
let duration_ms = r.duration.as_millis() as u64;
EventKind::ActionFailed {
step_id: exec_ctx.step_id,
action_name: name.to_string(),
call_id: call_id.to_string(),
error: error_msg,
duration_ms: if duration_ms > 0 {
duration_ms
} else {
execution_start.elapsed().as_millis() as u64
},
params_summary: params_summary.clone(),
}
} else {
Expand Down Expand Up @@ -1634,6 +1656,7 @@ async fn execute_single_action(
action_name: name.to_string(),
call_id: call_id.to_string(),
error: e.to_string(),
duration_ms: execution_start.elapsed().as_millis() as u64,
params_summary,
};
let result_json = serde_json::json!({
Expand Down Expand Up @@ -1714,11 +1737,13 @@ fn handle_emit_event(
let action_name = extract_string_kwarg(kwargs, "action_name").unwrap_or_default();
let call_id = extract_string_kwarg(kwargs, "call_id").unwrap_or_default();
let error = extract_string_kwarg(kwargs, "error").unwrap_or_default();
let duration_ms = extract_u64_kwarg(kwargs, "duration_ms").unwrap_or(0);
EventKind::ActionFailed {
step_id: StepId::new(),
action_name,
call_id,
error,
duration_ms,
params_summary: None,
}
}
Expand Down
44 changes: 29 additions & 15 deletions crates/ironclaw_engine/src/executor/scripting.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@

use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use std::time::{Duration, Instant};

use monty::{
ExcType, ExtFunctionResult, LimitedTracker, MontyException, MontyObject, MontyRun,
Expand Down Expand Up @@ -644,9 +644,11 @@ pub async fn execute_code_with_skills(
let ps = crate::types::event::summarize_params(&name, &params);

let handle = tokio::spawn(async move {
effects
let execution_start = Instant::now();
let result = effects
.execute_action(&name, params_clone, &lease_clone, &ctx)
.await
.await;
(result, execution_start.elapsed().as_millis() as u64)
});

pending_futures.insert(
Expand Down Expand Up @@ -976,7 +978,7 @@ pub fn code_hash(code: &str) -> String {
enum PendingFuture {
/// Tool action execution.
Tool {
handle: tokio::task::JoinHandle<Result<ActionResult, EngineError>>,
handle: tokio::task::JoinHandle<(Result<ActionResult, EngineError>, u64)>,
action_name: String,
call_id: String,
lease_id: crate::types::capability::LeaseId,
Expand Down Expand Up @@ -1021,6 +1023,7 @@ async fn preflight_action(
action_name: action_name.into(),
call_id: call_id.into(),
error: format!("no lease for action '{action_name}'"),
duration_ms: 0,
params_summary: None,
});
return PreflightResult::Denied(ExtFunctionResult::Error(MontyException::new(
Expand All @@ -1044,6 +1047,7 @@ async fn preflight_action(
action_name: action_name.into(),
call_id: call_id.into(),
error: reason.clone(),
duration_ms: 0,
params_summary: None,
});
return PreflightResult::Denied(ExtFunctionResult::Error(MontyException::new(
Expand Down Expand Up @@ -1522,7 +1526,7 @@ async fn handle_llm_query_batched_standalone(
/// Resolve a pending tool execution future.
#[allow(clippy::too_many_arguments)]
async fn resolve_tool_future(
handle: tokio::task::JoinHandle<Result<ActionResult, EngineError>>,
handle: tokio::task::JoinHandle<(Result<ActionResult, EngineError>, u64)>,
action_name: &str,
call_id: &str,
lease_id: crate::types::capability::LeaseId,
Expand All @@ -1534,7 +1538,7 @@ async fn resolve_tool_future(
events: &mut Vec<EventKind>,
) -> ExtFunctionResult {
match handle.await {
Ok(Ok(result)) => {
Ok((Ok(result), execution_duration_ms)) => {
// If the effect adapter wrapped a tool error as an Ok(ActionResult)
// with is_error=true (current convention in
// `EffectBridgeAdapter::execute_action_internal`), surface it as
Expand All @@ -1548,11 +1552,17 @@ async fn resolve_tool_future(
.and_then(|v| v.as_str())
.map(String::from)
.unwrap_or_else(|| result.output.to_string());
let duration_ms = result.duration.as_millis() as u64;
events.push(EventKind::ActionFailed {
step_id: context.step_id,
action_name: action_name.into(),
call_id: call_id.into(),
error: error_msg,
duration_ms: if duration_ms > 0 {
duration_ms
} else {
execution_duration_ms
},
params_summary,
});
} else {
Expand All @@ -1568,13 +1578,16 @@ async fn resolve_tool_future(
action_results.push(result);
ExtFunctionResult::Return(monty_val)
}
Ok(Err(EngineError::GatePaused {
gate_name,
action_name,
call_id,
resume_kind,
..
})) => {
Ok((
Err(EngineError::GatePaused {
gate_name,
action_name,
call_id,
resume_kind,
..
}),
_,
)) => {
let _ = leases.refund_use(lease_id).await;
events.push(EventKind::ApprovalRequested {
action_name,
Expand All @@ -1593,20 +1606,21 @@ async fn resolve_tool_future(
Some(format!("execution paused by gate '{gate_name}'")),
))
}
Ok(Err(e)) => {
Ok((Err(e), execution_duration_ms)) => {
events.push(EventKind::ActionFailed {
step_id: context.step_id,
action_name: action_name.into(),
call_id: call_id.into(),
error: e.to_string(),
duration_ms: execution_duration_ms,
params_summary,
});
action_results.push(ActionResult {
call_id: call_id.into(),
action_name: action_name.into(),
output: serde_json::json!({"error": e.to_string()}),
is_error: true,
duration: Duration::ZERO,
duration: Duration::from_millis(execution_duration_ms),
});
ExtFunctionResult::Error(MontyException::new(
ExcType::RuntimeError,
Expand Down
Loading
Loading