#![allow(dead_code)]
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use salvor_core::{Effect, Event, EventEnvelope, RunId};
use salvor_llm::Config;
use salvor_runtime::{Agent, AgentBuilder, ClockFn, RandomFn};
use salvor_tools::{DynTool, HandlerError, Sleep, ToolCtx, ToolError, ToolOutcome};
use serde_json::{Value, json};
use time::OffsetDateTime;
use time::macros::datetime;
use uuid::Uuid;
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, Request, Respond, ResponseTemplate};
pub fn fixed_run_id(tag: u8) -> RunId {
let mut bytes = [0u8; 16];
bytes[15] = tag;
bytes[6] = 0x40;
bytes[8] = 0x80;
RunId::from_uuid(Uuid::from_bytes(bytes))
}
pub fn fixed_clock() -> ClockFn {
Arc::new(|| datetime!(2026-07-14 12:00:00 UTC))
}
pub fn fixed_random() -> RandomFn {
Arc::new(|| 7)
}
#[derive(Clone)]
pub struct TestClock {
now: Arc<std::sync::Mutex<OffsetDateTime>>,
}
impl TestClock {
pub fn new(start: OffsetDateTime) -> Self {
Self {
now: Arc::new(std::sync::Mutex::new(start)),
}
}
pub fn injected(&self) -> ClockFn {
let now = self.now.clone();
Arc::new(move || *now.lock().expect("clock is not poisoned"))
}
pub fn read(&self) -> OffsetDateTime {
*self.now.lock().expect("clock is not poisoned")
}
pub fn set(&self, instant: OffsetDateTime) {
*self.now.lock().expect("clock is not poisoned") = instant;
}
}
pub struct NappingTool {
pub name: String,
pub effect: Effect,
pub wake_at: OffsetDateTime,
pub calls: Arc<AtomicUsize>,
}
impl NappingTool {
pub fn new(name: &str, effect: Effect, wake_at: OffsetDateTime) -> (Self, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
(
Self {
name: name.to_owned(),
effect,
wake_at,
calls: calls.clone(),
},
calls,
)
}
}
#[async_trait::async_trait]
impl DynTool for NappingTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"a test tool that parks its run on a timer"
}
fn effect(&self) -> Effect {
self.effect
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
_input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(ToolOutcome::Sleep(Sleep::until(self.wake_at)))
}
}
pub fn event_kinds(log: &[EventEnvelope]) -> Vec<&'static str> {
log.iter()
.map(|envelope| match &envelope.event {
Event::RunStarted { .. } => "RunStarted",
Event::ModelCallRequested { .. } => "ModelCallRequested",
Event::ModelCallCompleted { .. } => "ModelCallCompleted",
Event::ToolCallRequested { .. } => "ToolCallRequested",
Event::ToolCallCompleted { .. } => "ToolCallCompleted",
Event::NowObserved { .. } => "NowObserved",
Event::RandomObserved { .. } => "RandomObserved",
Event::Suspended { .. } => "Suspended",
Event::Resumed { .. } => "Resumed",
Event::SleepStarted { .. } => "SleepStarted",
Event::SleepCompleted {} => "SleepCompleted",
Event::BudgetExceeded { .. } => "BudgetExceeded",
Event::RunCompleted { .. } => "RunCompleted",
Event::RunFailed { .. } => "RunFailed",
Event::RunAbandoned { .. } => "RunAbandoned",
Event::GraphRunStarted { .. } => "GraphRunStarted",
Event::NodeEntered { .. } => "NodeEntered",
Event::NodeExited { .. } => "NodeExited",
Event::NodeSkipped { .. } => "NodeSkipped",
Event::BranchTaken { .. } => "BranchTaken",
Event::MapFannedOut { .. } => "MapFannedOut",
Event::MapIterationStarted { .. } => "MapIterationStarted",
Event::MapIterationJoined { .. } => "MapIterationJoined",
Event::FoldIterationStarted { .. } => "FoldIterationStarted",
Event::FoldIterationJoined { .. } => "FoldIterationJoined",
Event::FoldConverged { .. } => "FoldConverged",
})
.collect()
}
pub fn text_response(text: &str, input_tokens: u64, output_tokens: u64) -> Value {
json!({
"id": format!("msg_text_{input_tokens}_{output_tokens}"),
"model": "test-model",
"role": "assistant",
"content": [{"type": "text", "text": text}],
"stop_reason": "end_turn",
"usage": {"input_tokens": input_tokens, "output_tokens": output_tokens}
})
}
pub fn tool_use_response(
tool_use_id: &str,
tool: &str,
input: Value,
input_tokens: u64,
output_tokens: u64,
) -> Value {
json!({
"id": format!("msg_tool_{tool_use_id}"),
"model": "test-model",
"role": "assistant",
"content": [{"type": "tool_use", "id": tool_use_id, "name": tool, "input": input}],
"stop_reason": "tool_use",
"usage": {"input_tokens": input_tokens, "output_tokens": output_tokens}
})
}
pub struct ScriptedModel {
script: Vec<(usize, Value)>,
}
impl ScriptedModel {
pub async fn mount(script: Vec<(usize, Value)>) -> MockServer {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/messages"))
.respond_with(Self { script })
.mount(&server)
.await;
server
}
}
impl Respond for ScriptedModel {
fn respond(&self, request: &Request) -> ResponseTemplate {
let body: Value = match serde_json::from_slice(&request.body) {
Ok(body) => body,
Err(_) => return ResponseTemplate::new(400),
};
let count = body
.get("messages")
.and_then(Value::as_array)
.map_or(0, Vec::len);
for (expected, response) in &self.script {
if *expected == count {
return ResponseTemplate::new(200).set_body_json(response.clone());
}
}
ResponseTemplate::new(500).set_body_json(json!({
"error": {"type": "test_script", "message": format!("no scripted response for {count} messages")}
}))
}
}
pub struct ContentScriptedModel {
script: Vec<(String, Value)>,
}
impl ContentScriptedModel {
pub async fn mount(script: Vec<(&str, Value)>) -> MockServer {
let server = MockServer::start().await;
let script = script
.into_iter()
.map(|(needle, response)| (needle.to_owned(), response))
.collect();
Mock::given(method("POST"))
.and(path("/v1/messages"))
.respond_with(Self { script })
.mount(&server)
.await;
server
}
}
impl Respond for ContentScriptedModel {
fn respond(&self, request: &Request) -> ResponseTemplate {
let body = String::from_utf8_lossy(&request.body);
for (needle, response) in &self.script {
if body.contains(needle.as_str()) {
return ResponseTemplate::new(200).set_body_json(response.clone());
}
}
ResponseTemplate::new(500).set_body_json(json!({
"error": {"type": "test_script", "message": "no scripted response matched the request"}
}))
}
}
pub struct EchoTool {
pub name: String,
pub effect: Effect,
pub calls: Arc<AtomicUsize>,
}
impl EchoTool {
pub fn new(name: &str, effect: Effect) -> (Self, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
(
Self {
name: name.to_owned(),
effect,
calls: calls.clone(),
},
calls,
)
}
}
#[async_trait::async_trait]
impl DynTool for EchoTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"an echo test tool"
}
fn effect(&self) -> Effect {
self.effect
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(ToolOutcome::Output(json!({"published": input})))
}
}
pub struct ConstTool {
pub name: String,
pub effect: Effect,
pub value: Value,
pub calls: Arc<AtomicUsize>,
}
impl ConstTool {
pub fn new(name: &str, effect: Effect, value: Value) -> (Self, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
(
Self {
name: name.to_owned(),
effect,
value,
calls: calls.clone(),
},
calls,
)
}
}
#[async_trait::async_trait]
impl DynTool for ConstTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"a constant-value test tool"
}
fn effect(&self) -> Effect {
self.effect
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
_input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(ToolOutcome::Output(self.value.clone()))
}
}
pub struct PassTool {
pub name: String,
pub effect: Effect,
pub scores: Vec<Value>,
pub calls: Arc<AtomicUsize>,
}
impl PassTool {
pub fn new(name: &str, effect: Effect, scores: Vec<Value>) -> (Self, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
(
Self {
name: name.to_owned(),
effect,
scores,
calls: calls.clone(),
},
calls,
)
}
}
#[async_trait::async_trait]
impl DynTool for PassTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"a scripted fold-pass test tool"
}
fn effect(&self) -> Effect {
self.effect
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
self.calls.fetch_add(1, Ordering::SeqCst);
let pass = input.get("pass").and_then(Value::as_u64).unwrap_or(0);
let mut output = json!({"pass": pass + 1});
if let Some(score) = self.scores.get(pass as usize)
&& !score.is_null()
{
output["score"] = score.clone();
}
Ok(ToolOutcome::Output(output))
}
}
pub struct EnvelopePassTool {
pub name: String,
pub effect: Effect,
pub scores: Vec<Value>,
pub envelopes: Vec<bool>,
pub calls: Arc<AtomicUsize>,
}
impl EnvelopePassTool {
pub fn new(
name: &str,
effect: Effect,
scores: Vec<Value>,
envelopes: Vec<bool>,
) -> (Self, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
(
Self {
name: name.to_owned(),
effect,
scores,
envelopes,
calls: calls.clone(),
},
calls,
)
}
}
#[async_trait::async_trait]
impl DynTool for EnvelopePassTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"a scripted fold-pass test tool answering in an MCP result envelope"
}
fn effect(&self) -> Effect {
self.effect
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
self.calls.fetch_add(1, Ordering::SeqCst);
let pass = input.get("pass").and_then(Value::as_u64).unwrap_or(0);
let mut payload = json!({"pass": pass + 1});
if let Some(score) = self.scores.get(pass as usize)
&& !score.is_null()
{
payload["score"] = score.clone();
}
let wrapped = self.envelopes.get(pass as usize).copied().unwrap_or(true);
if !wrapped {
return Ok(ToolOutcome::Output(payload));
}
Ok(ToolOutcome::Output(json!({
"content": [{"type": "text", "text": payload.to_string()}],
"structuredContent": payload,
})))
}
}
pub struct SuspendingTool {
pub name: String,
pub reason: String,
pub input_schema: Value,
pub calls: Arc<AtomicUsize>,
}
impl SuspendingTool {
pub fn new(name: &str, reason: &str, input_schema: Value) -> (Self, Arc<AtomicUsize>) {
let calls = Arc::new(AtomicUsize::new(0));
(
Self {
name: name.to_owned(),
reason: reason.to_owned(),
input_schema,
calls: calls.clone(),
},
calls,
)
}
}
#[async_trait::async_trait]
impl DynTool for SuspendingTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"a test tool that always suspends"
}
fn effect(&self) -> Effect {
Effect::Read
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
_input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(ToolOutcome::Suspend(salvor_tools::Suspension::new(
self.reason.clone(),
self.input_schema.clone(),
)))
}
}
pub struct FailingTool {
pub name: String,
}
#[async_trait::async_trait]
impl DynTool for FailingTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"a failing test tool"
}
fn effect(&self) -> Effect {
Effect::Read
}
fn input_schema(&self) -> Value {
json!({"type": "object"})
}
async fn call_json(
&self,
_ctx: &ToolCtx,
_input: Value,
) -> Result<ToolOutcome<Value>, ToolError> {
Err(ToolError::Handler {
tool: self.name.clone(),
source: HandlerError::message("publish endpoint unreachable"),
})
}
}
pub fn agent_builder(server_uri: &str) -> AgentBuilder {
Agent::builder()
.model(
Config::new().with_base_url(server_uri).with_max_retries(0),
"test-model",
)
.system_prompt("You are a test agent.")
}