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