use std::{
fs,
path::PathBuf,
sync::atomic::{AtomicU64, Ordering},
time::{Duration, Instant},
};
use kcode_codex_terra::{CodexTerra, ErrorKind, ToolRun};
use kcode_k1_accounting::Accounting;
use rust_decimal::Decimal;
use serde_json::{Value, json};
static NEXT_FIXTURE: AtomicU64 = AtomicU64::new(1);
fn tool_run(input: String) -> ToolRun {
ToolRun {
input,
tool_name: "extract".into(),
tool_description: "Extract one integer".into(),
input_schema: json!({
"type":"object",
"properties":{"answer":{"type":"integer"}},
"required":["answer"],
"additionalProperties":false
}),
}
}
#[cfg(unix)]
fn fixture(body: &str) -> (PathBuf, PathBuf) {
use std::os::unix::fs::PermissionsExt;
let id = NEXT_FIXTURE.fetch_add(1, Ordering::Relaxed);
let base = format!("kcode-codex-terra-{}-{id}", std::process::id());
let script_path = std::env::temp_dir().join(format!("{base}.sh"));
let log_path = std::env::temp_dir().join(format!("{base}.jsonl"));
fs::write(
&script_path,
format!("#!/bin/sh\nLOG='{}'\n{}", log_path.to_string_lossy(), body),
)
.unwrap();
let mut permissions = fs::metadata(&script_path).unwrap().permissions();
permissions.set_mode(0o700);
fs::set_permissions(&script_path, permissions).unwrap();
(script_path, log_path)
}
#[cfg(unix)]
async fn incomplete_usage_is_protocol(body: &str) {
let (script, log) = fixture(body);
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
script.clone(),
std::env::temp_dir(),
Duration::from_secs(5),
)
.unwrap();
let error = client.run(tool_run("usage".into())).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::Protocol);
let entries = accounting.entries();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].usage["input tokens"].units, Decimal::from(9));
let _ = fs::remove_file(script);
let _ = fs::remove_file(log);
}
#[test]
fn local_construction_canary() {
let accounting = Accounting::new();
let start = Instant::now();
for _ in 0..10_000 {
CodexTerra::new(
accounting.clone(),
"/tmp/codex-safe",
"/tmp",
Duration::from_secs(1),
)
.unwrap();
}
let limit = if cfg!(debug_assertions) {
Duration::from_secs(5)
} else {
Duration::from_secs(1)
};
assert!(start.elapsed() < limit);
}
#[tokio::test]
async fn tool_name_domain_is_enforced() {
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
"/path/that/does/not/exist/codex",
std::env::temp_dir(),
Duration::from_secs(1),
)
.unwrap();
for name in ["a".to_owned(), "A0_-".repeat(16)] {
let mut run = tool_run("valid".into());
run.tool_name = name;
let error = client.run(run).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::Unavailable);
}
for name in [
String::new(),
"a".repeat(65),
"has.dot".into(),
"has space".into(),
"é".into(),
] {
let mut run = tool_run("invalid".into());
run.tool_name = name;
let error = client.run(run).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::InvalidInput);
}
assert!(accounting.entries().is_empty());
}
#[cfg(unix)]
#[tokio::test]
async fn local_run_canary() {
let body = r#"
read initialize
printf '%s\n' "$initialize" >> "$LOG"
printf '%s\n' '{"id":1,"result":{}}'
read initialized
read thread_start
printf '%s\n' "$thread_start" >> "$LOG"
printf '%s\n' '{"method":"thread/started","params":{"thread":{"id":"thread-1"}}}'
printf '%s\n' '{"id":2,"result":{"thread":{"id":"thread-1"}}}'
read turn_start
printf '%s\n' "$turn_start" >> "$LOG"
printf '%s\n' '{"method":"turn/started","params":{"threadId":"thread-1","turn":{"id":"turn-1"}}}'
printf '%s\n' '{"method":"item/started","params":{"threadId":"thread-1","turnId":"turn-1","item":{"type":"userMessage"}}}'
printf '%s\n' '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-1","turnId":"turn-1","tokenUsage":{"total":{"inputTokens":100,"cachedInputTokens":20,"cacheWriteInputTokens":10,"outputTokens":40,"reasoningOutputTokens":10},"last":{"inputTokens":100,"cachedInputTokens":20,"cacheWriteInputTokens":10,"outputTokens":40,"reasoningOutputTokens":10}}}}'
printf '%s\n' '{"id":77,"method":"item/tool/call","params":{"threadId":"thread-1","turnId":"turn-1","callId":"call-1","tool":"extract","arguments":{"answer":7}}}'
printf '%s\n' '{"id":3,"result":{"turn":{"id":"turn-1"}}}'
read tool_result
printf '%s\n' "$tool_result" >> "$LOG"
printf '%s\n' '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-1","turnId":"turn-1","tokenUsage":{"total":{"inputTokens":300101,"cachedInputTokens":100020,"cacheWriteInputTokens":50,"outputTokens":240,"reasoningOutputTokens":110},"last":{"inputTokens":300001,"cachedInputTokens":100000,"cacheWriteInputTokens":40,"outputTokens":200,"reasoningOutputTokens":100}}}}'
printf '%s\n' '{"method":"turn/completed","params":{"threadId":"thread-1","turn":{"id":"turn-1","status":"completed"}}}'
"#;
let (script, log) = fixture(body);
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
script.clone(),
std::env::temp_dir(),
Duration::from_secs(5),
)
.unwrap();
let input = "x".repeat(64 * 1024);
let start = Instant::now();
let result = client.run(tool_run(input.clone())).await.unwrap();
assert!(start.elapsed() < Duration::from_secs(5));
assert_eq!(result.arguments, json!({"answer":7}));
assert_eq!(result.thread_id, "thread-1");
assert_eq!(result.turn_id, "turn-1");
assert_eq!(result.usage["input tokens"].units, Decimal::from(200_081));
assert_eq!(
result.usage["input tokens"].price_cents,
Decimal::new(800_209, 4)
);
assert_eq!(
result.usage["cached input tokens"].units,
Decimal::from(100_020)
);
assert_eq!(result.usage["output tokens"].units, Decimal::from(130));
assert_eq!(result.usage["reasoning tokens"].units, Decimal::from(110));
let entries = accounting.entries();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].usage, result.usage);
let records = fs::read_to_string(&log)
.unwrap()
.lines()
.map(|line| serde_json::from_str::<Value>(line).unwrap())
.collect::<Vec<_>>();
assert!(records.iter().all(|record| record.get("jsonrpc").is_none()));
let thread = records
.iter()
.find(|record| record.get("method") == Some(&json!("thread/start")))
.unwrap();
assert_eq!(
thread.pointer("/params/model"),
Some(&json!("gpt-5.6-terra"))
);
assert_eq!(thread.pointer("/params/ephemeral"), Some(&json!(true)));
assert_eq!(thread.pointer("/params/sandbox"), Some(&json!("readOnly")));
assert_eq!(
thread.pointer("/params/dynamicTools/0/name"),
Some(&json!("extract"))
);
let turn = records
.iter()
.find(|record| record.get("method") == Some(&json!("turn/start")))
.unwrap();
assert_eq!(
turn.pointer("/params/input"),
Some(&json!([{"type":"text","text":input}]))
);
assert_eq!(
records.last().unwrap().pointer("/result/contentItems"),
Some(&json!([{"type":"inputText","text":"ok"}]))
);
let _ = fs::remove_file(script);
let _ = fs::remove_file(log);
}
#[cfg(unix)]
#[tokio::test]
async fn completion_with_only_first_round_usage_is_protocol() {
let body = r#"
read initialize
printf '%s\n' '{"id":1,"result":{}}'
read initialized
read thread_start
printf '%s\n' '{"id":2,"result":{"thread":{"id":"thread-usage"}}}'
read turn_start
printf '%s\n' '{"id":3,"result":{"turn":{"id":"turn-usage"}}}'
printf '%s\n' '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-usage","turnId":"turn-usage","tokenUsage":{"total":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2},"last":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2}}}}'
printf '%s\n' '{"id":77,"method":"item/tool/call","params":{"threadId":"thread-usage","turnId":"turn-usage","callId":"call-usage","tool":"extract","arguments":{"answer":7}}}'
read tool_result
printf '%s\n' '{"method":"turn/completed","params":{"threadId":"thread-usage","turn":{"id":"turn-usage","status":"completed"}}}'
"#;
incomplete_usage_is_protocol(body).await;
}
#[cfg(unix)]
#[tokio::test]
async fn completion_with_duplicate_only_second_usage_is_protocol() {
let body = r#"
read initialize
printf '%s\n' '{"id":1,"result":{}}'
read initialized
read thread_start
printf '%s\n' '{"id":2,"result":{"thread":{"id":"thread-usage"}}}'
read turn_start
printf '%s\n' '{"id":3,"result":{"turn":{"id":"turn-usage"}}}'
printf '%s\n' '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-usage","turnId":"turn-usage","tokenUsage":{"total":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2},"last":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2}}}}'
printf '%s\n' '{"id":77,"method":"item/tool/call","params":{"threadId":"thread-usage","turnId":"turn-usage","callId":"call-usage","tool":"extract","arguments":{"answer":7}}}'
read tool_result
printf '%s\n' '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-usage","turnId":"turn-usage","tokenUsage":{"total":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2},"last":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2}}}}'
printf '%s\n' '{"method":"turn/completed","params":{"threadId":"thread-usage","turn":{"id":"turn-usage","status":"completed"}}}'
"#;
incomplete_usage_is_protocol(body).await;
}
#[cfg(unix)]
#[tokio::test]
async fn error_path_records_usage_and_reaps() {
let body = r#"
read initialize
printf '%s\n' '{"id":1,"result":{}}'
read initialized
read thread_start
printf '%s\n' '{"id":2,"result":{"thread":{"id":"thread-error"}}}'
read turn_start
printf '%s\n' '{"id":3,"result":{"turn":{"id":"turn-error"}}}'
printf '%s\n' '{"method":"thread/tokenUsage/updated","params":{"threadId":"thread-error","turnId":"turn-error","tokenUsage":{"total":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2},"last":{"inputTokens":10,"cachedInputTokens":1,"cacheWriteInputTokens":0,"outputTokens":4,"reasoningOutputTokens":2}}}}'
printf '%s\n' "$$" > "$LOG"
printf '%s\n' '{"method":"unexpected/event","params":{}}'
exec sleep 30
"#;
let (script, log) = fixture(body);
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
script.clone(),
std::env::temp_dir(),
Duration::from_secs(5),
)
.unwrap();
let error = client.run(tool_run("fail".into())).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::Protocol);
let entries = accounting.entries();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].usage["input tokens"].units, Decimal::from(9));
let pid = fs::read_to_string(&log).unwrap();
let status = std::process::Command::new("sh")
.arg("-c")
.arg("kill -0 \"$1\" 2>/dev/null")
.arg("sh")
.arg(pid.trim())
.status()
.unwrap();
assert!(!status.success());
let _ = fs::remove_file(script);
let _ = fs::remove_file(log);
}
#[cfg(unix)]
#[tokio::test]
async fn wrong_response_id_is_protocol_error() {
let body = r#"
read initialize
printf '%s\n' '{"id":9,"result":{}}'
exec sleep 30
"#;
let (script, log) = fixture(body);
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
script.clone(),
std::env::temp_dir(),
Duration::from_secs(5),
)
.unwrap();
let error = client.run(tool_run("bad id".into())).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::Protocol);
assert!(accounting.entries().is_empty());
let _ = fs::remove_file(script);
let _ = fs::remove_file(log);
}
#[cfg(unix)]
#[tokio::test]
async fn timeout_after_turn_start_records_unknown_usage() {
let body = r#"
read initialize
printf '%s\n' '{"id":1,"result":{}}'
read initialized
read thread_start
printf '%s\n' '{"id":2,"result":{"thread":{"id":"thread-timeout"}}}'
read turn_start
printf '%s\n' '{"id":3,"result":{"turn":{"id":"turn-timeout"}}}'
sleep 5
"#;
let (script, log) = fixture(body);
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
script.clone(),
std::env::temp_dir(),
Duration::from_millis(100),
)
.unwrap();
let error = client.run(tool_run("wait".into())).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::Timeout);
assert_eq!(accounting.entries().len(), 1);
assert!(accounting.entries()[0].usage.is_empty());
let _ = fs::remove_file(script);
let _ = fs::remove_file(log);
}
#[tokio::test]
async fn startup_failure_records_no_event() {
let accounting = Accounting::new();
let client = CodexTerra::new(
accounting.clone(),
"/path/that/does/not/exist/codex",
std::env::temp_dir(),
Duration::from_secs(1),
)
.unwrap();
let error = client.run(tool_run("input".into())).await.unwrap_err();
assert_eq!(error.kind(), ErrorKind::Unavailable);
assert!(accounting.entries().is_empty());
}