use crate::agent::runloop::unified::inline_events::harness::{HarnessEventEmitter, tool_outcome_observation_event};
use serde_json::Value;
use vtcode_core::core::agent::events::{
ToolOutputPayload, error_item_completed_event, tool_invocation_completed_event, tool_output_completed_event,
tool_output_payload_from_value, tool_started_event,
};
use vtcode_core::exec::events::{ToolCallStatus, tool_outcome_from_status};
use vtcode_core::tools::registry::ToolExecutionError;
use super::status::{ToolExecutionStatus, ToolPipelineOutcome};
pub(crate) fn emit_tool_start_if_needed(
harness_emitter: Option<&HarnessEventEmitter>,
already_started: bool,
tool_item_id: &str,
tool_call_id: &str,
tool_name: &str,
args: &Value,
) -> bool {
if already_started {
return true;
}
if let Some(emitter) = harness_emitter {
let _ = emitter.emit(tool_started_event(tool_item_id.to_string(), tool_name, Some(args), Some(tool_call_id)));
return true;
}
false
}
pub(crate) fn emit_tool_outcome_observation(
harness_emitter: Option<&HarnessEventEmitter>,
tool_name: &str,
outcome: &ToolPipelineOutcome,
) {
let Some(emitter) = harness_emitter else {
return;
};
let duration_ms = outcome.total_duration.as_millis().min(u128::from(u64::MAX)) as u64;
let category = outcome
.last_error_category
.as_ref()
.map(|category| category.as_str().to_string());
let _ = emitter.emit(tool_outcome_observation_event(tool_name, outcome.attempts, duration_ms, category));
}
fn tool_error_output_payload(error: &ToolExecutionError) -> ToolOutputPayload {
let mut payload = tool_output_payload_from_value(&error.to_json_value());
let user_message = error.user_message();
if payload.aggregated_output.is_empty() {
payload.aggregated_output = user_message;
} else {
payload.aggregated_output = format!("{user_message}\n{}", payload.aggregated_output);
}
payload
}
#[allow(
clippy::too_many_arguments,
reason = "Intentional compatibility, platform, or test-only suppression."
)] pub(super) fn emit_tool_completion_status(
harness_emitter: Option<&HarnessEventEmitter>,
tool_started_emitted: bool,
tool_execution_started: bool,
tool_item_id: &str,
tool_call_id: &str,
tool_name: &str,
args: &Value,
status: ToolCallStatus,
exit_code: Option<i32>,
spool_path: Option<&str>,
aggregated_output: impl Into<String>,
) {
if !tool_started_emitted {
return;
}
if let Some(emitter) = harness_emitter {
let aggregated_output = aggregated_output.into();
let outcome = tool_outcome_from_status(&status);
let _ = emitter.emit(tool_invocation_completed_event(
tool_item_id.to_string(),
tool_name,
Some(args),
Some(tool_call_id),
status.clone(),
outcome,
));
if tool_execution_started {
let _ = emitter.emit(tool_output_completed_event(
tool_item_id.to_string(),
Some(tool_call_id),
status,
exit_code,
spool_path,
aggregated_output,
));
} else if !aggregated_output.is_empty() {
let _ = emitter.emit(error_item_completed_event(format!("{tool_item_id}:error"), aggregated_output));
}
}
}
pub(crate) fn emit_tool_completion_for_status(
harness_emitter: Option<&HarnessEventEmitter>,
tool_started_emitted: bool,
tool_execution_started: bool,
tool_item_id: &str,
tool_call_id: &str,
tool_name: &str,
args: &Value,
tool_status: &ToolExecutionStatus,
) {
let (status, exit_code, output_payload) = match tool_status {
ToolExecutionStatus::Success { output, command_success, .. } => (
if *command_success {
ToolCallStatus::Completed
} else {
ToolCallStatus::Failed
},
output
.get("exit_code")
.and_then(Value::as_i64)
.and_then(|code| i32::try_from(code).ok()),
tool_output_payload_from_value(output),
),
ToolExecutionStatus::Failure { error } => (ToolCallStatus::Failed, None, tool_error_output_payload(error)),
ToolExecutionStatus::Timeout { error } => (ToolCallStatus::Failed, None, tool_error_output_payload(error)),
ToolExecutionStatus::Cancelled => (
ToolCallStatus::Failed,
None,
ToolOutputPayload {
aggregated_output: "Tool execution cancelled".to_string(),
spool_path: None,
},
),
};
emit_tool_completion_status(
harness_emitter,
tool_started_emitted,
tool_execution_started,
tool_item_id,
tool_call_id,
tool_name,
args,
status,
exit_code,
output_payload.spool_path.as_deref(),
output_payload.aggregated_output,
);
if tool_execution_started
&& tool_started_emitted
&& let Some(emitter) = harness_emitter
&& let ToolExecutionStatus::Success { output, .. } = tool_status
{
let _ = emitter.emit_exec_session_output(tool_item_id, tool_name, args, output);
}
if tool_execution_started
&& tool_started_emitted
&& let Some(emitter) = harness_emitter
&& let ToolExecutionStatus::Success { output, command_success: true, .. } = tool_status
&& let Some(event) =
vtcode_core::core::agent::events::file_change_completed_event(tool_item_id, tool_name, args, output)
{
let _ = emitter.emit(event);
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::{Value, json};
#[tokio::test]
async fn verifier_poll_completion_updates_the_launch_through_canonical_output() {
use vtcode_core::exec::events::{InputOrigin, ThreadEvent, TurnStartedEvent};
use vtcode_memory::explanation::ExplanationScope;
let workspace = tempfile::tempdir().unwrap();
let emitter = HarnessEventEmitter::new_async(workspace.path(), "interactive-verifier", None)
.await
.unwrap();
emitter.begin_task_turn("turn", "Verify changes", InputOrigin::User).unwrap();
emitter.emit(ThreadEvent::TurnStarted(TurnStartedEvent::default())).unwrap();
let launch_args = json!({"cmd":"cargo check --locked"});
assert!(emit_tool_start_if_needed(
Some(&emitter),
false,
"launch",
"launch-call",
"exec_command",
&launch_args
));
emit_tool_completion_for_status(
Some(&emitter),
true,
true,
"launch",
"launch-call",
"exec_command",
&launch_args,
&ToolExecutionStatus::Success {
output: json!({"session_id":"running-verifier","output":"checking"}),
stdout: None,
modified_files: vec![],
command_success: true,
},
);
let pending = emitter.explanation(ExplanationScope::Task).await.unwrap();
assert_eq!(pending.verification.len(), 1);
assert_eq!(pending.verification[0].fact.status, "pending");
let poll_args = json!({"session_id":"running-verifier","chars":""});
assert!(emit_tool_start_if_needed(Some(&emitter), false, "poll", "poll-call", "write_stdin", &poll_args));
emit_tool_completion_for_status(
Some(&emitter),
true,
true,
"poll",
"poll-call",
"write_stdin",
&poll_args,
&ToolExecutionStatus::Success {
output: json!({"session_id":"running-verifier","exit_code":0,"output":"finished"}),
stdout: None,
modified_files: vec![],
command_success: true,
},
);
let verified = emitter.explanation(ExplanationScope::Task).await.unwrap();
assert_eq!(verified.verification.len(), 1);
assert_eq!(verified.verification[0].fact.status, "passed");
assert!(verified.verification[0].fresh);
assert_eq!(verified.verification[0].exit_code, Some(0));
let source = emitter
.evidence(verified.verification[0].fact.evidence.clone(), 0)
.await
.unwrap();
assert!(source.text.contains("launch:exec-session-output"));
assert!(source.text.contains("finished"));
emitter.finish().await.unwrap();
}
#[tokio::test]
async fn background_verifier_completion_keeps_raw_and_derived_facts_in_the_original_task() {
use vtcode_core::exec::events::{
InputOrigin, ItemCompletedEvent, ThreadEvent, ThreadItem, ThreadItemDetails, TurnStartedEvent,
};
use vtcode_memory::explanation::ExplanationScope;
let workspace = tempfile::tempdir().unwrap();
let emitter = HarnessEventEmitter::new_async(workspace.path(), "background-verifier", None)
.await
.unwrap();
let original = emitter
.begin_task_turn("old-turn", "Verify first task", InputOrigin::User)
.unwrap();
emitter.emit(ThreadEvent::TurnStarted(TurnStartedEvent::default())).unwrap();
let args = json!({"cmd":"cargo check --locked"});
assert!(emit_tool_start_if_needed(Some(&emitter), false, "launch", "call", "exec_command", &args));
emit_tool_completion_for_status(
Some(&emitter),
true,
true,
"launch",
"call",
"exec_command",
&args,
&ToolExecutionStatus::Success {
output: json!({"session_id":"running-verifier"}),
stdout: None,
modified_files: vec![],
command_success: true,
},
);
let current = emitter
.begin_task_turn("new-turn", "Inspect second task", InputOrigin::User)
.unwrap();
emitter.emit(ThreadEvent::TurnStarted(TurnStartedEvent::default())).unwrap();
emitter
.emit(ThreadEvent::ItemCompleted(ItemCompletedEvent {
item: ThreadItem {
id: "raw-exit".into(),
context: None,
details: ThreadItemDetails::Harness(Box::new(
serde_json::from_value(json!({
"event":"background_subprocess_completed", "task_id":"exec:running-verifier",
"session_id":"running-verifier", "exec_session_id":"running-verifier", "exit_code":0,
"message":"Earlier verifier exited",
}))
.unwrap(),
)),
},
}))
.unwrap();
let task = emitter.explanation(ExplanationScope::Task).await.unwrap();
assert_eq!(task.task_id.as_deref(), Some(current.as_str()));
assert!(task.actions.is_empty());
assert!(task.verification.is_empty());
let session = emitter.explanation(ExplanationScope::Session).await.unwrap();
assert_eq!(session.verification.len(), 1);
assert_eq!(session.verification[0].fact.status, "passed");
assert!(
session
.actions
.iter()
.all(|action| action.task_id.as_deref() == Some(original.as_str()))
);
emitter.finish().await.unwrap();
}
#[test]
fn aggregates_command_output_without_duplicates() {
let output = json!({
"output": "same",
"stdout": "same",
"stderr": "warn"
});
let payload = tool_output_payload_from_value(&output);
assert_eq!(payload.aggregated_output, "same\n[stderr]\nwarn");
assert_eq!(payload.spool_path, None);
}
#[test]
fn includes_content_when_command_stream_fields_absent() {
let output = json!({
"content": "file body",
"path": "README.md"
});
let payload = tool_output_payload_from_value(&output);
assert!(payload.aggregated_output.starts_with("file body\nStructured output:"));
assert!(payload.aggregated_output.contains("\"path\": \"README.md\""));
assert_eq!(payload.spool_path, None);
}
#[test]
fn prefers_spool_reference_over_inline_output() {
let output = json!({
"output": "preview",
"spool_path": ".vtcode/context/tool_outputs/run-1.txt"
});
let payload = tool_output_payload_from_value(&output);
assert_eq!(payload.aggregated_output, "");
assert_eq!(payload.spool_path.as_deref(), Some(".vtcode/context/tool_outputs/run-1.txt"));
}
#[test]
fn structured_list_output_emits_compact_summary() {
let output = json!({
"items": [
{"path": "vtcode-tui/src/app.rs", "type": "file"},
{"path": "vtcode-tui/src/core_tui", "type": "directory"}
],
"count": 2,
"total": 11
});
let payload = tool_output_payload_from_value(&output);
assert!(payload.aggregated_output.contains("Listed 11 items"));
assert!(payload.aggregated_output.contains("vtcode-tui/src/app.rs"));
}
#[test]
fn non_zero_command_marks_failed_completion() {
let status = ToolExecutionStatus::Success {
output: json!({
"stdout": "boom",
"exit_code": 1
}),
stdout: Some("boom".to_string()),
modified_files: vec![],
command_success: false,
};
let (event_status, exit_code, output_payload) = match &status {
ToolExecutionStatus::Success { output, command_success, .. } => (
if *command_success {
ToolCallStatus::Completed
} else {
ToolCallStatus::Failed
},
output
.get("exit_code")
.and_then(Value::as_i64)
.and_then(|code| i32::try_from(code).ok()),
tool_output_payload_from_value(output),
),
_ => panic!("success status expected"),
};
assert_eq!(event_status, ToolCallStatus::Failed);
assert_eq!(exit_code, Some(1));
assert_eq!(output_payload.aggregated_output, "boom");
assert_eq!(output_payload.spool_path, None);
}
#[test]
fn failure_payload_preserves_structured_error_context() {
let error = ToolExecutionError::new(
"write_file",
vtcode_core::tools::registry::ToolErrorType::ExecutionError,
"write failed",
)
.with_partial_state(true, false)
.with_surface("unified_runloop")
.with_attempt(2);
let payload = tool_error_output_payload(&error);
assert!(payload.aggregated_output.contains(&error.user_message()));
assert!(payload.aggregated_output.contains("\"partial_state_possible\": true"));
assert!(payload.aggregated_output.contains("\"surface\": \"unified_runloop\""));
assert!(payload.aggregated_output.contains("\"attempt\": 2"));
assert!(payload.aggregated_output.contains("\"original_error\": null"));
assert_eq!(payload.spool_path, None);
}
}