1use 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
23pub 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 pub fn new(dsl: Arc<HandleDslAuthority>) -> Self {
60 Self {
61 dsl,
62 cleanup_observer: RwLock::new(None),
63 }
64 }
65
66 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 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 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}