diff --git a/Cargo.lock b/Cargo.lock index 96d4a65e2..121ef5068 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -390,7 +390,6 @@ dependencies = [ name = "aionui-app" version = "0.1.41" dependencies = [ - "aion-config", "aionui-ai-agent", "aionui-api-types", "aionui-assets", diff --git a/crates/aionui-ai-agent/src/registry.rs b/crates/aionui-ai-agent/src/registry.rs index 2c7ed057c..62f16c4fe 100644 --- a/crates/aionui-ai-agent/src/registry.rs +++ b/crates/aionui-ai-agent/src/registry.rs @@ -5,8 +5,8 @@ //! rows all live there. The registry: //! //! - hydrates `select *` into memory at startup; -//! - probes each row's spawn command via `which()` so the `available` -//! field reflects PATH state right now (not a persisted column); +//! - projects startup availability from the latest persisted snapshot +//! without probing PATH; //! - exposes lookups the factory and routes use (`get`, //! `find_by_backend`, `list_by_agent_type`, etc.); //! - writes ACP handshake payloads back to the row through @@ -59,6 +59,7 @@ mod registry_tests; pub struct AgentRegistry { repo: Arc, by_id: RwLock>, + unavailable_reasons: RwLock>, /// MPSC sender shared with every forwarder in every `AcpAgentManager`. /// Draining happens in a single background task owned by this /// registry, so DB writes for the same (id, field) serialize. @@ -71,6 +72,7 @@ impl AgentRegistry { let this = Arc::new(Self { repo, by_id: RwLock::new(HashMap::new()), + unavailable_reasons: RwLock::new(HashMap::new()), catalog_tx: tx, }); @@ -144,11 +146,38 @@ impl AgentRegistry { return Ok(()); }; - if let Some((meta, _)) = decode_row(row) { + if let Some((mut meta, reason)) = decode_row(row, AvailabilityProjection::Cached) { + let existing_availability = self + .by_id + .read() + .await + .get(&meta.id) + .map(|existing| (existing.available, existing.resolved_command.clone())); + let reason = if let Some((available, resolved_command)) = existing_availability { + meta.available = available; + meta.resolved_command = resolved_command; + if meta.available { + None + } else { + self.unavailable_reasons.read().await.get(&meta.id).cloned().or(reason) + } + } else { + reason + }; + self.update_cached_unavailable_reason(&meta.id, reason).await; self.by_id.write().await.insert(meta.id.clone(), meta); } Ok(()) } + + async fn update_cached_unavailable_reason(&self, id: &str, reason: Option) { + let mut guard = self.unavailable_reasons.write().await; + if let Some(reason) = reason { + guard.insert(id.to_owned(), reason); + } else { + guard.remove(id); + } + } } impl AgentRegistry { @@ -159,8 +188,12 @@ impl AgentRegistry { tx: self.catalog_tx.clone(), } } - /// Reload every enabled row from the database and re-probe their - /// spawn commands on `$PATH`. + /// Reload every row from the database without probing spawn commands. + /// + /// Startup must be cheap and side-effect free: user-facing health checks + /// are explicit, so hydration projects availability from persisted + /// snapshots and only marks deterministic row states (disabled/no command) + /// as unavailable. pub async fn hydrate(&self) -> Result<(), AgentError> { let rows = self .repo @@ -169,11 +202,14 @@ impl AgentRegistry { .map_err(|e| AgentError::internal(format!("load agent_metadata: {e}")))?; let mut map = HashMap::with_capacity(rows.len()); + let mut reasons = HashMap::new(); for row in rows { - let Some((meta, reason)) = decode_row(row) else { + let Some((meta, reason)) = decode_row(row, AvailabilityProjection::Cached) else { continue; }; - log_probe_result(&meta, &reason); + if let Some(reason) = reason { + reasons.insert(meta.id.clone(), reason); + } map.insert(meta.id.clone(), meta); } // Snapshot the summary off the local map before transferring it @@ -181,6 +217,7 @@ impl AgentRegistry { // and we don't want that borrow to outlive the move. log_availability_summary(map.values(), "AgentRegistry hydrated"); *self.by_id.write().await = map; + *self.unavailable_reasons.write().await = reasons; Ok(()) } @@ -188,27 +225,43 @@ impl AgentRegistry { /// Useful after PATH has changed (e.g. `launchctl setenv`). pub async fn refresh_availability(&self) { let mut guard = self.by_id.write().await; + let mut reasons = HashMap::new(); for meta in guard.values_mut() { let (path, reason) = probe_with_reason(meta); meta.resolved_command = path; meta.available = meta.resolved_command.is_some() || is_internal_commandless_agent(meta); let reason = if meta.available { None } else { reason }; log_probe_result(meta, &reason); + if let Some(reason) = reason { + reasons.insert(meta.id.clone(), reason); + } } log_availability_summary(guard.values(), "AgentRegistry refresh_availability complete"); + *self.unavailable_reasons.write().await = reasons; } - /// Refetch every row from the repository, then re-resolve PATH. - /// - /// Called after any mutation that changed the set of rows on disk - /// (create/delete) or the spawn command of an existing row - /// (update). Pure refresh with no DB writes — just rebuilds the - /// in-memory snapshot so `list_all()` and `get()` return the latest - /// catalog state without waiting for the next process restart. - pub async fn invalidate_and_rehydrate(&self) -> Result<(), AgentError> { - self.hydrate().await?; - self.refresh_availability().await; - Ok(()) + /// Refetch and re-probe one row from the repository, leaving the rest of + /// the in-memory availability snapshot untouched. + pub async fn reload_one(&self, id: &str) -> Result, AgentError> { + let row = self + .repo + .get(id) + .await + .map_err(|e| AgentError::internal(format!("load agent_metadata '{id}': {e}")))?; + let Some(row) = row else { + self.by_id.write().await.remove(id); + self.unavailable_reasons.write().await.remove(id); + return Ok(None); + }; + let Some((meta, reason)) = decode_row(row, AvailabilityProjection::Probe) else { + self.by_id.write().await.remove(id); + self.unavailable_reasons.write().await.remove(id); + return Ok(None); + }; + log_probe_result(&meta, &reason); + self.update_cached_unavailable_reason(&meta.id, reason).await; + self.by_id.write().await.insert(meta.id.clone(), meta.clone()); + Ok(Some(meta)) } pub async fn get(&self, id: &str) -> Option { @@ -271,6 +324,7 @@ impl AgentRegistry { /// Management read model for settings surfaces that need to show /// official/custom rows even when unavailable. pub async fn list_management_rows(&self) -> Vec { + let reasons = self.unavailable_reasons.read().await.clone(); let mut rows: Vec = self .by_id .read() @@ -278,8 +332,9 @@ impl AgentRegistry { .values() .cloned() .map(|meta| { - let status = derive_management_status(&meta); - let diagnostics = derive_management_diagnostics(&meta, status); + let reason = reasons.get(&meta.id); + let status = derive_management_status(&meta, reason); + let diagnostics = derive_management_diagnostics(&meta, status, reason); let handshake = meta.handshake; AgentManagementRow { id: meta.id, @@ -325,6 +380,52 @@ impl AgentRegistry { rows } + pub async fn management_row_by_id(&self, id: &str) -> Option { + let reason = self.unavailable_reasons.read().await.get(id).cloned(); + let meta = self.by_id.read().await.get(id).cloned()?; + let status = derive_management_status(&meta, reason.as_ref()); + let diagnostics = derive_management_diagnostics(&meta, status, reason.as_ref()); + let handshake = meta.handshake.clone(); + Some(AgentManagementRow { + id: meta.id, + icon: meta.icon, + name: meta.name, + name_i18n: meta.name_i18n, + description: meta.description, + description_i18n: meta.description_i18n, + backend: meta.backend, + agent_type: meta.agent_type, + agent_source: meta.agent_source, + agent_source_info: meta.agent_source_info, + enabled: meta.enabled, + installed: meta.available, + command: meta.command, + args: meta.args, + env: Vec::new(), + native_skills_dirs: meta.native_skills_dirs, + behavior_policy: meta.behavior_policy, + yolo_id: meta.yolo_id, + config_options: handshake.config_options.clone(), + available_modes: handshake.available_modes.clone(), + available_models: handshake.available_models.clone(), + sort_order: meta.sort_order, + team_capable: meta.team_capable, + status, + last_check_status: meta.last_check_status, + last_check_kind: meta.last_check_kind, + last_check_error_code: diagnostics.error_code, + last_check_error_message: diagnostics.error_message, + last_check_error_details: diagnostics.details, + last_check_guidance: diagnostics.guidance, + last_check_latency_ms: meta.last_check_latency_ms, + last_check_at: meta.last_check_at, + last_success_at: meta.last_success_at, + last_failure_at: meta.last_failure_at, + has_command_override: meta.has_command_override, + env_override_key_count: meta.env_override_key_count, + }) + } + /// Like [`Self::list_all_including_hidden`] but pairs every row /// with a freshly-computed availability reason so callers (the /// `doctor` command, diagnostic UIs) can explain *why* a row is @@ -338,17 +439,14 @@ impl AgentRegistry { /// when `available = true` would just confuse the caller, so we /// suppress it here. pub async fn diagnostic_snapshot(&self) -> Vec<(AgentMetadata, Option)> { + let reasons = self.unavailable_reasons.read().await.clone(); let mut rows: Vec<(AgentMetadata, Option)> = self .by_id .read() .await .values() .map(|m| { - let reason = if m.available { - None - } else { - probe_resolved_command(m).err() - }; + let reason = if m.available { None } else { reasons.get(&m.id).cloned() }; (m.clone(), reason) }) .collect(); @@ -370,7 +468,7 @@ impl AgentRegistry { /// keeps both uninstalled CLIs and rows that most recently failed /// ACP/session admission out of visible legacy catalog reads. fn is_visible(meta: &AgentMetadata) -> bool { - meta.enabled && matches!(derive_management_status(meta), AgentManagementStatus::Online) + meta.enabled && matches!(derive_management_status(meta, None), AgentManagementStatus::Online) } /// Extract and trim a command override, filtering out empty strings. @@ -390,12 +488,23 @@ fn parse_env_override(raw: &Option) -> Option> { serde_json::from_str::>(s).ok() } -/// Turn a DB row into the public `AgentMetadata`, probing the command -/// on disk so `available` reflects the current PATH state. Returns -/// the probe reason alongside the row so the caller can log a single -/// uniform `(meta, reason)` line per agent without re-running the -/// probe. -fn decode_row(row: AgentMetadataRow) -> Option<(AgentMetadata, Option)> { +#[derive(Clone, Copy)] +enum AvailabilityProjection { + /// Use only deterministic row state plus persisted health snapshots. + Cached, + /// Resolve the spawn command against the current runtime environment. + Probe, +} + +/// Turn a DB row into the public `AgentMetadata`. +/// +/// Callers choose whether availability should come from persisted snapshots +/// or from a live command probe. Startup hydration uses cached projection; +/// explicit refresh and single-agent health-check reloads use probe. +fn decode_row( + row: AgentMetadataRow, + availability: AvailabilityProjection, +) -> Option<(AgentMetadata, Option)> { // Extract override fields before row is partially moved let command_override_raw = row.command_override.clone(); let env_override_raw = row.env_override.clone(); @@ -504,10 +613,10 @@ fn decode_row(row: AgentMetadataRow) -> Option<(AgentMetadata, Option apply_cached_availability(&mut meta), + AvailabilityProjection::Probe => apply_probe_availability(&mut meta), + }; Some((meta, reason)) } @@ -519,6 +628,95 @@ fn is_internal_commandless_agent(meta: &AgentMetadata) -> bool { meta.enabled && meta.command.is_none() && meta.agent_source == AgentSource::Internal } +fn apply_probe_availability(meta: &mut AgentMetadata) -> Option { + let (path, reason) = probe_with_reason(meta); + meta.resolved_command = path; + meta.available = meta.resolved_command.is_some() || is_internal_commandless_agent(meta); + if meta.available { None } else { reason } +} + +fn apply_cached_availability(meta: &mut AgentMetadata) -> Option { + if !meta.enabled { + meta.available = false; + return Some(UnavailableReason::Disabled); + } + if is_internal_commandless_agent(meta) { + meta.available = true; + return None; + } + if meta.command.as_deref().filter(|s| !s.is_empty()).is_none() { + meta.available = false; + return Some(UnavailableReason::NoCommand); + } + if cached_snapshot_indicates_missing(meta) { + meta.available = false; + return cached_unavailable_reason(meta); + } + if !has_availability_snapshot(meta) { + meta.available = false; + return None; + } + meta.available = true; + None +} + +fn has_availability_snapshot(meta: &AgentMetadata) -> bool { + meta.last_check_status.is_some() + || meta.last_check_kind.is_some() + || meta.last_check_error_code.is_some() + || meta.last_check_at.is_some() + || meta.last_success_at.is_some() + || meta.last_failure_at.is_some() +} + +fn cached_snapshot_indicates_missing(meta: &AgentMetadata) -> bool { + matches!( + meta.last_check_error_code.as_deref(), + Some( + "command_not_found" + | "command_missing" + | "primary_missing" + | "bridge_missing" + | "managed_runtime_unavailable" + | "no_command" + | "disabled" + ) + ) +} + +fn cached_unavailable_reason(meta: &AgentMetadata) -> Option { + match meta.last_check_error_code.as_deref()? { + "disabled" => Some(UnavailableReason::Disabled), + "no_command" => Some(UnavailableReason::NoCommand), + "bridge_missing" => meta + .agent_source_info + .bridge_binary + .clone() + .map(|bridge| UnavailableReason::BridgeMissing { bridge }), + "primary_missing" => meta + .agent_source_info + .binary_name + .clone() + .map(|binary| UnavailableReason::PrimaryMissing { binary }), + "command_not_found" | "command_missing" => Some(UnavailableReason::CommandMissing { + command: meta + .agent_source_info + .binary_name + .clone() + .or_else(|| meta.command.clone()) + .unwrap_or_else(|| "command".to_owned()), + }), + "managed_runtime_unavailable" => Some(UnavailableReason::ManagedRuntimeUnavailable { + resource: meta.backend.clone().unwrap_or_else(|| "runtime".to_owned()), + detail: meta + .last_check_error_message + .clone() + .unwrap_or_else(|| "managed runtime was unavailable during the last health check".to_owned()), + }), + _ => None, + } +} + /// Wrapper around [`probe_resolved_command`] that returns both the /// resolved path (if any) and the failure reason as a tuple, so the /// hydrate / refresh loops can persist the path and emit a single @@ -645,14 +843,21 @@ fn parse_last_check_kind(raw: Option<&str>) -> Option { }) } -fn derive_management_status(meta: &AgentMetadata) -> AgentManagementStatus { +fn derive_management_status(meta: &AgentMetadata, reason: Option<&UnavailableReason>) -> AgentManagementStatus { if !meta.available { - return AgentManagementStatus::Missing; + if reason.is_some() || has_availability_snapshot(meta) { + return AgentManagementStatus::Missing; + } + return AgentManagementStatus::Unchecked; + } + if is_internal_commandless_agent(meta) { + return AgentManagementStatus::Online; } match meta.last_check_status { Some(AgentSnapshotCheckStatus::Offline) => AgentManagementStatus::Offline, - _ => AgentManagementStatus::Online, + Some(AgentSnapshotCheckStatus::Online) => AgentManagementStatus::Online, + None => AgentManagementStatus::Unchecked, } } @@ -663,9 +868,13 @@ struct ManagementDiagnostics { guidance: Option, } -fn derive_management_diagnostics(meta: &AgentMetadata, status: AgentManagementStatus) -> ManagementDiagnostics { +fn derive_management_diagnostics( + meta: &AgentMetadata, + status: AgentManagementStatus, + reason: Option<&UnavailableReason>, +) -> ManagementDiagnostics { let derived_reason = if matches!(status, AgentManagementStatus::Missing) { - probe_resolved_command(meta).err() + reason.cloned() } else { None }; @@ -1198,10 +1407,9 @@ mod tests { } /// `diagnostic_snapshot` returns one entry per row, populates a - /// reason for every unavailable row, and leaves available rows - /// without one. The CI host doesn't have the seeded CLIs - /// installed, so the bridge/CLI rows are reliably unavailable - /// here — the assertion exploits that to lock the contract. + /// reason for rows known unavailable by probe/cache, and leaves + /// available or unchecked rows without one. Unchecked rows are + /// expected after startup hydration avoids live probing. #[tokio::test] async fn diagnostic_snapshot_pairs_rows_with_reasons() { let reg = registry().await; @@ -1212,6 +1420,7 @@ mod tests { match (meta.available, reason) { (true, None) => {} (false, Some(_)) => {} + (false, None) if matches!(derive_management_status(meta, None), AgentManagementStatus::Unchecked) => {} (true, Some(r)) => panic!("available row {} has unexpected reason {:?}", meta.id, r), (false, None) => panic!( "unavailable row {} (source={:?}) is missing a reason", @@ -1320,7 +1529,7 @@ mod tests { created_at: 0, updated_at: 0, }; - let (meta, _) = super::decode_row(row).expect("decodes"); + let (meta, _) = super::decode_row(row, super::AvailabilityProjection::Cached).expect("decodes"); assert_eq!(meta.command.as_deref(), Some("/opt/factory/bin/droid")); } @@ -1368,7 +1577,7 @@ mod tests { created_at: 0, updated_at: 0, }; - let (meta, reason) = super::decode_row(row).expect("decodes"); + let (meta, reason) = super::decode_row(row, super::AvailabilityProjection::Cached).expect("decodes"); assert_eq!(meta.agent_type, AgentType::Aionrs); assert_eq!(meta.agent_source, AgentSource::Internal); assert_eq!(meta.command, None); @@ -1423,7 +1632,7 @@ mod tests { created_at: 0, updated_at: 0, }; - let (meta, _) = super::decode_row(row).expect("decodes"); + let (meta, _) = super::decode_row(row, super::AvailabilityProjection::Cached).expect("decodes"); let names: Vec<&str> = meta.env.iter().map(|e| e.name.as_str()).collect(); assert!(names.contains(&"BASE")); assert!(names.contains(&"ANTHROPIC_API_KEY")); diff --git a/crates/aionui-ai-agent/src/registry_tests.rs b/crates/aionui-ai-agent/src/registry_tests.rs index ad52d1f41..87289883c 100644 --- a/crates/aionui-ai-agent/src/registry_tests.rs +++ b/crates/aionui-ai-agent/src/registry_tests.rs @@ -197,6 +197,7 @@ async fn management_rows_derive_missing_diagnostics_from_probe_reason() { let registry = AgentRegistry::new(repo); registry.hydrate().await.unwrap(); + registry.refresh_availability().await; let row = registry .list_management_rows() @@ -224,6 +225,57 @@ async fn management_rows_derive_missing_diagnostics_from_probe_reason() { ); } +#[tokio::test] +async fn management_rows_mark_unchecked_agents_unchecked_without_probe() { + let db = init_database_memory().await.unwrap(); + let repo: Arc = Arc::new(SqliteAgentMetadataRepository::new(db.pool().clone())); + + repo.upsert(&UpsertAgentMetadataParams { + id: "agent-unchecked-cli", + icon: None, + name: "Unchecked CLI Agent", + name_i18n: None, + description: None, + description_i18n: None, + backend: Some("custom"), + agent_type: "acp", + agent_source: "custom", + agent_source_info: Some(r#"{"binary_name":"unchecked-cli"}"#), + enabled: true, + command: Some("unchecked-cli"), + args: Some("[]"), + env: Some("[]"), + native_skills_dirs: None, + behavior_policy: None, + yolo_id: None, + agent_capabilities: None, + auth_methods: None, + config_options: None, + available_modes: None, + available_models: None, + available_commands: None, + sort_order: 100, + }) + .await + .unwrap(); + + let registry = AgentRegistry::new(repo); + registry.hydrate().await.unwrap(); + + let row = registry + .list_management_rows() + .await + .into_iter() + .find(|item| item.id == "agent-unchecked-cli") + .unwrap(); + + let row_json = serde_json::to_value(&row).unwrap(); + assert_eq!(row_json["status"].as_str(), Some("unchecked")); + assert!(!row.installed); + assert!(row.last_check_status.is_none()); + assert!(row.last_check_error_code.is_none()); +} + #[tokio::test] async fn hydrate_continues_when_agent_metadata_config_options_has_invalid_utf8() { let db = init_database_memory().await.unwrap(); diff --git a/crates/aionui-ai-agent/src/routes/agent.rs b/crates/aionui-ai-agent/src/routes/agent.rs index 98db9537a..ad6614cd4 100644 --- a/crates/aionui-ai-agent/src/routes/agent.rs +++ b/crates/aionui-ai-agent/src/routes/agent.rs @@ -5,7 +5,6 @@ //! Endpoints: //! //! - `GET /api/agents/management` — list diagnostics-first agent rows -//! - `POST /api/agents/refresh` — refresh agent list (e.g. after new agent is added to the system) //! - `POST /api/agents/custom/try-connect` — test custom agent configuration (e.g. ACP connection) use axum::Router; @@ -28,7 +27,6 @@ pub fn agent_routes(state: AgentRouterState) -> Router { Router::new() .route("/api/agents/logos", get(list_agent_logos)) .route("/api/agents/management", get(list_management_agents)) - .route("/api/agents/refresh", post(refresh_agents)) .route("/api/agents/{id}/health-check", post(health_check_by_id)) .route("/api/agents/provider-health-check", post(provider_health_check)) .route("/api/agents/{id}/enabled", patch(set_agent_enabled)) @@ -42,15 +40,6 @@ pub fn agent_routes(state: AgentRouterState) -> Router { .with_state(state) } -async fn refresh_agents( - State(state): State, - Extension(_user): Extension, -) -> Result>>, ApiError> { - Ok(Json(ApiResponse::ok( - state.service.refresh_agents().await.map_err(agent_error_to_api_error)?, - ))) -} - async fn list_agent_logos( State(state): State, Extension(_user): Extension, diff --git a/crates/aionui-ai-agent/src/services/agent.rs b/crates/aionui-ai-agent/src/services/agent.rs index c65342776..2de04fa5b 100644 --- a/crates/aionui-ai-agent/src/services/agent.rs +++ b/crates/aionui-ai-agent/src/services/agent.rs @@ -15,9 +15,7 @@ use std::path::PathBuf; use std::sync::Arc; -use aionui_api_types::{ - AgentLogoEntry, AgentManagementRow, AgentMetadata, ProviderHealthCheckRequest, ProviderHealthCheckResponse, -}; +use aionui_api_types::{AgentLogoEntry, AgentManagementRow, ProviderHealthCheckRequest, ProviderHealthCheckResponse}; use aionui_db::IProviderRepository; use aionui_realtime::EventBroadcaster; @@ -76,17 +74,6 @@ impl AgentService { // Agent operations impl AgentService { - pub async fn refresh_agents(&self) -> Result, AgentError> { - self.registry.refresh_availability().await; - Ok(self - .registry - .list_all() - .await - .into_iter() - .filter(|agent| agent.agent_type.supports_new_conversation()) - .collect()) - } - pub async fn list_management_agents(&self) -> Result, AgentError> { Ok(self.availability.list_management_rows().await) } diff --git a/crates/aionui-ai-agent/src/services/availability/mod.rs b/crates/aionui-ai-agent/src/services/availability/mod.rs index 09d7a275e..f81293419 100644 --- a/crates/aionui-ai-agent/src/services/availability/mod.rs +++ b/crates/aionui-ai-agent/src/services/availability/mod.rs @@ -53,17 +53,15 @@ impl AgentAvailabilityService { } pub async fn list_management_rows(&self) -> Vec { - self.registry.refresh_availability().await; self.registry.list_management_rows().await } pub async fn run_manual_health_check(&self, id: &str) -> Result { - self.registry.invalidate_and_rehydrate().await?; let meta = self .registry - .get(id) + .reload_one(id) .await - .ok_or_else(|| AgentError::not_found(format!("Agent '{id}' not found")))?; + .and_then(|row| row.ok_or_else(|| AgentError::not_found(format!("Agent '{id}' not found"))))?; if !meta.available { return self @@ -113,11 +111,7 @@ impl AgentAvailabilityService { } pub async fn management_row_by_id(&self, id: &str) -> Option { - self.registry - .list_management_rows() - .await - .into_iter() - .find(|row| row.id == id) + self.registry.management_row_by_id(id).await } async fn persist_snapshot(&self, id: &str, snapshot: &AvailabilitySnapshot) -> Result<(), AgentError> { @@ -156,7 +150,7 @@ impl AgentAvailabilityService { .update_availability_snapshot(id, ¶ms) .await .map_err(|error| AgentError::internal(format!("repo.update_availability_snapshot: {error}")))?; - self.registry.invalidate_and_rehydrate().await?; + self.registry.reload_one(id).await?; Ok(()) } } diff --git a/crates/aionui-ai-agent/src/services/custom.rs b/crates/aionui-ai-agent/src/services/custom.rs index 577431996..c7dc6037c 100644 --- a/crates/aionui-ai-agent/src/services/custom.rs +++ b/crates/aionui-ai-agent/src/services/custom.rs @@ -107,8 +107,8 @@ impl AgentService { if !removed { return Err(AgentError::not_found(format!("Agent '{id}' not found"))); } - if let Err(err) = self.registry().invalidate_and_rehydrate().await { - warn!(agent_id = %id, error = %err, "registry rehydrate failed after delete_custom_agent"); + if let Err(err) = self.registry().reload_one(id).await { + warn!(agent_id = %id, error = %err, "registry reload failed after delete_custom_agent"); } Ok(()) } @@ -123,8 +123,8 @@ impl AgentService { if !updated { return Err(AgentError::not_found(format!("Agent '{id}' not found"))); } - if let Err(err) = self.registry().invalidate_and_rehydrate().await { - warn!(agent_id = %id, error = %err, "registry rehydrate failed after set_agent_enabled"); + if let Err(err) = self.registry().reload_one(id).await { + warn!(agent_id = %id, error = %err, "registry reload failed after set_agent_enabled"); } self.registry() .get(id) @@ -195,9 +195,9 @@ impl AgentService { .map_err(|e| AgentError::internal(format!("repo.upsert: {e}")))?; self.registry() - .invalidate_and_rehydrate() + .reload_one(id) .await - .map_err(|e| AgentError::internal(format!("registry rehydrate: {e}")))?; + .map_err(|e| AgentError::internal(format!("registry reload: {e}")))?; self.registry() .get(id) diff --git a/crates/aionui-ai-agent/tests/agent_availability_integration.rs b/crates/aionui-ai-agent/tests/agent_availability_integration.rs index 01b6c4fb9..b53a1f7db 100644 --- a/crates/aionui-ai-agent/tests/agent_availability_integration.rs +++ b/crates/aionui-ai-agent/tests/agent_availability_integration.rs @@ -1,11 +1,18 @@ use std::sync::Arc; -use aionui_ai_agent::AgentRegistry; +use aionui_ai_agent::{AgentRegistry, AgentService}; use aionui_api_types::{AgentManagementStatus, AgentSnapshotCheckKind, AgentSnapshotCheckStatus}; use aionui_db::{ - IAgentMetadataRepository, SqliteAgentMetadataRepository, UpdateAgentAvailabilitySnapshotParams, - UpsertAgentMetadataParams, init_database_memory, + IAgentMetadataRepository, IProviderRepository, SqliteAgentMetadataRepository, SqliteProviderRepository, + UpdateAgentAvailabilitySnapshotParams, UpsertAgentMetadataParams, init_database_memory, }; +use aionui_realtime::EventBroadcaster; + +struct NoopBroadcaster; + +impl EventBroadcaster for NoopBroadcaster { + fn broadcast(&self, _msg: aionui_api_types::WebSocketMessage) {} +} fn custom_params<'a>( id: &'a str, @@ -41,6 +48,14 @@ fn custom_params<'a>( } } +fn agent_service( + registry: Arc, + provider_repo: Arc, + data_dir: std::path::PathBuf, +) -> Arc { + AgentService::new(registry, Arc::new(NoopBroadcaster), provider_repo, [0; 32], data_dir) +} + #[tokio::test] async fn management_rows_derive_missing_available_and_unavailable_statuses() { let db = init_database_memory().await.unwrap(); @@ -107,6 +122,7 @@ async fn management_rows_derive_missing_available_and_unavailable_statuses() { let registry = AgentRegistry::new(repo); registry.hydrate().await.unwrap(); + registry.refresh_availability().await; let rows = registry.list_management_rows().await; @@ -131,3 +147,215 @@ async fn management_rows_derive_missing_available_and_unavailable_statuses() { assert_eq!(available.last_check_kind, Some(AgentSnapshotCheckKind::Scheduled)); assert_eq!(available.last_check_latency_ms, Some(120)); } + +#[tokio::test] +async fn hydrate_uses_persisted_availability_without_reprobing_path() { + let db = init_database_memory().await.unwrap(); + let repo: Arc = Arc::new(SqliteAgentMetadataRepository::new(db.pool().clone())); + let temp = tempfile::tempdir().unwrap(); + let command_path = temp.path().join("startup-cached-agent-command"); + std::fs::write(&command_path, "#!/bin/sh\nexit 0\n").unwrap(); + let command = command_path.to_string_lossy().to_string(); + let source_info = serde_json::json!({ "binary_name": command }).to_string(); + + repo.upsert(&custom_params( + "agent-startup-cached", + "Startup Cached Agent", + &command, + &source_info, + )) + .await + .unwrap(); + repo.update_availability_snapshot( + "agent-startup-cached", + &UpdateAgentAvailabilitySnapshotParams { + last_check_status: Some("online"), + last_check_kind: Some("manual"), + last_check_error_code: None, + last_check_error_message: None, + last_check_guidance: None, + last_check_latency_ms: Some(42), + last_check_at: Some(1_750_000_100_000), + last_success_at: Some(1_750_000_100_000), + last_failure_at: None, + }, + ) + .await + .unwrap(); + + std::fs::remove_file(&command_path).unwrap(); + + let registry = AgentRegistry::new(repo); + registry.hydrate().await.unwrap(); + + let rows = registry.list_management_rows().await; + let cached = rows.iter().find(|row| row.id == "agent-startup-cached").unwrap(); + + assert_eq!(cached.status, AgentManagementStatus::Online); + assert!(cached.installed, "startup hydrate should not refresh PATH"); + assert_eq!(cached.last_check_status, Some(AgentSnapshotCheckStatus::Online)); +} + +#[tokio::test] +async fn management_list_uses_cached_availability_without_reprobing_path() { + let db = init_database_memory().await.unwrap(); + let repo: Arc = Arc::new(SqliteAgentMetadataRepository::new(db.pool().clone())); + let provider_repo: Arc = Arc::new(SqliteProviderRepository::new(db.pool().clone())); + let temp = tempfile::tempdir().unwrap(); + let command_path = temp.path().join("cached-agent-command"); + std::fs::write(&command_path, "#!/bin/sh\nexit 0\n").unwrap(); + let command = command_path.to_string_lossy().to_string(); + let source_info = serde_json::json!({ "binary_name": command }).to_string(); + + repo.upsert(&custom_params("agent-cached", "Cached Agent", &command, &source_info)) + .await + .unwrap(); + + let registry = AgentRegistry::new(repo); + registry.hydrate().await.unwrap(); + + std::fs::remove_file(&command_path).unwrap(); + + let service = agent_service(registry, provider_repo, temp.path().to_path_buf()); + let rows = service.list_management_agents().await.unwrap(); + let cached = rows.iter().find(|row| row.id == "agent-cached").unwrap(); + + assert_eq!(cached.status, AgentManagementStatus::Unchecked); + assert!(!cached.installed, "management list should not refresh PATH on read"); +} + +#[tokio::test] +async fn manual_health_check_does_not_refresh_unrelated_agents() { + let db = init_database_memory().await.unwrap(); + let repo: Arc = Arc::new(SqliteAgentMetadataRepository::new(db.pool().clone())); + let provider_repo: Arc = Arc::new(SqliteProviderRepository::new(db.pool().clone())); + let temp = tempfile::tempdir().unwrap(); + let unrelated_path = temp.path().join("unrelated-agent-command"); + std::fs::write(&unrelated_path, "#!/bin/sh\nexit 0\n").unwrap(); + let unrelated_command = unrelated_path.to_string_lossy().to_string(); + let unrelated_source_info = serde_json::json!({ "binary_name": unrelated_command }).to_string(); + + repo.upsert(&custom_params( + "agent-unrelated", + "Unrelated Agent", + &unrelated_command, + &unrelated_source_info, + )) + .await + .unwrap(); + repo.upsert(&custom_params( + "agent-target-missing", + "Target Missing Agent", + "aionui-definitely-missing-health-check-target", + r#"{"binary_name":"aionui-definitely-missing-health-check-target"}"#, + )) + .await + .unwrap(); + + let registry = AgentRegistry::new(repo); + registry.hydrate().await.unwrap(); + std::fs::remove_file(&unrelated_path).unwrap(); + + let service = agent_service(registry.clone(), provider_repo, temp.path().to_path_buf()); + service.health_check_agent_by_id("agent-target-missing").await.unwrap(); + + let rows = registry.list_management_rows().await; + let unrelated = rows.iter().find(|row| row.id == "agent-unrelated").unwrap(); + + assert_eq!(unrelated.status, AgentManagementStatus::Unchecked); + assert!( + !unrelated.installed, + "single-agent health check should not refresh unrelated agents" + ); +} + +#[tokio::test] +async fn custom_enabled_toggle_does_not_refresh_unrelated_agents() { + let db = init_database_memory().await.unwrap(); + let repo: Arc = Arc::new(SqliteAgentMetadataRepository::new(db.pool().clone())); + let provider_repo: Arc = Arc::new(SqliteProviderRepository::new(db.pool().clone())); + let temp = tempfile::tempdir().unwrap(); + let unrelated_path = temp.path().join("unrelated-agent-command"); + std::fs::write(&unrelated_path, "#!/bin/sh\nexit 0\n").unwrap(); + let unrelated_command = unrelated_path.to_string_lossy().to_string(); + let unrelated_source_info = serde_json::json!({ "binary_name": unrelated_command }).to_string(); + + repo.upsert(&custom_params( + "agent-unrelated", + "Unrelated Agent", + &unrelated_command, + &unrelated_source_info, + )) + .await + .unwrap(); + repo.upsert(&custom_params( + "agent-target-toggle", + "Target Toggle Agent", + "aionui-target-toggle-command", + r#"{"binary_name":"aionui-target-toggle-command"}"#, + )) + .await + .unwrap(); + + let registry = AgentRegistry::new(repo); + registry.hydrate().await.unwrap(); + std::fs::remove_file(&unrelated_path).unwrap(); + + let service = agent_service(registry.clone(), provider_repo, temp.path().to_path_buf()); + service.set_agent_enabled("agent-target-toggle", false).await.unwrap(); + + let rows = registry.list_management_rows().await; + let unrelated = rows.iter().find(|row| row.id == "agent-unrelated").unwrap(); + + assert_eq!(unrelated.status, AgentManagementStatus::Unchecked); + assert!( + !unrelated.installed, + "custom enabled toggle should not refresh unrelated agents" + ); +} + +#[tokio::test] +async fn custom_delete_does_not_refresh_unrelated_agents() { + let db = init_database_memory().await.unwrap(); + let repo: Arc = Arc::new(SqliteAgentMetadataRepository::new(db.pool().clone())); + let provider_repo: Arc = Arc::new(SqliteProviderRepository::new(db.pool().clone())); + let temp = tempfile::tempdir().unwrap(); + let unrelated_path = temp.path().join("unrelated-agent-command"); + std::fs::write(&unrelated_path, "#!/bin/sh\nexit 0\n").unwrap(); + let unrelated_command = unrelated_path.to_string_lossy().to_string(); + let unrelated_source_info = serde_json::json!({ "binary_name": unrelated_command }).to_string(); + + repo.upsert(&custom_params( + "agent-unrelated", + "Unrelated Agent", + &unrelated_command, + &unrelated_source_info, + )) + .await + .unwrap(); + repo.upsert(&custom_params( + "agent-target-delete", + "Target Delete Agent", + "aionui-target-delete-command", + r#"{"binary_name":"aionui-target-delete-command"}"#, + )) + .await + .unwrap(); + + let registry = AgentRegistry::new(repo); + registry.hydrate().await.unwrap(); + std::fs::remove_file(&unrelated_path).unwrap(); + + let service = agent_service(registry.clone(), provider_repo, temp.path().to_path_buf()); + service.delete_custom_agent("agent-target-delete").await.unwrap(); + + let rows = registry.list_management_rows().await; + let unrelated = rows.iter().find(|row| row.id == "agent-unrelated").unwrap(); + + assert_eq!(unrelated.status, AgentManagementStatus::Unchecked); + assert!( + !unrelated.installed, + "custom delete should not refresh unrelated agents" + ); + assert!(rows.iter().all(|row| row.id != "agent-target-delete")); +} diff --git a/crates/aionui-api-types/src/agent_discovery.rs b/crates/aionui-api-types/src/agent_discovery.rs index cce7fef1b..3c3a25b88 100644 --- a/crates/aionui-api-types/src/agent_discovery.rs +++ b/crates/aionui-api-types/src/agent_discovery.rs @@ -112,6 +112,7 @@ pub enum AgentManagementStatus { Missing, Online, Offline, + Unchecked, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] @@ -406,6 +407,8 @@ mod tests { fn agent_management_status_serializes_snake_case() { let value = serde_json::to_value(AgentManagementStatus::Offline).unwrap(); assert_eq!(value, json!("offline")); + let value = serde_json::to_value(AgentManagementStatus::Unchecked).unwrap(); + assert_eq!(value, json!("unchecked")); } } diff --git a/crates/aionui-app/Cargo.toml b/crates/aionui-app/Cargo.toml index cbbc721e1..89fe1df71 100644 --- a/crates/aionui-app/Cargo.toml +++ b/crates/aionui-app/Cargo.toml @@ -35,7 +35,6 @@ aionui-team-prompts.workspace = true aionui-cron.workspace = true aionui-assistant.workspace = true aionui-runtime.workspace = true -aion-config.workspace = true axum.workspace = true chrono.workspace = true dirs.workspace = true diff --git a/crates/aionui-app/assets/builtin-skills/aionui-config/SKILL.md b/crates/aionui-app/assets/builtin-skills/aionui-config/SKILL.md index f1113b1dd..b67c73bb8 100644 --- a/crates/aionui-app/assets/builtin-skills/aionui-config/SKILL.md +++ b/crates/aionui-app/assets/builtin-skills/aionui-config/SKILL.md @@ -567,17 +567,14 @@ python3 scripts/aionui_api.py get /api/settings/client # confirm — PUT retur the `/management` sub-path. Each row is rich: alongside `id`, `name`, `enabled` (toggled on), `installed` (spawn command resolvable on `$PATH`), `team_capable` (can run in a team), `backend`, `agent_type`, and a `status` of `online` / -`offline` / `missing`, it also carries `config_options`, `available_modes`, -`available_models` (when the engine advertises them), plus `last_check_*` -diagnostics. Check `installed` (and `status`) before binding an assistant to that -engine (via its `agent_id` — see *Picking the engine* above). - -> What the management row does **not** carry is the rest of the engine -> `handshake` — `agent_capabilities`, `auth_methods`, `available_commands`. For -> those, `POST /api/agents/refresh` re-scans agents and returns each one's full -> metadata (`available` + `handshake`). The at-a-glance modes/models are already -> on the management row; reach for `refresh` when you need capabilities, auth -> methods, or the command list. +`offline` / `missing` / `unchecked`, it also carries `config_options`, +`available_modes`, `available_models` (when the engine advertises them), plus +`last_check_*` diagnostics. `unchecked` means no persisted connectivity check has +run for that row yet. Check `installed` (and `status`) before binding an +assistant to that engine (via its `agent_id` — see *Picking the engine* above). + +The management row is the supported engine catalog surface. Do not call legacy +agent refresh endpoints; connectivity checks are explicit per-agent operations. --- diff --git a/crates/aionui-app/src/bootstrap/tracing_init.rs b/crates/aionui-app/src/bootstrap/tracing_init.rs index 41a7ba519..12b4825df 100644 --- a/crates/aionui-app/src/bootstrap/tracing_init.rs +++ b/crates/aionui-app/src/bootstrap/tracing_init.rs @@ -4,7 +4,11 @@ //! subscriber registration that should never be invoked from tests or //! external consumers of the library. -use std::path::{Path, PathBuf}; +use std::{ + fs::{File, OpenOptions}, + io::{self, Write}, + path::{Path, PathBuf}, +}; use chrono::Datelike; use tracing_subscriber::{EnvFilter, Layer, fmt, layer::SubscriberExt, util::SubscriberInitExt}; @@ -98,19 +102,7 @@ pub fn init_tracing(log_dir: &Path, log_level: Option<&str>) -> Result) -> Result) -> Result PathBuf { - let now = chrono::Local::now(); + dated_log_dir_for(log_root, LogDate::today()) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +struct LogDate { + year: i32, + month: u32, + day: u32, +} + +impl LogDate { + fn today() -> Self { + let now = chrono::Local::now(); + Self { + year: now.year(), + month: now.month(), + day: now.day(), + } + } + + fn file_name(self, suffix: &str) -> String { + format!("{:04}-{:02}-{:02}.{}", self.year, self.month, self.day, suffix) + } +} + +fn dated_log_dir_for(log_root: &Path, date: LogDate) -> PathBuf { log_root - .join(format!("{:04}", now.year())) - .join(format!("{:02}", now.month())) - .join(format!("{:02}", now.day())) + .join(format!("{:04}", date.year)) + .join(format!("{:02}", date.month)) + .join(format!("{:02}", date.day)) +} + +fn dated_log_file_path(log_root: &Path, date: LogDate, suffix: &str) -> PathBuf { + dated_log_dir_for(log_root, date).join(date.file_name(suffix)) +} + +struct DailyDatedLogWriter { + log_root: PathBuf, + filename_suffix: &'static str, + date_provider: Box LogDate + Send + Sync>, + active_date: Option, + active_file: Option, +} + +impl DailyDatedLogWriter { + fn new(log_root: PathBuf, filename_suffix: &'static str) -> Self { + Self::new_with_date_provider(log_root, filename_suffix, Box::new(LogDate::today)) + } + + fn new_with_date_provider( + log_root: PathBuf, + filename_suffix: &'static str, + date_provider: Box LogDate + Send + Sync>, + ) -> Self { + Self { + log_root, + filename_suffix, + date_provider, + active_date: None, + active_file: None, + } + } + + fn active_file(&mut self) -> io::Result<&mut File> { + let date = (self.date_provider)(); + if self.active_date != Some(date) { + let file_path = dated_log_file_path(&self.log_root, date, self.filename_suffix); + if let Some(parent) = file_path.parent() { + std::fs::create_dir_all(parent)?; + } + self.active_file = Some(OpenOptions::new().create(true).append(true).open(file_path)?); + self.active_date = Some(date); + } + + self.active_file + .as_mut() + .ok_or_else(|| io::Error::other("log file was not opened")) + } +} + +impl Write for DailyDatedLogWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.active_file()?.write_all(buf)?; + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + if let Some(file) = self.active_file.as_mut() { + file.flush()?; + } + Ok(()) + } } #[cfg(test)] @@ -246,4 +329,41 @@ mod tests { assert!(parts[1].chars().all(|ch| ch.is_ascii_digit())); assert!(parts[2].chars().all(|ch| ch.is_ascii_digit())); } + + #[test] + fn dated_file_writer_moves_new_day_files_into_matching_day_directory() { + let tmp = tempfile::tempdir().expect("temp dir"); + let first_day = LogDate { + year: 2026, + month: 7, + day: 2, + }; + let second_day = LogDate { + year: 2026, + month: 7, + day: 3, + }; + let days = std::sync::Arc::new(std::sync::Mutex::new(vec![second_day, first_day])); + let mut writer = DailyDatedLogWriter::new_with_date_provider( + tmp.path().to_path_buf(), + "aioncore.log", + Box::new({ + let days = std::sync::Arc::clone(&days); + move || days.lock().expect("date queue").pop().expect("date") + }), + ); + + std::io::Write::write_all(&mut writer, b"july 2\n").expect("write first day"); + std::io::Write::write_all(&mut writer, b"july 3\n").expect("write second day"); + std::io::Write::flush(&mut writer).expect("flush"); + + let first_path = tmp.path().join("2026/07/02/2026-07-02.aioncore.log"); + let second_path = tmp.path().join("2026/07/03/2026-07-03.aioncore.log"); + assert_eq!(std::fs::read_to_string(first_path).expect("first day log"), "july 2\n"); + assert_eq!( + std::fs::read_to_string(second_path).expect("second day log"), + "july 3\n" + ); + assert!(!tmp.path().join("2026/07/02/2026-07-03.aioncore.log").exists()); + } } diff --git a/crates/aionui-app/src/router/state.rs b/crates/aionui-app/src/router/state.rs index 0ff4c0e38..10e633f46 100644 --- a/crates/aionui-app/src/router/state.rs +++ b/crates/aionui-app/src/router/state.rs @@ -300,7 +300,6 @@ pub fn build_assistant_state(services: &AppServices) -> AssistantRouterState { #[async_trait::async_trait] impl AssistantAgentCatalogPort for RegistryAssistantAgentCatalog { async fn list_management_agents(&self) -> Result, AssistantError> { - self.registry.refresh_availability().await; Ok(self.registry.list_management_rows().await) } } diff --git a/crates/aionui-app/tests/acp_e2e.rs b/crates/aionui-app/tests/acp_e2e.rs index 4f343c1b8..dc0f1ec94 100644 --- a/crates/aionui-app/tests/acp_e2e.rs +++ b/crates/aionui-app/tests/acp_e2e.rs @@ -1,6 +1,6 @@ //! E2E integration tests for ACP management routes. //! -//! Tests cover: agents list, agents/refresh, agents/test, +//! Tests cover: agents list, legacy agents/refresh removal, agents/test, //! and session-bound routes (mode/model). mod common; @@ -35,17 +35,13 @@ async fn management_list_returns_array() { } #[tokio::test] -async fn refresh_agents_returns_array() { +async fn legacy_refresh_agents_endpoint_is_not_found() { let (mut app, services) = build_app().await; let (token, csrf) = setup_and_login(&mut app, &services, "user1", "pass123").await; let req = json_with_token("POST", "/api/agents/refresh", json!({}), &token, &csrf); let resp = app.oneshot(req).await.unwrap(); - assert_eq!(resp.status(), StatusCode::OK); - - let body = body_json(resp).await; - assert_eq!(body["success"], true); - assert!(body["data"].is_array()); + assert_eq!(resp.status(), StatusCode::NOT_FOUND); } #[tokio::test] @@ -106,7 +102,8 @@ async fn management_list_includes_missing_custom_agents() { }) .await .unwrap(); - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); + services.agent_registry.refresh_availability().await; let req = get_with_token("/api/agents/management", &token); let resp = app.oneshot(req).await.unwrap(); @@ -159,7 +156,7 @@ async fn management_list_marks_rows_with_unavailable_snapshot() { repo.update_availability_snapshot( "custom-unavailable-agent", &UpdateAgentAvailabilitySnapshotParams { - last_check_status: Some("unavailable"), + last_check_status: Some("offline"), last_check_kind: Some("scheduled"), last_check_error_code: Some("acp_init_failed"), last_check_error_message: Some("Synthetic unavailable snapshot"), @@ -172,7 +169,7 @@ async fn management_list_marks_rows_with_unavailable_snapshot() { ) .await .unwrap(); - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); let req = get_with_token("/api/agents/management", &token); let resp = app.oneshot(req).await.unwrap(); @@ -184,7 +181,7 @@ async fn management_list_marks_rows_with_unavailable_snapshot() { .iter() .find(|item| item["id"].as_str() == Some("custom-unavailable-agent")) .expect("management list should include unavailable rows"); - assert_eq!(row["status"], "online"); + assert_eq!(row["status"], "offline"); } #[tokio::test] @@ -232,7 +229,8 @@ async fn health_check_by_id_returns_missing_status_for_uninstalled_agent() { }) .await .unwrap(); - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); + services.agent_registry.refresh_availability().await; let req = json_with_token( "POST", diff --git a/crates/aionui-app/tests/agent_integration_e2e.rs b/crates/aionui-app/tests/agent_integration_e2e.rs index 7b5a83df2..a0a72d6f6 100644 --- a/crates/aionui-app/tests/agent_integration_e2e.rs +++ b/crates/aionui-app/tests/agent_integration_e2e.rs @@ -264,7 +264,8 @@ async fn management_endpoint_keeps_deprecated_runtime_rows_for_diagnostics() { ] { upsert_visible_agent_metadata(&services, id, agent_type).await; } - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); + services.agent_registry.refresh_availability().await; let req = get_with_token("/api/agents/management", &token); let resp = app.oneshot(req).await.unwrap(); @@ -402,7 +403,8 @@ async fn agent_logos_endpoint_includes_disabled_and_missing_rows() { }) .await .unwrap(); - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); + services.agent_registry.refresh_availability().await; let req = get_with_token("/api/agents/logos", &token); let resp = app.oneshot(req).await.unwrap(); @@ -624,7 +626,8 @@ async fn agent_overrides_roundtrip_and_management_summary() { let (mut app, services, _mock_tm) = build_app_with_mock_tasks().await; let (token, csrf) = setup_and_login(&mut app, &services, "admin", "Pass123!").await; upsert_visible_agent_metadata(&services, "ovr-agent", "acp").await; - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); + services.agent_registry.refresh_availability().await; // PUT overrides let body = json!({ @@ -672,7 +675,8 @@ async fn internal_aion_cli_rejects_overrides() { let (mut app, services, _mock_tm) = build_app_with_mock_tasks().await; let (token, csrf) = setup_and_login(&mut app, &services, "admin", "Pass123!").await; upsert_visible_agent_metadata(&services, "632f31d2", "aionrs").await; - services.agent_registry.invalidate_and_rehydrate().await.unwrap(); + services.agent_registry.hydrate().await.unwrap(); + services.agent_registry.refresh_availability().await; let command_body = json!({ "command_override": "irm https://claude.ai/install.ps1 | iex",