From 4dbb570ea615a7991bb25bd08ebefb183521c290 Mon Sep 17 00:00:00 2001 From: daniel Date: Tue, 1 Sep 2026 00:38:04 +0100 Subject: [PATCH 1/4] feat(skills): announce catalog changes --- src/protocols/acp.rs | 29 ++- src/protocols/acp/skill_catalog.rs | 290 +++++++++++++++++++++++++++++ src/protocols/acp/v2.rs | 21 ++- src/runtime.rs | 60 ++++-- 4 files changed, 375 insertions(+), 25 deletions(-) create mode 100644 src/protocols/acp/skill_catalog.rs diff --git a/src/protocols/acp.rs b/src/protocols/acp.rs index 1050713..0a1ec36 100644 --- a/src/protocols/acp.rs +++ b/src/protocols/acp.rs @@ -43,6 +43,7 @@ use tokio::{ time::timeout, }; +mod skill_catalog; pub mod v2; use crate::{ @@ -849,15 +850,18 @@ impl Server { let tasks = driver.tasks.clone(); let structured_completion = driver.structured_completion; let canonical_transcript = driver.canonical_transcript; + let skill_catalog = skill_catalog::SkillCatalogMonitor::new(&driver.skills); let (tx, rx) = mpsc::channel(8); let actor = SessionActor { session_id: session_id.clone(), + runtime: Arc::clone(&self.runtime), integration: Arc::clone(&self.integration), binding, driver: driver.driver, tasks: driver.tasks, background_jobs: background_jobs.clone(), structured_completion, + skill_catalog, adapter: driver.adapter, catalog, commands: rx, @@ -1077,12 +1081,14 @@ pub(super) async fn detach_compose_call( struct SessionActor { session_id: agentkit_acp::SessionId, + runtime: Arc, integration: Arc, binding: SessionBindingGuard, driver: LoopDriver, tasks: TaskManagerHandle, background_jobs: BackgroundJobs, structured_completion: bool, + skill_catalog: skill_catalog::SkillCatalogMonitor, adapter: SelectableAdapter, catalog: Vec, commands: mpsc::Receiver, @@ -1093,12 +1099,14 @@ struct SessionActor { async fn session_actor(actor: SessionActor) { let SessionActor { session_id, + runtime, integration, binding, mut driver, tasks, background_jobs, structured_completion, + mut skill_catalog, adapter, catalog, mut commands, @@ -1114,9 +1122,12 @@ async fn session_actor(actor: SessionActor) { biased; command = commands.recv() => match command { Some(Command::Prompt { request, reply }) => { + let skills = runtime.current_skills().await; let result = drive_prompt( &session_id, + &skills, &integration, + &mut skill_catalog, &mut driver, request, &tasks, @@ -1420,9 +1431,12 @@ fn record_acp_loop_failure( } } +#[allow(clippy::too_many_arguments)] async fn drive_prompt( session_id: &agentkit_acp::SessionId, + skills: &[agentkit_tool_skills::Skill], integration: &AcpIntegration, + skill_catalog: &mut skill_catalog::SkillCatalogMonitor, driver: &mut LoopDriver, request: PromptRequest, tasks: &TaskManagerHandle, @@ -1434,8 +1448,8 @@ async fn drive_prompt( } background_jobs.begin_turn(); let items = integration.input_port().prompt_to_items(&request)?; - driver - .submit_input(items) + skill_catalog + .submit(skills, items, |items| driver.submit_input(items)) .map_err(|error| record_acp_loop_failure(session_id, &error))?; let response = match drive_until_pause( session_id, @@ -2705,9 +2719,12 @@ pub(super) mod tests { let task_manager = AsyncTaskManager::new(); let tasks = task_manager.handle(); let background_jobs = BackgroundJobs::default(); + let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]); let response = drive_prompt( &acp_session_id, + &[], &integration, + &mut skill_catalog, &mut driver, PromptRequest::new( acp_session_id.clone(), @@ -2788,6 +2805,7 @@ pub(super) mod tests { .await .unwrap(); let background_jobs = BackgroundJobs::default(); + let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]); let request = PromptRequest::new( acp_session_id.clone(), vec![agentkit_acp::ContentBlock::Text( @@ -2796,7 +2814,9 @@ pub(super) mod tests { ); let prompt = drive_prompt( &acp_session_id, + &[], &integration, + &mut skill_catalog, &mut driver, request, &tasks, @@ -2928,14 +2948,19 @@ pub(super) mod tests { let (turn_states_tx, mut turn_states_rx) = mpsc::unbounded_channel(); let test_mcp = crate::tools::mcp::empty(); let mcp_events = test_mcp.subscribe(acp_session_id.to_string()); + let root = tempfile::tempdir().unwrap(); + let runtime = Runtime::new(root.path(), "gpt-5.4").unwrap(); + let skills = runtime.current_skills().await; let actor = tokio::spawn(session_actor(SessionActor { session_id: acp_session_id.clone(), + runtime, integration: Arc::clone(&integration), binding: SessionBindingGuard::new(Arc::clone(&integration), acp_session_id.clone()), driver, tasks: tasks.clone(), background_jobs: background_jobs.clone(), structured_completion: false, + skill_catalog: skill_catalog::SkillCatalogMonitor::new(&skills), adapter: SelectableAdapter::new(crate::ProviderKind::OpenAiSubscription, "gpt-5.4") .unwrap(), catalog: Vec::new(), diff --git a/src/protocols/acp/skill_catalog.rs b/src/protocols/acp/skill_catalog.rs new file mode 100644 index 0000000..50299bd --- /dev/null +++ b/src/protocols/acp/skill_catalog.rs @@ -0,0 +1,290 @@ +use std::collections::BTreeMap; + +use agentkit_core::Item; +use agentkit_tool_skills::Skill; +use serde_json::json; + +const MAX_NOTIFICATION_BYTES: usize = 2_048; +const MAX_REPORTED_CHANGES: usize = 64; + +type Catalog = BTreeMap; + +pub(super) struct SkillCatalogMonitor { + baseline: Catalog, +} + +impl SkillCatalogMonitor { + pub(super) fn new(skills: &[Skill]) -> Self { + Self { + baseline: catalog(skills), + } + } + + pub(super) fn submit( + &mut self, + skills: &[Skill], + mut user_items: Vec, + submit: impl FnOnce(Vec) -> Result<(), E>, + ) -> Result<(), E> { + let current = catalog(skills); + if let Some(notification) = notification(&self.baseline, ¤t) { + user_items.insert(0, Item::notification(notification)); + } + submit(user_items)?; + self.baseline = current; + Ok(()) + } +} + +fn catalog(skills: &[Skill]) -> Catalog { + skills + .iter() + .map(|skill| (skill.name.clone(), fingerprint(skill))) + .collect() +} + +fn fingerprint(skill: &Skill) -> blake3::Hash { + let mut resources = skill + .resources + .iter() + .map(|path| path.strip_prefix(&skill.base_dir).unwrap_or(path)) + .collect::>(); + resources.sort(); + let semantic = json!({ + "description": skill.description, + "body": skill.body, + "frontmatter": { + "license": skill.frontmatter.license, + "compatibility": skill.frontmatter.compatibility, + "metadata": skill.frontmatter.metadata, + "allowed_tools": skill.frontmatter.allowed_tools, + }, + }); + let mut fingerprint = blake3::Hasher::new(); + fingerprint.update( + &serde_json::to_vec(&semantic).expect("skill fingerprint data is JSON serializable"), + ); + for resource in resources { + let path = resource.as_os_str().as_encoded_bytes(); + fingerprint.update(&(path.len() as u64).to_le_bytes()); + fingerprint.update(path); + } + fingerprint.finalize() +} + +fn notification(previous: &Catalog, current: &Catalog) -> Option { + if previous == current { + return None; + } + + let added = current + .keys() + .filter(|name| !previous.contains_key(*name)) + .map(String::as_str) + .collect::>(); + let changed = current + .iter() + .filter(|(name, fingerprint)| previous.get(*name).is_some_and(|old| old != *fingerprint)) + .map(|(name, _)| name.as_str()) + .collect::>(); + let removed = previous + .keys() + .filter(|name| !current.contains_key(*name)) + .map(String::as_str) + .collect::>(); + let total = added.len() + changed.len() + removed.len(); + let omission_reserve = format!("; {total} more change(s) omitted").len(); + let content_limit = MAX_NOTIFICATION_BYTES.saturating_sub(omission_reserve); + let mut message = String::from("Skill catalog update (informational)."); + let mut reported = 0; + append_group(&mut message, "Added", &added, &mut reported, content_limit); + append_group( + &mut message, + "Changed", + &changed, + &mut reported, + content_limit, + ); + append_group( + &mut message, + "Removed", + &removed, + &mut reported, + content_limit, + ); + if reported < total { + message.push_str(&format!("; {} more change(s) omitted", total - reported)); + } + Some(message) +} + +fn append_group( + message: &mut String, + label: &str, + names: &[&str], + reported: &mut usize, + content_limit: usize, +) { + if names.is_empty() || *reported >= MAX_REPORTED_CHANGES { + return; + } + let heading = format!("; {label}: "); + if message.len() + heading.len() > content_limit { + return; + } + let heading_start = message.len(); + message.push_str(&heading); + let mut group_count = 0; + for name in names { + if *reported >= MAX_REPORTED_CHANGES { + break; + } + let separator = if group_count == 0 { "" } else { ", " }; + if message.len() + separator.len() + name.len() > content_limit { + break; + } + message.push_str(separator); + message.push_str(name); + group_count += 1; + *reported += 1; + } + if group_count == 0 { + message.truncate(heading_start); + } +} + +#[cfg(test)] +mod tests { + use std::path::PathBuf; + + use agentkit_core::{ItemKind, Part, TextPart}; + use agentkit_tool_skills::SkillFrontmatter; + + use super::*; + + fn skill(name: &str, body: &str) -> Skill { + let base_dir = PathBuf::from("/skills").join(name); + Skill { + name: name.into(), + description: format!("Description for {name}"), + location: base_dir.join("SKILL.md"), + base_dir, + body: body.into(), + resources: Vec::new(), + frontmatter: SkillFrontmatter::default(), + } + } + + fn notification_text(items: &[Item]) -> Option<&str> { + let item = items + .first() + .filter(|item| item.kind == ItemKind::Notification)?; + match item.parts.first() { + Some(Part::Text(TextPart { text, .. })) => Some(text), + _ => None, + } + } + + #[test] + fn fingerprint_covers_semantic_fields_and_sorted_resource_paths() { + let mut original = skill("semantic", "body"); + original.resources = vec![ + original.base_dir.join("z.txt"), + original.base_dir.join("a.txt"), + ]; + + let mut reordered = original.clone(); + reordered.resources.reverse(); + assert_eq!(fingerprint(&original), fingerprint(&reordered)); + + let mut changed = original.clone(); + changed.description.push_str(" changed"); + assert_ne!(fingerprint(&original), fingerprint(&changed)); + changed = original.clone(); + changed.body.push_str(" changed"); + assert_ne!(fingerprint(&original), fingerprint(&changed)); + changed = original.clone(); + changed + .frontmatter + .metadata + .insert("key".into(), "value".into()); + assert_ne!(fingerprint(&original), fingerprint(&changed)); + changed = original.clone(); + changed.resources.push(changed.base_dir.join("new.txt")); + assert_ne!(fingerprint(&original), fingerprint(&changed)); + } + + #[test] + fn baseline_and_unchanged_catalog_do_not_notify() { + let skills = vec![skill("existing", "body")]; + let mut monitor = SkillCatalogMonitor::new(&skills); + monitor + .submit( + &skills, + vec![Item::text(ItemKind::User, "hello")], + |items| { + assert_eq!(items.len(), 1); + assert_eq!(items[0].kind, ItemKind::User); + Ok::<_, ()>(()) + }, + ) + .unwrap(); + } + + #[test] + fn reports_added_changed_and_removed_names_in_the_user_batch() { + let before = vec![skill("removed", "old"), skill("changed", "old")]; + let after = vec![skill("added", "new"), skill("changed", "new")]; + let mut monitor = SkillCatalogMonitor::new(&before); + monitor + .submit(&after, vec![Item::text(ItemKind::User, "hello")], |items| { + assert_eq!(items.len(), 2); + assert_eq!(items[0].kind, ItemKind::Notification); + assert_eq!(items[1].kind, ItemKind::User); + let text = notification_text(&items).unwrap(); + assert!(text.starts_with("Skill catalog update (informational).")); + assert!(text.contains("Added: added")); + assert!(text.contains("Changed: changed")); + assert!(text.contains("Removed: removed")); + Ok::<_, ()>(()) + }) + .unwrap(); + } + + #[test] + fn notification_is_sorted_and_bounded() { + let after = (0..100) + .rev() + .map(|index| skill(&format!("skill-{index:03}"), "body")) + .collect::>(); + let mut monitor = SkillCatalogMonitor::new(&[]); + monitor + .submit(&after, vec![Item::text(ItemKind::User, "hello")], |items| { + let text = notification_text(&items).unwrap(); + assert!(text.len() <= MAX_NOTIFICATION_BYTES); + assert!(text.find("skill-000").unwrap() < text.find("skill-001").unwrap()); + assert!(text.contains("more change(s) omitted")); + Ok::<_, ()>(()) + }) + .unwrap(); + } + + #[test] + fn failed_submission_does_not_advance_baseline() { + let after = vec![skill("added", "body")]; + let mut monitor = SkillCatalogMonitor::new(&[]); + assert!( + monitor + .submit(&after, vec![Item::text(ItemKind::User, "first")], |_| Err( + () + )) + .is_err() + ); + monitor + .submit(&after, vec![Item::text(ItemKind::User, "retry")], |items| { + assert!(notification_text(&items).unwrap().contains("Added: added")); + Ok::<_, ()>(()) + }) + .unwrap(); + } +} diff --git a/src/protocols/acp/v2.rs b/src/protocols/acp/v2.rs index 283ca82..f5dbc27 100644 --- a/src/protocols/acp/v2.rs +++ b/src/protocols/acp/v2.rs @@ -34,7 +34,7 @@ use crate::{ use super::{ CancelBackgroundRequest, CancelBackgroundResponse, DetachComposeRequest, DetachComposeResponse, - SessionRegistry, + SessionRegistry, skill_catalog, }; const PAGE_SIZE: usize = 100; @@ -675,6 +675,7 @@ impl Server { let catalog = model_catalog(¤t).await; let config_options = v2_config_options(¤t, reasoning, &catalog); let canonical_transcript = driver.canonical_transcript; + let skill_catalog = skill_catalog::SkillCatalogMonitor::new(&driver.skills); let background_jobs = driver.background_jobs.clone(); let tasks = driver.tasks.clone(); let structured_completion = driver.structured_completion; @@ -683,6 +684,7 @@ impl Server { let busy = Arc::new(AtomicBool::new(false)); let actor = SessionActor { session_id: session_id.clone(), + runtime: Arc::clone(&self.runtime), integration: Arc::clone(&self.integration), handle: handle.clone(), busy: Arc::clone(&busy), @@ -692,6 +694,7 @@ impl Server { tasks: driver.tasks, background_jobs: background_jobs.clone(), structured_completion, + skill_catalog, adapter: driver.adapter, catalog, commands: rx, @@ -947,6 +950,7 @@ impl Server { struct SessionActor { session_id: wire::SessionId, + runtime: Arc, integration: Arc, handle: AcpSessionHandle, busy: Arc, @@ -956,6 +960,7 @@ struct SessionActor { tasks: TaskManagerHandle, background_jobs: BackgroundJobs, structured_completion: bool, + skill_catalog: skill_catalog::SkillCatalogMonitor, adapter: SelectableAdapter, catalog: Vec, commands: mpsc::Receiver, @@ -965,6 +970,7 @@ struct SessionActor { async fn session_actor(actor: SessionActor) { let SessionActor { session_id, + runtime, integration, handle, busy, @@ -974,6 +980,7 @@ async fn session_actor(actor: SessionActor) tasks, background_jobs, structured_completion, + mut skill_catalog, adapter, catalog, mut commands, @@ -985,10 +992,13 @@ async fn session_actor(actor: SessionActor) biased; command = commands.recv() => match command { Some(Command::Prompt(command)) => { + let skills = runtime.current_skills().await; let result = prepare_prompt( &session_id, + &skills, &integration, &handle, + &mut skill_catalog, &mut driver, command, &sink, @@ -1066,8 +1076,10 @@ async fn session_actor(actor: SessionActor) #[allow(clippy::too_many_arguments)] async fn prepare_prompt( session_id: &wire::SessionId, + skills: &[agentkit_tool_skills::Skill], integration: &AcpIntegration, handle: &AcpSessionHandle, + skill_catalog: &mut skill_catalog::SkillCatalogMonitor, driver: &mut LoopDriver, command: PromptCommand, sink: &impl AcpSessionUpdateSink, @@ -1089,8 +1101,8 @@ async fn prepare_prompt( } background_jobs.begin_turn(); let prepared = integration.prompt_to_items(&request).and_then(|items| { - driver - .submit_input(items) + skill_catalog + .submit(skills, items, |items| driver.submit_input(items)) .map_err(|error| map_loop_error(session_id, &error))?; integration.begin_prompt(session_id) }); @@ -2684,11 +2696,14 @@ mod tests { let task_manager = AsyncTaskManager::new(); let tasks = task_manager.handle(); let background_jobs = BackgroundJobs::default(); + let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]); let (result, ()) = tokio::join!( prepare_prompt( &session_id, + &[], &integration, &handle, + &mut skill_catalog, &mut driver, command, &sink, diff --git a/src/runtime.rs b/src/runtime.rs index 268f02b..71892a8 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -21,7 +21,7 @@ use agentkit_task_manager::{AsyncTaskManager, RoutingDecision, TaskManager, Task use agentkit_tool_compose::{ BackendRun, ComposeBackend, ComposeConfig, ComposeOutcome, ComposeTool, RunletBackend, }; -use agentkit_tool_skills::SkillRegistry; +use agentkit_tool_skills::{Skill, SkillRegistry}; use agentkit_tools_core::{ PermissionRequest, Tool, ToolContext, ToolError, ToolExecutionOutcome, ToolName, ToolRequest, ToolResult, ToolSource, ToolSpec, @@ -217,6 +217,7 @@ impl Drop for SessionClaim { pub(crate) struct AcpDriver { pub driver: LoopDriver, + pub skills: Vec, pub tasks: TaskManagerHandle, pub background_jobs: BackgroundJobs, pub structured_completion: bool, @@ -731,6 +732,17 @@ impl Runtime { ) } + pub(crate) async fn current_skills(&self) -> Vec { + let registry = build_skill_registry( + &self.root, + &self.skill_package_roots, + &self.skill_directories, + ) + .discover_skills() + .await; + registry.skills().into_iter().cloned().collect() + } + fn compose_with(&self, depth: usize, subagents: Subagents) -> ComposeOnly { self.compose_with_jobs( depth, @@ -1051,6 +1063,7 @@ impl Runtime { ) .map_err(AcpRuntimeError::Loop)?; let skills = self.fresh_skills(); + let skill_catalog = self.current_skills().await; let compactor = crate::compaction::automatic( adapter.clone(), self.agentkit_telemetry(), @@ -1089,6 +1102,7 @@ impl Runtime { .map_err(|error| AcpRuntimeError::Loop(error.to_string()))?; let driver = AcpDriver { driver, + skills: skill_catalog, tasks, background_jobs, structured_completion: self.base_depth > 0, @@ -1144,6 +1158,14 @@ fn build_skill_tools( package_roots: &[PathBuf], skill_directories: &[PathBuf], ) -> Arc { + Arc::new(build_skill_registry(root, package_roots, skill_directories)) +} + +fn build_skill_registry( + root: &Path, + package_roots: &[PathBuf], + skill_directories: &[PathBuf], +) -> SkillRegistry { let default_roots = default_skill_roots(root); let canonical_defaults = default_roots .iter() @@ -1159,26 +1181,24 @@ fn build_skill_tools( .collect::>(); let mut roots = default_roots; roots.extend(skill_directories.iter().cloned()); - Arc::new(SkillRegistry::from_paths(roots).with_filter( - move |skill: &agentkit_tool_skills::Skill| { - let Ok(base) = skill.base_dir.canonicalize() else { - return false; - }; - if canonical_defaults + SkillRegistry::from_paths(roots).with_filter(move |skill: &agentkit_tool_skills::Skill| { + let Ok(base) = skill.base_dir.canonicalize() else { + return false; + }; + if canonical_defaults + .iter() + .any(|directory| base.starts_with(directory)) + { + return true; + } + let Ok(location) = skill.location.canonicalize() else { + return false; + }; + canonical_plugin_skills.contains(&base) + && canonical_package_roots .iter() - .any(|directory| base.starts_with(directory)) - { - return true; - } - let Ok(location) = skill.location.canonicalize() else { - return false; - }; - canonical_plugin_skills.contains(&base) - && canonical_package_roots - .iter() - .any(|package| location.starts_with(package)) - }, - )) + .any(|package| location.starts_with(package)) + }) } fn default_skill_roots(root: &Path) -> Vec { From 782ac626e125b5e161a99681f2231e268d484f74 Mon Sep 17 00:00:00 2001 From: daniel Date: Tue, 1 Sep 2026 01:39:36 +0100 Subject: [PATCH 2/4] fix(skills): propagate catalog serialization errors --- src/protocols/acp.rs | 18 ++++-- src/protocols/acp/skill_catalog.rs | 100 +++++++++++++++++++---------- src/protocols/acp/v2.rs | 12 +++- 3 files changed, 89 insertions(+), 41 deletions(-) diff --git a/src/protocols/acp.rs b/src/protocols/acp.rs index 0a1ec36..0e35302 100644 --- a/src/protocols/acp.rs +++ b/src/protocols/acp.rs @@ -850,7 +850,8 @@ impl Server { let tasks = driver.tasks.clone(); let structured_completion = driver.structured_completion; let canonical_transcript = driver.canonical_transcript; - let skill_catalog = skill_catalog::SkillCatalogMonitor::new(&driver.skills); + let skill_catalog = skill_catalog::SkillCatalogMonitor::new(&driver.skills) + .map_err(|error| record_acp_runtime_failure(&session_id, "skill_catalog", error))?; let (tx, rx) = mpsc::channel(8); let actor = SessionActor { session_id: session_id.clone(), @@ -1450,7 +1451,14 @@ async fn drive_prompt( let items = integration.input_port().prompt_to_items(&request)?; skill_catalog .submit(skills, items, |items| driver.submit_input(items)) - .map_err(|error| record_acp_loop_failure(session_id, &error))?; + .map_err(|error| match error { + skill_catalog::SubmitError::Catalog(error) => { + record_acp_runtime_failure(session_id, "skill_catalog", error) + } + skill_catalog::SubmitError::Submit(error) => { + record_acp_loop_failure(session_id, &error) + } + })?; let response = match drive_until_pause( session_id, integration, @@ -2719,7 +2727,7 @@ pub(super) mod tests { let task_manager = AsyncTaskManager::new(); let tasks = task_manager.handle(); let background_jobs = BackgroundJobs::default(); - let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]); + let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]).unwrap(); let response = drive_prompt( &acp_session_id, &[], @@ -2805,7 +2813,7 @@ pub(super) mod tests { .await .unwrap(); let background_jobs = BackgroundJobs::default(); - let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]); + let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]).unwrap(); let request = PromptRequest::new( acp_session_id.clone(), vec![agentkit_acp::ContentBlock::Text( @@ -2960,7 +2968,7 @@ pub(super) mod tests { tasks: tasks.clone(), background_jobs: background_jobs.clone(), structured_completion: false, - skill_catalog: skill_catalog::SkillCatalogMonitor::new(&skills), + skill_catalog: skill_catalog::SkillCatalogMonitor::new(&skills).unwrap(), adapter: SelectableAdapter::new(crate::ProviderKind::OpenAiSubscription, "gpt-5.4") .unwrap(), catalog: Vec::new(), diff --git a/src/protocols/acp/skill_catalog.rs b/src/protocols/acp/skill_catalog.rs index 50299bd..0f3b195 100644 --- a/src/protocols/acp/skill_catalog.rs +++ b/src/protocols/acp/skill_catalog.rs @@ -2,22 +2,43 @@ use std::collections::BTreeMap; use agentkit_core::Item; use agentkit_tool_skills::Skill; -use serde_json::json; +use serde::Serialize; const MAX_NOTIFICATION_BYTES: usize = 2_048; const MAX_REPORTED_CHANGES: usize = 64; type Catalog = BTreeMap; +#[derive(Debug)] +pub(super) enum SubmitError { + Catalog(serde_json::Error), + Submit(E), +} + +#[derive(Serialize)] +struct Fingerprint<'a> { + description: &'a str, + body: &'a str, + frontmatter: FingerprintFrontmatter<'a>, +} + +#[derive(Serialize)] +struct FingerprintFrontmatter<'a> { + license: Option<&'a str>, + compatibility: Option<&'a str>, + metadata: &'a BTreeMap, + allowed_tools: Option<&'a str>, +} + pub(super) struct SkillCatalogMonitor { baseline: Catalog, } impl SkillCatalogMonitor { - pub(super) fn new(skills: &[Skill]) -> Self { - Self { - baseline: catalog(skills), - } + pub(super) fn new(skills: &[Skill]) -> serde_json::Result { + Ok(Self { + baseline: catalog(skills)?, + }) } pub(super) fn submit( @@ -25,51 +46,49 @@ impl SkillCatalogMonitor { skills: &[Skill], mut user_items: Vec, submit: impl FnOnce(Vec) -> Result<(), E>, - ) -> Result<(), E> { - let current = catalog(skills); + ) -> Result<(), SubmitError> { + let current = catalog(skills).map_err(SubmitError::Catalog)?; if let Some(notification) = notification(&self.baseline, ¤t) { user_items.insert(0, Item::notification(notification)); } - submit(user_items)?; + submit(user_items).map_err(SubmitError::Submit)?; self.baseline = current; Ok(()) } } -fn catalog(skills: &[Skill]) -> Catalog { +fn catalog(skills: &[Skill]) -> serde_json::Result { skills .iter() - .map(|skill| (skill.name.clone(), fingerprint(skill))) + .map(|skill| Ok((skill.name.clone(), fingerprint(skill)?))) .collect() } -fn fingerprint(skill: &Skill) -> blake3::Hash { +fn fingerprint(skill: &Skill) -> serde_json::Result { let mut resources = skill .resources .iter() .map(|path| path.strip_prefix(&skill.base_dir).unwrap_or(path)) .collect::>(); resources.sort(); - let semantic = json!({ - "description": skill.description, - "body": skill.body, - "frontmatter": { - "license": skill.frontmatter.license, - "compatibility": skill.frontmatter.compatibility, - "metadata": skill.frontmatter.metadata, - "allowed_tools": skill.frontmatter.allowed_tools, + let semantic = Fingerprint { + description: &skill.description, + body: &skill.body, + frontmatter: FingerprintFrontmatter { + license: skill.frontmatter.license.as_deref(), + compatibility: skill.frontmatter.compatibility.as_deref(), + metadata: &skill.frontmatter.metadata, + allowed_tools: skill.frontmatter.allowed_tools.as_deref(), }, - }); + }; let mut fingerprint = blake3::Hasher::new(); - fingerprint.update( - &serde_json::to_vec(&semantic).expect("skill fingerprint data is JSON serializable"), - ); + fingerprint.update(&serde_json::to_vec(&semantic)?); for resource in resources { let path = resource.as_os_str().as_encoded_bytes(); fingerprint.update(&(path.len() as u64).to_le_bytes()); fingerprint.update(path); } - fingerprint.finalize() + Ok(fingerprint.finalize()) } fn notification(previous: &Catalog, current: &Catalog) -> Option { @@ -195,29 +214,44 @@ mod tests { let mut reordered = original.clone(); reordered.resources.reverse(); - assert_eq!(fingerprint(&original), fingerprint(&reordered)); + assert_eq!( + fingerprint(&original).unwrap(), + fingerprint(&reordered).unwrap() + ); let mut changed = original.clone(); changed.description.push_str(" changed"); - assert_ne!(fingerprint(&original), fingerprint(&changed)); + assert_ne!( + fingerprint(&original).unwrap(), + fingerprint(&changed).unwrap() + ); changed = original.clone(); changed.body.push_str(" changed"); - assert_ne!(fingerprint(&original), fingerprint(&changed)); + assert_ne!( + fingerprint(&original).unwrap(), + fingerprint(&changed).unwrap() + ); changed = original.clone(); changed .frontmatter .metadata .insert("key".into(), "value".into()); - assert_ne!(fingerprint(&original), fingerprint(&changed)); + assert_ne!( + fingerprint(&original).unwrap(), + fingerprint(&changed).unwrap() + ); changed = original.clone(); changed.resources.push(changed.base_dir.join("new.txt")); - assert_ne!(fingerprint(&original), fingerprint(&changed)); + assert_ne!( + fingerprint(&original).unwrap(), + fingerprint(&changed).unwrap() + ); } #[test] fn baseline_and_unchanged_catalog_do_not_notify() { let skills = vec![skill("existing", "body")]; - let mut monitor = SkillCatalogMonitor::new(&skills); + let mut monitor = SkillCatalogMonitor::new(&skills).unwrap(); monitor .submit( &skills, @@ -235,7 +269,7 @@ mod tests { fn reports_added_changed_and_removed_names_in_the_user_batch() { let before = vec![skill("removed", "old"), skill("changed", "old")]; let after = vec![skill("added", "new"), skill("changed", "new")]; - let mut monitor = SkillCatalogMonitor::new(&before); + let mut monitor = SkillCatalogMonitor::new(&before).unwrap(); monitor .submit(&after, vec![Item::text(ItemKind::User, "hello")], |items| { assert_eq!(items.len(), 2); @@ -257,7 +291,7 @@ mod tests { .rev() .map(|index| skill(&format!("skill-{index:03}"), "body")) .collect::>(); - let mut monitor = SkillCatalogMonitor::new(&[]); + let mut monitor = SkillCatalogMonitor::new(&[]).unwrap(); monitor .submit(&after, vec![Item::text(ItemKind::User, "hello")], |items| { let text = notification_text(&items).unwrap(); @@ -272,7 +306,7 @@ mod tests { #[test] fn failed_submission_does_not_advance_baseline() { let after = vec![skill("added", "body")]; - let mut monitor = SkillCatalogMonitor::new(&[]); + let mut monitor = SkillCatalogMonitor::new(&[]).unwrap(); assert!( monitor .submit(&after, vec![Item::text(ItemKind::User, "first")], |_| Err( diff --git a/src/protocols/acp/v2.rs b/src/protocols/acp/v2.rs index f5dbc27..9dd485f 100644 --- a/src/protocols/acp/v2.rs +++ b/src/protocols/acp/v2.rs @@ -675,7 +675,8 @@ impl Server { let catalog = model_catalog(¤t).await; let config_options = v2_config_options(¤t, reasoning, &catalog); let canonical_transcript = driver.canonical_transcript; - let skill_catalog = skill_catalog::SkillCatalogMonitor::new(&driver.skills); + let skill_catalog = skill_catalog::SkillCatalogMonitor::new(&driver.skills) + .map_err(|error| AcpRuntimeError::Loop(format!("skill catalog error: {error}")))?; let background_jobs = driver.background_jobs.clone(); let tasks = driver.tasks.clone(); let structured_completion = driver.structured_completion; @@ -1103,7 +1104,12 @@ async fn prepare_prompt( let prepared = integration.prompt_to_items(&request).and_then(|items| { skill_catalog .submit(skills, items, |items| driver.submit_input(items)) - .map_err(|error| map_loop_error(session_id, &error))?; + .map_err(|error| match error { + skill_catalog::SubmitError::Catalog(error) => { + AcpRuntimeError::Loop(format!("skill catalog error: {error}")) + } + skill_catalog::SubmitError::Submit(error) => map_loop_error(session_id, &error), + })?; integration.begin_prompt(session_id) }); let user_message_id = match prepared { @@ -2696,7 +2702,7 @@ mod tests { let task_manager = AsyncTaskManager::new(); let tasks = task_manager.handle(); let background_jobs = BackgroundJobs::default(); - let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]); + let mut skill_catalog = skill_catalog::SkillCatalogMonitor::new(&[]).unwrap(); let (result, ()) = tokio::join!( prepare_prompt( &session_id, From 91b3648ffe069b3a1d92282852c93a2c12499e09 Mon Sep 17 00:00:00 2001 From: daniel Date: Tue, 1 Sep 2026 11:41:18 +0100 Subject: [PATCH 3/4] fix(skills): include activation details in reload notices --- src/protocols/acp/skill_catalog.rs | 121 ++++++++++++++++++++++------- src/runtime.rs | 19 ++++- src/runtime/tests.rs | 22 +++++- 3 files changed, 128 insertions(+), 34 deletions(-) diff --git a/src/protocols/acp/skill_catalog.rs b/src/protocols/acp/skill_catalog.rs index 0f3b195..3e823e8 100644 --- a/src/protocols/acp/skill_catalog.rs +++ b/src/protocols/acp/skill_catalog.rs @@ -48,7 +48,7 @@ impl SkillCatalogMonitor { submit: impl FnOnce(Vec) -> Result<(), E>, ) -> Result<(), SubmitError> { let current = catalog(skills).map_err(SubmitError::Catalog)?; - if let Some(notification) = notification(&self.baseline, ¤t) { + if let Some(notification) = notification(&self.baseline, ¤t, skills) { user_items.insert(0, Item::notification(notification)); } submit(user_items).map_err(SubmitError::Submit)?; @@ -91,20 +91,32 @@ fn fingerprint(skill: &Skill) -> serde_json::Result { Ok(fingerprint.finalize()) } -fn notification(previous: &Catalog, current: &Catalog) -> Option { +fn notification(previous: &Catalog, current: &Catalog, skills: &[Skill]) -> Option { if previous == current { return None; } + let descriptions = skills + .iter() + .map(|skill| (skill.name.as_str(), skill.description.as_str())) + .collect::>(); let added = current .keys() .filter(|name| !previous.contains_key(*name)) - .map(String::as_str) + .filter_map(|name| { + descriptions + .get(name.as_str()) + .map(|description| (name.as_str(), *description)) + }) .collect::>(); let changed = current .iter() .filter(|(name, fingerprint)| previous.get(*name).is_some_and(|old| old != *fingerprint)) - .map(|(name, _)| name.as_str()) + .filter_map(|(name, _)| { + descriptions + .get(name.as_str()) + .map(|description| (name.as_str(), *description)) + }) .collect::>(); let removed = previous .keys() @@ -112,58 +124,101 @@ fn notification(previous: &Catalog, current: &Catalog) -> Option { .map(String::as_str) .collect::>(); let total = added.len() + changed.len() + removed.len(); - let omission_reserve = format!("; {total} more change(s) omitted").len(); + let omission_reserve = format!("\n{total} more change(s) omitted").len(); let content_limit = MAX_NOTIFICATION_BYTES.saturating_sub(omission_reserve); let mut message = String::from("Skill catalog update (informational)."); let mut reported = 0; - append_group(&mut message, "Added", &added, &mut reported, content_limit); - append_group( + append_skill_group( &mut message, - "Changed", - &changed, + "Added skills", + &added, &mut reported, content_limit, ); - append_group( + append_skill_group( &mut message, - "Removed", - &removed, + "Changed skills", + &changed, &mut reported, content_limit, ); + append_removed_group(&mut message, &removed, &mut reported, content_limit); if reported < total { - message.push_str(&format!("; {} more change(s) omitted", total - reported)); + message.push_str(&format!("\n{} more change(s) omitted", total - reported)); } Some(message) } -fn append_group( +fn append_skill_group( message: &mut String, label: &str, - names: &[&str], + skills: &[(&str, &str)], reported: &mut usize, content_limit: usize, ) { - if names.is_empty() || *reported >= MAX_REPORTED_CHANGES { + if skills.is_empty() || *reported >= MAX_REPORTED_CHANGES { return; } - let heading = format!("; {label}: "); + let heading = format!("\n{label}:"); if message.len() + heading.len() > content_limit { return; } let heading_start = message.len(); message.push_str(&heading); let mut group_count = 0; - for name in names { - if *reported >= MAX_REPORTED_CHANGES { - break; + for (name, description) in skills + .iter() + .take(MAX_REPORTED_CHANGES.saturating_sub(*reported)) + { + const ROW_OVERHEAD: usize = "\n- name: \"\"\n description: \"\"".len(); + if message.len() + ROW_OVERHEAD + name.len() + description.len() > content_limit { + continue; } - let separator = if group_count == 0 { "" } else { ", " }; - if message.len() + separator.len() + name.len() > content_limit { - break; + let row = format!( + "\n- name: {}\n description: {}", + serde_json::Value::String((*name).to_owned()), + serde_json::Value::String((*description).to_owned()) + ); + if message.len() + row.len() > content_limit { + continue; + } + message.push_str(&row); + group_count += 1; + *reported += 1; + } + if group_count == 0 { + message.truncate(heading_start); + } +} + +fn append_removed_group( + message: &mut String, + names: &[&str], + reported: &mut usize, + content_limit: usize, +) { + if names.is_empty() || *reported >= MAX_REPORTED_CHANGES { + return; + } + const HEADING: &str = "\nRemoved skills:"; + if message.len() + HEADING.len() > content_limit { + return; + } + let heading_start = message.len(); + message.push_str(HEADING); + let mut group_count = 0; + for name in names + .iter() + .take(MAX_REPORTED_CHANGES.saturating_sub(*reported)) + { + let row = format!( + "\n- name: {}", + serde_json::Value::String((*name).to_owned()) + ); + if message.len() + row.len() > content_limit { + continue; } - message.push_str(separator); - message.push_str(name); + message.push_str(&row); group_count += 1; *reported += 1; } @@ -277,9 +332,13 @@ mod tests { assert_eq!(items[1].kind, ItemKind::User); let text = notification_text(&items).unwrap(); assert!(text.starts_with("Skill catalog update (informational).")); - assert!(text.contains("Added: added")); - assert!(text.contains("Changed: changed")); - assert!(text.contains("Removed: removed")); + assert!(text.contains( + "Added skills:\n- name: \"added\"\n description: \"Description for added\"" + )); + assert!(text.contains( + "Changed skills:\n- name: \"changed\"\n description: \"Description for changed\"" + )); + assert!(text.contains("Removed skills:\n- name: \"removed\"")); Ok::<_, ()>(()) }) .unwrap(); @@ -316,7 +375,11 @@ mod tests { ); monitor .submit(&after, vec![Item::text(ItemKind::User, "retry")], |items| { - assert!(notification_text(&items).unwrap().contains("Added: added")); + assert!( + notification_text(&items) + .unwrap() + .contains("Added skills:\n- name: \"added\"") + ); Ok::<_, ()>(()) }) .unwrap(); diff --git a/src/runtime.rs b/src/runtime.rs index 71892a8..4490b3a 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -785,11 +785,17 @@ impl Runtime { .register(Observed::new(ToolSearch::new(self.mcp.clone()))) .register(Observed::new(AuthTool::new(self.mcp.clone()))) .register(Observed::new(McpTool::new(self.mcp.clone()))); + let mut child_specs = children.specs(); let skill_tools = skills.tool_registry(); if let Some(skill_tool) = skill_tools.get(&ToolName::new("skill")) { + let frozen_skill_spec = skill_spec_with_open_name( + skill_tool + .current_spec() + .unwrap_or_else(|| skill_tool.spec().clone()), + ); children.register(observe_shared(skill_tool)); + child_specs.push(frozen_skill_spec); } - let child_specs = children.specs(); let compose = ComposeTool::wrap(children) .with_source(self.mcp.catalog().unadvertised()) .with_config( @@ -1763,6 +1769,17 @@ fn background_route(request: &ToolRequest) -> RoutingDecision { } } +fn skill_spec_with_open_name(mut spec: ToolSpec) -> ToolSpec { + if let Some(name_schema) = spec + .input_schema + .pointer_mut("/properties/name") + .and_then(Value::as_object_mut) + { + name_schema.remove("enum"); + } + spec +} + struct HiddenRunletBackend(Vec); #[async_trait] diff --git a/src/runtime/tests.rs b/src/runtime/tests.rs index 15dd79b..0ffa73e 100644 --- a/src/runtime/tests.rs +++ b/src/runtime/tests.rs @@ -604,7 +604,7 @@ fn project_skills_take_precedence_over_plugin_skills() { } #[tokio::test] -async fn compose_can_load_a_skill_repeatedly() { +async fn compose_can_load_a_skill_added_after_its_spec_is_frozen() { let root = tempfile::tempdir().unwrap(); write_skill( &root.path().join(".agents/skills/reusable"), @@ -614,6 +614,15 @@ async fn compose_can_load_a_skill_repeatedly() { ); let runtime = Runtime::new(root.path(), "gpt-5.4").unwrap(); let compose = runtime.compose(0); + let frozen_description = &compose.specs()[0].description; + assert!(frozen_description.contains("- name: reusable")); + assert!(!frozen_description.contains(r#""enum":["reusable"]"#)); + write_skill( + &root.path().join(".agents/skills/added-later"), + "added-later", + "Added after the Compose schema was frozen.", + "new instructions", + ); let source: Arc = Arc::new(compose.compose.clone()); let executor: Arc = Arc::new(BasicToolExecutor::new([source])); let permissions = Arc::new(AllowAllPermissions); @@ -644,7 +653,7 @@ async fn compose_can_load_a_skill_repeatedly() { ToolCallId::new("call"), ToolName::new("compose"), json!({ - "script": "first = skill({ name: \"reusable\" })\nsecond = skill({ name: \"reusable\" })\nreturn [first, second]" + "script": "first = skill({ name: \"reusable\" })\nsecond = skill({ name: \"reusable\" })\nthird = skill({ name: \"added-later\" })\nreturn [first, second, third]" }), session_id, turn_id, @@ -660,12 +669,17 @@ async fn compose_can_load_a_skill_repeatedly() { panic!("compose did not return structured output"); }; let loaded = loaded.as_array().expect("compose returned an array"); - assert_eq!(loaded.len(), 2); - assert!(loaded.iter().all(|skill| { + assert_eq!(loaded.len(), 3); + assert!(loaded[..2].iter().all(|skill| { skill .as_str() .is_some_and(|text| text.contains("full instructions")) })); + assert!( + loaded[2] + .as_str() + .is_some_and(|text| text.contains("new instructions")) + ); } #[test] From fe86613e8f9c137b7829a012f59b84fd65a796e3 Mon Sep 17 00:00:00 2001 From: daniel Date: Tue, 1 Sep 2026 11:56:37 +0100 Subject: [PATCH 4/4] chore: bump version to 0.1.121 --- Cargo.lock | 2 +- Cargo.toml | 2 +- macos/Config/Version.xcconfig | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 130e283..fb42aea 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2554,7 +2554,7 @@ dependencies = [ [[package]] name = "kit" -version = "0.1.120" +version = "0.1.121" dependencies = [ "a2a-protocol-client", "a2a-protocol-server", diff --git a/Cargo.toml b/Cargo.toml index 0651894..d6c4d7d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "kit" -version = "0.1.120" +version = "0.1.121" edition = "2024" rust-version = "1.94.0" publish = false diff --git a/macos/Config/Version.xcconfig b/macos/Config/Version.xcconfig index 5049afe..5ff06b2 100644 --- a/macos/Config/Version.xcconfig +++ b/macos/Config/Version.xcconfig @@ -1,2 +1,2 @@ // Generated by scripts/generate-macos-project.sh from Cargo.toml. -KIT_VERSION = 0.1.120 +KIT_VERSION = 0.1.121