use super::*;
use crate::config::McPaths;
use serde_json::json;
use std::time::Instant;
fn config(root: &Path) -> EffectiveConfig {
EffectiveConfig {
provider: None,
model: None,
no_color: true,
file_autocomplete_respects_gitignore: true,
custom_providers: BTreeMap::new(),
thinking_level: crate::ThinkingLevel::Low,
auth: None,
paths: McPaths::from_root(root.join("mc")),
}
}
#[test]
fn config_refresh_cancels_old_work_and_passes_current_selection() {
let temp = tempfile::tempdir().unwrap();
let mut selected = config(temp.path());
let settings = Settings::default();
let (sender, receiver) = crossbeam_channel::bounded(4);
let mut runtime = SummarizerRuntime {
sender: Some(sender),
mailbox: Arc::new(Mutex::new(None)),
session: None,
generation: 0,
cancellation: None,
config: Arc::new((selected.clone(), settings.clone())),
};
let session = SessionManager::new(selected.paths.sessions.clone())
.create()
.unwrap();
runtime.set_session(Some(session));
let old = receiver.recv().unwrap();
selected.provider = Some("replacement".into());
selected.model = Some("new-model".into());
selected.thinking_level = crate::ThinkingLevel::High;
runtime.refresh_config(&selected, &settings);
let current = receiver.recv().unwrap();
assert!(old.cancellation.is_canceled());
assert!(!current.cancellation.is_canceled());
assert!(current.config.0 == selected);
publish(
&runtime.mailbox,
old.generation,
SummarySnapshot {
session_id: old.session.unwrap().id().into(),
enabled: true,
running: false,
entries: vec!["old provider".into()],
error: None,
},
);
assert!(runtime.poll().is_none());
runtime.refresh_config(&selected, &settings);
assert!(receiver.try_recv().is_err());
assert!(!current.cancellation.is_canceled());
}
fn append(session: &Session, kind: &str, payload: serde_json::Value) {
session
.append(&SessionEvent::new(
kind,
session.id().into(),
Path::new(".").into(),
payload,
))
.unwrap();
}
#[test]
fn full_command_queue_cancels_old_session_and_rejects_stale_snapshots() {
let temp = tempfile::tempdir().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let first = manager.create().unwrap();
let second = manager.create().unwrap();
let (sender, receiver) = crossbeam_channel::bounded(1);
let mut runtime = SummarizerRuntime {
sender: Some(sender),
mailbox: Arc::new(Mutex::new(None)),
session: None,
generation: 0,
cancellation: None,
config: Arc::new((config(temp.path()), Settings::default())),
};
runtime.set_session(Some(first.clone()));
runtime.set_session(Some(second.clone()));
let old = receiver.recv().unwrap();
assert!(old.cancellation.is_canceled());
publish(
&runtime.mailbox,
old.generation,
SummarySnapshot {
session_id: first.id().into(),
enabled: true,
running: true,
entries: vec!["Stale".into()],
error: None,
},
);
let snapshot = runtime.poll().unwrap();
assert_eq!(snapshot.session_id, second.id());
assert!(snapshot.error.is_some());
assert!(snapshot.entries.is_empty());
runtime.set_enabled(true);
let retry = receiver.recv().unwrap();
assert_eq!(retry.session.unwrap().id(), second.id());
assert!(!retry.cancellation.is_canceled());
}
#[test]
fn failed_checkpoint_write_does_not_advance_visible_entries_or_cursors() {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("summaries");
crate::sessions::prepare_session_root(&root).unwrap();
let log = SessionManager::new(root.clone()).open("primary").unwrap();
std::fs::create_dir(log.path()).unwrap();
let lease =
crate::persistence::CrossProcessFileLock::try_acquire(&root.join("primary.observer"))
.unwrap()
.unwrap();
let mut active = Active {
command: Command {
generation: 1,
session: None,
enabled: None,
cancellation: AgentCancellation::default(),
config: Arc::new((config(temp.path()), Settings::default())),
},
log,
checkpoint: Checkpoint::default(),
error: None,
_lease: lease,
};
let checkpoint = Checkpoint {
entries: vec!["Not committed".into()],
cursors: BTreeMap::from([("primary".into(), 42)]),
..Checkpoint::default()
};
assert!(active.commit(checkpoint).is_err());
assert!(active.snapshot(false).entries.is_empty());
assert!(active.checkpoint.cursors.is_empty());
}
#[test]
fn parser_accepts_silence_and_multiple_entries_but_rejects_invalid_output() {
assert!(parse_entries("[]").unwrap().is_empty());
assert_eq!(
parse_entries(r#"[" Updated the parser.\nChecks passed. ","Found a missing guard."]"#)
.unwrap(),
vec![
"Updated the parser. Checks passed.",
"Found a missing guard."
]
);
assert_eq!(
parse_entries(r#"["a","b","c","d","e","f"]"#).unwrap().len(),
6
);
for raw in [
"not json",
"{}",
"[1]",
"[\"\"]",
"[\"a\",\"b\",\"c\",\"d\",\"e\",\"f\",\"g\"]",
] {
assert!(parse_entries(raw).is_err(), "accepted {raw}");
}
assert!(parse_entries(&serde_json::to_string(&vec!["é".repeat(601)]).unwrap()).is_err());
}
#[test]
fn repeated_outputs_are_skipped_against_recent_history_and_within_batch() {
let mut previous = vec!["Updated the parser.".into(), "Checks passed.".into()];
let entries = parse_entries(
r#"["Updated the parser.","Subagent child: Found a missing guard.","Checks passed.","Subagent child: Found a missing guard."]"#,
)
.unwrap();
append_unique_entries(&mut previous, entries.clone());
append_unique_entries(&mut previous, entries);
assert_eq!(
previous,
vec![
"Updated the parser.",
"Checks passed.",
"Subagent child: Found a missing guard."
]
);
}
#[test]
fn dedup_window_keeps_all_twelve_prior_entries_but_allows_older_outcomes() {
let mut previous = (0..13)
.map(|index| format!("Outcome {index}"))
.collect::<Vec<_>>();
append_unique_entries(
&mut previous,
vec![
"New outcome".into(),
"Outcome 1".into(),
"Outcome 0".into(),
"Outcome 12".into(),
],
);
assert_eq!(previous.len(), 15);
assert_eq!(&previous[13..], &["New outcome", "Outcome 0"]);
}
#[test]
fn batches_follow_child_links_skip_local_events_and_resume_without_repetition() {
let temp = tempfile::tempdir().unwrap();
let config = config(temp.path());
let primary = SessionManager::new(config.paths.sessions.clone())
.create()
.unwrap();
let child = SessionManager::new(config.paths.sessions.join("subagents"))
.create()
.unwrap();
append(&primary, "user_input", json!({"text":"Fix the parser"}));
append(
&primary,
"diagnostic",
json!({"text":"private diagnostic must not be sent"}),
);
append(
&child,
"assistant_output",
json!({"text":"Added the missing guard."}),
);
append(
&primary,
"tool_result",
json!({"result":{"tool_name":"subagents", "content":json!({"results":[{"session_id":child.id()}, {"session_id":"../escape"}]}).to_string()}}),
);
let (text, cursors, children) =
collect_activity(&primary, &config, &Checkpoint::default()).unwrap();
assert!(text.contains(&format!(
"Primary {} user_input: Fix the parser\n",
primary.id()
)));
assert!(text.contains(&format!(
"Subagent {} assistant_output: Added the missing guard.\n",
child.id()
)));
assert!(!text.contains("private diagnostic"));
assert_eq!(children, vec![child.id()]);
let checkpoint = Checkpoint {
enabled: true,
cursors,
children,
entries: vec![],
};
let log = SessionManager::new(config.paths.sessions.join("summaries"))
.open(primary.id())
.unwrap();
append(
&log,
"summary_checkpoint",
serde_json::to_value(&checkpoint).unwrap(),
);
let checkpoint = load(&log, false).unwrap();
let (text, _, _) = collect_activity(&primary, &config, &checkpoint).unwrap();
assert!(text.is_empty());
append(&child, "assistant_output", json!({"text":"Tests passed."}));
let (text, _, _) = collect_activity(&primary, &config, &checkpoint).unwrap();
assert_eq!(
text,
format!("Subagent {} assistant_output: Tests passed.\n", child.id())
);
assert!(!text.contains("Added the missing guard."));
}
#[test]
fn launched_child_is_discovered_during_provider_execution() {
use crate::providers::{Provider, ProviderEvent, ProviderRequest};
struct LiveProbe {
primary: Session,
config: EffectiveConfig,
observed: std::sync::atomic::AtomicBool,
}
impl Provider for LiveProbe {
fn stream_cancellable(
&self,
_request: ProviderRequest,
_cancellation: &AgentCancellation,
_on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let history = self.primary.read_events()?;
assert!(
!history
.iter()
.any(|event| event.event_type == "tool_result")
);
let (text, cursors, children) =
collect_activity(&self.primary, &self.config, &Checkpoint::default())?;
assert_eq!(children.len(), 1);
assert!(cursors.contains_key(&children[0]));
assert!(text.contains("Inspect live linkage"));
assert!(!text.contains("unrelated secret"));
self.observed
.store(true, std::sync::atomic::Ordering::SeqCst);
anyhow::bail!("probe stops child after observing live discovery")
}
}
for depth in 0..=1 {
let temp = tempfile::tempdir().unwrap();
let config = config(temp.path());
let parent_root = if depth == 0 {
config.paths.sessions.clone()
} else {
config.paths.sessions.join("subagents")
};
let primary = SessionManager::new(parent_root).create().unwrap();
append(
&primary,
"user_input",
json!({"text":"Delegate inspection"}),
);
let unrelated = SessionManager::new(config.paths.sessions.join("subagents"))
.create()
.unwrap();
append(&unrelated, "user_input", json!({"text":"unrelated secret"}));
let provider = Arc::new(LiveProbe {
primary: primary.clone(),
config: config.clone(),
observed: std::sync::atomic::AtomicBool::new(false),
});
let run = crate::subagents::SubagentRunConfig {
parent_agent: crate::agent::AgentSession::new(
"model",
&[],
&crate::skills::SkillDiscovery::default(),
),
provider: provider.clone(),
provider_override: None,
parent_tools: crate::tools::ToolRuntime::new(temp.path()).unwrap(),
parent_cwd: temp.path().into(),
cancellation: AgentCancellation::default(),
profiles: BTreeMap::new(),
subagent_profiles_prompt: None,
sessions_root: Some(config.paths.sessions.clone()),
parent_session_id: Some(primary.id().into()),
depth,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: None,
semantic_progress_timeout: Duration::from_secs(5),
schema_validation_max_retries: 0,
compaction: None,
};
crate::subagents::dispatch_subagents(
json!({"tasks":[{"intent":"Inspect live linkage"}]}),
run.clone(),
);
assert!(provider.observed.load(std::sync::atomic::Ordering::SeqCst));
provider
.observed
.store(false, std::sync::atomic::Ordering::SeqCst);
let mut invalid_run = run;
invalid_run.parent_session_id = Some("missing-parent".into());
let result = crate::subagents::dispatch_subagents(
json!({"tasks":[{"intent":"Inspect live linkage"}]}),
invalid_run,
);
assert!(result.content.contains("subagent session linkage failed"));
assert!(!provider.observed.load(std::sync::atomic::Ordering::SeqCst));
}
}
#[test]
fn busy_primary_does_not_starve_linked_children() {
let temp = tempfile::tempdir().unwrap();
let config = config(temp.path());
let primary = SessionManager::new(config.paths.sessions.clone())
.create()
.unwrap();
let child = SessionManager::new(config.paths.sessions.join("subagents"))
.create()
.unwrap();
for _ in 0..30 {
append(
&primary,
"assistant_chunk",
json!({"text": "x".repeat(2000)}),
);
}
append(
&child,
"assistant_chunk",
json!({"text": "Child is still working"}),
);
let checkpoint = Checkpoint {
children: vec![child.id().into()],
..Checkpoint::default()
};
let (text, cursors, _) = collect_activity(&primary, &config, &checkpoint).unwrap();
assert!(text.contains("Child is still working"));
assert!(cursors.contains_key(primary.id()));
assert!(cursors.contains_key(child.id()));
assert!(text.chars().count() <= BATCH_CHARS);
}
#[test]
fn summary_config_trims_identifiers_and_defaults_for_changed_provider() {
let temp = tempfile::tempdir().unwrap();
let mut config = config(temp.path());
config.provider = Some("openai-codex".into());
config.model = Some("active-model".into());
let mut settings = Settings::default();
assert_eq!(
summary_selection(&config, &settings).unwrap(),
("openai-codex".into(), "active-model".into())
);
settings.summarizer.provider = Some(" anthropic ".into());
assert_eq!(
summary_selection(&config, &settings).unwrap(),
(
"anthropic".into(),
crate::providers::default_model_for_provider("anthropic").into()
)
);
settings.summarizer.model = Some(" custom-model ".into());
assert_eq!(
summary_selection(&config, &settings).unwrap(),
("anthropic".into(), "custom-model".into())
);
settings.summarizer.provider = Some("custom".into());
settings.summarizer.model = None;
assert!(summary_selection(&config, &settings).is_err());
}
#[test]
fn custom_summary_prompt_replaces_builtin_instructions() {
let mut settings = Settings::default();
assert_eq!(summary_prompt(&settings), PROMPT);
settings.summarizer.prompt = Some("Return brief JSON observations.".into());
assert_eq!(summary_prompt(&settings), "Return brief JSON observations.");
}
#[test]
fn batch_budget_leaves_remaining_events_for_next_batch() {
let temp = tempfile::tempdir().unwrap();
let config = config(temp.path());
let session = SessionManager::new(config.paths.sessions.clone())
.create()
.unwrap();
for index in 0..30 {
append(
&session,
"assistant_output",
json!({"text":format!("{index}:{}", "x".repeat(1900))}),
);
}
let (first, cursors, children) =
collect_activity(&session, &config, &Checkpoint::default()).unwrap();
assert!(first.chars().count() <= BATCH_CHARS);
let (second, _, _) = collect_activity(
&session,
&config,
&Checkpoint {
cursors,
children,
..Checkpoint::default()
},
)
.unwrap();
assert!(!second.is_empty());
assert_ne!(first, second);
}
#[test]
fn checkpoint_restores_enable_entries_and_cursors_with_private_permissions() {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("summaries");
crate::sessions::prepare_session_root(&root).unwrap();
let log = SessionManager::new(root).open("primary").unwrap();
let checkpoint = Checkpoint {
enabled: true,
entries: vec!["Updated the parser.".into()],
cursors: BTreeMap::from([("primary".into(), 42)]),
children: vec![],
};
append(
&log,
"summary_checkpoint",
serde_json::to_value(&checkpoint).unwrap(),
);
let restored = load(&log, false).unwrap();
assert!(restored.enabled);
assert_eq!(restored.entries, checkpoint.entries);
assert_eq!(restored.cursors, checkpoint.cursors);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
assert_eq!(
std::fs::metadata(log.path()).unwrap().permissions().mode() & 0o777,
0o600
);
}
}
fn wait_snapshot(runtime: &mut SummarizerRuntime) -> SummarySnapshot {
let deadline = Instant::now() + Duration::from_secs(5);
loop {
if let Some(snapshot) = runtime.poll() {
return snapshot;
}
assert!(
Instant::now() < deadline,
"worker did not publish a snapshot"
);
thread::sleep(Duration::from_millis(10));
}
}
#[test]
fn controller_persists_toggles_and_isolates_session_switches() {
let temp = tempfile::tempdir().unwrap();
let config = config(temp.path());
let manager = SessionManager::new(config.paths.sessions.clone());
let first = manager.create().unwrap();
let second = manager.create().unwrap();
let mut runtime = SummarizerRuntime::new(config.clone(), Settings::default());
runtime.set_session(Some(first.clone()));
assert!(!wait_snapshot(&mut runtime).enabled);
runtime.set_enabled(true);
assert!(wait_snapshot(&mut runtime).enabled);
runtime.set_session(Some(second.clone()));
let snapshot = wait_snapshot(&mut runtime);
assert_eq!(snapshot.session_id, second.id());
assert!(!snapshot.enabled);
runtime.set_session(Some(first.clone()));
let snapshot = wait_snapshot(&mut runtime);
assert_eq!(snapshot.session_id, first.id());
assert!(snapshot.enabled);
runtime.set_enabled(false);
assert!(!wait_snapshot(&mut runtime).enabled);
let log = SessionManager::new(config.paths.sessions.join("summaries"))
.open(first.id())
.unwrap();
assert!(!load(&log, true).unwrap().enabled);
runtime.set_session(None);
assert!(runtime.poll().is_none());
runtime.shutdown();
}
#[cfg(unix)]
#[test]
fn summary_storage_rejects_symlinked_log() {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("summaries");
crate::sessions::prepare_session_root(&root).unwrap();
let outside = temp.path().join("outside");
std::fs::write(&outside, "unchanged").unwrap();
std::os::unix::fs::symlink(&outside, root.join("primary.jsonl")).unwrap();
let log = SessionManager::new(root).open("primary").unwrap();
assert!(load(&log, false).is_err());
assert_eq!(std::fs::read_to_string(outside).unwrap(), "unchanged");
}
struct TestProvider {
wait_for_cancel: bool,
output: String,
}
impl crate::providers::Provider for TestProvider {
fn stream_cancellable(
&self,
_request: ProviderRequest,
cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
if self.wait_for_cancel {
while !cancellation.is_canceled() {
thread::sleep(Duration::from_millis(1));
}
cancellation.check()?;
}
on_event(ProviderEvent::TextDelta(self.output.clone()))?;
on_event(ProviderEvent::Done)
}
}
#[test]
fn provider_stream_is_bounded_and_timeout_cancels_without_entries() {
let request = || ProviderRequest::new_without_tools("test", vec![]);
let provider = TestProvider {
wait_for_cancel: false,
output: "[\"Added a guard.\"]".into(),
};
assert_eq!(
stream_entries(
&provider,
request(),
&AgentCancellation::default(),
Duration::from_secs(1)
)
.unwrap(),
vec!["Added a guard."]
);
let provider = TestProvider {
wait_for_cancel: false,
output: "x".repeat(16_385),
};
assert!(
stream_entries(
&provider,
request(),
&AgentCancellation::default(),
Duration::from_secs(1)
)
.unwrap_err()
.to_string()
.contains("exceeds limit")
);
let provider = TestProvider {
wait_for_cancel: true,
output: "[]".into(),
};
assert!(
stream_entries(
&provider,
request(),
&AgentCancellation::default(),
Duration::from_millis(10)
)
.unwrap_err()
.to_string()
.contains("timed out")
);
let (token, handle) = AgentCancellation::default().child_token();
handle.cancel();
assert!(crate::cancellation::is_run_canceled(
&stream_entries(&provider, request(), &token, Duration::from_secs(1)).unwrap_err()
));
}