use std::collections::HashMap;
use std::sync::Arc;
use foundation_ai::agentic::testing::{mock_text, mock_tool_call, MockModelProvider};
use foundation_ai::agentic::tool_impl::ToolCallManager;
use foundation_ai::agentic::{
AgentConfig, AgentLoop, ContextConfig, ContextProvider, ErrorPolicy, KvMemoryStore,
LoopDetectorConfig, MemoryConfig, MemoryCoordinator, MemoryHierarchy, MessageApi,
SteeringQueues, TokenLedger,
};
use foundation_ai::types::{
MessageRole, Messages, ModelId, ModelOutput, ProviderRouter, SessionId, SessionRecord,
TextContent, UserModelContent,
};
use foundation_core::valtron::{TaskIterator, TaskStatus};
use foundation_db::{MemoryDocumentStore, MemoryStorage};
type TestMemStore = KvMemoryStore<MemoryStorage>;
type TestDocStore = MemoryDocumentStore;
struct Harness {
agent: AgentLoop<TestDocStore, TestMemStore>,
follow_up: Arc<concurrent_queue::ConcurrentQueue<Messages>>,
priority: Arc<concurrent_queue::ConcurrentQueue<Messages>>,
cancel: Arc<std::sync::atomic::AtomicU32>,
ledger: TokenLedger,
}
impl Harness {
fn set_budget(&self, budget: u64) {
self.ledger.set_budget(Some(budget));
}
}
fn harness_with(router: ProviderRouter, config: AgentConfig) -> Harness {
harness_with_tools(router, config, Vec::new())
}
fn harness_with_tools(
router: ProviderRouter,
config: AgentConfig,
tools: Vec<Arc<dyn foundation_ai::agentic::tool_impl::ToolImpl>>,
) -> Harness {
let session_id = SessionId::new();
let ledger = TokenLedger::new();
let memory_store = Arc::new(KvMemoryStore::new(MemoryStorage::new()));
let message_api = MessageApi::new(session_id.clone(), MemoryDocumentStore::new());
let context_provider = ContextProvider::new(
session_id.clone(),
message_api.clone(),
Arc::clone(&memory_store),
ledger.clone(),
Some("You are a helpful assistant.".into()),
ContextConfig::default(),
);
let queues = SteeringQueues::new();
let follow_up = Arc::clone(&queues.follow_up);
let priority = Arc::clone(&queues.priority);
let cancel = Arc::clone(&queues.cancel_signal);
let memory = MemoryHierarchy::new(
session_id.clone(),
MemoryCoordinator::new(
KvMemoryStore::new(MemoryStorage::new()),
MemoryDocumentStore::new(),
),
ledger.clone(),
MemoryConfig::default(),
);
let tool_manager = ToolCallManager::new(session_id.clone());
for tool in tools {
tool_manager.register(tool);
}
let ledger_handle = ledger.clone();
let agent = AgentLoop::new(
session_id.clone(),
context_provider,
tool_manager,
queues,
memory,
message_api,
ledger,
ErrorPolicy::new(),
router,
config,
);
Harness {
agent,
follow_up,
priority,
cancel,
ledger: ledger_handle,
}
}
fn config_for(model: &str) -> AgentConfig {
AgentConfig {
primary_model: ModelId::Name(model.into(), None),
..Default::default()
}
}
fn user_msg(text: &str) -> Messages {
Messages::User {
id: foundation_compact::ids::new_scru128(),
role: MessageRole::User,
content: UserModelContent::Text(TextContent {
content: text.into(),
signature: None,
}),
signature: None,
}
}
fn drive(h: &mut Harness) -> Vec<SessionRecord> {
let mut records = Vec::new();
for _ in 0..2_000 {
match h.agent.next_status() {
None => return records,
Some(TaskStatus::Ready(record)) => records.push(record),
Some(_) => {}
}
}
panic!("agent loop did not terminate within 2000 steps");
}
fn assistant_texts(records: &[SessionRecord]) -> Vec<String> {
records
.iter()
.filter_map(|r| match r {
SessionRecord::Conversation {
message: Messages::Assistant { content, .. },
} => match content {
ModelOutput::Text(t) => Some(t.content.clone()),
_ => None,
},
_ => None,
})
.collect()
}
fn has_failed_action(records: &[SessionRecord]) -> bool {
records
.iter()
.any(|r| matches!(r, SessionRecord::FailedAction { .. }))
}
fn summary_count(records: &[SessionRecord]) -> Option<u64> {
records.iter().find_map(|r| match r {
SessionRecord::Summary { message_count, .. } => Some(*message_count),
_ => None,
})
}
#[test]
fn follow_up_drives_a_full_generation_turn() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("hello from the model")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert_eq!(
assistant_texts(&records),
vec!["hello from the model".to_string()],
"the turn should emit the model's reply: {records:?}"
);
}
#[test]
fn ending_emits_summary_counting_messages() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("reply")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
let count = summary_count(&records).expect("a Summary record must be emitted");
assert!(
count > 0,
"Summary should count the turn's messages, got {count}: {records:?}"
);
}
#[test]
fn loop_terminates_and_yields_none() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("done")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
drive(&mut h);
assert!(
h.agent.next_status().is_none(),
"a completed loop must keep yielding None"
);
}
#[test]
fn provider_failure_emits_failed_action_not_silent_success() {
let mut mock = MockModelProvider::new();
mock.fail_with(|_| true, "provider exploded");
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert!(
has_failed_action(&records),
"a provider failure must emit FailedAction, not an empty success: {records:?}"
);
assert!(
assistant_texts(&records).is_empty(),
"a failed turn must not emit assistant text: {records:?}"
);
}
#[test]
fn tool_call_output_drives_the_tool_path() {
let mut mock = MockModelProvider::new();
mock.on_nth_call(1, vec![mock_tool_call("search", HashMap::new())]);
mock.on_any(vec![mock_text("done after tool")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("use a tool"));
let records = drive(&mut h);
assert!(
!records.is_empty(),
"a tool-calling turn must still produce records: {records:?}"
);
}
#[test]
fn priority_message_is_processed_before_follow_up() {
use std::sync::atomic::Ordering;
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("ack")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("second"));
let _ = h.priority.push(user_msg("first"));
h.cancel.store(1, Ordering::SeqCst);
let records = drive(&mut h);
assert!(
!records.is_empty(),
"the turn should run with both queues populated: {records:?}"
);
assert_eq!(
h.priority.len(),
0,
"the priority queue must be drained by the loop"
);
assert_eq!(
h.follow_up.len(),
0,
"the follow-up queue must also be drained before ending"
);
}
#[test]
fn max_outer_iterations_terminates_a_generating_loop() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("again")]);
let config = AgentConfig {
primary_model: ModelId::Name("mock".into(), None),
max_outer_iterations: 2,
..Default::default()
};
let mut h = harness_with(mock.into_router(), config);
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert!(
summary_count(&records).is_some(),
"a capped loop must still emit its Summary: {records:?}"
);
}
#[test]
fn assistant_reply_is_persisted_to_message_api() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("persisted reply")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert!(
!assistant_texts(&records).is_empty(),
"precondition: the turn produced a reply"
);
}
#[test]
fn registered_tool_executes_and_emits_its_result() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![mock_tool_call("echo", HashMap::new())]);
mock.on_any(vec![mock_text("finished")]);
let tool = Arc::new(MockTool::returning("echo", "tool output here"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("call the tool"));
let records = drive(&mut h);
let has_tool_result = records.iter().any(|r| {
matches!(
r,
SessionRecord::Conversation {
message: Messages::ToolResult { .. }
}
)
});
assert!(
has_tool_result,
"the executed tool's result must be emitted as a record: {records:?}"
);
}
#[test]
fn failing_tool_does_not_kill_the_turn() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![mock_tool_call("broken", HashMap::new())]);
mock.on_any(vec![mock_text("recovered")]);
let tool = Arc::new(MockTool::failing("broken", "tool exploded"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("call the broken tool"));
let records = drive(&mut h);
assert!(
!records.is_empty(),
"a failing tool must still produce records: {records:?}"
);
}
#[test]
fn unknown_tool_name_errors_without_panic() {
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![mock_tool_call("does_not_exist", HashMap::new())]);
mock.on_any(vec![mock_text("moved on")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("call a missing tool"));
let records = drive(&mut h);
assert!(
!records.is_empty(),
"an unknown tool must terminate the turn cleanly: {records:?}"
);
}
#[test]
fn multiple_tool_calls_all_execute() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(
0,
vec![
mock_tool_call("alpha", HashMap::new()),
mock_tool_call("beta", HashMap::new()),
],
);
mock.on_any(vec![mock_text("both done")]);
let tools: Vec<Arc<dyn foundation_ai::agentic::tool_impl::ToolImpl>> = vec![
Arc::new(MockTool::returning("alpha", "A")),
Arc::new(MockTool::returning("beta", "B")),
];
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), tools);
let _ = h.follow_up.push(user_msg("call both"));
let records = drive(&mut h);
let tool_results = records
.iter()
.filter(|r| {
matches!(
r,
SessionRecord::Conversation {
message: Messages::ToolResult { .. }
}
)
})
.count();
assert_eq!(
tool_results, 2,
"both tool calls must execute and emit results: {records:?}"
);
}
#[test]
fn repeated_failures_trip_the_breaker_and_terminate() {
let mut mock = MockModelProvider::new();
mock.fail_with(|_| true, "always fails");
let config = AgentConfig {
primary_model: ModelId::Name("primary".into(), None),
fallback_models: vec![ModelId::Name("fallback".into(), None)],
circuit_breaker_threshold: 2,
max_outer_iterations: 3,
..Default::default()
};
let mut h = harness_with(mock.into_router(), config);
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert!(
has_failed_action(&records),
"persistent provider failure must surface as FailedAction: {records:?}"
);
}
#[test]
fn empty_router_fails_cleanly() {
let mut h = harness_with(
ProviderRouter::builder().build(),
config_for("nobody-serves-this"),
);
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert!(
has_failed_action(&records),
"an unroutable model must emit FailedAction: {records:?}"
);
assert!(
assistant_texts(&records).is_empty(),
"an unroutable turn must not emit assistant text: {records:?}"
);
}
#[test]
fn unscripted_interaction_fails_loudly() {
let mock = MockModelProvider::new();
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
let records = drive(&mut h);
assert!(
has_failed_action(&records),
"an unscripted mock must fail loudly: {records:?}"
);
}
#[test]
fn max_inner_iterations_bounds_a_non_converging_tool_loop() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_tool_call("loop_forever", HashMap::new())]);
let config = AgentConfig {
primary_model: ModelId::Name("mock".into(), None),
max_inner_iterations: 3,
max_outer_iterations: 2,
..Default::default()
};
let tool = Arc::new(MockTool::returning("loop_forever", "again"));
let mut h = harness_with_tools(mock.into_router(), config, vec![tool]);
let _ = h.follow_up.push(user_msg("loop"));
let records = drive(&mut h);
assert!(
summary_count(&records).is_some(),
"a capped inner loop must still reach Ending and emit a Summary: {records:?}"
);
}
#[test]
fn assembled_interaction_carries_system_message_and_tools() {
use foundation_ai::agentic::testing::MockTool;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc as StdArc;
let saw_system = StdArc::new(AtomicBool::new(false));
let saw_user = StdArc::new(AtomicBool::new(false));
let saw_tool = StdArc::new(AtomicBool::new(false));
let (s, u, t) = (saw_system.clone(), saw_user.clone(), saw_tool.clone());
let mut mock = MockModelProvider::new();
mock.on(
move |mi| {
if mi.system_prompt.is_some() {
s.store(true, Ordering::SeqCst);
}
if mi.messages.iter().any(|m| {
matches!(m, Messages::User { content: foundation_ai::types::UserModelContent::Text(tc), .. } if tc.content.contains("find me"))
}) {
u.store(true, Ordering::SeqCst);
}
let shed = &mi.tools_shed;
if !shed.tools.is_empty() || shed.shed.is_some() {
t.store(true, Ordering::SeqCst);
}
true
},
vec![mock_text("ok")],
);
let tool = Arc::new(MockTool::returning("search", "results"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("find me something"));
drive(&mut h);
assert!(saw_system.load(Ordering::SeqCst), "interaction must carry the system prompt");
assert!(saw_user.load(Ordering::SeqCst), "interaction must carry the user message text");
assert!(saw_tool.load(Ordering::SeqCst), "interaction must carry the registered tools");
}
#[test]
fn tool_result_is_fed_back_into_next_assemble() {
use foundation_ai::agentic::testing::MockTool;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc as StdArc;
let saw_result = StdArc::new(AtomicBool::new(false));
let flag = saw_result.clone();
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![mock_tool_call("lookup", HashMap::new())]);
mock.on(
move |mi| {
if mi.messages.iter().any(|m| matches!(
m,
Messages::ToolResult { content: foundation_ai::types::UserModelContent::Text(t), .. }
if t.content.contains("TOOL_OUTPUT_MARKER")
)) {
flag.store(true, Ordering::SeqCst);
}
true
},
vec![mock_text("done")],
);
let tool = Arc::new(MockTool::returning("lookup", "TOOL_OUTPUT_MARKER"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("use the tool"));
drive(&mut h);
assert!(
saw_result.load(Ordering::SeqCst),
"the tool's result must be fed back into the next model interaction"
);
}
#[test]
fn abort_terminates_the_loop_before_generation() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("should not be reached")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hi"));
h.cancel.store(2, std::sync::atomic::Ordering::SeqCst);
let records = drive(&mut h);
assert!(
assistant_texts(&records).is_empty(),
"an aborted turn must not generate an assistant reply: {records:?}"
);
assert!(
summary_count(&records).is_some(),
"an aborted turn still emits its Summary and terminates: {records:?}"
);
}
#[test]
fn context_pressure_note_injected_when_over_threshold() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc as StdArc;
let saw_pressure = StdArc::new(AtomicBool::new(false));
let flag = saw_pressure.clone();
let mut mock = MockModelProvider::new();
mock.on(
move |mi| {
if mi.system_prompt.as_deref().unwrap_or("").contains("capacity") {
flag.store(true, Ordering::SeqCst);
}
true
},
vec![mock_text("ok")],
);
let config = AgentConfig {
primary_model: ModelId::Name("mock".into(), None),
context_pressure_threshold: 0.70,
..Default::default()
};
let mut h = harness_with(mock.into_router(), config);
h.set_budget(10);
let long = "word ".repeat(200);
let _ = h.follow_up.push(user_msg(&long));
drive(&mut h);
assert!(
saw_pressure.load(Ordering::SeqCst),
"a context over the pressure threshold must inject the pressure note"
);
}
#[test]
fn preflight_compression_drops_oldest_when_over_budget() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc as StdArc;
let sent = StdArc::new(AtomicUsize::new(usize::MAX));
let counter = sent.clone();
let mut mock = MockModelProvider::new();
mock.on(
move |mi| {
counter.store(mi.messages.len(), Ordering::SeqCst);
true
},
vec![mock_text("ok")],
);
let config = AgentConfig {
primary_model: ModelId::Name("mock".into(), None),
preflight_compression_threshold: 0.85,
context_pressure_threshold: 0.0, ..Default::default()
};
let mut h = harness_with(mock.into_router(), config);
h.set_budget(5);
for i in 0..8 {
let _ = h.follow_up.push(user_msg(&format!("history message number {i} with some words")));
}
drive(&mut h);
let n = sent.load(Ordering::SeqCst);
assert!(n != usize::MAX, "the mock must have been called");
assert!(
n < 8,
"an over-budget context must be compressed below the full history (sent {n} of 8)"
);
assert!(n >= 1, "compression must keep at least the newest message (sent {n})");
}
fn tool_call_with_deps(
call_id: &str,
name: &str,
depends_on: Vec<String>,
) -> Messages {
Messages::Assistant {
id: foundation_compact::ids::new_scru128(),
model: ModelId::Name("mock".into(), None),
timestamp: foundation_compact::SystemTime::UNIX_EPOCH,
usage: foundation_ai::agentic::testing::zero_usage(),
content: ModelOutput::ToolCall {
id: call_id.to_string(),
name: name.to_string(),
arguments: Some(HashMap::new()),
signature: None,
depends_on,
execution_hint: foundation_ai::types::ExecutionHint::Unspecified,
},
stop_reason: foundation_ai::types::StopReason::ToolUse,
provider: foundation_ai::types::ModelProviders::Custom("mock".into()),
error_detail: None,
signature: None,
metadata: None,
}
}
#[test]
fn a_self_referential_dependency_fails_the_workflow_without_panicking() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(
0,
vec![tool_call_with_deps("call_a", "echo", vec!["call_a".into()])],
);
mock.on_any(vec![mock_text("finished")]);
let tool = Arc::new(MockTool::returning("echo", "out"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("go"));
let records = drive(&mut h);
assert!(
has_failed_action(&records),
"a cyclic dependency must surface as a FailedAction record: {records:?}"
);
}
#[test]
fn a_dependency_on_an_unknown_call_fails_the_workflow() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(
0,
vec![tool_call_with_deps("call_a", "echo", vec!["ghost".into()])],
);
mock.on_any(vec![mock_text("finished")]);
let tool = Arc::new(MockTool::returning("echo", "out"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("go"));
let records = drive(&mut h);
assert!(
has_failed_action(&records),
"a dependency on an unknown call must surface as a FailedAction: {records:?}"
);
}
#[test]
fn a_workflow_failure_does_not_end_the_turn() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(
0,
vec![tool_call_with_deps("call_a", "echo", vec!["call_a".into()])],
);
mock.on_any(vec![mock_text("recovered")]);
let tool = Arc::new(MockTool::returning("echo", "out"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("go"));
let records = drive(&mut h);
assert!(
!records.is_empty(),
"the turn must still produce records after a workflow failure"
);
}
#[test]
fn well_ordered_dependencies_still_execute() {
use foundation_ai::agentic::testing::MockTool;
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![tool_call_with_deps("call_a", "echo", vec![])]);
mock.on_any(vec![mock_text("finished")]);
let tool = Arc::new(MockTool::returning("echo", "out"));
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![tool]);
let _ = h.follow_up.push(user_msg("go"));
let records = drive(&mut h);
assert!(
!has_failed_action(&records),
"a dependency-free call must not be reported as a workflow failure: {records:?}"
);
let has_tool_result = records.iter().any(|r| {
matches!(
r,
SessionRecord::Conversation {
message: Messages::ToolResult { .. }
}
)
});
assert!(has_tool_result, "the tool must actually run: {records:?}");
}
#[test]
fn a_hard_abort_ends_the_turn_before_the_next_generation() {
use foundation_ai::agentic::CancelCode;
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("should not be reached")]);
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![]);
let _ = h.follow_up.push(user_msg("go"));
CancelCode::Abort.store(&h.cancel);
let records = drive(&mut h);
let replies = assistant_texts(&records);
assert!(
replies.is_empty(),
"an aborted turn must not produce an assistant reply: {replies:?}"
);
}
#[test]
fn an_abort_resets_the_cancel_signal() {
use foundation_ai::agentic::CancelCode;
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("hi")]);
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![]);
let _ = h.follow_up.push(user_msg("go"));
CancelCode::Abort.store(&h.cancel);
let _ = drive(&mut h);
assert_eq!(
CancelCode::load(&h.cancel),
CancelCode::None,
"the abort must be reset after it is honoured"
);
}
#[test]
fn a_priority_message_is_injected_and_the_signal_cleared() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("answered")]);
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![]);
let _ = h.follow_up.push(user_msg("original"));
h.priority
.push(user_msg("urgent"))
.expect("priority queue accepts");
let records = drive(&mut h);
assert!(
h.priority.is_empty(),
"the priority queue must be drained, not left to re-fire"
);
assert!(
!records.is_empty(),
"the turn must still make progress after steering: {records:?}"
);
}
#[test]
fn the_outer_iteration_cap_terminates_a_turn() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("more")]);
let config = AgentConfig {
primary_model: ModelId::Name("mock".into(), None),
max_outer_iterations: 1,
..Default::default()
};
let mut h = harness_with_tools(mock.into_router(), config, vec![]);
let _ = h.follow_up.push(user_msg("go"));
let records = drive(&mut h);
let _ = records;
}
#[test]
fn a_zero_outer_iteration_cap_ends_immediately() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("never")]);
let config = AgentConfig {
primary_model: ModelId::Name("mock".into(), None),
max_outer_iterations: 0,
..Default::default()
};
let mut h = harness_with_tools(mock.into_router(), config, vec![]);
let _ = h.follow_up.push(user_msg("go"));
let records = drive(&mut h);
assert!(
assistant_texts(&records).is_empty(),
"a zero cap must not permit a generation: {records:?}"
);
}
#[test]
fn a_priority_message_during_generation_discards_and_reassembles() {
use foundation_ai::agentic::AgentProgress;
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("original answer")]);
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![]);
let _ = h.follow_up.push(user_msg("first question"));
let mut injected = false;
let mut saw_mid_gen_steering = false;
for _ in 0..2_000 {
match h.agent.next_status() {
None => break,
Some(TaskStatus::Pending(AgentProgress::Generating { .. })) => {
if !injected {
h.priority
.push(user_msg("urgent interrupt"))
.expect("priority queue accepts");
injected = true;
}
}
Some(TaskStatus::Pending(AgentProgress::Steering { source })) => {
if source == "mid_gen_priority" {
saw_mid_gen_steering = true;
}
}
Some(_) => {}
}
}
assert!(injected, "the test never reached the generating state");
assert!(
saw_mid_gen_steering,
"a priority message arriving mid-generation must be reported as \
mid_gen_priority steering, not deferred to the next boundary"
);
}
#[test]
fn mid_generation_steering_drains_the_priority_queue() {
use foundation_ai::agentic::AgentProgress;
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("answer")]);
let mut h = harness_with_tools(mock.into_router(), config_for("mock"), vec![]);
let _ = h.follow_up.push(user_msg("question"));
let mut injected = false;
for _ in 0..2_000 {
match h.agent.next_status() {
None => break,
Some(TaskStatus::Pending(AgentProgress::Generating { .. })) => {
if !injected {
h.priority.push(user_msg("urgent")).expect("push");
injected = true;
}
}
Some(_) => {}
}
}
assert!(
h.priority.is_empty(),
"the priority queue must be drained, or the loop re-interrupts forever"
);
}
fn retractions(records: &[SessionRecord]) -> Vec<String> {
records
.iter()
.filter_map(|r| match r {
SessionRecord::Retracted { reason, .. } => Some(reason.clone()),
_ => None,
})
.collect()
}
#[test]
fn a_vacuous_turn_is_retried_and_the_good_answer_replaces_it() {
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![mock_text(".")]);
mock.on_any(vec![mock_text("Hello. How can I help you?")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("hello"));
let records = drive(&mut h);
let texts = assistant_texts(&records);
assert!(
texts.iter().any(|t| t.contains("How can I help you")),
"the retry's answer should reach the caller: {texts:?}"
);
assert_eq!(
retractions(&records).len(),
1,
"the loop must withdraw the turn it threw away: {records:?}"
);
}
#[test]
fn the_withdrawn_text_is_not_left_in_front_of_the_answer() {
let mut mock = MockModelProvider::new();
mock.on_nth_call(0, vec![mock_text(".")]);
mock.on_any(vec![mock_text("Paris")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("What is the capital of France?"));
let records = drive(&mut h);
let surviving: Vec<String> = records
.iter()
.skip_while(|r| !matches!(r, SessionRecord::Retracted { .. }))
.filter_map(|r| match r {
SessionRecord::Conversation {
message: Messages::Assistant { content, .. },
} => match content {
ModelOutput::Text(t) => Some(t.content.clone()),
_ => None,
},
_ => None,
})
.collect();
assert_eq!(
surviving.concat(),
"Paris",
"only the retry's answer should follow the retraction: {records:?}"
);
}
#[test]
fn a_real_answer_is_never_retried() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("Paris")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("What is the capital of France?"));
let records = drive(&mut h);
assert!(
retractions(&records).is_empty(),
"a good answer must not be withdrawn: {records:?}"
);
assert_eq!(assistant_texts(&records), vec!["Paris".to_string()]);
}
#[test]
fn a_bare_number_is_kept_when_the_question_asked_for_one() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text("4")]);
let mut h = harness_with(mock.into_router(), config_for("mock"));
let _ = h.follow_up.push(user_msg("What is 2+2?"));
let records = drive(&mut h);
assert!(
retractions(&records).is_empty(),
"a correct numeric answer must not be retried: {records:?}"
);
assert_eq!(assistant_texts(&records), vec!["4".to_string()]);
}
#[test]
fn a_model_stuck_on_junk_stops_being_asked_and_never_fails_the_turn() {
let mut mock = MockModelProvider::new();
mock.on_any(vec![mock_text(".")]);
let config = AgentConfig {
..config_for("mock")
};
let mut h = harness_with(mock.into_router(), config);
let _ = h.follow_up.push(user_msg("hello"));
let records = drive(&mut h);
assert!(
!has_failed_action(&records),
"a weak answer must not be turned into no answer: {records:?}"
);
let attempts = retractions(&records).len();
assert!(
attempts <= LoopDetectorConfig::default().max_redirects,
"the ladder should be spent once ({attempts} retries against a budget of \
{}), not refilled for each outer pass: {records:?}",
LoopDetectorConfig::default().max_redirects
);
assert!(
assistant_texts(&records).iter().any(|t| t.contains('.')),
"the last attempt should still reach the caller: {records:?}"
);
}