Skip to main content

meerkat_mobkit/identity_first/
runtime.rs

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