use crate::agentloop::action::SelfHandler;
use crate::agentloop::runner::{LoopAbort, LoopInput, Session, run_loop};
use crate::agentloop::stop::{Outcome, TerminalStatus};
use crate::config::SwapPolicy;
use crate::intel::client::{IntelClient, IntelHealthReport};
use crate::json::frame;
use crate::mcp::client::McpClient;
use crate::obs::log::{Comp, Level, LogCtx, Logger};
use crate::subagent::protocol::{AgentMsg, ControlMsg, IntelActive, SpawnPayload, SwapIntel};
use crate::supervisor::budget::Budget;
use std::io::{self, BufReader, Stdin, Stdout};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
type PendingSwap = Arc<Mutex<Option<SwapIntel>>>;
pub(crate) type Up = Arc<Mutex<Stdout>>;
const ELICITATION_TIMEOUT: Duration = Duration::from_secs(300);
const GATED_TOOL_TIMEOUT: Duration = Duration::from_secs(24 * 3600);
struct NoSelfTools {
gated: Vec<String>,
bridge: Option<GatedBridge>,
}
struct GatedBridge {
up: Up,
replies: Arc<crate::subagent::replies::Replies>,
cancel: Arc<AtomicBool>,
timeout: Duration,
}
impl SelfHandler for NoSelfTools {
fn tools(&self) -> Vec<crate::wire::intel::ToolDef> {
Vec::new()
}
fn handle(&mut self, name: &str, args: &serde_json::Value) -> Option<(String, bool)> {
if !self.gated.iter().any(|g| g == name) {
return None;
}
let b = self.bridge.as_ref()?;
let id = b.replies.next_id();
send_up(
&b.up,
&AgentMsg::ToolRequest {
id,
name: name.to_string(),
args: args.clone(),
},
);
let deadline = std::time::Instant::now() + b.timeout;
match b.replies.wait(id, deadline, &b.cancel) {
Some(crate::subagent::replies::Reply::Tool { result, is_error }) => {
let text = match &result {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
};
Some((text, is_error))
}
_ => Some((
format!("tool {name:?} is gated by policy and the supervisor did not answer"),
true,
)),
}
}
}
pub fn run() -> i32 {
install_pdeathsig();
#[cfg(unix)]
if unsafe { libc::getppid() } == 1 {
return crate::exit::GENERIC;
}
let mut stdin = BufReader::new(io::stdin());
let payload = match read_spawn(&mut stdin) {
Ok(p) => p,
Err(e) => {
eprintln!("agentd subagent: bad spawn payload: {e}");
return crate::exit::USAGE;
}
};
#[cfg(feature = "tls")]
if let Some(path) = payload.tls_ca.as_deref()
&& let Err(e) = std::fs::read(path).and_then(|pem| crate::net::tls::install_extra_ca(&pem))
{
eprintln!("agentd subagent: --tls-ca {path}: {e}");
return crate::exit::USAGE;
}
#[cfg(feature = "aauth")]
if let Some(settings) = &payload.aauth
&& let Err(e) = crate::aauth::setup(settings, std::time::Duration::from_secs(30))
{
eprintln!("agentd subagent: aauth: {e}");
return crate::exit::USAGE;
}
let up: Up = Arc::new(Mutex::new(io::stdout()));
let log = build_logger(&payload);
let cancel = Arc::new(AtomicBool::new(false));
let paused = Arc::new(AtomicBool::new(false));
let (inject_tx, inject_rx) = std::sync::mpsc::channel::<String>();
let pending_swap: PendingSwap = Arc::new(Mutex::new(None));
let replies = Arc::new(crate::subagent::replies::Replies::new());
spawn_control_thread(
stdin,
Arc::clone(&up),
Arc::clone(&cancel),
Arc::clone(&paused),
inject_tx,
Arc::clone(&pending_swap),
Arc::clone(&replies),
log.ctx().clone(),
);
let inject_rx = Some(inject_rx);
send_up(&up, &AgentMsg::Ready);
log.info(
"loop.start",
serde_json::json!({"depth": payload.depth, "warm": payload.warm}),
);
let mut intel = match IntelClient::from_parts(
&payload.intelligence.uri,
payload.intelligence.token.clone(),
) {
Ok(c) => {
#[allow(unused_mut)]
let mut c = c
.with_headers(payload.intelligence.headers.clone())
.with_dialect(payload.intelligence.dialect.as_deref());
#[cfg(feature = "oauth")]
if let Some(aws) = &payload.intelligence.aws_auth
&& let Ok(s) = crate::auth::aws::SigV4Signer::from_spec(aws, "intelligence")
{
c = c.with_signer(Some(s as std::sync::Arc<dyn ::mcp::http::RequestSigner>));
}
c.set_trace_id(payload.telemetry.trace_id.clone());
install_intel_health_reporter(&mut c, &up);
c
}
Err(e) => {
return fail(
&up,
&log,
format!("intel: {e}"),
crate::exit::INTEL_UNAVAILABLE,
);
}
};
let mut servers = Vec::new();
for spec in &payload.mcp_servers {
let elicit: Arc<dyn ::mcp::inbound::Handler> =
Arc::new(crate::mcp::elicit::ElicitationBridge::new(
Arc::clone(&up),
Arc::clone(&replies),
Arc::clone(&cancel),
ELICITATION_TIMEOUT,
));
let connected = crate::mcp::from_spec(spec, Duration::from_secs(60))
.map(|c| c.with_elicitation(elicit))
.and_then(|mut c| c.initialize().map(|()| c));
match connected {
Ok(mut c) => {
log.info("mcp.connect", serde_json::json!({"server": spec.name}));
let mut meta = serde_json::json!({"agent/run_id": payload.telemetry.run_id});
if let Some(tid) = &payload.telemetry.trace_id {
meta["traceparent"] = crate::obs::trace::outbound_traceparent(tid).into();
}
c.set_tool_meta(meta);
servers.push(c);
}
Err(e) => {
return fail(
&up,
&log,
format!("mcp '{}': {e}", spec.name),
crate::exit::MCP_REQUIRED_DOWN,
);
}
}
}
if payload.role == crate::subagent::protocol::Role::Turn {
return crate::runtime::worker::run_turn_child(
&payload, &intel, &servers, &up, &cancel, &replies, &log,
);
}
let mut input = LoopInput {
instruction: payload.instruction.clone(),
output_contract: payload.output_contract.clone(),
seed: payload
.context_seed
.iter()
.map(|m| (m.role.clone(), m.content.clone()))
.collect(),
model: payload.intelligence.model.clone().unwrap_or_default(),
max_steps: payload.limits.max_steps,
max_tokens: payload.limits.max_tokens,
deadline: Instant::now() + Duration::from_millis(payload.limits.deadline_ms.max(1)),
cancel: Some(Arc::clone(&cancel)),
};
let mut orch = NoSelfTools {
gated: payload.gated_tools.clone(),
bridge: (!payload.gated_tools.is_empty()).then(|| GatedBridge {
up: Arc::clone(&up),
replies: Arc::clone(&replies),
cancel: Arc::clone(&cancel),
timeout: GATED_TOOL_TIMEOUT,
}),
};
if payload.warm {
intel.enable_alldown_backoff(crate::intel::client::AllDownPolicy::default());
return run_warm(
intel,
&servers,
&input,
&payload,
&mut orch,
&cancel,
&paused,
inject_rx
.as_ref()
.expect("a warm session keeps its inject stream"),
&pending_swap,
&up,
&log,
);
}
pause_wait(&paused, &cancel, &log);
apply_pending_swap(&pending_swap, &mut intel, &mut input.model, &up, &log);
match run_loop(&intel, &servers, &input, &mut orch, &log) {
Ok((outcome, usage)) => {
let code = crate::exit::once_exit(outcome.status, outcome.partial);
send_up(&up, &AgentMsg::Usage(usage));
send_up(&up, &AgentMsg::Result { outcome });
code
}
Err(LoopAbort::Intel(m)) => fail(
&up,
&log,
format!("intel: {m}"),
crate::exit::INTEL_UNAVAILABLE,
),
Err(LoopAbort::Mcp(m)) => fail(
&up,
&log,
format!("mcp: {m}"),
crate::exit::MCP_REQUIRED_DOWN,
),
}
}
#[allow(clippy::too_many_arguments)]
fn run_warm(
mut intel: IntelClient,
servers: &[McpClient],
input: &LoopInput,
payload: &SpawnPayload,
orch: &mut NoSelfTools,
cancel: &Arc<AtomicBool>,
paused: &Arc<AtomicBool>,
inject_rx: &Receiver<String>,
pending_swap: &PendingSwap,
up: &Up,
log: &Logger,
) -> i32 {
let mut session = match Session::prepare(servers, input, orch) {
Ok(s) => s,
Err(LoopAbort::Intel(m)) => {
return fail(
up,
log,
format!("intel: {m}"),
crate::exit::INTEL_UNAVAILABLE,
);
}
Err(LoopAbort::Mcp(m)) => {
return fail(up, log, format!("mcp: {m}"), crate::exit::MCP_REQUIRED_DOWN);
}
};
let limits = &payload.limits;
loop {
pause_wait(paused, cancel, log);
refresh_tools_if_changed(&mut session, orch, servers, log);
apply_pending_swap_warm(pending_swap, &mut intel, &mut session, up, log);
let pre_turn = session.transcript_len();
let deadline = Instant::now() + Duration::from_millis(limits.deadline_ms.max(1));
let mut budget = Budget::new(limits.max_steps, limits.max_tokens, deadline);
let (outcome, usage) = match session.run_turn(&intel, orch, log, &mut budget, Some(cancel))
{
Ok(ou) => ou,
Err(LoopAbort::Intel(m)) => {
return fail(
up,
log,
format!("intel: {m}"),
crate::exit::INTEL_UNAVAILABLE,
);
}
Err(LoopAbort::Mcp(m)) => {
return fail(up, log, format!("mcp: {m}"), crate::exit::MCP_REQUIRED_DOWN);
}
};
if outcome.status == TerminalStatus::Cancelled {
break;
}
if restart_turn_pending(pending_swap, session.model()) {
session.truncate_transcript(pre_turn);
log.info(
"intel.swap.restart_turn",
serde_json::json!({"discarded_turn": true}),
);
continue;
}
send_up(up, &AgentMsg::Usage(usage));
send_up(up, &AgentMsg::Turn { outcome });
if cancel.load(Ordering::Relaxed) {
break;
}
match wait_for_inject(inject_rx, cancel) {
Some(message) => {
log.info(
"subagent.inject",
serde_json::json!({"bytes": message.len()}),
);
session.deliver(&message);
}
None => break, }
}
let status = if cancel.load(Ordering::Relaxed) {
TerminalStatus::Cancelled
} else {
TerminalStatus::Completed
};
let code = crate::exit::once_exit(status, false);
send_up(
up,
&AgentMsg::Result {
outcome: Outcome {
status,
partial: false,
result: serde_json::Value::Null,
scheduled: Vec::new(),
subscriptions: Vec::new(),
},
},
);
code
}
fn wait_for_inject(rx: &Receiver<String>, cancel: &AtomicBool) -> Option<String> {
loop {
if cancel.load(Ordering::Relaxed) {
return None;
}
match rx.recv_timeout(Duration::from_millis(200)) {
Ok(message) => return Some(message),
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => return None,
}
}
}
fn pause_wait(paused: &AtomicBool, cancel: &AtomicBool, log: &Logger) {
if !paused.load(Ordering::Relaxed) || cancel.load(Ordering::Relaxed) {
return; }
log.info("loop.paused", serde_json::json!({}));
while paused.load(Ordering::Relaxed) && !cancel.load(Ordering::Relaxed) {
std::thread::sleep(Duration::from_millis(50));
}
log.info("loop.resumed", serde_json::json!({}));
}
fn rebuild_intel(swap: &SwapIntel, old: &IntelClient, log: &Logger) -> Option<IntelClient> {
match IntelClient::from_parts(&swap.uri, swap.token.clone()) {
Ok(mut c) => {
c.set_trace_id(old.trace_id().map(str::to_string));
if old.alldown_enabled() {
c.enable_alldown_backoff(crate::intel::client::AllDownPolicy::default());
}
Some(c)
}
Err(e) => {
log.warn(
"intel.swap.reject",
serde_json::json!({"err": e.to_string()}),
);
None
}
}
}
fn log_swap(
log: &Logger,
from_model: &str,
to_model: &str,
endpoint_change: bool,
policy: SwapPolicy,
) {
let kind = if from_model != to_model {
"model"
} else {
"endpoint"
};
log.info(
"intel.swap",
serde_json::json!({
"kind": kind,
"model_from": from_model,
"model_to": to_model,
"endpoint_change": endpoint_change,
"policy": policy.as_str(),
}),
);
}
fn apply_pending_swap(
pending: &PendingSwap,
intel: &mut IntelClient,
model: &mut String,
up: &Up,
log: &Logger,
) {
let Some(swap) = pending.lock().unwrap_or_else(|e| e.into_inner()).take() else {
return; };
let from_model = model.clone();
let to_model = swap.model.clone().unwrap_or_else(|| from_model.clone());
let endpoint_change = match rebuild_intel(&swap, intel, log) {
Some(mut c) => {
install_intel_health_reporter(&mut c, up);
*intel = c;
true
}
None => false,
};
*model = to_model.clone();
log_swap(log, &from_model, &to_model, endpoint_change, swap.policy);
}
fn apply_pending_swap_warm(
pending: &PendingSwap,
intel: &mut IntelClient,
session: &mut Session<'_>,
up: &Up,
log: &Logger,
) {
let Some(swap) = pending.lock().unwrap_or_else(|e| e.into_inner()).take() else {
return; };
let from_model = session.model().to_string();
let to_model = swap.model.clone().unwrap_or_else(|| from_model.clone());
let endpoint_change = match rebuild_intel(&swap, intel, log) {
Some(mut c) => {
install_intel_health_reporter(&mut c, up);
*intel = c;
true
}
None => false,
};
session.set_model(&to_model);
log_swap(log, &from_model, &to_model, endpoint_change, swap.policy);
}
fn restart_turn_pending(pending: &PendingSwap, current_model: &str) -> bool {
let guard = pending.lock().unwrap_or_else(|e| e.into_inner());
match guard.as_ref() {
Some(swap) if swap.policy == SwapPolicy::RestartTurn => {
swap.model.as_deref().is_some_and(|m| m != current_model)
}
_ => false,
}
}
fn fail(up: &Up, log: &Logger, error: String, code: i32) -> i32 {
log.error("loop.error", serde_json::json!({"err": error}));
send_up(up, &AgentMsg::Failed { error });
code
}
fn refresh_tools_if_changed(
session: &mut Session,
orch: &mut NoSelfTools,
servers: &[McpClient],
log: &Logger,
) {
use crate::wire::mcp::method;
let changed = servers.iter().any(|s| {
s.drain_notifications()
.iter()
.any(|n| n.method == method::NOTIFY_TOOLS_LIST_CHANGED)
});
if !changed {
return;
}
match session.refresh_tools(orch) {
Ok(()) => log.info(
"mcp.tools_refreshed",
serde_json::json!({"tools": session.tools_len()}),
),
Err(e) => {
let msg = match e {
LoopAbort::Intel(m) | LoopAbort::Mcp(m) => m,
};
log.warn("mcp.tools_refresh_failed", serde_json::json!({"err": msg}));
}
}
}
fn read_spawn(reader: &mut BufReader<Stdin>) -> Result<SpawnPayload, String> {
let bytes = frame::read_frame(reader)
.map_err(|e| e.to_string())?
.ok_or_else(|| "stdin closed before spawn payload".to_string())?;
match serde_json::from_slice::<ControlMsg>(&bytes).map_err(|e| e.to_string())? {
ControlMsg::Spawn(p) => Ok(*p),
other => Err(format!(
"first frame was not Spawn (got {})",
control_msg_label(&other)
)),
}
}
fn control_msg_label(msg: &ControlMsg) -> &'static str {
match msg {
ControlMsg::Spawn(_) => "spawn",
ControlMsg::Ping { .. } => "ping",
ControlMsg::Pause => "pause",
ControlMsg::Resume => "resume",
ControlMsg::Cancel { .. } => "cancel",
ControlMsg::Inject { .. } => "inject",
ControlMsg::SwapIntel(_) => "swap_intel",
ControlMsg::ToolResult { .. } => "tool_result",
ControlMsg::BudgetGrant { .. } => "budget_grant",
}
}
fn build_logger(payload: &SpawnPayload) -> Logger {
let t = &payload.telemetry;
let level = Level::parse(&t.log_level).unwrap_or(Level::Info);
Logger::new(
LogCtx {
run_id: t.run_id.clone(),
agent_id: t.agent_id.clone(),
agent_path: t.agent_path.clone(),
comp: Comp::Agent,
pid: std::process::id(),
trace_id: t.trace_id.clone(),
},
level,
)
.with_content(t.log_content)
}
pub(crate) fn send_up(up: &Up, msg: &AgentMsg) {
if let Ok(mut out) = up.lock() {
let _ = frame::write_frame(&mut *out, msg);
}
}
fn install_intel_health_reporter(intel: &mut IntelClient, up: &Up) {
let up = Arc::clone(up);
intel.set_health_reporter(Box::new(move |r: IntelHealthReport| {
let active = r.active.map(|(index, transport)| IntelActive {
index,
transport: transport.to_string(),
});
send_up(
&up,
&AgentMsg::IntelHealth {
all_down: r.all_down,
active,
},
);
}));
}
#[allow(clippy::too_many_arguments)]
fn spawn_control_thread(
mut stdin: BufReader<Stdin>,
up: Up,
cancel: Arc<AtomicBool>,
paused: Arc<AtomicBool>,
inject_tx: Sender<String>,
pending_swap: PendingSwap,
replies: Arc<crate::subagent::replies::Replies>,
ctx: LogCtx,
) {
let log = Logger::new(ctx, Level::Debug);
std::thread::Builder::new()
.name("subagent-control".into())
.spawn(move || {
while let Ok(Some(bytes)) = frame::read_frame(&mut stdin) {
match serde_json::from_slice::<ControlMsg>(&bytes) {
Ok(ControlMsg::ToolResult {
id,
result,
is_error,
}) => {
replies.deliver(
id,
crate::subagent::replies::Reply::Tool { result, is_error },
);
}
Ok(ControlMsg::BudgetGrant {
id,
ok,
wait_ms,
model,
reason,
}) => {
replies.deliver(
id,
crate::subagent::replies::Reply::Budget {
ok,
wait_ms,
model,
reason,
},
);
}
Ok(ControlMsg::Ping { seq }) => send_up(&up, &AgentMsg::Pong { seq }),
Ok(ControlMsg::Cancel { reason }) => {
log.info("subagent.cancel", serde_json::json!({"reason": reason}));
cancel.store(true, Ordering::Relaxed);
}
Ok(ControlMsg::Pause) => paused.store(true, Ordering::Relaxed),
Ok(ControlMsg::Resume) => paused.store(false, Ordering::Relaxed),
Ok(ControlMsg::Inject { message }) => {
let _ = inject_tx.send(message);
}
Ok(ControlMsg::SwapIntel(swap)) => {
log.info(
"subagent.swap_intel",
serde_json::json!({"endpoint_change": true}),
);
*pending_swap.lock().unwrap_or_else(|e| e.into_inner()) = Some(*swap);
}
Ok(ControlMsg::Spawn(_)) | Err(_) => { }
}
}
replies.close();
})
.ok();
}
#[cfg(target_os = "linux")]
fn install_pdeathsig() {
unsafe {
libc::prctl(
libc::PR_SET_PDEATHSIG,
libc::SIGKILL as libc::c_ulong,
0,
0,
0,
);
}
}
#[cfg(not(target_os = "linux"))]
fn install_pdeathsig() {
}
#[cfg(test)]
mod tests {
use super::*;
use crate::obs::log::{Comp, Level, LogCtx, Logger};
fn test_log() -> Logger {
Logger::new(
LogCtx {
run_id: "r".into(),
agent_id: "0".into(),
agent_path: "0".into(),
comp: Comp::Agent,
pid: 0,
trace_id: None,
},
Level::Info,
)
}
fn test_up() -> Up {
Arc::new(Mutex::new(io::stdout()))
}
#[test]
fn pause_wait_returns_immediately_when_not_paused() {
let paused = AtomicBool::new(false);
let cancel = AtomicBool::new(false);
let t = Instant::now();
pause_wait(&paused, &cancel, &test_log());
assert!(t.elapsed() < Duration::from_millis(40));
}
#[test]
fn pause_wait_cancel_wins_over_pause() {
let paused = AtomicBool::new(true);
let cancel = AtomicBool::new(true);
let t = Instant::now();
pause_wait(&paused, &cancel, &test_log());
assert!(t.elapsed() < Duration::from_millis(40));
}
#[test]
fn pause_wait_suspends_until_resume() {
let paused = Arc::new(AtomicBool::new(true));
let cancel = Arc::new(AtomicBool::new(false));
let p2 = Arc::clone(&paused);
let unblock = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(120));
p2.store(false, Ordering::Relaxed); });
let t = Instant::now();
pause_wait(&paused, &cancel, &test_log());
assert!(t.elapsed() >= Duration::from_millis(80));
assert!(!paused.load(Ordering::Relaxed));
unblock.join().unwrap();
}
#[test]
fn pause_wait_breaks_out_on_cancel_during_pause() {
let paused = Arc::new(AtomicBool::new(true));
let cancel = Arc::new(AtomicBool::new(false));
let c2 = Arc::clone(&cancel);
let canceller = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(120));
c2.store(true, Ordering::Relaxed); });
pause_wait(&paused, &cancel, &test_log());
assert!(cancel.load(Ordering::Relaxed));
assert!(paused.load(Ordering::Relaxed)); canceller.join().unwrap();
}
fn swap_to(uri: &str, model: Option<&str>, policy: SwapPolicy) -> SwapIntel {
SwapIntel {
uri: uri.into(),
token: None,
model: model.map(str::to_string),
policy,
}
}
#[test]
fn apply_pending_swap_rebuilds_client_and_adopts_model() {
let pending: PendingSwap = Arc::new(Mutex::new(None));
let mut intel = IntelClient::from_parts("https://old.example", None).unwrap();
let mut model = "old-model".to_string();
*pending.lock().unwrap() = Some(swap_to(
"https://a.example,https://b.example",
Some("new-model"),
SwapPolicy::FinishOnOld,
));
apply_pending_swap(&pending, &mut intel, &mut model, &test_up(), &test_log());
assert_eq!(model, "new-model");
assert_eq!(
intel.endpoint_count(),
2,
"client repointed to the new list"
);
assert!(pending.lock().unwrap().is_none());
apply_pending_swap(&pending, &mut intel, &mut model, &test_up(), &test_log());
assert_eq!(model, "new-model");
}
#[test]
fn apply_pending_swap_is_a_noop_when_nothing_pending() {
let pending: PendingSwap = Arc::new(Mutex::new(None));
let mut intel = IntelClient::from_parts("https://only.example", None).unwrap();
let mut model = "m".to_string();
apply_pending_swap(&pending, &mut intel, &mut model, &test_up(), &test_log());
assert_eq!(model, "m");
assert_eq!(intel.endpoint_count(), 1);
}
#[test]
fn restart_turn_pending_only_for_model_change_under_restart_policy() {
let pending: PendingSwap = Arc::new(Mutex::new(None));
assert!(!restart_turn_pending(&pending, "m"));
*pending.lock().unwrap() = Some(swap_to(
"https://a.example",
Some("big"),
SwapPolicy::FinishOnOld,
));
assert!(!restart_turn_pending(&pending, "small"));
*pending.lock().unwrap() = Some(swap_to(
"https://a.example",
Some("small"),
SwapPolicy::RestartTurn,
));
assert!(!restart_turn_pending(&pending, "small"));
*pending.lock().unwrap() = Some(swap_to(
"https://a.example",
Some("big"),
SwapPolicy::RestartTurn,
));
assert!(restart_turn_pending(&pending, "small"));
}
#[test]
fn read_spawn_error_never_echoes_a_swap_intel_token() {
let swap = ControlMsg::SwapIntel(Box::new(SwapIntel {
uri: "https://secret-host.example/secret-path".into(),
token: Some("super-secret-token".into()),
model: Some("m".into()),
policy: SwapPolicy::FinishOnOld,
}));
let mut buf = Vec::new();
frame::write_frame(&mut buf, &swap).unwrap();
let mut reader = BufReader::new(io::Cursor::new(buf));
let err = format!(
"first frame was not Spawn (got {})",
control_msg_label(&swap)
);
assert_eq!(err, "first frame was not Spawn (got swap_intel)");
assert!(!err.contains("super-secret-token"), "token leaked: {err}");
assert!(!err.contains("secret-host.example"), "uri leaked: {err}");
assert_eq!(
control_msg_label(&ControlMsg::Inject {
message: "do bad things".into()
}),
"inject"
);
let _ = &mut reader; }
#[test]
fn bad_swap_list_keeps_the_old_client() {
let pending: PendingSwap = Arc::new(Mutex::new(None));
let mut intel =
IntelClient::from_parts("https://old.example,https://old2.example", None).unwrap();
let mut model = "old".to_string();
*pending.lock().unwrap() = Some(swap_to("", Some("new"), SwapPolicy::FinishOnOld));
apply_pending_swap(&pending, &mut intel, &mut model, &test_up(), &test_log());
assert_eq!(intel.endpoint_count(), 2, "kept the old 2-endpoint client");
assert_eq!(model, "new");
}
}