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