Skip to main content

meerkat_runtime/handles/
interaction_stream.rs

1//! Runtime impl of [`meerkat_core::handles::InteractionStreamHandle`] (U6).
2//!
3//! Routes the interaction-stream lifecycle (`Reserved` / `Attached` /
4//! `Completed` / `Expired` / `ClosedEarly` / `Abandoned`) into the session's MeerkatMachine
5//! DSL `interaction_streams` substate map and fans the emitted
6//! `InteractionStreamCleanup` effects out to the installed shell-side
7//! observer, so the comms runtime's `interaction_stream_registry` becomes a
8//! pure projection of DSL truth — no shadow `state` field, no shell-side
9//! CAS, no TTL meaning hidden in registry maps.
10
11use std::sync::{Arc, RwLock, Weak};
12
13use meerkat_core::handles::{
14    DslTransitionError, InteractionStreamCleanupObserver, InteractionStreamHandle,
15};
16use meerkat_core::peer_correlation::{
17    InteractionStreamState as CoreInteractionStreamState, PeerCorrelationId,
18};
19
20use super::HandleDslAuthority;
21use crate::meerkat_machine::dsl as mm_dsl;
22
23/// Runtime-backed [`InteractionStreamHandle`] impl.
24///
25/// Every trait method routes to the corresponding DSL input on the session's
26/// shared MeerkatMachine authority. After the transition lands, emitted
27/// effects are scanned for `InteractionStreamCleanup` and dispatched to the
28/// installed [`InteractionStreamCleanupObserver`] (if any), closing the
29/// "terminal transition → effect → shell projection cleanup" loop.
30///
31/// The observer is held as a `Weak` reference because the canonical owner is
32/// the session's `CommsRuntime`, which in turn holds a strong handle pointer;
33/// storing the observer strongly would create a cycle preventing
34/// `CommsRuntime::drop` from firing on session teardown. `Weak::upgrade`
35/// returning `None` after teardown makes cleanup dispatch a no-op, which is
36/// the desired post-shutdown semantics.
37pub struct RuntimeInteractionStreamHandle {
38    dsl: Arc<HandleDslAuthority>,
39    cleanup_observer: RwLock<Option<Weak<dyn InteractionStreamCleanupObserver>>>,
40}
41
42impl std::fmt::Debug for RuntimeInteractionStreamHandle {
43    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
44        let observer_tag = self
45            .cleanup_observer
46            .read()
47            .ok()
48            .as_deref()
49            .and_then(|o| o.as_ref().map(|_| "<observer>"));
50        f.debug_struct("RuntimeInteractionStreamHandle")
51            .field("dsl", &self.dsl)
52            .field("cleanup_observer", &observer_tag)
53            .finish()
54    }
55}
56
57impl RuntimeInteractionStreamHandle {
58    /// Construct a handle backed by the session's shared DSL authority.
59    pub fn new(dsl: Arc<HandleDslAuthority>) -> Self {
60        Self {
61            dsl,
62            cleanup_observer: RwLock::new(None),
63        }
64    }
65
66    /// Construct a handle backed by an ephemeral DSL authority for tests and
67    /// minimal hosts that explicitly opt into machine-owned semantics.
68    pub fn ephemeral() -> Self {
69        Self::new(Arc::new(HandleDslAuthority::ephemeral()))
70    }
71
72    fn apply_input_and_dispatch_cleanup(
73        &self,
74        input: mm_dsl::MeerkatMachineInput,
75        context: &'static str,
76    ) -> Result<(), DslTransitionError> {
77        // Sample observer UNDER the DSL lock, dispatch OUTSIDE — same
78        // race-closing pattern as `session_context.rs`.
79        type CleanupTarget = Result<
80            (
81                PeerCorrelationId,
82                Option<meerkat_core::InteractionStreamAbandonReason>,
83            ),
84            String,
85        >;
86        let dispatch: Option<(
87            Arc<dyn InteractionStreamCleanupObserver>,
88            Vec<CleanupTarget>,
89        )> = self
90            .dsl
91            .apply_input_with_effects_and_sample(input, context, |effects| {
92                let observer_opt = self
93                    .cleanup_observer
94                    .read()
95                    .unwrap_or_else(std::sync::PoisonError::into_inner)
96                    .as_ref()
97                    .and_then(Weak::upgrade);
98                let observer = observer_opt?;
99                let targets: Vec<CleanupTarget> = effects
100                    .iter()
101                    .filter_map(|effect| match effect {
102                        mm_dsl::MeerkatMachineEffect::InteractionStreamCleanup {
103                            corr_id,
104                            abandon_reason,
105                        } => Some(match dsl_corr_id_to_core(corr_id.clone()) {
106                            Some(core_id) => Ok((core_id, (*abandon_reason).map(Into::into))),
107                            None => Err(corr_id.0.clone()),
108                        }),
109                        _ => None,
110                    })
111                    .collect();
112                Some((observer, targets))
113            })?;
114        if let Some((observer, targets)) = dispatch {
115            for target in targets {
116                match target {
117                    Ok((core_id, abandon_reason)) => {
118                        observer.on_interaction_stream_cleanup(core_id, abandon_reason);
119                    }
120                    Err(raw) => tracing::error!(
121                        raw = %raw,
122                        context = context,
123                        "InteractionStreamCleanup: DSL emitted a corr_id that is not a valid UUID — broken invariant; skipping observer dispatch"
124                    ),
125                }
126            }
127        }
128        Ok(())
129    }
130}
131
132fn dsl_corr_id_to_core(dsl_id: mm_dsl::PeerCorrelationId) -> Option<PeerCorrelationId> {
133    uuid::Uuid::parse_str(&dsl_id.0)
134        .ok()
135        .map(PeerCorrelationId::from_uuid)
136}
137
138impl InteractionStreamHandle for RuntimeInteractionStreamHandle {
139    fn reserved(&self, corr_id: PeerCorrelationId) -> Result<(), DslTransitionError> {
140        self.apply_input_and_dispatch_cleanup(
141            mm_dsl::MeerkatMachineInput::InteractionStreamReserved {
142                corr_id: corr_id.into(),
143            },
144            "InteractionStreamHandle::reserved",
145        )
146    }
147
148    fn attached(&self, corr_id: PeerCorrelationId) -> Result<(), DslTransitionError> {
149        self.apply_input_and_dispatch_cleanup(
150            mm_dsl::MeerkatMachineInput::InteractionStreamAttached {
151                corr_id: corr_id.into(),
152            },
153            "InteractionStreamHandle::attached",
154        )
155    }
156
157    fn completed(&self, corr_id: PeerCorrelationId) -> Result<(), DslTransitionError> {
158        self.apply_input_and_dispatch_cleanup(
159            mm_dsl::MeerkatMachineInput::InteractionStreamCompleted {
160                corr_id: corr_id.into(),
161            },
162            "InteractionStreamHandle::completed",
163        )
164    }
165
166    fn expired(&self, corr_id: PeerCorrelationId) -> Result<(), DslTransitionError> {
167        self.apply_input_and_dispatch_cleanup(
168            mm_dsl::MeerkatMachineInput::InteractionStreamExpired {
169                corr_id: corr_id.into(),
170            },
171            "InteractionStreamHandle::expired",
172        )
173    }
174
175    fn closed_early(&self, corr_id: PeerCorrelationId) -> Result<(), DslTransitionError> {
176        self.apply_input_and_dispatch_cleanup(
177            mm_dsl::MeerkatMachineInput::InteractionStreamClosedEarly {
178                corr_id: corr_id.into(),
179            },
180            "InteractionStreamHandle::closed_early",
181        )
182    }
183
184    fn abandoned(
185        &self,
186        corr_id: PeerCorrelationId,
187        reason: meerkat_core::InteractionStreamAbandonReason,
188    ) -> Result<(), DslTransitionError> {
189        self.apply_input_and_dispatch_cleanup(
190            mm_dsl::MeerkatMachineInput::InteractionStreamAbandoned {
191                corr_id: corr_id.into(),
192                reason: reason.into(),
193            },
194            "InteractionStreamHandle::abandoned",
195        )
196    }
197
198    fn state(&self, corr_id: PeerCorrelationId) -> Option<CoreInteractionStreamState> {
199        let dsl_key: mm_dsl::PeerCorrelationId = corr_id.into();
200        let snapshot = self.dsl.snapshot_state();
201        // Disjoint-set encoding (matches the DSL `interaction_stream_disjoint`
202        // discipline): a corr_id is in at most one of the two active sets.
203        // Terminal states (`Completed` / `Expired` / `ClosedEarly` /
204        // `Abandoned`) leave both
205        // sets and are never observable here — they surface only via the
206        // `InteractionStreamStateChanged` effect, like the peer-correlation
207        // sibling enums.
208        if snapshot.attached_interaction_streams.contains(&dsl_key) {
209            Some(CoreInteractionStreamState::Attached)
210        } else if snapshot.reserved_interaction_streams.contains(&dsl_key) {
211            Some(CoreInteractionStreamState::Reserved)
212        } else {
213            None
214        }
215    }
216
217    fn install_cleanup_observer(&self, observer: Arc<dyn InteractionStreamCleanupObserver>) {
218        *self
219            .cleanup_observer
220            .write()
221            .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Arc::downgrade(&observer));
222    }
223}
224
225#[cfg(test)]
226#[allow(clippy::expect_used, clippy::unwrap_used)]
227mod tests {
228    use super::*;
229    use std::sync::Mutex;
230
231    fn new_handle() -> RuntimeInteractionStreamHandle {
232        let mut authority = mm_dsl::MeerkatMachineAuthority::new();
233        authority
234            .apply_signal(mm_dsl::MeerkatMachineSignal::Initialize)
235            .unwrap();
236        mm_dsl::MeerkatMachineMutator::apply(
237            &mut authority,
238            mm_dsl::MeerkatMachineInput::RegisterSession {
239                session_id: mm_dsl::SessionId::from("interaction-stream-test".to_string()),
240            },
241        )
242        .unwrap();
243        let dsl = Arc::new(HandleDslAuthority::from_shared(Arc::new(Mutex::new(
244            authority,
245        ))));
246        RuntimeInteractionStreamHandle::new(dsl)
247    }
248
249    struct CleanupRecorder(
250        Mutex<
251            Vec<(
252                PeerCorrelationId,
253                Option<meerkat_core::InteractionStreamAbandonReason>,
254            )>,
255        >,
256    );
257
258    impl InteractionStreamCleanupObserver for CleanupRecorder {
259        fn on_interaction_stream_cleanup(
260            &self,
261            corr_id: PeerCorrelationId,
262            abandon_reason: Option<meerkat_core::InteractionStreamAbandonReason>,
263        ) {
264            self.0.lock().unwrap().push((corr_id, abandon_reason));
265        }
266    }
267
268    #[test]
269    fn abandoned_is_typed_terminal_from_reserved_or_attached() {
270        let handle = new_handle();
271        let recorder = Arc::new(CleanupRecorder(Mutex::new(Vec::new())));
272        handle.install_cleanup_observer(
273            recorder.clone() as Arc<dyn InteractionStreamCleanupObserver>
274        );
275
276        let reserved = PeerCorrelationId::new();
277        handle.reserved(reserved).unwrap();
278        handle
279            .abandoned(
280                reserved,
281                meerkat_core::InteractionStreamAbandonReason::SendFailed,
282            )
283            .unwrap();
284        assert_eq!(handle.state(reserved), None);
285
286        let attached = PeerCorrelationId::new();
287        handle.reserved(attached).unwrap();
288        handle.attached(attached).unwrap();
289        handle
290            .abandoned(
291                attached,
292                meerkat_core::InteractionStreamAbandonReason::TerminalDeliveryFailed,
293            )
294            .unwrap();
295        assert_eq!(handle.state(attached), None);
296        assert_eq!(
297            *recorder.0.lock().unwrap(),
298            vec![
299                (
300                    reserved,
301                    Some(meerkat_core::InteractionStreamAbandonReason::SendFailed),
302                ),
303                (
304                    attached,
305                    Some(meerkat_core::InteractionStreamAbandonReason::TerminalDeliveryFailed,),
306                ),
307            ]
308        );
309    }
310}