meerkat_mobkit/unified_runtime/
edge_reconcile.rs1use 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 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 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 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 let edges = self.reconcile_edges_from_members(active_snapshots).await;
77 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 if !report.mob.failures.is_empty() {
93 return Err(UnifiedRuntimeReconcileError::PartialFailure(Box::new(
94 report,
95 )));
96 }
97 Ok(report)
98 }
99
100 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 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 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 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 let raw_desired = edge_discovery.discover_edges(member_views).await;
167
168 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 for (a, b) in &desired {
188 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 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 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 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 !active_ids.contains(&a) || !active_ids.contains(&b) {
227 stale_pruned.push((a, b));
228 continue;
229 }
230 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 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}