use super::super::*;
use super::run;
use super::*;
#[tokio::test]
async fn section_with_alternating_blocks_executes_in_order() {
async fn completions(
State(calls): State<Arc<AtomicU32>>,
Json(_body): Json<Value>,
) -> Json<Value> {
let n = calls.fetch_add(1, Ordering::Relaxed) + 1;
Json(json!({
"choices": [{
"message": {
"role": "assistant",
"content": format!("reply-{n}")
}
}]
}))
}
let calls = Arc::new(AtomicU32::new(0));
let router = Router::new()
.route("/v1/chat/completions", post(completions))
.with_state(Arc::clone(&calls));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
axum::serve(listener, router).await.unwrap();
});
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Only\n\n\
```lua\nstore.append('order.txt', 'lua1\\n')\n```\n\n\
First ask.\n\n\
```lua\nstore.append('order.txt', 'lua2\\n')\n```\n\n\
Final ask.\n\n\
```lua\nstore.append('order.txt', 'lua3\\n')\nreturn reply\n```\n";
let store = StoreRef::memory();
let out = run(
&bound_for_model(md),
"",
&[],
&store,
RunOptions {
execution: EXECUTION,
observer: Arc::new(NullObserver),
client: Some(GatewayClient::new(
GatewayEndpoint::new(&format!("http://{addr}/v1")).expect("valid test endpoint"),
SecretString::new("test").expect("non-empty test key"),
)),
debug: None,
},
)
.await
.expect("alternating blocks must execute");
assert_eq!(out, "reply-2");
assert_eq!(calls.load(Ordering::Relaxed), 2);
assert_eq!(
store.read_lines("order.txt").expect("order log"),
"1| lua1\n2| lua2\n3| lua3"
);
let section = parse(md).entry().clone();
assert_eq!(section.blocks.len(), 5);
assert!(matches!(
§ion.blocks[1],
crate::parser::Block::Prose {
loop_capable: false,
..
}
));
assert!(matches!(
§ion.blocks[3],
crate::parser::Block::Prose {
loop_capable: true,
..
}
));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn execute_runs_named_section_as_subroutine() {
async fn completions(Json(_body): Json<Value>) -> Json<Value> {
Json(json!({
"choices": [{
"message": {
"role": "assistant",
"content": "research-reply"
}
}]
}))
}
let router = Router::new().route("/v1/chat/completions", post(completions));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
axum::serve(listener, router).await.unwrap();
});
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Research\n\n\
```lua\n\
local step = tasks['## Research']\n\
assert(step.name == 'Research')\n\
assert(step.has_prose == true)\n\
```\n\n\
Research {{ args }}.\n\n\
```lua\nstore.write('evidence.md', reply)\n```\n\n\
## Main\n\n\
```lua\n\
local research = tasks['## Research']\n\
assert(research.name == 'Research')\n\
assert(research.has_prose == true)\n\
local by_name = execute('## Research')\n\
local by_obj = execute(research)\n\
assert(by_name == 'research-reply')\n\
assert(by_obj == 'research-reply')\n\
assert(store.read('evidence.md') == 'research-reply')\n\
return by_name\n\
```\n";
let store = StoreRef::memory();
let out = run(
&bound_for_model(md),
"topic",
&[],
&store,
RunOptions {
execution: EXECUTION,
observer: Arc::new(NullObserver),
client: Some(GatewayClient::new(
GatewayEndpoint::new(&format!("http://{addr}/v1")).expect("valid test endpoint"),
SecretString::new("test").expect("non-empty test key"),
)),
debug: None,
},
)
.await
.expect("execute must run named section as subroutine");
assert_eq!(out, "research-reply");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn jump_transfers_control_and_clears_context() {
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Check\n\n\
```lua\n\
store.write('seen.txt', 'check')\n\
local help = tasks['## Help']\n\
assert(help.name == 'Help')\n\
assert(help.has_prose == false)\n\
jump(help)\n\
store.write('seen.txt', 'should-not-run')\n\
```\n\n\
## Accept\n\n\
```lua\nreturn 'accepted'\n```\n\n\
## Help\n\n\
```lua\n\
assert(reply == nil, 'jump must clear prior reply context')\n\
return 'helped:' .. store.read('seen.txt')\n\
```\n";
let store = StoreRef::memory();
let out = run(&fixture(md), "", &[], &store, silent())
.await
.expect("jump must transfer control");
assert_eq!(out, "helped:check");
assert_eq!(store.read("seen.txt").expect("seen"), "check");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn jump_inside_execute_is_rejected() {
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Main\n\n\
```lua\nreturn execute('## Sub')\n```\n\n\
## Sub\n\n\
```lua\n\
local m = tasks['## Main']\n\
jump(m)\n\
```\n";
let error = run_offline(md)
.await
.expect_err("jump inside execute must fail");
let rendered = format!("{error:?}");
assert!(
rendered.contains("not allowed inside execute"),
"expected a jump-rejection error, got: {rendered}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn execute_recursion_is_capped() {
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Loop\n\n\
```lua\nreturn execute('## Loop')\n```\n";
let error = run_offline(md)
.await
.expect_err("unbounded execute recursion must fail");
let rendered = format!("{error:?}");
assert!(
rendered.contains("recursion exceeded cap"),
"expected a recursion-cap error, got: {rendered}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fanout_returns_structured_results() {
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Parent\n\n\
```lua\n\
local r = fanout('### Worker', '### Items')\n\
assert(r[1].text == 'alpha-1')\n\
assert(r[1].ok == true)\n\
assert(r[1].item == 'alpha')\n\
assert(r[1].exhausted == false)\n\
assert(r[2].text == 'beta-2')\n\
assert(r[2].ok == true)\n\
assert(r[2].item == 'beta')\n\
assert(r[2].exhausted == false)\n\
assert(tostring(r[1]) == r[1].text)\n\
assert(table.concat(r, ',') == 'alpha-1,beta-2')\n\
return 'ok'\n\
```\n\n\
### Worker\n\n\
```lua\nreturn item .. '-' .. sys.taskid\n```\n\n\
Do work.\n\n\
### Items\n\n\
- alpha\n\
- beta\n";
let out = run(&fixture(md), "", &[], &StoreRef::memory(), silent())
.await
.expect("fanout must return structured results");
assert_eq!(out, "ok");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fanout_exhausted_arm_exposes_failure_metadata() {
let gateway =
ScriptedGateway::start(vec![resp_tool_call("call_x", "echo", "{\"value\":\"x\"}")]).await;
let addr = gateway.addr();
let md = "---\nname: t\ndescription: d\npromptforge: 1\nmax_tool_iterations: 2\n---\n\n\
# Test prompt\n\n```lua shared\n\
tools.need('echo', 'echo tool')\n\
models.always('writer', 'A general model for tests')\n```\n\n\
## Parent\n\n\
```lua\n\
local r = fanout('### Worker', '### Items')\n\
assert(r[1].ok == false)\n\
assert(r[1].exhausted == true)\n\
assert(r[1].item == 'alpha')\n\
assert(r[1].text:find('tool loop exhausted', 1, true))\n\
assert(tostring(r[1]) == r[1].text)\n\
return 'ok'\n\
```\n\n\
### Worker\n\n\
```lua\ntools.add('echo')\n```\n\n\
Loop forever on {{ item }}.\n\n\
### Items\n\n\
- alpha\n";
let prompt = bound_with_tools(md, Vec::new());
let out = run(
&prompt,
"",
&[Arc::new(EchoTool) as Arc<dyn Tool>],
&StoreRef::memory(),
RunOptions {
execution: EXECUTION,
observer: Arc::new(NullObserver),
client: Some(GatewayClient::new(
GatewayEndpoint::new(&format!("http://{addr}/v1")).expect("valid test endpoint"),
SecretString::new("test").expect("non-empty test key"),
)),
debug: None,
},
)
.await
.expect("soft-degraded fanout must still return structured results");
assert_eq!(out, "ok");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn sys_exposes_section_metadata() {
let gateway = ScriptedGateway::start(vec![resp_text_finish("alpha-answer", "stop")]).await;
let addr = gateway.addr();
let md = "---\nname: t\ndescription: d\npromptforge: 1\n---\n\n\
## Alpha\n\n\
```lua\n\
assert(sys.section_name == 'Alpha')\n\
assert(sys.execution == 'execute-test')\n\
assert(sys.section_count == 2)\n\
local ok = pcall(function() return sys.reply_finish_reason end)\n\
assert(not ok, 'reply_finish_reason must be absent before prose')\n\
```\n\n\
Write one fact.\n\n\
```lua\n\
assert(sys.reply_finish_reason == 'stop')\n\
assert(reply == 'alpha-answer')\n\
```\n\n\
## Beta\n\n\
```lua\n\
assert(sys.section_name == 'Beta')\n\
assert(sys.section_count == 2)\n\
assert(sys.execution == 'execute-test')\n\
return 'done'\n\
```\n";
let out = run(
&bound_for_model(md),
"",
&[],
&StoreRef::memory(),
RunOptions {
execution: EXECUTION,
observer: Arc::new(NullObserver),
client: Some(GatewayClient::new(
GatewayEndpoint::new(&format!("http://{addr}/v1")).expect("valid test endpoint"),
SecretString::new("test").expect("non-empty test key"),
)),
debug: None,
},
)
.await
.expect("sys must expose section metadata");
assert_eq!(out, "done");
}
#[test]
fn advance_turn_saturates_and_never_wraps_the_stored_counter() {
let turns = AtomicU32::new(0);
assert_eq!(advance_turn(&turns), 1);
assert_eq!(advance_turn(&turns), 2);
assert_eq!(turns.load(Ordering::Relaxed), 2);
let maxed = AtomicU32::new(u32::MAX);
assert_eq!(advance_turn(&maxed), u32::MAX);
assert_eq!(
maxed.load(Ordering::Relaxed),
u32::MAX,
"the stored counter must not wrap to zero"
);
let near = AtomicU32::new(u32::MAX - 1);
assert_eq!(advance_turn(&near), u32::MAX);
assert_eq!(advance_turn(&near), u32::MAX);
assert_eq!(near.load(Ordering::Relaxed), u32::MAX);
}
#[test]
fn now_rfc3339_checked_produces_a_parseable_timestamp() {
let now = now_rfc3339_checked().expect("formatting the current time must succeed");
assert!(!now.is_empty(), "a formatted timestamp is never empty");
assert!(now.contains('T'), "RFC 3339 has a T separator: {now}");
assert!(
now.ends_with('Z') || now.contains('+'),
"RFC 3339 UTC has a zone designator: {now}"
);
}