use super::*;
#[tokio::test]
async fn accepted_trigger_spawns_sub_agent_and_writes_trigger_result_audit() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn = Some(faux_stream_fn("sub-agent done"));
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = executor
.handle_trigger(sample_trigger("k-spawn", "trace-spawn"))
.await;
let completed = wait_for_event(&events, 5, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerCompleted {
trace_id, summary, ..
} if trace_id == "trace-spawn" => Some(summary.clone()),
_ => None,
})
})
.await;
assert!(completed.is_some(), "must emit TriggerCompleted");
let entries = session.entries().await.unwrap();
let record = entries
.iter()
.find_map(|e| match e {
SessionTreeEntry::Custom {
custom_type, data, ..
} if custom_type == "trigger_result" => Some(data.clone()),
_ => None,
})
.expect("trigger_result audit must exist");
let data = record.expect("trigger_result must carry data");
assert_eq!(data["trace_id"].as_str(), Some("trace-spawn"));
assert_eq!(data["success"].as_bool(), Some(true));
assert_eq!(
data["summary"].as_str(),
Some("sub-agent done"),
"summary must be the sub-agent's final assistant text"
);
assert!(data["branch_id"].is_null(), "5a in-memory: branch_id null");
}
#[tokio::test]
async fn event_ordering_handled_then_started_then_completed() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn = Some(faux_stream_fn("ok"));
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = executor
.handle_trigger(sample_trigger("k-order", "trace-order"))
.await;
wait_for_event(&events, 5, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerCompleted { trace_id, .. } if trace_id == "trace-order" => {
Some(())
}
_ => None,
})
})
.await
.expect("must complete");
let evs = events.lock().unwrap().clone();
let mut handled_idx = None;
let mut started_idx = None;
let mut completed_idx = None;
for (i, e) in evs.iter().enumerate() {
match e {
TriggerEvent::TriggerHandled {
trace_id,
state: theway_daemon::trigger_engine::types::TriggerState::Accepted,
..
} if trace_id == "trace-order" => handled_idx = Some(i),
TriggerEvent::TriggerExecutionStarted { trace_id, .. } if trace_id == "trace-order" => {
started_idx = Some(i)
}
TriggerEvent::TriggerCompleted { trace_id, .. } if trace_id == "trace-order" => {
completed_idx = Some(i)
}
_ => {}
}
}
let h = handled_idx.expect("TriggerHandled(Accepted)");
let s = started_idx.expect("TriggerExecutionStarted");
let c = completed_idx.expect("TriggerCompleted");
assert!(
h < s && s < c,
"expected Handled({h}) < Started({s}) < Completed({c}); events={:?}",
evs.iter().map(std::mem::discriminant).collect::<Vec<_>>()
);
}
#[tokio::test]
async fn pump_non_blocking_second_trigger_audited_while_first_runs() {
use tokio::sync::Mutex as TokioMutex;
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn: Option<StreamFn>;
let stream_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let stream_count_for_fn = stream_count.clone();
let stream_fn_val: StreamFn = Arc::new(move |_, _, _| {
let (stream, mut sender) = AssistantMessageEventStream::new();
let nth = stream_count_for_fn.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
tokio::spawn(async move {
if nth == 0 {
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
}
let msg = AssistantMessage {
role: AssistantRole::Assistant,
content: vec![ContentBlock::text("done")],
api: theway_llm_provider::Api::from("faux"),
provider: theway_llm_provider::Provider::from("faux"),
model: "faux".into(),
response_model: None,
response_id: None,
diagnostics: None,
usage: Usage::default(),
stop_reason: StopReason::Stop,
error_message: None,
timestamp: 0,
};
sender.push(AssistantMessageEvent::Start {
partial: msg.clone(),
});
sender.push(AssistantMessageEvent::Done {
reason: DoneReason::Stop,
message: msg,
});
});
stream
});
stream_fn = Some(stream_fn_val);
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = TokioMutex::new(()); let t0 = std::time::Instant::now();
let _ = executor
.handle_trigger(sample_trigger("k-slow", "trace-slow"))
.await;
let elapsed_1 = t0.elapsed();
assert!(
elapsed_1 < std::time::Duration::from_millis(200),
"handle_trigger must return promptly; took {elapsed_1:?}"
);
let t1 = std::time::Instant::now();
let _ = executor
.handle_trigger(sample_trigger("k-fast", "trace-fast"))
.await;
let elapsed_2 = t1.elapsed();
assert!(
elapsed_2 < std::time::Duration::from_millis(200),
"second handle_trigger must not block on first sub-agent; took {elapsed_2:?}"
);
wait_for_event(&events, 2, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerHandled { trace_id, .. } if trace_id == "trace-fast" => Some(()),
_ => None,
})
})
.await
.expect("second trigger must reach TriggerHandled within 2s");
}
#[tokio::test]
async fn running_snapshot_lists_in_flight_trigger_with_preview() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn: Option<StreamFn>;
let stream_fn_val: StreamFn = Arc::new(|_, _, _| {
let (stream, mut sender) = AssistantMessageEventStream::new();
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let msg = AssistantMessage {
role: AssistantRole::Assistant,
content: vec![ContentBlock::text("done")],
api: theway_llm_provider::Api::from("faux"),
provider: theway_llm_provider::Provider::from("faux"),
model: "faux".into(),
response_model: None,
response_id: None,
diagnostics: None,
usage: Usage::default(),
stop_reason: StopReason::Stop,
error_message: None,
timestamp: 0,
};
sender.push(AssistantMessageEvent::Start {
partial: msg.clone(),
});
sender.push(AssistantMessageEvent::Done {
reason: DoneReason::Stop,
message: msg,
});
});
stream
});
stream_fn = Some(stream_fn_val);
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let _ = executor
.handle_trigger(sample_trigger("k-running", "trace-running"))
.await;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
let mut found = None;
while std::time::Instant::now() < deadline {
let snap = executor.notification_status_snapshot();
if let Some(rt) = snap.running.iter().find(|r| r.trace_id == "trace-running") {
found = Some(rt.clone());
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
let rt = found.expect("running snapshot must include in-flight trigger");
assert_eq!(rt.source_label, "MCP github");
assert_eq!(rt.event_label, "pr merged");
assert!(
rt.prompt_preview.contains("MCP github") && rt.prompt_preview.contains("pr merged"),
"preview must reflect default prompt mapping, got {:?}",
rt.prompt_preview
);
tokio::time::sleep(std::time::Duration::from_millis(800)).await;
let snap_after = executor.notification_status_snapshot();
assert!(
snap_after
.running
.iter()
.all(|r| r.trace_id != "trace-running"),
"running snapshot must drop completed triggers"
);
}
#[tokio::test]
async fn abort_trigger_cancels_in_flight_sub_agent_and_emits_failed() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn: Option<StreamFn>;
let stream_fn_val: StreamFn = Arc::new(|_, _, _| {
let (stream, mut sender) = AssistantMessageEventStream::new();
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_secs(30)).await;
let msg = AssistantMessage {
role: AssistantRole::Assistant,
content: vec![ContentBlock::text("done")],
api: theway_llm_provider::Api::from("faux"),
provider: theway_llm_provider::Provider::from("faux"),
model: "faux".into(),
response_model: None,
response_id: None,
diagnostics: None,
usage: Usage::default(),
stop_reason: StopReason::Stop,
error_message: None,
timestamp: 0,
};
sender.push(AssistantMessageEvent::Start {
partial: msg.clone(),
});
sender.push(AssistantMessageEvent::Done {
reason: DoneReason::Stop,
message: msg,
});
});
stream
});
stream_fn = Some(stream_fn_val);
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = executor
.handle_trigger(sample_trigger("k-abort", "trace-abort"))
.await;
wait_for_event(&events, 2, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerExecutionStarted { trace_id, .. } if trace_id == "trace-abort" => {
Some(())
}
_ => None,
})
})
.await
.expect("ExecutionStarted must fire before abort");
executor.abort_trigger("trace-abort");
let reason = wait_for_event(&events, 3, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerFailed { trace_id, reason } if trace_id == "trace-abort" => {
Some(reason.clone())
}
_ => None,
})
})
.await
.expect("TriggerFailed must arrive within 3s of abort");
assert_eq!(
reason, "aborted",
"abort_trigger must emit TriggerFailed with reason \"aborted\""
);
let entries = session.entries().await.unwrap();
let record = entries
.iter()
.find_map(|e| match e {
SessionTreeEntry::Custom {
custom_type, data, ..
} if custom_type == "trigger_result" => Some(data.clone()),
_ => None,
})
.expect("trigger_result must be written even on abort");
let data = record.expect("data");
assert_eq!(data["success"].as_bool(), Some(false));
assert_eq!(data["trace_id"].as_str(), Some("trace-abort"));
}
#[tokio::test]
async fn non_accepted_states_do_not_spawn_sub_agent() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn = Some(faux_stream_fn("done"));
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = executor
.handle_trigger(sample_trigger("k-dedup-test", "trace-1"))
.await;
let _ = executor
.handle_trigger(sample_trigger("k-dedup-test", "trace-2"))
.await;
wait_for_event(&events, 5, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerCompleted { trace_id, .. } if trace_id == "trace-1" => Some(()),
_ => None,
})
})
.await
.expect("first trigger completes");
let evs = events.lock().unwrap().clone();
let spawned_for_trace_2 = evs.iter().any(|e| {
matches!(
e,
TriggerEvent::TriggerExecutionStarted { trace_id, .. } if trace_id == "trace-2"
)
});
assert!(
!spawned_for_trace_2,
"Deduped trigger must NOT spawn a sub-agent; got events: {:?}",
evs.iter()
.filter_map(|e| match e {
TriggerEvent::TriggerExecutionStarted { trace_id, .. } => Some(trace_id.clone()),
_ => None,
})
.collect::<Vec<_>>()
);
}
#[tokio::test]
async fn trigger_result_audit_records_failure_reason_for_resume_archaeology() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn: Option<StreamFn>;
let stream_fn_val: StreamFn = Arc::new(|_, _, _| {
let (stream, mut sender) = AssistantMessageEventStream::new();
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_secs(30)).await;
let msg = AssistantMessage {
role: AssistantRole::Assistant,
content: vec![ContentBlock::text("done")],
api: theway_llm_provider::Api::from("faux"),
provider: theway_llm_provider::Provider::from("faux"),
model: "faux".into(),
response_model: None,
response_id: None,
diagnostics: None,
usage: Usage::default(),
stop_reason: StopReason::Stop,
error_message: None,
timestamp: 0,
};
sender.push(AssistantMessageEvent::Start {
partial: msg.clone(),
});
sender.push(AssistantMessageEvent::Done {
reason: DoneReason::Stop,
message: msg,
});
});
stream
});
stream_fn = Some(stream_fn_val);
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = executor
.handle_trigger(sample_trigger("k-reason", "trace-reason"))
.await;
wait_for_event(&events, 2, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerExecutionStarted { trace_id, .. }
if trace_id == "trace-reason" =>
{
Some(())
}
_ => None,
})
})
.await
.expect("ExecutionStarted first");
executor.abort_trigger("trace-reason");
wait_for_event(&events, 3, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerFailed { trace_id, .. } if trace_id == "trace-reason" => Some(()),
_ => None,
})
})
.await
.expect("TriggerFailed within 3s");
let entries = session.entries().await.unwrap();
let data = entries
.iter()
.find_map(|e| match e {
SessionTreeEntry::Custom {
custom_type, data, ..
} if custom_type == "trigger_result" => data.clone(),
_ => None,
})
.expect("trigger_result data");
assert_eq!(data["success"].as_bool(), Some(false));
assert_eq!(
data["reason"].as_str(),
Some("aborted"),
"trigger_result must persist failure reason so jsonl-only readers see WHY: {data:?}"
);
assert!(
data["cost_usd"].is_null(),
"5a does not measure cost — null is honest; 0.0 was misleading"
);
}
#[tokio::test]
async fn trigger_result_summary_truncation_handles_multibyte_codepoint_via_production_path() {
let storage = Arc::new(MemorySessionStorage::new());
let session = Session::new(storage.clone() as Arc<dyn SessionStorage>);
let trigger_runtime = TriggerRuntimeConfig::default();
let before_trigger: Option<BeforeTriggerHook> = None;
let on_trigger_prompt: Option<OnTriggerPromptHook> = None;
let before_trigger_action: Option<BeforeTriggerActionHook> = None;
let stream_fn = {
let huge_text: &'static str = Box::leak(("你".repeat(1366)).into_boxed_str());
Some(faux_stream_fn(huge_text))
};
let harness = AgentHarness::new(AgentHarnessOptions::new(faux_model(), session.clone()));
let executor = Arc::new(TriggerExecutor::new(
harness.agent_arc(),
session.clone(),
trigger_runtime,
before_trigger,
on_trigger_prompt,
before_trigger_action,
stream_fn,
None,
None,
));
let events = Arc::new(std::sync::Mutex::new(Vec::<TriggerEvent>::new()));
let sink = events.clone();
let _unsub = executor.subscribe(Arc::new(move |ev| {
sink.lock().unwrap().push(ev);
}));
let _ = executor
.handle_trigger(sample_trigger("k-utf8-trunc", "trace-utf8-trunc"))
.await;
wait_for_event(&events, 5, |evs| {
evs.iter().find_map(|e| match e {
TriggerEvent::TriggerCompleted {
trace_id, summary, ..
} if trace_id == "trace-utf8-trunc" => Some(summary.clone()),
_ => None,
})
})
.await
.expect("TriggerCompleted must fire — pre-fix code would panic and abort the task");
let entries = session.entries().await.unwrap();
let data = entries
.iter()
.find_map(|e| match e {
SessionTreeEntry::Custom {
custom_type, data, ..
} if custom_type == "trigger_result" => data.clone(),
_ => None,
})
.expect("trigger_result audit");
let summary = data["summary"]
.as_str()
.expect("summary must be a string (proves valid UTF-8 round-trip through serde_json)");
assert!(
summary.ends_with("…[truncated]"),
"summary must be capped with truncation marker; got len={} ending={:?}",
summary.len(),
summary.chars().rev().take(15).collect::<String>(),
);
assert!(
summary.len() <= 4096,
"final summary (including truncation marker) must respect 4 KiB cap; got {}",
summary.len()
);
let body_only = summary.trim_end_matches("…[truncated]");
assert!(
body_only.chars().all(|c| c == '你'),
"truncation MUST land on a char boundary; got non-你 chars in body"
);
}