meerkat-mobkit 0.8.49

Companion orchestration platform for the Meerkat multi-agent runtime
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
//! Customizer tools that follow an identity across meerkat-side rebuilds (#563).
//!
//! The tools a host returns from `AgentCustomizer::customize_build` used to
//! reach a member only as the per-spawn `external_tools` overlay of the one
//! spawn that ran the customizer. meerkat-mob keeps that overlay in memory
//! and never persists it, so a restart restore, an adopted occupant
//! (`MemberAlreadyExists`), a respawn or a delivery-time repair rebuilt the
//! member without them: the host's handlers stayed registered, but the live
//! agent no longer advertised the tools.
//!
//! Instead, every identity member carries ONE stable, dynamic dispatcher
//! ([`IdentityCustomizerTools`]) for its whole life. Each successful
//! `customize_build` publishes its result into that dispatcher, swapping the
//! tool list and the dispatcher that serves it (for gateway hosts, the
//! handler scope) together. meerkat composes per-spawn overlays through a
//! dynamic composite that re-reads `tools()` on every turn, so:
//!
//! - an occupant that is adopted instead of respawned already holds the
//!   identity's dispatcher, and publishing re-attaches the current tools
//!   live, without a respawn;
//! - [`CustomizerToolsSpawnCustomizer`], installed on the mob whenever an
//!   agent customizer exists, attaches the same dispatcher to every
//!   meerkat-side build of a registered identity member (restart restore,
//!   explicit resume, respawn, delivery-time repair). It is synchronous and
//!   makes no host call: it attaches the stable dispatcher rather than
//!   re-running the host.
//!
//! Members that are not roster identities (helpers, forks, flow-provisioned
//! or raw-spawned members under their own ids) are never registered and get
//! no customizer tools: `customize_build` is a contract over roster
//! identities (it takes a `DurableAgentSpec`). A raw spawn that targets a
//! roster identity's member id does get that identity's current tools.

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;

/// One `customize_build` result: the dispatcher serving the published tools
/// (`None` when the build declared none) and the publication generation.
#[derive(Clone)]
struct Published {
    generation: u64,
    dispatcher: Option<Arc<dyn AgentToolDispatcher>>,
}

/// The stable, dynamic customizer-tool dispatcher of one identity member.
///
/// Every read takes one snapshot of the current publication, so the tools a
/// turn sees and the dispatcher that serves a call come from the same
/// `customize_build` result. A call to a tool the current publication does
/// not advertise is refused typed (`ToolError::NotFound`), never routed to a
/// dispatcher (or handler scope) from an earlier publication.
///
/// Execution plans follow the same rule. A plan is resolved by the dispatcher
/// that serves the tool (so its own owner fencing and its argument-sensitive
/// mode and deadlines are kept), then witnessed against this publication; a
/// later publication, even with identical tool metadata, refuses that plan
/// (`ExecutionOwnerChanged`) and the caller resolves again.
pub struct IdentityCustomizerTools {
    member_id: MobAgentIdentity,
    current: RwLock<Published>,
    /// Distinguishes this dispatcher's plan witnesses from every other
    /// execution authority in the process.
    execution_authority_id: u64,
}

