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::bridge::SessionBridge;
19use super::contracts::{
20    AgentCustomizer, ContinuityStore, LeaseProvider, RosterProvider, TopologyProvider,
21};
22use super::types::{
23    AgentAddressability, AgentBuildContext, AgentIdentity, AgentRuntimeId, AgentRuntimeServices,
24    CheckpointVersion, ContinuityGeneration, ContinuityHealth, ContinuityRecord,
25    ContinuityStoreError, DispatchInput, DurabilityPolicy, DurableAgentSpec, FencingToken,
26    IdentityLifecycleState, IdentityStatus, LeaseGrant, LeaseInfo, ManagedPeerEdge, NotAddressable,
27    RosterContext, SessionSnapshot,
28};
29
30const MANAGED_PEER_RECONCILE_CONCURRENCY: usize = 64;
31const MATERIALIZATION_FAILURE_BACKOFF: Duration = Duration::from_secs(30);
32
33fn durable_spec_uses_external_binding(spec: &DurableAgentSpec) -> bool {
34    matches!(spec.backend, Some(meerkat_mob::MobBackendKind::External))
35        || matches!(
36            spec.binding.as_ref(),
37            Some(meerkat_contracts::WireRuntimeBinding::External { .. })
38        )
39}
40
41// ---------------------------------------------------------------------------
42// Error types
43// ---------------------------------------------------------------------------
44
45/// Errors from identity-first runtime operations.
46#[derive(Debug)]
47pub enum IdentityRuntimeError {
48    /// Target identity is not registered/active.
49    UnknownIdentity(AgentIdentity),
50    /// send() rejected: target is InternalOnly.
51    NotAddressable(NotAddressable),
52    /// Operation rejected: no active lease for this identity.
53    NoActiveLease(AgentIdentity),
54    /// Operation rejected: lease was lost.
55    LeaseLost(AgentIdentity),
56    /// Operation rejected: identity is not in a state that permits this operation.
57    InvalidState {
58        identity: AgentIdentity,
59        state: IdentityLifecycleState,
60        operation: &'static str,
61    },
62    /// Continuity store error.
63    Store(ContinuityStoreError),
64    /// Lease provider error.
65    Lease(super::types::LeaseError),
66    /// Duplicate identities in roster.
67    DuplicateIdentity(AgentIdentity),
68    /// Stale fencing token on checkpoint.
69    StaleFencingToken {
70        identity: AgentIdentity,
71        presented: FencingToken,
72        current: FencingToken,
73    },
74    /// Stale checkpoint version.
75    StaleCheckpointVersion {
76        identity: AgentIdentity,
77        presented: CheckpointVersion,
78        current: CheckpointVersion,
79    },
80    /// Generic I/O or internal error.
81    Internal(String),
82}
83
84impl std::fmt::Display for IdentityRuntimeError {
85    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
86        match self {
87            Self::UnknownIdentity(id) => write!(f, "unknown identity: {id}"),
88            Self::NotAddressable(err) => write!(f, "{err}"),
89            Self::NoActiveLease(id) => write!(f, "no active lease for {id}"),
90            Self::LeaseLost(id) => write!(f, "lease lost for {id}"),
91            Self::InvalidState {
92                identity,
93                state,
94                operation,
95            } => write!(
96                f,
97                "cannot {operation} identity {identity} in state {state:?}"
98            ),
99            Self::Store(err) => write!(f, "continuity store: {err}"),
100            Self::Lease(err) => write!(f, "lease provider: {err}"),
101            Self::DuplicateIdentity(id) => write!(f, "duplicate identity in roster: {id}"),
102            Self::StaleFencingToken {
103                identity,
104                presented,
105                current,
106            } => write!(
107                f,
108                "stale fencing token for {identity}: presented {presented}, current {current}"
109            ),
110            Self::StaleCheckpointVersion {
111                identity,
112                presented,
113                current,
114            } => write!(
115                f,
116                "stale checkpoint version for {identity}: presented {presented}, current {current}"
117            ),
118            Self::Internal(msg) => write!(f, "internal: {msg}"),
119        }
120    }
121}
122
123impl std::error::Error for IdentityRuntimeError {}
124
125impl From<ContinuityStoreError> for IdentityRuntimeError {
126    fn from(err: ContinuityStoreError) -> Self {
127        match err {
128            ContinuityStoreError::StaleFencingToken {
129                identity,
130                presented,
131                current,
132            } => Self::StaleFencingToken {
133                identity,
134                presented,
135                current,
136            },
137            ContinuityStoreError::StaleCheckpointVersion {
138                identity,
139                presented,
140                current,
141            } => Self::StaleCheckpointVersion {
142                identity,
143                presented,
144                current,
145            },
146            other => Self::Store(other),
147        }
148    }
149}
150
151// ---------------------------------------------------------------------------
152// Per-identity runtime state
153// ---------------------------------------------------------------------------
154
155/// Tracks the live state for a single identity within the runtime.
156#[derive(Debug, Clone)]
157pub(crate) struct IdentityEntry {
158    pub spec: DurableAgentSpec,
159    pub state: IdentityLifecycleState,
160    pub continuity: Option<ContinuityRecord>,
161    pub lease: Option<LeaseEntry>,
162    pub checkpoint_version: CheckpointVersion,
163    /// Whether a durable runtime_store is available (affects dispatch ack semantics).
164    pub has_runtime_store: bool,
165}
166
167/// Tracks a held lease for an identity.
168#[derive(Debug, Clone)]
169pub(crate) struct LeaseEntry {
170    pub fencing_token: FencingToken,
171    pub ttl: Duration,
172    pub acquired_at: Instant,
173}
174
175impl LeaseEntry {
176    pub fn is_expired(&self) -> bool {
177        self.acquired_at.elapsed() > self.ttl
178    }
179
180    pub fn ttl_remaining(&self) -> Duration {
181        self.ttl.saturating_sub(self.acquired_at.elapsed())
182    }
183
184    pub fn is_healthy(&self) -> bool {
185        // Healthy if more than 20% TTL remains
186        let remaining = self.ttl_remaining();
187        remaining > self.ttl / 5
188    }
189}
190
191// ---------------------------------------------------------------------------
192// Identity-scoped events
193// ---------------------------------------------------------------------------
194
195/// Events emitted for a specific identity, used by `subscribe()`.
196#[derive(Debug, Clone)]
197pub enum IdentityEvent {
198    /// Lifecycle state changed.
199    StateChanged {
200        identity: AgentIdentity,
201        new_state: IdentityLifecycleState,
202    },
203    /// Lease acquired or renewed.
204    LeaseUpdated {
205        identity: AgentIdentity,
206        fencing_token: FencingToken,
207    },
208    /// Lease lost.
209    LeaseLost { identity: AgentIdentity },
210    /// Checkpoint completed.
211    CheckpointCompleted {
212        identity: AgentIdentity,
213        version: CheckpointVersion,
214    },
215    /// Resume could not reuse a persisted runtime binding and materialization
216    /// fresh-spawned a member instead.
217    ResumeFallback {
218        identity: AgentIdentity,
219        reason: super::bridge::ResumeFallbackReason,
220    },
221}
222
223/// Per-identity event channel capacity.
224const IDENTITY_EVENT_CHANNEL_CAPACITY: usize = 64;
225const DEFAULT_LEASE_RENEWAL_MAX_POLL_INTERVAL: Duration = Duration::from_mins(1);
226const DEFAULT_LEASE_RENEWAL_MIN_POLL_INTERVAL: Duration = Duration::from_millis(10);
227/// Base delay for the lease-renewal failure backoff. The renewal tick's
228/// normal cadence is TTL-derived (down to a 10ms floor), so a lease provider
229/// that errors persistently would otherwise retry — and warn — at that floor
230/// rate. On failure we back off from this base, doubling toward the max poll
231/// interval, so a backend outage can't spin the renewal task.
232const LEASE_RENEWAL_FAILURE_BACKOFF_BASE: Duration = Duration::from_secs(1);
233
234fn lease_renewal_failure_backoff(
235    consecutive_failures: u32,
236    max_poll_interval: Duration,
237) -> Duration {
238    LEASE_RENEWAL_FAILURE_BACKOFF_BASE
239        .saturating_mul(1u32 << consecutive_failures.min(6))
240        .min(max_poll_interval)
241}
242
243// ---------------------------------------------------------------------------
244// IdentityRuntime
245// ---------------------------------------------------------------------------
246
247/// Configuration for the identity-first runtime.
248pub struct IdentityRuntimeConfig {
249    pub continuity_store: Arc<dyn ContinuityStore>,
250    pub lease_provider: Arc<dyn LeaseProvider>,
251    pub runtime_instance_id: String,
252    pub has_runtime_store: bool,
253    pub durability_policy: DurabilityPolicy,
254    /// Optional session bridge for real session delivery. When `None`,
255    /// delivery operations validate invariants but do not forward to
256    /// the Meerkat session pipeline (useful for tests).
257    pub bridge: Option<Arc<dyn SessionBridge>>,
258    /// Default timeout for wait_for_output / wait_for_output_containing.
259    /// Defaults to 90 seconds if not set.
260    pub default_timeout: Option<Duration>,
261}
262
263#[derive(Clone)]
264pub struct IdentityFirstRuntimeContext {
265    pub runtime: Arc<IdentityRuntime>,
266    pub roster_provider: Arc<dyn RosterProvider>,
267    pub topology_provider: Option<Arc<dyn TopologyProvider>>,
268    pub customizer: Option<Arc<dyn AgentCustomizer>>,
269    mob_definition: Option<meerkat_mob::MobDefinition>,
270    lazy_materialization: bool,
271}
272
273impl IdentityFirstRuntimeContext {
274    pub fn new(
275        runtime: Arc<IdentityRuntime>,
276        roster_provider: Arc<dyn RosterProvider>,
277        topology_provider: Option<Arc<dyn TopologyProvider>>,
278        customizer: Option<Arc<dyn AgentCustomizer>>,
279        mob_definition: Option<meerkat_mob::MobDefinition>,
280    ) -> Self {
281        Self::new_with_lazy_materialization(
282            runtime,
283            roster_provider,
284            topology_provider,
285            customizer,
286            mob_definition,
287            false,
288        )
289    }
290
291    pub fn new_with_lazy_materialization(
292        runtime: Arc<IdentityRuntime>,
293        roster_provider: Arc<dyn RosterProvider>,
294        topology_provider: Option<Arc<dyn TopologyProvider>>,
295        customizer: Option<Arc<dyn AgentCustomizer>>,
296        mob_definition: Option<meerkat_mob::MobDefinition>,
297        lazy_materialization: bool,
298    ) -> Self {
299        Self {
300            runtime,
301            roster_provider,
302            topology_provider,
303            customizer,
304            mob_definition,
305            lazy_materialization,
306        }
307    }
308
309    pub async fn refresh_desired_topology(
310        &self,
311    ) -> Result<super::orchestrator::RestoreFlowResult, IdentityRuntimeError> {
312        let roster = self
313            .roster_provider
314            .roster(&RosterContext {
315                mob_definition: self.mob_definition.clone(),
316                previous_identities: Vec::new(),
317            })
318            .await
319            .map_err(|err| IdentityRuntimeError::Internal(format!("roster provider: {err}")))?;
320
321        if self.lazy_materialization {
322            super::orchestrator::lazy_register_flow(
323                &self.runtime,
324                &roster,
325                self.topology_provider.as_deref(),
326            )
327            .await
328        } else {
329            super::orchestrator::restore_flow(
330                &self.runtime,
331                &roster,
332                self.topology_provider.as_deref(),
333                self.customizer.as_deref(),
334            )
335            .await
336        }
337    }
338}
339
340/// The identity-first runtime tracks active identities and enforces delivery,
341/// ownership, and lifecycle invariants.
342pub struct IdentityRuntime {
343    entries: RwLock<BTreeMap<AgentIdentity, IdentityEntry>>,
344    event_channels: RwLock<BTreeMap<AgentIdentity, broadcast::Sender<IdentityEvent>>>,
345    continuity_store: Arc<dyn ContinuityStore>,
346    lease_provider: Arc<dyn LeaseProvider>,
347    runtime_instance_id: String,
348    has_runtime_store: bool,
349    durability_policy: DurabilityPolicy,
350    bridge: Option<Arc<dyn SessionBridge>>,
351    runtime_services: AgentRuntimeServices,
352    managed_peer_edges: RwLock<BTreeSet<(AgentIdentity, AgentIdentity)>>,
353    managed_peer_reconcile_lock: Mutex<()>,
354    desired_peer_edges: RwLock<Vec<ManagedPeerEdge>>,
355    materialization_locks: RwLock<BTreeMap<AgentIdentity, Arc<Mutex<()>>>>,
356    best_effort_materialization_locks: RwLock<BTreeMap<AgentIdentity, Arc<Mutex<()>>>>,
357    lifecycle_locks: RwLock<BTreeMap<AgentIdentity, Arc<Mutex<()>>>>,
358    customizer: RwLock<Option<Arc<dyn AgentCustomizer>>>,
359    lease_renewal_notify: Notify,
360    default_timeout: Duration,
361    materialization_failure_backoff: RwLock<BTreeMap<AgentIdentity, MaterializationFailureBackoff>>,
362    error_hook: StdRwLock<Option<crate::unified_runtime::ErrorHook>>,
363}
364
365#[derive(Debug, Clone)]
366struct MaterializationFailureBackoff {
367    suppress_until: Instant,
368    error: String,
369}
370
371impl IdentityRuntime {
372    /// Create a new identity runtime with the given configuration.
373    pub fn new(config: IdentityRuntimeConfig) -> Self {
374        Self {
375            entries: RwLock::new(BTreeMap::new()),
376            event_channels: RwLock::new(BTreeMap::new()),
377            continuity_store: config.continuity_store,
378            lease_provider: config.lease_provider,
379            runtime_instance_id: config.runtime_instance_id,
380            has_runtime_store: config.has_runtime_store,
381            durability_policy: config.durability_policy,
382            bridge: config.bridge,
383            runtime_services: AgentRuntimeServices::empty(),
384            managed_peer_edges: RwLock::new(BTreeSet::new()),
385            managed_peer_reconcile_lock: Mutex::new(()),
386            desired_peer_edges: RwLock::new(Vec::new()),
387            materialization_locks: RwLock::new(BTreeMap::new()),
388            best_effort_materialization_locks: RwLock::new(BTreeMap::new()),
389            lifecycle_locks: RwLock::new(BTreeMap::new()),
390            customizer: RwLock::new(None),
391            lease_renewal_notify: Notify::new(),
392            default_timeout: config.default_timeout.unwrap_or(Duration::from_secs(90)),
393            materialization_failure_backoff: RwLock::new(BTreeMap::new()),
394            error_hook: StdRwLock::new(None),
395        }
396    }
397
398    pub fn with_runtime_services(mut self, runtime_services: AgentRuntimeServices) -> Self {
399        self.runtime_services = runtime_services;
400        self
401    }
402
403    pub(crate) fn runtime_services(&self) -> AgentRuntimeServices {
404        self.runtime_services.clone()
405    }
406
407    pub async fn set_agent_customizer(&self, customizer: Option<Arc<dyn AgentCustomizer>>) {
408        *self.customizer.write().await = customizer;
409    }
410
411    /// Attach a best-effort operational error hook used for alerting.
412    pub fn set_error_hook(&self, hook: Option<crate::unified_runtime::ErrorHook>) {
413        match self.error_hook.write() {
414            Ok(mut stored_hook) => *stored_hook = hook,
415            Err(err) => {
416                tracing::warn!(
417                    error = %err,
418                    "identity runtime error hook lock poisoned; dropping hook update"
419                );
420            }
421        }
422    }
423
424    /// Spawn a background supervisor that renews active identity leases before
425    /// they reach their TTL deadline.
426    pub fn spawn_lease_renewal_task(self: Arc<Self>) -> JoinHandle<()> {
427        self.spawn_lease_renewal_task_with_poll_interval(DEFAULT_LEASE_RENEWAL_MAX_POLL_INTERVAL)
428    }
429
430    /// Spawn a lease renewal supervisor with a caller-provided maximum poll
431    /// interval. Embedders can use this for shorter external lease TTLs; tests
432    /// use it to exercise renewal without waiting on wall-clock TTLs.
433    pub fn spawn_lease_renewal_task_with_poll_interval(
434        self: Arc<Self>,
435        max_poll_interval: Duration,
436    ) -> JoinHandle<()> {
437        let max_poll_interval = max_poll_interval.max(DEFAULT_LEASE_RENEWAL_MIN_POLL_INTERVAL);
438        tokio::spawn(async move {
439            let mut consecutive_failures: u32 = 0;
440            loop {
441                let base = self.lease_renewal_sleep_interval(max_poll_interval).await;
442                // While the provider is failing, hold off at least the backoff
443                // delay so a persistent outage retries at a bounded rate
444                // instead of the TTL-derived floor (down to 10ms).
445                let sleep = if consecutive_failures > 0 {
446                    base.max(lease_renewal_failure_backoff(
447                        consecutive_failures,
448                        max_poll_interval,
449                    ))
450                } else {
451                    base
452                };
453                tokio::select! {
454                    () = tokio::time::sleep(sleep) => {}
455                    () = self.lease_renewal_notify.notified() => {}
456                }
457                match self.renew_due_leases_once().await {
458                    Ok(_) => consecutive_failures = 0,
459                    Err(err) => {
460                        // Warn once, then debounce to debug so a backend outage
461                        // can't flood the log at the renewal cadence.
462                        if consecutive_failures == 0 {
463                            tracing::warn!(
464                                error = %err,
465                                "identity-first proactive lease renewal tick failed; backing off"
466                            );
467                        } else {
468                            tracing::debug!(
469                                error = %err,
470                                consecutive_failures,
471                                "identity-first lease renewal still failing; backing off"
472                            );
473                        }
474                        consecutive_failures = consecutive_failures.saturating_add(1);
475                    }
476                }
477            }
478        })
479    }
480
481    async fn lease_renewal_sleep_interval(&self, max_poll_interval: Duration) -> Duration {
482        let entries = self.entries.read().await;
483        entries
484            .values()
485            .filter(|entry| entry.state == IdentityLifecycleState::Active)
486            .filter_map(|entry| entry.lease.as_ref())
487            .map(|lease| (lease.ttl / 10).max(DEFAULT_LEASE_RENEWAL_MIN_POLL_INTERVAL))
488            .min()
489            .unwrap_or(max_poll_interval)
490            .min(max_poll_interval)
491    }
492
493    /// Renew every active lease that has entered the runtime's renewal window.
494    pub async fn renew_due_leases_once(&self) -> Result<usize, IdentityRuntimeError> {
495        let due = {
496            let entries = self.entries.read().await;
497            entries
498                .iter()
499                .filter(|(_, entry)| entry.state == IdentityLifecycleState::Active)
500                .filter_map(|(identity, entry)| {
501                    entry
502                        .lease
503                        .as_ref()
504                        .filter(|lease| !lease.is_healthy())
505                        .map(|_| identity.clone())
506                })
507                .collect::<Vec<_>>()
508        };
509
510        let mut renewed = 0;
511        let mut first_error = None;
512        for identity in due {
513            let lifecycle_lock = self.lifecycle_lock_for(&identity).await;
514            let _lifecycle_guard = lifecycle_lock.lock().await;
515            match self.ensure_active_lease(&identity).await {
516                Ok(_) => renewed += 1,
517                Err(err) => {
518                    if first_error.is_none() {
519                        first_error = Some(err);
520                    }
521                }
522            }
523        }
524        if let Some(err) = first_error {
525            Err(err)
526        } else {
527            Ok(renewed)
528        }
529    }
530
531    async fn release_uninstalled_materialize_lease(&self, grant: &LeaseGrant) -> Option<String> {
532        self.lease_provider
533            .release_leases(std::slice::from_ref(grant))
534            .await
535            .err()
536            .map(|err| err.to_string())
537    }
538
539    pub async fn set_desired_peer_edges(&self, edges: Vec<ManagedPeerEdge>) {
540        *self.desired_peer_edges.write().await = edges;
541    }
542
543    pub async fn desired_peer_edges(&self) -> Vec<ManagedPeerEdge> {
544        self.desired_peer_edges.read().await.clone()
545    }
546
547    async fn registered_identities(&self) -> Vec<AgentIdentity> {
548        self.entries.read().await.keys().cloned().collect()
549    }
550
551    async fn reachable_peer_identities(&self, identity: &AgentIdentity) -> Vec<AgentIdentity> {
552        self.desired_peer_edges
553            .read()
554            .await
555            .iter()
556            .filter_map(|edge| {
557                if edge.a() == identity {
558                    Some(edge.b().clone())
559                } else if edge.b() == identity {
560                    Some(edge.a().clone())
561                } else {
562                    None
563                }
564            })
565            .collect::<BTreeSet<_>>()
566            .into_iter()
567            .collect()
568    }
569
570    #[must_use]
571    pub fn has_session_bridge(&self) -> bool {
572        self.bridge.is_some()
573    }
574
575    /// Apply identity-first managed topology to the concrete mob graph.
576    ///
577    /// Topology providers return stable logical identities. The mob comms graph
578    /// is keyed by active runtime member IDs, so this resolves each endpoint
579    /// through continuity records before calling the same-mob bridge wire APIs.
580    pub async fn reconcile_managed_peer_edges(
581        &self,
582        desired_edges: &[ManagedPeerEdge],
583    ) -> Result<(), IdentityRuntimeError> {
584        let _guard = self.managed_peer_reconcile_lock.lock().await;
585        let Some(bridge) = self.bridge.clone() else {
586            return Ok(());
587        };
588
589        let active_runtimes: BTreeMap<AgentIdentity, AgentRuntimeId> = {
590            let entries = self.entries.read().await;
591            entries
592                .iter()
593                .filter_map(|(identity, entry)| {
594                    if entry.state != IdentityLifecycleState::Active {
595                        return None;
596                    }
597                    entry
598                        .continuity
599                        .as_ref()
600                        .map(|record| (identity.clone(), record.agent_runtime_id.clone()))
601                })
602                .collect()
603        };
604        let runtime_identities: BTreeMap<AgentRuntimeId, AgentIdentity> = active_runtimes
605            .iter()
606            .map(|(identity, runtime_id)| (runtime_id.clone(), identity.clone()))
607            .collect();
608        let current_logical_edges: Option<BTreeSet<(AgentIdentity, AgentIdentity)>> =
609            match bridge.current_member_wires().await {
610                Ok(current_runtime_edges) => Some(
611                    current_runtime_edges
612                        .iter()
613                        .filter_map(|(runtime_a, runtime_b)| {
614                            let a = runtime_identities.get(runtime_a)?;
615                            let b = runtime_identities.get(runtime_b)?;
616                            if a <= b {
617                                Some((a.clone(), b.clone()))
618                            } else {
619                                Some((b.clone(), a.clone()))
620                            }
621                        })
622                        .collect(),
623                ),
624                Err(err) => {
625                    tracing::debug!(
626                        error = %err,
627                        "identity-first topology reconcile could not inspect current member wires"
628                    );
629                    None
630                }
631            };
632
633        let desired: BTreeSet<(AgentIdentity, AgentIdentity)> = desired_edges
634            .iter()
635            .map(|edge| (edge.a().clone(), edge.b().clone()))
636            .collect();
637
638        let managed_snapshot = self.managed_peer_edges.read().await.clone();
639        let edge_is_managed_and_live = |edge: &(AgentIdentity, AgentIdentity)| {
640            // Managed-but-missing live edges are retried deliberately so tolerant topology restores self-heal.
641            managed_snapshot.contains(edge)
642                && current_logical_edges
643                    .as_ref()
644                    .is_none_or(|edges| edges.contains(edge))
645        };
646        let retained_logical_edges: Vec<(AgentIdentity, AgentIdentity)> = desired
647            .iter()
648            .filter(|edge| !edge_is_managed_and_live(edge))
649            .filter(|edge| {
650                current_logical_edges
651                    .as_ref()
652                    .is_some_and(|edges| edges.contains(*edge))
653            })
654            .filter(|(a, b)| active_runtimes.contains_key(a) && active_runtimes.contains_key(b))
655            .cloned()
656            .collect();
657        let to_wire: Vec<(AgentIdentity, AgentIdentity, AgentRuntimeId, AgentRuntimeId)> = desired
658            .iter()
659            .filter(|edge| !edge_is_managed_and_live(edge))
660            .filter(|edge| {
661                current_logical_edges
662                    .as_ref()
663                    .is_none_or(|edges| !edges.contains(*edge))
664            })
665            .filter_map(|(a, b)| {
666                let runtime_a = active_runtimes.get(a)?;
667                let runtime_b = active_runtimes.get(b)?;
668                Some((a.clone(), b.clone(), runtime_a.clone(), runtime_b.clone()))
669            })
670            .collect();
671
672        let stale: Vec<(AgentIdentity, AgentIdentity)> = managed_snapshot
673            .iter()
674            .filter(|edge| !desired.contains(*edge))
675            .cloned()
676            .collect();
677        let to_unwire: Vec<(AgentIdentity, AgentIdentity, AgentRuntimeId, AgentRuntimeId)> = stale
678            .iter()
679            .filter_map(|(a, b)| {
680                let runtime_a = active_runtimes.get(a)?;
681                let runtime_b = active_runtimes.get(b)?;
682                if current_logical_edges
683                    .as_ref()
684                    .is_some_and(|edges| !edges.contains(&(a.clone(), b.clone())))
685                {
686                    return None;
687                }
688                Some((a.clone(), b.clone(), runtime_a.clone(), runtime_b.clone()))
689            })
690            .collect();
691
692        let wire_logical_edges = to_wire
693            .iter()
694            .map(|(a, b, _, _)| (a.clone(), b.clone()))
695            .collect::<Vec<_>>();
696        let wire_runtime_edges = to_wire
697            .iter()
698            .map(|(_, _, runtime_a, runtime_b)| (runtime_a.clone(), runtime_b.clone()))
699            .collect::<Vec<_>>();
700        if !wire_runtime_edges.is_empty() {
701            bridge
702                .wire_peers_batch(&wire_runtime_edges)
703                .await
704                .map_err(|e| {
705                    IdentityRuntimeError::Internal(format!("bridge wire_peers_batch: {e}"))
706                })?;
707        }
708
709        let unwire_results =
710            stream::iter(to_unwire.into_iter().map(|(a, b, runtime_a, runtime_b)| {
711                let bridge = bridge.clone();
712                async move {
713                    let result = bridge
714                        .unwire_peer(&runtime_a, &runtime_b)
715                        .await
716                        .map_err(|e| format!("{e}"));
717                    (a, b, result)
718                }
719            }))
720            .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
721            .collect::<Vec<_>>()
722            .await;
723
724        let mut managed = self.managed_peer_edges.write().await;
725        for (a, b) in retained_logical_edges {
726            managed.insert((a, b));
727        }
728        for (a, b) in wire_logical_edges {
729            managed.insert((a, b));
730        }
731
732        for (a, b) in stale {
733            let key = (a.clone(), b.clone());
734            if !active_runtimes.contains_key(&a)
735                || !active_runtimes.contains_key(&b)
736                || current_logical_edges
737                    .as_ref()
738                    .is_some_and(|edges| !edges.contains(&key))
739            {
740                managed.remove(&key);
741            }
742        }
743        for (a, b, result) in unwire_results {
744            result
745                .map_err(|e| IdentityRuntimeError::Internal(format!("bridge unwire_peer: {e}")))?;
746            managed.remove(&(a, b));
747        }
748
749        Ok(())
750    }
751
752    /// Emit an event for the given identity. Best-effort — no error if no subscribers.
753    async fn emit_event(&self, identity: &AgentIdentity, event: IdentityEvent) {
754        let channels = self.event_channels.read().await;
755        if let Some(tx) = channels.get(identity) {
756            let _ = tx.send(event);
757        }
758    }
759
760    fn emit_error(&self, event: crate::unified_runtime::types::ErrorEvent) {
761        let hook = match self.error_hook.read() {
762            Ok(stored_hook) => stored_hook.clone(),
763            Err(err) => {
764                tracing::warn!(
765                    error = %err,
766                    "identity runtime error hook lock poisoned; dropping error event"
767                );
768                None
769            }
770        };
771        if let Some(hook) = hook {
772            tokio::spawn(async move {
773                let () = hook(event).await;
774            });
775        }
776    }
777
778    async fn materialization_backoff_error(&self, identity: &AgentIdentity) -> Option<String> {
779        let backoffs = self.materialization_failure_backoff.read().await;
780        let backoff = backoffs.get(identity)?;
781        if Instant::now() < backoff.suppress_until {
782            Some(backoff.error.clone())
783        } else {
784            None
785        }
786    }
787
788    async fn clear_materialization_backoff(&self, identity: &AgentIdentity) {
789        self.materialization_failure_backoff
790            .write()
791            .await
792            .remove(identity);
793    }
794
795    async fn record_best_effort_materialization_failure(
796        &self,
797        identity: &AgentIdentity,
798        initiator: Option<&AgentIdentity>,
799        operation: &'static str,
800        err: &IdentityRuntimeError,
801    ) {
802        let error = err.to_string();
803        let suppress_until = Instant::now() + MATERIALIZATION_FAILURE_BACKOFF;
804        self.materialization_failure_backoff.write().await.insert(
805            identity.clone(),
806            MaterializationFailureBackoff {
807                suppress_until,
808                error: error.clone(),
809            },
810        );
811        self.emit_error(
812            crate::unified_runtime::types::ErrorEvent::IdentityMaterializationFailure {
813                identity: identity.to_string(),
814                initiator: initiator.map(ToString::to_string),
815                operation: operation.to_string(),
816                error,
817            },
818        );
819    }
820
821    // -----------------------------------------------------------------------
822    // Registration / activation
823    // -----------------------------------------------------------------------
824
825    /// Register an identity entry in the runtime (called during restore flow).
826    pub async fn register(
827        &self,
828        spec: DurableAgentSpec,
829        state: IdentityLifecycleState,
830        continuity: Option<ContinuityRecord>,
831        lease: Option<LeaseGrant>,
832    ) {
833        let identity = spec.identity.clone();
834        let cpv = continuity
835            .as_ref()
836            .map(|r| r.checkpoint_version)
837            .unwrap_or(CheckpointVersion::new(0));
838        let lease_entry = lease.map(|g| LeaseEntry {
839            fencing_token: g.fencing_token,
840            ttl: g.ttl,
841            acquired_at: Instant::now(),
842        });
843        let entry = IdentityEntry {
844            spec,
845            state,
846            continuity,
847            lease: lease_entry,
848            checkpoint_version: cpv,
849            has_runtime_store: self.has_runtime_store,
850        };
851        let has_active_lease =
852            entry.state == IdentityLifecycleState::Active && entry.lease.is_some();
853        self.entries.write().await.insert(identity.clone(), entry);
854
855        // Create event channel for this identity
856        let (tx, _) = broadcast::channel(IDENTITY_EVENT_CHANNEL_CAPACITY);
857        self.event_channels
858            .write()
859            .await
860            .insert(identity.clone(), tx);
861        if has_active_lease {
862            self.lease_renewal_notify.notify_one();
863        }
864    }
865
866    async fn materialization_lock_for(&self, identity: &AgentIdentity) -> Arc<Mutex<()>> {
867        if let Some(lock) = self.materialization_locks.read().await.get(identity) {
868            return lock.clone();
869        }
870        let mut locks = self.materialization_locks.write().await;
871        locks
872            .entry(identity.clone())
873            .or_insert_with(|| Arc::new(Mutex::new(())))
874            .clone()
875    }
876
877    async fn best_effort_materialization_lock_for(
878        &self,
879        identity: &AgentIdentity,
880    ) -> Arc<Mutex<()>> {
881        if let Some(lock) = self
882            .best_effort_materialization_locks
883            .read()
884            .await
885            .get(identity)
886        {
887            return lock.clone();
888        }
889        let mut locks = self.best_effort_materialization_locks.write().await;
890        locks
891            .entry(identity.clone())
892            .or_insert_with(|| Arc::new(Mutex::new(())))
893            .clone()
894    }
895
896    async fn lifecycle_lock_for(&self, identity: &AgentIdentity) -> Arc<Mutex<()>> {
897        if let Some(lock) = self.lifecycle_locks.read().await.get(identity) {
898            return lock.clone();
899        }
900        let mut locks = self.lifecycle_locks.write().await;
901        locks
902            .entry(identity.clone())
903            .or_insert_with(|| Arc::new(Mutex::new(())))
904            .clone()
905    }
906
907    /// Materialize a dormant identity into a concrete mob member/session.
908    ///
909    /// This is the lazy counterpart to eager `restore_flow`: it performs the
910    /// expensive bridge create/resume and snapshot load only when an identity is
911    /// actually touched. Parallel calls for one identity coalesce on a
912    /// per-identity lock and re-check state after acquiring it.
913    pub async fn materialize(
914        &self,
915        identity: &AgentIdentity,
916    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
917        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
918        let _lifecycle_guard = lifecycle_lock.lock().await;
919        let lock = self.materialization_lock_for(identity).await;
920        let _guard = lock.lock().await;
921
922        let (spec, continuity, state) = {
923            let entries = self.entries.read().await;
924            let entry = entries
925                .get(identity)
926                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
927            if entry.state == IdentityLifecycleState::Active {
928                let continuity = entry.continuity.clone().ok_or_else(|| {
929                    IdentityRuntimeError::Internal(format!(
930                        "active identity {identity} has no continuity record"
931                    ))
932                })?;
933                drop(entries);
934                self.clear_materialization_backoff(identity).await;
935                return Ok(continuity);
936            }
937            (entry.spec.clone(), entry.continuity.clone(), entry.state)
938        };
939        let original_continuity = continuity.clone();
940        let continuity = if durable_spec_uses_external_binding(&spec) {
941            None
942        } else {
943            continuity
944        };
945
946        match state {
947            IdentityLifecycleState::Dormant | IdentityLifecycleState::Uninitialized => {}
948            IdentityLifecycleState::Broken
949            | IdentityLifecycleState::Retiring
950            | IdentityLifecycleState::Suspended => {
951                return Err(IdentityRuntimeError::InvalidState {
952                    identity: identity.clone(),
953                    state,
954                    operation: "materialize",
955                });
956            }
957            IdentityLifecycleState::Active => unreachable!("active handled above"),
958        }
959
960        let lease_results = self
961            .lease_provider
962            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
963            .await
964            .map_err(IdentityRuntimeError::Lease)?;
965        let grant = match lease_results.get(identity) {
966            Some(super::types::LeaseAcquireResult::Acquired(grant)) => grant.clone(),
967            _ => return Err(IdentityRuntimeError::NoActiveLease(identity.clone())),
968        };
969        if let Some(record) = continuity.as_ref()
970            && let Err(err) = self
971                .continuity_store
972                .upsert_continuity_record(record, grant.fencing_token)
973                .await
974        {
975            let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
976            return Err(IdentityRuntimeError::Internal(format!(
977                "continuity upsert before materialize: {err}{}",
978                cleanup_error
979                    .as_ref()
980                    .map(|e| format!("; lease cleanup failed: {e}"))
981                    .unwrap_or_default(),
982            )));
983        }
984
985        let active_peers = self.entries.read().await.keys().cloned().collect();
986        let managed_edges = self.desired_peer_edges.read().await.clone();
987        let build_context = AgentBuildContext {
988            identity: identity.clone(),
989            active_peers,
990            managed_edges,
991            runtime_services: self.runtime_services(),
992        };
993        let mut draft = super::types::AgentBuildDraft {
994            model: None,
995            system_prompt: None,
996            additional_instructions: spec.additional_instructions.clone(),
997            labels: spec.labels.clone(),
998            app_context: spec.context.clone(),
999            external_tools: Vec::new(),
1000            local_external_tools: Default::default(),
1001        };
1002        if let Some(customizer) = self.customizer.read().await.clone()
1003            && let Err(err) = customizer
1004                .customize_build(&build_context, &spec, &mut draft)
1005                .await
1006        {
1007            let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1008            return Err(IdentityRuntimeError::Internal(format!(
1009                "customizer: {err}{}",
1010                cleanup_error
1011                    .as_ref()
1012                    .map(|e| format!("; lease cleanup failed: {e}"))
1013                    .unwrap_or_default(),
1014            )));
1015        }
1016
1017        let mut abandoned_session_registrations: Vec<SessionId> = Vec::new();
1018        let mut record = if let Some(mut record) = continuity {
1019            let snapshot = match self
1020                .continuity_store
1021                .load_session_snapshot(&record.session_id)
1022                .await
1023            {
1024                Ok(snapshot) => snapshot,
1025                Err(err) => {
1026                    let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1027                    return Err(IdentityRuntimeError::Internal(format!(
1028                        "load session snapshot before materialize: {err}{}",
1029                        cleanup_error
1030                            .as_ref()
1031                            .map(|e| format!("; lease cleanup failed: {e}"))
1032                            .unwrap_or_default(),
1033                    )));
1034                }
1035            };
1036
1037            if let Some(bridge) = self.bridge.as_ref() {
1038                if let Err(err) = bridge
1039                    .register_session_runtime_state(
1040                        &record.session_id,
1041                        identity,
1042                        record.generation,
1043                        record.checkpoint_version,
1044                        grant.fencing_token,
1045                    )
1046                    .await
1047                {
1048                    let unregister_error = Self::unregister_bridge_session_runtime_states(
1049                        bridge.as_ref(),
1050                        std::slice::from_ref(&record.session_id),
1051                    )
1052                    .await;
1053                    let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1054                    return Err(IdentityRuntimeError::Internal(format!(
1055                        "bridge register_session_runtime_state: {err}{}{}",
1056                        unregister_error
1057                            .as_ref()
1058                            .map(|e| format!("; unregister session failed: {e}"))
1059                            .unwrap_or_default(),
1060                        cleanup_error
1061                            .as_ref()
1062                            .map(|e| format!("; lease cleanup failed: {e}"))
1063                            .unwrap_or_default(),
1064                    )));
1065                }
1066                let registered_session_id = record.session_id.clone();
1067                let snapshot = snapshot.unwrap_or(SessionSnapshot { data: Vec::new() });
1068                let outcome = bridge
1069                    .resume_session(
1070                        identity,
1071                        &record.agent_runtime_id,
1072                        &spec,
1073                        &draft,
1074                        &record.session_id,
1075                        &snapshot,
1076                    )
1077                    .await;
1078                let outcome = match outcome {
1079                    Ok(outcome) => outcome,
1080                    Err(err) => {
1081                        let unregister_error = Self::unregister_bridge_session_runtime_states(
1082                            bridge.as_ref(),
1083                            std::slice::from_ref(&registered_session_id),
1084                        )
1085                        .await;
1086                        let cleanup_error =
1087                            bridge.retire_member(&record.agent_runtime_id).await.err();
1088                        let lease_cleanup_error =
1089                            self.release_uninstalled_materialize_lease(&grant).await;
1090                        let detail = format!(
1091                            "bridge resume_session: {err}{}{}{}",
1092                            unregister_error
1093                                .as_ref()
1094                                .map(|e| format!("; unregister session failed: {e}"))
1095                                .unwrap_or_default(),
1096                            cleanup_error
1097                                .as_ref()
1098                                .map(|e| format!("; cleanup retire failed: {e}"))
1099                                .unwrap_or_default(),
1100                            lease_cleanup_error
1101                                .as_ref()
1102                                .map(|e| format!("; lease cleanup failed: {e}"))
1103                                .unwrap_or_default(),
1104                        );
1105                        return Err(IdentityRuntimeError::Internal(detail));
1106                    }
1107                };
1108                if let Some(reason) = outcome.fallback_reason().cloned() {
1109                    tracing::warn!(
1110                        %identity,
1111                        reason = ?reason,
1112                        "lazy identity materialization fresh-spawned after typed resume fallback"
1113                    );
1114                    self.emit_event(
1115                        identity,
1116                        IdentityEvent::ResumeFallback {
1117                            identity: identity.clone(),
1118                            reason,
1119                        },
1120                    )
1121                    .await;
1122                }
1123                let effective_session_id = outcome.session_id().clone();
1124                if effective_session_id != registered_session_id {
1125                    abandoned_session_registrations.push(registered_session_id);
1126                }
1127                record.session_id = effective_session_id;
1128            }
1129            record
1130        } else {
1131            let new_runtime_id =
1132                AgentRuntimeId::parse(&format!("rt:{identity}:0")).map_err(|err| {
1133                    IdentityRuntimeError::Internal(format!("failed to mint runtime id: {err}"))
1134                })?;
1135            let mut record = ContinuityRecord {
1136                identity: identity.clone(),
1137                agent_runtime_id: new_runtime_id,
1138                session_id: meerkat_core::types::SessionId::new(),
1139                generation: ContinuityGeneration::new(0),
1140                checkpoint_version: CheckpointVersion::new(0),
1141            };
1142            if let Err(err) = self
1143                .continuity_store
1144                .upsert_continuity_record(&record, grant.fencing_token)
1145                .await
1146            {
1147                let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1148                return Err(IdentityRuntimeError::Internal(format!(
1149                    "continuity upsert before materialize create: {err}{}",
1150                    cleanup_error
1151                        .as_ref()
1152                        .map(|e| format!("; lease cleanup failed: {e}"))
1153                        .unwrap_or_default(),
1154                )));
1155            }
1156            if let Some(bridge) = self.bridge.as_ref() {
1157                let provisional_session_id = record.session_id.clone();
1158                if let Err(err) = bridge
1159                    .register_session_runtime_state(
1160                        &record.session_id,
1161                        identity,
1162                        record.generation,
1163                        record.checkpoint_version,
1164                        grant.fencing_token,
1165                    )
1166                    .await
1167                    .map_err(|err| {
1168                        IdentityRuntimeError::Internal(format!(
1169                            "bridge register_session_runtime_state: {err}"
1170                        ))
1171                    })
1172                {
1173                    let unregister_error = Self::unregister_bridge_session_runtime_states(
1174                        bridge.as_ref(),
1175                        std::slice::from_ref(&provisional_session_id),
1176                    )
1177                    .await;
1178                    let delete_error = self
1179                        .continuity_store
1180                        .delete_continuity_record(identity, grant.fencing_token)
1181                        .await
1182                        .err();
1183                    let cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1184                    if let Some(delete_error) = delete_error {
1185                        return Err(IdentityRuntimeError::Internal(format!(
1186                            "{err}{}; tentative continuity cleanup failed: {delete_error}{}",
1187                            unregister_error
1188                                .as_ref()
1189                                .map(|e| format!("; unregister session failed: {e}"))
1190                                .unwrap_or_default(),
1191                            cleanup_error
1192                                .as_ref()
1193                                .map(|e| format!("; lease cleanup failed: {e}"))
1194                                .unwrap_or_default(),
1195                        )));
1196                    }
1197                    if let Some(cleanup_error) = cleanup_error {
1198                        return Err(IdentityRuntimeError::Internal(format!(
1199                            "{err}{}; lease cleanup failed: {cleanup_error}",
1200                            unregister_error
1201                                .as_ref()
1202                                .map(|e| format!("; unregister session failed: {e}"))
1203                                .unwrap_or_default(),
1204                        )));
1205                    }
1206                    if let Some(unregister_error) = unregister_error {
1207                        return Err(IdentityRuntimeError::Internal(format!(
1208                            "{err}; unregister session failed: {unregister_error}"
1209                        )));
1210                    }
1211                    return Err(err);
1212                }
1213                let created_session_id = bridge
1214                    .create_session(
1215                        identity,
1216                        &record.agent_runtime_id,
1217                        &spec,
1218                        &draft,
1219                        &record.session_id,
1220                    )
1221                    .await
1222                    .map_err(|err| {
1223                        IdentityRuntimeError::Internal(format!("bridge create_session: {err}"))
1224                    });
1225                match created_session_id {
1226                    Ok(session_id) => {
1227                        if session_id != provisional_session_id {
1228                            abandoned_session_registrations.push(provisional_session_id);
1229                        }
1230                        record.session_id = session_id;
1231                    }
1232                    Err(err) => {
1233                        let unregister_error = Self::unregister_bridge_session_runtime_states(
1234                            bridge.as_ref(),
1235                            std::slice::from_ref(&provisional_session_id),
1236                        )
1237                        .await;
1238                        let cleanup_error =
1239                            bridge.retire_member(&record.agent_runtime_id).await.err();
1240                        let delete_error = self
1241                            .continuity_store
1242                            .delete_continuity_record(identity, grant.fencing_token)
1243                            .await
1244                            .err();
1245                        let lease_cleanup_error =
1246                            self.release_uninstalled_materialize_lease(&grant).await;
1247                        if unregister_error.is_some()
1248                            || cleanup_error.is_some()
1249                            || delete_error.is_some()
1250                            || lease_cleanup_error.is_some()
1251                        {
1252                            return Err(IdentityRuntimeError::Internal(format!(
1253                                "{err}{}{}{}{}",
1254                                unregister_error
1255                                    .as_ref()
1256                                    .map(|e| format!("; unregister session failed: {e}"))
1257                                    .unwrap_or_default(),
1258                                cleanup_error
1259                                    .as_ref()
1260                                    .map(|e| format!("; cleanup retire failed: {e}"))
1261                                    .unwrap_or_default(),
1262                                delete_error
1263                                    .as_ref()
1264                                    .map(|e| format!("; tentative continuity cleanup failed: {e}"))
1265                                    .unwrap_or_default(),
1266                                lease_cleanup_error
1267                                    .as_ref()
1268                                    .map(|e| format!("; lease cleanup failed: {e}"))
1269                                    .unwrap_or_default(),
1270                            )));
1271                        }
1272                        return Err(err);
1273                    }
1274                }
1275            }
1276            record
1277        };
1278
1279        if let Err(err) = self
1280            .continuity_store
1281            .upsert_continuity_record(&record, grant.fencing_token)
1282            .await
1283        {
1284            let unregister_error = if let Some(bridge) = self.bridge.as_ref() {
1285                let mut sessions_to_unregister = abandoned_session_registrations.clone();
1286                sessions_to_unregister.push(record.session_id.clone());
1287                Self::unregister_bridge_session_runtime_states(
1288                    bridge.as_ref(),
1289                    &sessions_to_unregister,
1290                )
1291                .await
1292            } else {
1293                None
1294            };
1295            let cleanup_error = if let Some(bridge) = self.bridge.as_ref() {
1296                bridge.retire_member(&record.agent_runtime_id).await.err()
1297            } else {
1298                None
1299            };
1300            let restore_error = self
1301                .restore_continuity_after_materialize_failure(
1302                    identity,
1303                    original_continuity.as_ref(),
1304                    &grant,
1305                )
1306                .await;
1307            let lease_cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1308            if unregister_error.is_some()
1309                || cleanup_error.is_some()
1310                || restore_error.is_some()
1311                || lease_cleanup_error.is_some()
1312            {
1313                return Err(IdentityRuntimeError::Internal(format!(
1314                    "continuity upsert after materialize: {err}{}{}{}{}",
1315                    unregister_error
1316                        .as_ref()
1317                        .map(|e| format!("; unregister session failed: {e}"))
1318                        .unwrap_or_default(),
1319                    cleanup_error
1320                        .as_ref()
1321                        .map(|e| format!("; cleanup retire failed: {e}"))
1322                        .unwrap_or_default(),
1323                    restore_error
1324                        .as_ref()
1325                        .map(|e| format!("; continuity rollback failed: {e}"))
1326                        .unwrap_or_default(),
1327                    lease_cleanup_error
1328                        .as_ref()
1329                        .map(|e| format!("; lease cleanup failed: {e}"))
1330                        .unwrap_or_default(),
1331                )));
1332            }
1333            return Err(IdentityRuntimeError::Internal(format!(
1334                "continuity upsert after materialize: {err}"
1335            )));
1336        }
1337        if let Some(bridge) = self.bridge.as_ref() {
1338            let register_result = bridge
1339                .register_session_runtime_state(
1340                    &record.session_id,
1341                    identity,
1342                    record.generation,
1343                    record.checkpoint_version,
1344                    grant.fencing_token,
1345                )
1346                .await;
1347            let effective_checkpoint_version = match register_result {
1348                Ok(version) => version,
1349                Err(err) => {
1350                    let mut sessions_to_unregister = abandoned_session_registrations.clone();
1351                    sessions_to_unregister.push(record.session_id.clone());
1352                    let unregister_error = Self::unregister_bridge_session_runtime_states(
1353                        bridge.as_ref(),
1354                        &sessions_to_unregister,
1355                    )
1356                    .await;
1357                    let cleanup_error = bridge.retire_member(&record.agent_runtime_id).await.err();
1358                    let restore_error = self
1359                        .restore_continuity_after_materialize_failure(
1360                            identity,
1361                            original_continuity.as_ref(),
1362                            &grant,
1363                        )
1364                        .await;
1365                    let lease_cleanup_error =
1366                        self.release_uninstalled_materialize_lease(&grant).await;
1367                    return Err(IdentityRuntimeError::Internal(format!(
1368                        "bridge register actual session runtime state: {err}{}{}{}{}",
1369                        unregister_error
1370                            .as_ref()
1371                            .map(|e| format!("; unregister session failed: {e}"))
1372                            .unwrap_or_default(),
1373                        cleanup_error
1374                            .as_ref()
1375                            .map(|e| format!("; cleanup retire failed: {e}"))
1376                            .unwrap_or_default(),
1377                        restore_error
1378                            .as_ref()
1379                            .map(|e| format!("; continuity rollback failed: {e}"))
1380                            .unwrap_or_default(),
1381                        lease_cleanup_error
1382                            .as_ref()
1383                            .map(|e| format!("; lease cleanup failed: {e}"))
1384                            .unwrap_or_default(),
1385                    )));
1386                }
1387            };
1388            record.checkpoint_version = effective_checkpoint_version;
1389            if let Some(err) = Self::unregister_bridge_session_runtime_states(
1390                bridge.as_ref(),
1391                &abandoned_session_registrations,
1392            )
1393            .await
1394            {
1395                let actual_unregister_error = Self::unregister_bridge_session_runtime_states(
1396                    bridge.as_ref(),
1397                    std::slice::from_ref(&record.session_id),
1398                )
1399                .await;
1400                let unregister_error = actual_unregister_error
1401                    .map(|actual_err| format!("{err}; actual session: {actual_err}"))
1402                    .unwrap_or(err);
1403                let cleanup_error = bridge.retire_member(&record.agent_runtime_id).await.err();
1404                let restore_error = self
1405                    .restore_continuity_after_materialize_failure(
1406                        identity,
1407                        original_continuity.as_ref(),
1408                        &grant,
1409                    )
1410                    .await;
1411                let lease_cleanup_error = self.release_uninstalled_materialize_lease(&grant).await;
1412                return Err(IdentityRuntimeError::Internal(format!(
1413                    "bridge unregister abandoned session runtime state: {unregister_error}{}{}{}",
1414                    cleanup_error
1415                        .as_ref()
1416                        .map(|e| format!("; cleanup retire failed: {e}"))
1417                        .unwrap_or_default(),
1418                    restore_error
1419                        .as_ref()
1420                        .map(|e| format!("; continuity rollback failed: {e}"))
1421                        .unwrap_or_default(),
1422                    lease_cleanup_error
1423                        .as_ref()
1424                        .map(|e| format!("; lease cleanup failed: {e}"))
1425                        .unwrap_or_default(),
1426                )));
1427            }
1428        }
1429
1430        {
1431            let mut entries = self.entries.write().await;
1432            let entry = entries
1433                .get_mut(identity)
1434                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1435            entry.continuity = Some(record.clone());
1436            entry.lease = Some(Self::lease_entry_from_grant(&grant));
1437            entry.state = IdentityLifecycleState::Active;
1438            entry.checkpoint_version = record.checkpoint_version;
1439        }
1440        self.emit_event(
1441            identity,
1442            IdentityEvent::StateChanged {
1443                identity: identity.clone(),
1444                new_state: IdentityLifecycleState::Active,
1445            },
1446        )
1447        .await;
1448        let desired_edges = self.desired_peer_edges.read().await.clone();
1449        if !desired_edges.is_empty()
1450            && let Err(err) = self.reconcile_managed_peer_edges(&desired_edges).await
1451        {
1452            tracing::warn!(
1453                identity = %identity,
1454                error = %err,
1455                "identity materialized with topology reconcile warning"
1456            );
1457        }
1458        self.clear_materialization_backoff(identity).await;
1459        Ok(record)
1460    }
1461
1462    async fn best_effort_materialize_identity(
1463        &self,
1464        identity: AgentIdentity,
1465        initiator: Option<&AgentIdentity>,
1466        operation: &'static str,
1467    ) -> Option<ContinuityRecord> {
1468        let attempt_lock = self.best_effort_materialization_lock_for(&identity).await;
1469        let _attempt_guard = attempt_lock.lock().await;
1470
1471        if let Some(error) = self.materialization_backoff_error(&identity).await {
1472            tracing::debug!(
1473                identity = %identity,
1474                initiator = initiator.map(ToString::to_string).as_deref(),
1475                error = %error,
1476                "identity best-effort materialization skipped due to materialization backoff"
1477            );
1478            return None;
1479        }
1480
1481        match self.materialize(&identity).await {
1482            Ok(record) => {
1483                self.clear_materialization_backoff(&identity).await;
1484                Some(record)
1485            }
1486            Err(err) => {
1487                tracing::warn!(
1488                    identity = %identity,
1489                    initiator = initiator.map(ToString::to_string).as_deref(),
1490                    error = %err,
1491                    "identity best-effort materialization skipped identity after materialization failure"
1492                );
1493                self.record_best_effort_materialization_failure(
1494                    &identity, initiator, operation, &err,
1495                )
1496                .await;
1497                None
1498            }
1499        }
1500    }
1501
1502    async fn materialize_all_records(
1503        &self,
1504    ) -> Vec<(
1505        AgentIdentity,
1506        Result<ContinuityRecord, IdentityRuntimeError>,
1507    )> {
1508        let identities = self.registered_identities().await;
1509        stream::iter(identities.into_iter().map(|identity| async move {
1510            let result = self.materialize(&identity).await;
1511            (identity, result)
1512        }))
1513        .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
1514        .collect::<Vec<_>>()
1515        .await
1516    }
1517
1518    /// Materialize all identities currently registered with the runtime.
1519    ///
1520    /// Fleet hydration is best-effort: one member that cannot build is skipped
1521    /// and surfaced through logs/error hooks rather than aborting unrelated
1522    /// members.
1523    pub async fn materialize_all(&self) -> Result<Vec<ContinuityRecord>, IdentityRuntimeError> {
1524        let identities = self.registered_identities().await;
1525        let records = stream::iter(identities.into_iter().map(|identity| async move {
1526            self.best_effort_materialize_identity(identity, None, "materialize_all")
1527                .await
1528        }))
1529        .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
1530        .filter_map(async move |record| record)
1531        .collect::<Vec<_>>()
1532        .await;
1533
1534        let desired_edges = self.desired_peer_edges.read().await.clone();
1535        if !desired_edges.is_empty()
1536            && let Err(err) = self.reconcile_managed_peer_edges(&desired_edges).await
1537        {
1538            tracing::warn!(
1539                error = %err,
1540                "identity materialize_all completed with topology reconcile warning"
1541            );
1542        }
1543
1544        Ok(records)
1545    }
1546
1547    /// Materialize all identities and fail if any registered identity cannot
1548    /// hydrate. Flow admission uses this strict variant so a run is not accepted
1549    /// with a partially materialized identity-first fleet.
1550    pub async fn materialize_all_required(
1551        &self,
1552    ) -> Result<Vec<ContinuityRecord>, IdentityRuntimeError> {
1553        let results = self.materialize_all_records().await;
1554        let mut records = Vec::with_capacity(results.len());
1555        let mut failures = Vec::new();
1556
1557        for (identity, result) in results {
1558            match result {
1559                Ok(record) => records.push(record),
1560                Err(err) => failures.push(format!("{identity}: {err}")),
1561            }
1562        }
1563
1564        if !failures.is_empty() {
1565            return Err(IdentityRuntimeError::Internal(format!(
1566                "identity-first required materialization failed for {} identities: {}",
1567                failures.len(),
1568                failures.join("; ")
1569            )));
1570        }
1571
1572        let desired_edges = self.desired_peer_edges.read().await.clone();
1573        if !desired_edges.is_empty() {
1574            self.reconcile_managed_peer_edges(&desired_edges).await?;
1575        }
1576
1577        Ok(records)
1578    }
1579
1580    pub(crate) async fn best_effort_background_warm_identity(&self, identity: AgentIdentity) {
1581        self.best_effort_materialize_identity(identity, None, "background_warm")
1582            .await;
1583    }
1584
1585    /// Ensure an active identity's desired peer neighborhood exists in the
1586    /// concrete mob graph before ordinary communication starts.
1587    pub async fn materialize_reachable_peers(
1588        &self,
1589        identity: &AgentIdentity,
1590    ) -> Result<Vec<ContinuityRecord>, IdentityRuntimeError> {
1591        let peers = self.reachable_peer_identities(identity).await;
1592        let records = stream::iter(peers.into_iter().map(|peer| async move {
1593            self.best_effort_materialize_identity(
1594                peer,
1595                Some(identity),
1596                "materialize_reachable_peers",
1597            )
1598            .await
1599        }))
1600        .buffer_unordered(MANAGED_PEER_RECONCILE_CONCURRENCY)
1601        .filter_map(async move |record| record)
1602        .collect::<Vec<_>>()
1603        .await;
1604
1605        let desired_edges = self.desired_peer_edges.read().await.clone();
1606        if !desired_edges.is_empty()
1607            && let Err(err) = self.reconcile_managed_peer_edges(&desired_edges).await
1608        {
1609            tracing::warn!(
1610                identity = %identity,
1611                error = %err,
1612                "identity peer materialization completed with topology reconcile warning"
1613            );
1614        }
1615
1616        Ok(records)
1617    }
1618
1619    // -----------------------------------------------------------------------
1620    // Subscribe — REQ-06
1621    // -----------------------------------------------------------------------
1622
1623    /// Subscribe to identity-scoped events.
1624    ///
1625    /// Returns a broadcast receiver that yields `IdentityEvent` items for
1626    /// state changes, lease updates, lease loss, and checkpoint completions.
1627    pub async fn subscribe(
1628        &self,
1629        identity: &AgentIdentity,
1630    ) -> Result<broadcast::Receiver<IdentityEvent>, IdentityRuntimeError> {
1631        let channels = self.event_channels.read().await;
1632        let tx = channels
1633            .get(identity)
1634            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1635        Ok(tx.subscribe())
1636    }
1637
1638    /// Update the spec for an existing identity (used during reconciliation).
1639    pub async fn update_spec(&self, spec: DurableAgentSpec) -> Result<(), IdentityRuntimeError> {
1640        let mut entries = self.entries.write().await;
1641        let entry = entries
1642            .get_mut(&spec.identity)
1643            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(spec.identity.clone()))?;
1644        entry.spec = spec;
1645        Ok(())
1646    }
1647
1648    /// Update the lease for an identity.
1649    pub async fn update_lease(
1650        &self,
1651        identity: &AgentIdentity,
1652        grant: LeaseGrant,
1653    ) -> Result<(), IdentityRuntimeError> {
1654        let fencing_token = grant.fencing_token;
1655        let mut entries = self.entries.write().await;
1656        let entry = entries
1657            .get_mut(identity)
1658            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1659        entry.lease = Some(LeaseEntry {
1660            fencing_token,
1661            ttl: grant.ttl,
1662            acquired_at: Instant::now(),
1663        });
1664        drop(entries);
1665        self.lease_renewal_notify.notify_one();
1666        self.emit_event(
1667            identity,
1668            IdentityEvent::LeaseUpdated {
1669                identity: identity.clone(),
1670                fencing_token,
1671            },
1672        )
1673        .await;
1674        Ok(())
1675    }
1676
1677    /// Mark a lease as lost for an identity (INV-02).
1678    pub async fn mark_lease_lost(
1679        &self,
1680        identity: &AgentIdentity,
1681    ) -> Result<(), IdentityRuntimeError> {
1682        let mut entries = self.entries.write().await;
1683        let entry = entries
1684            .get_mut(identity)
1685            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1686        entry.lease = None;
1687        drop(entries);
1688        self.emit_event(
1689            identity,
1690            IdentityEvent::LeaseLost {
1691                identity: identity.clone(),
1692            },
1693        )
1694        .await;
1695        Ok(())
1696    }
1697
1698    /// Remove an identity from the runtime.
1699    #[allow(dead_code)]
1700    pub(crate) async fn remove(&self, identity: &AgentIdentity) -> Option<IdentityEntry> {
1701        self.event_channels.write().await.remove(identity);
1702        self.entries.write().await.remove(identity)
1703    }
1704
1705    /// Set the lifecycle state for an identity.
1706    pub async fn set_state(
1707        &self,
1708        identity: &AgentIdentity,
1709        state: IdentityLifecycleState,
1710    ) -> Result<(), IdentityRuntimeError> {
1711        let mut entries = self.entries.write().await;
1712        let entry = entries
1713            .get_mut(identity)
1714            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1715        entry.state = state;
1716        drop(entries);
1717        if state == IdentityLifecycleState::Active {
1718            self.lease_renewal_notify.notify_one();
1719        }
1720        self.emit_event(
1721            identity,
1722            IdentityEvent::StateChanged {
1723                identity: identity.clone(),
1724                new_state: state,
1725            },
1726        )
1727        .await;
1728        Ok(())
1729    }
1730
1731    // -----------------------------------------------------------------------
1732    // Lease checking (INV-01, INV-02)
1733    // -----------------------------------------------------------------------
1734
1735    /// Check that the identity has an active, non-expired lease.
1736    /// Returns the fencing token if valid.
1737    fn check_lease(entry: &IdentityEntry) -> Result<FencingToken, IdentityRuntimeError> {
1738        match &entry.lease {
1739            Some(lease) if !lease.is_expired() => Ok(lease.fencing_token),
1740            Some(_) => Err(IdentityRuntimeError::LeaseLost(entry.spec.identity.clone())),
1741            None => Err(IdentityRuntimeError::NoActiveLease(
1742                entry.spec.identity.clone(),
1743            )),
1744        }
1745    }
1746
1747    fn lease_entry_from_grant(grant: &LeaseGrant) -> LeaseEntry {
1748        LeaseEntry {
1749            fencing_token: grant.fencing_token,
1750            ttl: grant.ttl,
1751            acquired_at: Instant::now(),
1752        }
1753    }
1754
1755    async fn ensure_active_lease(
1756        &self,
1757        identity: &AgentIdentity,
1758    ) -> Result<FencingToken, IdentityRuntimeError> {
1759        let (grant, continuity) = {
1760            let entries = self.entries.read().await;
1761            let entry = entries
1762                .get(identity)
1763                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1764            let lease = match &entry.lease {
1765                Some(lease) if lease.is_healthy() => return Ok(lease.fencing_token),
1766                Some(lease) => lease,
1767                None => return Err(IdentityRuntimeError::NoActiveLease(identity.clone())),
1768            };
1769            (
1770                LeaseGrant {
1771                    identity: identity.clone(),
1772                    fencing_token: lease.fencing_token,
1773                    ttl: lease.ttl,
1774                },
1775                entry.continuity.clone(),
1776            )
1777        };
1778
1779        let renewed = self
1780            .lease_provider
1781            .renew_leases(std::slice::from_ref(&grant))
1782            .await
1783            .map_err(IdentityRuntimeError::Lease)?;
1784        let renewed_grant = match renewed.get(identity) {
1785            Some(super::types::LeaseRenewResult::Renewed(grant)) => grant.clone(),
1786            Some(super::types::LeaseRenewResult::Lost { .. }) | None => {
1787                self.mark_lease_lost(identity).await?;
1788                return Err(IdentityRuntimeError::LeaseLost(identity.clone()));
1789            }
1790        };
1791
1792        if let Some(record) = continuity.as_ref() {
1793            self.continuity_store
1794                .upsert_continuity_record(record, renewed_grant.fencing_token)
1795                .await
1796                .map_err(IdentityRuntimeError::Store)?;
1797            if let Some(bridge) = self.bridge.as_ref() {
1798                bridge
1799                    .register_session_runtime_state(
1800                        &record.session_id,
1801                        identity,
1802                        record.generation,
1803                        record.checkpoint_version,
1804                        renewed_grant.fencing_token,
1805                    )
1806                    .await
1807                    .map_err(|err| {
1808                        IdentityRuntimeError::Internal(format!(
1809                            "bridge refresh session runtime state after lease renewal: {err}"
1810                        ))
1811                    })?;
1812            }
1813        }
1814
1815        let fencing_token = renewed_grant.fencing_token;
1816        let mut entries = self.entries.write().await;
1817        let entry = entries
1818            .get_mut(identity)
1819            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1820        match entry.lease.as_ref() {
1821            Some(current) if current.fencing_token == grant.fencing_token => {
1822                entry.lease = Some(Self::lease_entry_from_grant(&renewed_grant));
1823            }
1824            Some(current) => return Ok(current.fencing_token),
1825            None => return Err(IdentityRuntimeError::NoActiveLease(identity.clone())),
1826        }
1827        drop(entries);
1828        self.lease_renewal_notify.notify_one();
1829
1830        self.emit_event(
1831            identity,
1832            IdentityEvent::LeaseUpdated {
1833                identity: identity.clone(),
1834                fencing_token,
1835            },
1836        )
1837        .await;
1838        Ok(fencing_token)
1839    }
1840
1841    async fn mark_lifecycle_in_progress(
1842        &self,
1843        identity: &AgentIdentity,
1844        state: IdentityLifecycleState,
1845    ) -> Result<IdentityEntry, IdentityRuntimeError> {
1846        let mut entries = self.entries.write().await;
1847        let entry = entries
1848            .get_mut(identity)
1849            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1850        let snapshot = entry.clone();
1851        entry.state = state;
1852        entry.lease = None;
1853        Ok(snapshot)
1854    }
1855
1856    async fn restore_entry(&self, identity: &AgentIdentity, entry: IdentityEntry) {
1857        self.entries.write().await.insert(identity.clone(), entry);
1858    }
1859
1860    async fn restore_entry_with_grant(
1861        &self,
1862        identity: &AgentIdentity,
1863        mut entry: IdentityEntry,
1864        grant: &LeaseGrant,
1865    ) {
1866        let restore_live_lease = entry.state == IdentityLifecycleState::Active;
1867        if let Some(record) = entry.continuity.as_ref() {
1868            if let Err(err) = self
1869                .continuity_store
1870                .upsert_continuity_record(record, grant.fencing_token)
1871                .await
1872            {
1873                tracing::warn!(
1874                    %identity,
1875                    error = %err,
1876                    "failed to advance restored continuity fencing token after lifecycle failure"
1877                );
1878                entry.state = IdentityLifecycleState::Broken;
1879            } else if let Some(bridge) = self.bridge.as_ref()
1880                && let Err(err) = bridge
1881                    .register_session_runtime_state(
1882                        &record.session_id,
1883                        identity,
1884                        record.generation,
1885                        record.checkpoint_version,
1886                        grant.fencing_token,
1887                    )
1888                    .await
1889            {
1890                tracing::warn!(
1891                    %identity,
1892                    error = %err,
1893                    "failed to refresh restored session runtime state after lifecycle failure"
1894                );
1895                entry.state = IdentityLifecycleState::Broken;
1896            }
1897            entry.lease = restore_live_lease.then(|| Self::lease_entry_from_grant(grant));
1898        }
1899        self.restore_entry(identity, entry).await;
1900        if restore_live_lease {
1901            self.lease_renewal_notify.notify_one();
1902        }
1903    }
1904
1905    pub(crate) async fn refresh_active_restore_grant(
1906        &self,
1907        identity: &AgentIdentity,
1908        grant: &LeaseGrant,
1909    ) -> Result<(), IdentityRuntimeError> {
1910        let record = {
1911            let entries = self.entries.read().await;
1912            let entry = entries
1913                .get(identity)
1914                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1915            if entry.state != IdentityLifecycleState::Active {
1916                return Err(IdentityRuntimeError::InvalidState {
1917                    identity: identity.clone(),
1918                    state: entry.state,
1919                    operation: "refresh_active_restore_grant",
1920                });
1921            }
1922            entry.continuity.clone()
1923        };
1924
1925        if let Some(record) = record.as_ref()
1926            && let Err(err) = self
1927                .continuity_store
1928                .upsert_continuity_record(record, grant.fencing_token)
1929                .await
1930        {
1931            let mut entries = self.entries.write().await;
1932            if let Some(entry) = entries.get_mut(identity) {
1933                entry.state = IdentityLifecycleState::Broken;
1934                entry.lease = None;
1935            }
1936            drop(entries);
1937            self.emit_event(
1938                identity,
1939                IdentityEvent::StateChanged {
1940                    identity: identity.clone(),
1941                    new_state: IdentityLifecycleState::Broken,
1942                },
1943            )
1944            .await;
1945            return Err(IdentityRuntimeError::Store(err));
1946        }
1947
1948        let mut entries = self.entries.write().await;
1949        let entry = entries
1950            .get_mut(identity)
1951            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
1952        if entry.state != IdentityLifecycleState::Active {
1953            return Err(IdentityRuntimeError::InvalidState {
1954                identity: identity.clone(),
1955                state: entry.state,
1956                operation: "refresh_active_restore_grant",
1957            });
1958        }
1959        entry.lease = Some(Self::lease_entry_from_grant(grant));
1960        drop(entries);
1961        self.emit_event(
1962            identity,
1963            IdentityEvent::LeaseUpdated {
1964                identity: identity.clone(),
1965                fencing_token: grant.fencing_token,
1966            },
1967        )
1968        .await;
1969        Ok(())
1970    }
1971
1972    async fn restore_broken_entry_with_fenced_store(
1973        &self,
1974        identity: &AgentIdentity,
1975        mut entry: IdentityEntry,
1976        grant: &LeaseGrant,
1977    ) {
1978        entry.state = IdentityLifecycleState::Broken;
1979        entry.lease = None;
1980        if let Some(record) = entry.continuity.as_ref()
1981            && let Err(err) = self
1982                .continuity_store
1983                .upsert_continuity_record(record, grant.fencing_token)
1984                .await
1985        {
1986            tracing::warn!(
1987                %identity,
1988                error = %err,
1989                "failed to preserve fenced continuity record for broken identity"
1990            );
1991        }
1992        self.restore_entry(identity, entry).await;
1993    }
1994
1995    async fn mark_rebind_failure_broken(
1996        &self,
1997        identity: &AgentIdentity,
1998        mut entry: IdentityEntry,
1999        grant: &LeaseGrant,
2000        rebound_record: &ContinuityRecord,
2001    ) {
2002        entry.state = IdentityLifecycleState::Broken;
2003        entry.lease = None;
2004        entry.checkpoint_version = rebound_record.checkpoint_version;
2005        entry.continuity = Some(rebound_record.clone());
2006        if let Err(err) = self
2007            .continuity_store
2008            .upsert_continuity_record(rebound_record, grant.fencing_token)
2009            .await
2010        {
2011            tracing::warn!(
2012                %identity,
2013                session_id = %rebound_record.session_id,
2014                error = %err,
2015                "failed to preserve rebound continuity after live respawn rebind failure"
2016            );
2017        }
2018        self.restore_entry(identity, entry).await;
2019    }
2020
2021    async fn restore_entry_after_reset_bridge_failure(
2022        &self,
2023        identity: &AgentIdentity,
2024        entry: IdentityEntry,
2025        grant: &LeaseGrant,
2026        force_broken: bool,
2027    ) -> Option<ContinuityStoreError> {
2028        let delete_error = if entry.continuity.is_none() {
2029            self.continuity_store
2030                .delete_continuity_record(identity, grant.fencing_token)
2031                .await
2032                .err()
2033        } else {
2034            None
2035        };
2036        if force_broken || delete_error.is_some() {
2037            self.restore_broken_entry_with_fenced_store(identity, entry, grant)
2038                .await;
2039        } else {
2040            self.restore_entry_with_grant(identity, entry, grant).await;
2041        }
2042        delete_error
2043    }
2044
2045    async fn restore_continuity_after_materialize_failure(
2046        &self,
2047        identity: &AgentIdentity,
2048        previous: Option<&ContinuityRecord>,
2049        grant: &LeaseGrant,
2050    ) -> Option<ContinuityStoreError> {
2051        match previous {
2052            Some(record) => self
2053                .continuity_store
2054                .upsert_continuity_record(record, grant.fencing_token)
2055                .await
2056                .err(),
2057            None => self
2058                .continuity_store
2059                .delete_continuity_record(identity, grant.fencing_token)
2060                .await
2061                .err(),
2062        }
2063    }
2064
2065    async fn unregister_bridge_session_runtime_states(
2066        bridge: &dyn SessionBridge,
2067        session_ids: &[SessionId],
2068    ) -> Option<String> {
2069        let mut errors = Vec::new();
2070        let mut seen = BTreeSet::new();
2071        for session_id in session_ids {
2072            if !seen.insert(session_id.to_string()) {
2073                continue;
2074            }
2075            if let Err(err) = bridge.unregister_session_runtime_state(session_id).await {
2076                errors.push(format!("{session_id}: {err}"));
2077            }
2078        }
2079        (!errors.is_empty()).then(|| errors.join("; "))
2080    }
2081
2082    async fn advance_existing_continuity_fence(
2083        &self,
2084        identity: &AgentIdentity,
2085        entry: &IdentityEntry,
2086        grant: &LeaseGrant,
2087    ) -> Result<(), IdentityRuntimeError> {
2088        if let Some(record) = entry.continuity.as_ref() {
2089            self.continuity_store
2090                .upsert_continuity_record(record, grant.fencing_token)
2091                .await
2092                .map_err(IdentityRuntimeError::Store)?;
2093        }
2094        let _ = identity;
2095        Ok(())
2096    }
2097
2098    async fn refresh_existing_session_runtime_state(
2099        &self,
2100        identity: &AgentIdentity,
2101        record: &ContinuityRecord,
2102        grant: &LeaseGrant,
2103    ) -> Result<CheckpointVersion, IdentityRuntimeError> {
2104        let Some(bridge) = self.bridge.as_ref() else {
2105            return Ok(record.checkpoint_version);
2106        };
2107        bridge
2108            .register_session_runtime_state(
2109                &record.session_id,
2110                identity,
2111                record.generation,
2112                record.checkpoint_version,
2113                grant.fencing_token,
2114            )
2115            .await
2116            .map_err(|err| {
2117                IdentityRuntimeError::Internal(format!(
2118                    "bridge refresh session runtime state: {err}"
2119                ))
2120            })
2121    }
2122
2123    // -----------------------------------------------------------------------
2124    // Delivery: send() — REQ-01, REQ-03
2125    // -----------------------------------------------------------------------
2126
2127    /// Send conversational content to an addressable identity.
2128    ///
2129    /// Enforces:
2130    /// - Identity must be registered and active
2131    /// - Identity must be Addressable (REQ-03)
2132    /// - Lease must be held (INV-01)
2133    /// - Lease must not be lost (INV-02)
2134    ///
2135    /// Returns the fencing token for the delivery (caller uses it for checkpoint).
2136    pub async fn send(
2137        &self,
2138        identity: &AgentIdentity,
2139        content: &meerkat_core::ContentInput,
2140    ) -> Result<FencingToken, IdentityRuntimeError> {
2141        self.send_with_mode(identity, content, HandlingMode::Queue)
2142            .await
2143    }
2144
2145    /// Send conversational content using an explicit turn handling mode.
2146    ///
2147    /// This is the identity-first counterpart to the mob member send path used
2148    /// by the console. Ordinary API callers can keep using [`Self::send`],
2149    /// which preserves queue semantics.
2150    pub async fn send_with_mode(
2151        &self,
2152        identity: &AgentIdentity,
2153        content: &meerkat_core::ContentInput,
2154        handling_mode: HandlingMode,
2155    ) -> Result<FencingToken, IdentityRuntimeError> {
2156        let should_materialize = {
2157            let entries = self.entries.read().await;
2158            let entry = entries
2159                .get(identity)
2160                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2161
2162            // REQ-03: reject send to InternalOnly
2163            if entry.spec.addressability == AgentAddressability::InternalOnly {
2164                return Err(IdentityRuntimeError::NotAddressable(NotAddressable {
2165                    identity: identity.clone(),
2166                    addressability: entry.spec.addressability,
2167                }));
2168            }
2169            entry.state == IdentityLifecycleState::Dormant
2170                || entry.state == IdentityLifecycleState::Uninitialized
2171        };
2172        if should_materialize {
2173            self.materialize(identity).await?;
2174        }
2175        // Live steers are latency-sensitive operator input for an already
2176        // active turn. Ordinary sends may hydrate the reachable topology first,
2177        // but a steer must reach the current session boundary before the tool
2178        // turn resumes; background/full-fleet materialization owns the peers.
2179        if handling_mode != HandlingMode::Steer {
2180            self.materialize_reachable_peers(identity).await?;
2181        }
2182
2183        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2184        let _lifecycle_guard = lifecycle_lock.lock().await;
2185        {
2186            let entries = self.entries.read().await;
2187            let entry = entries
2188                .get(identity)
2189                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2190            if entry.state != IdentityLifecycleState::Active {
2191                return Err(IdentityRuntimeError::InvalidState {
2192                    identity: identity.clone(),
2193                    state: entry.state,
2194                    operation: "send",
2195                });
2196            }
2197        }
2198
2199        let mut token = self.ensure_active_lease(identity).await?;
2200        let runtime_id = {
2201            let entries = self.entries.read().await;
2202            let entry = entries
2203                .get(identity)
2204                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2205            entry
2206                .continuity
2207                .as_ref()
2208                .map(|c| c.agent_runtime_id.clone())
2209        };
2210
2211        // Deliver through the session bridge when available.
2212        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id) {
2213            let delivered_session_id = bridge
2214                .deliver_with_mode(rid, content, handling_mode)
2215                .await
2216                .map_err(|e| IdentityRuntimeError::Internal(format!("bridge deliver: {e}")))?;
2217            if let Some(rebound_token) = self
2218                .reconcile_delivered_session_locked(identity, delivered_session_id)
2219                .await?
2220            {
2221                token = rebound_token;
2222            }
2223        }
2224
2225        Ok(token)
2226    }
2227
2228    // -----------------------------------------------------------------------
2229    // Delivery: dispatch() — REQ-02
2230    // -----------------------------------------------------------------------
2231
2232    /// Dispatch internal content to any identity (Addressable or InternalOnly).
2233    ///
2234    /// Enforces:
2235    /// - Identity must be registered and active
2236    /// - Lease must be held (INV-01)
2237    /// - Lease must not be lost (INV-02)
2238    ///
2239    /// Returns (fencing_token, is_durable) where is_durable indicates whether
2240    /// the dispatch is backed by a runtime_store (REQ-04).
2241    pub async fn dispatch(
2242        &self,
2243        identity: &AgentIdentity,
2244        input: &DispatchInput,
2245    ) -> Result<(FencingToken, bool), IdentityRuntimeError> {
2246        let should_materialize = {
2247            let entries = self.entries.read().await;
2248            let entry = entries
2249                .get(identity)
2250                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2251            entry.state == IdentityLifecycleState::Dormant
2252                || entry.state == IdentityLifecycleState::Uninitialized
2253        };
2254        if should_materialize {
2255            self.materialize(identity).await?;
2256        }
2257        self.materialize_reachable_peers(identity).await?;
2258
2259        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2260        let _lifecycle_guard = lifecycle_lock.lock().await;
2261        {
2262            let entries = self.entries.read().await;
2263            let entry = entries
2264                .get(identity)
2265                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2266            if entry.state != IdentityLifecycleState::Active {
2267                return Err(IdentityRuntimeError::InvalidState {
2268                    identity: identity.clone(),
2269                    state: entry.state,
2270                    operation: "dispatch",
2271                });
2272            }
2273        }
2274
2275        let mut token = self.ensure_active_lease(identity).await?;
2276        let (is_durable, runtime_id) = {
2277            let entries = self.entries.read().await;
2278            let entry = entries
2279                .get(identity)
2280                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2281            // REQ-04: durability depends on runtime_store
2282            let is_durable = entry.has_runtime_store;
2283
2284            let runtime_id = entry
2285                .continuity
2286                .as_ref()
2287                .map(|c| c.agent_runtime_id.clone());
2288
2289            (is_durable, runtime_id)
2290        };
2291
2292        // Deliver through the session bridge when available.
2293        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id) {
2294            let delivered_session_id = bridge
2295                .deliver(rid, &input.content)
2296                .await
2297                .map_err(|e| IdentityRuntimeError::Internal(format!("bridge dispatch: {e}")))?;
2298            if let Some(rebound_token) = self
2299                .reconcile_delivered_session_locked(identity, delivered_session_id)
2300                .await?
2301            {
2302                token = rebound_token;
2303            }
2304        }
2305
2306        Ok((token, is_durable))
2307    }
2308
2309    // -----------------------------------------------------------------------
2310    // Status: status() — REQ-07
2311    // -----------------------------------------------------------------------
2312
2313    /// Return the full identity status for the given identity.
2314    pub async fn status(
2315        &self,
2316        identity: &AgentIdentity,
2317    ) -> Result<IdentityStatus, IdentityRuntimeError> {
2318        let entries = self.entries.read().await;
2319        let entry = entries
2320            .get(identity)
2321            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2322
2323        let lease_info = entry.lease.as_ref().map(|l| LeaseInfo {
2324            fencing_token: l.fencing_token,
2325            ttl_remaining: l.ttl_remaining(),
2326            healthy: l.is_healthy(),
2327        });
2328
2329        let continuity_health = Some(ContinuityHealth {
2330            store_reachable: true, // tracked per-store in production
2331            durability_policy: self.durability_policy.clone(),
2332            last_checkpoint_version: if entry.checkpoint_version.get() > 0 {
2333                Some(entry.checkpoint_version)
2334            } else {
2335                None
2336            },
2337        });
2338
2339        Ok(IdentityStatus {
2340            identity: identity.clone(),
2341            state: entry.state,
2342            agent_runtime_id: entry
2343                .continuity
2344                .as_ref()
2345                .map(|c| c.agent_runtime_id.clone()),
2346            session_id: entry.continuity.as_ref().map(|c| c.session_id.clone()),
2347            profile: Some(entry.spec.profile.clone()),
2348            runtime_mode: entry.spec.runtime_mode_override,
2349            addressability: entry.spec.addressability,
2350            display_name: entry.spec.display_name.clone(),
2351            labels: entry.spec.labels.clone(),
2352            generation: entry.continuity.as_ref().map(|c| c.generation),
2353            checkpoint_version: if entry.checkpoint_version.get() > 0 {
2354                Some(entry.checkpoint_version)
2355            } else {
2356                None
2357            },
2358            lease: lease_info,
2359            continuity_health,
2360        })
2361    }
2362
2363    /// Return statuses for every registered identity without materializing
2364    /// dormant members.
2365    pub async fn statuses(&self) -> Vec<IdentityStatus> {
2366        let identities = self
2367            .entries
2368            .read()
2369            .await
2370            .keys()
2371            .cloned()
2372            .collect::<Vec<_>>();
2373        let mut statuses = Vec::with_capacity(identities.len());
2374        for identity in identities {
2375            if let Ok(status) = self.status(&identity).await {
2376                statuses.push(status);
2377            }
2378        }
2379        statuses
2380    }
2381
2382    // -----------------------------------------------------------------------
2383    // Lifecycle: retire() — REQ-08
2384    // -----------------------------------------------------------------------
2385
2386    /// Retire an identity. Validates lease ownership and retires the mob member.
2387    pub async fn retire(
2388        &self,
2389        identity: &AgentIdentity,
2390    ) -> Result<FencingToken, IdentityRuntimeError> {
2391        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2392        let _lifecycle_guard = lifecycle_lock.lock().await;
2393        self.ensure_active_lease(identity).await?;
2394        let registered_entry = self
2395            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Retiring)
2396            .await?;
2397        let _previous_token = match Self::check_lease(&registered_entry) {
2398            Ok(token) => token,
2399            Err(err) => {
2400                self.restore_entry(identity, registered_entry).await;
2401                return Err(err);
2402            }
2403        };
2404        let runtime_id = registered_entry
2405            .continuity
2406            .as_ref()
2407            .map(|c| c.agent_runtime_id.clone());
2408        let session_id = registered_entry
2409            .continuity
2410            .as_ref()
2411            .map(|c| c.session_id.clone());
2412
2413        let acquire_result = match self
2414            .lease_provider
2415            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
2416            .await
2417        {
2418            Ok(result) => result,
2419            Err(err) => {
2420                self.restore_entry(identity, registered_entry).await;
2421                return Err(IdentityRuntimeError::Lease(err));
2422            }
2423        };
2424
2425        let grant = match acquire_result.get(identity) {
2426            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
2427            _ => {
2428                self.restore_entry(identity, registered_entry).await;
2429                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
2430            }
2431        };
2432        if let Err(err) = self
2433            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
2434            .await
2435        {
2436            let mut broken_entry = registered_entry;
2437            broken_entry.state = IdentityLifecycleState::Broken;
2438            broken_entry.lease = None;
2439            self.restore_entry(identity, broken_entry).await;
2440            return Err(err);
2441        }
2442
2443        // Retire the mob member through the session bridge when available.
2444        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id)
2445            && let Err(err) = bridge.retire_member(rid).await
2446        {
2447            self.restore_entry_with_grant(identity, registered_entry, &grant)
2448                .await;
2449            return Err(IdentityRuntimeError::Internal(format!(
2450                "bridge retire: {err}"
2451            )));
2452        }
2453        if let (Some(bridge), Some(session_id)) = (&self.bridge, &session_id)
2454            && let Err(err) = bridge.unregister_session_runtime_state(session_id).await
2455        {
2456            let mut broken_entry = registered_entry;
2457            broken_entry.state = IdentityLifecycleState::Broken;
2458            broken_entry.lease = None;
2459            self.restore_entry(identity, broken_entry).await;
2460            return Err(IdentityRuntimeError::Internal(format!(
2461                "bridge unregister retired session: {err}"
2462            )));
2463        }
2464
2465        Ok(grant.fencing_token)
2466    }
2467
2468    // -----------------------------------------------------------------------
2469    // Lifecycle: respawn() — REQ-09
2470    // -----------------------------------------------------------------------
2471
2472    /// Respawn: non-destructive recovery.
2473    ///
2474    /// 1. Fence the current owner
2475    /// 2. Attempt final checkpoint
2476    /// 3. Reactivate from authoritative continuity with same record + runtime ID
2477    /// 4. ContinuityGeneration does NOT advance
2478    pub async fn respawn(
2479        &self,
2480        identity: &AgentIdentity,
2481    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
2482        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2483        let _lifecycle_guard = lifecycle_lock.lock().await;
2484        let registered_entry = self
2485            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Suspended)
2486            .await?;
2487
2488        // Fence the old owner by re-acquiring the lease
2489        let acquire_result = match self
2490            .lease_provider
2491            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
2492            .await
2493        {
2494            Ok(result) => result,
2495            Err(err) => {
2496                self.restore_entry(identity, registered_entry).await;
2497                return Err(IdentityRuntimeError::Lease(err));
2498            }
2499        };
2500
2501        let grant = match acquire_result.get(identity) {
2502            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
2503            _ => {
2504                self.restore_entry(identity, registered_entry).await;
2505                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
2506            }
2507        };
2508
2509        if let Err(err) = self
2510            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
2511            .await
2512        {
2513            let mut broken_entry = registered_entry;
2514            broken_entry.state = IdentityLifecycleState::Broken;
2515            broken_entry.lease = None;
2516            self.restore_entry(identity, broken_entry).await;
2517            return Err(err);
2518        }
2519
2520        // Resolve current continuity state
2521        let resolved = match self
2522            .continuity_store
2523            .resolve_many(std::slice::from_ref(identity))
2524            .await
2525        {
2526            Ok(resolved) => resolved,
2527            Err(err) => {
2528                self.restore_entry_with_grant(identity, registered_entry, &grant)
2529                    .await;
2530                return Err(IdentityRuntimeError::Store(err));
2531            }
2532        };
2533
2534        let record = match resolved.get(identity) {
2535            Some(super::types::ContinuityResolveState::Ready { record }) => record.clone(),
2536            Some(super::types::ContinuityResolveState::Broken { failure }) => {
2537                self.restore_entry_with_grant(identity, registered_entry, &grant)
2538                    .await;
2539                return Err(IdentityRuntimeError::Internal(format!(
2540                    "broken continuity for {identity}: {}",
2541                    failure.detail
2542                )));
2543            }
2544            Some(super::types::ContinuityResolveState::Uninitialized) => {
2545                self.restore_entry_with_grant(identity, registered_entry, &grant)
2546                    .await;
2547                return Err(IdentityRuntimeError::Internal(format!(
2548                    "cannot respawn uninitialized identity {identity}"
2549                )));
2550            }
2551            None => {
2552                self.restore_entry_with_grant(identity, registered_entry, &grant)
2553                    .await;
2554                return Err(IdentityRuntimeError::Store(
2555                    ContinuityStoreError::NotFound {
2556                        identity: identity.clone(),
2557                    },
2558                ));
2559            }
2560        };
2561
2562        let effective_checkpoint_version = match self
2563            .refresh_existing_session_runtime_state(identity, &record, &grant)
2564            .await
2565        {
2566            Ok(version) => version,
2567            Err(err) => {
2568                self.restore_entry_with_grant(identity, registered_entry, &grant)
2569                    .await;
2570                return Err(err);
2571            }
2572        };
2573        let mut record = record;
2574        record.checkpoint_version = effective_checkpoint_version;
2575
2576        // Update runtime state: same record, new lease, back to Active
2577        let mut entries = self.entries.write().await;
2578        if let Some(entry) = entries.get_mut(identity) {
2579            entry.continuity = Some(record.clone());
2580            entry.lease = Some(LeaseEntry {
2581                fencing_token: grant.fencing_token,
2582                ttl: grant.ttl,
2583                acquired_at: Instant::now(),
2584            });
2585            entry.state = IdentityLifecycleState::Active;
2586            entry.checkpoint_version = record.checkpoint_version;
2587        }
2588
2589        Ok(record)
2590    }
2591
2592    /// Rebind continuity to the concrete session created by a lower-level
2593    /// member respawn. This keeps identity-first status aligned when a control
2594    /// surface refreshes the mob member outside the identity runtime bridge.
2595    pub async fn rebind_session_after_live_respawn(
2596        &self,
2597        identity: &AgentIdentity,
2598        session_id: SessionId,
2599    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
2600        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2601        let _lifecycle_guard = lifecycle_lock.lock().await;
2602        self.rebind_session_after_live_respawn_locked(identity, session_id)
2603            .await
2604    }
2605
2606    async fn reconcile_delivered_session_locked(
2607        &self,
2608        identity: &AgentIdentity,
2609        delivered_session_id: SessionId,
2610    ) -> Result<Option<FencingToken>, IdentityRuntimeError> {
2611        let current_session_id = {
2612            let entries = self.entries.read().await;
2613            let entry = entries
2614                .get(identity)
2615                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2616            entry
2617                .continuity
2618                .as_ref()
2619                .map(|record| record.session_id.clone())
2620        };
2621
2622        let Some(current_session_id) = current_session_id else {
2623            return Ok(None);
2624        };
2625        if current_session_id == delivered_session_id {
2626            return Ok(None);
2627        }
2628
2629        tracing::warn!(
2630            %identity,
2631            old_session_id = %current_session_id,
2632            new_session_id = %delivered_session_id,
2633            "identity bridge delivery returned a rotated session; rebinding continuity"
2634        );
2635        self.rebind_session_after_live_respawn_locked(identity, delivered_session_id)
2636            .await?;
2637
2638        let entries = self.entries.read().await;
2639        let entry = entries
2640            .get(identity)
2641            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
2642        Ok(entry.lease.as_ref().map(|lease| lease.fencing_token))
2643    }
2644
2645    async fn rebind_session_after_live_respawn_locked(
2646        &self,
2647        identity: &AgentIdentity,
2648        session_id: SessionId,
2649    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
2650        let registered_entry = self
2651            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Suspended)
2652            .await?;
2653
2654        let acquire_result = match self
2655            .lease_provider
2656            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
2657            .await
2658        {
2659            Ok(result) => result,
2660            Err(err) => {
2661                self.restore_entry(identity, registered_entry).await;
2662                return Err(IdentityRuntimeError::Lease(err));
2663            }
2664        };
2665
2666        let grant = match acquire_result.get(identity) {
2667            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
2668            _ => {
2669                self.restore_entry(identity, registered_entry).await;
2670                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
2671            }
2672        };
2673
2674        let mut record = match registered_entry.continuity.as_ref() {
2675            Some(record) => record.clone(),
2676            None => {
2677                self.restore_entry_with_grant(identity, registered_entry, &grant)
2678                    .await;
2679                return Err(IdentityRuntimeError::UnknownIdentity(identity.clone()));
2680            }
2681        };
2682        if let Err(err) = self
2683            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
2684            .await
2685        {
2686            if let Some(bridge) = self.bridge.as_ref()
2687                && let Err(unregister_err) =
2688                    bridge.unregister_session_runtime_state(&session_id).await
2689            {
2690                tracing::warn!(
2691                    %identity,
2692                    session_id = %session_id,
2693                    error = %unregister_err,
2694                    "failed to unregister rebound session after continuity fence failure"
2695                );
2696            }
2697            let mut broken_entry = registered_entry;
2698            broken_entry.state = IdentityLifecycleState::Broken;
2699            broken_entry.lease = None;
2700            self.restore_entry(identity, broken_entry).await;
2701            return Err(err);
2702        }
2703        let previous_session_id = record.session_id.clone();
2704        record.session_id = session_id;
2705
2706        if let Err(err) = self
2707            .continuity_store
2708            .upsert_continuity_record(&record, grant.fencing_token)
2709            .await
2710        {
2711            if let Some(bridge) = self.bridge.as_ref()
2712                && let Err(unregister_err) = bridge
2713                    .unregister_session_runtime_state(&record.session_id)
2714                    .await
2715            {
2716                tracing::warn!(
2717                    %identity,
2718                    session_id = %record.session_id,
2719                    error = %unregister_err,
2720                    "failed to unregister rebound session after continuity upsert failure"
2721                );
2722            }
2723            self.mark_rebind_failure_broken(identity, registered_entry, &grant, &record)
2724                .await;
2725            return Err(IdentityRuntimeError::Store(err));
2726        }
2727
2728        if let Some(bridge) = self.bridge.as_ref() {
2729            match bridge
2730                .register_session_runtime_state(
2731                    &record.session_id,
2732                    identity,
2733                    record.generation,
2734                    record.checkpoint_version,
2735                    grant.fencing_token,
2736                )
2737                .await
2738            {
2739                Ok(version) => record.checkpoint_version = version,
2740                Err(err) => {
2741                    if let Err(unregister_err) = bridge
2742                        .unregister_session_runtime_state(&record.session_id)
2743                        .await
2744                    {
2745                        tracing::warn!(
2746                            %identity,
2747                            session_id = %record.session_id,
2748                            error = %unregister_err,
2749                            "failed to unregister rebound session after bridge register failure"
2750                        );
2751                    }
2752                    self.mark_rebind_failure_broken(identity, registered_entry, &grant, &record)
2753                        .await;
2754                    return Err(IdentityRuntimeError::Internal(format!(
2755                        "bridge rebind respawned session runtime state: {err}"
2756                    )));
2757                }
2758            }
2759            if previous_session_id != record.session_id
2760                && let Err(err) = bridge
2761                    .unregister_session_runtime_state(&previous_session_id)
2762                    .await
2763            {
2764                tracing::warn!(
2765                    %identity,
2766                    session_id = %previous_session_id,
2767                    error = %err,
2768                    "failed to unregister previous session after live respawn rebind"
2769                );
2770            }
2771        }
2772
2773        if let Err(err) = self
2774            .continuity_store
2775            .upsert_continuity_record(&record, grant.fencing_token)
2776            .await
2777        {
2778            if let Some(bridge) = self.bridge.as_ref()
2779                && let Err(unregister_err) = bridge
2780                    .unregister_session_runtime_state(&record.session_id)
2781                    .await
2782            {
2783                tracing::warn!(
2784                    %identity,
2785                    session_id = %record.session_id,
2786                    error = %unregister_err,
2787                    "failed to unregister rebound session after final continuity upsert failure"
2788                );
2789            }
2790            self.mark_rebind_failure_broken(identity, registered_entry, &grant, &record)
2791                .await;
2792            return Err(IdentityRuntimeError::Store(err));
2793        }
2794
2795        self.register(
2796            registered_entry.spec,
2797            IdentityLifecycleState::Active,
2798            Some(record.clone()),
2799            Some(grant),
2800        )
2801        .await;
2802        Ok(record)
2803    }
2804
2805    // -----------------------------------------------------------------------
2806    // Lifecycle: reset() — REQ-10
2807    // -----------------------------------------------------------------------
2808
2809    /// Reset: destructive continuity reset.
2810    ///
2811    /// 1. Fence old owner
2812    /// 2. Advance ContinuityGeneration
2813    /// 3. Create fresh continuity under the same AgentIdentity
2814    /// 4. Old-owner late writes rejected by stale fencing token
2815    pub async fn reset(
2816        &self,
2817        identity: &AgentIdentity,
2818    ) -> Result<ContinuityRecord, IdentityRuntimeError> {
2819        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
2820        let _lifecycle_guard = lifecycle_lock.lock().await;
2821        let registered_entry = self
2822            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Suspended)
2823            .await?;
2824
2825        // INV-05: fence the old owner first
2826        let acquire_result = match self
2827            .lease_provider
2828            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
2829            .await
2830        {
2831            Ok(result) => result,
2832            Err(err) => {
2833                self.restore_entry(identity, registered_entry).await;
2834                return Err(IdentityRuntimeError::Lease(err));
2835            }
2836        };
2837
2838        let grant = match acquire_result.get(identity) {
2839            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
2840            _ => {
2841                self.restore_entry(identity, registered_entry).await;
2842                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
2843            }
2844        };
2845        if let Err(err) = self
2846            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
2847            .await
2848        {
2849            let mut broken_entry = registered_entry;
2850            broken_entry.state = IdentityLifecycleState::Broken;
2851            broken_entry.lease = None;
2852            self.restore_entry(identity, broken_entry).await;
2853            return Err(err);
2854        }
2855
2856        // Resolve to get current generation
2857        let resolved = match self
2858            .continuity_store
2859            .resolve_many(std::slice::from_ref(identity))
2860            .await
2861        {
2862            Ok(resolved) => resolved,
2863            Err(err) => {
2864                self.restore_entry_with_grant(identity, registered_entry.clone(), &grant)
2865                    .await;
2866                return Err(IdentityRuntimeError::Store(err));
2867            }
2868        };
2869
2870        let current_gen = match resolved.get(identity) {
2871            Some(super::types::ContinuityResolveState::Ready { record }) => record.generation,
2872            Some(super::types::ContinuityResolveState::Uninitialized) => {
2873                ContinuityGeneration::new(0)
2874            }
2875            _ => ContinuityGeneration::new(0),
2876        };
2877
2878        // Advance generation
2879        let new_gen = ContinuityGeneration::new(current_gen.get() + 1);
2880        let new_session_id = meerkat_core::types::SessionId::new();
2881        let new_runtime_id = AgentRuntimeId::parse(&format!("rt:{identity}:{}", new_gen.get()))
2882            .map_err(|e| {
2883                IdentityRuntimeError::Internal(format!("failed to mint runtime id: {e}"))
2884            })?;
2885
2886        let new_record = ContinuityRecord {
2887            identity: identity.clone(),
2888            agent_runtime_id: new_runtime_id,
2889            session_id: new_session_id,
2890            generation: new_gen,
2891            checkpoint_version: CheckpointVersion::new(0),
2892        };
2893
2894        // Bridge: retire old mob member and create fresh session for the new identity.
2895        if let Some(bridge) = &self.bridge {
2896            if let Err(err) = self
2897                .continuity_store
2898                .upsert_continuity_record(&new_record, grant.fencing_token)
2899                .await
2900            {
2901                self.restore_entry_with_grant(identity, registered_entry, &grant)
2902                    .await;
2903                return Err(IdentityRuntimeError::Store(err));
2904            }
2905
2906            let old_runtime_id = registered_entry
2907                .continuity
2908                .as_ref()
2909                .map(|c| c.agent_runtime_id.clone());
2910            let old_session_id = registered_entry
2911                .continuity
2912                .as_ref()
2913                .map(|c| c.session_id.clone());
2914
2915            let spec = registered_entry.spec.clone();
2916            let draft = super::types::AgentBuildDraft {
2917                model: None,
2918                system_prompt: None,
2919                additional_instructions: spec.additional_instructions.clone(),
2920                labels: spec.labels.clone(),
2921                app_context: spec.context.clone(),
2922                external_tools: Vec::new(),
2923                local_external_tools: Default::default(),
2924            };
2925            let session_id = bridge
2926                .create_session(
2927                    identity,
2928                    &new_record.agent_runtime_id,
2929                    &spec,
2930                    &draft,
2931                    &new_record.session_id,
2932                )
2933                .await
2934                .map_err(|e| {
2935                    IdentityRuntimeError::Internal(format!(
2936                        "bridge create_session after reset: {e}"
2937                    ))
2938                });
2939            let session_id = match session_id {
2940                Ok(session_id) => session_id,
2941                Err(err) => {
2942                    let cleanup_error = bridge
2943                        .retire_member(&new_record.agent_runtime_id)
2944                        .await
2945                        .err();
2946                    let delete_error = self
2947                        .restore_entry_after_reset_bridge_failure(
2948                            identity,
2949                            registered_entry.clone(),
2950                            &grant,
2951                            cleanup_error.is_some(),
2952                        )
2953                        .await;
2954                    if cleanup_error.is_some() || delete_error.is_some() {
2955                        return Err(IdentityRuntimeError::Internal(format!(
2956                            "{err}{}{}",
2957                            cleanup_error
2958                                .as_ref()
2959                                .map(|e| format!("; cleanup retire failed: {e}"))
2960                                .unwrap_or_default(),
2961                            delete_error
2962                                .as_ref()
2963                                .map(|e| format!("; tentative continuity cleanup failed: {e}"))
2964                                .unwrap_or_default()
2965                        )));
2966                    }
2967                    return Err(err);
2968                }
2969            };
2970            // Update the record with the actual session ID
2971            let mut new_record = new_record;
2972            new_record.session_id = session_id;
2973
2974            if let Err(err) = self
2975                .continuity_store
2976                .upsert_continuity_record(&new_record, grant.fencing_token)
2977                .await
2978            {
2979                let unregister_error = Self::unregister_bridge_session_runtime_states(
2980                    bridge.as_ref(),
2981                    std::slice::from_ref(&new_record.session_id),
2982                )
2983                .await;
2984                let cleanup_error = bridge
2985                    .retire_member(&new_record.agent_runtime_id)
2986                    .await
2987                    .err();
2988                let delete_error = self
2989                    .restore_entry_after_reset_bridge_failure(
2990                        identity,
2991                        registered_entry.clone(),
2992                        &grant,
2993                        unregister_error.is_some() || cleanup_error.is_some(),
2994                    )
2995                    .await;
2996                if unregister_error.is_some() || cleanup_error.is_some() || delete_error.is_some() {
2997                    return Err(IdentityRuntimeError::Internal(format!(
2998                        "continuity upsert actual session after reset: {err}{}{}{}",
2999                        unregister_error
3000                            .as_ref()
3001                            .map(|e| format!("; unregister session failed: {e}"))
3002                            .unwrap_or_default(),
3003                        cleanup_error
3004                            .as_ref()
3005                            .map(|e| format!("; cleanup retire failed: {e}"))
3006                            .unwrap_or_default(),
3007                        delete_error
3008                            .as_ref()
3009                            .map(|e| format!("; tentative continuity cleanup failed: {e}"))
3010                            .unwrap_or_default(),
3011                    )));
3012                }
3013                return Err(IdentityRuntimeError::Store(err));
3014            }
3015
3016            let register_result = bridge
3017                .register_session_runtime_state(
3018                    &new_record.session_id,
3019                    identity,
3020                    new_record.generation,
3021                    new_record.checkpoint_version,
3022                    grant.fencing_token,
3023                )
3024                .await;
3025            let effective_checkpoint_version = match register_result {
3026                Ok(version) => version,
3027                Err(err) => {
3028                    let unregister_error = Self::unregister_bridge_session_runtime_states(
3029                        bridge.as_ref(),
3030                        std::slice::from_ref(&new_record.session_id),
3031                    )
3032                    .await;
3033                    let cleanup_error = bridge
3034                        .retire_member(&new_record.agent_runtime_id)
3035                        .await
3036                        .err();
3037                    let mut detail =
3038                        format!("bridge register actual session runtime state after reset: {err}");
3039                    if let Some(unregister_error) = unregister_error.as_ref() {
3040                        detail
3041                            .push_str(&format!("; unregister session failed: {unregister_error}"));
3042                    }
3043                    if let Some(cleanup_error) = cleanup_error.as_ref() {
3044                        detail.push_str(&format!("; cleanup retire failed: {cleanup_error}"));
3045                    }
3046                    let delete_error = self
3047                        .restore_entry_after_reset_bridge_failure(
3048                            identity,
3049                            registered_entry.clone(),
3050                            &grant,
3051                            unregister_error.is_some() || cleanup_error.is_some(),
3052                        )
3053                        .await;
3054                    if let Some(delete_error) = delete_error {
3055                        return Err(IdentityRuntimeError::Internal(format!(
3056                            "{detail}; tentative continuity cleanup failed: {delete_error}"
3057                        )));
3058                    }
3059                    return Err(IdentityRuntimeError::Internal(detail));
3060                }
3061            };
3062            new_record.checkpoint_version = effective_checkpoint_version;
3063
3064            if let Some(old_id) = old_runtime_id.as_ref()
3065                && old_id != &new_record.agent_runtime_id
3066            {
3067                if let Err(err) = bridge.retire_member(old_id).await {
3068                    let unregister_error = bridge
3069                        .unregister_session_runtime_state(&new_record.session_id)
3070                        .await
3071                        .err();
3072                    let cleanup_error = bridge
3073                        .retire_member(&new_record.agent_runtime_id)
3074                        .await
3075                        .err();
3076                    self.restore_broken_entry_with_fenced_store(
3077                        identity,
3078                        registered_entry.clone(),
3079                        &grant,
3080                    )
3081                    .await;
3082                    let detail = format!(
3083                        "bridge retire old member after reset: {err}{}{}",
3084                        unregister_error
3085                            .as_ref()
3086                            .map(|e| format!("; unregister session failed: {e}"))
3087                            .unwrap_or_default(),
3088                        cleanup_error
3089                            .as_ref()
3090                            .map(|e| format!("; cleanup retire failed: {e}"))
3091                            .unwrap_or_default(),
3092                    );
3093                    return Err(IdentityRuntimeError::Internal(detail));
3094                }
3095                if let Some(old_session_id) = old_session_id.as_ref()
3096                    && old_session_id != &new_record.session_id
3097                    && let Err(err) = bridge
3098                        .unregister_session_runtime_state(old_session_id)
3099                        .await
3100                {
3101                    let unregister_error = bridge
3102                        .unregister_session_runtime_state(&new_record.session_id)
3103                        .await
3104                        .err();
3105                    let cleanup_error = bridge
3106                        .retire_member(&new_record.agent_runtime_id)
3107                        .await
3108                        .err();
3109                    self.restore_broken_entry_with_fenced_store(
3110                        identity,
3111                        registered_entry.clone(),
3112                        &grant,
3113                    )
3114                    .await;
3115                    let detail = format!(
3116                        "bridge unregister old session after reset: {err}{}{}",
3117                        unregister_error
3118                            .as_ref()
3119                            .map(|e| format!("; unregister new session failed: {e}"))
3120                            .unwrap_or_default(),
3121                        cleanup_error
3122                            .as_ref()
3123                            .map(|e| format!("; cleanup retire failed: {e}"))
3124                            .unwrap_or_default(),
3125                    );
3126                    return Err(IdentityRuntimeError::Internal(detail));
3127                }
3128            }
3129
3130            if let Err(err) = self
3131                .continuity_store
3132                .upsert_continuity_record(&new_record, grant.fencing_token)
3133                .await
3134            {
3135                let unregister_error = bridge
3136                    .unregister_session_runtime_state(&new_record.session_id)
3137                    .await
3138                    .err();
3139                let cleanup_error = bridge
3140                    .retire_member(&new_record.agent_runtime_id)
3141                    .await
3142                    .err();
3143                let rollback_error = self
3144                    .restore_continuity_after_materialize_failure(
3145                        identity,
3146                        registered_entry.continuity.as_ref(),
3147                        &grant,
3148                    )
3149                    .await;
3150                let mut entries = self.entries.write().await;
3151                if let Some(entry) = entries.get_mut(identity) {
3152                    entry.state = IdentityLifecycleState::Broken;
3153                }
3154                if unregister_error.is_some() || cleanup_error.is_some() || rollback_error.is_some()
3155                {
3156                    return Err(IdentityRuntimeError::Internal(format!(
3157                        "continuity upsert after reset: {err}{}{}{}",
3158                        unregister_error
3159                            .as_ref()
3160                            .map(|e| format!("; unregister session failed: {e}"))
3161                            .unwrap_or_default(),
3162                        cleanup_error
3163                            .as_ref()
3164                            .map(|e| format!("; cleanup retire failed: {e}"))
3165                            .unwrap_or_default(),
3166                        rollback_error
3167                            .as_ref()
3168                            .map(|e| format!("; continuity rollback failed: {e}"))
3169                            .unwrap_or_default()
3170                    )));
3171                }
3172                return Err(IdentityRuntimeError::Store(err));
3173            }
3174
3175            // Update runtime state
3176            let mut entries = self.entries.write().await;
3177            let Some(entry) = entries.get_mut(identity) else {
3178                let _ = bridge
3179                    .unregister_session_runtime_state(&new_record.session_id)
3180                    .await;
3181                let _ = bridge.retire_member(&new_record.agent_runtime_id).await;
3182                if registered_entry.continuity.is_none() {
3183                    let _ = self
3184                        .continuity_store
3185                        .delete_continuity_record(identity, grant.fencing_token)
3186                        .await;
3187                }
3188                return Err(IdentityRuntimeError::UnknownIdentity(identity.clone()));
3189            };
3190            entry.continuity = Some(new_record.clone());
3191            entry.lease = Some(Self::lease_entry_from_grant(&grant));
3192            entry.state = IdentityLifecycleState::Active;
3193            entry.checkpoint_version = new_record.checkpoint_version;
3194            return Ok(new_record);
3195        }
3196
3197        // Persist the new record (fencing token from new lease protects against old writes)
3198        if let Err(err) = self
3199            .continuity_store
3200            .upsert_continuity_record(&new_record, grant.fencing_token)
3201            .await
3202        {
3203            self.restore_entry_with_grant(identity, registered_entry, &grant)
3204                .await;
3205            return Err(IdentityRuntimeError::Store(err));
3206        }
3207
3208        // No bridge — update runtime state only (validation mode)
3209        let mut entries = self.entries.write().await;
3210        let entry = entries
3211            .get_mut(identity)
3212            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3213        entry.continuity = Some(new_record.clone());
3214        entry.lease = Some(Self::lease_entry_from_grant(&grant));
3215        entry.state = IdentityLifecycleState::Active;
3216        entry.checkpoint_version = CheckpointVersion::new(0);
3217
3218        Ok(new_record)
3219    }
3220
3221    // -----------------------------------------------------------------------
3222    // Lifecycle: delete_identity() — REQ-11
3223    // -----------------------------------------------------------------------
3224
3225    /// Delete an identity: removes continuity record.
3226    ///
3227    /// 1. Fence old owner
3228    /// 2. Remove ContinuityRecord
3229    /// 3. Future bootstrap treats identity as Uninitialized
3230    pub async fn delete_identity(
3231        &self,
3232        identity: &AgentIdentity,
3233    ) -> Result<(), IdentityRuntimeError> {
3234        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3235        let _lifecycle_guard = lifecycle_lock.lock().await;
3236        let registered_entry = self
3237            .mark_lifecycle_in_progress(identity, IdentityLifecycleState::Retiring)
3238            .await?;
3239        let runtime_id = registered_entry
3240            .continuity
3241            .as_ref()
3242            .map(|c| c.agent_runtime_id.clone());
3243        let session_id = registered_entry
3244            .continuity
3245            .as_ref()
3246            .map(|c| c.session_id.clone());
3247
3248        // INV-05: fence the old owner first
3249        let acquire_result = match self
3250            .lease_provider
3251            .acquire_leases(std::slice::from_ref(identity), &self.runtime_instance_id)
3252            .await
3253        {
3254            Ok(result) => result,
3255            Err(err) => {
3256                self.restore_entry(identity, registered_entry).await;
3257                return Err(IdentityRuntimeError::Lease(err));
3258            }
3259        };
3260
3261        let grant = match acquire_result.get(identity) {
3262            Some(super::types::LeaseAcquireResult::Acquired(g)) => g.clone(),
3263            _ => {
3264                self.restore_entry(identity, registered_entry).await;
3265                return Err(IdentityRuntimeError::NoActiveLease(identity.clone()));
3266            }
3267        };
3268        if let Err(err) = self
3269            .advance_existing_continuity_fence(identity, &registered_entry, &grant)
3270            .await
3271        {
3272            let mut broken_entry = registered_entry;
3273            broken_entry.state = IdentityLifecycleState::Broken;
3274            broken_entry.lease = None;
3275            self.restore_entry(identity, broken_entry).await;
3276            return Err(err);
3277        }
3278
3279        // Retire the mob member through the session bridge before removing
3280        // the continuity record. This ensures the mob actor is cleaned up.
3281        if let (Some(bridge), Some(rid)) = (&self.bridge, &runtime_id)
3282            && let Err(err) = bridge.retire_member(rid).await
3283        {
3284            self.restore_entry_with_grant(identity, registered_entry, &grant)
3285                .await;
3286            return Err(IdentityRuntimeError::Internal(format!(
3287                "bridge retire before delete: {err}"
3288            )));
3289        }
3290
3291        if let (Some(bridge), Some(session_id)) = (&self.bridge, &session_id)
3292            && let Some(err) = Self::unregister_bridge_session_runtime_states(
3293                bridge.as_ref(),
3294                std::slice::from_ref(session_id),
3295            )
3296            .await
3297        {
3298            self.restore_broken_entry_with_fenced_store(identity, registered_entry, &grant)
3299                .await;
3300            return Err(IdentityRuntimeError::Internal(format!(
3301                "bridge unregister session before delete: {err}"
3302            )));
3303        }
3304
3305        // Remove authoritative continuity record from the store
3306        if let Err(err) = self
3307            .continuity_store
3308            .delete_continuity_record(identity, grant.fencing_token)
3309            .await
3310        {
3311            let mut entries = self.entries.write().await;
3312            if let Some(entry) = entries.get_mut(identity) {
3313                entry.state = IdentityLifecycleState::Broken;
3314            }
3315            return Err(IdentityRuntimeError::Store(err));
3316        }
3317
3318        // Remove from runtime tracking
3319        self.event_channels.write().await.remove(identity);
3320        self.entries.write().await.remove(identity);
3321
3322        Ok(())
3323    }
3324
3325    // -----------------------------------------------------------------------
3326    // Checkpoint — REQ-14, REQ-15, REQ-16, REQ-17
3327    // -----------------------------------------------------------------------
3328
3329    /// Save a checkpoint snapshot. Enforces version ordering and fencing.
3330    pub async fn checkpoint(
3331        &self,
3332        identity: &AgentIdentity,
3333        snapshot: &SessionSnapshot,
3334    ) -> Result<CheckpointVersion, IdentityRuntimeError> {
3335        let lifecycle_lock = self.lifecycle_lock_for(identity).await;
3336        let _lifecycle_guard = lifecycle_lock.lock().await;
3337        {
3338            let entries = self.entries.read().await;
3339            let entry = entries
3340                .get(identity)
3341                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3342            if entry.state != IdentityLifecycleState::Active {
3343                return Err(IdentityRuntimeError::InvalidState {
3344                    identity: identity.clone(),
3345                    state: entry.state,
3346                    operation: "checkpoint",
3347                });
3348            }
3349        }
3350
3351        let token = self.ensure_active_lease(identity).await?;
3352        let (record, new_version) = {
3353            let entries = self.entries.read().await;
3354            let entry = entries
3355                .get(identity)
3356                .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3357            let record = entry
3358                .continuity
3359                .as_ref()
3360                .ok_or_else(|| {
3361                    IdentityRuntimeError::Internal(format!("no continuity record for {identity}"))
3362                })?
3363                .clone();
3364
3365            let new_version = CheckpointVersion::new(entry.checkpoint_version.get() + 1);
3366            (record, new_version)
3367        };
3368
3369        // REQ-15 + REQ-16: store enforces version ordering and fencing
3370        self.continuity_store
3371            .save_session_snapshot(
3372                identity,
3373                &record.session_id,
3374                record.generation,
3375                new_version,
3376                token,
3377                snapshot,
3378            )
3379            .await?;
3380
3381        // Update local checkpoint version
3382        {
3383            let mut entries = self.entries.write().await;
3384            if let Some(entry) = entries.get_mut(identity) {
3385                entry.checkpoint_version = new_version;
3386            }
3387        }
3388
3389        self.emit_event(
3390            identity,
3391            IdentityEvent::CheckpointCompleted {
3392                identity: identity.clone(),
3393                version: new_version,
3394            },
3395        )
3396        .await;
3397
3398        Ok(new_version)
3399    }
3400
3401    // -----------------------------------------------------------------------
3402    // Roster inspection — REQ-32
3403    // -----------------------------------------------------------------------
3404
3405    /// Return all active identities with their specs and status.
3406    pub async fn roster_inspect(
3407        &self,
3408    ) -> BTreeMap<AgentIdentity, (DurableAgentSpec, IdentityStatus)> {
3409        let entries = self.entries.read().await;
3410        let mut result = BTreeMap::new();
3411        for (identity, entry) in entries.iter() {
3412            let lease_info = entry.lease.as_ref().map(|l| LeaseInfo {
3413                fencing_token: l.fencing_token,
3414                ttl_remaining: l.ttl_remaining(),
3415                healthy: l.is_healthy(),
3416            });
3417            let continuity_health = Some(ContinuityHealth {
3418                store_reachable: true,
3419                durability_policy: self.durability_policy.clone(),
3420                last_checkpoint_version: if entry.checkpoint_version.get() > 0 {
3421                    Some(entry.checkpoint_version)
3422                } else {
3423                    None
3424                },
3425            });
3426            let status = IdentityStatus {
3427                identity: identity.clone(),
3428                state: entry.state,
3429                agent_runtime_id: entry
3430                    .continuity
3431                    .as_ref()
3432                    .map(|c| c.agent_runtime_id.clone()),
3433                session_id: entry.continuity.as_ref().map(|c| c.session_id.clone()),
3434                profile: Some(entry.spec.profile.clone()),
3435                runtime_mode: entry.spec.runtime_mode_override,
3436                addressability: entry.spec.addressability,
3437                display_name: entry.spec.display_name.clone(),
3438                labels: entry.spec.labels.clone(),
3439                generation: entry.continuity.as_ref().map(|c| c.generation),
3440                checkpoint_version: if entry.checkpoint_version.get() > 0 {
3441                    Some(entry.checkpoint_version)
3442                } else {
3443                    None
3444                },
3445                lease: lease_info,
3446                continuity_health,
3447            };
3448            result.insert(identity.clone(), (entry.spec.clone(), status));
3449        }
3450        result
3451    }
3452
3453    // -----------------------------------------------------------------------
3454    // Roster uniqueness validation — INV-06
3455    // -----------------------------------------------------------------------
3456
3457    /// Validate that a roster contains no duplicate identities.
3458    pub fn validate_roster_uniqueness(
3459        specs: &[DurableAgentSpec],
3460    ) -> Result<(), IdentityRuntimeError> {
3461        let mut seen = std::collections::BTreeSet::new();
3462        for spec in specs {
3463            if !seen.insert(&spec.identity) {
3464                return Err(IdentityRuntimeError::DuplicateIdentity(
3465                    spec.identity.clone(),
3466                ));
3467            }
3468        }
3469        Ok(())
3470    }
3471
3472    // -----------------------------------------------------------------------
3473    // Accessors for internal state
3474    // -----------------------------------------------------------------------
3475
3476    /// Get the current entries (read-only snapshot).
3477    #[allow(dead_code)]
3478    pub(crate) async fn entries(&self) -> BTreeMap<AgentIdentity, IdentityEntry> {
3479        self.entries.read().await.clone()
3480    }
3481
3482    /// Check if an identity is registered.
3483    pub async fn contains(&self, identity: &AgentIdentity) -> bool {
3484        self.entries.read().await.contains_key(identity)
3485    }
3486
3487    /// Check if an identity is registered AND in Active state.
3488    pub async fn is_active(&self, identity: &AgentIdentity) -> bool {
3489        self.entries
3490            .read()
3491            .await
3492            .get(identity)
3493            .is_some_and(|e| e.state == IdentityLifecycleState::Active)
3494    }
3495
3496    /// Get the continuity store reference.
3497    pub fn continuity_store(&self) -> &Arc<dyn ContinuityStore> {
3498        &self.continuity_store
3499    }
3500
3501    /// Get the lease provider reference.
3502    pub fn lease_provider(&self) -> &Arc<dyn LeaseProvider> {
3503        &self.lease_provider
3504    }
3505
3506    /// Get the runtime instance ID.
3507    pub fn runtime_instance_id(&self) -> &str {
3508        &self.runtime_instance_id
3509    }
3510
3511    /// Get the durability policy.
3512    pub fn durability_policy(&self) -> &DurabilityPolicy {
3513        &self.durability_policy
3514    }
3515
3516    /// Get whether a runtime store is configured.
3517    pub fn has_runtime_store(&self) -> bool {
3518        self.has_runtime_store
3519    }
3520
3521    /// Get the session bridge reference, if configured.
3522    pub fn bridge(&self) -> Option<&Arc<dyn SessionBridge>> {
3523        self.bridge.as_ref()
3524    }
3525
3526    // -----------------------------------------------------------------------
3527    // Convenience methods
3528    // -----------------------------------------------------------------------
3529
3530    /// Send plain text to an addressable identity.
3531    pub async fn send_text(
3532        &self,
3533        identity: &AgentIdentity,
3534        text: impl Into<String>,
3535    ) -> Result<FencingToken, IdentityRuntimeError> {
3536        self.send(identity, &meerkat_core::ContentInput::Text(text.into()))
3537            .await
3538    }
3539
3540    /// Dispatch plain text with system origin.
3541    pub async fn dispatch_text(
3542        &self,
3543        identity: &AgentIdentity,
3544        text: impl Into<String>,
3545    ) -> Result<(FencingToken, bool), IdentityRuntimeError> {
3546        self.dispatch(identity, &DispatchInput::system(text)).await
3547    }
3548
3549    /// Execute the restore flow for the given roster.
3550    pub async fn restore_flow(
3551        &self,
3552        roster: &[DurableAgentSpec],
3553        topology_provider: Option<&dyn super::contracts::TopologyProvider>,
3554        customizer: Option<&dyn super::contracts::AgentCustomizer>,
3555    ) -> Result<super::orchestrator::RestoreFlowResult, IdentityRuntimeError> {
3556        super::orchestrator::restore_flow(self, roster, topology_provider, customizer).await
3557    }
3558
3559    /// Resolve the AgentRuntimeId for a registered identity.
3560    pub async fn runtime_id_for(
3561        &self,
3562        identity: &AgentIdentity,
3563    ) -> Result<AgentRuntimeId, IdentityRuntimeError> {
3564        let entries = self.entries.read().await;
3565        let entry = entries
3566            .get(identity)
3567            .ok_or_else(|| IdentityRuntimeError::UnknownIdentity(identity.clone()))?;
3568        entry
3569            .continuity
3570            .as_ref()
3571            .map(|c| c.agent_runtime_id.clone())
3572            .ok_or_else(|| {
3573                IdentityRuntimeError::Internal(format!("no continuity record for {identity}"))
3574            })
3575    }
3576
3577    /// Inspect the current execution state of an identity via the bridge.
3578    pub async fn inspect(
3579        &self,
3580        identity: &AgentIdentity,
3581    ) -> Result<super::bridge::MemberInspection, IdentityRuntimeError> {
3582        let runtime_id = self.runtime_id_for(identity).await?;
3583        let bridge = self
3584            .bridge
3585            .as_ref()
3586            .ok_or_else(|| IdentityRuntimeError::Internal("no bridge configured".to_string()))?;
3587        bridge
3588            .inspect_member(&runtime_id)
3589            .await
3590            .map_err(|e| IdentityRuntimeError::Internal(format!("inspect: {e}")))
3591    }
3592
3593    /// The configured default timeout for wait operations.
3594    pub fn default_timeout(&self) -> Duration {
3595        self.default_timeout
3596    }
3597
3598    /// Poll until the identity produces an output_preview, or timeout.
3599    pub async fn wait_for_output(
3600        &self,
3601        identity: &AgentIdentity,
3602        timeout: Duration,
3603    ) -> Result<String, IdentityRuntimeError> {
3604        let deadline = Instant::now() + timeout;
3605        loop {
3606            if let Ok(inspection) = self.inspect(identity).await
3607                && let Some(preview) = inspection.output_preview
3608            {
3609                return Ok(preview);
3610            }
3611            if Instant::now() >= deadline {
3612                return Err(IdentityRuntimeError::Internal(format!(
3613                    "timed out waiting for output from {identity}"
3614                )));
3615            }
3616            tokio::time::sleep(Duration::from_millis(500)).await;
3617        }
3618    }
3619
3620    /// Poll until output_preview contains the given substring, or timeout.
3621    pub async fn wait_for_output_containing(
3622        &self,
3623        identity: &AgentIdentity,
3624        needle: &str,
3625        timeout: Duration,
3626    ) -> Result<String, IdentityRuntimeError> {
3627        let deadline = Instant::now() + timeout;
3628        loop {
3629            if let Ok(inspection) = self.inspect(identity).await
3630                && let Some(ref preview) = inspection.output_preview
3631                && preview.contains(needle)
3632            {
3633                return Ok(preview.clone());
3634            }
3635            if Instant::now() >= deadline {
3636                return Err(IdentityRuntimeError::Internal(format!(
3637                    "timed out waiting for output containing '{needle}' from {identity}"
3638                )));
3639            }
3640            tokio::time::sleep(Duration::from_millis(500)).await;
3641        }
3642    }
3643}
3644
3645/// Wire two identities across mobs, resolving runtime IDs from both IdentityRuntimes.
3646///
3647/// This is a convenience function for same-process cross-mob scenarios where both
3648/// IdentityRuntimes are available. It resolves the AgentRuntimeId for each identity
3649/// and delegates to `UnifiedRuntime::wire_cross_mob()`.
3650pub async fn wire_cross_mob_by_identity(
3651    local_irt: &IdentityRuntime,
3652    local_identity: &AgentIdentity,
3653    remote_irt: &IdentityRuntime,
3654    remote_identity: &AgentIdentity,
3655    local_unified: &crate::UnifiedRuntime,
3656    remote_mob_id: &str,
3657) -> Result<(), IdentityRuntimeError> {
3658    let local_rt = local_irt.runtime_id_for(local_identity).await?;
3659    let remote_rt = remote_irt.runtime_id_for(remote_identity).await?;
3660    Box::pin(local_unified.wire_cross_mob(local_rt.as_str(), remote_rt.as_str(), remote_mob_id))
3661        .await
3662        .map_err(|e| IdentityRuntimeError::Internal(format!("wire_cross_mob: {e}")))
3663}
3664
3665#[cfg(test)]
3666mod lease_renewal_backoff_tests {
3667    use super::*;
3668
3669    /// Regression: a lease provider that errors persistently must not spin the
3670    /// renewal task at the TTL-derived floor (down to 10ms). The failure
3671    /// backoff grows from a 1s base and caps at the max poll interval.
3672    #[test]
3673    fn lease_renewal_failure_backoff_grows_and_caps() {
3674        let max = Duration::from_mins(1);
3675        assert_eq!(
3676            lease_renewal_failure_backoff(0, max),
3677            LEASE_RENEWAL_FAILURE_BACKOFF_BASE
3678        );
3679        assert_eq!(
3680            lease_renewal_failure_backoff(1, max),
3681            LEASE_RENEWAL_FAILURE_BACKOFF_BASE * 2
3682        );
3683        assert_eq!(lease_renewal_failure_backoff(6, max), max);
3684        // Saturates at the cap for arbitrarily many failures (no shift overflow).
3685        assert_eq!(lease_renewal_failure_backoff(99, max), max);
3686        assert!(lease_renewal_failure_backoff(2, max) > lease_renewal_failure_backoff(1, max));
3687    }
3688}