1use std::sync::Arc;
4
5use meerkat_core::handles::{DslTransitionError, PeerCommsHandle};
6use meerkat_core::interaction::{
7 PeerIngressAdmission, PeerIngressDequeueAuthority, PeerIngressDequeueFacts,
8 PeerIngressEnvelopeFacts, PeerIngressPlainEventFacts, PeerIngressReceiveAuthority,
9 PeerIngressReceiveFacts,
10};
11
12use super::HandleDslAuthority;
13use crate::meerkat_machine::dsl as mm_dsl;
14
15#[derive(Debug)]
20pub struct RuntimePeerCommsHandle {
21 dsl: Arc<HandleDslAuthority>,
22}
23
24impl RuntimePeerCommsHandle {
25 pub fn new(dsl: Arc<HandleDslAuthority>) -> Self {
27 Self { dsl }
28 }
29
30 #[doc(hidden)]
35 pub fn generated_install_factory(
36 dsl: Arc<HandleDslAuthority>,
37 ) -> Result<meerkat_core::handles::GeneratedPeerCommsInstallFactory, String> {
38 let owner = dsl.generated_authority_owner_token();
39 crate::protocol_comms_trust_reconcile::generated_peer_comms_install_factory(
40 Arc::new(Self::new(dsl)),
41 owner,
42 )
43 }
44
45 pub fn install_generated_on(
50 dsl: Arc<HandleDslAuthority>,
51 target: &(dyn meerkat_core::handles::PeerCommsInstallTarget + '_),
52 ) -> Result<(), String> {
53 Self::generated_install_factory(dsl)?.install_on_target(target)
54 }
55
56 pub fn ephemeral() -> Self {
60 Self::new(Arc::new(HandleDslAuthority::ephemeral()))
61 }
62}
63
64fn lifecycle_to_dsl(
65 kind: meerkat_core::comms::PeerLifecycleKind,
66) -> Result<mm_dsl::PeerIngressLifecycleClass, DslTransitionError> {
67 match kind {
68 meerkat_core::comms::PeerLifecycleKind::PeerAdded => {
69 Ok(mm_dsl::PeerIngressLifecycleClass::PeerAdded)
70 }
71 meerkat_core::comms::PeerLifecycleKind::PeerRetired => {
72 Ok(mm_dsl::PeerIngressLifecycleClass::PeerRetired)
73 }
74 meerkat_core::comms::PeerLifecycleKind::PeerUnwired => {
75 Ok(mm_dsl::PeerIngressLifecycleClass::PeerUnwired)
76 }
77 meerkat_core::comms::PeerLifecycleKind::Dismiss => Err(DslTransitionError::guard_rejected(
82 "PeerCommsHandle::classify_external_envelope",
83 "mob.dismiss is a drain-lifecycle terminal and is not classifiable as peer ingress",
84 )),
85 }
86}
87
88fn lifecycle_from_dsl(
89 kind: mm_dsl::PeerIngressLifecycleClass,
90) -> meerkat_core::comms::PeerLifecycleKind {
91 match kind {
92 mm_dsl::PeerIngressLifecycleClass::PeerAdded => {
93 meerkat_core::comms::PeerLifecycleKind::PeerAdded
94 }
95 mm_dsl::PeerIngressLifecycleClass::PeerRetired => {
96 meerkat_core::comms::PeerLifecycleKind::PeerRetired
97 }
98 mm_dsl::PeerIngressLifecycleClass::PeerUnwired => {
99 meerkat_core::comms::PeerLifecycleKind::PeerUnwired
100 }
101 }
102}
103
104fn response_status_to_dsl(
105 status: meerkat_core::ResponseStatus,
106) -> mm_dsl::PeerIngressResponseStatus {
107 match status {
108 meerkat_core::ResponseStatus::Accepted => mm_dsl::PeerIngressResponseStatus::Accepted,
109 meerkat_core::ResponseStatus::Completed => mm_dsl::PeerIngressResponseStatus::Completed,
110 meerkat_core::ResponseStatus::Failed => mm_dsl::PeerIngressResponseStatus::Failed,
111 }
112}
113
114fn lifecycle_subject(
124 params: &serde_json::Value,
125 context: &'static str,
126) -> Result<String, DslTransitionError> {
127 let parsed: meerkat_contracts::CommsPeerLifecycleParams =
128 serde_json::from_value(params.clone()).map_err(|error| {
129 DslTransitionError::guard_rejected(
130 context,
131 format!("malformed peer lifecycle params: {error}"),
132 )
133 })?;
134 let subject = match parsed.peer_spec {
135 Some(spec) => spec.name,
136 None => parsed.peer,
137 };
138 if subject.is_empty() {
139 return Err(DslTransitionError::guard_rejected(
140 context,
141 "peer lifecycle params carry an empty peer subject",
142 ));
143 }
144 Ok(subject)
145}
146
147fn terminality_from_dsl(
148 terminality: mm_dsl::PeerIngressResponseTerminality,
149) -> meerkat_core::TerminalityClass {
150 match terminality {
151 mm_dsl::PeerIngressResponseTerminality::Progress => {
152 meerkat_core::TerminalityClass::Progress
153 }
154 mm_dsl::PeerIngressResponseTerminality::TerminalCompleted => {
155 meerkat_core::TerminalityClass::Terminal {
156 disposition: meerkat_core::TerminalDisposition::Completed,
157 }
158 }
159 mm_dsl::PeerIngressResponseTerminality::TerminalFailed => {
160 meerkat_core::TerminalityClass::Terminal {
161 disposition: meerkat_core::TerminalDisposition::Failed,
162 }
163 }
164 }
165}
166
167fn phase_from_dsl(
168 phase: mm_dsl::PeerIngressAuthorityPhaseClass,
169) -> meerkat_core::PeerIngressAuthorityPhase {
170 phase.into()
171}
172
173fn kind_to_dsl(kind: meerkat_core::PeerIngressKind) -> mm_dsl::PeerIngressAdmittedKind {
174 match kind {
175 meerkat_core::PeerIngressKind::Message => mm_dsl::PeerIngressAdmittedKind::Message,
176 meerkat_core::PeerIngressKind::Request => mm_dsl::PeerIngressAdmittedKind::Request,
177 meerkat_core::PeerIngressKind::Response => mm_dsl::PeerIngressAdmittedKind::Response,
178 meerkat_core::PeerIngressKind::Ack => mm_dsl::PeerIngressAdmittedKind::Ack,
179 meerkat_core::PeerIngressKind::PlainEvent => mm_dsl::PeerIngressAdmittedKind::PlainEvent,
180 }
181}
182
183fn auth_to_dsl(auth: meerkat_core::PeerIngressAuthDecision) -> mm_dsl::PeerIngressAuthClass {
184 match auth {
185 meerkat_core::PeerIngressAuthDecision::Required => mm_dsl::PeerIngressAuthClass::Required,
186 meerkat_core::PeerIngressAuthDecision::Exempt(
187 meerkat_core::PeerIngressAuthExemption::SupervisorBridge,
188 ) => mm_dsl::PeerIngressAuthClass::SupervisorBridgeExempt,
189 }
190}
191
192fn request_intent_class_for(intent: &str) -> mm_dsl::PeerIngressRequestClass {
198 match intent {
199 "mob.peer_added" => mm_dsl::PeerIngressRequestClass::MobPeerAdded,
200 "mob.peer_retired" => mm_dsl::PeerIngressRequestClass::MobPeerRetired,
201 "mob.peer_unwired" => mm_dsl::PeerIngressRequestClass::MobPeerUnwired,
202 "supervisor.bridge" => mm_dsl::PeerIngressRequestClass::SupervisorBridge,
203 _ => mm_dsl::PeerIngressRequestClass::Other,
204 }
205}
206
207fn external_envelope_signal(
208 facts: &PeerIngressEnvelopeFacts,
209) -> Result<mm_dsl::MeerkatMachineSignal, DslTransitionError> {
210 const CONTEXT: &str = "PeerCommsHandle::classify_external_envelope";
211 let (
212 envelope_kind,
213 request_intent,
214 lifecycle_kind,
215 lifecycle_peer_param,
216 response_status,
217 in_reply_to,
218 ) = match &facts.kind {
219 meerkat_core::PeerIngressEnvelopeKind::Message { .. } => (
220 mm_dsl::PeerIngressEnvelopeClass::Message,
221 String::new(),
222 mm_dsl::PeerIngressLifecycleClass::PeerAdded,
223 None,
224 mm_dsl::PeerIngressResponseStatus::Accepted,
225 String::new(),
226 ),
227 meerkat_core::PeerIngressEnvelopeKind::Request { intent, params } => {
228 let lifecycle_peer_param = match request_intent_class_for(intent) {
233 mm_dsl::PeerIngressRequestClass::MobPeerAdded
234 | mm_dsl::PeerIngressRequestClass::MobPeerRetired
235 | mm_dsl::PeerIngressRequestClass::MobPeerUnwired => {
236 Some(lifecycle_subject(params, CONTEXT)?)
237 }
238 mm_dsl::PeerIngressRequestClass::SupervisorBridge
239 | mm_dsl::PeerIngressRequestClass::Other => None,
240 };
241 (
242 mm_dsl::PeerIngressEnvelopeClass::Request,
243 intent.clone(),
244 mm_dsl::PeerIngressLifecycleClass::PeerAdded,
245 lifecycle_peer_param,
246 mm_dsl::PeerIngressResponseStatus::Accepted,
247 String::new(),
248 )
249 }
250 meerkat_core::PeerIngressEnvelopeKind::Lifecycle { kind, params } => (
251 mm_dsl::PeerIngressEnvelopeClass::Lifecycle,
252 String::new(),
253 lifecycle_to_dsl(*kind)?,
254 Some(lifecycle_subject(params, CONTEXT)?),
255 mm_dsl::PeerIngressResponseStatus::Accepted,
256 String::new(),
257 ),
258 meerkat_core::PeerIngressEnvelopeKind::Response {
259 in_reply_to: reply_to,
260 status,
261 ..
262 } => (
263 mm_dsl::PeerIngressEnvelopeClass::Response,
264 String::new(),
265 mm_dsl::PeerIngressLifecycleClass::PeerAdded,
266 None,
267 response_status_to_dsl(*status),
268 reply_to.clone(),
269 ),
270 meerkat_core::PeerIngressEnvelopeKind::Ack {
271 in_reply_to: reply_to,
272 } => (
273 mm_dsl::PeerIngressEnvelopeClass::Ack,
274 String::new(),
275 mm_dsl::PeerIngressLifecycleClass::PeerAdded,
276 None,
277 mm_dsl::PeerIngressResponseStatus::Accepted,
278 reply_to.clone(),
279 ),
280 };
281
282 let request_intent_class = request_intent_class_for(&request_intent);
283
284 Ok(mm_dsl::MeerkatMachineSignal::ClassifyExternalEnvelope {
285 item_id: facts.item_id.clone(),
286 from_peer: facts.from_peer.clone(),
287 from_peer_id: mm_dsl::PeerId(facts.from_peer_id.to_string()),
288 envelope_kind,
289 request_intent,
290 request_intent_class,
291 lifecycle_kind,
292 lifecycle_peer_param,
293 response_status,
294 in_reply_to,
295 })
296}
297
298struct PeerIngressClassifiedEffect {
299 class: mm_dsl::PeerIngressInputClass,
300 actionable: bool,
301 kind: mm_dsl::PeerIngressAdmittedKind,
302 auth: mm_dsl::PeerIngressAuthClass,
303 from_peer_id: Option<mm_dsl::PeerId>,
304 lifecycle_kind: Option<mm_dsl::PeerIngressLifecycleClass>,
305 lifecycle_peer: Option<String>,
306 request_id: Option<String>,
307 response_terminality: Option<mm_dsl::PeerIngressResponseTerminality>,
308}
309
310fn classified_effect(
311 effects: Vec<mm_dsl::MeerkatMachineEffect>,
312 context: &'static str,
313) -> Result<PeerIngressClassifiedEffect, DslTransitionError> {
314 effects
315 .into_iter()
316 .find_map(|effect| match effect {
317 mm_dsl::MeerkatMachineEffect::PeerIngressClassified {
318 class,
319 actionable,
320 kind,
321 auth,
322 from_peer_id,
323 lifecycle_kind,
324 lifecycle_peer,
325 request_id,
326 response_terminality,
327 } => Some(PeerIngressClassifiedEffect {
328 class,
329 actionable,
330 kind,
331 auth,
332 from_peer_id,
333 lifecycle_kind,
334 lifecycle_peer,
335 request_id,
336 response_terminality,
337 }),
338 _ => None,
339 })
340 .ok_or_else(|| {
341 DslTransitionError::guard_rejected(
342 context,
343 "machine transition did not emit PeerIngressClassified",
344 )
345 })
346}
347
348fn canonical_peer_id_from_effect(
351 from_peer_id: Option<&mm_dsl::PeerId>,
352 context: &'static str,
353) -> Result<Option<meerkat_core::comms::PeerId>, DslTransitionError> {
354 from_peer_id
355 .map(|peer_id| {
356 meerkat_core::comms::PeerId::parse(peer_id.0.as_str()).map_err(|error| {
357 DslTransitionError::guard_rejected(
358 context,
359 format!("machine emitted malformed canonical sender peer id: {error}"),
360 )
361 })
362 })
363 .transpose()
364}
365
366struct PeerIngressReceiveResolvedEffect {
367 outcome: mm_dsl::PeerIngressReceiveOutcomeClass,
368 admission_diagnostic: Option<mm_dsl::PeerIngressAdmissionDiagnosticClass>,
369 phase: mm_dsl::PeerIngressAuthorityPhaseClass,
370}
371
372fn receive_resolved_effect(
373 effects: Vec<mm_dsl::MeerkatMachineEffect>,
374 context: &'static str,
375) -> Result<PeerIngressReceiveResolvedEffect, DslTransitionError> {
376 effects
377 .into_iter()
378 .find_map(|effect| match effect {
379 mm_dsl::MeerkatMachineEffect::PeerIngressReceiveResolved {
380 outcome,
381 admission_diagnostic,
382 phase,
383 } => Some(PeerIngressReceiveResolvedEffect {
384 outcome,
385 admission_diagnostic,
386 phase,
387 }),
388 _ => None,
389 })
390 .ok_or_else(|| {
391 DslTransitionError::guard_rejected(
392 context,
393 "machine transition did not emit PeerIngressReceiveResolved",
394 )
395 })
396}
397
398fn dequeue_resolved_effect(
399 effects: Vec<mm_dsl::MeerkatMachineEffect>,
400 context: &'static str,
401) -> Result<mm_dsl::PeerIngressAuthorityPhaseClass, DslTransitionError> {
402 effects
403 .into_iter()
404 .find_map(|effect| match effect {
405 mm_dsl::MeerkatMachineEffect::PeerIngressDequeueResolved { phase } => Some(phase),
406 _ => None,
407 })
408 .ok_or_else(|| {
409 DslTransitionError::guard_rejected(
410 context,
411 "machine transition did not emit PeerIngressDequeueResolved",
412 )
413 })
414}
415
416fn classification_from_effect(
417 effect: &PeerIngressClassifiedEffect,
418) -> meerkat_core::PeerIngressClassification {
419 let class = match effect.class {
420 mm_dsl::PeerIngressInputClass::ActionableMessage => {
421 meerkat_core::PeerInputClass::ActionableMessage
422 }
423 mm_dsl::PeerIngressInputClass::ActionableRequest => {
424 meerkat_core::PeerInputClass::ActionableRequest
425 }
426 mm_dsl::PeerIngressInputClass::ResponseProgress => {
427 meerkat_core::PeerInputClass::ResponseProgress
428 }
429 mm_dsl::PeerIngressInputClass::ResponseTerminal => {
430 meerkat_core::PeerInputClass::ResponseTerminal
431 }
432 mm_dsl::PeerIngressInputClass::PeerLifecycleAdded => {
433 meerkat_core::PeerInputClass::PeerLifecycleAdded
434 }
435 mm_dsl::PeerIngressInputClass::PeerLifecycleRetired => {
436 meerkat_core::PeerInputClass::PeerLifecycleRetired
437 }
438 mm_dsl::PeerIngressInputClass::PeerLifecycleUnwired => {
439 meerkat_core::PeerInputClass::PeerLifecycleUnwired
440 }
441 mm_dsl::PeerIngressInputClass::SilentRequest => meerkat_core::PeerInputClass::SilentRequest,
442 mm_dsl::PeerIngressInputClass::Ack => meerkat_core::PeerInputClass::Ack,
443 mm_dsl::PeerIngressInputClass::PlainEvent => meerkat_core::PeerInputClass::PlainEvent,
444 };
445 let kind = match effect.kind {
446 mm_dsl::PeerIngressAdmittedKind::Message => meerkat_core::PeerIngressKind::Message,
447 mm_dsl::PeerIngressAdmittedKind::Request => meerkat_core::PeerIngressKind::Request,
448 mm_dsl::PeerIngressAdmittedKind::Response => meerkat_core::PeerIngressKind::Response,
449 mm_dsl::PeerIngressAdmittedKind::Ack => meerkat_core::PeerIngressKind::Ack,
450 mm_dsl::PeerIngressAdmittedKind::PlainEvent => meerkat_core::PeerIngressKind::PlainEvent,
451 };
452 let auth = match effect.auth {
453 mm_dsl::PeerIngressAuthClass::Required => meerkat_core::PeerIngressAuthDecision::Required,
454 mm_dsl::PeerIngressAuthClass::SupervisorBridgeExempt => {
455 meerkat_core::PeerIngressAuthDecision::Exempt(
456 meerkat_core::PeerIngressAuthExemption::SupervisorBridge,
457 )
458 }
459 };
460
461 meerkat_core::PeerIngressClassification {
462 class,
463 actionable: effect.actionable,
464 kind,
465 auth,
466 lifecycle_kind: effect.lifecycle_kind.map(lifecycle_from_dsl),
467 response_terminality: effect.response_terminality.map(terminality_from_dsl),
468 }
469}
470
471impl PeerCommsHandle for RuntimePeerCommsHandle {
472 fn classify_external_envelope(
473 &self,
474 facts: PeerIngressEnvelopeFacts,
475 ) -> Result<PeerIngressAdmission, DslTransitionError> {
476 let context = "PeerCommsHandle::classify_external_envelope";
477 let effects = self
478 .dsl
479 .apply_signal_with_effects(external_envelope_signal(&facts)?, context)?;
480 let effect = classified_effect(effects, context)?;
481 let classification = classification_from_effect(&effect);
482 let from_peer_id = canonical_peer_id_from_effect(effect.from_peer_id.as_ref(), context)?
486 .ok_or_else(|| {
487 DslTransitionError::guard_rejected(
488 context,
489 "machine classification did not echo the canonical sender peer id",
490 )
491 })?;
492 Ok(PeerIngressAdmission {
493 rendered_text: meerkat_core::render_peer_ingress_admitted_text(&facts, &classification),
494 classification,
495 from_peer_id: Some(from_peer_id),
496 lifecycle_peer: effect.lifecycle_peer,
501 request_id: effect.request_id,
502 })
503 }
504
505 fn classify_plain_event(
506 &self,
507 facts: PeerIngressPlainEventFacts,
508 ) -> Result<PeerIngressAdmission, DslTransitionError> {
509 let context = "PeerCommsHandle::classify_plain_event";
510 let effects = self.dsl.apply_signal_with_effects(
511 mm_dsl::MeerkatMachineSignal::ClassifyPlainEvent {
512 source_name: facts.source_name.clone(),
513 },
514 context,
515 )?;
516 let effect = classified_effect(effects, context)?;
517 Ok(PeerIngressAdmission {
518 classification: classification_from_effect(&effect),
519 from_peer_id: canonical_peer_id_from_effect(effect.from_peer_id.as_ref(), context)?,
522 lifecycle_peer: effect.lifecycle_peer,
523 request_id: effect.request_id,
524 rendered_text: meerkat_core::interaction::format_external_event_projection(
525 &facts.source_name,
526 Some(&facts.body),
527 ),
528 })
529 }
530
531 fn resolve_peer_ingress_receive(
532 &self,
533 facts: PeerIngressReceiveFacts,
534 ) -> Result<PeerIngressReceiveAuthority, DslTransitionError> {
535 let context = "PeerCommsHandle::resolve_peer_ingress_receive";
536 let effects = self.dsl.apply_input_with_effects(
537 mm_dsl::MeerkatMachineInput::ResolvePeerIngressReceive {
538 kind: kind_to_dsl(facts.kind),
539 auth_required: facts.auth_required,
540 auth_exempt: facts.auth_exempt,
541 trusted: facts.trusted,
542 queued_work_present: facts.queued_work_present,
543 queue_closed: facts.queue_closed,
544 queue_capacity_available: facts.queue_capacity_available,
545 },
546 context,
547 )?;
548 let effect = receive_resolved_effect(effects, context)?;
549 Ok(PeerIngressReceiveAuthority {
550 outcome: effect.outcome.into(),
551 admission_diagnostic: effect.admission_diagnostic.map(Into::into),
552 authority_phase: phase_from_dsl(effect.phase),
553 })
554 }
555
556 fn resolve_peer_ingress_dequeue(
557 &self,
558 facts: PeerIngressDequeueFacts,
559 ) -> Result<PeerIngressDequeueAuthority, DslTransitionError> {
560 let context = "PeerCommsHandle::resolve_peer_ingress_dequeue";
561 let effects = self.dsl.apply_input_with_effects(
562 mm_dsl::MeerkatMachineInput::ResolvePeerIngressDequeue {
563 kind: kind_to_dsl(facts.kind),
564 auth: auth_to_dsl(facts.auth),
565 queued_work_remaining: facts.queued_work_remaining,
566 },
567 context,
568 )?;
569 let phase = dequeue_resolved_effect(effects, context)?;
570 Ok(PeerIngressDequeueAuthority {
571 authority_phase: phase_from_dsl(phase),
572 })
573 }
574
575 fn set_peer_ingress_context(&self, keep_alive: bool) -> Result<(), DslTransitionError> {
576 self.dsl.apply_input(
578 mm_dsl::MeerkatMachineInput::SetPeerIngressContext { keep_alive },
579 "PeerCommsHandle::set_peer_ingress_context",
580 )
581 }
582
583 fn install_generated_peer_comms_on_target(
584 &self,
585 expected_owner: &meerkat_core::comms::GeneratedPeerCommsOwnerToken,
586 target: &(dyn meerkat_core::handles::PeerCommsInstallTarget + '_),
587 ) -> Result<(), String> {
588 let target_endpoint = target.generated_peer_comms_target_endpoint()?;
589 let mut guard = self
590 .dsl
591 .inner
592 .lock()
593 .unwrap_or_else(std::sync::PoisonError::into_inner);
594 let previous_endpoint = guard.state().local_endpoint.clone();
595 let publish_input = mm_dsl::MeerkatMachineInput::PublishLocalEndpoint {
596 endpoint: mm_dsl::PeerEndpoint::from(&target_endpoint),
597 };
598 crate::meerkat_machine_types::MeerkatMachineFieldlessRuntimeInternalInput::reject_raw_dsl_input(
599 &publish_input,
600 )
601 .map_err(|reason| {
602 DslTransitionError::no_matching(
603 "PeerCommsHandle::install_generated_peer_comms_on_target",
604 reason,
605 )
606 .to_string()
607 })?;
608 mm_dsl::MeerkatMachineMutator::apply(&mut *guard, publish_input).map_err(|error| {
609 super::map_kernel_error(
610 error,
611 "PeerCommsHandle::install_generated_peer_comms_on_target",
612 )
613 .to_string()
614 })?;
615 let approved_peer_id = guard
616 .state()
617 .local_endpoint
618 .as_ref()
619 .map(|endpoint| endpoint.peer_id.clone())
620 .ok_or_else(|| {
621 "generated peer-comms install target authority did not publish local endpoint"
622 .to_string()
623 })?;
624 let approved_peer_id = meerkat_core::comms::PeerId::parse(approved_peer_id.as_str())
625 .map_err(|error| {
626 format!("generated peer-comms install target peer_id invalid: {error}")
627 })?;
628 let install = crate::protocol_comms_trust_reconcile::generated_peer_comms_install(
629 Arc::new(Self::new(Arc::clone(&self.dsl))),
630 guard.generated_authority_owner_token(),
631 approved_peer_id,
632 )?;
633 if !install.owner_token().same_owner(expected_owner) {
634 restore_local_endpoint(&mut guard, previous_endpoint)?;
635 return Err(
636 "generated peer-comms install came from a different MeerkatMachine owner"
637 .to_string(),
638 );
639 }
640 if let Err(error) = target.install_generated_peer_comms_handle(install) {
641 let rollback = restore_local_endpoint(&mut guard, previous_endpoint);
642 if let Err(rollback_error) = rollback {
643 return Err(format!(
644 "{error}; generated local endpoint rollback failed: {rollback_error}"
645 ));
646 }
647 return Err(error);
648 }
649 Ok(())
650 }
651}
652
653fn restore_local_endpoint(
654 guard: &mut mm_dsl::MeerkatMachineAuthority,
655 previous_endpoint: Option<mm_dsl::PeerEndpoint>,
656) -> Result<(), String> {
657 let input = match previous_endpoint {
658 Some(endpoint) => mm_dsl::MeerkatMachineInput::PublishLocalEndpoint { endpoint },
659 None => mm_dsl::MeerkatMachineInput::ClearLocalEndpoint,
660 };
661 crate::meerkat_machine_types::MeerkatMachineFieldlessRuntimeInternalInput::reject_raw_dsl_input(
662 &input,
663 )
664 .map_err(|reason| {
665 DslTransitionError::no_matching("PeerCommsHandle::restore_local_endpoint", reason)
666 .to_string()
667 })?;
668 mm_dsl::MeerkatMachineMutator::apply(guard, input)
669 .map(|_| ())
670 .map_err(|error| {
671 super::map_kernel_error(error, "PeerCommsHandle::restore_local_endpoint").to_string()
672 })
673}
674
675#[cfg(test)]
676mod tests {
677 use super::*;
678 use std::collections::BTreeSet;
679 use std::sync::Mutex;
680
681 struct RejectingPeerCommsInstallTarget {
682 descriptor: meerkat_core::comms::TrustedPeerDescriptor,
683 notify: Arc<tokio::sync::Notify>,
684 }
685
686 impl RejectingPeerCommsInstallTarget {
687 fn new(descriptor: meerkat_core::comms::TrustedPeerDescriptor) -> Self {
688 Self {
689 descriptor,
690 notify: Arc::new(tokio::sync::Notify::new()),
691 }
692 }
693 }
694
695 #[async_trait::async_trait]
696 impl meerkat_core::agent::CommsRuntime for RejectingPeerCommsInstallTarget {
697 async fn drain_messages(&self) -> Vec<String> {
698 Vec::new()
699 }
700
701 fn inbox_notify(&self) -> Arc<tokio::sync::Notify> {
702 Arc::clone(&self.notify)
703 }
704
705 fn peer_id(&self) -> Option<meerkat_core::comms::PeerId> {
706 Some(self.descriptor.peer_id)
707 }
708
709 fn public_key_bytes(&self) -> Option<[u8; 32]> {
710 Some(self.descriptor.pubkey)
711 }
712
713 fn comms_name(&self) -> Option<String> {
714 Some(self.descriptor.name.as_str().to_string())
715 }
716
717 fn advertised_address(&self) -> Option<String> {
718 Some(self.descriptor.address.to_string())
719 }
720 }
721
722 impl meerkat_core::handles::PeerCommsInstallTarget for RejectingPeerCommsInstallTarget {
723 fn install_generated_peer_comms_handle(
724 &self,
725 _install: meerkat_core::handles::GeneratedPeerCommsInstall,
726 ) -> Result<(), String> {
727 Err("target rejected install".to_string())
728 }
729 }
730
731 fn peer_descriptor(name: &str, seed: u8) -> meerkat_core::comms::TrustedPeerDescriptor {
732 let pubkey = [seed; 32];
733 let peer_id = meerkat_core::comms::PeerId::from_ed25519_pubkey(&pubkey);
734 meerkat_core::comms::TrustedPeerDescriptor::unsigned_with_pubkey(
735 name,
736 peer_id.to_string(),
737 pubkey,
738 format!("inproc://{name}"),
739 )
740 .expect("valid peer descriptor")
741 }
742
743 fn handle_for_phase(phase: mm_dsl::MeerkatPhase) -> RuntimePeerCommsHandle {
744 let state = mm_dsl::MeerkatMachineState {
745 lifecycle_phase: phase,
746 session_id: Some(mm_dsl::SessionId("session-1".to_string())),
747 ..Default::default()
748 };
749 let authority = Arc::new(Mutex::new(
750 mm_dsl::MeerkatMachineAuthority::recover_from_state(state)
751 .expect("test MeerkatMachine state must be recoverable"),
752 ));
753 RuntimePeerCommsHandle::new(Arc::new(HandleDslAuthority::from_shared(authority)))
754 }
755
756 #[test]
757 fn peer_comms_install_rejection_restores_generated_local_endpoint() {
758 for previous in [None, Some(peer_descriptor("previous", 7))] {
759 let rejected = peer_descriptor("rejected", 9);
760 let expected_previous = previous.as_ref().map(mm_dsl::PeerEndpoint::from);
761 let state = mm_dsl::MeerkatMachineState {
762 lifecycle_phase: mm_dsl::MeerkatPhase::Attached,
763 session_id: Some(mm_dsl::SessionId("session-1".to_string())),
764 local_endpoint: expected_previous.clone(),
765 ..Default::default()
766 };
767 let authority = Arc::new(Mutex::new(
768 mm_dsl::MeerkatMachineAuthority::recover_from_state(state)
769 .expect("test MeerkatMachine state must be recoverable"),
770 ));
771 let dsl = Arc::new(HandleDslAuthority::from_shared(authority));
772 let factory = RuntimePeerCommsHandle::generated_install_factory(Arc::clone(&dsl))
773 .expect("generated peer-comms install factory");
774 let target = RejectingPeerCommsInstallTarget::new(rejected);
775 let target: &dyn meerkat_core::handles::PeerCommsInstallTarget = ⌖
776
777 let error = factory
778 .install_on_target(target)
779 .expect_err("target rejection should fail install");
780
781 assert!(
782 error.contains("target rejected install"),
783 "unexpected install error: {error}"
784 );
785 assert_eq!(dsl.snapshot_state().local_endpoint, expected_previous);
786 }
787 }
788
789 #[test]
790 fn runtime_peer_comms_handle_classifies_from_dsl_silent_intents() {
791 let state = mm_dsl::MeerkatMachineState {
792 lifecycle_phase: mm_dsl::MeerkatPhase::Attached,
793 session_id: Some(mm_dsl::SessionId("session-1".to_string())),
794 silent_intent_overrides: BTreeSet::from(["probe.silent".to_string()]),
795 ..Default::default()
796 };
797 let authority = Arc::new(Mutex::new(
798 mm_dsl::MeerkatMachineAuthority::recover_from_state(state)
799 .expect("test MeerkatMachine state must be recoverable"),
800 ));
801 let handle =
802 RuntimePeerCommsHandle::new(Arc::new(HandleDslAuthority::from_shared(authority)));
803
804 let sender_peer_id = meerkat_core::comms::PeerId::new();
805 let admission = handle
806 .classify_external_envelope(PeerIngressEnvelopeFacts {
807 item_id: "request-1".to_string(),
808 from_peer: "peer-1".to_string(),
809 from_peer_id: sender_peer_id,
810 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
811 intent: "probe.silent".to_string(),
812 params: serde_json::json!({}),
813 },
814 })
815 .expect("attached session should classify peer ingress");
816
817 assert_eq!(
818 admission.classification.class,
819 meerkat_core::PeerInputClass::SilentRequest
820 );
821 assert_eq!(
822 admission.classification.auth,
823 meerkat_core::PeerIngressAuthDecision::Required
824 );
825 assert_eq!(admission.request_id.as_deref(), Some("request-1"));
826 assert_eq!(
829 admission.from_peer_id,
830 Some(sender_peer_id),
831 "machine classification must echo the canonical sender peer id"
832 );
833 }
834
835 #[test]
836 fn machine_emits_actionable_bit_matching_grouping_for_classified_envelopes() {
837 let handle = handle_for_phase(mm_dsl::MeerkatPhase::Attached);
841
842 let message = handle
844 .classify_external_envelope(PeerIngressEnvelopeFacts {
845 item_id: "m1".to_string(),
846 from_peer: "peer-1".to_string(),
847 from_peer_id: meerkat_core::comms::PeerId::new(),
848 kind: meerkat_core::PeerIngressEnvelopeKind::Message {
849 body: "hi".to_string(),
850 },
851 })
852 .expect("attached should classify message");
853 assert_eq!(
854 message.classification.class,
855 meerkat_core::PeerInputClass::ActionableMessage
856 );
857 assert!(message.classification.actionable);
858
859 let request = handle
861 .classify_external_envelope(PeerIngressEnvelopeFacts {
862 item_id: "r1".to_string(),
863 from_peer: "peer-1".to_string(),
864 from_peer_id: meerkat_core::comms::PeerId::new(),
865 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
866 intent: "do.work".to_string(),
867 params: serde_json::json!({}),
868 },
869 })
870 .expect("attached should classify request");
871 assert_eq!(
872 request.classification.class,
873 meerkat_core::PeerInputClass::ActionableRequest
874 );
875 assert!(request.classification.actionable);
876
877 let response = handle
881 .classify_external_envelope(PeerIngressEnvelopeFacts {
882 item_id: "resp-1".to_string(),
883 from_peer: "peer-1".to_string(),
884 from_peer_id: meerkat_core::comms::PeerId::new(),
885 kind: meerkat_core::PeerIngressEnvelopeKind::Response {
886 in_reply_to: uuid::Uuid::new_v4().to_string(),
887 status: meerkat_core::ResponseStatus::Completed,
888 result: serde_json::json!({}),
889 },
890 })
891 .expect("attached should classify response");
892 assert_eq!(
893 response.classification.class,
894 meerkat_core::PeerInputClass::ResponseTerminal
895 );
896 assert!(response.classification.actionable);
897
898 let lifecycle = handle
900 .classify_external_envelope(PeerIngressEnvelopeFacts {
901 item_id: "lc-1".to_string(),
902 from_peer: "orchestrator".to_string(),
903 from_peer_id: meerkat_core::comms::PeerId::new(),
904 kind: meerkat_core::PeerIngressEnvelopeKind::Lifecycle {
905 kind: meerkat_core::comms::PeerLifecycleKind::PeerAdded,
906 params: serde_json::json!({ "peer": "worker-1" }),
907 },
908 })
909 .expect("attached should classify lifecycle");
910 assert_eq!(
911 lifecycle.classification.class,
912 meerkat_core::PeerInputClass::PeerLifecycleAdded
913 );
914 assert!(!lifecycle.classification.actionable);
915
916 let ack = handle
918 .classify_external_envelope(PeerIngressEnvelopeFacts {
919 item_id: "ack-1".to_string(),
920 from_peer: "peer-1".to_string(),
921 from_peer_id: meerkat_core::comms::PeerId::new(),
922 kind: meerkat_core::PeerIngressEnvelopeKind::Ack {
923 in_reply_to: uuid::Uuid::new_v4().to_string(),
924 },
925 })
926 .expect("attached should classify ack");
927 assert_eq!(ack.classification.class, meerkat_core::PeerInputClass::Ack);
928 assert!(!ack.classification.actionable);
929
930 let plain = handle
932 .classify_plain_event(PeerIngressPlainEventFacts {
933 source_name: "external".to_string(),
934 body: "event".to_string(),
935 })
936 .expect("attached should classify plain event");
937 assert_eq!(
938 plain.classification.class,
939 meerkat_core::PeerInputClass::PlainEvent
940 );
941 assert!(plain.classification.actionable);
942
943 let silent_state = mm_dsl::MeerkatMachineState {
945 lifecycle_phase: mm_dsl::MeerkatPhase::Attached,
946 session_id: Some(mm_dsl::SessionId("session-1".to_string())),
947 silent_intent_overrides: BTreeSet::from(["probe.silent".to_string()]),
948 ..Default::default()
949 };
950 let silent_authority = Arc::new(Mutex::new(
951 mm_dsl::MeerkatMachineAuthority::recover_from_state(silent_state).expect("recoverable"),
952 ));
953 let silent_handle = RuntimePeerCommsHandle::new(Arc::new(HandleDslAuthority::from_shared(
954 silent_authority,
955 )));
956 let silent = silent_handle
957 .classify_external_envelope(PeerIngressEnvelopeFacts {
958 item_id: "s1".to_string(),
959 from_peer: "peer-1".to_string(),
960 from_peer_id: meerkat_core::comms::PeerId::new(),
961 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
962 intent: "probe.silent".to_string(),
963 params: serde_json::json!({}),
964 },
965 })
966 .expect("attached should classify silent request");
967 assert_eq!(
968 silent.classification.class,
969 meerkat_core::PeerInputClass::SilentRequest
970 );
971 assert!(!silent.classification.actionable);
972 }
973
974 #[test]
975 fn runtime_peer_comms_handle_lifecycle_subject_is_parsed_typed_at_ingress() {
976 let state = mm_dsl::MeerkatMachineState {
981 lifecycle_phase: mm_dsl::MeerkatPhase::Attached,
982 session_id: Some(mm_dsl::SessionId("session-1".to_string())),
983 ..Default::default()
984 };
985 let authority = Arc::new(Mutex::new(
986 mm_dsl::MeerkatMachineAuthority::recover_from_state(state)
987 .expect("test MeerkatMachine state must be recoverable"),
988 ));
989 let handle =
990 RuntimePeerCommsHandle::new(Arc::new(HandleDslAuthority::from_shared(authority)));
991
992 let with_param = handle
993 .classify_external_envelope(PeerIngressEnvelopeFacts {
994 item_id: "request-param".to_string(),
995 from_peer: "orchestrator".to_string(),
996 from_peer_id: meerkat_core::comms::PeerId::new(),
997 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
998 intent: "mob.peer_added".to_string(),
999 params: serde_json::json!({ "peer": "worker-1" }),
1000 },
1001 })
1002 .expect("machine should classify lifecycle request");
1003 assert_eq!(with_param.lifecycle_peer.as_deref(), Some("worker-1"));
1004
1005 let with_spec = handle
1008 .classify_external_envelope(PeerIngressEnvelopeFacts {
1009 item_id: "request-spec".to_string(),
1010 from_peer: "orchestrator".to_string(),
1011 from_peer_id: meerkat_core::comms::PeerId::new(),
1012 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
1013 intent: "mob.peer_added".to_string(),
1014 params: serde_json::json!({
1015 "peer": "worker-1",
1016 "peer_spec": {
1017 "name": "mob/worker-1",
1018 "peer_id": "pid-worker-1",
1019 "address": "inproc://mob/worker-1"
1020 }
1021 }),
1022 },
1023 })
1024 .expect("machine should classify lifecycle request with peer_spec");
1025 assert_eq!(with_spec.lifecycle_peer.as_deref(), Some("mob/worker-1"));
1026
1027 let without_param = handle.classify_external_envelope(PeerIngressEnvelopeFacts {
1029 item_id: "request-missing-subject".to_string(),
1030 from_peer: "orchestrator".to_string(),
1031 from_peer_id: meerkat_core::comms::PeerId::new(),
1032 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
1033 intent: "mob.peer_retired".to_string(),
1034 params: serde_json::json!({}),
1035 },
1036 });
1037 let error = without_param
1038 .expect_err("lifecycle request without a peer subject must be rejected at ingress");
1039 assert_eq!(error.context, "PeerCommsHandle::classify_external_envelope");
1040
1041 let empty_param = handle.classify_external_envelope(PeerIngressEnvelopeFacts {
1043 item_id: "lifecycle-empty".to_string(),
1044 from_peer: "orchestrator".to_string(),
1045 from_peer_id: meerkat_core::comms::PeerId::new(),
1046 kind: meerkat_core::PeerIngressEnvelopeKind::Lifecycle {
1047 kind: meerkat_core::comms::PeerLifecycleKind::PeerUnwired,
1048 params: serde_json::json!({ "peer": "" }),
1049 },
1050 });
1051 let error = empty_param
1052 .expect_err("lifecycle event with an empty peer subject must be rejected at ingress");
1053 assert_eq!(error.context, "PeerCommsHandle::classify_external_envelope");
1054
1055 let malformed = handle.classify_external_envelope(PeerIngressEnvelopeFacts {
1057 item_id: "lifecycle-malformed".to_string(),
1058 from_peer: "orchestrator".to_string(),
1059 from_peer_id: meerkat_core::comms::PeerId::new(),
1060 kind: meerkat_core::PeerIngressEnvelopeKind::Lifecycle {
1061 kind: meerkat_core::comms::PeerLifecycleKind::PeerAdded,
1062 params: serde_json::json!({ "peer": "worker-1", "unexpected": true }),
1063 },
1064 });
1065 malformed.expect_err("malformed lifecycle params must be rejected at ingress");
1066 }
1067
1068 #[test]
1069 fn runtime_peer_comms_handle_classifies_idle_lifecycle_without_opening_peer_work() {
1070 let handle = handle_for_phase(mm_dsl::MeerkatPhase::Idle);
1071
1072 let retired_notice = handle
1073 .classify_external_envelope(PeerIngressEnvelopeFacts {
1074 item_id: "lifecycle-retired".to_string(),
1075 from_peer: "orchestrator".to_string(),
1076 from_peer_id: meerkat_core::comms::PeerId::new(),
1077 kind: meerkat_core::PeerIngressEnvelopeKind::Lifecycle {
1078 kind: meerkat_core::comms::PeerLifecycleKind::PeerRetired,
1079 params: serde_json::json!({ "peer": "worker-1" }),
1080 },
1081 })
1082 .expect("idle live session should classify mob lifecycle notices");
1083 assert_eq!(
1084 retired_notice.classification.class,
1085 meerkat_core::PeerInputClass::PeerLifecycleRetired
1086 );
1087 assert_eq!(
1088 retired_notice.classification.lifecycle_kind,
1089 Some(meerkat_core::comms::PeerLifecycleKind::PeerRetired)
1090 );
1091 assert_eq!(retired_notice.lifecycle_peer.as_deref(), Some("worker-1"));
1092 assert_eq!(retired_notice.request_id, None);
1093
1094 let added_request = handle
1095 .classify_external_envelope(PeerIngressEnvelopeFacts {
1096 item_id: "request-added".to_string(),
1097 from_peer: "orchestrator".to_string(),
1098 from_peer_id: meerkat_core::comms::PeerId::new(),
1099 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
1100 intent: "mob.peer_added".to_string(),
1101 params: serde_json::json!({ "peer": "worker-2" }),
1102 },
1103 })
1104 .expect("idle live session should classify lifecycle requests");
1105 assert_eq!(
1106 added_request.classification.class,
1107 meerkat_core::PeerInputClass::PeerLifecycleAdded
1108 );
1109 assert_eq!(added_request.lifecycle_peer.as_deref(), Some("worker-2"));
1110 assert_eq!(added_request.request_id.as_deref(), Some("request-added"));
1111
1112 let work_admission = handle.classify_external_envelope(PeerIngressEnvelopeFacts {
1113 item_id: "message-1".to_string(),
1114 from_peer: "peer-1".to_string(),
1115 from_peer_id: meerkat_core::comms::PeerId::new(),
1116 kind: meerkat_core::PeerIngressEnvelopeKind::Message {
1117 body: "wake up".to_string(),
1118 },
1119 });
1120 assert!(
1121 work_admission.is_err(),
1122 "idle lifecycle admission must not reopen normal peer work ingress"
1123 );
1124 }
1125
1126 #[test]
1127 fn runtime_peer_comms_handle_classifies_idle_supervisor_bridge() {
1128 let handle = handle_for_phase(mm_dsl::MeerkatPhase::Idle);
1129
1130 let admission = handle
1131 .classify_external_envelope(PeerIngressEnvelopeFacts {
1132 item_id: "supervisor-request".to_string(),
1133 from_peer: "mob/__mob_supervisor__".to_string(),
1134 from_peer_id: meerkat_core::comms::PeerId::new(),
1135 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
1136 intent: "supervisor.bridge".to_string(),
1137 params: serde_json::json!({}),
1138 },
1139 })
1140 .expect("idle session should classify supervisor bridge requests");
1141 assert_eq!(
1142 admission.classification.class,
1143 meerkat_core::PeerInputClass::ActionableRequest
1144 );
1145 assert_eq!(
1146 admission.classification.auth,
1147 meerkat_core::PeerIngressAuthDecision::Exempt(
1148 meerkat_core::PeerIngressAuthExemption::SupervisorBridge
1149 )
1150 );
1151
1152 let state = mm_dsl::MeerkatMachineState {
1153 lifecycle_phase: mm_dsl::MeerkatPhase::Idle,
1154 session_id: Some(mm_dsl::SessionId("session-1".to_string())),
1155 silent_intent_overrides: BTreeSet::from(["supervisor.bridge".to_string()]),
1156 ..Default::default()
1157 };
1158 let authority = Arc::new(Mutex::new(
1159 mm_dsl::MeerkatMachineAuthority::recover_from_state(state)
1160 .expect("test MeerkatMachine state must be recoverable"),
1161 ));
1162 let silent_handle =
1163 RuntimePeerCommsHandle::new(Arc::new(HandleDslAuthority::from_shared(authority)));
1164
1165 let silent = silent_handle
1166 .classify_external_envelope(PeerIngressEnvelopeFacts {
1167 item_id: "supervisor-silent".to_string(),
1168 from_peer: "mob/__mob_supervisor__".to_string(),
1169 from_peer_id: meerkat_core::comms::PeerId::new(),
1170 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
1171 intent: "supervisor.bridge".to_string(),
1172 params: serde_json::json!({}),
1173 },
1174 })
1175 .expect("idle session should classify silent supervisor bridge requests");
1176 assert_eq!(
1177 silent.classification.class,
1178 meerkat_core::PeerInputClass::SilentRequest
1179 );
1180 assert_eq!(
1181 silent.classification.auth,
1182 meerkat_core::PeerIngressAuthDecision::Exempt(
1183 meerkat_core::PeerIngressAuthExemption::SupervisorBridge
1184 )
1185 );
1186 }
1187
1188 #[test]
1189 fn runtime_peer_comms_handle_drains_terminal_cleanup_without_reopening_topology_adds() {
1190 for phase in [mm_dsl::MeerkatPhase::Retired, mm_dsl::MeerkatPhase::Stopped] {
1191 let handle = handle_for_phase(phase);
1192
1193 let retired_notice = handle
1194 .classify_external_envelope(PeerIngressEnvelopeFacts {
1195 item_id: "lifecycle-retired".to_string(),
1196 from_peer: "orchestrator".to_string(),
1197 from_peer_id: meerkat_core::comms::PeerId::new(),
1198 kind: meerkat_core::PeerIngressEnvelopeKind::Lifecycle {
1199 kind: meerkat_core::comms::PeerLifecycleKind::PeerRetired,
1200 params: serde_json::json!({ "peer": "worker-1" }),
1201 },
1202 })
1203 .expect("terminal sessions should drain peer-retired cleanup notices");
1204 assert_eq!(
1205 retired_notice.classification.class,
1206 meerkat_core::PeerInputClass::PeerLifecycleRetired
1207 );
1208
1209 let unwired_request = handle
1210 .classify_external_envelope(PeerIngressEnvelopeFacts {
1211 item_id: "request-unwired".to_string(),
1212 from_peer: "orchestrator".to_string(),
1213 from_peer_id: meerkat_core::comms::PeerId::new(),
1214 kind: meerkat_core::PeerIngressEnvelopeKind::Request {
1215 intent: "mob.peer_unwired".to_string(),
1216 params: serde_json::json!({ "peer": "worker-2" }),
1217 },
1218 })
1219 .expect("terminal sessions should drain peer-unwired cleanup requests");
1220 assert_eq!(
1221 unwired_request.classification.class,
1222 meerkat_core::PeerInputClass::PeerLifecycleUnwired
1223 );
1224
1225 let added_notice = handle.classify_external_envelope(PeerIngressEnvelopeFacts {
1226 item_id: "lifecycle-added".to_string(),
1227 from_peer: "orchestrator".to_string(),
1228 from_peer_id: meerkat_core::comms::PeerId::new(),
1229 kind: meerkat_core::PeerIngressEnvelopeKind::Lifecycle {
1230 kind: meerkat_core::comms::PeerLifecycleKind::PeerAdded,
1231 params: serde_json::json!({ "peer": "worker-3" }),
1232 },
1233 });
1234 assert!(
1235 added_notice.is_err(),
1236 "terminal cleanup admission must not accept new peer topology"
1237 );
1238 }
1239 }
1240
1241 #[test]
1242 fn runtime_peer_comms_handle_resolves_receive_authority_from_dsl() {
1243 let handle = handle_for_phase(mm_dsl::MeerkatPhase::Attached);
1244
1245 let admitted = handle
1246 .resolve_peer_ingress_receive(PeerIngressReceiveFacts {
1247 kind: meerkat_core::PeerIngressKind::Request,
1248 current_phase: meerkat_core::PeerIngressAuthorityPhase::Absent,
1249 auth_required: true,
1250 auth_exempt: false,
1251 trusted: true,
1252 queued_work_present: false,
1253 queue_closed: false,
1254 queue_capacity_available: true,
1255 })
1256 .expect("trusted receive should resolve");
1257 assert_eq!(
1258 admitted.outcome,
1259 meerkat_core::PeerIngressReceiveOutcome::Admitted
1260 );
1261 assert_eq!(
1262 admitted.admission_diagnostic,
1263 Some(meerkat_core::PeerIngressAdmissionDiagnostic::TrustedAtAdmission)
1264 );
1265 assert_eq!(
1266 admitted.authority_phase,
1267 meerkat_core::PeerIngressAuthorityPhase::Received
1268 );
1269
1270 let dropped = handle
1271 .resolve_peer_ingress_receive(PeerIngressReceiveFacts {
1272 kind: meerkat_core::PeerIngressKind::Request,
1273 current_phase: meerkat_core::PeerIngressAuthorityPhase::Absent,
1274 auth_required: true,
1275 auth_exempt: false,
1276 trusted: false,
1277 queued_work_present: false,
1278 queue_closed: false,
1279 queue_capacity_available: true,
1280 })
1281 .expect("untrusted receive should resolve as a typed drop");
1282 assert_eq!(
1283 dropped.outcome,
1284 meerkat_core::PeerIngressReceiveOutcome::DroppedUntrustedSender
1285 );
1286 assert_eq!(
1287 dropped.authority_phase,
1288 meerkat_core::PeerIngressAuthorityPhase::Dropped
1289 );
1290 }
1291
1292 #[test]
1293 fn runtime_peer_comms_handle_resolves_dequeue_phase_from_dsl() {
1294 let handle = handle_for_phase(mm_dsl::MeerkatPhase::Attached);
1295
1296 handle
1297 .resolve_peer_ingress_receive(PeerIngressReceiveFacts {
1298 kind: meerkat_core::PeerIngressKind::Request,
1299 current_phase: meerkat_core::PeerIngressAuthorityPhase::Absent,
1300 auth_required: true,
1301 auth_exempt: false,
1302 trusted: true,
1303 queued_work_present: false,
1304 queue_closed: false,
1305 queue_capacity_available: true,
1306 })
1307 .expect("trusted receive should seed Received phase");
1308
1309 let retained = handle
1310 .resolve_peer_ingress_dequeue(PeerIngressDequeueFacts {
1311 kind: meerkat_core::PeerIngressKind::Request,
1312 auth: meerkat_core::PeerIngressAuthDecision::Required,
1313 queued_work_remaining: true,
1314 })
1315 .expect("dequeue with queued work should resolve");
1316 assert_eq!(
1317 retained.authority_phase,
1318 meerkat_core::PeerIngressAuthorityPhase::Received
1319 );
1320
1321 let delivered = handle
1322 .resolve_peer_ingress_dequeue(PeerIngressDequeueFacts {
1323 kind: meerkat_core::PeerIngressKind::Request,
1324 auth: meerkat_core::PeerIngressAuthDecision::Required,
1325 queued_work_remaining: false,
1326 })
1327 .expect("empty dequeue should resolve");
1328 assert_eq!(
1329 delivered.authority_phase,
1330 meerkat_core::PeerIngressAuthorityPhase::Delivered
1331 );
1332 }
1333
1334 #[test]
1335 fn runtime_signal_builder_does_not_preselect_lifecycle_subject() {
1336 let source = include_str!("peer_comms.rs");
1337 let signal_builder = source
1338 .split("fn external_envelope_signal")
1339 .nth(1)
1340 .expect("signal builder should exist")
1341 .split("struct PeerIngressClassifiedEffect")
1342 .next()
1343 .expect("classified effect should follow signal builder");
1344 let forbidden_helper = ["peer", "lifecycle", "subject"].join("_");
1345
1346 assert!(
1347 signal_builder.contains("lifecycle_peer_param"),
1348 "runtime should pass the parsed lifecycle peer candidate"
1349 );
1350 assert!(
1351 !signal_builder.contains(&forbidden_helper),
1352 "runtime must not call the lifecycle subject selector before the machine"
1353 );
1354 assert!(
1355 !signal_builder.contains("facts.from_peer.as_str()"),
1356 "fallback peer must remain a machine input fact, not a preselected subject"
1357 );
1358 assert!(
1359 !signal_builder.contains("unwrap_or"),
1360 "runtime must not choose a lifecycle subject fallback before the machine"
1361 );
1362 }
1363}