Skip to main content

hydro_lang/sim/runtime/
observation.rs

1//! Top-level observation hooks ([`ObservationHook`]): hooks with no tick DFIR that are
2//! their own scheduling unit — running one *is* releasing data (e.g. `assume_ordering`
3//! on a non-tick stream, or a hooked top-level `fold`). Each is its own
4//! `SimObservation` in the scheduler. Their scripted decision/status types and
5//! [`ScriptableHook`] impls live alongside them.
6
7use std::cell::RefCell;
8use std::collections::VecDeque;
9use std::hash::Hash;
10use std::rc::Rc;
11
12use bolero::generator::bolero_generator::driver::object::Borrowed;
13use bolero::{ValueGenerator, produce};
14use dfir_rs::rustc_hash::FxHashMap;
15use dfir_rs::util::unsync::mpsc::Sender;
16
17use super::{
18    HookLocationMeta, ObservationHook, RuntimeHook, ScriptDecision, ScriptableHook,
19    ScriptableObservationHook, TruncatedLabeledVecDebug, TruncatedVecDebug, abort,
20    describe_keyed_pending, keyed_buffer_len, log_release,
21};
22
23/// Top-level (outside-tick) `assume_ordering` hooks release elements **one at a
24/// time** rather than shuffling the entire batch. This is the key mechanism for
25/// simulating causality in feedback cycles: when data flows through a network hop
26/// and cycles back (e.g. via `forward_ref`), the cycled-back result can arrive
27/// and interleave with elements that are still pending in the input queue.
28///
29/// For example, given input `[1, 2, 3]` where each element is mapped and sent
30/// through a network cycle, the simulator can explore orderings like
31/// `[1, map(1), 2, map(2), 3, ...]` -- the cycled-back `map(1)` arrives before `2`.
32///
33/// This is the only place where such causality is observable, because top-level
34/// unbounded streams are "maximally async" -- any non-atomic stream can be
35/// arbitrarily decoupled. The in-tick variants (`StreamOrderHook`,
36/// `KeyedStreamOrderHook`) shuffle the full batch instead, since within a tick
37/// all data is available simultaneously.
38///
39/// The `sim_top_level_assume_ordering_*` tests are regression tests ensuring
40/// this one-at-a-time release behavior correctly explores causal interleavings.
41pub struct TopLevelStreamOrderHook<T> {
42    pub input: Rc<RefCell<VecDeque<T>>>,
43    pub to_release: Option<Vec<T>>,
44    pub output: Sender<T>,
45    pub location: HookLocationMeta,
46    pub format_item_debug: fn(&T) -> Option<String>,
47}
48
49impl<T> RuntimeHook for TopLevelStreamOrderHook<T> {
50    fn has_pending_input(&self) -> bool {
51        !self.input.borrow().is_empty()
52    }
53
54    fn only_one_possible_decision(&self) -> bool {
55        // A sole buffered element has exactly one possible release; ordering only
56        // becomes a choice with two or more.
57        self.input.borrow().len() <= 1
58    }
59
60    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
61        if let Some(to_release) = self.to_release.take() {
62            if !to_release.is_empty()
63                && let Some(log_writer) = log_writer
64            {
65                let HookLocationMeta {
66                    location: batch_location,
67                    line,
68                    caret_indent,
69                } = self.location;
70                let note_str = format!(
71                    "^ observed non-deterministic order: {:?}",
72                    TruncatedVecDebug(
73                        RefCell::new(Some(to_release.iter())),
74                        8,
75                        self.format_item_debug
76                    )
77                );
78
79                let _ = writeln!(log_writer);
80                log_release(
81                    log_writer,
82                    batch_location,
83                    line,
84                    caret_indent,
85                    &note_str,
86                    colored::Color::Green,
87                );
88            }
89
90            for item in to_release {
91                self.output.try_send(item).unwrap();
92            }
93        } else {
94            panic!("No decision to release");
95        }
96    }
97
98    fn location_meta(&self) -> HookLocationMeta {
99        self.location
100    }
101}
102
103impl<T> ObservationHook for TopLevelStreamOrderHook<T> {
104    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
105        let mut current_input = self.input.borrow_mut();
106        // Instead of a full shuffle, we only release one element at a time
107        // in order to handle possible feedback cycles.
108        let idx = (0..current_input.len()).generate(driver).unwrap();
109        let item = current_input.remove(idx).unwrap();
110        self.to_release = Some(vec![item]);
111    }
112}
113
114/// A scripted decision for a top-level `assume_ordering` observation.
115#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
116pub enum TopLevelOrderingDecision<T> {
117    Next(T),
118}
119
120impl<T> ScriptDecision for TopLevelOrderingDecision<T>
121where
122    T: serde::Serialize + serde::de::DeserializeOwned,
123{
124    fn describe(&self) -> String {
125        "next(value)".to_owned()
126    }
127}
128
129#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
130pub struct OrderingStatus {
131    pub buffered: usize,
132}
133
134impl<T> ScriptableHook for TopLevelStreamOrderHook<T>
135where
136    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
137{
138    type Decision = TopLevelOrderingDecision<T>;
139    type Status = OrderingStatus;
140
141    fn is_honorable(&self, decision: &Self::Decision) -> Result<bool, String> {
142        let TopLevelOrderingDecision::Next(expected) = decision;
143        Ok(self.input.borrow().iter().any(|item| item == expected))
144    }
145
146    fn apply(&mut self, decision: Self::Decision) {
147        let mut input = self.input.borrow_mut();
148        let TopLevelOrderingDecision::Next(expected) = decision;
149        let index = input.iter().position(|item| item == &expected).unwrap();
150        self.to_release = Some(vec![input.remove(index).unwrap()]);
151    }
152
153    fn implicit(&mut self) {
154        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
155        // hooks; a top-level observation consists of exactly this hook, so it can never
156        // be forced to run without a scripted decision.
157        abort!("implicit decision invoked on a top-level ordering hook");
158    }
159
160    fn status(&self) -> Self::Status {
161        OrderingStatus {
162            buffered: self.input.borrow().len(),
163        }
164    }
165
166    fn describe_pending(&self) -> Option<String> {
167        let input = self.input.borrow();
168        (!input.is_empty()).then(|| {
169            format!(
170                "{} buffered ordering item(s): {:?}",
171                input.len(),
172                TruncatedVecDebug(RefCell::new(Some(input.iter())), 8, self.format_item_debug)
173            )
174        })
175    }
176}
177
178impl<T> ScriptableObservationHook for TopLevelStreamOrderHook<T> where
179    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq
180{
181}
182
183/// Hook for top-level folds. Selects a non-empty subset of buffered inputs to release,
184/// always permuting them to explore all orderings. Unselected elements remain
185/// in the buffer for future releases (modeling delayed/lossy inputs).
186pub struct TopLevelFoldHook<T> {
187    pub input: Rc<RefCell<VecDeque<T>>>,
188    pub to_release: Option<Vec<T>>,
189    pub output: Sender<Vec<T>>,
190    pub location: HookLocationMeta,
191    pub format_item_debug: fn(&T) -> Option<String>,
192}
193
194impl<T> RuntimeHook for TopLevelFoldHook<T> {
195    fn has_pending_input(&self) -> bool {
196        !self.input.borrow().is_empty()
197    }
198
199    fn only_one_possible_decision(&self) -> bool {
200        // The subset must be non-empty, so a sole buffered element is forced; subset
201        // choice and permutation only appear with two or more.
202        self.input.borrow().len() <= 1
203    }
204
205    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
206        if let Some(to_release) = self.to_release.take() {
207            if !to_release.is_empty()
208                && let Some(log_writer) = log_writer
209            {
210                let HookLocationMeta {
211                    location: batch_location,
212                    line,
213                    caret_indent,
214                } = self.location;
215                let note_str = format!(
216                    "^ fold input batch (permuted): {:?}",
217                    TruncatedVecDebug(
218                        RefCell::new(Some(to_release.iter())),
219                        8,
220                        self.format_item_debug
221                    )
222                );
223
224                let _ = writeln!(log_writer);
225                log_release(
226                    log_writer,
227                    batch_location,
228                    line,
229                    caret_indent,
230                    &note_str,
231                    colored::Color::Green,
232                );
233            }
234
235            self.output.try_send(to_release).unwrap();
236        } else {
237            panic!("No decision to release");
238        }
239    }
240
241    fn location_meta(&self) -> HookLocationMeta {
242        self.location
243    }
244}
245
246impl<T> ObservationHook for TopLevelFoldHook<T> {
247    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
248        let mut current_input = self.input.borrow_mut();
249
250        // Select a non-empty subset: for each element, decide include/exclude.
251        // Only force inclusion on the last element if nothing was selected yet.
252        let mut selected = Vec::new();
253        let mut remaining = VecDeque::new();
254
255        let len = current_input.len();
256        for (i, item) in current_input.drain(..).enumerate() {
257            let is_last = i == len - 1;
258            let must_include = is_last && selected.is_empty();
259            if must_include || produce().generate(driver).unwrap() {
260                selected.push(item);
261            } else {
262                remaining.push_back(item);
263            }
264        }
265
266        // Put unselected elements back
267        *current_input = remaining;
268
269        // Always permute selected elements (Fisher-Yates) to explore all orderings.
270        // Even if commutativity is claimed via manual_proof!, the simulator is
271        // conservative and does not trust it — it still explores permutations.
272        {
273            let slen = selected.len();
274            for i in (1..slen).rev() {
275                let j = (0..=i).generate(driver).unwrap();
276                selected.swap(i, j);
277            }
278        }
279
280        self.to_release = Some(selected);
281    }
282}
283
284/// Scripting a top-level fold releases exactly **one named element per decision**
285/// (like a top-level `assume_ordering`), so intermediate fold states become observable
286/// exactly at the script's release points. The autonomous subset-and-permute path is
287/// never used: it is only sound when the fuzzer explores every subset split.
288impl<T> ScriptableHook for TopLevelFoldHook<T>
289where
290    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
291{
292    type Decision = TopLevelOrderingDecision<T>;
293    type Status = OrderingStatus;
294
295    fn is_honorable(&self, decision: &Self::Decision) -> Result<bool, String> {
296        let TopLevelOrderingDecision::Next(expected) = decision;
297        Ok(self.input.borrow().iter().any(|item| item == expected))
298    }
299
300    fn apply(&mut self, decision: Self::Decision) {
301        let TopLevelOrderingDecision::Next(expected) = decision;
302        let mut input = self.input.borrow_mut();
303        let index = input.iter().position(|item| item == &expected).unwrap();
304        self.to_release = Some(vec![input.remove(index).unwrap()]);
305    }
306
307    fn implicit(&mut self) {
308        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
309        // hooks; a top-level observation consists of exactly this hook, so it can never
310        // be forced to run without a scripted decision.
311        abort!("implicit decision invoked on a top-level fold hook");
312    }
313
314    fn status(&self) -> Self::Status {
315        OrderingStatus {
316            buffered: self.input.borrow().len(),
317        }
318    }
319
320    fn describe_pending(&self) -> Option<String> {
321        let input = self.input.borrow();
322        (!input.is_empty()).then(|| {
323            format!(
324                "{} buffered fold input(s): {:?}",
325                input.len(),
326                TruncatedVecDebug(RefCell::new(Some(input.iter())), 8, self.format_item_debug)
327            )
328        })
329    }
330}
331
332impl<T> ScriptableObservationHook for TopLevelFoldHook<T> where
333    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq
334{
335}
336
337/// Keyed variant of [`TopLevelStreamOrderHook`]. Same one-at-a-time release
338/// strategy to simulate causal interleavings -- see the comment on
339/// [`TopLevelStreamOrderHook`] for the full explanation.
340pub struct TopLevelKeyedStreamOrderHook<K: Hash + Eq + Clone, V> {
341    pub input: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
342    pub to_release: Option<Vec<(K, V)>>,
343    pub output: Sender<(K, V)>,
344    pub location: HookLocationMeta,
345    pub format_item_debug: fn(&(K, V)) -> Option<String>,
346}
347
348impl<K: Hash + Eq + Clone, V> RuntimeHook for TopLevelKeyedStreamOrderHook<K, V> {
349    fn has_pending_input(&self) -> bool {
350        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
351        !self.input.borrow().values().all(|q| q.is_empty())
352    }
353
354    fn only_one_possible_decision(&self) -> bool {
355        // A sole buffered element (across all keys) has exactly one possible release.
356        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
357        let total: usize = self.input.borrow().values().map(|q| q.len()).sum();
358        total <= 1
359    }
360
361    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
362        if let Some(to_release) = self.to_release.take() {
363            if !to_release.is_empty()
364                && let Some(log_writer) = log_writer
365            {
366                let HookLocationMeta {
367                    location: batch_location,
368                    line,
369                    caret_indent,
370                } = self.location;
371                let note_str = format!(
372                    "^ observed non-deterministic order: {:?}",
373                    TruncatedVecDebug(
374                        RefCell::new(Some(to_release.iter())),
375                        8,
376                        self.format_item_debug
377                    )
378                );
379
380                let _ = writeln!(log_writer);
381                log_release(
382                    log_writer,
383                    batch_location,
384                    line,
385                    caret_indent,
386                    &note_str,
387                    colored::Color::Green,
388                );
389            }
390
391            for item in to_release {
392                self.output.try_send(item).unwrap();
393            }
394        } else {
395            panic!("No decision to release");
396        }
397    }
398
399    fn location_meta(&self) -> HookLocationMeta {
400        self.location
401    }
402}
403
404impl<K: Hash + Eq + Clone, V> ObservationHook for TopLevelKeyedStreamOrderHook<K, V> {
405    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
406        let mut current_input = self.input.borrow_mut();
407
408        // Collect non-empty keys with their queue lengths
409        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
410        let nonempty_keys: Vec<(K, usize)> = current_input
411            .iter()
412            .filter(|(_, q)| !q.is_empty())
413            .map(|(k, q)| (k.clone(), q.len()))
414            .collect();
415
416        // Pick which key to release from
417        let key_idx = (0..nonempty_keys.len()).generate(driver).unwrap();
418        let (key, queue_len) = &nonempty_keys[key_idx];
419
420        // Pick which item from that key's queue
421        let item_idx = (0..*queue_len).generate(driver).unwrap();
422        let item = current_input
423            .get_mut(key)
424            .unwrap()
425            .remove(item_idx)
426            .unwrap();
427
428        self.to_release = Some(vec![(key.clone(), item)]);
429    }
430}
431
432/// Keyed variant of [`TopLevelStreamOrderHook`]'s scripting: a top-level keyed
433/// `assume_ordering` decision names one `(key, value)` entry to release; values within a
434/// key may be matched at any buffered position (the input is unordered within each key).
435impl<K, V> ScriptableHook for TopLevelKeyedStreamOrderHook<K, V>
436where
437    K: Hash + Eq + Clone + serde::Serialize + serde::de::DeserializeOwned,
438    V: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
439{
440    type Decision = TopLevelOrderingDecision<(K, V)>;
441    type Status = OrderingStatus;
442
443    fn is_honorable(&self, decision: &Self::Decision) -> Result<bool, String> {
444        let TopLevelOrderingDecision::Next((key, expected)) = decision;
445        Ok(self
446            .input
447            .borrow()
448            .get(key)
449            .is_some_and(|queue| queue.iter().any(|item| item == expected)))
450    }
451
452    fn apply(&mut self, decision: Self::Decision) {
453        let TopLevelOrderingDecision::Next((key, expected)) = decision;
454        let mut input = self.input.borrow_mut();
455        let queue = input.get_mut(&key).unwrap();
456        let index = queue.iter().position(|item| item == &expected).unwrap();
457        let item = queue.remove(index).unwrap();
458        self.to_release = Some(vec![(key, item)]);
459    }
460
461    fn implicit(&mut self) {
462        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
463        // hooks; a top-level observation consists of exactly this hook, so it can never
464        // be forced to run without a scripted decision.
465        abort!("implicit decision invoked on a top-level keyed ordering hook");
466    }
467
468    fn status(&self) -> Self::Status {
469        OrderingStatus {
470            buffered: keyed_buffer_len(&self.input.borrow()),
471        }
472    }
473
474    fn describe_pending(&self) -> Option<String> {
475        describe_keyed_pending(&self.input.borrow())
476    }
477}
478
479impl<K, V> ScriptableObservationHook for TopLevelKeyedStreamOrderHook<K, V>
480where
481    K: Hash + Eq + Clone + serde::Serialize + serde::de::DeserializeOwned,
482    V: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
483{
484}
485
486/// Top-level variant of [`PartiallyOrderedStreamHook`](super::PartiallyOrderedStreamHook). Same one-at-a-time release
487/// strategy as [`TopLevelKeyedStreamOrderHook`], but always takes from the FRONT
488/// of the chosen key's queue to preserve within-key order.
489pub struct TopLevelPartiallyOrderedStreamHook<K: Hash + Eq + Clone, V> {
490    pub input: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
491    pub to_release: Option<Vec<(K, V)>>,
492    pub output: Sender<(K, V)>,
493    pub location: HookLocationMeta,
494    pub format_item_debug: fn(&(K, V)) -> Option<String>,
495}
496
497impl<K: Hash + Eq + Clone, V> RuntimeHook for TopLevelPartiallyOrderedStreamHook<K, V> {
498    fn has_pending_input(&self) -> bool {
499        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
500        !self.input.borrow().values().all(|q| q.is_empty())
501    }
502
503    fn only_one_possible_decision(&self) -> bool {
504        // Within a key the front element is forced, so the only choice is which key
505        // releases next: a single non-empty key is fully forced regardless of depth.
506        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
507        let nonempty_keys = self
508            .input
509            .borrow()
510            .values()
511            .filter(|q| !q.is_empty())
512            .count();
513        nonempty_keys <= 1
514    }
515
516    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
517        if let Some(to_release) = self.to_release.take() {
518            if !to_release.is_empty()
519                && let Some(log_writer) = log_writer
520            {
521                let HookLocationMeta {
522                    location: batch_location,
523                    line,
524                    caret_indent,
525                } = self.location;
526                let note_str = format!(
527                    "^ observed partially-ordered interleaving: {:?}",
528                    TruncatedVecDebug(
529                        RefCell::new(Some(to_release.iter())),
530                        8,
531                        self.format_item_debug
532                    )
533                );
534
535                let _ = writeln!(log_writer);
536                log_release(
537                    log_writer,
538                    batch_location,
539                    line,
540                    caret_indent,
541                    &note_str,
542                    colored::Color::Green,
543                );
544            }
545
546            for item in to_release {
547                self.output.try_send(item).unwrap();
548            }
549        } else {
550            panic!("No decision to release");
551        }
552    }
553
554    fn location_meta(&self) -> HookLocationMeta {
555        self.location
556    }
557}
558
559impl<K: Hash + Eq + Clone, V> ObservationHook for TopLevelPartiallyOrderedStreamHook<K, V> {
560    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
561        let mut current_input = self.input.borrow_mut();
562
563        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
564        let nonempty_keys: Vec<K> = current_input
565            .iter()
566            .filter(|(_, q)| !q.is_empty())
567            .map(|(k, _)| k.clone())
568            .collect();
569
570        // Pick which key to release from
571        let key_idx = (0..nonempty_keys.len()).generate(driver).unwrap();
572        let key = &nonempty_keys[key_idx];
573
574        // Always take from the front to preserve within-key order
575        let item = current_input.get_mut(key).unwrap().pop_front().unwrap();
576
577        self.to_release = Some(vec![(key.clone(), item)]);
578    }
579}
580
581/// Scripting for a top-level `entries_partially_ordered`: each decision names one
582/// `(key, value)` entry to release, and the value must be the *front* of that key's
583/// buffered queue since within-key order is preserved.
584impl<K, V> ScriptableHook for TopLevelPartiallyOrderedStreamHook<K, V>
585where
586    K: Hash + Eq + Clone + serde::Serialize + serde::de::DeserializeOwned,
587    V: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
588{
589    type Decision = TopLevelOrderingDecision<(K, V)>;
590    type Status = OrderingStatus;
591
592    fn is_honorable(&self, decision: &Self::Decision) -> Result<bool, String> {
593        let TopLevelOrderingDecision::Next((key, expected)) = decision;
594        match self.input.borrow().get(key).and_then(VecDeque::front) {
595            Some(front) if front == expected => Ok(true),
596            // Within-key order is preserved, so a mismatching front can never be
597            // released ahead of the named value: the decision is permanently stuck.
598            Some(_) => Err(
599                "next: the named value did not match the front of its key's buffered queue \
600                 (within-key order is preserved by this operator)"
601                    .to_owned(),
602            ),
603            None => Ok(false),
604        }
605    }
606
607    fn apply(&mut self, decision: Self::Decision) {
608        let TopLevelOrderingDecision::Next((key, _expected)) = decision;
609        let mut input = self.input.borrow_mut();
610        let item = input.get_mut(&key).unwrap().pop_front().unwrap();
611        self.to_release = Some(vec![(key, item)]);
612    }
613
614    fn implicit(&mut self) {
615        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
616        // hooks; a top-level observation consists of exactly this hook, so it can never
617        // be forced to run without a scripted decision.
618        abort!("implicit decision invoked on a top-level partially-ordered hook");
619    }
620
621    fn status(&self) -> Self::Status {
622        OrderingStatus {
623            buffered: keyed_buffer_len(&self.input.borrow()),
624        }
625    }
626
627    fn describe_pending(&self) -> Option<String> {
628        describe_keyed_pending(&self.input.borrow())
629    }
630}
631
632impl<K, V> ScriptableObservationHook for TopLevelPartiallyOrderedStreamHook<K, V>
633where
634    K: Hash + Eq + Clone + serde::Serialize + serde::de::DeserializeOwned,
635    V: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
636{
637}
638
639/// Top-level merge-ordered hook. Releases one element at a time, picking from
640/// the front of either the first or second input queue. This preserves per-input
641/// order while allowing feedback cycles to deliver elements between releases.
642pub struct TopLevelMergeOrderedHook<T> {
643    pub first: Rc<RefCell<VecDeque<T>>>,
644    pub second: Rc<RefCell<VecDeque<T>>>,
645    pub to_release: Option<Vec<T>>,
646    pub release_source: Option<&'static str>,
647    pub output: Sender<T>,
648    pub location: HookLocationMeta,
649    pub format_item_debug: fn(&T) -> Option<String>,
650}
651
652impl<T> RuntimeHook for TopLevelMergeOrderedHook<T> {
653    fn has_pending_input(&self) -> bool {
654        !self.first.borrow().is_empty() || !self.second.borrow().is_empty()
655    }
656
657    fn only_one_possible_decision(&self) -> bool {
658        // Each side's front element is forced, so the only choice is which side
659        // releases next: it exists only when both sides are non-empty.
660        self.first.borrow().is_empty() || self.second.borrow().is_empty()
661    }
662
663    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
664        if let Some(to_release) = self.to_release.take() {
665            let source = self.release_source.take();
666            if !to_release.is_empty()
667                && let Some(log_writer) = log_writer
668            {
669                let HookLocationMeta {
670                    location: batch_location,
671                    line,
672                    caret_indent,
673                } = self.location;
674                let source_label = source.unwrap_or("?");
675
676                let labeled_iter = to_release.iter().map(|item| (source_label, item));
677
678                let note_str = format!(
679                    "^ observed non-deterministic merge order: {:?}",
680                    TruncatedLabeledVecDebug(
681                        RefCell::new(Some(labeled_iter)),
682                        8,
683                        self.format_item_debug
684                    )
685                );
686
687                let _ = writeln!(log_writer);
688                log_release(
689                    log_writer,
690                    batch_location,
691                    line,
692                    caret_indent,
693                    &note_str,
694                    colored::Color::Green,
695                );
696            }
697
698            for item in to_release {
699                self.output.try_send(item).unwrap();
700            }
701        } else {
702            panic!("No decision to release");
703        }
704    }
705
706    fn location_meta(&self) -> HookLocationMeta {
707        self.location
708    }
709}
710
711impl<T> ObservationHook for TopLevelMergeOrderedHook<T> {
712    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
713        let first_empty = self.first.borrow().is_empty();
714        let second_empty = self.second.borrow().is_empty();
715
716        let (item, source) = if first_empty {
717            (self.second.borrow_mut().pop_front().unwrap(), "r")
718        } else if second_empty {
719            (self.first.borrow_mut().pop_front().unwrap(), "l")
720        } else {
721            let take_second: bool = produce().generate(driver).unwrap();
722            if take_second {
723                (self.second.borrow_mut().pop_front().unwrap(), "r")
724            } else {
725                (self.first.borrow_mut().pop_front().unwrap(), "l")
726            }
727        };
728
729        self.to_release = Some(vec![item]);
730        self.release_source = Some(source);
731    }
732}
733
734/// A scripted decision for a top-level `merge_ordered` observation: which input the next
735/// released element comes from. The element may be named ([`Self::First`] /
736/// [`Self::Second`], asserting it equals the front of that input's buffer) or released
737/// positionally ([`Self::FirstNext`] / [`Self::SecondNext`], releasing the front without
738/// asserting its value). `S` selects *which* front for keyed merges (the key), and is `()`
739/// for plain streams.
740#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
741pub enum MergeDecision<T, S = ()> {
742    /// Release the front of the first (left) input's buffer, which must equal this value.
743    First(T),
744    /// Release the front of the second (right) input's buffer, which must equal this value.
745    Second(T),
746    /// Release the front of the first (left) input's buffer, whatever it is.
747    FirstNext(S),
748    /// Release the front of the second (right) input's buffer, whatever it is.
749    SecondNext(S),
750}
751
752impl<T, S> ScriptDecision for MergeDecision<T, S>
753where
754    T: serde::Serialize + serde::de::DeserializeOwned,
755    S: serde::Serialize + serde::de::DeserializeOwned,
756{
757    fn describe(&self) -> String {
758        match self {
759            MergeDecision::First(_) => "next_first(value)".to_owned(),
760            MergeDecision::Second(_) => "next_second(value)".to_owned(),
761            MergeDecision::FirstNext(_) => "advance_first(..)".to_owned(),
762            MergeDecision::SecondNext(_) => "advance_second(..)".to_owned(),
763        }
764    }
765}
766
767/// The pending-input view a top-level merge hook reports to its test-side handle (see
768/// [`ScriptableHook::status`]), used by `pause_until_*` predicates.
769#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
770pub struct MergeStatus {
771    /// The number of elements buffered from the first (left) input.
772    pub first_buffered: usize,
773    /// The number of elements buffered from the second (right) input.
774    pub second_buffered: usize,
775}
776
777impl<T> ScriptableHook for TopLevelMergeOrderedHook<T>
778where
779    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
780{
781    type Decision = MergeDecision<T>;
782    type Status = MergeStatus;
783
784    fn is_honorable(&self, decision: &MergeDecision<T>) -> Result<bool, String> {
785        let (buffer, expected) = match decision {
786            MergeDecision::First(expected) => (&self.first, Some(expected)),
787            MergeDecision::Second(expected) => (&self.second, Some(expected)),
788            MergeDecision::FirstNext(()) => (&self.first, None),
789            MergeDecision::SecondNext(()) => (&self.second, None),
790        };
791        match (buffer.borrow().front(), expected) {
792            (Some(_), None) => Ok(true),
793            (Some(front), Some(expected)) if front == expected => Ok(true),
794            // Per-input order is preserved, so a mismatching front can never be released
795            // ahead of the named value: the decision is permanently stuck.
796            (Some(_), Some(_)) => Err(
797                "next_first/next_second: the named value did not match the front of that \
798                 input's buffer (per-input order is preserved by merge_ordered)"
799                    .to_owned(),
800            ),
801            (None, _) => Ok(false),
802        }
803    }
804
805    fn apply(&mut self, decision: MergeDecision<T>) {
806        let (item, source) = match decision {
807            MergeDecision::First(_) | MergeDecision::FirstNext(()) => {
808                (self.first.borrow_mut().pop_front().unwrap(), "l")
809            }
810            MergeDecision::Second(_) | MergeDecision::SecondNext(()) => {
811                (self.second.borrow_mut().pop_front().unwrap(), "r")
812            }
813        };
814        self.to_release = Some(vec![item]);
815        self.release_source = Some(source);
816    }
817
818    fn implicit(&mut self) {
819        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
820        // hooks; a top-level observation consists of exactly this hook, so it can never
821        // be forced to run without a scripted decision.
822        abort!("implicit decision invoked on a top-level merge-ordered hook");
823    }
824
825    fn status(&self) -> MergeStatus {
826        MergeStatus {
827            first_buffered: self.first.borrow().len(),
828            second_buffered: self.second.borrow().len(),
829        }
830    }
831
832    fn describe_pending(&self) -> Option<String> {
833        let first = self.first.borrow();
834        let second = self.second.borrow();
835        (!first.is_empty() || !second.is_empty()).then(|| {
836            format!(
837                "{} + {} buffered item(s) (first + second input)",
838                first.len(),
839                second.len()
840            )
841        })
842    }
843}
844
845impl<T> ScriptableObservationHook for TopLevelMergeOrderedHook<T> where
846    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq
847{
848}
849
850/// Keyed variant of [`TopLevelMergeOrderedHook`]. Releases one element at a
851/// time, picking from the front of some key's queue in either the first or the
852/// second input. This preserves per-input order within each key while allowing
853/// arbitrary interleaving both across the two inputs and across keys (which is
854/// unconstrained for keyed streams), and lets feedback cycles deliver elements
855/// between releases.
856pub struct TopLevelKeyedMergeOrderedHook<K: Hash + Eq + Clone, V> {
857    pub first: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
858    pub second: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
859    pub to_release: Option<Vec<(K, V)>>,
860    pub release_source: Option<&'static str>,
861    pub output: Sender<(K, V)>,
862    pub location: HookLocationMeta,
863    pub format_item_debug: fn(&(K, V)) -> Option<String>,
864}
865
866impl<K: Hash + Eq + Clone, V> RuntimeHook for TopLevelKeyedMergeOrderedHook<K, V> {
867    fn has_pending_input(&self) -> bool {
868        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
869        let first_nonempty = !self.first.borrow().values().all(|q| q.is_empty());
870        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
871        let second_nonempty = !self.second.borrow().values().all(|q| q.is_empty());
872        first_nonempty || second_nonempty
873    }
874
875    fn only_one_possible_decision(&self) -> bool {
876        // Each (side, key) queue's front element is forced, so the only choice is which
877        // queue releases next.
878        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
879        let first_count = self
880            .first
881            .borrow()
882            .values()
883            .filter(|q| !q.is_empty())
884            .count();
885        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
886        let second_count = self
887            .second
888            .borrow()
889            .values()
890            .filter(|q| !q.is_empty())
891            .count();
892        first_count + second_count <= 1
893    }
894
895    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
896        if let Some(to_release) = self.to_release.take() {
897            let source = self.release_source.take();
898            if !to_release.is_empty()
899                && let Some(log_writer) = log_writer
900            {
901                let HookLocationMeta {
902                    location: batch_location,
903                    line,
904                    caret_indent,
905                } = self.location;
906                let source_label = source.unwrap_or("?");
907
908                let labeled_iter = to_release.iter().map(|item| (source_label, item));
909
910                let note_str = format!(
911                    "^ observed non-deterministic merge order: {:?}",
912                    TruncatedLabeledVecDebug(
913                        RefCell::new(Some(labeled_iter)),
914                        8,
915                        self.format_item_debug
916                    )
917                );
918
919                let _ = writeln!(log_writer);
920                log_release(
921                    log_writer,
922                    batch_location,
923                    line,
924                    caret_indent,
925                    &note_str,
926                    colored::Color::Green,
927                );
928            }
929
930            for item in to_release {
931                self.output.try_send(item).unwrap();
932            }
933        } else {
934            panic!("No decision to release");
935        }
936    }
937
938    fn location_meta(&self) -> HookLocationMeta {
939        self.location
940    }
941}
942
943impl<K: Hash + Eq + Clone, V> ObservationHook for TopLevelKeyedMergeOrderedHook<K, V> {
944    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
945        // Collect candidates: for each non-empty key queue in either input, we
946        // can release its front element. `false` = first input, `true` = second.
947        let mut candidates: Vec<(bool, K)> = Vec::new();
948        {
949            #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
950            for (k, q) in self.first.borrow().iter() {
951                if !q.is_empty() {
952                    candidates.push((false, k.clone()));
953                }
954            }
955            #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
956            for (k, q) in self.second.borrow().iter() {
957                if !q.is_empty() {
958                    candidates.push((true, k.clone()));
959                }
960            }
961        }
962
963        let idx = (0..candidates.len()).generate(driver).unwrap();
964        let (take_second, key) = &candidates[idx];
965        let take_second = *take_second;
966
967        let item = if take_second {
968            self.second
969                .borrow_mut()
970                .get_mut(key)
971                .unwrap()
972                .pop_front()
973                .unwrap()
974        } else {
975            self.first
976                .borrow_mut()
977                .get_mut(key)
978                .unwrap()
979                .pop_front()
980                .unwrap()
981        };
982
983        self.to_release = Some(vec![(key.clone(), item)]);
984        self.release_source = Some(if take_second { "r" } else { "l" });
985    }
986}
987
988impl<K, V> ScriptableHook for TopLevelKeyedMergeOrderedHook<K, V>
989where
990    K: Hash + Eq + Clone + serde::Serialize + serde::de::DeserializeOwned,
991    V: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
992{
993    type Decision = MergeDecision<(K, V), K>;
994    type Status = MergeStatus;
995
996    fn is_honorable(&self, decision: &MergeDecision<(K, V), K>) -> Result<bool, String> {
997        let (buffer, key, expected) = match decision {
998            MergeDecision::First((key, expected)) => (&self.first, key, Some(expected)),
999            MergeDecision::Second((key, expected)) => (&self.second, key, Some(expected)),
1000            MergeDecision::FirstNext(key) => (&self.first, key, None),
1001            MergeDecision::SecondNext(key) => (&self.second, key, None),
1002        };
1003        match (buffer.borrow().get(key).and_then(VecDeque::front), expected) {
1004            (Some(_), None) => Ok(true),
1005            (Some(front), Some(expected)) if front == expected => Ok(true),
1006            // Per-input order within a key is preserved, so a mismatching front can
1007            // never be released ahead of the named value: the decision is permanently
1008            // stuck.
1009            (Some(_), Some(_)) => Err(
1010                "next_first/next_second: the named value did not match the front of its \
1011                 key's buffer in that input (per-input within-key order is preserved by \
1012                 merge_ordered)"
1013                    .to_owned(),
1014            ),
1015            (None, _) => Ok(false),
1016        }
1017    }
1018
1019    fn apply(&mut self, decision: MergeDecision<(K, V), K>) {
1020        let (buffer, key, source) = match decision {
1021            MergeDecision::First((key, _)) => (&self.first, key, "l"),
1022            MergeDecision::Second((key, _)) => (&self.second, key, "r"),
1023            MergeDecision::FirstNext(key) => (&self.first, key, "l"),
1024            MergeDecision::SecondNext(key) => (&self.second, key, "r"),
1025        };
1026        let item = buffer
1027            .borrow_mut()
1028            .get_mut(&key)
1029            .unwrap()
1030            .pop_front()
1031            .unwrap();
1032        self.to_release = Some(vec![(key, item)]);
1033        self.release_source = Some(source);
1034    }
1035
1036    fn implicit(&mut self) {
1037        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
1038        // hooks; a top-level observation consists of exactly this hook, so it can never
1039        // be forced to run without a scripted decision.
1040        abort!("implicit decision invoked on a top-level keyed merge-ordered hook");
1041    }
1042
1043    fn status(&self) -> MergeStatus {
1044        MergeStatus {
1045            first_buffered: keyed_buffer_len(&self.first.borrow()),
1046            second_buffered: keyed_buffer_len(&self.second.borrow()),
1047        }
1048    }
1049
1050    fn describe_pending(&self) -> Option<String> {
1051        let first = keyed_buffer_len(&self.first.borrow());
1052        let second = keyed_buffer_len(&self.second.borrow());
1053        (first + second > 0).then(|| {
1054            format!(
1055                "{} + {} buffered item(s) (first + second input)",
1056                first, second
1057            )
1058        })
1059    }
1060}
1061
1062impl<K, V> ScriptableObservationHook for TopLevelKeyedMergeOrderedHook<K, V>
1063where
1064    K: Hash + Eq + Clone + serde::Serialize + serde::de::DeserializeOwned,
1065    V: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
1066{
1067}