use std::collections::BTreeMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use meerkat_core::types::{ToolCallView, ToolDef};
use meerkat_core::{
AgentToolDispatcher, EphemeralToolBindingFingerprint, ResolvedToolExecutionPlan,
ToolCatalogCapabilities, ToolCatalogEntry, ToolDispatchContext, ToolDispatchOutcome, ToolError,
ToolExecutionOwnerWitness, ToolExecutionResolutionContext, ToolExecutionResolutionError,
ToolUnavailableReason,
};
use meerkat_mob::{
AgentIdentity as MobAgentIdentity, MobError, SpawnCustomizationContext, SpawnMemberCustomizer,
SpawnMemberSpec,
};
use super::types::AgentIdentity;
#[derive(Clone)]
struct Published {
generation: u64,
dispatcher: Option<Arc<dyn AgentToolDispatcher>>,
}
pub struct IdentityCustomizerTools {
member_id: MobAgentIdentity,
current: RwLock<Published>,
execution_authority_id: u64,
}
static NEXT_EXECUTION_AUTHORITY_ID: AtomicU64 = AtomicU64::new(1);
impl IdentityCustomizerTools {
fn new(member_id: MobAgentIdentity) -> Self {
Self {
member_id,
current: RwLock::new(Published {
generation: 0,
dispatcher: None,
}),
execution_authority_id: NEXT_EXECUTION_AUTHORITY_ID.fetch_add(1, Ordering::Relaxed),
}
}
pub fn member_id(&self) -> &MobAgentIdentity {
&self.member_id
}
fn snapshot(&self) -> Published {
self.current
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub fn publish(&self, dispatcher: Option<Arc<dyn AgentToolDispatcher>>) -> u64 {
let mut current = self
.current
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
current.generation = current.generation.saturating_add(1);
current.dispatcher = dispatcher;
current.generation
}
pub fn generation(&self) -> u64 {
self.snapshot().generation
}
fn owner_for(published: &Published, tool_name: &str) -> Option<Arc<dyn AgentToolDispatcher>> {
let dispatcher = published.dispatcher.as_ref()?;
dispatcher
.tools()
.iter()
.any(|tool| tool.name.as_ref() == tool_name)
.then(|| Arc::clone(dispatcher))
}
fn binding_epoch_in(published: &Published, tool_name: &str) -> u64 {
let inner = Self::owner_for(published, tool_name)
.map_or(0, |owner| owner.execution_binding_epoch(tool_name));
(published.generation << 32) | (inner & 0xFFFF_FFFF)
}
fn binding_fingerprint_in(
published: &Published,
tool_name: &str,
) -> Result<EphemeralToolBindingFingerprint, ToolExecutionResolutionError> {
let not_found = || ToolExecutionResolutionError::NotFound {
tool_name: tool_name.to_string(),
};
let owner = Self::owner_for(published, tool_name).ok_or_else(not_found)?;
let catalog = owner.tool_catalog();
let entry = catalog
.iter()
.find(|entry| entry.tool.name == tool_name)
.ok_or_else(not_found)?;
let child = owner.execution_binding_fingerprint(tool_name)?;
Ok(
meerkat_core::ephemeral_tool_catalog_binding_fingerprint(entry)
.with_live_authority(0, Self::binding_epoch_in(published, tool_name))
.with_dependency(&child),
)
}
fn execution_authority_key(&self) -> String {
format!("identity-customizer-tools:{}", self.execution_authority_id)
}
fn not_advertised(&self, name: &str) -> ToolError {
tracing::debug!(
member_id = %self.member_id,
tool = name,
"customizer tool call refused: the current customize_build publication does not \
advertise it"
);
ToolError::NotFound {
name: name.to_string(),
}
}
}
#[async_trait::async_trait]
impl AgentToolDispatcher for IdentityCustomizerTools {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
match self.snapshot().dispatcher {
Some(dispatcher) => dispatcher.tools(),
None => Arc::from([]),
}
}
fn tool_catalog_capabilities(&self) -> ToolCatalogCapabilities {
match self.snapshot().dispatcher {
Some(dispatcher) => dispatcher.tool_catalog_capabilities(),
None => ToolCatalogCapabilities {
exact_catalog: true,
may_require_catalog_control_plane: false,
},
}
}
fn tool_catalog(&self) -> Arc<[ToolCatalogEntry]> {
match self.snapshot().dispatcher {
Some(dispatcher) => dispatcher.tool_catalog(),
None => Arc::from([]),
}
}
fn tool_mutation_class(&self, tool_name: &str) -> meerkat_core::ToolMutationClass {
match Self::owner_for(&self.snapshot(), tool_name) {
Some(owner) => owner.tool_mutation_class(tool_name),
None => meerkat_core::ToolMutationClass::Unknown,
}
}
fn live_bridge_effect_kind(&self, tool_name: &str) -> meerkat_core::LiveBridgeEffectKind {
match Self::owner_for(&self.snapshot(), tool_name) {
Some(owner) => owner.live_bridge_effect_kind(tool_name),
None => meerkat_core::LiveBridgeEffectKind::ExternalIo,
}
}
fn execution_binding_epoch(&self, tool_name: &str) -> u64 {
Self::binding_epoch_in(&self.snapshot(), tool_name)
}
fn execution_binding_fingerprint(
&self,
tool_name: &str,
) -> Result<EphemeralToolBindingFingerprint, ToolExecutionResolutionError> {
Self::binding_fingerprint_in(&self.snapshot(), tool_name)
}
fn resolve_execution_plan(
&self,
call: ToolCallView<'_>,
dispatch_context: &ToolDispatchContext,
resolution_context: &ToolExecutionResolutionContext,
) -> Result<ResolvedToolExecutionPlan, ToolExecutionResolutionError> {
let published = self.snapshot();
let Some(owner) = Self::owner_for(&published, call.name) else {
let _ = self.not_advertised(call.name);
return Err(ToolExecutionResolutionError::NotFound {
tool_name: call.name.to_string(),
});
};
let before = Self::binding_fingerprint_in(&published, call.name)?;
let plan = owner.resolve_execution_plan(call, dispatch_context, resolution_context)?;
let current = self.snapshot();
if current.generation != published.generation
|| Self::binding_fingerprint_in(¤t, call.name).as_ref() != Ok(&before)
{
return Err(ToolExecutionResolutionError::Unavailable {
tool_name: call.name.to_string(),
reason: ToolUnavailableReason::ExecutionOwnerChanged,
});
}
let witness = ToolExecutionOwnerWitness::new(
self.execution_authority_key(),
published.generation.to_string(),
before,
)?;
plan.with_owner_witness(witness)
}
fn validate_resolved_execution_plan(
&self,
call: ToolCallView<'_>,
resolution_context: &ToolExecutionResolutionContext,
plan: &ResolvedToolExecutionPlan,
) -> Result<(), ToolExecutionResolutionError> {
match Self::owner_for(&self.snapshot(), call.name) {
Some(owner) => owner.validate_resolved_execution_plan(call, resolution_context, plan),
None => Err(ToolExecutionResolutionError::NotFound {
tool_name: call.name.to_string(),
}),
}
}
fn pending_catalog_sources(&self) -> Arc<[String]> {
match self.snapshot().dispatcher {
Some(dispatcher) => dispatcher.pending_catalog_sources(),
None => Arc::from([]),
}
}
async fn dispatch(&self, call: ToolCallView<'_>) -> Result<ToolDispatchOutcome, ToolError> {
match Self::owner_for(&self.snapshot(), call.name) {
Some(owner) => owner.dispatch(call).await,
None => Err(self.not_advertised(call.name)),
}
}
async fn dispatch_with_context(
&self,
call: ToolCallView<'_>,
context: &ToolDispatchContext,
) -> Result<ToolDispatchOutcome, ToolError> {
match Self::owner_for(&self.snapshot(), call.name) {
Some(owner) => owner.dispatch_with_context(call, context).await,
None => Err(self.not_advertised(call.name)),
}
}
async fn dispatch_resolved_with_context(
&self,
call: ToolCallView<'_>,
context: &ToolDispatchContext,
plan: &ResolvedToolExecutionPlan,
) -> Result<ToolDispatchOutcome, ToolError> {
let published = self.snapshot();
let Some(owner) = Self::owner_for(&published, call.name) else {
return Err(self.not_advertised(call.name));
};
let owner_changed =
|| ToolError::unavailable(call.name, ToolUnavailableReason::ExecutionOwnerChanged);
let witness = plan
.owner_witness(&self.execution_authority_key())
.ok_or_else(owner_changed)?;
let current =
Self::binding_fingerprint_in(&published, call.name).map_err(|_| owner_changed())?;
if witness.owner_key() != published.generation.to_string()
|| witness.binding_fingerprint() != ¤t
{
return Err(owner_changed());
}
owner
.dispatch_resolved_with_context(call, context, plan)
.await
}
}
#[derive(Default)]
pub struct CustomizerToolRegistry {
entries: RwLock<BTreeMap<MobAgentIdentity, Arc<IdentityCustomizerTools>>>,
}
impl CustomizerToolRegistry {
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
}
fn member_id(identity: &AgentIdentity) -> MobAgentIdentity {
crate::member_comms_id::mob_member_id(identity.as_str())
}
pub fn ensure(&self, identity: &AgentIdentity) -> Arc<IdentityCustomizerTools> {
let member_id = Self::member_id(identity);
if let Some(existing) = self
.entries
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&member_id)
{
return Arc::clone(existing);
}
let mut entries = self
.entries
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Arc::clone(
entries
.entry(member_id.clone())
.or_insert_with(|| Arc::new(IdentityCustomizerTools::new(member_id))),
)
}
pub fn ensure_member(&self, member_id: &MobAgentIdentity) -> Arc<IdentityCustomizerTools> {
if let Some(existing) = self.for_member(member_id) {
return existing;
}
let mut entries = self
.entries
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Arc::clone(
entries
.entry(member_id.clone())
.or_insert_with(|| Arc::new(IdentityCustomizerTools::new(member_id.clone()))),
)
}
pub fn for_member(&self, member_id: &MobAgentIdentity) -> Option<Arc<IdentityCustomizerTools>> {
self.entries
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(member_id)
.cloned()
}
pub fn publish(
&self,
identity: &AgentIdentity,
dispatcher: Option<Arc<dyn AgentToolDispatcher>>,
) -> Arc<IdentityCustomizerTools> {
let entry = self.ensure(identity);
entry.publish(dispatcher);
entry
}
}
pub async fn prepublish(
registry: &CustomizerToolRegistry,
roster: &[super::types::DurableAgentSpec],
customizer: &dyn super::contracts::AgentCustomizer,
runtime_services: super::types::AgentRuntimeServices,
active_peers: &[AgentIdentity],
managed_edges: &[super::types::ManagedPeerEdge],
) -> BTreeMap<AgentIdentity, super::types::CustomizerToolsPending> {
let mut pending = BTreeMap::new();
for spec in roster {
let build_context = super::types::AgentBuildContext {
identity: spec.identity.clone(),
active_peers: active_peers.to_vec(),
managed_edges: managed_edges.to_vec(),
runtime_services: runtime_services.clone(),
};
let mut draft = super::types::AgentBuildDraft {
model: None,
system_prompt: None,
additional_instructions: spec.additional_instructions.clone(),
labels: spec.labels.clone(),
app_context: spec.context.clone(),
external_tools: Vec::new(),
local_external_tools: Default::default(),
provider_params: None,
compaction_curator: Default::default(),
};
match customizer
.customize_build(&build_context, spec, &mut draft)
.await
{
Ok(()) => {
registry.publish(&spec.identity, draft.local_external_tools.dispatcher());
}
Err(error) => {
let reason = format!("pre-activation customize_build failed: {error}");
tracing::warn!(
identity = %spec.identity,
%reason,
"restored member's customizer tools are not published; it advertises none \
until its materialization publishes them"
);
registry.ensure(&spec.identity);
pending.insert(
spec.identity.clone(),
super::types::CustomizerToolsPending { reason },
);
}
}
}
pending
}
fn same_dispatcher(a: &Arc<dyn AgentToolDispatcher>, b: &Arc<IdentityCustomizerTools>) -> bool {
std::ptr::eq(Arc::as_ptr(a).cast::<()>(), Arc::as_ptr(b).cast::<()>())
}
pub struct CustomizerToolsSpawnCustomizer {
registry: Arc<CustomizerToolRegistry>,
}
impl CustomizerToolsSpawnCustomizer {
pub fn new(registry: Arc<CustomizerToolRegistry>) -> Self {
Self { registry }
}
}
impl SpawnMemberCustomizer for CustomizerToolsSpawnCustomizer {
fn customize_spawn(
&self,
_ctx: &SpawnCustomizationContext,
spec: &mut SpawnMemberSpec,
) -> Result<(), MobError> {
self.apply(spec);
Ok(())
}
}
impl CustomizerToolsSpawnCustomizer {
fn apply(&self, spec: &mut SpawnMemberSpec) {
let entry = match self.registry.for_member(&spec.identity) {
Some(entry) => entry,
None if spec
.labels
.as_ref()
.and_then(crate::member_comms_id::durable_identity_label)
.is_some() =>
{
self.registry.ensure_member(&spec.identity)
}
None => return,
};
spec.external_tools = Some(match spec.external_tools.take() {
None => entry,
Some(existing) if same_dispatcher(&existing, &entry) => existing,
Some(existing) => {
crate::tool_compose::ComposedExternalTools::over(entry, Some(existing))
}
});
}
}
pub struct ComposedSpawnMemberCustomizer {
customizers: Vec<Arc<dyn SpawnMemberCustomizer>>,
}
impl ComposedSpawnMemberCustomizer {
pub fn over(
existing: Option<Arc<dyn SpawnMemberCustomizer>>,
added: Arc<dyn SpawnMemberCustomizer>,
) -> Arc<dyn SpawnMemberCustomizer> {
match existing {
None => added,
Some(existing) => Arc::new(Self {
customizers: vec![existing, added],
}),
}
}
}
impl SpawnMemberCustomizer for ComposedSpawnMemberCustomizer {
fn customize_spawn(
&self,
ctx: &SpawnCustomizationContext,
spec: &mut SpawnMemberSpec,
) -> Result<(), MobError> {
for customizer in &self.customizers {
customizer.customize_spawn(ctx, spec)?;
}
Ok(())
}
}
#[cfg(test)]
#[path = "customizer_tools_tests.rs"]
mod tests;