use std::collections::BTreeMap;
use std::fmt;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use anyhow::{Context, Result, bail};
use futures::future::{FutureExt, Shared, TryFutureExt, join_all};
use futures::stream::{FuturesUnordered, StreamExt};
use mj_client::session::BoxFuture;
use mj_core::config::{Config, HarnessKind, HarnessProfile, SubagentConfig};
use mj_core::worker_launch::ProfileConfig;
use tokio_util::sync::CancellationToken;
#[derive(Clone, PartialEq)]
struct ProfilesKey {
profiles: BTreeMap<String, HarnessProfile>,
subagents: SubagentConfig,
}
impl ProfilesKey {
fn of(config: &Config) -> Self {
Self {
profiles: config.profiles.clone(),
subagents: config.subagents.clone(),
}
}
fn matches(&self, config: &Config) -> bool {
self.profiles.len() == config.profiles.len()
&& self.profiles.iter().all(|(id, profile)| {
config.profiles.get(id).is_some_and(|updated| {
profile.enabled == updated.enabled
&& profile.discovery_inputs() == updated.discovery_inputs()
})
})
&& self.subagents.eligible_profiles == config.subagents.eligible_profiles
}
fn candidates(&self, parent: &str) -> Vec<(String, HarnessKind)> {
self.profiles
.iter()
.filter(|(id, profile)| {
profile.enabled && self.subagents.profile_is_eligible(parent, id)
})
.map(|(id, profile)| (id.clone(), profile.kind))
.collect()
}
fn warm_set(&self) -> Vec<String> {
self.profiles
.iter()
.filter(|(_, profile)| profile.enabled)
.map(|(id, _)| id.clone())
.collect()
}
}
pub(crate) type Probe = dyn Fn(String) -> BoxFuture<'static, Result<ProfileConfig>> + Send + Sync;
type ModelProbe = dyn Fn(String, String) -> BoxFuture<'static, Result<ProfileConfig>> + Send + Sync;
#[derive(Clone, Debug)]
struct DiscoveryFailure(Arc<str>);
impl fmt::Display for DiscoveryFailure {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.0)
}
}
impl std::error::Error for DiscoveryFailure {}
type Attempt = Shared<BoxFuture<'static, Result<ProfileConfig, DiscoveryFailure>>>;
enum Entry {
Ready(ProfileConfig),
Pending(Attempt),
}
impl Entry {
fn belongs_to(&self, attempt: &Attempt) -> bool {
matches!(self, Self::Pending(current) if current.ptr_eq(attempt))
}
}
#[derive(Default)]
struct Inner {
generation: u64,
key: Option<ProfilesKey>,
entries: BTreeMap<String, Entry>,
models: BTreeMap<(String, String), Entry>,
fingerprints: BTreeMap<String, String>,
}
impl Inner {
fn adopted(&self) -> Result<&ProfilesKey> {
self.key
.as_ref()
.context("the profile catalogue has not adopted a configuration")
}
fn generation_for(&self, profiles: &[String]) -> Result<u64> {
let key = self.adopted()?;
for profile in profiles {
if !key.profiles.contains_key(profile) {
bail!("the adopted configuration holds no '{profile}' profile");
}
}
Ok(self.generation)
}
}
pub(crate) struct ProfileCatalog {
cancellation: CancellationToken,
probe: Arc<Probe>,
model_probe: Arc<ModelProbe>,
inner: Mutex<Inner>,
check_home: bool,
}
impl ProfileCatalog {
pub(crate) fn new(cancellation: CancellationToken) -> Arc<Self> {
let mut catalog = Self::build(
cancellation,
Arc::new(|profile| {
Box::pin(crate::controller::profile_config::discover(
profile, None, false,
))
}),
Arc::new(|profile, model| {
Box::pin(crate::controller::profile_config::discover(
profile,
Some(model),
false,
))
}),
);
Arc::get_mut(&mut catalog)
.expect("new catalogue has one owner")
.check_home = true;
catalog
}
fn build(
cancellation: CancellationToken,
probe: Arc<Probe>,
model_probe: Arc<ModelProbe>,
) -> Arc<Self> {
Arc::new(Self {
cancellation,
probe,
model_probe,
inner: Mutex::new(Inner::default()),
check_home: false,
})
}
pub(crate) async fn model_capabilities(
&self,
profile: String,
model: String,
) -> Result<ProfileConfig> {
let generation = self.lock().generation_for(std::slice::from_ref(&profile))?;
let generation = self.check_inputs(generation, &profile).await?;
let key = (profile.clone(), model.clone());
let entry = {
let mut inner = self.lock();
anyhow::ensure!(
inner.generation == generation,
"profile configuration changed during discovery"
);
if let Some(Entry::Ready(choices)) = inner.entries.get(&profile)
&& choices.model.as_ref() == Some(&model)
{
return Ok(choices.clone());
}
match inner.models.get(&key) {
Some(Entry::Ready(choices)) => return Ok(choices.clone()),
Some(Entry::Pending(attempt)) => attempt.clone(),
None => {
let attempt = (self.model_probe)(profile, model)
.map_err(|error| DiscoveryFailure(format!("{error:#}").into()))
.boxed()
.shared();
inner
.models
.insert(key.clone(), Entry::Pending(attempt.clone()));
attempt
}
}
};
let result = tokio::select! {
_ = self.cancellation.cancelled() => bail!("profile discovery cancelled by daemon shutdown"),
result = entry.clone() => result,
};
let mut inner = self.lock();
anyhow::ensure!(
inner.generation == generation,
"profile configuration changed during discovery"
);
let owns_entry = inner
.models
.get(&key)
.is_some_and(|current| current.belongs_to(&entry));
match result {
Ok(choices) => {
if owns_entry {
inner.models.insert(key, Entry::Ready(choices.clone()));
}
Ok(choices)
}
Err(error) => {
if owns_entry {
inner.models.remove(&key);
}
Err(error.into())
}
}
}
pub(crate) fn sync(self: &Arc<Self>, config: &Config) {
let Some((generation, key)) = self.claim_pass(config) else {
return;
};
let catalog = self.clone();
tokio::spawn(async move { catalog.warm_pass(generation, key).await });
}
fn claim_pass(&self, config: &Config) -> Option<(u64, ProfilesKey)> {
let mut inner = self.lock();
if inner.key.as_ref().is_some_and(|key| key.matches(config)) {
return None;
}
let key = ProfilesKey::of(config);
inner.generation = inner.generation.wrapping_add(1);
let unchanged = |id: &String| {
inner
.key
.as_ref()
.and_then(|old| old.profiles.get(id))
.zip(key.profiles.get(id))
.is_some_and(|(old, new)| {
new.enabled && old.discovery_inputs() == new.discovery_inputs()
})
};
let retained: std::collections::BTreeSet<_> = inner
.entries
.keys()
.chain(inner.models.keys().map(|(id, _)| id))
.filter(|id| unchanged(id))
.cloned()
.collect();
inner.entries.retain(|id, _| retained.contains(id));
inner.models.retain(|(id, _), _| retained.contains(id));
inner.fingerprints.retain(|id, _| retained.contains(id));
inner.key = Some(key.clone());
Some((inner.generation, key))
}
pub(crate) fn candidates(&self, parent: &str) -> Result<Vec<(String, HarnessKind)>> {
let inner = self.lock();
Ok(inner.adopted()?.candidates(parent))
}
pub(crate) async fn capabilities(&self, profiles: &[String]) -> Result<Vec<ProfileConfig>> {
let generation = {
let inner = self.lock();
inner.generation_for(profiles)?
};
let mut generation = generation;
for profile in profiles {
generation = self.check_inputs(generation, profile).await?;
}
let mut discoveries: Vec<BoxFuture<'_, Result<ProfileConfig>>> =
Vec::with_capacity(profiles.len());
for profile in profiles {
match self.entry(generation, profile) {
Some(Entry::Ready(config)) => discoveries.push(async move { Ok(config) }.boxed()),
Some(Entry::Pending(attempt)) => {
let catalog = self;
let profile = profile.clone();
discoveries.push(
async move {
match attempt.clone().await {
Ok(config) => {
catalog.remember(generation, &profile, &attempt, &config);
Ok(config)
}
Err(failure) => {
catalog.forget(generation, &profile, &attempt);
Err(anyhow::anyhow!(
"could not discover the capabilities of the \
'{profile}' profile: {failure}"
))
}
}
}
.boxed(),
);
}
None => bail!(
"the profile catalogue adopted a new configuration before it could \
answer for the '{profile}' profile"
),
}
}
let results = tokio::select! {
_ = self.cancellation.cancelled() => bail!(
"the server stopped before the profile catalogue discovered what \
the sub-agent tools need"
),
results = join_all(discoveries) => results,
};
anyhow::ensure!(
self.is_current(generation),
"profile configuration changed during discovery"
);
results.into_iter().collect()
}
pub(crate) async fn options_for(
&self,
config: &Config,
parent: &str,
model: Option<String>,
) -> Result<mj_core::subagent::SubagentOptions> {
crate::controller::profile_config::subagent_options_with(
config,
parent,
model,
|id, model| {
let live = self
.lock()
.key
.as_ref()
.and_then(|key| key.profiles.get(&id))
.zip(config.profiles.get(&id))
.is_some_and(|(live, draft)| {
live.enabled && live.discovery_inputs() == draft.discovery_inputs()
});
async move {
if live {
match model {
Some(model) => self.model_capabilities(id, model).await,
None => Ok(self.capabilities(&[id]).await?.remove(0)),
}
} else {
crate::controller::profile_config::discover_for(config.clone(), id, model)
.await
}
}
},
)
.await
}
async fn check_inputs(&self, generation: u64, id: &str) -> Result<u64> {
if !self.check_home {
return Ok(generation);
}
let profile = {
let inner = self.lock();
anyhow::ensure!(
inner.generation == generation,
"profile configuration changed during discovery"
);
inner
.adopted()?
.profiles
.get(id)
.context("profile unavailable")?
.clone()
};
let profile_id = id.to_owned();
let fingerprint = tokio::task::spawn_blocking(move || {
crate::controller::profile_config::discovery_fingerprint(&profile_id, &profile)
})
.await
.context("profile fingerprint task panicked")??;
let mut inner = self.lock();
anyhow::ensure!(
inner.generation == generation,
"profile configuration changed during discovery"
);
if inner
.fingerprints
.get(id)
.is_some_and(|old| old != &fingerprint)
{
inner.entries.remove(id);
inner.models.retain(|(profile, _), _| profile != id);
inner.generation = inner.generation.wrapping_add(1);
}
inner.fingerprints.insert(id.to_owned(), fingerprint);
Ok(inner.generation)
}
fn entry(&self, generation: u64, profile: &str) -> Option<Entry> {
let mut inner = self.lock();
if inner.generation != generation {
return None;
}
Some(match inner.entries.get(profile) {
Some(Entry::Ready(config)) => Entry::Ready(config.clone()),
Some(Entry::Pending(attempt)) => Entry::Pending(attempt.clone()),
None => {
let attempt = (self.probe)(profile.to_owned())
.map_err(|error| DiscoveryFailure(format!("{error:#}").into()))
.boxed()
.shared();
inner
.entries
.insert(profile.to_owned(), Entry::Pending(attempt.clone()));
Entry::Pending(attempt)
}
})
}
fn remember(&self, generation: u64, profile: &str, attempt: &Attempt, config: &ProfileConfig) {
let mut inner = self.lock();
if inner.generation != generation
|| !inner
.entries
.get(profile)
.is_some_and(|entry| entry.belongs_to(attempt))
{
return;
}
inner
.entries
.insert(profile.to_owned(), Entry::Ready(config.clone()));
}
fn forget(&self, generation: u64, profile: &str, attempt: &Attempt) {
let mut inner = self.lock();
if inner.generation != generation
|| !inner
.entries
.get(profile)
.is_some_and(|entry| entry.belongs_to(attempt))
{
return;
}
inner.entries.remove(profile);
}
async fn warm_pass(self: Arc<Self>, mut generation: u64, key: ProfilesKey) {
let mut discoveries = FuturesUnordered::new();
for profile in key.warm_set() {
if self.cancellation.is_cancelled() {
break;
}
match self.check_inputs(generation, &profile).await {
Ok(current) => generation = current,
Err(error) => {
tracing::warn!(profile, error = %error, "profile catalogue input discovery failed");
continue;
}
}
match self.entry(generation, &profile) {
None => break,
Some(Entry::Ready(_)) => {}
Some(Entry::Pending(attempt)) => {
discoveries.push(async move {
let result = attempt.clone().await;
(profile, attempt, result)
});
}
}
}
loop {
let next = tokio::select! {
_ = self.cancellation.cancelled() => break,
next = discoveries.next() => next,
};
let Some((profile, attempt, result)) = next else {
break;
};
if !self.is_current(generation) {
break;
}
match result {
Ok(config) => self.remember(generation, &profile, &attempt, &config),
Err(failure) => {
self.forget(generation, &profile, &attempt);
tracing::warn!(
profile,
error = %failure,
"profile catalog discovery failed; list_profiles will retry it on demand"
);
}
}
}
}
fn is_current(&self, generation: u64) -> bool {
self.lock().generation == generation
}
fn lock(&self) -> MutexGuard<'_, Inner> {
self.inner.lock().unwrap_or_else(PoisonError::into_inner)
}
#[cfg(test)]
pub(crate) fn with_probe(probe: Arc<Probe>) -> Arc<Self> {
let model_probe = probe.clone();
Self::build(
CancellationToken::new(),
probe,
Arc::new(move |profile, _model| model_probe(profile)),
)
}
#[cfg(test)]
pub(crate) async fn sync_now(self: &Arc<Self>, config: &Config) {
if let Some((generation, key)) = self.claim_pass(config) {
self.clone().warm_pass(generation, key).await;
}
}
}
#[cfg(test)]
pub(crate) fn test_config(profiles: &[(&str, HarnessKind)], eligible: &[&str]) -> Config {
Config {
profiles: profiles
.iter()
.map(|(id, kind)| {
(
(*id).to_owned(),
HarnessProfile {
enabled: true,
kind: *kind,
home: std::path::PathBuf::from("/home/agent").join(id),
environment: Default::default(),
context_window_bytes: None,
subagents: Default::default(),
guardian_review_model: None,
},
)
})
.collect(),
subagents: SubagentConfig {
eligible_profiles: eligible.iter().map(|id| ((*id).to_owned(), true)).collect(),
..SubagentConfig::default()
},
..Config::default()
}
}
#[cfg(test)]
pub(crate) fn counting_probe(calls: Arc<std::sync::atomic::AtomicUsize>) -> Arc<Probe> {
Arc::new(move |profile| {
calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Box::pin(async move { Ok(test_choices(&profile)) })
})
}
#[cfg(test)]
fn test_choices(profile: &str) -> ProfileConfig {
ProfileConfig {
model: Some(format!("{profile}-model")),
models: vec![mj_core::acp::SessionConfigChoice {
value: format!("{profile}-model"),
name: format!("{profile} model"),
description: None,
}],
efforts: Vec::new(),
observed_at: 1_700_000_000,
}
}
#[cfg(test)]
fn flag_probe(
fails: Arc<std::sync::atomic::AtomicBool>,
calls: Arc<std::sync::atomic::AtomicUsize>,
) -> Arc<Probe> {
Arc::new(move |profile| {
let fails = fails.clone();
let calls = calls.clone();
Box::pin(async move {
calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if fails.load(std::sync::atomic::Ordering::SeqCst) {
anyhow::bail!("harness is not installed")
}
Ok(test_choices(&profile))
})
})
}
#[cfg(test)]
fn gated_probe(
calls: Arc<std::sync::atomic::AtomicUsize>,
started: tokio::sync::mpsc::UnboundedSender<String>,
gate: tokio::sync::watch::Receiver<bool>,
) -> Arc<Probe> {
Arc::new(move |profile| {
let calls = calls.clone();
let started = started.clone();
let mut gate = gate.clone();
Box::pin(async move {
calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let _ = started.send(profile.clone());
if !*gate.borrow_and_update() {
let _ = gate.changed().await;
}
Ok(test_choices(&profile))
})
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
fn calls() -> Arc<AtomicUsize> {
Arc::new(AtomicUsize::new(0))
}
fn fails() -> Arc<std::sync::atomic::AtomicBool> {
Arc::new(std::sync::atomic::AtomicBool::new(true))
}
#[tokio::test]
async fn a_delayed_failed_waiter_cannot_remove_or_replace_the_retry() {
let calls = calls();
let probes_fail = fails();
let catalog = ProfileCatalog::with_probe(flag_probe(probes_fail.clone(), calls.clone()));
let config = test_config(&[("parent", HarnessKind::Codex)], &[]);
let (generation, _) = catalog.claim_pass(&config).unwrap();
let Some(Entry::Pending(failed)) = catalog.entry(generation, "parent") else {
panic!("first attempt");
};
assert!(failed.clone().await.is_err());
catalog.forget(generation, "parent", &failed);
probes_fail.store(false, Ordering::SeqCst);
let Some(Entry::Pending(retry)) = catalog.entry(generation, "parent") else {
panic!("retry");
};
catalog.forget(generation, "parent", &failed);
catalog.remember(generation, "parent", &failed, &test_choices("stale"));
let choices = catalog.capabilities(&["parent".into()]).await.unwrap();
assert_eq!(choices[0].model.as_deref(), Some("parent-model"));
assert_eq!(calls.load(Ordering::SeqCst), 2, "the retry is shared");
catalog.forget(generation, "parent", &failed);
catalog.forget(generation, "parent", &retry);
catalog.capabilities(&["parent".into()]).await.unwrap();
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"late waiters cannot discard ready entries"
);
}
#[tokio::test]
async fn provider_file_changes_retire_pending_model_replies_and_populate_only_affected_entries()
{
let home = tempfile::tempdir().unwrap();
let calls = calls();
let model_calls = calls.clone();
let (started_tx, mut started) = tokio::sync::mpsc::unbounded_channel();
let (release, gate) = tokio::sync::watch::channel(false);
let mut catalog = ProfileCatalog::build(
CancellationToken::new(),
counting_probe(Arc::new(AtomicUsize::new(0))),
Arc::new(move |profile, _model| {
let number = model_calls.fetch_add(1, Ordering::SeqCst);
let mut gate = gate.clone();
let started = started_tx.clone();
Box::pin(async move {
if number == 0 {
started.send(()).unwrap();
gate.changed().await.unwrap();
}
let mut choices = test_choices(&profile);
choices.observed_at = number as i64;
Ok(choices)
})
}),
);
Arc::get_mut(&mut catalog).unwrap().check_home = true;
let mut config = test_config(&[("parent", HarnessKind::Codex)], &[]);
config.profiles.get_mut("parent").unwrap().home = home.path().into();
catalog.sync_now(&config).await;
let pending = {
let catalog = catalog.clone();
tokio::spawn(async move {
catalog
.model_capabilities("parent".into(), "other".into())
.await
})
};
started.recv().await.unwrap();
std::fs::write(
home.path().join("config.toml"),
"model_provider = 'changed'\n",
)
.unwrap();
let fresh = catalog
.model_capabilities("parent".into(), "other".into())
.await
.unwrap();
assert_eq!(fresh.observed_at, 1);
release.send(true).unwrap();
assert!(
pending.await.unwrap().is_err(),
"retired provider reply is rejected"
);
assert_eq!(
catalog
.model_capabilities("parent".into(), "other".into())
.await
.unwrap()
.observed_at,
1
);
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"new provider is populated once and stale reply cannot replace it"
);
}
#[tokio::test]
async fn draft_eligibility_uses_warmed_entries_without_adopting_draft_configuration() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let live = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Codex),
],
&[],
);
catalog.sync_now(&live).await;
let mut draft = live.clone();
draft
.subagents
.eligible_profiles
.insert("helper".into(), true);
let options = catalog.options_for(&draft, "parent", None).await.unwrap();
assert_eq!(options.models.len(), 2);
assert_eq!(
catalog.candidates("parent").unwrap(),
vec![("parent".into(), HarnessKind::Codex)]
);
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"draft discovery shares warm live entries"
);
}
#[tokio::test]
async fn default_edits_and_cached_models_reuse_probes_while_installation_changes_invalidate() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let mut config = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
],
&["helper"],
);
catalog.sync_now(&config).await;
let (left, right) = tokio::join!(
catalog.model_capabilities("parent".into(), "another".into()),
catalog.model_capabilities("parent".into(), "another".into()),
);
left.unwrap();
right.unwrap();
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"concurrent model misses share one probe"
);
for effort in ["low", "high"] {
config.profiles.get_mut("parent").unwrap().subagents =
mj_core::subagent::SubagentPolicy::SingleModel {
model: "parent-model".into(),
effort: Some(effort.into()),
};
config
.profiles
.get_mut("parent")
.unwrap()
.context_window_bytes = Some(1234);
config.subagents.max_concurrent = 3;
catalog.sync_now(&config).await;
catalog
.options_for(&config, "parent", Some("parent-model".into()))
.await
.unwrap();
catalog
.model_capabilities("parent".into(), "another".into())
.await
.unwrap();
}
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"defaults, effort edits, reopening and revisiting a model do not probe"
);
config.profiles.get_mut("parent").unwrap().home = "/another/home".into();
catalog.sync_now(&config).await;
assert_eq!(
calls.load(Ordering::SeqCst),
4,
"only changed installation is probed"
);
catalog
.model_capabilities("parent".into(), "another".into())
.await
.unwrap();
assert_eq!(
calls.load(Ordering::SeqCst),
5,
"model efforts for changed installation are rediscovered"
);
}
#[tokio::test]
async fn a_warm_pass_discovers_every_enabled_profile_for_every_parent() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let config = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
("private", HarnessKind::Kimi),
],
&["helper"],
);
catalog.sync_now(&config).await;
assert_eq!(
catalog
.candidates("parent")
.expect("a configuration is adopted"),
vec![
("helper".to_owned(), HarnessKind::Claude),
("parent".to_owned(), HarnessKind::Codex)
],
"a parent may delegate to itself and to the profiles listed as eligible"
);
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"every enabled profile is discovered once, whether or not a parent may offer it"
);
let choices = catalog
.capabilities(&["parent".to_owned(), "helper".to_owned()])
.await
.expect("the pass discovered both");
assert_eq!(choices[0].model.as_deref(), Some("parent-model"));
assert_eq!(choices[1].model.as_deref(), Some("helper-model"));
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"an answer from the warm catalogue discovers nothing"
);
}
#[tokio::test]
async fn a_warm_pass_ignores_the_deprecated_subagent_enabled_flag() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let mut config = test_config(&[("parent", HarnessKind::Codex)], &[]);
config.subagents.enabled = false;
catalog.sync_now(&config).await;
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the deprecated global switch no longer suppresses warming"
);
assert_eq!(
catalog
.candidates("parent")
.expect("a configuration is adopted"),
vec![("parent".to_owned(), HarnessKind::Codex)],
"a parent may always delegate to itself"
);
}
#[tokio::test]
async fn a_configuration_change_invalidates_and_re_warms_the_catalogue() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let config = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
],
&["helper"],
);
catalog.sync_now(&config).await;
catalog.sync_now(&config).await;
assert_eq!(calls.load(Ordering::SeqCst), 2);
let changed = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
("sibling", HarnessKind::Kimi),
],
&["helper", "sibling"],
);
catalog.sync_now(&changed).await;
assert_eq!(
catalog
.candidates("parent")
.expect("the replacement configuration is adopted")
.iter()
.map(|(id, _)| id.as_str())
.collect::<Vec<_>>(),
vec!["helper", "parent", "sibling"]
);
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"unchanged profiles are retained while the added profile is discovered"
);
let choices = catalog
.capabilities(&["sibling".to_owned()])
.await
.expect("the new pass discovered the profile it added");
assert_eq!(choices[0].model.as_deref(), Some("sibling-model"));
}
#[tokio::test]
async fn a_superseded_pass_starts_no_harness() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let first = test_config(&[("parent", HarnessKind::Codex)], &[]);
let second = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
],
&["helper"],
);
let (generation, key) = catalog
.claim_pass(&first)
.expect("the first pass is claimed");
catalog
.claim_pass(&second)
.expect("the change claims a new pass");
catalog.clone().warm_pass(generation, key).await;
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"a superseded pass must stop before it starts a harness"
);
}
#[tokio::test]
async fn a_discovery_of_a_superseded_generation_is_not_published() {
let calls = calls();
let probes_fail = fails();
let catalog = ProfileCatalog::with_probe(flag_probe(probes_fail.clone(), calls.clone()));
let first = test_config(&[("parent", HarnessKind::Codex)], &[]);
let (stale, _key) = catalog
.claim_pass(&first)
.expect("the first pass is claimed");
let stale_entry = catalog
.entry(stale, "parent")
.expect("the superseded pass planned its discovery");
let mut second = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
],
&["helper"],
);
second.profiles.get_mut("parent").unwrap().home = "/replacement/home".into();
catalog.sync_now(&second).await;
probes_fail.store(false, Ordering::SeqCst);
let Entry::Pending(attempt) = stale_entry else {
panic!("the pending discovery of the superseded generation is at hand")
};
let config = attempt.clone().await.expect("the probe itself succeeds");
catalog.remember(stale, "parent", &attempt, &config);
let before = calls.load(Ordering::SeqCst);
let choices = catalog
.capabilities(&["parent".to_owned()])
.await
.expect("the adopted configuration answers");
assert_eq!(choices[0].model.as_deref(), Some("parent-model"));
assert_eq!(
calls.load(Ordering::SeqCst),
before + 1,
"the superseded discovery was dropped, so this call discovers the profile itself"
);
}
#[tokio::test]
async fn a_failed_discovery_is_not_cached_and_the_next_call_retries_it() {
let calls = calls();
let probes_fail = fails();
let catalog = ProfileCatalog::with_probe(flag_probe(probes_fail.clone(), calls.clone()));
let config = test_config(&[("parent", HarnessKind::Codex)], &[]);
catalog.sync_now(&config).await;
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the pass tries the profile once"
);
let error = catalog
.capabilities(&["parent".to_owned()])
.await
.expect_err("the harness is still broken");
assert!(
format!("{error:#}").contains("harness is not installed"),
"{error:#}"
);
probes_fail.store(false, Ordering::SeqCst);
let choices = catalog
.capabilities(&["parent".to_owned()])
.await
.expect("the retry succeeds");
assert_eq!(choices[0].model.as_deref(), Some("parent-model"));
let discoveries = calls.load(Ordering::SeqCst);
catalog
.capabilities(&["parent".to_owned()])
.await
.expect("the retry's answer is kept");
assert_eq!(
calls.load(Ordering::SeqCst),
discoveries,
"a discovery the call waited for serves the calls after it"
);
}
#[tokio::test]
async fn a_call_waits_for_the_discovery_the_background_pass_is_running() {
let calls = calls();
let (started_tx, mut started) = tokio::sync::mpsc::unbounded_channel();
let (release, gate) = tokio::sync::watch::channel(false);
let catalog =
ProfileCatalog::with_probe(gated_probe(calls.clone(), started_tx, gate.clone()));
let config = test_config(
&[
("parent", HarnessKind::Codex),
("helper", HarnessKind::Claude),
],
&["helper"],
);
catalog.sync(&config);
let mut in_flight = Vec::new();
while in_flight.len() < 2 {
in_flight.push(started.recv().await.expect("the pass starts its probes"));
}
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"the pass discovers each enabled profile once: {in_flight:?}"
);
let call = {
let catalog = catalog.clone();
tokio::spawn(async move {
catalog
.capabilities(&["parent".to_owned(), "helper".to_owned()])
.await
})
};
tokio::task::yield_now().await;
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"the call must not start a second discovery for either profile"
);
release.send(true).expect("the gate is held by the test");
let choices = call
.await
.expect("the call task is not cancelled")
.expect("the pass's discoveries answer the call");
assert_eq!(choices[0].model.as_deref(), Some("parent-model"));
assert_eq!(choices[1].model.as_deref(), Some("helper-model"));
catalog
.capabilities(&["parent".to_owned()])
.await
.expect("the pass's discovery is now the catalogue's answer");
assert_eq!(
calls.load(Ordering::SeqCst),
2,
"the discovery the call waited for is kept"
);
}
#[tokio::test]
async fn an_answer_before_a_configuration_is_adopted_reports_that() {
let calls = calls();
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let error = catalog
.candidates("parent")
.expect_err("nothing has been adopted");
assert!(
format!("{error:#}").contains("has not adopted a configuration"),
"{error:#}"
);
let error = catalog
.capabilities(&["parent".to_owned()])
.await
.expect_err("nothing has been adopted");
assert!(
format!("{error:#}").contains("has not adopted a configuration"),
"{error:#}"
);
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"a call that cannot be answered starts no harness"
);
}
}