Skip to main content

meerkat_mobkit/identity_first/
runtime.rs

1//! Identity-first runtime: delivery, status, lifecycle, and ownership enforcement.
2//!
3//! This module implements the behavioral core of identity-first continuity:
4//! - Delivery: `send()` and `dispatch()` with addressability and lease enforcement
5//! - Status: `status()` returning `IdentityStatus`
6//! - Lifecycle: `retire()`, `respawn()`, `reset()`, `delete_identity()`
7//! - Ownership: lease tracking, fencing, and invariant enforcement
8
9use std::collections::{BTreeMap, BTreeSet};
10use std::sync::{Arc, RwLock as StdRwLock};
11use std::time::{Duration, Instant};
12
13use futures::stream::{self, StreamExt};
14use meerkat_core::types::{HandlingMode, SessionId};
15use tokio::sync::{Mutex, Notify, RwLock, broadcast};
16use tokio::task::JoinHandle;
17
18use super::agent_memory::{
19    AgentMemoryError, AgentMemoryForgetResult, AgentMemoryRecallRequest, AgentMemoryRecord,
20    AgentMemoryRuntimeInjector, NewAgentMemory,
21};
22use super::bridge::SessionBridge;
23use super::contracts::{
24    AgentCustomizer, ContinuityStore, LeaseProvider, RosterProvider, TopologyProvider,
25};
26use super::types::{
27    AgentAddressability, AgentBuildContext, AgentIdentity, AgentRuntimeId, AgentRuntimeServices,
28    CheckpointVersion, ContinuityGeneration, ContinuityHealth, ContinuityRecord,
29    ContinuityStoreError, DispatchInput, DurabilityPolicy, DurableAgentSpec, FencingToken,
30    IdentityLifecycleState, IdentityStatus, LeaseGrant, LeaseInfo, ManagedPeerEdge, NotAddressable,
31    RosterContext, SessionSnapshot, TopologyContext,
32};
33use crate::memory::records::{
34    ManifestTier, MemoryId, MemoryKind, MemoryScope, NewMemoryRecord, RecordMeta, UsageEvent,
35};
36
37const MANAGED_PEER_RECONCILE_CONCURRENCY: usize = 64;
38const MATERIALIZATION_FAILURE_BACKOFF: Duration = Duration::from_secs(30);
39fn durable_spec_uses_external_binding(spec: &DurableAgentSpec) -> bool {
40    matches!(spec.backend, Some(meerkat_mob::MobBackendKind::External))
41        || matches!(
42            spec.binding.as_ref(),
43            Some(meerkat_contracts::WireRuntimeBinding::External { .. })
44        )
45}
46
47// ---------------------------------------------------------------------------
48// Error types
49// ---------------------------------------------------------------------------
50
51/// Errors from identity-first runtime operations.
52#[derive(Debug)]
53pub enum IdentityRuntimeError {
54    /// Target identity is not registered/active.
55    UnknownIdentity(AgentIdentity),
56    /// send() rejected: target is InternalOnly.
57    NotAddressable(NotAddressable),
58    /// Operation rejected: no active lease for this identity.
59    NoActiveLease(AgentIdentity),
60    /// Fail-closed single-embodiment guard: the identity's durable lease is
61    /// held by another live runtime instance — a second live embodiment is
62    /// refused, loudly and with the holder named. This is a POLICY point,
63    /// not a structural assumption: multi-bind / forked-session semantics
64    /// (future work) relax exactly this arm.
65    AlreadyEmbodied {
66        identity: AgentIdentity,
67        holder: String,
68    },
69    /// Operation rejected: lease was lost.
70    LeaseLost(AgentIdentity),
71    /// Operation rejected: identity is not in a state that permits this operation.
72    InvalidState {
73        identity: AgentIdentity,
74        state: IdentityLifecycleState,
75        operation: &'static str,
76    },
77    /// Continuity store error.
78    Store(ContinuityStoreError),
79    /// Lease provider error.
80    Lease(super::types::LeaseError),
81    /// Duplicate identities in roster.
82    DuplicateIdentity(AgentIdentity),
83    /// Stale fencing token on checkpoint.
84    StaleFencingToken {
85        identity: AgentIdentity,
86        presented: FencingToken,
87        current: FencingToken,
88    },
89    /// Stale checkpoint version.
90    StaleCheckpointVersion {
91        identity: AgentIdentity,
92        presented: CheckpointVersion,
93        current: CheckpointVersion,
94    },
95    /// Generic I/O or internal error.
96    Internal(String),
97}
98
99impl std::fmt::Display for IdentityRuntimeError {
100    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
101        match self {
102            Self::UnknownIdentity(id) => write!(f, "unknown identity: {id}"),
103            Self::NotAddressable(err) => write!(f, "{err}"),
104            Self::NoActiveLease(id) => write!(f, "no active lease for {id}"),
105            Self::AlreadyEmbodied { identity, holder } => write!(
106                f,
107                "identity {identity} is already embodied by runtime instance '{holder}' \
108                 (single-embodiment guard: refusing a second live bind)"
109            ),
110            Self::LeaseLost(id) => write!(f, "lease lost for {id}"),
111            Self::InvalidState {
112                identity,
113                state,
114                operation,
115            } => write!(
116                f,
117                "cannot {operation} identity {identity} in state {state:?}"
118            ),
119            Self::Store(err) => write!(f, "continuity store: {err}"),
120            Self::Lease(err) => write!(f, "lease provider: {err}"),
121            Self::DuplicateIdentity(id) => write!(f, "duplicate identity in roster: {id}"),
122            Self::StaleFencingToken {
123                identity,
124                presented,
125                current,
126            } => write!(
127                f,
128                "stale fencing token for {identity}: presented {presented}, current {current}"
129            ),
130            Self::StaleCheckpointVersion {
131                identity,
132                presented,
133                current,
134            } => write!(
135                f,
136                "stale checkpoint version for {identity}: presented {presented}, current {current}"
137            ),
138            Self::Internal(msg) => write!(f, "internal: {msg}"),
139        }
140    }
141}
142
143impl std::error::Error for IdentityRuntimeError {}
144
145impl From<ContinuityStoreError> for IdentityRuntimeError {
146    fn from(err: ContinuityStoreError) -> Self {
147        match err {
148            ContinuityStoreError::StaleFencingToken {
149                identity,
150                presented,
151                current,
152            } => Self::StaleFencingToken {
153                identity,
154                presented,
155                current,
156            },
157            ContinuityStoreError::StaleCheckpointVersion {
158                identity,
159                presented,
160                current,
161            } => Self::StaleCheckpointVersion {
162                identity,
163                presented,
164                current,
165            },
166            other => Self::Store(other),
167        }
168    }
169}
170
171// ---------------------------------------------------------------------------
172// Per-identity runtime state
173// ---------------------------------------------------------------------------
174
175/// Tracks the live state for a single identity within the runtime.
176#[derive(Debug, Clone)]
177pub(crate) struct IdentityEntry {
178    pub spec: DurableAgentSpec,
179    pub state: IdentityLifecycleState,
180    pub continuity: Option<ContinuityRecord>,
181    pub lease: Option<LeaseEntry>,
182    pub checkpoint_version: CheckpointVersion,
183    /// Whether a durable runtime_store is available (affects dispatch ack semantics).
184    pub has_runtime_store: bool,
185}
186
187/// Tracks a held lease for an identity.
188#[derive(Debug, Clone)]
189pub(crate) struct LeaseEntry {
190    pub fencing_token: FencingToken,
191    pub ttl: Duration,
192    pub acquired_at: Instant,
193}
194
195impl LeaseEntry {
196    pub fn is_expired(&self) -> bool {
197        self.acquired_at.elapsed() > self.ttl
198    }
199
200    pub fn ttl_remaining(&self) -> Duration {
201        self.ttl.saturating_sub(self.acquired_at.elapsed())
202    }
203
204    pub fn is_healthy(&self) -> bool {
205        // Healthy if more than 20% TTL remains
206        let remaining = self.ttl_remaining();
207        remaining > self.ttl / 5
208    }
209}
210
211// ---------------------------------------------------------------------------
212// Identity-scoped events
213// ---------------------------------------------------------------------------
214
215/// Events emitted for a specific identity, used by `subscribe()`.
216#[derive(Debug, Clone)]
217pub enum IdentityEvent {
218    /// Lifecycle state changed.
219    StateChanged {
220        identity: AgentIdentity,
221        new_state: IdentityLifecycleState,
222    },
223    /// Lease acquired or renewed.
224    LeaseUpdated {
225        identity: AgentIdentity,
226        fencing_token: FencingToken,
227    },
228    /// Lease lost.
229    LeaseLost { identity: AgentIdentity },
230    /// Checkpoint completed.
231    CheckpointCompleted {
232        identity: AgentIdentity,
233        version: CheckpointVersion,
234    },
235    /// Resume could not reuse a persisted runtime binding and materialization
236    /// fresh-spawned a member instead.
237    ResumeFallback {
238        identity: AgentIdentity,
239        reason: super::bridge::ResumeFallbackReason,
240    },
241}
242
243/// Per-identity event channel capacity.
244const IDENTITY_EVENT_CHANNEL_CAPACITY: usize = 64;
245const DEFAULT_LEASE_RENEWAL_MAX_POLL_INTERVAL: Duration = Duration::from_mins(1);
246const DEFAULT_LEASE_RENEWAL_MIN_POLL_INTERVAL: Duration = Duration::from_millis(10);
247/// Base delay for the lease-renewal failure backoff. The renewal tick's
248/// normal cadence is TTL-derived (down to a 10ms floor), so a lease provider
249/// that errors persistently would otherwise retry — and warn — at that floor
250/// rate. On failure we back off from this base, doubling toward the max poll
251/// interval, so a backend outage can't spin the renewal task.
252const LEASE_RENEWAL_FAILURE_BACKOFF_BASE: Duration = Duration::from_secs(1);
253
254fn lease_renewal_failure_backoff(
255    consecutive_failures: u32,
256    max_poll_interval: Duration,
257) -> Duration {
258    LEASE_RENEWAL_FAILURE_BACKOFF_BASE
259        .saturating_mul(1u32 << consecutive_failures.min(6))
260        .min(max_poll_interval)
261}
262
263// ---------------------------------------------------------------------------
264// IdentityRuntime
265// ---------------------------------------------------------------------------
266
267/// Configuration for the identity-first runtime.
268pub struct IdentityRuntimeConfig {
269    pub continuity_store: Arc<dyn ContinuityStore>,
270    pub lease_provider: Arc<dyn LeaseProvider>,
271    pub runtime_instance_id: String,
272    pub has_runtime_store: bool,
273    pub durability_policy: DurabilityPolicy,
274    /// Optional session bridge for real session delivery. When `None`,
275    /// delivery operations validate invariants but do not forward to
276    /// the Meerkat session pipeline (useful for tests).
277    pub bridge: Option<Arc<dyn SessionBridge>>,
278    /// Default timeout for wait_for_output / wait_for_output_containing.
279    /// Defaults to 90 seconds if not set.
280    pub default_timeout: Option<Duration>,
281}
282
283#[derive(Clone)]
284pub struct IdentityFirstRuntimeContext {
285    pub runtime: Arc<IdentityRuntime>,
286    pub roster_provider: Arc<dyn RosterProvider>,
287    pub topology_provider: Option<Arc<dyn TopologyProvider>>,
288    pub customizer: Option<Arc<dyn AgentCustomizer>>,
289    mob_definition: Option<meerkat_mob::MobDefinition>,
290    lazy_materialization: bool,
291}
292
293impl IdentityFirstRuntimeContext {
294    pub(crate) async fn topology_snapshot_inputs(
295        &self,
296    ) -> Result<(Vec<DurableAgentSpec>, Vec<ManagedPeerEdge>), IdentityRuntimeError> {
297        let previous_identities = self.runtime.registered_identities().await;
298        let roster = self
299            .roster_provider
300            .roster(&RosterContext {
301                mob_definition: self.mob_definition.clone(),
302                previous_identities,
303            })
304            .await
305            .map_err(|error| IdentityRuntimeError::Internal(format!("roster provider: {error}")))?;
306        let identities = roster
307            .iter()
308            .map(|spec| spec.identity.clone())
309            .collect::<Vec<_>>();
310        let declared = match self.topology_provider.as_deref() {
311            Some(provider) => provider
312                .compute_edges(
313                    &identities,
314                    &TopologyContext {
315                        roster: roster.clone(),
316                    },
317                )
318                .await
319                .map_err(|error| {
320                    IdentityRuntimeError::Internal(format!("topology provider: {error}"))
321                })?,
322            None => self.runtime.desired_peer_edges.read().await.clone(),
323        };
324        Ok((roster, declared))
325    }
326
327    pub fn new(
328        runtime: Arc<IdentityRuntime>,
329        roster_provider: Arc<dyn RosterProvider>,
330        topology_provider: Option<Arc<dyn TopologyProvider>>,
331        customizer: Option<Arc<dyn AgentCustomizer>>,
332        mob_definition: Option<meerkat_mob::MobDefinition>,
333    ) -> Self {
334        Self::new_with_lazy_materialization(
335            runtime,
336            roster_provider,
337            topology_provider,
338            customizer,
339            mob_definition,
340            false,
341        )
342    }
343
344    pub fn new_with_lazy_materialization(
345        runtime: Arc<IdentityRuntime>,
346        roster_provider: Arc<dyn RosterProvider>,
347        topology_provider: Option<Arc<dyn TopologyProvider>>,
348        customizer: Option<Arc<dyn AgentCustomizer>>,
349        mob_definition: Option<meerkat_mob::MobDefinition>,
350        lazy_materialization: bool,
351    ) -> Self {
352        runtime.set_reset_roster_provider_context(
353            Some(roster_provider.clone()),
354            mob_definition.clone(),
355        );
356        Self {
357            runtime,
358            roster_provider,
359            topology_provider,
360            customizer,
361            mob_definition,
362            lazy_materialization,
363        }
364    }
365
366    pub async fn refresh_desired_topology(
367        &self,
368    ) -> Result<super::orchestrator::RestoreFlowResult, IdentityRuntimeError> {
369        let roster = self
370            .roster_provider
371            .roster(&RosterContext {
372                mob_definition: self.mob_definition.clone(),
373                previous_identities: Vec::new(),
374            })
375            .await
376            .map_err(|err| IdentityRuntimeError::Internal(format!("roster provider: {err}")))?;
377
378        if self.lazy_materialization {
379            super::orchestrator::lazy_register_flow(
380                &self.runtime,
381                &roster,
382                self.topology_provider.as_deref(),
383            )
384            .await
385        } else {
386            super::orchestrator::restore_flow(
387                &self.runtime,
388                &roster,
389                self.topology_provider.as_deref(),
390                self.customizer.as_deref(),
391            )
392            .await
393        }
394    }
395
396    /// Background repair for Broken identities. A rejected resume degrades the
397    /// identity to Broken while preserving the durable session; the documented
398    /// contract is "the next reconcile retries the resume" — but without this
399    /// task, nothing runs that reconcile: delivery refuses Broken identities
400    /// (REQ-13 fails loudly), `materialize` refuses the Broken state, and the
401    /// only retries were a manual `mobkit/reconcile_identity` RPC or a process
402    /// restart. HomeCore 0.7.23 sat with 14 preserved-but-parked identities
403    /// because of exactly that gap.
404    ///
405    /// The loop sleeps, and only when at least one identity is Broken re-runs
406    /// [`Self::refresh_desired_topology`] — the same idempotent flow the
407    /// reconcile RPC runs (eager: `restore_flow` retries the resume; lazy:
408    /// `lazy_register_flow` re-registers a store-Ready identity as Dormant so
409    /// on-demand materialization retries). Backoff doubles while identities
410    /// stay Broken — persistent causes (e.g. an upstream store regression)
411    /// produce bounded log noise, and transient causes (disk full, lock
412    /// contention) heal without a restart.
413    pub fn spawn_broken_identity_repair_task(
414        self: Arc<Self>,
415        policy: ContinuityRepairPolicy,
416    ) -> JoinHandle<()> {
417        tokio::spawn(async move {
418            let mut backoff = policy.initial_backoff;
419            loop {
420                tokio::time::sleep(backoff).await;
421                let broken = self.runtime.broken_identities().await;
422                if broken.is_empty() {
423                    backoff = policy.initial_backoff;
424                    continue;
425                }
426                tracing::info!(
427                    broken = broken.len(),
428                    "continuity repair: retrying restore for Broken identities"
429                );
430                if let Err(err) = self.refresh_desired_topology().await {
431                    tracing::warn!(
432                        error = %err,
433                        "continuity repair reconcile failed; backing off"
434                    );
435                    backoff = (backoff * 2).min(policy.max_backoff);
436                    continue;
437                }
438                let still_broken = self.runtime.broken_identities().await;
439                let healed = broken
440                    .iter()
441                    .filter(|id| !still_broken.contains(id))
442                    .count();
443                if healed > 0 {
444                    tracing::info!(
445                        healed,
446                        still_broken = still_broken.len(),
447                        "continuity repair healed identities"
448                    );
449                }
450                backoff = if still_broken.is_empty() {
451                    policy.initial_backoff
452                } else {
453                    (backoff * 2).min(policy.max_backoff)
454                };
455            }
456        })
457    }
458}
459
460/// Retry cadence for [`IdentityFirstRuntimeContext::spawn_broken_identity_repair_task`].
461#[derive(Debug, Clone, Copy, PartialEq, Eq)]
462pub struct ContinuityRepairPolicy {
463    /// Delay before the first check and after every fully-healed pass.
464    pub initial_backoff: Duration,
465    /// Ceiling for the doubling backoff while identities stay Broken.
466    pub max_backoff: Duration,
467}
468
469impl Default for ContinuityRepairPolicy {
470    fn default() -> Self {
471        Self {
472            initial_backoff: Duration::from_secs(30),
473            max_backoff: Duration::from_mins(10),
474        }
475    }
476}
477
478/// The identity-first runtime tracks active identities and enforces delivery,
479/// ownership, and lifecycle invariants.
480pub struct IdentityRuntime {
481    entries: RwLock<BTreeMap<AgentIdentity, IdentityEntry>>,
482    event_channels: RwLock<BTreeMap<AgentIdentity, broadcast::Sender<IdentityEvent>>>,
483    continuity_store: Arc<dyn ContinuityStore>,
484    lease_provider: Arc<dyn LeaseProvider>,
485    runtime_instance_id: String,
486    has_runtime_store: bool,
487    durability_policy: DurabilityPolicy,
488    bridge: Option<Arc<dyn SessionBridge>>,
489    reset_roster_source: StdRwLock<Option<ResetRosterSource>>,
490    runtime_services: AgentRuntimeServices,
491    managed_peer_edges: RwLock<BTreeSet<(AgentIdentity, AgentIdentity)>>,
492    managed_peer_reconcile_lock: Mutex<()>,
493    desired_peer_edges: RwLock<Vec<ManagedPeerEdge>>,
494    topology_controller: StdRwLock<Option<crate::topology_control::TopologyController>>,
495    materialization_locks: RwLock<BTreeMap<AgentIdentity, Arc<Mutex<()>>>>,
496    best_effort_materialization_locks: RwLock<BTreeMap<AgentIdentity, Arc<Mutex<()>>>>,
497    lifecycle_locks: RwLock<BTreeMap<AgentIdentity, Arc<Mutex<()>>>>,
498    customizer: RwLock<Option<Arc<dyn AgentCustomizer>>>,
499    agent_memory: RwLock<Option<AgentMemoryRuntimeInjector>>,
500    lease_renewal_notify: Notify,
501    default_timeout: Duration,
502    materialization_failure_backoff: RwLock<BTreeMap<AgentIdentity, MaterializationFailureBackoff>>,
503    error_hook: StdRwLock<Option<crate::unified_runtime::ErrorHook>>,
504}
505
506#[derive(Clone)]
507struct ResetRosterSource {
508    provider: Arc<dyn RosterProvider>,
509    mob_definition: Option<meerkat_mob::MobDefinition>,
510}
511
512#[derive(Debug, Clone)]
513struct MaterializationFailureBackoff {
514    suppress_until: Instant,
515    error: String,
516}
517
518impl IdentityRuntime {
519    /// Create a new identity runtime with the given configuration.
520    pub fn new(config: IdentityRuntimeConfig) -> Self {
521        Self {
522            entries: RwLock::new(BTreeMap::new()),
523            event_channels: RwLock::new(BTreeMap::new()),
524            continuity_store: config.continuity_store,
525            lease_provider: config.lease_provider,
526            runtime_instance_id: config.runtime_instance_id,
527            has_runtime_store: config.has_runtime_store,
528            durability_policy: config.durability_policy,
529            bridge: config.bridge,
530            reset_roster_source: StdRwLock::new(None),
531            runtime_services: AgentRuntimeServices::empty(),
532            managed_peer_edges: RwLock::new(BTreeSet::new()),
533            managed_peer_reconcile_lock: Mutex::new(()),
534            desired_peer_edges: RwLock::new(Vec::new()),
535            topology_controller: StdRwLock::new(None),
536            materialization_locks: RwLock::new(BTreeMap::new()),
537            best_effort_materialization_locks: RwLock::new(BTreeMap::new()),
538            lifecycle_locks: RwLock::new(BTreeMap::new()),
539            customizer: RwLock::new(None),
540            agent_memory: RwLock::new(None),
541            lease_renewal_notify: Notify::new(),
542            default_timeout: config.default_timeout.unwrap_or(Duration::from_secs(90)),
543            materialization_failure_backoff: RwLock::new(BTreeMap::new()),
544            error_hook: StdRwLock::new(None),
545        }
546    }
547
548    pub fn with_runtime_services(mut self, runtime_services: AgentRuntimeServices) -> Self {
549        self.runtime_services = runtime_services;
550        self
551    }
552
553    pub fn with_reset_roster_provider(self, provider: Arc<dyn RosterProvider>) -> Self {
554        self.set_reset_roster_provider(Some(provider));
555        self
556    }
557
558    pub fn with_reset_roster_provider_context(
559        self,
560        provider: Arc<dyn RosterProvider>,
561        mob_definition: Option<meerkat_mob::MobDefinition>,
562    ) -> Self {
563        self.set_reset_roster_provider_context(Some(provider), mob_definition);
564        self
565    }
566
567    pub(crate) fn runtime_services(&self) -> AgentRuntimeServices {
568        self.runtime_services.clone()
569    }
570
571    pub async fn set_agent_customizer(&self, customizer: Option<Arc<dyn AgentCustomizer>>) {
572        *self.customizer.write().await = customizer;
573    }
574
575    pub async fn set_agent_memory(&self, injector: Option<AgentMemoryRuntimeInjector>) {
576        *self.agent_memory.write().await = injector;
577    }
578
579    pub async fn agent_memory_supports_recall(&self) -> bool {
580        self.agent_memory.read().await.is_some()
581    }
582
583    pub async fn agent_memory_supports_remember(&self) -> bool {
584        self.agent_memory
585            .read()
586            .await
587            .as_ref()
588            .is_some_and(|injector| injector.provider().supports_remember())
589    }
590
591    pub async fn agent_memory_supports_forget(&self) -> bool {
592        self.agent_memory
593            .read()
594            .await
595            .as_ref()
596            .is_some_and(|injector| injector.provider().supports_forget())
597    }
598
599    pub async fn agent_memory_supports_update(&self) -> bool {
600        self.agent_memory
601            .read()
602            .await
603            .as_ref()
604            .is_some_and(|injector| injector.provider().supports_supersede())
605    }
606
607    pub async fn agent_memory_supports_manifest(&self) -> bool {
608        self.agent_memory
609            .read()
610            .await
611            .as_ref()
612            .is_some_and(|injector| injector.provider().supports_manifest())
613    }
614
615    pub async fn remember_agent_memory(
616        &self,
617        realm: &str,
618        identity: &AgentIdentity,
619        memory: NewAgentMemory,
620    ) -> Result<AgentMemoryRecord, AgentMemoryError> {
621        self.status(identity)
622            .await
623            .map_err(|err| AgentMemoryError::InvalidConfig(err.to_string()))?;
624        let provider = self
625            .agent_memory
626            .read()
627            .await
628            .as_ref()
629            .map(AgentMemoryRuntimeInjector::provider)
630            .ok_or_else(|| {
631                AgentMemoryError::InvalidConfig("agent memory is not configured".to_string())
632            })?;
633        provider.remember(realm, identity, memory).await
634    }
635
636    pub async fn forget_agent_memory(
637        &self,
638        realm: &str,
639        identity: &AgentIdentity,
640        memory_id: &str,
641    ) -> Result<AgentMemoryForgetResult, AgentMemoryError> {
642        self.status(identity)
643            .await
644            .map_err(|err| AgentMemoryError::InvalidConfig(err.to_string()))?;
645        let provider = self
646            .agent_memory
647            .read()
648            .await
649            .as_ref()
650            .map(AgentMemoryRuntimeInjector::provider)
651            .ok_or_else(|| {
652                AgentMemoryError::InvalidConfig("agent memory is not configured".to_string())
653            })?;
654        provider.forget(realm, identity, memory_id).await
655    }
656
657    /// Supersede `memory_id` within its lineage (the D4 fix): the new
658    /// title/body/tags become the active record; the prior stays
659    /// retrievable with provenance.
660    pub async fn update_agent_memory(
661        &self,
662        realm: &str,
663        identity: &AgentIdentity,
664        memory_id: &str,
665        memory: NewAgentMemory,
666    ) -> Result<MemoryId, AgentMemoryError> {
667        self.status(identity)
668            .await
669            .map_err(|err| AgentMemoryError::InvalidConfig(err.to_string()))?;
670        let provider = self
671            .agent_memory
672            .read()
673            .await
674            .as_ref()
675            .map(AgentMemoryRuntimeInjector::provider)
676            .ok_or_else(|| {
677                AgentMemoryError::InvalidConfig("agent memory is not configured".to_string())
678            })?;
679        let scope = MemoryScope::Identity {
680            realm: realm.to_string(),
681            identity: identity.as_str().to_string(),
682        };
683        let record = NewMemoryRecord {
684            kind: MemoryKind::Fact,
685            title: memory.title,
686            description: String::new(),
687            body: memory.body,
688            tags: memory.tags,
689            evidence: Vec::new(),
690            verification: None,
691        };
692        provider.supersede(&scope, memory_id, record).await
693    }
694
695    /// Tiered metadata manifest for the identity's own scope (§8.3).
696    pub async fn manifest_agent_memory(
697        &self,
698        realm: &str,
699        identity: &AgentIdentity,
700        tier: ManifestTier,
701    ) -> Result<Vec<RecordMeta>, AgentMemoryError> {
702        self.status(identity)
703            .await
704            .map_err(|err| AgentMemoryError::InvalidConfig(err.to_string()))?;
705        let provider = self
706            .agent_memory
707            .read()
708            .await
709            .as_ref()
710            .map(AgentMemoryRuntimeInjector::provider)
711            .ok_or_else(|| {
712                AgentMemoryError::InvalidConfig("agent memory is not configured".to_string())
713            })?;
714        let scope = MemoryScope::Identity {
715            realm: realm.to_string(),
716            identity: identity.as_str().to_string(),
717        };
718        provider.manifest(&[scope], tier).await
719    }
720
721    pub async fn recall_agent_memory(
722        &self,
723        request: AgentMemoryRecallRequest,
724    ) -> Result<Vec<AgentMemoryRecord>, AgentMemoryError> {
725        self.status(&request.identity)
726            .await
727            .map_err(|err| AgentMemoryError::InvalidConfig(err.to_string()))?;
728        let provider = self
729            .agent_memory
730            .read()
731            .await
732            .as_ref()
733            .map(AgentMemoryRuntimeInjector::provider)
734            .ok_or_else(|| {
735                AgentMemoryError::InvalidConfig("agent memory is not configured".to_string())
736            })?;
737        let records = provider.recall(request).await?;
738        // §9.2: explicit recall reads mark usage mechanically. Telemetry
739        // never fails the read — providers without usage support
740        // (markdown) return Unsupported, which is downgraded here.
741        if !records.is_empty() {
742            let ids: Vec<MemoryId> = records
743                .iter()
744                .map(|record| record.memory_id.clone())
745                .collect();
746            if let Err(err) = provider.mark_usage(&ids, UsageEvent::ExplicitRecall).await {
747                tracing::debug!(error = %err, "agent memory explicit-recall usage marking skipped");
748            }
749        }
750        Ok(records)
751    }
752
753    /// Attach the roster provider reset should consult for current specs.
754    pub fn set_reset_roster_provider(&self, provider: Option<Arc<dyn RosterProvider>>) {
755        self.set_reset_roster_provider_context(provider, None);
756    }
757
758    /// Attach the roster provider and context reset should consult for current specs.
759    pub fn set_reset_roster_provider_context(
760        &self,
761        provider: Option<Arc<dyn RosterProvider>>,
762        mob_definition: Option<meerkat_mob::MobDefinition>,
763    ) {
764        let source = provider.map(|provider| ResetRosterSource {
765            provider,
766            mob_definition,
767        });
768        match self.reset_roster_source.write() {
769            Ok(mut stored_source) => *stored_source = source,
770            Err(err) => {
771                tracing::warn!(
772                    error = %err,
773                    "identity runtime reset roster source lock poisoned; dropping provider update"
774                );
775            }
776        }
777    }
778
779    async fn adopt_current_roster_spec_for_reset(&self, identity: &AgentIdentity) {
780        let source = match self.reset_roster_source.read() {
781            Ok(stored_source) => stored_source.clone(),
782            Err(err) => {
783                tracing::warn!(
784                    error = %err,
785                    "reset: roster source lock poisoned; rebuilding on stored spec"
786                );
787                None
788            }
789        };
790        if let Some(source) = source {
791            self.adopt_roster_spec_with_context(
792                &source.provider,
793                identity,
794                source.mob_definition.clone(),
795            )
796            .await;
797        }
798    }
799
800    /// Attach a best-effort operational error hook used for alerting.
801    pub fn set_error_hook(&self, hook: Option<crate::unified_runtime::ErrorHook>) {
802        match self.error_hook.write() {
803            Ok(mut stored_hook) => *stored_hook = hook,
804            Err(err) => {
805                tracing::warn!(
806                    error = %err,
807                    "identity runtime error hook lock poisoned; dropping hook update"
808                );
809            }
810        }
811    }
812
813    /// Spawn a background supervisor that renews active identity leases before
814    /// they reach their TTL deadline.
815    pub fn spawn_lease_renewal_task(self: Arc<Self>) -> JoinHandle<()> {
816        self.spawn_lease_renewal_task_with_poll_interval(DEFAULT_LEASE_RENEWAL_MAX_POLL_INTERVAL)
817    }
818
819    /// Spawn a lease renewal supervisor with a caller-provided maximum poll
820    /// interval. Embedders can use this for shorter external lease TTLs; tests
821    /// use it to exercise renewal without waiting on wall-clock TTLs.
822    pub fn spawn_lease_renewal_task_with_poll_interval(
823        self: Arc<Self>,
824        max_poll_interval: Duration,
825    ) -> JoinHandle<()> {
826        let max_poll_interval = max_poll_interval.max(DEFAULT_LEASE_RENEWAL_MIN_POLL_INTERVAL);
827        tokio::spawn(async move {
828            let mut consecutive_failures: u32 = 0;
829            loop {
830                let base = self.lease_renewal_sleep_interval(max_poll_interval).await;
831                // While the provider is failing, hold off at least the backoff
832                // delay so a persistent outage retries at a bounded rate
833                // instead of the TTL-derived floor (down to 10ms).
834                let sleep = if consecutive_failures > 0 {
835                    base.max(lease_renewal_failure_backoff(
836                        consecutive_failures,
837                        max_poll_interval,
838                    ))
839                } else {
840                    base
841                };
842                tokio::select! {
843                    () = tokio::time::sleep(sleep) => {}
844                    () = self.lease_renewal_notify.notified() => {}
845                }
846                match self.renew_due_leases_once().await {
847                    Ok(_) => consecutive_failures = 0,
848                    Err(err) => {
849                        // Warn once, then debounce to debug so a backend outage
850                        // can't flood the log at the renewal cadence.
851                        if consecutive_failures == 0 {
852                            tracing::warn!(
853                                error = %err,
854                                "identity-first proactive lease renewal tick failed; backing off"
855                            );
856                        } else {
857                            tracing::debug!(
858                                error = %err,
859                                consecutive_failures,
860                                "identity-first lease renewal still failing; backing off"
861                            );
862                        }
863                        consecutive_failures = consecutive_failures.saturating_add(1);
864                    }
865                }
866            }
867        })
868    }
869
870    async fn lease_renewal_sleep_interval(&self, max_poll_interval: Duration) -> Duration {
871        let entries = self.entries.read().await;
872        entries
873            .values()
874            .filter(|entry| entry.state == IdentityLifecycleState::Active)
875            .filter_map(|entry| entry.lease.as_ref())
876            .map(|lease| (lease.ttl / 10).max(DEFAULT_LEASE_RENEWAL_MIN_POLL_INTERVAL))
877            .min()
878            .unwrap_or(max_poll_interval)
879            .min(max_poll_interval)
880    }
881
882    /// Renew every active lease that has entered the runtime's renewal window.
883    pub async fn renew_due_leases_once(&self) -> Result<usize, IdentityRuntimeError> {
884        let due = {
885            let entries = self.entries.read().await;
886            entries
887                .iter()
888                .filter(|(_, entry)| entry.state == IdentityLifecycleState::Active)
889                .filter_map(|(identity, entry)| {
890                    entry
891                        .lease
892                        .as_ref()
893                        .filter(|lease| !lease.is_healthy())
894                        .map(|_| identity.clone())
895                })
896                .collect::<Vec<_>>()
897        };
898
899        let mut renewed = 0;
900        let mut first_error = None;
901        for identity in due {
902            let lifecycle_lock = self.lifecycle_lock_for(&identity).await;
903            let _lifecycle_guard = lifecycle_lock.lock().await;
904            match self.ensure_active_lease(&identity).await {
905                Ok(_) => renewed += 1,
906                Err(err) => {
907                    if first_error.is_none() {
908                        first_error = Some(err);
909                    }
910                }
911            }
912        }
913        if let Some(err) = first_error {
914            Err(err)
915        } else {
916            Ok(renewed)
917        }
918    }
919
920    async fn release_uninstalled_materialize_lease(&self, grant: &LeaseGrant) -> Option<String> {
921        self.lease_provider
922            .release_leases(std::slice::from_ref(grant))
923            .await
924            .err()
925            .map(|err| err.to_string())
926    }
927
928    pub async fn set_desired_peer_edges(&self, edges: Vec<ManagedPeerEdge>) {
929        *self.desired_peer_edges.write().await = edges;
930    }
931
932    pub(crate) fn set_topology_controller(
933        &self,
934        controller: crate::topology_control::TopologyController,
935    ) {
936        *self
937            .topology_controller
938            .write()
939            .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(controller);
940    }
941
942    fn topology_controller(&self) -> Option<crate::topology_control::TopologyController> {
943        self.topology_controller
944            .read()
945            .unwrap_or_else(std::sync::PoisonError::into_inner)
946            .clone()
947    }
948
949    pub async fn desired_peer_edges(&self) -> Vec<ManagedPeerEdge> {
950        let declared = self.desired_peer_edges.read().await.clone();
951        match self.topology_controller() {
952            Some(controller) => {
953                // Materialization/console reads must not observe the target
954                // intent while a topology transaction is between WAL and
955                // terminal commit/rollback.
956                let _admission = controller.mutation_guard().await;
957                match controller.compose_managed_peer_edges(&declared).await {
958                    Ok(edges) => edges,
959                    Err(error) => {
960                        tracing::error!(error = %error, "failed to compose identity topology overlay");
961                        Vec::new()
962                    }
963                }
964            }
965            None => declared,
966        }
967    }
968
969    async fn registered_identities(&self) -> Vec<AgentIdentity> {
970        self.entries.read().await.keys().cloned().collect()
971    }
972
973    async fn reachable_peer_identities(&self, identity: &AgentIdentity) -> Vec<AgentIdentity> {
974        self.desired_peer_edges()
975            .await
976            .iter()
977            .filter_map(|edge| {
978                if edge.a() == identity {
979                    Some(edge.b().clone())
980                } else if edge.b() == identity {
981                    Some(edge.a().clone())
982                } else {
983                    None
984                }
985            })
986            .collect::<BTreeSet<_>>()
987            .into_iter()
988            .collect()
989    }
990
991    #[must_use]
992    pub fn has_session_bridge(&self) -> bool {
993        self.bridge.is_some()
994    }
995
996    pub(crate) async fn logical_peer_edges(
997        &self,
998    ) -> Result<Vec<ManagedPeerEdge>, IdentityRuntimeError> {
999        let Some(bridge) = self.bridge.as_ref() else {
1000            return Ok(Vec::new());
1001        };
1002        let runtime_identities: BTreeMap<AgentRuntimeId, AgentIdentity> = self
1003            .entries
1004            .read()
1005            .await
1006            .iter()
1007            .filter_map(|(identity, entry)| {
1008                entry
1009                    .continuity
1010                    .as_ref()
1011                    .map(|record| (record.agent_runtime_id.clone(), identity.clone()))
1012            })
1013            .collect();
1014        let runtime_edges = bridge.current_member_wires().await.map_err(|error| {
1015            IdentityRuntimeError::Internal(format!("bridge current_member_wires: {error}"))
1016        })?;
1017        Ok(runtime_edges
1018            .into_iter()
1019            .filter_map(|(runtime_a, runtime_b)| {
1020                let a = runtime_identities.get(&runtime_a)?.clone();
1021                let b = runtime_identities.get(&runtime_b)?.clone();
1022                ManagedPeerEdge::new(a, b).ok()
1023            })
1024            .collect())
1025    }
1026
1027    pub(crate) async fn managed_peer_edges_snapshot(
1028        &self,
1029    ) -> BTreeSet<(AgentIdentity, AgentIdentity)> {
1030        self.managed_peer_edges.read().await.clone()
1031    }
1032
1033    /// Logical identity actuator used only while the shared topology
1034    /// controller's admission lock is already held by TopologyRuntimeHandle.
1035    pub(crate) async fn mutate_managed_peer_edge_admitted(
1036        &self,
1037        action: crate::topology_control::TopologyAction,
1038        edge: &ManagedPeerEdge,
1039    ) -> Result<(), IdentityRuntimeError> {
1040        let _guard = self.managed_peer_reconcile_lock.lock().await;
1041        let Some(bridge) = self.bridge.clone() else {
1042            return Ok(());
1043        };
1044        let (runtime_a, runtime_b) = {
1045            let entries = self.entries.read().await;
1046            let resolve = |identity: &AgentIdentity| {
1047                entries
1048                    .get(identity)
1049                    .filter(|entry| entry.state == IdentityLifecycleState::Active)
1050                    .and_then(|entry| entry.continuity.as_ref())
1051                    .map(|record| record.agent_runtime_id.clone())
1052                    .ok_or_else(|| {
1053                        IdentityRuntimeError::Internal(format!(
1054                            "topology endpoint is not active: {identity}"
1055                        ))
1056                    })
1057            };
1058            (resolve(edge.a())?, resolve(edge.b())?)
1059        };
1060        let key = (edge.a().clone(), edge.b().clone());
1061        let current = bridge.current_member_wires().await.map_err(|error| {
1062            IdentityRuntimeError::Internal(format!("bridge current_member_wires: {error}"))
1063        })?;
1064        let actual = current.iter().any(|(a, b)| {
1065            (a == &runtime_a && b == &runtime_b) || (a == &runtime_b && b == &runtime_a)
1066        });
1067        let any_half = bridge
1068            .current_member_wires_any_half()
1069            .await
1070            .map_err(|error| {
1071                IdentityRuntimeError::Internal(format!(
1072                    "bridge current_member_wires_any_half: {error}"
1073                ))
1074            })?
1075            .iter()
1076            .any(|(a, b)| {
1077                (a == &runtime_a && b == &runtime_b) || (a == &runtime_b && b == &runtime_a)
1078            });
1079        match action {
1080            crate::topology_control::TopologyAction::Connect if !actual => bridge
1081                .wire_peers_batch(&[(runtime_a, runtime_b)])
1082                .await
1083                .map_err(|error| {
1084                    IdentityRuntimeError::Internal(format!("bridge wire_peers_batch: {error}"))
1085                })?,
1086            crate::topology_control::TopologyAction::Reconnect => {
1087                if any_half {
1088                    bridge
1089                        .unwire_peer(&runtime_a, &runtime_b)
1090                        .await
1091                        .map_err(|error| {
1092                            IdentityRuntimeError::Internal(format!("bridge unwire_peer: {error}"))
1093                        })?;
1094                }
1095                bridge
1096                    .wire_peers_batch(&[(runtime_a, runtime_b)])
1097                    .await
1098                    .map_err(|error| {
1099                        IdentityRuntimeError::Internal(format!("bridge wire_peers_batch: {error}"))
1100                    })?;
1101            }
1102            crate::topology_control::TopologyAction::Disconnect if any_half => bridge
1103                .unwire_peer(&runtime_a, &runtime_b)
1104                .await
1105                .map_err(|error| {
1106                    IdentityRuntimeError::Internal(format!("bridge unwire_peer: {error}"))
1107                })?,
1108            _ => {}
1109        }
1110        let mut managed = self.managed_peer_edges.write().await;
1111        if matches!(action, crate::topology_control::TopologyAction::Disconnect) {
1112            managed.remove(&key);
1113        } else {
1114            managed.insert(key);
1115        }
1116        Ok(())
1117    }
1118
1119    /// Apply identity-first managed topology to the concrete mob graph.
1120    ///
1121    /// Topology providers return stable logical identities. The mob comms graph
1122    /// is keyed by active runtime member IDs, so this resolves each endpoint
1123    /// through continuity records before calling the same-mob bridge wire APIs.
1124    pub async fn reconcile_managed_peer_edges(
1125        &self,
1126        desired_edges: &[ManagedPeerEdge],
1127    ) -> Result<(), IdentityRuntimeError> {
1128        let topology_controller = self.topology_controller();
1129        let _topology_guard = match topology_controller.as_ref() {
1130            Some(controller) => Some(controller.mutation_guard().await),
1131            None => None,
1132        };
1133        if let Some(controller) = topology_controller.as_ref() {
1134            controller
1135                .prepare_pending_recovery()
1136                .await
1137                .map_err(|error| {
1138                    IdentityRuntimeError::Internal(format!("topology recovery journal: {error}"))
1139                })?;
1140        }
1141        let _guard = self.managed_peer_reconcile_lock.lock().await;
1142        let composed_edges;
1143        let desired_edges = if let Some(controller) = topology_controller.as_ref() {
1144            composed_edges = controller
1145                .compose_managed_peer_edges(desired_edges)
1146                .await
1147                .map_err(|error| {
1148                    IdentityRuntimeError::Internal(format!("topology overlay: {error}"))
1149                })?;
1150            composed_edges.as_slice()
1151        } else {
1152            desired_edges
1153        };
1154        let Some(bridge) = self.bridge.clone() else {
1155            if let Some(controller) = topology_controller.as_ref() {
1156                controller
1157                    .finalize_recovered_pending(true)
1158                    .await
1159                    .map_err(|error| {
1160                        IdentityRuntimeError::Internal(format!(
1161                            "topology recovery receipt: {error}"
1162                        ))
1163                    })?;
1164            }
1165            return Ok(());
1166        };
1167
1168        let active_runtimes: BTreeMap<AgentIdentity, AgentRuntimeId> = {
1169            let entries = self.entries.read().await;
1170            entries
1171                .iter()
1172                .filter_map(|(identity, entry)| {
1173                    if entry.state != IdentityLifecycleState::Active {
1174                        return None;
1175                    }
1176                    entry
1177                        .continuity
1178                        .as_ref()
1179                        .map(|record| (identity.clone(), record.agent_runtime_id.clone()))
1180                })
1181                .collect()
1182        };
1183        let runtime_identities: BTreeMap<AgentRuntimeId, AgentIdentity> = active_runtimes
1184            .iter()
1185            .map(|(identity, runtime_id)| (runtime_id.clone(), identity.clone()))
1186            .collect();
1187        let current_logical_edges: Option<BTreeSet<(AgentIdentity, AgentIdentity)>> =
1188            match bridge.current_member_wires().await {
1189                Ok(current_runtime_edges) => Some(
1190                    current_runtime_edges
1191                        .iter()
1192                        .filter_map(|(runtime_a, runtime_b)| {
1193                            let a = runtime_identities.get(runtime_a)?;
1194                            let b = runtime_identities.get(runtime_b)?;
1195                            if a <= b {
1196                                Some((a.clone(), b.clone()))
1197                            } else {
1198                                Some((b.clone(), a.clone()))
1199                            }
1200                        })
1201                        .collect(),
1202                ),
1203                Err(err) => {
1204                    tracing::debug!(
1205                        error = %err,
1206                        "identity-first topology reconcile could not inspect current member wires"
1207                    );
1208                    None
1209                }
1210            };
1211        let current_any_half_edges: Option<BTreeSet<(AgentIdentity, AgentIdentity)>> =
1212            match bridge.current_member_wires_any_half().await {
1213                Ok(current_runtime_edges) => Some(
1214                    current_runtime_edges
1215                        .iter()
1216                        .filter_map(|(runtime_a, runtime_b)| {
1217                            let a = runtime_identities.get(runtime_a)?;
1218                            let b = runtime_identities.get(runtime_b)?;
1219                            if a <= b {
1220                                Some((a.clone(), b.clone()))
1221                            } else {
1222                                Some((b.clone(), a.clone()))
1223                            }
1224                        })
1225                        .collect(),
1226                ),
1227                Err(error) => {
1228                    tracing::debug!(
1229                        %error,
1230                        "identity-first topology reconcile could not inspect orphan wire halves"
1231                    );
1232                    None
1233                }
1234            };
1235
1236        let desired: BTreeSet<(AgentIdentity, AgentIdentity)> = desired_edges
1237            .iter()
1238            .map(|edge| (edge.a().clone(), edge.b().clone()))
1239            .collect();
1240
1241        let managed_snapshot = self.managed_peer_edges.read().await.clone();
1242        let edge_is_managed_and_live = |edge: &(AgentIdentity, AgentIdentity)| {
1243            // Managed-but-missing live edges are retried deliberately so tolerant topology restores self-heal.
1244            managed_snapshot.contains(edge)
1245                && current_logical_edges
1246                    .as_ref()
1247                    .is_none_or(|edges| edges.contains(edge))
1248        };
1249        let retained_logical_edges: Vec<(AgentIdentity, AgentIdentity)> = desired
1250            .iter()
1251            .filter(|edge| !edge_is_managed_and_live(edge))
1252            .filter(|edge| {
1253                current_logical_edges
1254                    .as_ref()
1255                    .is_some_and(|edges| edges.contains(*edge))
1256            })
1257            .filter(|(a, b)| active_runtimes.contains_key(a) && active_runtimes.contains_key(b))
1258            .cloned()
1259            .collect();
1260        let to_wire: Vec<(AgentIdentity, AgentIdentity, AgentRuntimeId, AgentRuntimeId)> = desired
1261            .iter()
1262            .filter(|edge| !edge_is_managed_and_live(edge))
1263            .filter(|edge| {
1264                current_logical_edges
1265                    .as_ref()
1266                    .is_none_or(|edges| !edges.contains(*edge))
1267            })
1268            .filter_map(|(a, b)| {
1269                let runtime_a = active_runtimes.get(a)?;
1270                let runtime_b = active_runtimes.get(b)?;
1271                Some((a.clone(), b.clone(), runtime_a.clone(), runtime_b.clone()))
1272            })
1273            .collect();
1274
1275        let stale: Vec<(AgentIdentity, AgentIdentity)> = managed_snapshot
1276            .iter()
1277            .filter(|edge| !desired.contains(*edge))
1278            .cloned()
1279            .collect();
1280        let to_unwire: Vec<(AgentIdentity, AgentIdentity, AgentRuntimeId, AgentRuntimeId)> = stale
1281            .iter()
1282            .filter_map(|(a, b)| {
1283                let runtime_a = active_runtimes.get(a)?;
1284                let runtime_b = active_runtimes.get(b)?;
1285                if current_any_half_edges
1286                    .as_ref()
1287                    .is_some_and(|edges| !edges.contains(&(a.clone(), b.clone())))
1288                {
1289                    return None;
1290                }
1291                Some((a.clone(), b.clone(), runtime_a.clone(), runtime_b.clone()))
1292            })
1293            .collect();
1294
1295        let wire_logical_edges = to_wire
1296            .iter()
1297            .map(|(a, b, _, _)| (a.clone(), b.clone()))
1298            .collect::<Vec<_>>();
1299        let wire_runtime_edges = to_wire
1300            .iter()
1301            .map(|(_, _, runtime_a, runtime_b)| (runtime_a.clone(), runtime_b.clone()))
1302            .collect::<Vec<_>>();
1303        if !wire_runtime_edges.is_empty() {
1304            bridge
1305                .wire_peers_batch(&wire_runtime_edges)
1306                .await
1307                .map_err(|e| {
1308                    IdentityRuntimeError::Internal(format!("bridge wire_peers_batch: {e}"))
1309                })?;
1310        }
1311
1312        let unwire_results =
1313            stream::iter(to_unwire.into_iter().map(|(a, b, runtime_a, runtime_b)| {
1314                let bridge = bridge.clone();
1315                async move {
1316                    let result = bridge
1317                        .unwire_peer(&runtime_a, &runtime_b)
1318                        .await
1319                        .map_err(|e| format!("{e}"));
1320                    (a, b, result)
1321                }
1322            }))
1323            .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
1324            .collect::<Vec<_>>()
1325            .await;
1326
1327        let mut managed = self.managed_peer_edges.write().await;
1328        for (a, b) in retained_logical_edges {
1329            managed.insert((a, b));
1330        }
1331        for (a, b) in wire_logical_edges {
1332            managed.insert((a, b));
1333        }
1334
1335        for (a, b) in stale {
1336            let key = (a.clone(), b.clone());
1337            if !active_runtimes.contains_key(&a)
1338                || !active_runtimes.contains_key(&b)
1339                || current_logical_edges
1340                    .as_ref()
1341                    .is_some_and(|edges| !edges.contains(&key))
1342            {
1343                managed.remove(&key);
1344            }
1345        }
1346        for (a, b, result) in unwire_results {
1347            result
1348                .map_err(|e| IdentityRuntimeError::Internal(format!("bridge unwire_peer: {e}")))?;
1349            managed.remove(&(a, b));
1350        }
1351
1352        if let Some(controller) = topology_controller.as_ref() {
1353            controller
1354                .finalize_recovered_pending(true)
1355                .await
1356                .map_err(|error| {
1357                    IdentityRuntimeError::Internal(format!("topology recovery receipt: {error}"))
1358                })?;
1359        }
1360        Ok(())
1361    }
1362
1363    /// Emit an event for the given identity. Best-effort — no error if no subscribers.
1364    async fn emit_event(&self, identity: &AgentIdentity, event: IdentityEvent) {
1365        let channels = self.event_channels.read().await;
1366        if let Some(tx) = channels.get(identity) {
1367            let _ = tx.send(event);
1368        }
1369    }
1370
1371    fn emit_error(&self, event: crate::unified_runtime::types::ErrorEvent) {
1372        let hook = match self.error_hook.read() {
1373            Ok(stored_hook) => stored_hook.clone(),
1374            Err(err) => {
1375                tracing::warn!(
1376                    error = %err,
1377                    "identity runtime error hook lock poisoned; dropping error event"
1378                );
1379                None
1380            }
1381        };
1382        if let Some(hook) = hook {
1383            tokio::spawn(async move {
1384                let () = hook(event).await;
1385            });
1386        }
1387    }
1388
1389    async fn materialization_backoff_error(&self, identity: &AgentIdentity) -> Option<String> {
1390        let backoffs = self.materialization_failure_backoff.read().await;
1391        let backoff = backoffs.get(identity)?;
1392        if Instant::now() < backoff.suppress_until {
1393            Some(backoff.error.clone())
1394        } else {
1395            None
1396        }
1397    }
1398
1399    async fn clear_materialization_backoff(&self, identity: &AgentIdentity) {
1400        self.materialization_failure_backoff
1401            .write()
1402            .await
1403            .remove(identity);
1404    }
1405
1406    async fn record_best_effort_materialization_failure(
1407        &self,
1408        identity: &AgentIdentity,
1409        initiator: Option<&AgentIdentity>,
1410        operation: &'static str,
1411        err: &IdentityRuntimeError,
1412    ) {
1413        let error = err.to_string();
1414        let suppress_until = Instant::now() + MATERIALIZATION_FAILURE_BACKOFF;
1415        self.materialization_failure_backoff.write().await.insert(
1416            identity.clone(),
1417            MaterializationFailureBackoff {
1418                suppress_until,
1419                error: error.clone(),
1420            },
1421        );
1422        self.emit_error(
1423            crate::unified_runtime::types::ErrorEvent::IdentityMaterializationFailure {
1424                identity: identity.to_string(),
1425                initiator: initiator.map(ToString::to_string),
1426                operation: operation.to_string(),
1427                error,
1428            },
1429        );
1430    }
1431
1432    // -----------------------------------------------------------------------
1433    // Registration / activation
1434    // -----------------------------------------------------------------------
1435
1436    /// Register an identity entry in the runtime (called during restore flow).
1437    pub async fn register(
1438        &self,
1439        spec: DurableAgentSpec,
1440        state: IdentityLifecycleState,
1441        continuity: Option<ContinuityRecord>,
1442        lease: Option<LeaseGrant>,
1443    ) {
1444        let identity = spec.identity.clone();
1445        let cpv = continuity
1446            .as_ref()
1447            .map(|r| r.checkpoint_version)
1448            .unwrap_or(CheckpointVersion::new(0));
1449        let lease_entry = lease.map(|g| LeaseEntry {
1450            fencing_token: g.fencing_token,
1451            ttl: g.ttl,
1452            acquired_at: Instant::now(),
1453        });
1454        let entry = IdentityEntry {
1455            spec,
1456            state,
1457            continuity,
1458            lease: lease_entry,
1459            checkpoint_version: cpv,
1460            has_runtime_store: self.has_runtime_store,
1461        };
1462        let has_active_lease =
1463            entry.state == IdentityLifecycleState::Active && entry.lease.is_some();
1464        self.entries.write().await.insert(identity.clone(), entry);
1465
1466        // Create event channel for this identity
1467        let (tx, _) = broadcast::channel(IDENTITY_EVENT_CHANNEL_CAPACITY);
1468        self.event_channels
1469            .write()
1470            .await
1471            .insert(identity.clone(), tx);
1472        if has_active_lease {
1473            self.lease_renewal_notify.notify_one();
1474        }
1475    }
1476
1477    async fn materialization_lock_for(&self, identity: &AgentIdentity) -> Arc<Mutex<()>> {
1478        if let Some(lock) = self.materialization_locks.read().await.get(identity) {
1479            return lock.clone();
1480        }
1481        let mut locks = self.materialization_locks.write().await;
1482        locks
1483            .entry(identity.clone())
1484            .or_insert_with(|| Arc::new(Mutex::new(())))
1485            .clone()
1486    }
1487
1488    async fn best_effort_materialization_lock_for(
1489        &self,
1490        identity: &AgentIdentity,
1491    ) -> Arc<Mutex<()>> {
1492        if let Some(lock) = self
1493            .best_effort_materialization_locks
1494            .read()
1495            .await
1496            .get(identity)
1497        {
1498            return lock.clone();
1499        }
1500        let mut locks = self.best_effort_materialization_locks.write().await;
1501        locks
1502            .entry(identity.clone())
1503            .or_insert_with(|| Arc::new(Mutex::new(())))
1504            .clone()
1505    }
1506
1507    async fn lifecycle_lock_for(&self, identity: &AgentIdentity) -> Arc<Mutex<()>> {
1508        if let Some(lock) = self.lifecycle_locks.read().await.get(identity) {
1509            return lock.clone();
1510        }
1511        let mut locks = self.lifecycle_locks.write().await;
1512        locks
1513            .entry(identity.clone())
1514            .or_insert_with(|| Arc::new(Mutex::new(())))
1515            .clone()
1516    }
1517
1518    /// Materialize a dormant identity into a concrete mob member/session.
1519    ///
1520    /// This is the lazy counterpart to eager `restore_flow`: it performs the
1521    /// expensive bridge create/resume and snapshot load only when an identity is
1522    /// actually touched. Parallel calls for one identity coalesce on a
1523    /// per-identity lock and re-check state after acquiring it.
1524    pub async fn materialize(
1525        &self,
1526        identity: &AgentIdentity,
1527    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
1528        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
1529        let _lifecycle_guard = lifecycle_lock.lock().await;
1530        let lock = self.materialization_lock_for(identity).await;
1531        let _guard = lock.lock().await;
1532
1533        let (spec, continuity, state) = {
1534            let entries = self.entries.read().await;
1535            let entry = entries
1536                .get(identity)
1537                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1538            if entry.state == IdentityLifecycleState::Active {
1539                let continuity = entry.continuity.clone().ok_or_else(|| {
1540                    IdentityRuntimeError::Internal(format!(
1541                        "active identity {identity} has no continuity record"
1542                    ))
1543                })?;
1544                drop(entries);
1545                self.clear_materialization_backoff(identity).await;
1546                return Ok(continuity);
1547            }
1548            (entry.spec.clone(), entry.continuity.clone(), entry.state)
1549        };
1550        let original_continuity = continuity.clone();
1551        let continuity = if durable_spec_uses_external_binding(&spec) {
1552            None
1553        } else {
1554            continuity
1555        };
1556
1557        match state {
1558            IdentityLifecycleState::Dormant | IdentityLifecycleState::Uninitialized => {}
1559            IdentityLifecycleState::Broken
1560            | IdentityLifecycleState::Retiring
1561            | IdentityLifecycleState::Suspended => {
1562                return Err(IdentityRuntimeError::InvalidState {
1563                    identity: identity.clone(),
1564                    state,
1565                    operation: "materialize",
1566                });
1567            }
1568            IdentityLifecycleState::Active => unreachable!("active handled above"),
1569        }
1570
1571        let lease_results = self
1572            .lease_provider
1573            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
1574            .await
1575            .map_err(IdentityRuntimeError::Lease)?;
1576        let grant = match lease_results.get(identity) {
1577            Some(super::types::LeaseAcquireResult::Acquired(grant)) => grant.clone(),
1578            Some(super::types::LeaseAcquireResult::AlreadyHeld { holder, .. }) => {
1579                tracing::error!(
1580                    %identity,
1581                    holder = %holder,
1582                    "single-embodiment guard: refusing to materialize an identity whose \
1583                     durable lease is held by another live runtime instance"
1584                );
1585                return Err(IdentityRuntimeError::AlreadyEmbodied {
1586                    identity: identity.clone(),
1587                    holder: holder.clone(),
1588                });
1589            }
1590            None => return Err(IdentityRuntimeError::NoActiveLease(identity.clone())),
1591        };
1592        if let Some(record) = continuity.as_ref()
1593            && let Err(err) = self
1594                .continuity_store
1595                .upsert_continuity_record(record, grant.fencing_token)
1596                .await
1597        {
1598            let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1599            return Err(IdentityRuntimeError::Internal(format!(
1600                "continuity upsert before materialize: {err}{}",
1601                cleanup_error
1602                    .as_ref()
1603                    .map(|e| format!("; lease cleanup failed: {e}"))
1604                    .unwrap_or_default(),
1605            )));
1606        }
1607
1608        let active_peers = self.entries.read().await.keys().cloned().collect();
1609        let managed_edges = self.desired_peer_edges.read().await.clone();
1610        let build_context = AgentBuildContext {
1611            identity: identity.clone(),
1612            active_peers,
1613            managed_edges,
1614            runtime_services: self.runtime_services(),
1615        };
1616        let mut draft = super::types::AgentBuildDraft {
1617            model: None,
1618            system_prompt: None,
1619            additional_instructions: spec.additional_instructions.clone(),
1620            labels: spec.labels.clone(),
1621            app_context: spec.context.clone(),
1622            external_tools: Vec::new(),
1623            local_external_tools: Default::default(),
1624        };
1625        if let Some(customizer) = self.customizer.read().await.clone()
1626            && let Err(err) = customizer
1627                .customize_build(&build_context, &spec, &mut draft)
1628                .await
1629        {
1630            let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1631            return Err(IdentityRuntimeError::Internal(format!(
1632                "customizer: {err}{}",
1633                cleanup_error
1634                    .as_ref()
1635                    .map(|e| format!("; lease cleanup failed: {e}"))
1636                    .unwrap_or_default(),
1637            )));
1638        }
1639
1640        let mut abandoned_session_registrations: Vec<SessionId> = Vec::new();
1641        let mut record = if let Some(mut record) = continuity {
1642            let snapshot = match self
1643                .continuity_store
1644                .load_session_snapshot(&record.session_id)
1645                .await
1646            {
1647                Ok(snapshot) => snapshot,
1648                Err(err) => {
1649                    let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1650                    return Err(IdentityRuntimeError::Internal(format!(
1651                        "load session snapshot before materialize: {err}{}",
1652                        cleanup_error
1653                            .as_ref()
1654                            .map(|e| format!("; lease cleanup failed: {e}"))
1655                            .unwrap_or_default(),
1656                    )));
1657                }
1658            };
1659
1660            if let Some(bridge) = self.bridge.as_ref() {
1661                if let Err(err) = bridge
1662                    .register_session_runtime_state(
1663                        &record.session_id,
1664                        identity,
1665                        record.generation,
1666                        record.checkpoint_version,
1667                        grant.fencing_token,
1668                    )
1669                    .await
1670                {
1671                    let unregister_error = Self::unregister_bridge_session_runtime_states(
1672                        bridge.as_ref(),
1673                        std::slice::from_ref(&record.session_id),
1674                    )
1675                    .await;
1676                    let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1677                    return Err(IdentityRuntimeError::Internal(format!(
1678                        "bridge register_session_runtime_state: {err}{}{}",
1679                        unregister_error
1680                            .as_ref()
1681                            .map(|e| format!("; unregister session failed: {e}"))
1682                            .unwrap_or_default(),
1683                        cleanup_error
1684                            .as_ref()
1685                            .map(|e| format!("; lease cleanup failed: {e}"))
1686                            .unwrap_or_default(),
1687                    )));
1688                }
1689                let registered_session_id = record.session_id.clone();
1690                let snapshot = snapshot.unwrap_or(SessionSnapshot { data: Vec::new() });
1691                let outcome = bridge
1692                    .resume_session(
1693                        identity,
1694                        &record.agent_runtime_id,
1695                        &spec,
1696                        &draft,
1697                        &record.session_id,
1698                        &snapshot,
1699                    )
1700                    .await;
1701                let outcome = match outcome {
1702                    Ok(outcome) => outcome,
1703                    Err(err) => {
1704                        // A rejected resume must NEVER abandon the durable
1705                        // session (the transcript is the only copy). Keep the
1706                        // identity → session binding intact — no retire, no
1707                        // continuity rebind — mark the identity Broken with
1708                        // the error attached so this send fails loudly, and
1709                        // let the next reconcile retry the resume. Unregister
1710                        // only the session-runtime-state bookkeeping (it is
1711                        // re-registered on retry).
1712                        tracing::error!(
1713                            %identity,
1714                            session_id = %registered_session_id,
1715                            error = %err,
1716                            "materialize resume rejected; marking identity Broken and preserving \
1717                             the durable session for reconcile retry"
1718                        );
1719                        let unregister_error = Self::unregister_bridge_session_runtime_states(
1720                            bridge.as_ref(),
1721                            std::slice::from_ref(&registered_session_id),
1722                        )
1723                        .await;
1724                        let lease_cleanup_error =
1725                            self.release_uninstalled_materialize_lease(&grant).await;
1726                        {
1727                            let mut entries = self.entries.write().await;
1728                            if let Some(entry) = entries.get_mut(identity) {
1729                                entry.state = IdentityLifecycleState::Broken;
1730                                entry.lease = None;
1731                            }
1732                        }
1733                        self.emit_event(
1734                            identity,
1735                            IdentityEvent::StateChanged {
1736                                identity: identity.clone(),
1737                                new_state: IdentityLifecycleState::Broken,
1738                            },
1739                        )
1740                        .await;
1741                        let detail = format!(
1742                            "bridge resume_session rejected (identity degraded, durable session \
1743                             preserved for reconcile retry): {err}{}{}",
1744                            unregister_error
1745                                .as_ref()
1746                                .map(|e| format!("; unregister session failed: {e}"))
1747                                .unwrap_or_default(),
1748                            lease_cleanup_error
1749                                .as_ref()
1750                                .map(|e| format!("; lease cleanup failed: {e}"))
1751                                .unwrap_or_default(),
1752                        );
1753                        return Err(IdentityRuntimeError::Internal(detail));
1754                    }
1755                };
1756                if let Some(reason) = outcome.fallback_reason().cloned() {
1757                    tracing::warn!(
1758                        %identity,
1759                        reason = ?reason,
1760                        "lazy identity materialization fresh-spawned after typed resume fallback"
1761                    );
1762                    self.emit_event(
1763                        identity,
1764                        IdentityEvent::ResumeFallback {
1765                            identity: identity.clone(),
1766                            reason,
1767                        },
1768                    )
1769                    .await;
1770                }
1771                let effective_session_id = outcome.session_id().clone();
1772                if effective_session_id != registered_session_id {
1773                    // §8.4 trigger (b): the resume fallback abandoned the
1774                    // registered session — harvest it detached (materialize
1775                    // is a hot path; the session store read stays valid).
1776                    if let Some(injector) = self.agent_memory.read().await.as_ref() {
1777                        let abandoned_key = registered_session_id.to_string();
1778                        injector.note_session_generation(
1779                            identity,
1780                            &abandoned_key,
1781                            record.generation.get(),
1782                        );
1783                        injector.spawn_rotation_distillation(
1784                            identity,
1785                            &abandoned_key,
1786                            crate::memory::distiller::DistillCause::ResumeFallback,
1787                        );
1788                    }
1789                    abandoned_session_registrations.push(registered_session_id);
1790                }
1791                record.session_id = effective_session_id;
1792            }
1793            record
1794        } else {
1795            let new_runtime_id =
1796                AgentRuntimeId::parse(&format!("rt:{identity}:0")).map_err(|err| {
1797                    IdentityRuntimeError::Internal(format!("failed to mint runtime id: {err}"))
1798                })?;
1799            let mut record = ContinuityRecord {
1800                identity: identity.clone(),
1801                agent_runtime_id: new_runtime_id,
1802                session_id: meerkat_core::types::SessionId::new(),
1803                generation: ContinuityGeneration::new(0),
1804                checkpoint_version: CheckpointVersion::new(0),
1805            };
1806            if let Err(err) = self
1807                .continuity_store
1808                .upsert_continuity_record(&record, grant.fencing_token)
1809                .await
1810            {
1811                let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1812                return Err(IdentityRuntimeError::Internal(format!(
1813                    "continuity upsert before materialize create: {err}{}",
1814                    cleanup_error
1815                        .as_ref()
1816                        .map(|e| format!("; lease cleanup failed: {e}"))
1817                        .unwrap_or_default(),
1818                )));
1819            }
1820            if let Some(bridge) = self.bridge.as_ref() {
1821                let provisional_session_id = record.session_id.clone();
1822                if let Err(err) = bridge
1823                    .register_session_runtime_state(
1824                        &record.session_id,
1825                        identity,
1826                        record.generation,
1827                        record.checkpoint_version,
1828                        grant.fencing_token,
1829                    )
1830                    .await
1831                    .map_err(|err| {
1832                        IdentityRuntimeError::Internal(format!(
1833                            "bridge register_session_runtime_state: {err}"
1834                        ))
1835                    })
1836                {
1837                    let unregister_error = Self::unregister_bridge_session_runtime_states(
1838                        bridge.as_ref(),
1839                        std::slice::from_ref(&provisional_session_id),
1840                    )
1841                    .await;
1842                    let delete_error = self
1843                        .continuity_store
1844                        .delete_continuity_record(identity, grant.fencing_token)
1845                        .await
1846                        .err();
1847                    let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1848                    if let Some(delete_error) = delete_error {
1849                        return Err(IdentityRuntimeError::Internal(format!(
1850                            "{err}{}; tentative continuity cleanup failed: {delete_error}{}",
1851                            unregister_error
1852                                .as_ref()
1853                                .map(|e| format!("; unregister session failed: {e}"))
1854                                .unwrap_or_default(),
1855                            cleanup_error
1856                                .as_ref()
1857                                .map(|e| format!("; lease cleanup failed: {e}"))
1858                                .unwrap_or_default(),
1859                        )));
1860                    }
1861                    if let Some(cleanup_error) = cleanup_error {
1862                        return Err(IdentityRuntimeError::Internal(format!(
1863                            "{err}{}; lease cleanup failed: {cleanup_error}",
1864                            unregister_error
1865                                .as_ref()
1866                                .map(|e| format!("; unregister session failed: {e}"))
1867                                .unwrap_or_default(),
1868                        )));
1869                    }
1870                    if let Some(unregister_error) = unregister_error {
1871                        return Err(IdentityRuntimeError::Internal(format!(
1872                            "{err}; unregister session failed: {unregister_error}"
1873                        )));
1874                    }
1875                    return Err(err);
1876                }
1877                let created_session_id = bridge
1878                    .create_session(
1879                        identity,
1880                        &record.agent_runtime_id,
1881                        &spec,
1882                        &draft,
1883                        &record.session_id,
1884                    )
1885                    .await
1886                    .map_err(|err| {
1887                        IdentityRuntimeError::Internal(format!("bridge create_session: {err}"))
1888                    });
1889                match created_session_id {
1890                    Ok(session_id) => {
1891                        if session_id != provisional_session_id {
1892                            abandoned_session_registrations.push(provisional_session_id);
1893                        }
1894                        record.session_id = session_id;
1895                    }
1896                    Err(err) => {
1897                        let unregister_error = Self::unregister_bridge_session_runtime_states(
1898                            bridge.as_ref(),
1899                            std::slice::from_ref(&provisional_session_id),
1900                        )
1901                        .await;
1902                        let cleanup_error =
1903                            bridge.retire_member(&record.agent_runtime_id).await.err();
1904                        let delete_error = self
1905                            .continuity_store
1906                            .delete_continuity_record(identity, grant.fencing_token)
1907                            .await
1908                            .err();
1909                        let lease_cleanup_error =
1910                            self.release_uninstalled_materialize_lease(&grant).await;
1911                        if unregister_error.is_some()
1912                            || cleanup_error.is_some()
1913                            || delete_error.is_some()
1914                            || lease_cleanup_error.is_some()
1915                        {
1916                            return Err(IdentityRuntimeError::Internal(format!(
1917                                "{err}{}{}{}{}",
1918                                unregister_error
1919                                    .as_ref()
1920                                    .map(|e| format!("; unregister session failed: {e}"))
1921                                    .unwrap_or_default(),
1922                                cleanup_error
1923                                    .as_ref()
1924                                    .map(|e| format!("; cleanup retire failed: {e}"))
1925                                    .unwrap_or_default(),
1926                                delete_error
1927                                    .as_ref()
1928                                    .map(|e| format!("; tentative continuity cleanup failed: {e}"))
1929                                    .unwrap_or_default(),
1930                                lease_cleanup_error
1931                                    .as_ref()
1932                                    .map(|e| format!("; lease cleanup failed: {e}"))
1933                                    .unwrap_or_default(),
1934                            )));
1935                        }
1936                        return Err(err);
1937                    }
1938                }
1939            }
1940            record
1941        };
1942
1943        if let Err(err) = self
1944            .continuity_store
1945            .upsert_continuity_record(&record, grant.fencing_token)
1946            .await
1947        {
1948            let unregister_error = if let Some(bridge) = self.bridge.as_ref() {
1949                let mut sessions_to_unregister = abandoned_session_registrations.clone();
1950                sessions_to_unregister.push(record.session_id.clone());
1951                Self::unregister_bridge_session_runtime_states(
1952                    bridge.as_ref(),
1953                    &sessions_to_unregister,
1954                )
1955                .await
1956            } else {
1957                None
1958            };
1959            let cleanup_error = if let Some(bridge) = self.bridge.as_ref() {
1960                bridge.retire_member(&record.agent_runtime_id).await.err()
1961            } else {
1962                None
1963            };
1964            let restore_error = self
1965                .restore_continuity_after_materialize_failure(
1966                    identity,
1967                    original_continuity.as_ref(),
1968                    &grant,
1969                )
1970                .await;
1971            let lease_cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1972            if unregister_error.is_some()
1973                || cleanup_error.is_some()
1974                || restore_error.is_some()
1975                || lease_cleanup_error.is_some()
1976            {
1977                return Err(IdentityRuntimeError::Internal(format!(
1978                    "continuity upsert after materialize: {err}{}{}{}{}",
1979                    unregister_error
1980                        .as_ref()
1981                        .map(|e| format!("; unregister session failed: {e}"))
1982                        .unwrap_or_default(),
1983                    cleanup_error
1984                        .as_ref()
1985                        .map(|e| format!("; cleanup retire failed: {e}"))
1986                        .unwrap_or_default(),
1987                    restore_error
1988                        .as_ref()
1989                        .map(|e| format!("; continuity rollback failed: {e}"))
1990                        .unwrap_or_default(),
1991                    lease_cleanup_error
1992                        .as_ref()
1993                        .map(|e| format!("; lease cleanup failed: {e}"))
1994                        .unwrap_or_default(),
1995                )));
1996            }
1997            return Err(IdentityRuntimeError::Internal(format!(
1998                "continuity upsert after materialize: {err}"
1999            )));
2000        }
2001        if let Some(bridge) = self.bridge.as_ref() {
2002            let register_result = bridge
2003                .register_session_runtime_state(
2004                    &record.session_id,
2005                    identity,
2006                    record.generation,
2007                    record.checkpoint_version,
2008                    grant.fencing_token,
2009                )
2010                .await;
2011            let effective_checkpoint_version = match register_result {
2012                Ok(version) => version,
2013                Err(err) => {
2014                    let mut sessions_to_unregister = abandoned_session_registrations.clone();
2015                    sessions_to_unregister.push(record.session_id.clone());
2016                    let unregister_error = Self::unregister_bridge_session_runtime_states(
2017                        bridge.as_ref(),
2018                        &sessions_to_unregister,
2019                    )
2020                    .await;
2021                    let cleanup_error = bridge.retire_member(&record.agent_runtime_id).await.err();
2022                    let restore_error = self
2023                        .restore_continuity_after_materialize_failure(
2024                            identity,
2025                            original_continuity.as_ref(),
2026                            &grant,
2027                        )
2028                        .await;
2029                    let lease_cleanup_error =
2030                        self.release_uninstalled_materialize_lease(&grant).await;
2031                    return Err(IdentityRuntimeError::Internal(format!(
2032                        "bridge register actual session runtime state: {err}{}{}{}{}",
2033                        unregister_error
2034                            .as_ref()
2035                            .map(|e| format!("; unregister session failed: {e}"))
2036                            .unwrap_or_default(),
2037                        cleanup_error
2038                            .as_ref()
2039                            .map(|e| format!("; cleanup retire failed: {e}"))
2040                            .unwrap_or_default(),
2041                        restore_error
2042                            .as_ref()
2043                            .map(|e| format!("; continuity rollback failed: {e}"))
2044                            .unwrap_or_default(),
2045                        lease_cleanup_error
2046                            .as_ref()
2047                            .map(|e| format!("; lease cleanup failed: {e}"))
2048                            .unwrap_or_default(),
2049                    )));
2050                }
2051            };
2052            record.checkpoint_version = effective_checkpoint_version;
2053            if let Some(err) = Self::unregister_bridge_session_runtime_states(
2054                bridge.as_ref(),
2055                &abandoned_session_registrations,
2056            )
2057            .await
2058            {
2059                let actual_unregister_error = Self::unregister_bridge_session_runtime_states(
2060                    bridge.as_ref(),
2061                    std::slice::from_ref(&record.session_id),
2062                )
2063                .await;
2064                let unregister_error = actual_unregister_error
2065                    .map(|actual_err| format!("{err}; actual session: {actual_err}"))
2066                    .unwrap_or(err);
2067                let cleanup_error = bridge.retire_member(&record.agent_runtime_id).await.err();
2068                let restore_error = self
2069                    .restore_continuity_after_materialize_failure(
2070                        identity,
2071                        original_continuity.as_ref(),
2072                        &grant,
2073                    )
2074                    .await;
2075                let lease_cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
2076                return Err(IdentityRuntimeError::Internal(format!(
2077                    "bridge unregister abandoned session runtime state: {unregister_error}{}{}{}",
2078                    cleanup_error
2079                        .as_ref()
2080                        .map(|e| format!("; cleanup retire failed: {e}"))
2081                        .unwrap_or_default(),
2082                    restore_error
2083                        .as_ref()
2084                        .map(|e| format!("; continuity rollback failed: {e}"))
2085                        .unwrap_or_default(),
2086                    lease_cleanup_error
2087                        .as_ref()
2088                        .map(|e| format!("; lease cleanup failed: {e}"))
2089                        .unwrap_or_default(),
2090                )));
2091            }
2092        }
2093
2094        {
2095            let mut entries = self.entries.write().await;
2096            let entry = entries
2097                .get_mut(identity)
2098                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2099            entry.continuity = Some(record.clone());
2100            entry.lease = Some(Self::lease_entry_from_grant(&grant));
2101            entry.state = IdentityLifecycleState::Active;
2102            entry.checkpoint_version = record.checkpoint_version;
2103        }
2104        self.emit_event(
2105            identity,
2106            IdentityEvent::StateChanged {
2107                identity: identity.clone(),
2108                new_state: IdentityLifecycleState::Active,
2109            },
2110        )
2111        .await;
2112        let desired_edges = self.desired_peer_edges.read().await.clone();
2113        if !desired_edges.is_empty()
2114            && let Err(err) = self.reconcile_managed_peer_edges(&desired_edges).await
2115        {
2116            tracing::warn!(
2117                identity = %identity,
2118                error = %err,
2119                "identity materialized with topology reconcile warning"
2120            );
2121        }
2122        self.clear_materialization_backoff(identity).await;
2123        Ok(record)
2124    }
2125
2126    async fn best_effort_materialize_identity(
2127        &self,
2128        identity: AgentIdentity,
2129        initiator: Option<&AgentIdentity>,
2130        operation: &'static str,
2131    ) -> Option<ContinuityRecord> {
2132        let attempt_lock = self.best_effort_materialization_lock_for(&identity).await;
2133        let _attempt_guard = attempt_lock.lock().await;
2134
2135        if let Some(error) = self.materialization_backoff_error(&identity).await {
2136            tracing::debug!(
2137                identity = %identity,
2138                initiator = initiator.map(ToString::to_string).as_deref(),
2139                error = %error,
2140                "identity best-effort materialization skipped due to materialization backoff"
2141            );
2142            return None;
2143        }
2144
2145        match self.materialize(&identity).await {
2146            Ok(record) => {
2147                self.clear_materialization_backoff(&identity).await;
2148                Some(record)
2149            }
2150            Err(err) => {
2151                tracing::warn!(
2152                    identity = %identity,
2153                    initiator = initiator.map(ToString::to_string).as_deref(),
2154                    error = %err,
2155                    "identity best-effort materialization skipped identity after materialization failure"
2156                );
2157                self.record_best_effort_materialization_failure(
2158                    &identity, initiator, operation, &err,
2159                )
2160                .await;
2161                None
2162            }
2163        }
2164    }
2165
2166    async fn materialize_all_records(
2167        &self,
2168    ) -> Vec<(
2169        AgentIdentity,
2170        Result<ContinuityRecord, IdentityRuntimeError>,
2171    )> {
2172        let identities = self.registered_identities().await;
2173        stream::iter(identities.into_iter().map(|identity| async move {
2174            let result = self.materialize(&identity).await;
2175            (identity, result)
2176        }))
2177        .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
2178        .collect::<Vec<_>>()
2179        .await
2180    }
2181
2182    /// Materialize all identities currently registered with the runtime.
2183    ///
2184    /// Fleet hydration is best-effort: one member that cannot build is skipped
2185    /// and surfaced through logs/error hooks rather than aborting unrelated
2186    /// members.
2187    pub async fn materialize_all(&self) -> Result<Vec<ContinuityRecord>, IdentityRuntimeError> {
2188        let identities = self.registered_identities().await;
2189        let records = stream::iter(identities.into_iter().map(|identity| async move {
2190            self.best_effort_materialize_identity(identity, None, "materialize_all")
2191                .await
2192        }))
2193        .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
2194        .filter_map(async move |record| record)
2195        .collect::<Vec<_>>()
2196        .await;
2197
2198        let desired_edges = self.desired_peer_edges.read().await.clone();
2199        if !desired_edges.is_empty()
2200            && let Err(err) = self.reconcile_managed_peer_edges(&desired_edges).await
2201        {
2202            tracing::warn!(
2203                error = %err,
2204                "identity materialize_all completed with topology reconcile warning"
2205            );
2206        }
2207
2208        Ok(records)
2209    }
2210
2211    /// Materialize all identities and fail if any registered identity cannot
2212    /// hydrate. Flow admission uses this strict variant so a run is not accepted
2213    /// with a partially materialized identity-first fleet.
2214    pub async fn materialize_all_required(
2215        &self,
2216    ) -> Result<Vec<ContinuityRecord>, IdentityRuntimeError> {
2217        let results = self.materialize_all_records().await;
2218        let mut records = Vec::with_capacity(results.len());
2219        let mut failures = Vec::new();
2220
2221        for (identity, result) in results {
2222            match result {
2223                Ok(record) => records.push(record),
2224                Err(err) => failures.push(format!("{identity}: {err}")),
2225            }
2226        }
2227
2228        if !failures.is_empty() {
2229            return Err(IdentityRuntimeError::Internal(format!(
2230                "identity-first required materialization failed for {} identities: {}",
2231                failures.len(),
2232                failures.join("; ")
2233            )));
2234        }
2235
2236        let desired_edges = self.desired_peer_edges.read().await.clone();
2237        if !desired_edges.is_empty() {
2238            self.reconcile_managed_peer_edges(&desired_edges).await?;
2239        }
2240
2241        Ok(records)
2242    }
2243
2244    pub(crate) async fn best_effort_background_warm_identity(&self, identity: AgentIdentity) {
2245        self.best_effort_materialize_identity(identity, None, "background_warm")
2246            .await;
2247    }
2248
2249    /// Ensure an active identity's desired peer neighborhood exists in the
2250    /// concrete mob graph before ordinary communication starts.
2251    pub async fn materialize_reachable_peers(
2252        &self,
2253        identity: &AgentIdentity,
2254    ) -> Result<Vec<ContinuityRecord>, IdentityRuntimeError> {
2255        let peers = self.reachable_peer_identities(identity).await;
2256        let records = stream::iter(peers.into_iter().map(|peer| async move {
2257            self.best_effort_materialize_identity(
2258                peer,
2259                Some(identity),
2260                "materialize_reachable_peers",
2261            )
2262            .await
2263        }))
2264        .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
2265        .filter_map(async move |record| record)
2266        .collect::<Vec<_>>()
2267        .await;
2268
2269        let desired_edges = self.desired_peer_edges.read().await.clone();
2270        if !desired_edges.is_empty()
2271            && let Err(err) = self.reconcile_managed_peer_edges(&desired_edges).await
2272        {
2273            tracing::warn!(
2274                identity = %identity,
2275                error = %err,
2276                "identity peer materialization completed with topology reconcile warning"
2277            );
2278        }
2279
2280        Ok(records)
2281    }
2282
2283    // -----------------------------------------------------------------------
2284    // Subscribe — REQ-06
2285    // -----------------------------------------------------------------------
2286
2287    /// Subscribe to identity-scoped events.
2288    ///
2289    /// Returns a broadcast receiver that yields `IdentityEvent` items for
2290    /// state changes, lease updates, lease loss, and checkpoint completions.
2291    pub async fn subscribe(
2292        &self,
2293        identity: &AgentIdentity,
2294    ) -> Result<broadcast::Receiver<IdentityEvent>, IdentityRuntimeError> {
2295        let channels = self.event_channels.read().await;
2296        let tx = channels
2297            .get(identity)
2298            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2299        Ok(tx.subscribe())
2300    }
2301
2302    /// Update the spec for an existing identity (used during reconciliation).
2303    pub async fn update_spec(&self, spec: DurableAgentSpec) -> Result<(), IdentityRuntimeError> {
2304        let mut entries = self.entries.write().await;
2305        let entry = entries
2306            .get_mut(&spec.identity)
2307            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(spec.identity.clone()))?;
2308        entry.spec = spec;
2309        Ok(())
2310    }
2311
2312    /// Adopt the roster's CURRENT spec for `identity` into the in-memory entry,
2313    /// so a subsequent [`reset`](Self::reset) rebuilds the regenerated session
2314    /// on the current profile instead of carrying the stored one forward.
2315    ///
2316    /// Best-effort: logs and leaves the stored spec in place if the roster
2317    /// can't be resolved or no longer lists the identity (reset's primary job —
2318    /// the destructive continuity reset — must not fail because the roster
2319    /// provider hiccuped). The runtime stays roster-agnostic; the provider is
2320    /// supplied by the caller (the reset RPC handler), which owns it.
2321    pub async fn adopt_roster_spec(
2322        &self,
2323        roster_provider: &Arc<dyn RosterProvider>,
2324        identity: &AgentIdentity,
2325    ) {
2326        self.adopt_roster_spec_with_context(roster_provider, identity, None)
2327            .await;
2328    }
2329
2330    async fn adopt_roster_spec_with_context(
2331        &self,
2332        roster_provider: &Arc<dyn RosterProvider>,
2333        identity: &AgentIdentity,
2334        mob_definition: Option<meerkat_mob::MobDefinition>,
2335    ) {
2336        match roster_provider
2337            .roster(&RosterContext {
2338                mob_definition,
2339                previous_identities: Vec::new(),
2340            })
2341            .await
2342        {
2343            Ok(specs) => {
2344                if let Some(spec) = specs.into_iter().find(|s| &s.identity == identity)
2345                    && let Err(err) = self.update_spec(spec).await
2346                {
2347                    tracing::warn!(
2348                        identity = %identity,
2349                        error = %err,
2350                        "reset: failed to adopt current roster spec; rebuilding on stored spec",
2351                    );
2352                }
2353            }
2354            Err(err) => {
2355                tracing::warn!(
2356                    identity = %identity,
2357                    error = %err,
2358                    "reset: roster provider failed; rebuilding on stored spec",
2359                );
2360            }
2361        }
2362    }
2363
2364    /// Update the lease for an identity.
2365    pub async fn update_lease(
2366        &self,
2367        identity: &AgentIdentity,
2368        grant: LeaseGrant,
2369    ) -> Result<(), IdentityRuntimeError> {
2370        let fencing_token = grant.fencing_token;
2371        let mut entries = self.entries.write().await;
2372        let entry = entries
2373            .get_mut(identity)
2374            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2375        entry.lease = Some(LeaseEntry {
2376            fencing_token,
2377            ttl: grant.ttl,
2378            acquired_at: Instant::now(),
2379        });
2380        drop(entries);
2381        self.lease_renewal_notify.notify_one();
2382        self.emit_event(
2383            identity,
2384            IdentityEvent::LeaseUpdated {
2385                identity: identity.clone(),
2386                fencing_token,
2387            },
2388        )
2389        .await;
2390        Ok(())
2391    }
2392
2393    /// Mark a lease as lost for an identity (INV-02).
2394    pub async fn mark_lease_lost(
2395        &self,
2396        identity: &AgentIdentity,
2397    ) -> Result<(), IdentityRuntimeError> {
2398        let mut entries = self.entries.write().await;
2399        let entry = entries
2400            .get_mut(identity)
2401            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2402        entry.lease = None;
2403        drop(entries);
2404        self.emit_event(
2405            identity,
2406            IdentityEvent::LeaseLost {
2407                identity: identity.clone(),
2408            },
2409        )
2410        .await;
2411        Ok(())
2412    }
2413
2414    /// Remove an identity from the runtime.
2415    #[allow(dead_code)]
2416    pub(crate) async fn remove(&self, identity: &AgentIdentity) -> Option<IdentityEntry> {
2417        self.event_channels.write().await.remove(identity);
2418        self.entries.write().await.remove(identity)
2419    }
2420
2421    /// Set the lifecycle state for an identity.
2422    pub async fn set_state(
2423        &self,
2424        identity: &AgentIdentity,
2425        state: IdentityLifecycleState,
2426    ) -> Result<(), IdentityRuntimeError> {
2427        let mut entries = self.entries.write().await;
2428        let entry = entries
2429            .get_mut(identity)
2430            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2431        entry.state = state;
2432        drop(entries);
2433        if state == IdentityLifecycleState::Active {
2434            self.lease_renewal_notify.notify_one();
2435        }
2436        self.emit_event(
2437            identity,
2438            IdentityEvent::StateChanged {
2439                identity: identity.clone(),
2440                new_state: state,
2441            },
2442        )
2443        .await;
2444        Ok(())
2445    }
2446
2447    // -----------------------------------------------------------------------
2448    // Lease checking (INV-01, INV-02)
2449    // -----------------------------------------------------------------------
2450
2451    /// Check that the identity has an active, non-expired lease.
2452    /// Returns the fencing token if valid.
2453    fn check_lease(entry: &IdentityEntry) -> Result<FencingToken, IdentityRuntimeError> {
2454        match &entry.lease {
2455            Some(lease) if !lease.is_expired() => Ok(lease.fencing_token),
2456            Some(_) => Err(IdentityRuntimeError::LeaseLost(entry.spec.identity.clone())),
2457            None => Err(IdentityRuntimeError::NoActiveLease(
2458                entry.spec.identity.clone(),
2459            )),
2460        }
2461    }
2462
2463    fn lease_entry_from_grant(grant: &LeaseGrant) -> LeaseEntry {
2464        LeaseEntry {
2465            fencing_token: grant.fencing_token,
2466            ttl: grant.ttl,
2467            acquired_at: Instant::now(),
2468        }
2469    }
2470
2471    async fn ensure_active_lease(
2472        &self,
2473        identity: &AgentIdentity,
2474    ) -> Result<FencingToken, IdentityRuntimeError> {
2475        let (grant, continuity) = {
2476            let entries = self.entries.read().await;
2477            let entry = entries
2478                .get(identity)
2479                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2480            let lease = match &entry.lease {
2481                Some(lease) if lease.is_healthy() => return Ok(lease.fencing_token),
2482                Some(lease) => lease,
2483                None => return Err(IdentityRuntimeError::NoActiveLease(identity.clone())),
2484            };
2485            (
2486                LeaseGrant {
2487                    identity: identity.clone(),
2488                    fencing_token: lease.fencing_token,
2489                    ttl: lease.ttl,
2490                },
2491                entry.continuity.clone(),
2492            )
2493        };
2494
2495        let renewed = self
2496            .lease_provider
2497            .renew_leases(std::slice::from_ref(&grant))
2498            .await
2499            .map_err(IdentityRuntimeError::Lease)?;
2500        let renewed_grant = match renewed.get(identity) {
2501            Some(super::types::LeaseRenewResult::Renewed(grant)) => grant.clone(),
2502            Some(super::types::LeaseRenewResult::Lost { .. }) | None => {
2503                self.mark_lease_lost(identity).await?;
2504                return Err(IdentityRuntimeError::LeaseLost(identity.clone()));
2505            }
2506        };
2507
2508        if let Some(record) = continuity.as_ref() {
2509            self.continuity_store
2510                .upsert_continuity_record(record, renewed_grant.fencing_token)
2511                .await
2512                .map_err(IdentityRuntimeError::Store)?;
2513            if let Some(bridge) = self.bridge.as_ref() {
2514                bridge
2515                    .register_session_runtime_state(
2516                        &record.session_id,
2517                        identity,
2518                        record.generation,
2519                        record.checkpoint_version,
2520                        renewed_grant.fencing_token,
2521                    )
2522                    .await
2523                    .map_err(|err| {
2524                        IdentityRuntimeError::Internal(format!(
2525                            "bridge refresh session runtime state after lease renewal: {err}"
2526                        ))
2527                    })?;
2528            }
2529        }
2530
2531        let fencing_token = renewed_grant.fencing_token;
2532        let mut entries = self.entries.write().await;
2533        let entry = entries
2534            .get_mut(identity)
2535            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2536        match entry.lease.as_ref() {
2537            Some(current) if current.fencing_token == grant.fencing_token => {
2538                entry.lease = Some(Self::lease_entry_from_grant(&renewed_grant));
2539            }
2540            Some(current) => return Ok(current.fencing_token),
2541            None => return Err(IdentityRuntimeError::NoActiveLease(identity.clone())),
2542        }
2543        drop(entries);
2544        self.lease_renewal_notify.notify_one();
2545
2546        self.emit_event(
2547            identity,
2548            IdentityEvent::LeaseUpdated {
2549                identity: identity.clone(),
2550                fencing_token,
2551            },
2552        )
2553        .await;
2554        Ok(fencing_token)
2555    }
2556
2557    async fn mark_lifecycle_in_progress(
2558        &self,
2559        identity: &AgentIdentity,
2560        state: IdentityLifecycleState,
2561    ) -> Result<IdentityEntry, IdentityRuntimeError> {
2562        let mut entries = self.entries.write().await;
2563        let entry = entries
2564            .get_mut(identity)
2565            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2566        let snapshot = entry.clone();
2567        entry.state = state;
2568        entry.lease = None;
2569        Ok(snapshot)
2570    }
2571
2572    async fn restore_entry(&self, identity: &AgentIdentity, entry: IdentityEntry) {
2573        self.entries.write().await.insert(identity.clone(), entry);
2574    }
2575
2576    async fn restore_entry_with_grant(
2577        &self,
2578        identity: &AgentIdentity,
2579        mut entry: IdentityEntry,
2580        grant: &LeaseGrant,
2581    ) {
2582        let restore_live_lease = entry.state == IdentityLifecycleState::Active;
2583        if let Some(record) = entry.continuity.as_ref() {
2584            if let Err(err) = self
2585                .continuity_store
2586                .upsert_continuity_record(record, grant.fencing_token)
2587                .await
2588            {
2589                tracing::warn!(
2590                    %identity,
2591                    error = %err,
2592                    "failed to advance restored continuity fencing token after lifecycle failure"
2593                );
2594                entry.state = IdentityLifecycleState::Broken;
2595            } else if let Some(bridge) = self.bridge.as_ref()
2596                && let Err(err) = bridge
2597                    .register_session_runtime_state(
2598                        &record.session_id,
2599                        identity,
2600                        record.generation,
2601                        record.checkpoint_version,
2602                        grant.fencing_token,
2603                    )
2604                    .await
2605            {
2606                tracing::warn!(
2607                    %identity,
2608                    error = %err,
2609                    "failed to refresh restored session runtime state after lifecycle failure"
2610                );
2611                entry.state = IdentityLifecycleState::Broken;
2612            }
2613            entry.lease = restore_live_lease.then(|| Self::lease_entry_from_grant(grant));
2614        }
2615        self.restore_entry(identity, entry).await;
2616        if restore_live_lease {
2617            self.lease_renewal_notify.notify_one();
2618        }
2619    }
2620
2621    pub(crate) async fn refresh_active_restore_grant(
2622        &self,
2623        identity: &AgentIdentity,
2624        grant: &LeaseGrant,
2625    ) -> Result<(), IdentityRuntimeError> {
2626        let record = {
2627            let entries = self.entries.read().await;
2628            let entry = entries
2629                .get(identity)
2630                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2631            if entry.state != IdentityLifecycleState::Active {
2632                return Err(IdentityRuntimeError::InvalidState {
2633                    identity: identity.clone(),
2634                    state: entry.state,
2635                    operation: "refresh_active_restore_grant",
2636                });
2637            }
2638            entry.continuity.clone()
2639        };
2640
2641        if let Some(record) = record.as_ref()
2642            && let Err(err) = self
2643                .continuity_store
2644                .upsert_continuity_record(record, grant.fencing_token)
2645                .await
2646        {
2647            let mut entries = self.entries.write().await;
2648            if let Some(entry) = entries.get_mut(identity) {
2649                entry.state = IdentityLifecycleState::Broken;
2650                entry.lease = None;
2651            }
2652            drop(entries);
2653            self.emit_event(
2654                identity,
2655                IdentityEvent::StateChanged {
2656                    identity: identity.clone(),
2657                    new_state: IdentityLifecycleState::Broken,
2658                },
2659            )
2660            .await;
2661            return Err(IdentityRuntimeError::Store(err));
2662        }
2663
2664        let mut entries = self.entries.write().await;
2665        let entry = entries
2666            .get_mut(identity)
2667            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2668        if entry.state != IdentityLifecycleState::Active {
2669            return Err(IdentityRuntimeError::InvalidState {
2670                identity: identity.clone(),
2671                state: entry.state,
2672                operation: "refresh_active_restore_grant",
2673            });
2674        }
2675        entry.lease = Some(Self::lease_entry_from_grant(grant));
2676        drop(entries);
2677        self.emit_event(
2678            identity,
2679            IdentityEvent::LeaseUpdated {
2680                identity: identity.clone(),
2681                fencing_token: grant.fencing_token,
2682            },
2683        )
2684        .await;
2685        Ok(())
2686    }
2687
2688    async fn restore_broken_entry_with_fenced_store(
2689        &self,
2690        identity: &AgentIdentity,
2691        mut entry: IdentityEntry,
2692        grant: &LeaseGrant,
2693    ) {
2694        entry.state = IdentityLifecycleState::Broken;
2695        entry.lease = None;
2696        if let Some(record) = entry.continuity.as_ref()
2697            && let Err(err) = self
2698                .continuity_store
2699                .upsert_continuity_record(record, grant.fencing_token)
2700                .await
2701        {
2702            tracing::warn!(
2703                %identity,
2704                error = %err,
2705                "failed to preserve fenced continuity record for broken identity"
2706            );
2707        }
2708        self.restore_entry(identity, entry).await;
2709    }
2710
2711    async fn mark_rebind_failure_broken(
2712        &self,
2713        identity: &AgentIdentity,
2714        mut entry: IdentityEntry,
2715        grant: &LeaseGrant,
2716        rebound_record: &ContinuityRecord,
2717    ) {
2718        entry.state = IdentityLifecycleState::Broken;
2719        entry.lease = None;
2720        entry.checkpoint_version = rebound_record.checkpoint_version;
2721        entry.continuity = Some(rebound_record.clone());
2722        if let Err(err) = self
2723            .continuity_store
2724            .upsert_continuity_record(rebound_record, grant.fencing_token)
2725            .await
2726        {
2727            tracing::warn!(
2728                %identity,
2729                session_id = %rebound_record.session_id,
2730                error = %err,
2731                "failed to preserve rebound continuity after live respawn rebind failure"
2732            );
2733        }
2734        self.restore_entry(identity, entry).await;
2735    }
2736
2737    async fn restore_entry_after_reset_bridge_failure(
2738        &self,
2739        identity: &AgentIdentity,
2740        entry: IdentityEntry,
2741        grant: &LeaseGrant,
2742        force_broken: bool,
2743    ) -> Option<ContinuityStoreError> {
2744        let delete_error = if entry.continuity.is_none() {
2745            self.continuity_store
2746                .delete_continuity_record(identity, grant.fencing_token)
2747                .await
2748                .err()
2749        } else {
2750            None
2751        };
2752        if force_broken || delete_error.is_some() {
2753            self.restore_broken_entry_with_fenced_store(identity, entry, grant)
2754                .await;
2755        } else {
2756            self.restore_entry_with_grant(identity, entry, grant).await;
2757        }
2758        delete_error
2759    }
2760
2761    async fn restore_continuity_after_materialize_failure(
2762        &self,
2763        identity: &AgentIdentity,
2764        previous: Option<&ContinuityRecord>,
2765        grant: &LeaseGrant,
2766    ) -> Option<ContinuityStoreError> {
2767        match previous {
2768            Some(record) => self
2769                .continuity_store
2770                .upsert_continuity_record(record, grant.fencing_token)
2771                .await
2772                .err(),
2773            None => self
2774                .continuity_store
2775                .delete_continuity_record(identity, grant.fencing_token)
2776                .await
2777                .err(),
2778        }
2779    }
2780
2781    async fn unregister_bridge_session_runtime_states(
2782        bridge: &dyn SessionBridge,
2783        session_ids: &[SessionId],
2784    ) -> Option<String> {
2785        let mut errors = Vec::new();
2786        let mut seen = BTreeSet::new();
2787        for session_id in session_ids {
2788            if !seen.insert(session_id.to_string()) {
2789                continue;
2790            }
2791            if let Err(err) = bridge.unregister_session_runtime_state(session_id).await {
2792                errors.push(format!("{session_id}: {err}"));
2793            }
2794        }
2795        (!errors.is_empty()).then(|| errors.join("; "))
2796    }
2797
2798    async fn advance_existing_continuity_fence(
2799        &self,
2800        identity: &AgentIdentity,
2801        entry: &IdentityEntry,
2802        grant: &LeaseGrant,
2803    ) -> Result<(), IdentityRuntimeError> {
2804        if let Some(record) = entry.continuity.as_ref() {
2805            self.continuity_store
2806                .upsert_continuity_record(record, grant.fencing_token)
2807                .await
2808                .map_err(IdentityRuntimeError::Store)?;
2809        }
2810        let _ = identity;
2811        Ok(())
2812    }
2813
2814    async fn refresh_existing_session_runtime_state(
2815        &self,
2816        identity: &AgentIdentity,
2817        record: &ContinuityRecord,
2818        grant: &LeaseGrant,
2819    ) -> Result<CheckpointVersion, IdentityRuntimeError> {
2820        let Some(bridge) = self.bridge.as_ref() else {
2821            return Ok(record.checkpoint_version);
2822        };
2823        bridge
2824            .register_session_runtime_state(
2825                &record.session_id,
2826                identity,
2827                record.generation,
2828                record.checkpoint_version,
2829                grant.fencing_token,
2830            )
2831            .await
2832            .map_err(|err| {
2833                IdentityRuntimeError::Internal(format!(
2834                    "bridge refresh session runtime state: {err}"
2835                ))
2836            })
2837    }
2838
2839    // -----------------------------------------------------------------------
2840    // Delivery: send() — REQ-01, REQ-03
2841    // -----------------------------------------------------------------------
2842
2843    /// Send conversational content to an addressable identity.
2844    ///
2845    /// Enforces:
2846    /// - Identity must be registered and active
2847    /// - Identity must be Addressable (REQ-03)
2848    /// - Lease must be held (INV-01)
2849    /// - Lease must not be lost (INV-02)
2850    ///
2851    /// Returns the fencing token for the delivery (caller uses it for checkpoint).
2852    pub async fn send(
2853        &self,
2854        identity: &AgentIdentity,
2855        content: &meerkat_core::ContentInput,
2856    ) -> Result<FencingToken, IdentityRuntimeError> {
2857        self.send_with_mode(identity, content, HandlingMode::Queue)
2858            .await
2859    }
2860
2861    /// Send conversational content using an explicit turn handling mode.
2862    ///
2863    /// This is the identity-first counterpart to the mob member send path used
2864    /// by the console. Ordinary API callers can keep using [`Self::send`],
2865    /// which preserves queue semantics.
2866    pub async fn send_with_mode(
2867        &self,
2868        identity: &AgentIdentity,
2869        content: &meerkat_core::ContentInput,
2870        handling_mode: HandlingMode,
2871    ) -> Result<FencingToken, IdentityRuntimeError> {
2872        self.send_with_mode_and_interaction(identity, content, handling_mode, None)
2873            .await
2874    }
2875
2876    /// [`Self::send_with_mode`] with a host-minted interaction id (meerkat
2877    /// 0.7.25 ask 15 addendum). The id rides `WorkSpec` into runtime
2878    /// admission, so the turn's live events and its committed transcript
2879    /// messages carry the same identity the console stamped on its frames —
2880    /// the exact live↔history join the console dedup needs. Only UUID-form
2881    /// ids thread; others are delivered without one.
2882    pub async fn send_with_mode_and_interaction(
2883        &self,
2884        identity: &AgentIdentity,
2885        content: &meerkat_core::ContentInput,
2886        handling_mode: HandlingMode,
2887        interaction_id: Option<&str>,
2888    ) -> Result<FencingToken, IdentityRuntimeError> {
2889        let should_materialize = {
2890            let entries = self.entries.read().await;
2891            let entry = entries
2892                .get(identity)
2893                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2894
2895            // REQ-03: reject send to InternalOnly
2896            if entry.spec.addressability == AgentAddressability::InternalOnly {
2897                return Err(IdentityRuntimeError::NotAddressable(NotAddressable {
2898                    identity: identity.clone(),
2899                    addressability: entry.spec.addressability,
2900                }));
2901            }
2902            entry.state == IdentityLifecycleState::Dormant
2903                || entry.state == IdentityLifecycleState::Uninitialized
2904        };
2905        if should_materialize {
2906            self.materialize(identity).await?;
2907        }
2908        // Live steers are latency-sensitive operator input for an already
2909        // active turn. Ordinary sends may hydrate the reachable topology first,
2910        // but a steer must reach the current session boundary before the tool
2911        // turn resumes; background/full-fleet materialization owns the peers.
2912        if handling_mode != HandlingMode::Steer {
2913            self.materialize_reachable_peers(identity).await?;
2914        }
2915
2916        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2917        let _lifecycle_guard = lifecycle_lock.lock().await;
2918        {
2919            let entries = self.entries.read().await;
2920            let entry = entries
2921                .get(identity)
2922                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2923            if entry.state != IdentityLifecycleState::Active {
2924                return Err(IdentityRuntimeError::InvalidState {
2925                    identity: identity.clone(),
2926                    state: entry.state,
2927                    operation: "send",
2928                });
2929            }
2930        }
2931
2932        let mut token = self.ensure_active_lease(identity).await?;
2933        let (runtime_id, memory_session_key, memory_generation) = {
2934            let entries = self.entries.read().await;
2935            let entry = entries
2936                .get(identity)
2937                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2938            (
2939                entry
2940                    .continuity
2941                    .as_ref()
2942                    .map(|c| c.agent_runtime_id.clone()),
2943                // Scopes the injector's cross-turn dedup + cumulative budget.
2944                entry.continuity.as_ref().map(|c| c.session_id.to_string()),
2945                entry.continuity.as_ref().map(|c| c.generation.get()),
2946            )
2947        };
2948        // Steer is latency-sensitive live operator input: it bypasses both
2949        // memory injection and inbound defanging by design. Every other send
2950        // is defanged first (§9.1 anti-spoofing — even with injection off,
2951        // forged memory envelopes are an inbound threat) and only then
2952        // considered for ambient injection.
2953        // Ask 1: the user content and the ambient memory recall travel as
2954        // SEPARATE bodies — `content_to_deliver` is the (defanged) user
2955        // message, `injected_context` is the recall assembled as its own
2956        // typed injected-context body. They are never fused into one text.
2957        let (content_to_deliver, injected_context) = if handling_mode == HandlingMode::Steer {
2958            (content.clone(), Vec::new())
2959        } else {
2960            match self.agent_memory.read().await.clone() {
2961                Some(injector) => {
2962                    // §10.1 taint hook: authoritative session attribution
2963                    // ahead of the async observe stream — the run this send
2964                    // triggers belongs to this session. The generation bind
2965                    // feeds the Distiller's EvidenceRefs (§8.4).
2966                    if let Some(session_key) = memory_session_key.as_deref() {
2967                        injector.note_current_session(identity, session_key);
2968                        if let Some(generation) = memory_generation {
2969                            injector.note_session_generation(identity, session_key, generation);
2970                        }
2971                    }
2972                    let defanged = injector.defang_inbound(identity, content);
2973                    let injected_context = injector
2974                        .inject_for_turn(identity, memory_session_key.as_deref(), &defanged)
2975                        .await
2976                        .map_err(|err| {
2977                            IdentityRuntimeError::Internal(format!("agent memory recall: {err}"))
2978                        })?;
2979                    (defanged, injected_context)
2980                }
2981                None => (content.clone(), Vec::new()),
2982            }
2983        };
2984
2985        // Deliver through the session bridge when available.
2986        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id) {
2987            let delivered_session_id = bridge
2988                .deliver_with_mode_and_context(
2989                    rid,
2990                    &content_to_deliver,
2991                    &injected_context,
2992                    handling_mode,
2993                    interaction_id,
2994                )
2995                .await
2996                .map_err(|e| IdentityRuntimeError::Internal(format!("bridge deliver: {e}")))?;
2997            if let Some(rebound_token) = self
2998                .reconcile_delivered_session_locked(identity, delivered_session_id)
2999                .await?
3000            {
3001                token = rebound_token;
3002            }
3003        }
3004
3005        Ok(token)
3006    }
3007
3008    // -----------------------------------------------------------------------
3009    // Delivery: dispatch() — REQ-02
3010    // -----------------------------------------------------------------------
3011
3012    /// Dispatch internal content to any identity (Addressable or InternalOnly).
3013    ///
3014    /// Enforces:
3015    /// - Identity must be registered and active
3016    /// - Lease must be held (INV-01)
3017    /// - Lease must not be lost (INV-02)
3018    ///
3019    /// Returns (fencing_token, is_durable) where is_durable indicates whether
3020    /// the dispatch is backed by a runtime_store (REQ-04).
3021    pub async fn dispatch(
3022        &self,
3023        identity: &AgentIdentity,
3024        input: &DispatchInput,
3025    ) -> Result<(FencingToken, bool), IdentityRuntimeError> {
3026        let should_materialize = {
3027            let entries = self.entries.read().await;
3028            let entry = entries
3029                .get(identity)
3030                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3031            entry.state == IdentityLifecycleState::Dormant
3032                || entry.state == IdentityLifecycleState::Uninitialized
3033        };
3034        if should_materialize {
3035            self.materialize(identity).await?;
3036        }
3037        self.materialize_reachable_peers(identity).await?;
3038
3039        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3040        let _lifecycle_guard = lifecycle_lock.lock().await;
3041        {
3042            let entries = self.entries.read().await;
3043            let entry = entries
3044                .get(identity)
3045                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3046            if entry.state != IdentityLifecycleState::Active {
3047                return Err(IdentityRuntimeError::InvalidState {
3048                    identity: identity.clone(),
3049                    state: entry.state,
3050                    operation: "dispatch",
3051                });
3052            }
3053        }
3054
3055        let mut token = self.ensure_active_lease(identity).await?;
3056        let (is_durable, runtime_id) = {
3057            let entries = self.entries.read().await;
3058            let entry = entries
3059                .get(identity)
3060                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3061            // REQ-04: durability depends on runtime_store
3062            let is_durable = entry.has_runtime_store;
3063
3064            let runtime_id = entry
3065                .continuity
3066                .as_ref()
3067                .map(|c| c.agent_runtime_id.clone());
3068
3069            (is_durable, runtime_id)
3070        };
3071
3072        // Deliver through the session bridge when available.
3073        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id) {
3074            let delivered_session_id = bridge
3075                .deliver(rid, &input.content)
3076                .await
3077                .map_err(|e| IdentityRuntimeError::Internal(format!("bridge dispatch: {e}")))?;
3078            if let Some(rebound_token) = self
3079                .reconcile_delivered_session_locked(identity, delivered_session_id)
3080                .await?
3081            {
3082                token = rebound_token;
3083            }
3084        }
3085
3086        Ok((token, is_durable))
3087    }
3088
3089    // -----------------------------------------------------------------------
3090    // Status: status() — REQ-07
3091    // -----------------------------------------------------------------------
3092
3093    /// Return the full identity status for the given identity.
3094    pub async fn status(
3095        &self,
3096        identity: &AgentIdentity,
3097    ) -> Result<IdentityStatus, IdentityRuntimeError> {
3098        let entries = self.entries.read().await;
3099        let entry = entries
3100            .get(identity)
3101            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3102
3103        let lease_info = entry.lease.as_ref().map(|l| LeaseInfo {
3104            fencing_token: l.fencing_token,
3105            ttl_remaining: l.ttl_remaining(),
3106            healthy: l.is_healthy(),
3107        });
3108
3109        let continuity_health = Some(ContinuityHealth {
3110            store_reachable: true, // tracked per-store in production
3111            durability_policy: self.durability_policy.clone(),
3112            last_checkpoint_version: if entry.checkpoint_version.get() > 0 {
3113                Some(entry.checkpoint_version)
3114            } else {
3115                None
3116            },
3117        });
3118
3119        Ok(IdentityStatus {
3120            identity: identity.clone(),
3121            state: entry.state,
3122            agent_runtime_id: entry
3123                .continuity
3124                .as_ref()
3125                .map(|c| c.agent_runtime_id.clone()),
3126            session_id: entry.continuity.as_ref().map(|c| c.session_id.clone()),
3127            profile: Some(entry.spec.profile.clone()),
3128            runtime_mode: entry.spec.runtime_mode_override,
3129            addressability: entry.spec.addressability,
3130            display_name: entry.spec.display_name.clone(),
3131            labels: entry.spec.labels.clone(),
3132            generation: entry.continuity.as_ref().map(|c| c.generation),
3133            checkpoint_version: if entry.checkpoint_version.get() > 0 {
3134                Some(entry.checkpoint_version)
3135            } else {
3136                None
3137            },
3138            lease: lease_info,
3139            continuity_health,
3140        })
3141    }
3142
3143    /// Return statuses for every registered identity without materializing
3144    /// dormant members.
3145    pub async fn statuses(&self) -> Vec<IdentityStatus> {
3146        let identities = self
3147            .entries
3148            .read()
3149            .await
3150            .keys()
3151            .cloned()
3152            .collect::<Vec<_>>();
3153        let mut statuses = Vec::with_capacity(identities.len());
3154        for identity in identities {
3155            if let Ok(status) = self.status(&identity).await {
3156                statuses.push(status);
3157            }
3158        }
3159        statuses
3160    }
3161
3162    // -----------------------------------------------------------------------
3163    // Lifecycle: retire() — REQ-08
3164    // -----------------------------------------------------------------------
3165
3166    /// Retire an identity. Validates lease ownership and retires the mob member.
3167    pub async fn retire(
3168        &self,
3169        identity: &AgentIdentity,
3170    ) -> Result<FencingToken, IdentityRuntimeError> {
3171        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3172        let _lifecycle_guard = lifecycle_lock.lock().await;
3173        self.ensure_active_lease(identity).await?;
3174        let registered_entry = self
3175            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Retiring)
3176            .await?;
3177        let _previous_token = match Self::check_lease(&registered_entry) {
3178            Ok(token) => token,
3179            Err(err) => {
3180                self.restore_entry(identity, registered_entry).await;
3181                return Err(err);
3182            }
3183        };
3184        let runtime_id = registered_entry
3185            .continuity
3186            .as_ref()
3187            .map(|c| c.agent_runtime_id.clone());
3188        let session_id = registered_entry
3189            .continuity
3190            .as_ref()
3191            .map(|c| c.session_id.clone());
3192
3193        let acquire_result = match self
3194            .lease_provider
3195            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
3196            .await
3197        {
3198            Ok(result) => result,
3199            Err(err) => {
3200                self.restore_entry(identity, registered_entry).await;
3201                return Err(IdentityRuntimeError::Lease(err));
3202            }
3203        };
3204
3205        let grant = match acquire_result.get(identity) {
3206            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
3207            _ => {
3208                self.restore_entry(identity, registered_entry).await;
3209                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
3210            }
3211        };
3212        if let Err(err) = self
3213            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
3214            .await
3215        {
3216            let mut broken_entry = registered_entry;
3217            broken_entry.state = IdentityLifecycleState::Broken;
3218            broken_entry.lease = None;
3219            self.restore_entry(identity, broken_entry).await;
3220            return Err(err);
3221        }
3222
3223        // §8.4 trigger (b): distill the outgoing session's tail BEFORE the
3224        // member retires. Best-effort and bounded — retirement proceeds at
3225        // the distiller's pre-rotation timeout.
3226        if let Some(injector) = self.agent_memory.read().await.clone() {
3227            if let Some(session_id) = session_id.as_ref() {
3228                injector
3229                    .distill_before_rotation(
3230                        identity,
3231                        &session_id.to_string(),
3232                        crate::memory::distiller::DistillCause::Retire,
3233                    )
3234                    .await;
3235            }
3236            // §8.5 exit interview: queue the retired identity's store for
3237            // the next dream's harvest sub-phase.
3238            injector
3239                .note_identity_retired(
3240                    identity,
3241                    session_id
3242                        .as_ref()
3243                        .map(std::string::ToString::to_string)
3244                        .as_deref(),
3245                    "retire",
3246                )
3247                .await;
3248        }
3249
3250        // Retire the mob member through the session bridge when available.
3251        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id)
3252            && let Err(err) = bridge.retire_member(rid).await
3253        {
3254            self.restore_entry_with_grant(identity, registered_entry, &grant)
3255                .await;
3256            return Err(IdentityRuntimeError::Internal(format!(
3257                "bridge retire: {err}"
3258            )));
3259        }
3260        if let (Some(bridge), Some(session_id)) = (&self.bridge, &session_id)
3261            && let Err(err) = bridge.unregister_session_runtime_state(session_id).await
3262        {
3263            let mut broken_entry = registered_entry;
3264            broken_entry.state = IdentityLifecycleState::Broken;
3265            broken_entry.lease = None;
3266            self.restore_entry(identity, broken_entry).await;
3267            return Err(IdentityRuntimeError::Internal(format!(
3268                "bridge unregister retired session: {err}"
3269            )));
3270        }
3271
3272        Ok(grant.fencing_token)
3273    }
3274
3275    // -----------------------------------------------------------------------
3276    // Lifecycle: respawn() — REQ-09
3277    // -----------------------------------------------------------------------
3278
3279    /// Respawn: non-destructive recovery.
3280    ///
3281    /// 1. Fence the current owner
3282    /// 2. Attempt final checkpoint
3283    /// 3. Reactivate from authoritative continuity with same record + runtime ID
3284    /// 4. ContinuityGeneration does NOT advance
3285    pub async fn respawn(
3286        &self,
3287        identity: &AgentIdentity,
3288    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
3289        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3290        let _lifecycle_guard = lifecycle_lock.lock().await;
3291        let registered_entry = self
3292            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Suspended)
3293            .await?;
3294
3295        // Fence the old owner by re-acquiring the lease
3296        let acquire_result = match self
3297            .lease_provider
3298            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
3299            .await
3300        {
3301            Ok(result) => result,
3302            Err(err) => {
3303                self.restore_entry(identity, registered_entry).await;
3304                return Err(IdentityRuntimeError::Lease(err));
3305            }
3306        };
3307
3308        let grant = match acquire_result.get(identity) {
3309            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
3310            _ => {
3311                self.restore_entry(identity, registered_entry).await;
3312                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
3313            }
3314        };
3315
3316        if let Err(err) = self
3317            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
3318            .await
3319        {
3320            let mut broken_entry = registered_entry;
3321            broken_entry.state = IdentityLifecycleState::Broken;
3322            broken_entry.lease = None;
3323            self.restore_entry(identity, broken_entry).await;
3324            return Err(err);
3325        }
3326
3327        // Resolve current continuity state
3328        let resolved = match self
3329            .continuity_store
3330            .resolve_many(std::slice::from_ref(identity))
3331            .await
3332        {
3333            Ok(resolved) => resolved,
3334            Err(err) => {
3335                self.restore_entry_with_grant(identity, registered_entry.clone(), &grant)
3336                    .await;
3337                return Err(IdentityRuntimeError::Store(err));
3338            }
3339        };
3340
3341        let record = match resolved.get(identity) {
3342            Some(super::types::ContinuityResolveState::Ready { record }) => record.clone(),
3343            Some(super::types::ContinuityResolveState::Broken { failure }) => {
3344                self.restore_entry_with_grant(identity, registered_entry, &grant)
3345                    .await;
3346                return Err(IdentityRuntimeError::Internal(format!(
3347                    "broken continuity for {identity}: {}",
3348                    failure.detail
3349                )));
3350            }
3351            Some(super::types::ContinuityResolveState::Uninitialized) => {
3352                self.restore_entry_with_grant(identity, registered_entry, &grant)
3353                    .await;
3354                return Err(IdentityRuntimeError::Internal(format!(
3355                    "cannot respawn uninitialized identity {identity}"
3356                )));
3357            }
3358            None => {
3359                self.restore_entry_with_grant(identity, registered_entry, &grant)
3360                    .await;
3361                return Err(IdentityRuntimeError::Store(
3362                    ContinuityStoreError::NotFound {
3363                        identity: identity.clone(),
3364                    },
3365                ));
3366            }
3367        };
3368
3369        // §8.4 trigger (b): respawn is a recovery boundary — harvest the
3370        // session's window before the runtime refreshes (the SessionId does
3371        // not rotate here; the cursor stays valid). Bounded; respawn
3372        // proceeds at the pre-rotation timeout.
3373        if let Some(injector) = self.agent_memory.read().await.clone() {
3374            let session_key = record.session_id.to_string();
3375            injector.note_session_generation(identity, &session_key, record.generation.get());
3376            injector
3377                .distill_before_rotation(
3378                    identity,
3379                    &session_key,
3380                    crate::memory::distiller::DistillCause::Respawn,
3381                )
3382                .await;
3383        }
3384
3385        let effective_checkpoint_version = match self
3386            .refresh_existing_session_runtime_state(identity, &record, &grant)
3387            .await
3388        {
3389            Ok(version) => version,
3390            Err(err) => {
3391                self.restore_entry_with_grant(identity, registered_entry, &grant)
3392                    .await;
3393                return Err(err);
3394            }
3395        };
3396        let mut record = record;
3397        record.checkpoint_version = effective_checkpoint_version;
3398
3399        // Update runtime state: same record, new lease, back to Active
3400        let mut entries = self.entries.write().await;
3401        if let Some(entry) = entries.get_mut(identity) {
3402            entry.continuity = Some(record.clone());
3403            entry.lease = Some(LeaseEntry {
3404                fencing_token: grant.fencing_token,
3405                ttl: grant.ttl,
3406                acquired_at: Instant::now(),
3407            });
3408            entry.state = IdentityLifecycleState::Active;
3409            entry.checkpoint_version = record.checkpoint_version;
3410        }
3411
3412        Ok(record)
3413    }
3414
3415    /// Rebind continuity to the concrete session created by a lower-level
3416    /// member respawn. This keeps identity-first status aligned when a control
3417    /// surface refreshes the mob member outside the identity runtime bridge.
3418    pub async fn rebind_session_after_live_respawn(
3419        &self,
3420        identity: &AgentIdentity,
3421        session_id: SessionId,
3422    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
3423        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3424        let _lifecycle_guard = lifecycle_lock.lock().await;
3425        self.rebind_session_after_live_respawn_locked(identity, session_id)
3426            .await
3427    }
3428
3429    async fn reconcile_delivered_session_locked(
3430        &self,
3431        identity: &AgentIdentity,
3432        delivered_session_id: SessionId,
3433    ) -> Result<Option<FencingToken>, IdentityRuntimeError> {
3434        let current_session_id = {
3435            let entries = self.entries.read().await;
3436            let entry = entries
3437                .get(identity)
3438                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3439            entry
3440                .continuity
3441                .as_ref()
3442                .map(|record| record.session_id.clone())
3443        };
3444
3445        let Some(current_session_id) = current_session_id else {
3446            return Ok(None);
3447        };
3448        if current_session_id == delivered_session_id {
3449            return Ok(None);
3450        }
3451
3452        tracing::warn!(
3453            %identity,
3454            old_session_id = %current_session_id,
3455            new_session_id = %delivered_session_id,
3456            "identity bridge delivery returned a rotated session; rebinding continuity"
3457        );
3458        self.rebind_session_after_live_respawn_locked(identity, delivered_session_id)
3459            .await?;
3460
3461        let entries = self.entries.read().await;
3462        let entry = entries
3463            .get(identity)
3464            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3465        Ok(entry.lease.as_ref().map(|lease| lease.fencing_token))
3466    }
3467
3468    async fn rebind_session_after_live_respawn_locked(
3469        &self,
3470        identity: &AgentIdentity,
3471        session_id: SessionId,
3472    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
3473        let registered_entry = self
3474            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Suspended)
3475            .await?;
3476
3477        let acquire_result = match self
3478            .lease_provider
3479            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
3480            .await
3481        {
3482            Ok(result) => result,
3483            Err(err) => {
3484                self.restore_entry(identity, registered_entry).await;
3485                return Err(IdentityRuntimeError::Lease(err));
3486            }
3487        };
3488
3489        let grant = match acquire_result.get(identity) {
3490            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
3491            _ => {
3492                self.restore_entry(identity, registered_entry).await;
3493                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
3494            }
3495        };
3496
3497        let mut record = match registered_entry.continuity.as_ref() {
3498            Some(record) => record.clone(),
3499            None => {
3500                self.restore_entry_with_grant(identity, registered_entry, &grant)
3501                    .await;
3502                return Err(IdentityRuntimeError::UnknownIdentity(identity.clone()));
3503            }
3504        };
3505        if let Err(err) = self
3506            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
3507            .await
3508        {
3509            if let Some(bridge) = self.bridge.as_ref()
3510                && let Err(unregister_err) =
3511                    bridge.unregister_session_runtime_state(&session_id).await
3512            {
3513                tracing::warn!(
3514                    %identity,
3515                    session_id = %session_id,
3516                    error = %unregister_err,
3517                    "failed to unregister rebound session after continuity fence failure"
3518                );
3519            }
3520            let mut broken_entry = registered_entry;
3521            broken_entry.state = IdentityLifecycleState::Broken;
3522            broken_entry.lease = None;
3523            self.restore_entry(identity, broken_entry).await;
3524            return Err(err);
3525        }
3526        let previous_session_id = record.session_id.clone();
3527        record.session_id = session_id;
3528
3529        if let Err(err) = self
3530            .continuity_store
3531            .upsert_continuity_record(&record, grant.fencing_token)
3532            .await
3533        {
3534            if let Some(bridge) = self.bridge.as_ref()
3535                && let Err(unregister_err) = bridge
3536                    .unregister_session_runtime_state(&record.session_id)
3537                    .await
3538            {
3539                tracing::warn!(
3540                    %identity,
3541                    session_id = %record.session_id,
3542                    error = %unregister_err,
3543                    "failed to unregister rebound session after continuity upsert failure"
3544                );
3545            }
3546            self.mark_rebind_failure_broken(identity, registered_entry, &grant, &record)
3547                .await;
3548            return Err(IdentityRuntimeError::Store(err));
3549        }
3550
3551        if let Some(bridge) = self.bridge.as_ref() {
3552            match bridge
3553                .register_session_runtime_state(
3554                    &record.session_id,
3555                    identity,
3556                    record.generation,
3557                    record.checkpoint_version,
3558                    grant.fencing_token,
3559                )
3560                .await
3561            {
3562                Ok(version) => record.checkpoint_version = version,
3563                Err(err) => {
3564                    if let Err(unregister_err) = bridge
3565                        .unregister_session_runtime_state(&record.session_id)
3566                        .await
3567                    {
3568                        tracing::warn!(
3569                            %identity,
3570                            session_id = %record.session_id,
3571                            error = %unregister_err,
3572                            "failed to unregister rebound session after bridge register failure"
3573                        );
3574                    }
3575                    self.mark_rebind_failure_broken(identity, registered_entry, &grant, &record)
3576                        .await;
3577                    return Err(IdentityRuntimeError::Internal(format!(
3578                        "bridge rebind respawned session runtime state: {err}"
3579                    )));
3580                }
3581            }
3582            if previous_session_id != record.session_id
3583                && let Err(err) = bridge
3584                    .unregister_session_runtime_state(&previous_session_id)
3585                    .await
3586            {
3587                tracing::warn!(
3588                    %identity,
3589                    session_id = %previous_session_id,
3590                    error = %err,
3591                    "failed to unregister previous session after live respawn rebind"
3592                );
3593            }
3594        }
3595
3596        if let Err(err) = self
3597            .continuity_store
3598            .upsert_continuity_record(&record, grant.fencing_token)
3599            .await
3600        {
3601            if let Some(bridge) = self.bridge.as_ref()
3602                && let Err(unregister_err) = bridge
3603                    .unregister_session_runtime_state(&record.session_id)
3604                    .await
3605            {
3606                tracing::warn!(
3607                    %identity,
3608                    session_id = %record.session_id,
3609                    error = %unregister_err,
3610                    "failed to unregister rebound session after final continuity upsert failure"
3611                );
3612            }
3613            self.mark_rebind_failure_broken(identity, registered_entry, &grant, &record)
3614                .await;
3615            return Err(IdentityRuntimeError::Store(err));
3616        }
3617
3618        self.register(
3619            registered_entry.spec,
3620            IdentityLifecycleState::Active,
3621            Some(record.clone()),
3622            Some(grant),
3623        )
3624        .await;
3625        Ok(record)
3626    }
3627
3628    // -----------------------------------------------------------------------
3629    // Lifecycle: reset() — REQ-10
3630    // -----------------------------------------------------------------------
3631
3632    /// Reset: destructive continuity reset.
3633    ///
3634    /// 1. Fence old owner
3635    /// 2. Advance ContinuityGeneration
3636    /// 3. Create fresh continuity under the same AgentIdentity
3637    /// 4. Old-owner late writes rejected by stale fencing token
3638    pub async fn reset(
3639        &self,
3640        identity: &AgentIdentity,
3641    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
3642        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3643        let _lifecycle_guard = lifecycle_lock.lock().await;
3644        // Re-profile before snapshotting the lifecycle entry so the rebuilt
3645        // session uses the current roster spec, not the old checkpoint spec.
3646        self.adopt_current_roster_spec_for_reset(identity).await;
3647        let registered_entry = self
3648            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Suspended)
3649            .await?;
3650
3651        // INV-05: fence the old owner first
3652        let acquire_result = match self
3653            .lease_provider
3654            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
3655            .await
3656        {
3657            Ok(result) => result,
3658            Err(err) => {
3659                self.restore_entry(identity, registered_entry).await;
3660                return Err(IdentityRuntimeError::Lease(err));
3661            }
3662        };
3663
3664        let grant = match acquire_result.get(identity) {
3665            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
3666            _ => {
3667                self.restore_entry(identity, registered_entry).await;
3668                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
3669            }
3670        };
3671        if let Err(err) = self
3672            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
3673            .await
3674        {
3675            let mut broken_entry = registered_entry;
3676            broken_entry.state = IdentityLifecycleState::Broken;
3677            broken_entry.lease = None;
3678            self.restore_entry(identity, broken_entry).await;
3679            return Err(err);
3680        }
3681
3682        // Resolve to get current generation
3683        let resolved = match self
3684            .continuity_store
3685            .resolve_many(std::slice::from_ref(identity))
3686            .await
3687        {
3688            Ok(resolved) => resolved,
3689            Err(err) => {
3690                self.restore_entry_with_grant(identity, registered_entry.clone(), &grant)
3691                    .await;
3692                return Err(IdentityRuntimeError::Store(err));
3693            }
3694        };
3695
3696        let current_gen = match resolved.get(identity) {
3697            Some(super::types::ContinuityResolveState::Ready { record }) => record.generation,
3698            Some(super::types::ContinuityResolveState::Uninitialized) => {
3699                ContinuityGeneration::new(0)
3700            }
3701            _ => ContinuityGeneration::new(0),
3702        };
3703
3704        // Advance generation
3705        let new_gen = ContinuityGeneration::new(current_gen.get() + 1);
3706        let new_session_id = meerkat_core::types::SessionId::new();
3707        let new_runtime_id = AgentRuntimeId::parse(&format!("rt:{identity}:{}", new_gen.get()))
3708            .map_err(|e| {
3709                IdentityRuntimeError::Internal(format!("failed to mint runtime id: {e}"))
3710            })?;
3711
3712        let new_record = ContinuityRecord {
3713            identity: identity.clone(),
3714            agent_runtime_id: new_runtime_id,
3715            session_id: new_session_id,
3716            generation: new_gen,
3717            checkpoint_version: CheckpointVersion::new(0),
3718        };
3719        let spec = registered_entry.spec.clone();
3720        let mut draft = super::types::AgentBuildDraft {
3721            model: None,
3722            system_prompt: None,
3723            additional_instructions: spec.additional_instructions.clone(),
3724            labels: spec.labels.clone(),
3725            app_context: spec.context.clone(),
3726            external_tools: Vec::new(),
3727            local_external_tools: Default::default(),
3728        };
3729        if self.bridge.is_some() {
3730            let active_peers = self.entries.read().await.keys().cloned().collect();
3731            let managed_edges = self.desired_peer_edges.read().await.clone();
3732            let build_context = AgentBuildContext {
3733                identity: identity.clone(),
3734                active_peers,
3735                managed_edges,
3736                runtime_services: self.runtime_services(),
3737            };
3738            if let Some(customizer) = self.customizer.read().await.clone()
3739                && let Err(err) = customizer
3740                    .customize_build(&build_context, &spec, &mut draft)
3741                    .await
3742            {
3743                self.restore_entry_with_grant(identity, registered_entry, &grant)
3744                    .await;
3745                return Err(IdentityRuntimeError::Internal(format!(
3746                    "customizer after reset: {err}"
3747                )));
3748            }
3749        }
3750
3751        // Bridge: retire old mob member and create fresh session for the new identity.
3752        if let Some(bridge) = &self.bridge {
3753            if let Err(err) = self
3754                .continuity_store
3755                .upsert_continuity_record(&new_record, grant.fencing_token)
3756                .await
3757            {
3758                self.restore_entry_with_grant(identity, registered_entry, &grant)
3759                    .await;
3760                return Err(IdentityRuntimeError::Store(err));
3761            }
3762
3763            let old_runtime_id = registered_entry
3764                .continuity
3765                .as_ref()
3766                .map(|c| c.agent_runtime_id.clone());
3767            let old_session_id = registered_entry
3768                .continuity
3769                .as_ref()
3770                .map(|c| c.session_id.clone());
3771
3772            let session_id = bridge
3773                .create_session(
3774                    identity,
3775                    &new_record.agent_runtime_id,
3776                    &spec,
3777                    &draft,
3778                    &new_record.session_id,
3779                )
3780                .await
3781                .map_err(|e| {
3782                    IdentityRuntimeError::Internal(format!(
3783                        "bridge create_session after reset: {e}"
3784                    ))
3785                });
3786            let session_id = match session_id {
3787                Ok(session_id) => session_id,
3788                Err(err) => {
3789                    let cleanup_error = bridge
3790                        .retire_member(&new_record.agent_runtime_id)
3791                        .await
3792                        .err();
3793                    let delete_error = self
3794                        .restore_entry_after_reset_bridge_failure(
3795                            identity,
3796                            registered_entry.clone(),
3797                            &grant,
3798                            cleanup_error.is_some(),
3799                        )
3800                        .await;
3801                    if cleanup_error.is_some() || delete_error.is_some() {
3802                        return Err(IdentityRuntimeError::Internal(format!(
3803                            "{err}{}{}",
3804                            cleanup_error
3805                                .as_ref()
3806                                .map(|e| format!("; cleanup retire failed: {e}"))
3807                                .unwrap_or_default(),
3808                            delete_error
3809                                .as_ref()
3810                                .map(|e| format!("; tentative continuity cleanup failed: {e}"))
3811                                .unwrap_or_default()
3812                        )));
3813                    }
3814                    return Err(err);
3815                }
3816            };
3817            // Update the record with the actual session ID
3818            let mut new_record = new_record;
3819            new_record.session_id = session_id;
3820            tracing::debug!(
3821                identity = %identity,
3822                runtime_id = %new_record.agent_runtime_id,
3823                session_id = %new_record.session_id,
3824                "reset bridge create_session completed",
3825            );
3826
3827            if let Err(err) = self
3828                .continuity_store
3829                .upsert_continuity_record(&new_record, grant.fencing_token)
3830                .await
3831            {
3832                let unregister_error = Self::unregister_bridge_session_runtime_states(
3833                    bridge.as_ref(),
3834                    std::slice::from_ref(&new_record.session_id),
3835                )
3836                .await;
3837                let cleanup_error = bridge
3838                    .retire_member(&new_record.agent_runtime_id)
3839                    .await
3840                    .err();
3841                let delete_error = self
3842                    .restore_entry_after_reset_bridge_failure(
3843                        identity,
3844                        registered_entry.clone(),
3845                        &grant,
3846                        unregister_error.is_some() || cleanup_error.is_some(),
3847                    )
3848                    .await;
3849                if unregister_error.is_some() || cleanup_error.is_some() || delete_error.is_some() {
3850                    return Err(IdentityRuntimeError::Internal(format!(
3851                        "continuity upsert actual session after reset: {err}{}{}{}",
3852                        unregister_error
3853                            .as_ref()
3854                            .map(|e| format!("; unregister session failed: {e}"))
3855                            .unwrap_or_default(),
3856                        cleanup_error
3857                            .as_ref()
3858                            .map(|e| format!("; cleanup retire failed: {e}"))
3859                            .unwrap_or_default(),
3860                        delete_error
3861                            .as_ref()
3862                            .map(|e| format!("; tentative continuity cleanup failed: {e}"))
3863                            .unwrap_or_default(),
3864                    )));
3865                }
3866                return Err(IdentityRuntimeError::Store(err));
3867            }
3868
3869            let register_result = bridge
3870                .register_session_runtime_state(
3871                    &new_record.session_id,
3872                    identity,
3873                    new_record.generation,
3874                    new_record.checkpoint_version,
3875                    grant.fencing_token,
3876                )
3877                .await;
3878            let effective_checkpoint_version = match register_result {
3879                Ok(version) => version,
3880                Err(err) => {
3881                    let unregister_error = Self::unregister_bridge_session_runtime_states(
3882                        bridge.as_ref(),
3883                        std::slice::from_ref(&new_record.session_id),
3884                    )
3885                    .await;
3886                    let cleanup_error = bridge
3887                        .retire_member(&new_record.agent_runtime_id)
3888                        .await
3889                        .err();
3890                    let mut detail =
3891                        format!("bridge register actual session runtime state after reset: {err}");
3892                    if let Some(unregister_error) = unregister_error.as_ref() {
3893                        detail
3894                            .push_str(&format!("; unregister session failed: {unregister_error}"));
3895                    }
3896                    if let Some(cleanup_error) = cleanup_error.as_ref() {
3897                        detail.push_str(&format!("; cleanup retire failed: {cleanup_error}"));
3898                    }
3899                    let delete_error = self
3900                        .restore_entry_after_reset_bridge_failure(
3901                            identity,
3902                            registered_entry.clone(),
3903                            &grant,
3904                            unregister_error.is_some() || cleanup_error.is_some(),
3905                        )
3906                        .await;
3907                    if let Some(delete_error) = delete_error {
3908                        return Err(IdentityRuntimeError::Internal(format!(
3909                            "{detail}; tentative continuity cleanup failed: {delete_error}"
3910                        )));
3911                    }
3912                    return Err(IdentityRuntimeError::Internal(detail));
3913                }
3914            };
3915            new_record.checkpoint_version = effective_checkpoint_version;
3916            tracing::debug!(
3917                identity = %identity,
3918                runtime_id = %new_record.agent_runtime_id,
3919                session_id = %new_record.session_id,
3920                checkpoint_version = new_record.checkpoint_version.get(),
3921                "reset bridge session runtime state registered",
3922            );
3923
3924            let cleanup_old_runtime_id = old_runtime_id
3925                .as_ref()
3926                .filter(|old_id| *old_id != &new_record.agent_runtime_id)
3927                .cloned();
3928            let cleanup_old_session_id = old_session_id
3929                .as_ref()
3930                .filter(|old_session_id| *old_session_id != &new_record.session_id)
3931                .cloned();
3932            self.spawn_old_bridge_cleanup_after_reset(
3933                bridge.clone(),
3934                cleanup_old_runtime_id,
3935                cleanup_old_session_id,
3936            );
3937            tracing::debug!(
3938                identity = %identity,
3939                runtime_id = %new_record.agent_runtime_id,
3940                session_id = %new_record.session_id,
3941                "reset old bridge cleanup scheduled",
3942            );
3943
3944            if let Err(err) = self
3945                .continuity_store
3946                .upsert_continuity_record(&new_record, grant.fencing_token)
3947                .await
3948            {
3949                tracing::warn!(
3950                    identity = %identity,
3951                    runtime_id = %new_record.agent_runtime_id,
3952                    session_id = %new_record.session_id,
3953                    error = %err,
3954                    "reset final continuity upsert failed after bridge materialization; rolling back new generation",
3955                );
3956                let unregister_error = bridge
3957                    .unregister_session_runtime_state(&new_record.session_id)
3958                    .await
3959                    .err();
3960                let cleanup_error = bridge
3961                    .retire_member(&new_record.agent_runtime_id)
3962                    .await
3963                    .err();
3964                let rollback_error = self
3965                    .restore_continuity_after_materialize_failure(
3966                        identity,
3967                        registered_entry.continuity.as_ref(),
3968                        &grant,
3969                    )
3970                    .await;
3971                let mut entries = self.entries.write().await;
3972                if let Some(entry) = entries.get_mut(identity) {
3973                    entry.state = IdentityLifecycleState::Broken;
3974                }
3975                if unregister_error.is_some() || cleanup_error.is_some() || rollback_error.is_some()
3976                {
3977                    return Err(IdentityRuntimeError::Internal(format!(
3978                        "continuity upsert after reset: {err}{}{}{}",
3979                        unregister_error
3980                            .as_ref()
3981                            .map(|e| format!("; unregister session failed: {e}"))
3982                            .unwrap_or_default(),
3983                        cleanup_error
3984                            .as_ref()
3985                            .map(|e| format!("; cleanup retire failed: {e}"))
3986                            .unwrap_or_default(),
3987                        rollback_error
3988                            .as_ref()
3989                            .map(|e| format!("; continuity rollback failed: {e}"))
3990                            .unwrap_or_default()
3991                    )));
3992                }
3993                return Err(IdentityRuntimeError::Store(err));
3994            }
3995
3996            // Update runtime state
3997            let mut entries = self.entries.write().await;
3998            let Some(entry) = entries.get_mut(identity) else {
3999                tracing::warn!(
4000                    identity = %identity,
4001                    runtime_id = %new_record.agent_runtime_id,
4002                    session_id = %new_record.session_id,
4003                    "reset entry disappeared after bridge materialization; rolling back new generation",
4004                );
4005                let _ = bridge
4006                    .unregister_session_runtime_state(&new_record.session_id)
4007                    .await;
4008                let _ = bridge.retire_member(&new_record.agent_runtime_id).await;
4009                if registered_entry.continuity.is_none() {
4010                    let _ = self
4011                        .continuity_store
4012                        .delete_continuity_record(identity, grant.fencing_token)
4013                        .await;
4014                }
4015                return Err(IdentityRuntimeError::UnknownIdentity(identity.clone()));
4016            };
4017            entry.continuity = Some(new_record.clone());
4018            entry.lease = Some(Self::lease_entry_from_grant(&grant));
4019            entry.state = IdentityLifecycleState::Active;
4020            entry.checkpoint_version = new_record.checkpoint_version;
4021            tracing::debug!(
4022                identity = %identity,
4023                runtime_id = %new_record.agent_runtime_id,
4024                session_id = %new_record.session_id,
4025                "reset completed",
4026            );
4027            drop(entries);
4028            // §10.1: reset is the deliberate clean-slate boundary — clear
4029            // session taint explicitly (rotation clears implicitly; this
4030            // also drops pending pre-attribution taint). §8.4: distill the
4031            // outgoing session DETACHED (never on the reset critical path;
4032            // the session store outlives the member, so the read stays
4033            // valid after teardown) with the reset boundary marked first so
4034            // every distillate lands Quarantined pending steward review.
4035            if let Some(injector) = self.agent_memory.read().await.as_ref() {
4036                injector.clear_taint_for_identity(identity);
4037                injector.note_session_generation(
4038                    identity,
4039                    &new_record.session_id.to_string(),
4040                    new_record.generation.get(),
4041                );
4042                if let Some(old_continuity) = registered_entry.continuity.as_ref() {
4043                    let old_session_key = old_continuity.session_id.to_string();
4044                    injector.note_reset_boundary(&old_session_key);
4045                    injector.note_session_generation(
4046                        identity,
4047                        &old_session_key,
4048                        old_continuity.generation.get(),
4049                    );
4050                    injector.spawn_rotation_distillation(
4051                        identity,
4052                        &old_session_key,
4053                        crate::memory::distiller::DistillCause::Reset,
4054                    );
4055                }
4056            }
4057            return Ok(new_record);
4058        }
4059
4060        // Persist the new record (fencing token from new lease protects against old writes)
4061        if let Err(err) = self
4062            .continuity_store
4063            .upsert_continuity_record(&new_record, grant.fencing_token)
4064            .await
4065        {
4066            self.restore_entry_with_grant(identity, registered_entry, &grant)
4067                .await;
4068            return Err(IdentityRuntimeError::Store(err));
4069        }
4070
4071        // No bridge — update runtime state only (validation mode)
4072        let mut entries = self.entries.write().await;
4073        let entry = entries
4074            .get_mut(identity)
4075            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
4076        entry.continuity = Some(new_record.clone());
4077        entry.lease = Some(Self::lease_entry_from_grant(&grant));
4078        entry.state = IdentityLifecycleState::Active;
4079        entry.checkpoint_version = CheckpointVersion::new(0);
4080        drop(entries);
4081        if let Some(injector) = self.agent_memory.read().await.as_ref() {
4082            injector.clear_taint_for_identity(identity);
4083            // No-bridge (validation) reset: same §8.4 boundary semantics.
4084            if let Some(old_continuity) = registered_entry.continuity.as_ref() {
4085                let old_session_key = old_continuity.session_id.to_string();
4086                injector.note_reset_boundary(&old_session_key);
4087                injector.spawn_rotation_distillation(
4088                    identity,
4089                    &old_session_key,
4090                    crate::memory::distiller::DistillCause::Reset,
4091                );
4092            }
4093        }
4094
4095        Ok(new_record)
4096    }
4097
4098    // -----------------------------------------------------------------------
4099    // Lifecycle: delete_identity() — REQ-11
4100    // -----------------------------------------------------------------------
4101
4102    /// Delete an identity: removes continuity record.
4103    ///
4104    /// 1. Fence old owner
4105    /// 2. Remove ContinuityRecord
4106    /// 3. Future bootstrap treats identity as Uninitialized
4107    pub async fn delete_identity(
4108        &self,
4109        identity: &AgentIdentity,
4110    ) -> Result<(), IdentityRuntimeError> {
4111        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
4112        let _lifecycle_guard = lifecycle_lock.lock().await;
4113        let registered_entry = self
4114            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Retiring)
4115            .await?;
4116        let runtime_id = registered_entry
4117            .continuity
4118            .as_ref()
4119            .map(|c| c.agent_runtime_id.clone());
4120        let session_id = registered_entry
4121            .continuity
4122            .as_ref()
4123            .map(|c| c.session_id.clone());
4124
4125        // INV-05: fence the old owner first
4126        let acquire_result = match self
4127            .lease_provider
4128            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
4129            .await
4130        {
4131            Ok(result) => result,
4132            Err(err) => {
4133                self.restore_entry(identity, registered_entry).await;
4134                return Err(IdentityRuntimeError::Lease(err));
4135            }
4136        };
4137
4138        let grant = match acquire_result.get(identity) {
4139            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
4140            _ => {
4141                self.restore_entry(identity, registered_entry).await;
4142                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
4143            }
4144        };
4145        if let Err(err) = self
4146            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
4147            .await
4148        {
4149            let mut broken_entry = registered_entry;
4150            broken_entry.state = IdentityLifecycleState::Broken;
4151            broken_entry.lease = None;
4152            self.restore_entry(identity, broken_entry).await;
4153            return Err(err);
4154        }
4155
4156        // §8.4 trigger (b): delete is the identity's LAST boundary — harvest
4157        // the outgoing session before teardown (its exit-interview analog,
4158        // §8.5, is the steward's; the distillate is what it will read).
4159        // Bounded; deletion proceeds at the pre-rotation timeout.
4160        if let Some(injector) = self.agent_memory.read().await.clone() {
4161            if let Some(session_id) = session_id.as_ref() {
4162                let session_key = session_id.to_string();
4163                injector
4164                    .distill_before_rotation(
4165                        identity,
4166                        &session_key,
4167                        crate::memory::distiller::DistillCause::Delete,
4168                    )
4169                    .await;
4170                // Ask 2 GC: delete permanently abandons this session id, and
4171                // its knowledge has just been distilled into the identity-keyed
4172                // MobKit store (which the exit-interview harvest below reads —
4173                // a DIFFERENT store from the meerkat session-memory scope). So
4174                // reclaiming the orphaned meerkat scope here frees dead
4175                // re-embed weight without starving any downstream read.
4176                injector
4177                    .drop_orphaned_session_scope(
4178                        &session_key,
4179                        crate::memory::distiller::DistillCause::Delete,
4180                    )
4181                    .await;
4182            }
4183            // §8.5 exit interview (delete is the identity's LAST boundary).
4184            injector
4185                .note_identity_retired(
4186                    identity,
4187                    session_id
4188                        .as_ref()
4189                        .map(std::string::ToString::to_string)
4190                        .as_deref(),
4191                    "delete",
4192                )
4193                .await;
4194        }
4195
4196        // Retire the mob member through the session bridge before removing
4197        // the continuity record. This ensures the mob actor is cleaned up.
4198        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id)
4199            && let Err(err) = bridge.retire_member(rid).await
4200        {
4201            self.restore_entry_with_grant(identity, registered_entry, &grant)
4202                .await;
4203            return Err(IdentityRuntimeError::Internal(format!(
4204                "bridge retire before delete: {err}"
4205            )));
4206        }
4207
4208        if let (Some(bridge), Some(session_id)) = (&self.bridge, &session_id)
4209            && let Some(err) = Self::unregister_bridge_session_runtime_states(
4210                bridge.as_ref(),
4211                std::slice::from_ref(session_id),
4212            )
4213            .await
4214        {
4215            self.restore_broken_entry_with_fenced_store(identity, registered_entry, &grant)
4216                .await;
4217            return Err(IdentityRuntimeError::Internal(format!(
4218                "bridge unregister session before delete: {err}"
4219            )));
4220        }
4221
4222        // Remove authoritative continuity record from the store
4223        if let Err(err) = self
4224            .continuity_store
4225            .delete_continuity_record(identity, grant.fencing_token)
4226            .await
4227        {
4228            let mut entries = self.entries.write().await;
4229            if let Some(entry) = entries.get_mut(identity) {
4230                entry.state = IdentityLifecycleState::Broken;
4231            }
4232            return Err(IdentityRuntimeError::Store(err));
4233        }
4234
4235        // Remove from runtime tracking
4236        self.event_channels.write().await.remove(identity);
4237        self.entries.write().await.remove(identity);
4238
4239        Ok(())
4240    }
4241
4242    // -----------------------------------------------------------------------
4243    // Checkpoint — REQ-14, REQ-15, REQ-16, REQ-17
4244    // -----------------------------------------------------------------------
4245
4246    /// Save a checkpoint snapshot. Enforces version ordering and fencing.
4247    pub async fn checkpoint(
4248        &self,
4249        identity: &AgentIdentity,
4250        snapshot: &SessionSnapshot,
4251    ) -> Result<CheckpointVersion, IdentityRuntimeError> {
4252        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
4253        let _lifecycle_guard = lifecycle_lock.lock().await;
4254        {
4255            let entries = self.entries.read().await;
4256            let entry = entries
4257                .get(identity)
4258                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
4259            if entry.state != IdentityLifecycleState::Active {
4260                return Err(IdentityRuntimeError::InvalidState {
4261                    identity: identity.clone(),
4262                    state: entry.state,
4263                    operation: "checkpoint",
4264                });
4265            }
4266        }
4267
4268        let token = self.ensure_active_lease(identity).await?;
4269        let (record, new_version) = {
4270            let entries = self.entries.read().await;
4271            let entry = entries
4272                .get(identity)
4273                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
4274            let record = entry
4275                .continuity
4276                .as_ref()
4277                .ok_or_else(|| {
4278                    IdentityRuntimeError::Internal(format!("no continuity record for {identity}"))
4279                })?
4280                .clone();
4281
4282            let new_version = CheckpointVersion::new(entry.checkpoint_version.get() + 1);
4283            (record, new_version)
4284        };
4285
4286        // REQ-15 + REQ-16: store enforces version ordering and fencing
4287        self.continuity_store
4288            .save_session_snapshot(
4289                identity,
4290                &record.session_id,
4291                record.generation,
4292                new_version,
4293                token,
4294                snapshot,
4295            )
4296            .await?;
4297
4298        // Update local checkpoint version
4299        {
4300            let mut entries = self.entries.write().await;
4301            if let Some(entry) = entries.get_mut(identity) {
4302                entry.checkpoint_version = new_version;
4303            }
4304        }
4305
4306        self.emit_event(
4307            identity,
4308            IdentityEvent::CheckpointCompleted {
4309                identity: identity.clone(),
4310                version: new_version,
4311            },
4312        )
4313        .await;
4314
4315        Ok(new_version)
4316    }
4317
4318    // -----------------------------------------------------------------------
4319    // Roster inspection — REQ-32
4320    // -----------------------------------------------------------------------
4321
4322    /// Return all active identities with their specs and status.
4323    pub async fn roster_inspect(
4324        &self,
4325    ) -> BTreeMap<AgentIdentity, (DurableAgentSpec, IdentityStatus)> {
4326        let entries = self.entries.read().await;
4327        let mut result = BTreeMap::new();
4328        for (identity, entry) in entries.iter() {
4329            let lease_info = entry.lease.as_ref().map(|l| LeaseInfo {
4330                fencing_token: l.fencing_token,
4331                ttl_remaining: l.ttl_remaining(),
4332                healthy: l.is_healthy(),
4333            });
4334            let continuity_health = Some(ContinuityHealth {
4335                store_reachable: true,
4336                durability_policy: self.durability_policy.clone(),
4337                last_checkpoint_version: if entry.checkpoint_version.get() > 0 {
4338                    Some(entry.checkpoint_version)
4339                } else {
4340                    None
4341                },
4342            });
4343            let status = IdentityStatus {
4344                identity: identity.clone(),
4345                state: entry.state,
4346                agent_runtime_id: entry
4347                    .continuity
4348                    .as_ref()
4349                    .map(|c| c.agent_runtime_id.clone()),
4350                session_id: entry.continuity.as_ref().map(|c| c.session_id.clone()),
4351                profile: Some(entry.spec.profile.clone()),
4352                runtime_mode: entry.spec.runtime_mode_override,
4353                addressability: entry.spec.addressability,
4354                display_name: entry.spec.display_name.clone(),
4355                labels: entry.spec.labels.clone(),
4356                generation: entry.continuity.as_ref().map(|c| c.generation),
4357                checkpoint_version: if entry.checkpoint_version.get() > 0 {
4358                    Some(entry.checkpoint_version)
4359                } else {
4360                    None
4361                },
4362                lease: lease_info,
4363                continuity_health,
4364            };
4365            result.insert(identity.clone(), (entry.spec.clone(), status));
4366        }
4367        result
4368    }
4369
4370    // -----------------------------------------------------------------------
4371    // Roster uniqueness validation — INV-06
4372    // -----------------------------------------------------------------------
4373
4374    /// Validate that a roster contains no duplicate identities.
4375    pub fn validate_roster_uniqueness(
4376        specs: &[DurableAgentSpec],
4377    ) -> Result<(), IdentityRuntimeError> {
4378        let mut seen = std::collections::BTreeSet::new();
4379        for spec in specs {
4380            if !seen.insert(&spec.identity) {
4381                return Err(IdentityRuntimeError::DuplicateIdentity(
4382                    spec.identity.clone(),
4383                ));
4384            }
4385        }
4386        Ok(())
4387    }
4388
4389    // -----------------------------------------------------------------------
4390    // Accessors for internal state
4391    // -----------------------------------------------------------------------
4392
4393    /// Get the current entries (read-only snapshot).
4394    #[allow(dead_code)]
4395    pub(crate) async fn entries(&self) -> BTreeMap<AgentIdentity, IdentityEntry> {
4396        self.entries.read().await.clone()
4397    }
4398
4399    /// Check if an identity is registered.
4400    pub async fn contains(&self, identity: &AgentIdentity) -> bool {
4401        self.entries.read().await.contains_key(identity)
4402    }
4403
4404    /// Check if an identity is registered AND in Active state.
4405    pub async fn is_active(&self, identity: &AgentIdentity) -> bool {
4406        self.entries
4407            .read()
4408            .await
4409            .get(identity)
4410            .is_some_and(|e| e.state == IdentityLifecycleState::Active)
4411    }
4412
4413    /// Resolve a member alias to the durable identity that OWNS it, if any.
4414    ///
4415    /// Accepts both a generated runtime alias (`rt:<identity>:<generation>`)
4416    /// and a plain identity string; returns the identity only when it is
4417    /// actually registered in this runtime. Member-scoped RPCs use this to
4418    /// route lifecycle mutations of identity-owned members through the
4419    /// identity authority — a classic `handle.retire()`/`respawn()` on such
4420    /// a member would mutate it behind the IdentityRuntime's back (stale
4421    /// continuity binding, generation drift), which is the doctrine's
4422    /// "mob plane must not mangle durables" rule.
4423    pub async fn owned_identity_for_member_alias(&self, alias: &str) -> Option<AgentIdentity> {
4424        let identity = alias
4425            .strip_prefix("rt:")
4426            .and_then(|rest| rest.rsplit_once(':'))
4427            .filter(|(identity, generation)| {
4428                !identity.is_empty()
4429                    && !generation.is_empty()
4430                    && generation.chars().all(|ch| ch.is_ascii_digit())
4431            })
4432            .and_then(|(identity, _)| AgentIdentity::parse(identity).ok())
4433            .or_else(|| AgentIdentity::parse(alias).ok())?;
4434        self.contains(&identity).await.then_some(identity)
4435    }
4436
4437    /// Identities currently in the Broken lifecycle state.
4438    pub async fn broken_identities(&self) -> Vec<AgentIdentity> {
4439        self.entries
4440            .read()
4441            .await
4442            .iter()
4443            .filter(|(_, entry)| entry.state == IdentityLifecycleState::Broken)
4444            .map(|(identity, _)| identity.clone())
4445            .collect()
4446    }
4447
4448    /// Get the continuity store reference.
4449    pub fn continuity_store(&self) -> &Arc<dyn ContinuityStore> {
4450        &self.continuity_store
4451    }
4452
4453    /// Get the lease provider reference.
4454    pub fn lease_provider(&self) -> &Arc<dyn LeaseProvider> {
4455        &self.lease_provider
4456    }
4457
4458    /// Get the runtime instance ID.
4459    pub fn runtime_instance_id(&self) -> &str {
4460        &self.runtime_instance_id
4461    }
4462
4463    /// Get the durability policy.
4464    pub fn durability_policy(&self) -> &DurabilityPolicy {
4465        &self.durability_policy
4466    }
4467
4468    /// Get whether a runtime store is configured.
4469    pub fn has_runtime_store(&self) -> bool {
4470        self.has_runtime_store
4471    }
4472
4473    /// Get the session bridge reference, if configured.
4474    pub fn bridge(&self) -> Option<&Arc<dyn SessionBridge>> {
4475        self.bridge.as_ref()
4476    }
4477
4478    // -----------------------------------------------------------------------
4479    // Convenience methods
4480    // -----------------------------------------------------------------------
4481
4482    /// Send plain text to an addressable identity.
4483    pub async fn send_text(
4484        &self,
4485        identity: &AgentIdentity,
4486        text: impl Into<String>,
4487    ) -> Result<FencingToken, IdentityRuntimeError> {
4488        self.send(identity, &meerkat_core::ContentInput::Text(text.into()))
4489            .await
4490    }
4491
4492    /// Dispatch plain text with system origin.
4493    pub async fn dispatch_text(
4494        &self,
4495        identity: &AgentIdentity,
4496        text: impl Into<String>,
4497    ) -> Result<(FencingToken, bool), IdentityRuntimeError> {
4498        self.dispatch(identity, &DispatchInput::system(text)).await
4499    }
4500
4501    /// Execute the restore flow for the given roster.
4502    pub async fn restore_flow(
4503        &self,
4504        roster: &[DurableAgentSpec],
4505        topology_provider: Option<&dyn super::contracts::TopologyProvider>,
4506        customizer: Option<&dyn super::contracts::AgentCustomizer>,
4507    ) -> Result<super::orchestrator::RestoreFlowResult, IdentityRuntimeError> {
4508        super::orchestrator::restore_flow(self, roster, topology_provider, customizer).await
4509    }
4510
4511    /// Resolve the AgentRuntimeId for a registered identity.
4512    pub async fn runtime_id_for(
4513        &self,
4514        identity: &AgentIdentity,
4515    ) -> Result<AgentRuntimeId, IdentityRuntimeError> {
4516        let entries = self.entries.read().await;
4517        let entry = entries
4518            .get(identity)
4519            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
4520        entry
4521            .continuity
4522            .as_ref()
4523            .map(|c| c.agent_runtime_id.clone())
4524            .ok_or_else(|| {
4525                IdentityRuntimeError::Internal(format!("no continuity record for {identity}"))
4526            })
4527    }
4528
4529    /// Inspect the current execution state of an identity via the bridge.
4530    pub async fn inspect(
4531        &self,
4532        identity: &AgentIdentity,
4533    ) -> Result<super::bridge::MemberInspection, IdentityRuntimeError> {
4534        let runtime_id = self.runtime_id_for(identity).await?;
4535        let bridge = self
4536            .bridge
4537            .as_ref()
4538            .ok_or_else(|| IdentityRuntimeError::Internal("no bridge configured".to_string()))?;
4539        bridge
4540            .inspect_member(&runtime_id)
4541            .await
4542            .map_err(|e| IdentityRuntimeError::Internal(format!("inspect: {e}")))
4543    }
4544
4545    /// The configured default timeout for wait operations.
4546    pub fn default_timeout(&self) -> Duration {
4547        self.default_timeout
4548    }
4549
4550    fn spawn_old_bridge_cleanup_after_reset(
4551        &self,
4552        bridge: Arc<dyn SessionBridge>,
4553        old_runtime_id: Option<AgentRuntimeId>,
4554        old_session_id: Option<SessionId>,
4555    ) {
4556        if old_runtime_id.is_none() && old_session_id.is_none() {
4557            return;
4558        }
4559        let runtime_instance_id = self.runtime_instance_id.clone();
4560        let timeout = self.default_timeout;
4561        tokio::spawn(async move {
4562            if let Some(old_runtime_id) = old_runtime_id {
4563                tracing::debug!(
4564                    runtime_instance_id = %runtime_instance_id,
4565                    runtime_id = %old_runtime_id,
4566                    "skipping old bridge member retire after reset; reset commits the new generation and only clears stale session projection",
4567                );
4568            }
4569
4570            if let Some(old_session_id) = old_session_id {
4571                match tokio::time::timeout(
4572                    timeout,
4573                    bridge.unregister_session_runtime_state(&old_session_id),
4574                )
4575                .await
4576                {
4577                    Ok(Ok(())) => {}
4578                    Ok(Err(err)) => {
4579                        tracing::warn!(
4580                            runtime_instance_id = %runtime_instance_id,
4581                            session_id = %old_session_id,
4582                            error = %err,
4583                            "failed to unregister old bridge session after reset; continuing with new generation",
4584                        );
4585                    }
4586                    Err(_) => {
4587                        tracing::warn!(
4588                            runtime_instance_id = %runtime_instance_id,
4589                            session_id = %old_session_id,
4590                            timeout_ms = timeout.as_millis(),
4591                            "timed out unregistering old bridge session after reset; continuing with new generation",
4592                        );
4593                    }
4594                }
4595            }
4596        });
4597    }
4598
4599    /// Poll until the identity produces an output_preview, or timeout.
4600    pub async fn wait_for_output(
4601        &self,
4602        identity: &AgentIdentity,
4603        timeout: Duration,
4604    ) -> Result<String, IdentityRuntimeError> {
4605        let deadline = Instant::now() + timeout;
4606        loop {
4607            if let Ok(inspection) = self.inspect(identity).await
4608                && let Some(preview) = inspection.output_preview
4609            {
4610                return Ok(preview);
4611            }
4612            if Instant::now() >= deadline {
4613                return Err(IdentityRuntimeError::Internal(format!(
4614                    "timed out waiting for output from {identity}"
4615                )));
4616            }
4617            tokio::time::sleep(Duration::from_millis(500)).await;
4618        }
4619    }
4620
4621    /// Poll until output_preview contains the given substring, or timeout.
4622    pub async fn wait_for_output_containing(
4623        &self,
4624        identity: &AgentIdentity,
4625        needle: &str,
4626        timeout: Duration,
4627    ) -> Result<String, IdentityRuntimeError> {
4628        let deadline = Instant::now() + timeout;
4629        loop {
4630            if let Ok(inspection) = self.inspect(identity).await
4631                && let Some(ref preview) = inspection.output_preview
4632                && preview.contains(needle)
4633            {
4634                return Ok(preview.clone());
4635            }
4636            if Instant::now() >= deadline {
4637                return Err(IdentityRuntimeError::Internal(format!(
4638                    "timed out waiting for output containing '{needle}' from {identity}"
4639                )));
4640            }
4641            tokio::time::sleep(Duration::from_millis(500)).await;
4642        }
4643    }
4644}
4645
4646/// Wire two identities across mobs, resolving runtime IDs from both IdentityRuntimes.
4647///
4648/// This is a convenience function for same-process cross-mob scenarios where both
4649/// IdentityRuntimes are available. It resolves the AgentRuntimeId for each identity
4650/// and delegates to `UnifiedRuntime::wire_cross_mob()`.
4651pub async fn wire_cross_mob_by_identity(
4652    local_irt: &IdentityRuntime,
4653    local_identity: &AgentIdentity,
4654    remote_irt: &IdentityRuntime,
4655    remote_identity: &AgentIdentity,
4656    local_unified: &crate::UnifiedRuntime,
4657    remote_mob_id: &str,
4658) -> Result<(), IdentityRuntimeError> {
4659    let local_rt = local_irt.runtime_id_for(local_identity).await?;
4660    let remote_rt = remote_irt.runtime_id_for(remote_identity).await?;
4661    Box::pin(local_unified.wire_cross_mob(local_rt.as_str(), remote_rt.as_str(), remote_mob_id))
4662        .await
4663        .map_err(|e| IdentityRuntimeError::Internal(format!("wire_cross_mob: {e}")))
4664}
4665
4666#[cfg(test)]
4667mod reset_reprofile_tests {
4668    use super::*;
4669    use std::sync::Arc;
4670    use std::sync::atomic::{AtomicUsize, Ordering};
4671    use tokio::sync::{Mutex as AsyncMutex, RwLock as AsyncRwLock};
4672
4673    use super::super::bridge::{BridgeError, MemberInspection, ResumeSessionOutcome};
4674    use super::super::contracts::RosterProvider;
4675    use super::super::local_lease::LocalLeaseProvider;
4676    use super::super::local_store::LocalContinuityStore;
4677    use super::super::types::{AgentBuildDraft, RosterError, SessionSnapshot};
4678
4679    struct MutableRoster {
4680        specs: AsyncRwLock<Vec<DurableAgentSpec>>,
4681    }
4682
4683    impl MutableRoster {
4684        fn new(specs: Vec<DurableAgentSpec>) -> Self {
4685            Self {
4686                specs: AsyncRwLock::new(specs),
4687            }
4688        }
4689
4690        async fn set(&self, specs: Vec<DurableAgentSpec>) {
4691            *self.specs.write().await = specs;
4692        }
4693    }
4694
4695    #[async_trait::async_trait]
4696    impl RosterProvider for MutableRoster {
4697        async fn roster(
4698            &self,
4699            _context: &RosterContext,
4700        ) -> Result<Vec<DurableAgentSpec>, RosterError> {
4701            Ok(self.specs.read().await.clone())
4702        }
4703    }
4704
4705    #[derive(Default)]
4706    struct RecordingBridge {
4707        create_profiles: AsyncMutex<Vec<String>>,
4708        create_delay: Duration,
4709        creates_in_flight: AtomicUsize,
4710        max_creates_in_flight: AtomicUsize,
4711        retired_runtime_ids: AsyncMutex<Vec<String>>,
4712        hanging_retire_runtime_ids: AsyncMutex<BTreeSet<String>>,
4713        failing_unregister_session_ids: AsyncMutex<BTreeSet<String>>,
4714    }
4715
4716    impl RecordingBridge {
4717        async fn create_profiles(&self) -> Vec<String> {
4718            self.create_profiles.lock().await.clone()
4719        }
4720
4721        fn max_creates_in_flight(&self) -> usize {
4722            self.max_creates_in_flight.load(Ordering::SeqCst)
4723        }
4724
4725        async fn retired_runtime_ids(&self) -> Vec<String> {
4726            self.retired_runtime_ids.lock().await.clone()
4727        }
4728
4729        async fn hang_retire_for(&self, runtime_id: &AgentRuntimeId) {
4730            self.hanging_retire_runtime_ids
4731                .lock()
4732                .await
4733                .insert(runtime_id.to_string());
4734        }
4735
4736        async fn fail_unregister_for(&self, session_id: &SessionId) {
4737            self.failing_unregister_session_ids
4738                .lock()
4739                .await
4740                .insert(session_id.to_string());
4741        }
4742    }
4743
4744    #[async_trait::async_trait]
4745    impl SessionBridge for RecordingBridge {
4746        async fn create_session(
4747            &self,
4748            _identity: &AgentIdentity,
4749            _runtime_id: &AgentRuntimeId,
4750            spec: &DurableAgentSpec,
4751            _draft: &AgentBuildDraft,
4752            session_id: &SessionId,
4753        ) -> Result<SessionId, BridgeError> {
4754            let in_flight = self.creates_in_flight.fetch_add(1, Ordering::SeqCst) + 1;
4755            self.max_creates_in_flight
4756                .fetch_max(in_flight, Ordering::SeqCst);
4757            tokio::time::sleep(self.create_delay).await;
4758            self.creates_in_flight.fetch_sub(1, Ordering::SeqCst);
4759            self.create_profiles
4760                .lock()
4761                .await
4762                .push(spec.profile.to_string());
4763            Ok(session_id.clone())
4764        }
4765
4766        async fn resume_session(
4767            &self,
4768            _identity: &AgentIdentity,
4769            _runtime_id: &AgentRuntimeId,
4770            _spec: &DurableAgentSpec,
4771            _draft: &AgentBuildDraft,
4772            _session_id: &SessionId,
4773            _snapshot: &SessionSnapshot,
4774        ) -> Result<ResumeSessionOutcome, BridgeError> {
4775            Err(BridgeError::Mob(
4776                "resume not used in reset test".to_string(),
4777            ))
4778        }
4779
4780        async fn deliver(
4781            &self,
4782            _runtime_id: &AgentRuntimeId,
4783            _content: &meerkat_core::ContentInput,
4784        ) -> Result<SessionId, BridgeError> {
4785            Err(BridgeError::Mob(
4786                "deliver not used in reset test".to_string(),
4787            ))
4788        }
4789
4790        async fn checkpoint_session(
4791            &self,
4792            _runtime_id: &AgentRuntimeId,
4793            _session_id: &SessionId,
4794        ) -> Result<SessionSnapshot, BridgeError> {
4795            Err(BridgeError::Mob(
4796                "checkpoint not used in reset test".to_string(),
4797            ))
4798        }
4799
4800        async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError> {
4801            self.retired_runtime_ids
4802                .lock()
4803                .await
4804                .push(runtime_id.to_string());
4805            if self
4806                .hanging_retire_runtime_ids
4807                .lock()
4808                .await
4809                .contains(runtime_id.as_str())
4810            {
4811                futures::future::pending::<()>().await;
4812            }
4813            Ok(())
4814        }
4815
4816        async fn inspect_member(
4817            &self,
4818            _runtime_id: &AgentRuntimeId,
4819        ) -> Result<MemberInspection, BridgeError> {
4820            Err(BridgeError::Mob(
4821                "inspect not used in reset test".to_string(),
4822            ))
4823        }
4824
4825        async fn unregister_session_runtime_state(
4826            &self,
4827            session_id: &SessionId,
4828        ) -> Result<(), BridgeError> {
4829            if self
4830                .failing_unregister_session_ids
4831                .lock()
4832                .await
4833                .contains(&session_id.to_string())
4834            {
4835                return Err(BridgeError::Mob("old session still draining".to_string()));
4836            }
4837            Ok(())
4838        }
4839    }
4840
4841    fn durable_spec(identity: AgentIdentity, profile: &str) -> DurableAgentSpec {
4842        DurableAgentSpec {
4843            identity,
4844            profile: meerkat_mob::ProfileName::from(profile),
4845            addressability: AgentAddressability::Addressable,
4846            display_name: None,
4847            labels: BTreeMap::new(),
4848            context: None,
4849            additional_instructions: Vec::new(),
4850            initial_message: None,
4851            runtime_mode_override: None,
4852            backend: None,
4853            binding: None,
4854        }
4855    }
4856
4857    #[tokio::test]
4858    async fn restore_flow_bounds_parallel_member_creation() -> Result<(), Box<dyn std::error::Error>>
4859    {
4860        let specs = (0..8)
4861            .map(|index| {
4862                Ok(durable_spec(
4863                    AgentIdentity::parse(&format!("domain:restore-{index}"))?,
4864                    "domain",
4865                ))
4866            })
4867            .collect::<Result<Vec<_>, Box<dyn std::error::Error>>>()?;
4868        let bridge = Arc::new(RecordingBridge {
4869            create_delay: Duration::from_millis(25),
4870            ..RecordingBridge::default()
4871        });
4872        let runtime = IdentityRuntime::new(IdentityRuntimeConfig {
4873            continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
4874            lease_provider: Arc::new(LocalLeaseProvider::new()),
4875            runtime_instance_id: "parallel-restore-test".to_string(),
4876            has_runtime_store: true,
4877            durability_policy: DurabilityPolicy::SyncWriteThrough,
4878            bridge: Some(bridge.clone()),
4879            default_timeout: None,
4880        });
4881
4882        let result = super::super::orchestrator::restore_flow(&runtime, &specs, None, None).await?;
4883
4884        assert_eq!(result.outcomes.len(), specs.len());
4885        assert!(
4886            bridge.max_creates_in_flight() > 1,
4887            "member creation should no longer be serial"
4888        );
4889        assert!(
4890            bridge.max_creates_in_flight()
4891                <= super::super::orchestrator::IDENTITY_RESTORE_CONCURRENCY,
4892            "restore concurrency must remain bounded"
4893        );
4894        Ok(())
4895    }
4896
4897    #[tokio::test]
4898    async fn reset_reprofiles_session_from_runtime_configured_roster_provider()
4899    -> Result<(), Box<dyn std::error::Error>> {
4900        let identity = AgentIdentity::parse("domain:security")?;
4901        let roster = Arc::new(MutableRoster::new(vec![durable_spec(
4902            identity.clone(),
4903            "domain",
4904        )]));
4905        let bridge = Arc::new(RecordingBridge::default());
4906        let runtime = Arc::new(
4907            IdentityRuntime::new(IdentityRuntimeConfig {
4908                continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
4909                lease_provider: Arc::new(LocalLeaseProvider::new()),
4910                runtime_instance_id: "reset-reprofile-test".to_string(),
4911                has_runtime_store: true,
4912                durability_policy: DurabilityPolicy::SyncWriteThrough,
4913                bridge: Some(bridge.clone()),
4914                default_timeout: None,
4915            })
4916            .with_reset_roster_provider(roster.clone()),
4917        );
4918
4919        super::super::orchestrator::restore_flow(
4920            &runtime,
4921            &roster
4922                .roster(&RosterContext {
4923                    mob_definition: None,
4924                    previous_identities: Vec::new(),
4925                })
4926                .await?,
4927            None,
4928            None,
4929        )
4930        .await?;
4931        roster
4932            .set(vec![durable_spec(identity.clone(), "security")])
4933            .await;
4934
4935        let record = runtime.reset(&identity).await?;
4936
4937        assert_eq!(record.generation.get(), 1);
4938        assert_eq!(
4939            bridge.create_profiles().await,
4940            vec!["domain".to_string(), "security".to_string()]
4941        );
4942        let status = runtime.status(&identity).await?;
4943        assert_eq!(
4944            status.profile.map(|profile| profile.to_string()).as_deref(),
4945            Some("security")
4946        );
4947        Ok(())
4948    }
4949
4950    #[tokio::test]
4951    async fn reset_reprofiles_session_from_identity_first_context_roster_provider()
4952    -> Result<(), Box<dyn std::error::Error>> {
4953        let identity = AgentIdentity::parse("domain:security")?;
4954        let roster = Arc::new(MutableRoster::new(vec![durable_spec(
4955            identity.clone(),
4956            "domain",
4957        )]));
4958        let bridge = Arc::new(RecordingBridge::default());
4959        let runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
4960            continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
4961            lease_provider: Arc::new(LocalLeaseProvider::new()),
4962            runtime_instance_id: "reset-reprofile-context-test".to_string(),
4963            has_runtime_store: true,
4964            durability_policy: DurabilityPolicy::SyncWriteThrough,
4965            bridge: Some(bridge.clone()),
4966            default_timeout: None,
4967        }));
4968        let context =
4969            IdentityFirstRuntimeContext::new(runtime.clone(), roster.clone(), None, None, None);
4970
4971        context.refresh_desired_topology().await?;
4972        roster
4973            .set(vec![durable_spec(identity.clone(), "security")])
4974            .await;
4975
4976        let record = runtime.reset(&identity).await?;
4977
4978        assert_eq!(record.generation.get(), 1);
4979        assert_eq!(
4980            bridge.create_profiles().await,
4981            vec!["domain".to_string(), "security".to_string()]
4982        );
4983        let status = runtime.status(&identity).await?;
4984        assert_eq!(
4985            status.profile.map(|profile| profile.to_string()).as_deref(),
4986            Some("security")
4987        );
4988        Ok(())
4989    }
4990
4991    #[tokio::test]
4992    async fn reset_does_not_retire_old_generation_during_cleanup()
4993    -> Result<(), Box<dyn std::error::Error>> {
4994        let identity = AgentIdentity::parse("domain:security")?;
4995        let roster = Arc::new(MutableRoster::new(vec![durable_spec(
4996            identity.clone(),
4997            "domain",
4998        )]));
4999        let bridge = Arc::new(RecordingBridge::default());
5000        let runtime = Arc::new(
5001            IdentityRuntime::new(IdentityRuntimeConfig {
5002                continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5003                lease_provider: Arc::new(LocalLeaseProvider::new()),
5004                runtime_instance_id: "reset-skips-old-retire-test".to_string(),
5005                has_runtime_store: true,
5006                durability_policy: DurabilityPolicy::SyncWriteThrough,
5007                bridge: Some(bridge.clone()),
5008                default_timeout: Some(Duration::from_millis(50)),
5009            })
5010            .with_reset_roster_provider(roster.clone()),
5011        );
5012
5013        super::super::orchestrator::restore_flow(
5014            &runtime,
5015            &roster
5016                .roster(&RosterContext {
5017                    mob_definition: None,
5018                    previous_identities: Vec::new(),
5019                })
5020                .await?,
5021            None,
5022            None,
5023        )
5024        .await?;
5025
5026        let old_runtime_id = AgentRuntimeId::parse("rt:domain:security:0")?;
5027        bridge.hang_retire_for(&old_runtime_id).await;
5028        roster
5029            .set(vec![durable_spec(identity.clone(), "security")])
5030            .await;
5031
5032        let record = tokio::time::timeout(Duration::from_secs(1), runtime.reset(&identity))
5033            .await
5034            .map_err(|_| "reset timed out waiting for old generation retirement")??;
5035
5036        assert_eq!(record.generation.get(), 1);
5037        assert_eq!(
5038            bridge.create_profiles().await,
5039            vec!["domain".to_string(), "security".to_string()]
5040        );
5041        assert!(
5042            !bridge
5043                .retired_runtime_ids()
5044                .await
5045                .contains(&old_runtime_id.to_string()),
5046            "reset cleanup must not call the cancellation-unsafe mob-member retire path"
5047        );
5048        let status = runtime.status(&identity).await?;
5049        assert_eq!(
5050            status.profile.map(|profile| profile.to_string()).as_deref(),
5051            Some("security")
5052        );
5053        Ok(())
5054    }
5055
5056    #[tokio::test]
5057    async fn reset_returns_when_old_session_unregister_fails_after_new_generation()
5058    -> Result<(), Box<dyn std::error::Error>> {
5059        let identity = AgentIdentity::parse("domain:security")?;
5060        let roster = Arc::new(MutableRoster::new(vec![durable_spec(
5061            identity.clone(),
5062            "domain",
5063        )]));
5064        let bridge = Arc::new(RecordingBridge::default());
5065        let runtime = Arc::new(
5066            IdentityRuntime::new(IdentityRuntimeConfig {
5067                continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
5068                lease_provider: Arc::new(LocalLeaseProvider::new()),
5069                runtime_instance_id: "reset-unregister-failure-test".to_string(),
5070                has_runtime_store: true,
5071                durability_policy: DurabilityPolicy::SyncWriteThrough,
5072                bridge: Some(bridge.clone()),
5073                default_timeout: Some(Duration::from_millis(50)),
5074            })
5075            .with_reset_roster_provider(roster.clone()),
5076        );
5077
5078        super::super::orchestrator::restore_flow(
5079            &runtime,
5080            &roster
5081                .roster(&RosterContext {
5082                    mob_definition: None,
5083                    previous_identities: Vec::new(),
5084                })
5085                .await?,
5086            None,
5087            None,
5088        )
5089        .await?;
5090
5091        let old_status = runtime.status(&identity).await?;
5092        let Some(old_session_id) = old_status.session_id else {
5093            return Err("initial session id missing".into());
5094        };
5095        bridge.fail_unregister_for(&old_session_id).await;
5096        roster
5097            .set(vec![durable_spec(identity.clone(), "security")])
5098            .await;
5099
5100        let record = tokio::time::timeout(Duration::from_secs(1), runtime.reset(&identity))
5101            .await
5102            .map_err(|_| "reset timed out waiting for old session unregister cleanup")??;
5103
5104        assert_eq!(record.generation.get(), 1);
5105        assert_eq!(
5106            bridge.create_profiles().await,
5107            vec!["domain".to_string(), "security".to_string()]
5108        );
5109        let status = runtime.status(&identity).await?;
5110        assert_eq!(
5111            status.profile.map(|profile| profile.to_string()).as_deref(),
5112            Some("security")
5113        );
5114        Ok(())
5115    }
5116}
5117
5118#[cfg(test)]
5119mod lease_renewal_backoff_tests {
5120    use super::*;
5121
5122    /// Regression: a lease provider that errors persistently must not spin the
5123    /// renewal task at the TTL-derived floor (down to 10ms). The failure
5124    /// backoff grows from a 1s base and caps at the max poll interval.
5125    #[test]
5126    fn lease_renewal_failure_backoff_grows_and_caps() {
5127        let max = Duration::from_mins(1);
5128        assert_eq!(
5129            lease_renewal_failure_backoff(0, max),
5130            LEASE_RENEWAL_FAILURE_BACKOFF_BASE
5131        );
5132        assert_eq!(
5133            lease_renewal_failure_backoff(1, max),
5134            LEASE_RENEWAL_FAILURE_BACKOFF_BASE * 2
5135        );
5136        assert_eq!(lease_renewal_failure_backoff(6, max), max);
5137        // Saturates at the cap for arbitrarily many failures (no shift overflow).
5138        assert_eq!(lease_renewal_failure_backoff(99, max), max);
5139        assert!(lease_renewal_failure_backoff(2, max) > lease_renewal_failure_backoff(1, max));
5140    }
5141}