13th provider (16th registry variant). Also reconciles a naming collision the matrix flagged: this fork is house-branded "xai" (cf. xai.dev metadata, KIGI_CODE_XAI_API_KEY legacy env), so XAI_API_KEY + method xai.api_key were the GENERIC house BYOK, not x.ai/Grok. The provider table wants xai/XAI_API_KEY for Grok. Resolution (user-approved): XAI_API_KEY now keys the x.ai/Grok provider; the house BYOK primary env moves to KIGI_API_KEY, keeping XAI_API_KEY and KIGI_CODE_XAI_API_KEY as back-compat fallbacks (read_xai_api_key_env checks KIGI_API_KEY first). The xai.api_key method id is unchanged (persisted-session compat); the platform method id is the bare "xai", distinct from it. xAI spec: api.x.ai/v1, Bearer, OpenAI listing + ChatCompletions, Passthrough (docs confirm stream_options.include_usage accepted). /v1/models is minimal (ids only) and requires auth, so it doubles as the key validator (401 on bad key, no override) and metadata comes from models.dev enrichment. Live ids match the models.dev "xai" keys byte-for-byte, so restrict_to_enriched keeps the 5 tool-calling chat models (grok-4.5/4.3/4.20-0309-*/build-0.1) and drops the grok-imagine-* generators + the non-tool multi-agent model. Snapshot regenerated to include the xai provider (was stale; gen script already listed it in TARGETS). Env migration is comprehensive to avoid keying the xai platform (which would trigger a live api.x.ai fetch) or leaving house-key reads stranded: routed the trace CLI resolver + acp_agent/auth.json bridge + paste-key ext handler through the new primary; moved all leader/pager/e2e harness setters to KIGI_API_KEY; made every house-key isolation test unset KIGI_API_KEY too; updated user-facing hints to name KIGI_API_KEY. Tests: e2e proves enrichment-supplied context (wire carries none), non-vacuous tool_call restriction, bare-id round-trip under xai/, Passthrough; validation tests hit /models (401 reject, 200 accept); house_env_var_takes_precedence_over_xai pins the new precedence. Registry at 16; picker 17 rows; 4 auth arrays + xai. Review (16 findings, all fixed): caught a missed else-branch env clear in the paste-key handler (would leak the house key past a clear) and a non-hermetic credential-priority test; both fixed.
401 lines
17 KiB
Rust
401 lines
17 KiB
Rust
//! Leader soak: an in-process leader server fronting a REAL `MvpAgent`, hammered
|
|
//! by churning `LeaderClient`s until a time budget expires. Asserts the leader
|
|
//! neither leaks memory nor accumulates zombie clients, and that no response is
|
|
//! ever dropped on a live-client send (`leader.response.send_failed`).
|
|
//!
|
|
//! Duration is bounded by `LEADER_SOAK_SECS` (default 10s so an ad-hoc
|
|
//! `--ignored` run stays quick). RSS growth is bounded by
|
|
//! `LEADER_SOAK_MAX_RSS_GROWTH_MB` (default 1024). On-demand today — no CI
|
|
//! lane runs it; a real soak is the long form:
|
|
//!
|
|
//! ```bash
|
|
//! LEADER_SOAK_SECS=1200 cargo test -p kigi-shell --test test_leader_soak -- --ignored --nocapture
|
|
//! ```
|
|
|
|
#![cfg(unix)]
|
|
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use agent_client_protocol as acp;
|
|
use kigi_acp_lib::{
|
|
AcpAgentGatewayReceiver as GatewayReceiver, AcpAgentGatewaySender as GatewaySender,
|
|
LineBufferedRead,
|
|
};
|
|
use kigi_shell::agent::config::Config as AgentConfig;
|
|
use kigi_shell::agent::mvp_agent::MvpAgent;
|
|
use kigi_shell::leader::{
|
|
ClientCapabilities, ClientMode, LeaderClient, LeaderServerControlState, LeaderServerMetadata,
|
|
run_leader_server,
|
|
};
|
|
use tempfile::TempDir;
|
|
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
|
use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
|
|
use tokio_util::sync::CancellationToken;
|
|
|
|
const SIMPLEX_BUF: usize = 8 * 1024 * 1024;
|
|
|
|
fn env_u64(key: &str, default: u64) -> u64 {
|
|
std::env::var(key)
|
|
.ok()
|
|
.and_then(|v| v.parse().ok())
|
|
.unwrap_or(default)
|
|
}
|
|
|
|
/// Resident set size of THIS process (leader server + agent are in-process).
|
|
/// Copied from `kigi-codebase-graph/tests/memory_integration.rs`.
|
|
fn rss_bytes() -> Option<usize> {
|
|
#[cfg(target_os = "linux")]
|
|
{
|
|
let status = std::fs::read_to_string("/proc/self/status").ok()?;
|
|
for line in status.lines() {
|
|
if let Some(val) = line.strip_prefix("VmRSS:") {
|
|
let kb: usize = val.trim().trim_end_matches(" kB").trim().parse().ok()?;
|
|
return Some(kb * 1024);
|
|
}
|
|
}
|
|
None
|
|
}
|
|
|
|
#[cfg(target_os = "macos")]
|
|
{
|
|
use std::process::Command;
|
|
let output = Command::new("ps")
|
|
.args(["-o", "rss=", "-p", &std::process::id().to_string()])
|
|
.output()
|
|
.ok()?;
|
|
let kb: usize = String::from_utf8_lossy(&output.stdout)
|
|
.trim()
|
|
.parse()
|
|
.ok()?;
|
|
Some(kb * 1024)
|
|
}
|
|
|
|
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
|
|
{
|
|
None
|
|
}
|
|
}
|
|
|
|
/// `leader.response.send_failed` entries written by THIS process.
|
|
fn send_failed_count() -> usize {
|
|
let Some(bytes) = kigi_log::unified_log::snapshot_log() else {
|
|
return 0;
|
|
};
|
|
String::from_utf8_lossy(&bytes)
|
|
.lines()
|
|
.filter(|line| {
|
|
serde_json::from_str::<serde_json::Value>(line).is_ok_and(|entry| {
|
|
entry["msg"] == "leader.response.send_failed" && entry["pid"] == std::process::id()
|
|
})
|
|
})
|
|
.count()
|
|
}
|
|
|
|
/// Send one JSON-RPC request through a `LeaderClient` and await the response
|
|
/// with the matching id, skipping interleaved notifications.
|
|
async fn rpc(client: &mut LeaderClient, payload: String, id: u64, what: &str) -> serde_json::Value {
|
|
client
|
|
.send(payload)
|
|
.unwrap_or_else(|e| panic!("{what}: send failed: {e}"));
|
|
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
|
|
loop {
|
|
let remaining = deadline
|
|
.checked_duration_since(tokio::time::Instant::now())
|
|
.unwrap_or_else(|| panic!("{what}: timed out waiting for response id {id}"));
|
|
let msg = tokio::time::timeout(remaining, client.recv())
|
|
.await
|
|
.unwrap_or_else(|_| panic!("{what}: timed out waiting for response id {id}"))
|
|
.unwrap_or_else(|| panic!("{what}: connection closed awaiting response id {id}"));
|
|
let json: serde_json::Value = match serde_json::from_str(&msg) {
|
|
Ok(v) => v,
|
|
Err(_) => continue,
|
|
};
|
|
if json["id"] == id && (json.get("result").is_some() || json.get("error").is_some()) {
|
|
assert!(
|
|
json.get("error").is_none(),
|
|
"{what}: error response: {json}"
|
|
);
|
|
return json;
|
|
}
|
|
}
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
#[ignore = "leader soak; run with --ignored (LEADER_SOAK_SECS bounds the duration)"]
|
|
async fn leader_soak_churning_clients_no_leaks_no_zombies() {
|
|
let _ = rustls::crypto::ring::default_provider().install_default();
|
|
|
|
let server = kigi_test_support::MockInferenceServer::start()
|
|
.await
|
|
.unwrap();
|
|
let kigi_home = TempDir::new().unwrap();
|
|
let workdir = TempDir::new().unwrap();
|
|
|
|
// SAFETY: single-threaded current-thread runtime; set before any agent
|
|
// code reads these process-globals (same pattern as session_load_perf).
|
|
unsafe {
|
|
std::env::set_var("KIGI_SHARE_DIR", kigi_home.path());
|
|
std::env::set_var("KIGI_CODE_BASE_URL", server.url());
|
|
std::env::set_var("KIGI_API_BASE_URL", server.url());
|
|
std::env::set_var("KIGI_API_KEY", "test-key-for-ci");
|
|
std::env::set_var("KIGI_TELEMETRY_ENABLED", "false");
|
|
std::env::set_var("KIGI_FEEDBACK_ENABLED", "false");
|
|
std::env::set_var("KIGI_TRACE_UPLOAD", "false");
|
|
}
|
|
|
|
let sock_path = kigi_home.path().join("leader-soak.sock");
|
|
let soak_secs = env_u64("LEADER_SOAK_SECS", 10);
|
|
let max_growth_mb = env_u64("LEADER_SOAK_MAX_RSS_GROWTH_MB", 1024);
|
|
let send_failed_before = send_failed_count();
|
|
|
|
let local = tokio::task::LocalSet::new();
|
|
local
|
|
.run_until(async {
|
|
// ── Leader server (survives client churn) ────────────────────
|
|
let (acp_tx, mut acp_rx) = tokio::sync::mpsc::unbounded_channel::<String>();
|
|
let (response_tx, response_rx) = tokio::sync::mpsc::unbounded_channel::<String>();
|
|
let cancel = CancellationToken::new();
|
|
let client_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let control_state = LeaderServerControlState::new(LeaderServerMetadata {
|
|
pid: std::process::id(),
|
|
socket_path: sock_path.clone(),
|
|
lock_path: sock_path.with_extension("lock"),
|
|
socket_suffix: String::new(),
|
|
leader_binary_version: env!("CARGO_PKG_VERSION").to_string(),
|
|
});
|
|
let cancel_for_server = cancel.clone();
|
|
let sock_for_server = sock_path.clone();
|
|
let client_count_for_server = client_count.clone();
|
|
tokio::task::spawn_local(async move {
|
|
let _ = run_leader_server(
|
|
sock_for_server,
|
|
acp_tx,
|
|
response_rx,
|
|
cancel_for_server,
|
|
true,
|
|
client_count_for_server,
|
|
Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
|
kigi_shell::agent::activity::AgentActivity::default(),
|
|
tokio::sync::watch::channel(true).1,
|
|
tokio::sync::watch::channel(kigi_shell::leader::ShutdownReason::Manual).0,
|
|
None,
|
|
control_state,
|
|
)
|
|
.await;
|
|
});
|
|
|
|
// ── Real agent behind it ──────────────────────────────────────
|
|
// Copied from `run_leader`'s agent-spawn + IPC/stdout bridge
|
|
// blocks in src/agent/app.rs (inside its LocalSet body); kept as
|
|
// a deliberate copy so production stays untouched. Second copy of
|
|
// the same wiring: kigi-tui/src/app/leader_cluster/mod.rs
|
|
// (`spawn_leader_generation`) — keep the two copies behaviorally
|
|
// identical.
|
|
let (agent_in_read, agent_in_write) = tokio::io::simplex(SIMPLEX_BUF);
|
|
let (agent_out_read, agent_out_write) = tokio::io::simplex(SIMPLEX_BUF);
|
|
|
|
tokio::task::spawn_local(async move {
|
|
let agent_config = AgentConfig::default();
|
|
let auth_manager = Arc::new(agent_config.create_auth_manager());
|
|
let (gw_tx, gw_rx) = tokio::sync::mpsc::unbounded_channel();
|
|
let gateway = GatewaySender::new(gw_tx);
|
|
let agent = MvpAgent::new(gateway, &agent_config, auth_manager, None)
|
|
.expect("valid agent config");
|
|
let incoming = LineBufferedRead::spawn_local(agent_in_read.compat());
|
|
let (conn, handle_io) = acp::AgentSideConnection::new(
|
|
agent,
|
|
agent_out_write.compat_write(),
|
|
incoming,
|
|
|fut| {
|
|
tokio::task::spawn_local(fut);
|
|
},
|
|
);
|
|
tokio::task::spawn_local(
|
|
GatewayReceiver::new(gw_rx, conn)
|
|
.with_on_meta(kigi_file_utils::trace_context::span_from_meta_traceparent)
|
|
.run(),
|
|
);
|
|
let _ = handle_io.await;
|
|
});
|
|
|
|
// Leader → agent stdin.
|
|
tokio::task::spawn_local(async move {
|
|
let mut agent_in_write = agent_in_write;
|
|
while let Some(msg) = acp_rx.recv().await {
|
|
if agent_in_write.write_all(msg.as_bytes()).await.is_err()
|
|
|| agent_in_write.write_all(b"\n").await.is_err()
|
|
{
|
|
break;
|
|
}
|
|
}
|
|
});
|
|
// Agent stdout → leader responses.
|
|
let response_tx_for_agent = response_tx.clone();
|
|
tokio::task::spawn_local(async move {
|
|
let mut reader = BufReader::new(agent_out_read);
|
|
let mut line = String::new();
|
|
loop {
|
|
line.clear();
|
|
match reader.read_line(&mut line).await {
|
|
Ok(0) => break,
|
|
Ok(_) => {
|
|
let msg = line.trim_end_matches(['\r', '\n']).to_string();
|
|
if !msg.is_empty() {
|
|
let _ = response_tx_for_agent.send(msg);
|
|
}
|
|
}
|
|
Err(_) => break,
|
|
}
|
|
}
|
|
});
|
|
|
|
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
|
|
while !sock_path.exists() && tokio::time::Instant::now() < deadline {
|
|
tokio::time::sleep(Duration::from_millis(20)).await;
|
|
}
|
|
assert!(sock_path.exists(), "leader socket never bound");
|
|
|
|
// ── One-time initialize + authenticate through the leader ────
|
|
let mut bootstrap = LeaderClient::connect(
|
|
sock_path.clone(),
|
|
"soak-bootstrap",
|
|
ClientMode::Stdio,
|
|
ClientCapabilities::default(),
|
|
)
|
|
.await
|
|
.expect("bootstrap connect");
|
|
rpc(
|
|
&mut bootstrap,
|
|
r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":1,"clientCapabilities":{"fs":{"readTextFile":false,"writeTextFile":false},"terminal":false},"_meta":{"startupHints":{"nonInteractive":true,"skipGitStatus":true,"skipProjectLayout":true},"clientType":"soak","clientVersion":"0.0.0-test"}}}"#.to_string(),
|
|
1,
|
|
"initialize",
|
|
)
|
|
.await;
|
|
rpc(
|
|
&mut bootstrap,
|
|
r#"{"jsonrpc":"2.0","id":2,"method":"authenticate","params":{"methodId":"xai.api_key","_meta":{"headless":true}}}"#.to_string(),
|
|
2,
|
|
"authenticate",
|
|
)
|
|
.await;
|
|
|
|
let rss_baseline = rss_bytes();
|
|
let soak_deadline = tokio::time::Instant::now() + Duration::from_secs(soak_secs);
|
|
let workdir_str = workdir.path().to_string_lossy().to_string();
|
|
let mut cycles: u64 = 0;
|
|
let mut turns: u64 = 0;
|
|
|
|
// ── Churn: 10 fresh clients per cycle, 2 sessions each, one
|
|
// scripted turn per session, then all disconnect ───────────────
|
|
while tokio::time::Instant::now() < soak_deadline {
|
|
cycles += 1;
|
|
let mut clients = Vec::new();
|
|
for i in 0..10u64 {
|
|
let client = LeaderClient::connect(
|
|
sock_path.clone(),
|
|
"soak-client",
|
|
ClientMode::Stdio,
|
|
ClientCapabilities::default(),
|
|
)
|
|
.await
|
|
.unwrap_or_else(|e| panic!("cycle {cycles} client {i} connect: {e}"));
|
|
clients.push(client);
|
|
}
|
|
|
|
for (i, client) in clients.iter_mut().enumerate() {
|
|
for s in 0..2u64 {
|
|
let new_id = 100 + s;
|
|
let resp = rpc(
|
|
client,
|
|
format!(
|
|
r#"{{"jsonrpc":"2.0","id":{new_id},"method":"session/new","params":{{"cwd":"{workdir_str}","mcpServers":[]}}}}"#
|
|
),
|
|
new_id,
|
|
"session/new",
|
|
)
|
|
.await;
|
|
let sid = resp["result"]["sessionId"]
|
|
.as_str()
|
|
.unwrap_or_else(|| panic!("no sessionId in {resp}"))
|
|
.to_string();
|
|
|
|
let prompt_id = 200 + s;
|
|
rpc(
|
|
client,
|
|
format!(
|
|
r#"{{"jsonrpc":"2.0","id":{prompt_id},"method":"session/prompt","params":{{"sessionId":"{sid}","prompt":[{{"type":"text","text":"soak c{i} s{s} cycle {cycles}"}}]}}}}"#
|
|
),
|
|
prompt_id,
|
|
"session/prompt",
|
|
)
|
|
.await;
|
|
turns += 1;
|
|
}
|
|
}
|
|
|
|
// Churn: everyone disconnects; the roster must drain fully.
|
|
for client in clients {
|
|
client.cancel();
|
|
}
|
|
let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(30);
|
|
while client_count.load(std::sync::atomic::Ordering::Relaxed) > 1 {
|
|
assert!(
|
|
tokio::time::Instant::now() < drain_deadline,
|
|
"cycle {cycles}: roster kept {} zombie clients after churn",
|
|
client_count.load(std::sync::atomic::Ordering::Relaxed)
|
|
);
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
}
|
|
}
|
|
|
|
eprintln!("[soak] {cycles} cycles, {turns} turns in {soak_secs}s budget");
|
|
assert!(cycles > 0, "soak budget too small to complete one cycle");
|
|
|
|
// ── Convergence: only the bootstrap client remains, and the
|
|
// leader still serves a healthy round-trip ────────────────────
|
|
assert_eq!(
|
|
client_count.load(std::sync::atomic::Ordering::Relaxed),
|
|
1,
|
|
"roster must converge to the bootstrap client after churn"
|
|
);
|
|
let resp = rpc(
|
|
&mut bootstrap,
|
|
format!(
|
|
r#"{{"jsonrpc":"2.0","id":900,"method":"session/new","params":{{"cwd":"{workdir_str}","mcpServers":[]}}}}"#
|
|
),
|
|
900,
|
|
"post-soak session/new",
|
|
)
|
|
.await;
|
|
assert!(resp["result"]["sessionId"].is_string());
|
|
|
|
// ── No response was ever dropped on a live-client send ────────
|
|
assert_eq!(
|
|
send_failed_count(),
|
|
send_failed_before,
|
|
"leader.response.send_failed must not occur during the soak"
|
|
);
|
|
|
|
// ── RSS bound ─────────────────────────────────────────────────
|
|
if let (Some(before), Some(after)) = (rss_baseline, rss_bytes()) {
|
|
let growth_mb = after.saturating_sub(before) as f64 / (1024.0 * 1024.0);
|
|
eprintln!(
|
|
"[soak] rss: {:.1} MB -> {:.1} MB (growth {growth_mb:.1} MB)",
|
|
before as f64 / (1024.0 * 1024.0),
|
|
after as f64 / (1024.0 * 1024.0),
|
|
);
|
|
assert!(
|
|
growth_mb < max_growth_mb as f64,
|
|
"leader RSS grew {growth_mb:.1} MB over the soak (bound {max_growth_mb} MB)"
|
|
);
|
|
} else {
|
|
eprintln!("[soak] rss measurement unavailable on this platform; bound skipped");
|
|
}
|
|
|
|
bootstrap.cancel();
|
|
cancel.cancel();
|
|
})
|
|
.await;
|
|
}
|