Skip to main content

meerkat_runtime/handles/
peer_comms.rs

1//! Runtime impl of [`meerkat_core::handles::PeerCommsHandle`].
2
3use 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/// Runtime-backed [`PeerCommsHandle`] impl.
16///
17/// Routes every trait method to the corresponding DSL signal / input on a
18/// dedicated per-session MeerkatMachine DSL authority.
19#[derive(Debug)]
20pub struct RuntimePeerCommsHandle {
21    dsl: Arc<HandleDslAuthority>,
22}
23
24impl RuntimePeerCommsHandle {
25    /// Construct a handle backed by the session's shared DSL authority.
26    pub fn new(dsl: Arc<HandleDslAuthority>) -> Self {
27        Self { dsl }
28    }
29
30    /// Install a generated peer-comms handle and its matching trust owner.
31    ///
32    /// The owner token stays on the generated install path; callers receive no
33    /// reusable token that could be copied into a handwritten handle.
34    #[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    /// Install a generated peer-comms handle and its matching trust owner.
46    ///
47    /// The owner token stays on the generated install path; callers receive no
48    /// reusable token that could be copied into a handwritten handle.
49    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    /// Construct a handle backed by an ephemeral DSL authority.
57    ///
58    /// See [`RuntimeTurnStateHandle::ephemeral`].
59    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        // Dismissal is a supervisor/runtime drain-lifecycle terminal, not an
78        // ingress-classifiable topology notice. It has no peer-ingress
79        // classification route, so reject it fail-closed here rather than
80        // laundering it into a topology class.
81        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
114/// Parse the typed peer lifecycle subject fail-closed at machine ingress
115/// (K15).
116///
117/// The wire owner is `meerkat_contracts::CommsPeerLifecycleParams`:
118/// `peer_spec` is the canonical typed peer identity when the sender has it;
119/// `peer` is the required presentation subject otherwise. Malformed params,
120/// a missing subject, or an empty subject are rejected here with a typed
121/// error — a lifecycle notice is never silently attributed to the sender
122/// name.
123fn 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
192/// Classify a peer-ingress request intent into the closed routing class the
193/// machine guards on. Intents outside the closed set (including the empty
194/// intent carried by non-request envelopes) map to `Other`; the raw intent
195/// string is still threaded through the signal for the open-set
196/// `silent_intent_overrides` membership check inside the machine.
197fn 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            // K15: lifecycle-classed requests must carry a typed peer
229            // subject; the subject is parsed fail-closed here at ingress
230            // (never defaulted from the sender name). Non-lifecycle request
231            // params remain opaque payload.
232            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
348/// Convert the machine-echoed canonical sender peer id back to the core
349/// domain `PeerId`, fail-closed on malformed output.
350fn 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        // R084: the canonical sender peer id on the admission is the
483        // machine-echoed effect fact, never the shell-local input copy.
484        // Envelope classification must echo it; fail closed otherwise.
485        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            // #96/#106: the machine owns the typed `Option<PeerId>` lifecycle
497            // subject; project it to the display-string carried by the core
498            // admission fact at this single boundary (the downstream comms layer
499            // keeps the name-shaped `Option<String>` contract).
500            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            // Plain events have no peer sender identity; the machine emits
520            // `None` on its plain-event classification transitions.
521            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        // intra-machine: no route; dispatcher not applicable (handle targets the meerkat DSL directly, not a CompositionDispatcher seam)
577        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 = &target;
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        // R084: the machine echoes the canonical sender peer id on its
827        // classification effect; the admission carries the machine fact.
828        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        // The machine is the authority on the actionable grouping; assert the
838        // emitted `actionable` bit matches the documented 7-of-12 grouping for
839        // every class the MeerkatMachine PeerIngress classifier can emit.
840        let handle = handle_for_phase(mm_dsl::MeerkatPhase::Attached);
841
842        // ActionableMessage -> actionable
843        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        // ActionableRequest -> actionable
860        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        // ResponseTerminal -> actionable (responses need a pending request, but
878        // the classification grouping is independent of correlation; the
879        // machine emits the bit on the class). Verify via a completed response.
880        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        // PeerLifecycleAdded -> NOT actionable
899        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        // Ack -> NOT actionable
917        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        // PlainEvent -> actionable
931        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        // SilentRequest -> NOT actionable
944        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        // K15: the lifecycle subject is parsed fail-closed from the typed
977        // `CommsPeerLifecycleParams` wire contract at machine ingress.
978        // `peer_spec` is the canonical identity when present; a missing or
979        // empty subject is a typed rejection — never a sender-name default.
980        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        // `peer_spec` is the canonical typed identity and wins over the
1006        // presentation `peer` field.
1007        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        // Missing subject: typed rejection at ingress, not sender attribution.
1028        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        // Empty subject: typed rejection at ingress, not sender attribution.
1042        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        // Malformed params (unknown field / wrong shape): typed rejection.
1056        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}