Skip to main content

meerkat_mobkit/unified_runtime/
edge_reconcile.rs

1//! Edge topology reconciliation for distributed runtime nodes.
2
3use std::collections::{BTreeMap, BTreeSet};
4
5use futures::stream::{self, StreamExt};
6use meerkat_mob::SpawnMemberSpec;
7use meerkat_mob::ids::AgentIdentity;
8use meerkat_mob::runtime::MobMemberListEntry;
9use meerkat_mob::runtime::reconcile::ReconcileOptions;
10
11use crate::runtime::RuntimeRoute;
12use crate::unified_runtime::types::meerkat_reconcile_report_to_wire;
13
14use super::edge_types::{DesiredPeerEdge, EdgeMemberView, EdgeReconcileFailure};
15use super::types::{
16    UnifiedRuntimeReconcileEdgesReport, UnifiedRuntimeReconcileError,
17    UnifiedRuntimeReconcileReport, UnifiedRuntimeReconcileRoutingReport,
18};
19use super::{
20    ROSTER_ROUTE_CHANNEL, ROSTER_ROUTE_PREFIX, ROSTER_ROUTE_SINK, ROSTER_ROUTE_TARGET_MODULE,
21    UnifiedRuntime,
22};
23
24const EDGE_RECONCILE_CONCURRENCY: usize = 64;
25
26impl UnifiedRuntime {
27    pub async fn reconcile(
28        &self,
29        desired_specs: Vec<SpawnMemberSpec>,
30    ) -> Result<UnifiedRuntimeReconcileReport, UnifiedRuntimeReconcileError> {
31        self.mob_runtime
32            .set_baseline_member_specs(desired_specs.clone())
33            .await;
34        // Spec member ids arrive in the public alias space (identity-first
35        // identities like `domain:billing` contain `:`); encode them into
36        // comms-safe roster ids at the mobkit→meerkat-mob boundary (meerkat
37        // 0.7 `MemberCommsName` is fail-closed). The projections below —
38        // reconcile report, console identity registration, routing — decode
39        // back to the alias space.
40        let desired_specs: Vec<SpawnMemberSpec> = desired_specs
41            .into_iter()
42            .map(|mut spec| {
43                spec.identity = crate::member_comms_id::mob_member_id(spec.identity.as_str());
44                spec
45            })
46            .collect();
47        // 1. Member reconcile (meerkat 0.6 native path)
48        let mob_id = self.mob_handle().mob_id().to_string();
49        let meerkat_report = self
50            .mob_handle()
51            .reconcile(desired_specs, ReconcileOptions { retire_stale: true })
52            .await
53            .map_err(|err| UnifiedRuntimeReconcileError::Mob(err.into()))?;
54        let mob = meerkat_reconcile_report_to_wire(&mob_id, meerkat_report);
55        // 2. Refresh active members. Roster ids are comms-safe encodings; the
56        // agent-event ingest decodes them back to the alias space, so console
57        // registration is keyed by the public alias. Register as a fallback
58        // only: spawn/reserve paths register the durable console identity for
59        // the same alias key, and a reconcile must never clobber that with
60        // the alias self-mapping.
61        let active_snapshots = self.mob_handle().list_members_including_retiring().await;
62        for member in &active_snapshots {
63            let alias = crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str())
64                .into_owned();
65            self.console_events
66                .register_runtime_identity_fallback(alias.clone(), alias)
67                .await;
68        }
69        let active_member_ids = active_snapshots
70            .iter()
71            .map(|m| {
72                crate::member_comms_id::runtime_alias_str(m.agent_identity.as_str()).into_owned()
73            })
74            .collect::<Vec<_>>();
75        // 3 + 4. Edge discovery + dynamic edge reconcile
76        let edges = self.reconcile_edges_from_members(active_snapshots).await;
77        // 5. Routing reconcile
78        let routing = self.reconcile_routing_wiring(active_member_ids).await?;
79        let report = UnifiedRuntimeReconcileReport {
80            mob,
81            edges,
82            routing,
83        };
84        if let Some(hook) = &self.post_reconcile_hook {
85            hook(report.clone()).await;
86        }
87        // Meerkat 0.6's `MobHandle::reconcile` returns Ok even on per-identity
88        // failures (they're collected into `report.mob.failures`). Re-lift any
89        // non-empty failures list into an `Err` so callers using `?` see the
90        // same propagation behavior they had pre-cleanup, while still carrying
91        // the full report for inspection via `PartialFailure`.
92        if !report.mob.failures.is_empty() {
93            return Err(UnifiedRuntimeReconcileError::PartialFailure(Box::new(
94                report,
95            )));
96        }
97        Ok(report)
98    }
99
100    /// Reconcile dynamic peer edges using fresh roster state.
101    ///
102    /// Refreshes the roster, runs edge discovery if configured, diffs
103    /// desired vs managed edges, and calls wire/unwire as needed.
104    pub async fn reconcile_edges(&self) -> UnifiedRuntimeReconcileEdgesReport {
105        let active_members = self.mob_handle().list_members_including_retiring().await;
106        let report = self.reconcile_edges_from_members(active_members).await;
107        if !report.is_complete() {
108            self.fire_error(super::types::ErrorEvent::ReconcileIncomplete {
109                failures: report.failures.len(),
110                skipped: report.skipped_missing_members.len(),
111            });
112        }
113        report
114    }
115
116    pub(super) async fn reconcile_edges_from_members(
117        &self,
118        active_members: Vec<MobMemberListEntry>,
119    ) -> UnifiedRuntimeReconcileEdgesReport {
120        let edge_discovery = match &self.edge_discovery {
121            Some(d) => d,
122            None => return UnifiedRuntimeReconcileEdgesReport::default(),
123        };
124
125        // Everything in this function speaks the public alias space: roster
126        // ids are comms-safe encodings (meerkat 0.7 `MemberCommsName`), so
127        // decode snapshots on entry and encode only at the wire/unwire calls.
128        let alias_of =
129            |id: &str| -> String { crate::member_comms_id::runtime_alias_str(id).into_owned() };
130
131        let active_ids: BTreeSet<String> = active_members
132            .iter()
133            .map(|m| alias_of(m.agent_identity.as_str()))
134            .collect();
135
136        // Build current wiring map from snapshots
137        let mut current_edges: BTreeSet<(String, String)> = BTreeSet::new();
138        for member in &active_members {
139            for peer in &member.wired_to {
140                let mut a = alias_of(member.agent_identity.as_str());
141                let mut b = alias_of(peer.as_str());
142                if a > b {
143                    std::mem::swap(&mut a, &mut b);
144                }
145                current_edges.insert((a, b));
146            }
147        }
148
149        // Project to EdgeMemberView for the policy closure — it only needs
150        // identity/role/labels/wired_to, not meerkat's private runtime fields.
151        let member_views: Vec<EdgeMemberView> = active_members
152            .into_iter()
153            .map(|m| EdgeMemberView {
154                agent_identity: alias_of(m.agent_identity.as_str()),
155                role: m.role.to_string(),
156                wired_to: m
157                    .wired_to
158                    .iter()
159                    .map(|peer| alias_of(peer.as_str()))
160                    .collect(),
161                labels: m.labels,
162            })
163            .collect();
164
165        // Run edge discovery
166        let raw_desired = edge_discovery.discover_edges(member_views).await;
167
168        // Deduplicate and defensively validate (DesiredPeerEdge enforces
169        // invariants at construction, but we still canonicalize the key set)
170        let desired: BTreeSet<(String, String)> = raw_desired
171            .iter()
172            .map(|e| {
173                let (a, b) = e.endpoints();
174                (a.to_string(), b.to_string())
175            })
176            .collect();
177
178        let mut report = UnifiedRuntimeReconcileEdgesReport {
179            desired_edges: raw_desired,
180            ..Default::default()
181        };
182
183        let managed_snapshot = self.managed_dynamic_edges.read().await.clone();
184        let mut to_wire = Vec::new();
185
186        // Classify desired edges
187        for (a, b) in &desired {
188            // Skip if either endpoint is missing from the active roster
189            if !active_ids.contains(a) || !active_ids.contains(b) {
190                if let Ok(edge) = DesiredPeerEdge::new(a.clone(), b.clone()) {
191                    report.skipped_missing_members.push(edge);
192                }
193                continue;
194            }
195            let key = (a.clone(), b.clone());
196            if managed_snapshot.contains(&key) {
197                // Managed by us — check if the actual edge still exists in the
198                // mob graph. If an out-of-band unwire() removed it, re-wire.
199                if current_edges.contains(&key) {
200                    if let Ok(edge) = DesiredPeerEdge::new(a.clone(), b.clone()) {
201                        report.retained_edges.push(edge);
202                    }
203                } else {
204                    to_wire.push((a.clone(), b.clone(), "wire (heal)"));
205                }
206            } else if current_edges.contains(&key) {
207                // Exists but not managed by us (static or external) — don't claim
208                if let Ok(edge) = DesiredPeerEdge::new(a.clone(), b.clone()) {
209                    report.preexisting_edges.push(edge);
210                }
211            } else {
212                to_wire.push((a.clone(), b.clone(), "wire"));
213            }
214        }
215
216        // Unwire managed edges that are no longer desired
217        let mut stale_pruned = Vec::new();
218        let mut to_unwire = Vec::new();
219        for (a, b) in managed_snapshot
220            .iter()
221            .filter(|key| !desired.contains(*key))
222            .cloned()
223        {
224            let key = (a.clone(), b.clone());
225            // If either endpoint is gone, just prune from managed set
226            if !active_ids.contains(&a) || !active_ids.contains(&b) {
227                stale_pruned.push((a, b));
228                continue;
229            }
230            // If the edge is already gone from the mob graph (out-of-band
231            // unwire/reset), just drop ownership — don't attempt unwire.
232            if !current_edges.contains(&key) {
233                stale_pruned.push((a, b));
234                continue;
235            }
236            to_unwire.push((a, b));
237        }
238
239        let handle = self.mob_handle();
240        let wire_operations = to_wire
241            .iter()
242            .map(|(a, b, operation)| ((a.clone(), b.clone()), (*operation).to_string()))
243            .collect::<BTreeMap<_, _>>();
244        let mut wire_successes = Vec::new();
245        let mut wire_failures = Vec::new();
246        if !to_wire.is_empty() {
247            let batch_edges = to_wire
248                .iter()
249                .map(|(a, b, _)| {
250                    (
251                        crate::member_comms_id::mob_member_id(a.as_str()),
252                        crate::member_comms_id::mob_member_id(b.as_str()),
253                    )
254                })
255                .collect::<Vec<_>>();
256            // The batch report carries roster ids; decode back to aliases and
257            // re-canonicalize (alias-space ordering can differ from roster-id
258            // ordering) so keys line up with `to_wire`.
259            let alias_edge_key = |a: &AgentIdentity, b: &AgentIdentity| -> (String, String) {
260                let mut a = crate::member_comms_id::runtime_alias_str(a.as_str()).into_owned();
261                let mut b = crate::member_comms_id::runtime_alias_str(b.as_str()).into_owned();
262                if a > b {
263                    std::mem::swap(&mut a, &mut b);
264                }
265                (a, b)
266            };
267            match handle.wire_members_batch(batch_edges).await {
268                Ok(batch_report) => {
269                    let mut seen = BTreeSet::new();
270                    for edge in batch_report.wired {
271                        let key = alias_edge_key(&edge.a, &edge.b);
272                        seen.insert(key.clone());
273                        wire_successes.push((key.0, key.1, true));
274                    }
275                    for edge in batch_report.already_wired {
276                        let key = alias_edge_key(&edge.a, &edge.b);
277                        seen.insert(key.clone());
278                        wire_successes.push((key.0, key.1, false));
279                    }
280                    for (a, b, operation) in to_wire {
281                        let key = if a <= b { (a, b) } else { (b, a) };
282                        if !seen.contains(&key) {
283                            wire_failures.push((
284                                key.0,
285                                key.1,
286                                operation.to_string(),
287                                "wire_members_batch omitted edge from report".to_string(),
288                            ));
289                        }
290                    }
291                }
292                Err(err) => {
293                    let error = err.to_string();
294                    for (a, b, operation) in to_wire {
295                        wire_failures.push((a, b, operation.to_string(), error.clone()));
296                    }
297                }
298            }
299        }
300
301        let handle = self.mob_handle();
302        let unwire_results = stream::iter(to_unwire.into_iter().map(|(a, b)| {
303            let handle = handle.clone();
304            async move {
305                let result = handle
306                    .unwire(
307                        crate::member_comms_id::mob_member_id(a.as_str()),
308                        crate::member_comms_id::mob_member_id(b.as_str()),
309                    )
310                    .await
311                    .map_err(|err| format!("{err}"));
312                (a, b, result)
313            }
314        }))
315        .buffer_unordered(EDGE_RECONCILE_CONCURRENCY)
316        .collect::<Vec<_>>()
317        .await;
318
319        let mut managed_edges = self.managed_dynamic_edges.write().await;
320        for (a, b, newly_wired) in wire_successes {
321            managed_edges.insert((a.clone(), b.clone()));
322            if let Ok(edge) = DesiredPeerEdge::new(a.clone(), b.clone()) {
323                if newly_wired {
324                    report.wired_edges.push(edge);
325                } else {
326                    report.retained_edges.push(edge);
327                }
328            }
329        }
330        for (a, b, operation, error) in wire_failures {
331            if let Ok(edge) = DesiredPeerEdge::new(a.clone(), b.clone()) {
332                report.failures.push(EdgeReconcileFailure {
333                    edge,
334                    operation: wire_operations.get(&(a, b)).cloned().unwrap_or(operation),
335                    error,
336                });
337            }
338        }
339        for (a, b) in stale_pruned {
340            managed_edges.remove(&(a.clone(), b.clone()));
341            if let Ok(edge) = DesiredPeerEdge::new(a, b) {
342                report.pruned_stale_managed_edges.push(edge);
343            }
344        }
345        for (a, b, result) in unwire_results {
346            match result {
347                Ok(()) => {
348                    managed_edges.remove(&(a.clone(), b.clone()));
349                    if let Ok(edge) = DesiredPeerEdge::new(a, b) {
350                        report.unwired_edges.push(edge);
351                    }
352                }
353                Err(error) => {
354                    if let Ok(edge) = DesiredPeerEdge::new(a, b) {
355                        report.failures.push(EdgeReconcileFailure {
356                            edge,
357                            operation: "unwire".into(),
358                            error,
359                        });
360                    }
361                }
362            }
363        }
364
365        report
366    }
367
368    pub(super) async fn reconcile_routing_wiring(
369        &self,
370        mut active_members: Vec<String>,
371    ) -> Result<UnifiedRuntimeReconcileRoutingReport, UnifiedRuntimeReconcileError> {
372        active_members.sort();
373        active_members.dedup();
374
375        let mut rt = self.module_runtime.lock().await;
376        let router_module_loaded = rt
377            .loaded_modules()
378            .iter()
379            .any(|module_id| module_id == "router");
380        let mut added_route_keys = Vec::new();
381        let mut removed_route_keys = Vec::new();
382
383        if router_module_loaded {
384            let managed_routes: Vec<RuntimeRoute> = rt
385                .list_runtime_routes()
386                .into_iter()
387                .filter(|route| route.route_key.starts_with(ROSTER_ROUTE_PREFIX))
388                .collect();
389            let active_member_set = active_members.iter().cloned().collect::<BTreeSet<_>>();
390            for route in &managed_routes {
391                if !active_member_set.contains(&route.recipient) {
392                    rt.delete_runtime_route(&route.route_key)
393                        .map_err(UnifiedRuntimeReconcileError::RouteMutation)?;
394                    removed_route_keys.push(route.route_key.clone());
395                }
396            }
397
398            let existing_managed_recipients = managed_routes
399                .into_iter()
400                .map(|route| route.recipient)
401                .collect::<BTreeSet<_>>();
402            for member_id in &active_members {
403                if existing_managed_recipients.contains(member_id) {
404                    continue;
405                }
406                let route_key = format!("{ROSTER_ROUTE_PREFIX}{member_id}");
407                rt.add_runtime_route(RuntimeRoute {
408                    route_key: route_key.clone(),
409                    recipient: member_id.clone(),
410                    channel: Some(ROSTER_ROUTE_CHANNEL.to_string()),
411                    sink: ROSTER_ROUTE_SINK.to_string(),
412                    target_module: ROSTER_ROUTE_TARGET_MODULE.to_string(),
413                    retry_max: None,
414                    backoff_ms: None,
415                    rate_limit_per_minute: None,
416                })
417                .map_err(UnifiedRuntimeReconcileError::RouteMutation)?;
418                added_route_keys.push(route_key);
419            }
420        }
421
422        added_route_keys.sort();
423        removed_route_keys.sort();
424
425        Ok(UnifiedRuntimeReconcileRoutingReport {
426            router_module_loaded,
427            active_members,
428            added_route_keys,
429            removed_route_keys,
430        })
431    }
432}