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