/// Process-wide source of execution authority ids for identity dispatchers.
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),
        }
    }

    /// The roster member id this dispatcher belongs to.
    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()
    }

    /// Publish one `customize_build` result. `None` clears the tools (the
    /// build declared none). Returns the new publication generation.
    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
    }

    /// How many times a result has been published (0: none yet).
    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))
    }

    /// The binding epoch of `tool_name` in `published`: the publication
    /// generation over the serving dispatcher's own epoch.
    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)
    }

    /// The binding fingerprint of `tool_name` in `published`: the serving
    /// dispatcher's catalog entry under this publication's epoch, with the
    /// serving dispatcher's own fingerprint as its dependency.
    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(),
            // An empty publication is its own exact (empty) registry.
            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,
        }
    }

    /// A mutable authority: every publication advances the epoch, including
    /// an identical-metadata replacement (a new handler scope behind the same
    /// tool names), so a binding fingerprinted before it is never mistaken
    /// for the current one.
    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)
    }

    /// Resolve through the dispatcher that serves the tool, then witness the
    /// publication it was resolved against. The serving dispatcher's plan is
    /// kept whole: its own owner witnesses, its mode and its deadlines.
    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)?;
        // A publication (or an owner rebinding) during resolution: the plan
        // may belong to neither binding.
        let current = self.snapshot();
        if current.generation != published.generation
            || Self::binding_fingerprint_in(&current, 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)
    }

    /// Validate against the dispatcher that serves the tool now.
    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));
        };
        // Only a plan resolved against THIS publication: an earlier one, even
        // with identical tool metadata, may name another handler scope.
        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() != &current
        {
            return Err(owner_changed());
        }
        owner
            .dispatch_resolved_with_context(call, context, plan)
            .await
    }
}

/// The customizer-tool dispatchers of every registered identity member,
/// keyed by roster member id (the deterministic `mob_member_id(identity)`).
#[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())
    }

    /// Register `identity` (idempotent) and return its stable dispatcher.
    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))),
        )
    }

    /// Register a roster member id directly (idempotent).
    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()))),
        )
    }

    /// The registered dispatcher for a roster member id, if any.
    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()
    }

    /// Publish one `customize_build` result for `identity` (registering it
    /// if needed) and return the identity's stable dispatcher, which is what
    /// every spawn of that identity attaches.
    pub fn publish(
        &self,
        identity: &AgentIdentity,
        dispatcher: Option<Arc<dyn AgentToolDispatcher>>,
    ) -> Arc<IdentityCustomizerTools> {
        let entry = self.ensure(identity);
        entry.publish(dispatcher);
        entry
    }
}

/// Run the host customizer for each roster identity and publish its tools
/// into `registry`, BEFORE meerkat builds any member (#563). Returns the
/// identities whose `customize_build` failed, with the typed reason; they stay
/// registered with nothing published until their materialization publishes.
/// Each failure is logged per identity.
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::<()>())
}

/// Attaches a registered identity member's customizer-tool dispatcher to
/// every meerkat-side build of it: restart restore and explicit resume
/// (`SpawnSource::Resume`), respawn, delivery-time repair, and raw spawns
/// that target the identity's member id. Members that are not registered
/// identities are left untouched.
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 {
    /// The customizer body, factored off the trait so tests can drive it:
    /// `SpawnCustomizationContext` is `#[non_exhaustive]` and only
    /// meerkat-mob constructs it. The attachment does not depend on the spawn
    /// source: every build of a registered identity member gets it.
    fn apply(&self, spec: &mut SpawnMemberSpec) {
        // A member MobKit built for an identity carries the
        // runtime-authoritative `agent_identity` label (raw member creation
        // may not supply it), and the roster entry keeps it. That marks a
        // restored identity member even when meerkat restores the mob before
        // MobKit has seen the roster, so it is registered here, at its first
        // build, and receives the identity's tools when they are published.
        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,
            // MobKit's own identity spawn already attached it.
            Some(existing) if same_dispatcher(&existing, &entry) => existing,
            // A raw spawn of the identity's member id that brought its own
            // overlay: the customizer tools win name collisions, the rest
            // stays reachable.
            Some(existing) => {
                crate::tool_compose::ComposedExternalTools::over(entry, Some(existing))
            }
        });
    }
}

/// Runs several spawn customizers in order. meerkat-mob has a single
/// `SpawnMemberCustomizer` slot; composing keeps a later installer from
/// silently replacing an earlier one.
pub struct ComposedSpawnMemberCustomizer {
    customizers: Vec<Arc<dyn SpawnMemberCustomizer>>,
}

impl ComposedSpawnMemberCustomizer {
    /// `existing` (if any) runs first, then `added`.
    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;