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, log_release,
20};
21
22/// Top-level (outside-tick) `assume_ordering` hooks release elements **one at a
23/// time** rather than shuffling the entire batch. This is the key mechanism for
24/// simulating causality in feedback cycles: when data flows through a network hop
25/// and cycles back (e.g. via `forward_ref`), the cycled-back result can arrive
26/// and interleave with elements that are still pending in the input queue.
27///
28/// For example, given input `[1, 2, 3]` where each element is mapped and sent
29/// through a network cycle, the simulator can explore orderings like
30/// `[1, map(1), 2, map(2), 3, ...]` -- the cycled-back `map(1)` arrives before `2`.
31///
32/// This is the only place where such causality is observable, because top-level
33/// unbounded streams are "maximally async" -- any non-atomic stream can be
34/// arbitrarily decoupled. The in-tick variants (`StreamOrderHook`,
35/// `KeyedStreamOrderHook`) shuffle the full batch instead, since within a tick
36/// all data is available simultaneously.
37///
38/// The `sim_top_level_assume_ordering_*` tests are regression tests ensuring
39/// this one-at-a-time release behavior correctly explores causal interleavings.
40pub struct TopLevelStreamOrderHook<T> {
41    pub input: Rc<RefCell<VecDeque<T>>>,
42    pub to_release: Option<Vec<T>>,
43    pub output: Sender<T>,
44    pub location: HookLocationMeta,
45    pub format_item_debug: fn(&T) -> Option<String>,
46}
47
48impl<T> RuntimeHook for TopLevelStreamOrderHook<T> {
49    fn has_pending_input(&self) -> bool {
50        !self.input.borrow().is_empty()
51    }
52
53    fn only_one_possible_decision(&self) -> bool {
54        // A sole buffered element has exactly one possible release; ordering only
55        // becomes a choice with two or more.
56        self.input.borrow().len() <= 1
57    }
58
59    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
60        if let Some(to_release) = self.to_release.take() {
61            if !to_release.is_empty()
62                && let Some(log_writer) = log_writer
63            {
64                let HookLocationMeta {
65                    location: batch_location,
66                    line,
67                    caret_indent,
68                } = self.location;
69                let note_str = format!(
70                    "^ observed non-deterministic order: {:?}",
71                    TruncatedVecDebug(
72                        RefCell::new(Some(to_release.iter())),
73                        8,
74                        self.format_item_debug
75                    )
76                );
77
78                let _ = writeln!(log_writer);
79                log_release(
80                    log_writer,
81                    batch_location,
82                    line,
83                    caret_indent,
84                    &note_str,
85                    colored::Color::Green,
86                );
87            }
88
89            for item in to_release {
90                self.output.try_send(item).unwrap();
91            }
92        } else {
93            panic!("No decision to release");
94        }
95    }
96
97    fn location_meta(&self) -> HookLocationMeta {
98        self.location
99    }
100}
101
102impl<T> ObservationHook for TopLevelStreamOrderHook<T> {
103    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
104        let mut current_input = self.input.borrow_mut();
105        // Instead of a full shuffle, we only release one element at a time
106        // in order to handle possible feedback cycles.
107        let idx = (0..current_input.len()).generate(driver).unwrap();
108        let item = current_input.remove(idx).unwrap();
109        self.to_release = Some(vec![item]);
110    }
111}
112
113/// A scripted decision for a top-level `assume_ordering` observation.
114#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
115pub enum TopLevelOrderingDecision<T> {
116    Next(T),
117}
118
119impl<T> ScriptDecision for TopLevelOrderingDecision<T>
120where
121    T: serde::Serialize + serde::de::DeserializeOwned,
122{
123    fn describe(&self) -> String {
124        "next(value)".to_owned()
125    }
126}
127
128#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
129pub struct OrderingStatus {
130    pub buffered: usize,
131}
132
133impl<T> ScriptableHook for TopLevelStreamOrderHook<T>
134where
135    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
136{
137    type Decision = TopLevelOrderingDecision<T>;
138    type Status = OrderingStatus;
139
140    fn is_honorable(&self, decision: &Self::Decision) -> Result<bool, String> {
141        let TopLevelOrderingDecision::Next(expected) = decision;
142        Ok(self.input.borrow().iter().any(|item| item == expected))
143    }
144
145    fn apply(&mut self, decision: Self::Decision) {
146        let mut input = self.input.borrow_mut();
147        let TopLevelOrderingDecision::Next(expected) = decision;
148        let index = input.iter().position(|item| item == &expected).unwrap();
149        self.to_release = Some(vec![input.remove(index).unwrap()]);
150    }
151
152    fn implicit(&mut self) {
153        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
154        // hooks; a top-level observation consists of exactly this hook, so it can never
155        // be forced to run without a scripted decision.
156        abort!("implicit decision invoked on a top-level ordering hook");
157    }
158
159    fn status(&self) -> Self::Status {
160        OrderingStatus {
161            buffered: self.input.borrow().len(),
162        }
163    }
164
165    fn describe_pending(&self) -> Option<String> {
166        let input = self.input.borrow();
167        (!input.is_empty()).then(|| {
168            format!(
169                "{} buffered ordering item(s): {:?}",
170                input.len(),
171                TruncatedVecDebug(RefCell::new(Some(input.iter())), 8, self.format_item_debug)
172            )
173        })
174    }
175}
176
177impl<T> ScriptableObservationHook for TopLevelStreamOrderHook<T> where
178    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq
179{
180}
181
182/// Hook for top-level folds. Selects a non-empty subset of buffered inputs to release,
183/// always permuting them to explore all orderings. Unselected elements remain
184/// in the buffer for future releases (modeling delayed/lossy inputs).
185pub struct TopLevelFoldHook<T> {
186    pub input: Rc<RefCell<VecDeque<T>>>,
187    pub to_release: Option<Vec<T>>,
188    pub output: Sender<Vec<T>>,
189    pub location: HookLocationMeta,
190    pub format_item_debug: fn(&T) -> Option<String>,
191}
192
193impl<T> RuntimeHook for TopLevelFoldHook<T> {
194    fn has_pending_input(&self) -> bool {
195        !self.input.borrow().is_empty()
196    }
197
198    fn only_one_possible_decision(&self) -> bool {
199        // The subset must be non-empty, so a sole buffered element is forced; subset
200        // choice and permutation only appear with two or more.
201        self.input.borrow().len() <= 1
202    }
203
204    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
205        if let Some(to_release) = self.to_release.take() {
206            if !to_release.is_empty()
207                && let Some(log_writer) = log_writer
208            {
209                let HookLocationMeta {
210                    location: batch_location,
211                    line,
212                    caret_indent,
213                } = self.location;
214                let note_str = format!(
215                    "^ fold input batch (permuted): {:?}",
216                    TruncatedVecDebug(
217                        RefCell::new(Some(to_release.iter())),
218                        8,
219                        self.format_item_debug
220                    )
221                );
222
223                let _ = writeln!(log_writer);
224                log_release(
225                    log_writer,
226                    batch_location,
227                    line,
228                    caret_indent,
229                    &note_str,
230                    colored::Color::Green,
231                );
232            }
233
234            self.output.try_send(to_release).unwrap();
235        } else {
236            panic!("No decision to release");
237        }
238    }
239
240    fn location_meta(&self) -> HookLocationMeta {
241        self.location
242    }
243}
244
245impl<T> ObservationHook for TopLevelFoldHook<T> {
246    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
247        let mut current_input = self.input.borrow_mut();
248
249        // Select a non-empty subset: for each element, decide include/exclude.
250        // Only force inclusion on the last element if nothing was selected yet.
251        let mut selected = Vec::new();
252        let mut remaining = VecDeque::new();
253
254        let len = current_input.len();
255        for (i, item) in current_input.drain(..).enumerate() {
256            let is_last = i == len - 1;
257            let must_include = is_last && selected.is_empty();
258            if must_include || produce().generate(driver).unwrap() {
259                selected.push(item);
260            } else {
261                remaining.push_back(item);
262            }
263        }
264
265        // Put unselected elements back
266        *current_input = remaining;
267
268        // Always permute selected elements (Fisher-Yates) to explore all orderings.
269        // Even if commutativity is claimed via manual_proof!, the simulator is
270        // conservative and does not trust it — it still explores permutations.
271        {
272            let slen = selected.len();
273            for i in (1..slen).rev() {
274                let j = (0..=i).generate(driver).unwrap();
275                selected.swap(i, j);
276            }
277        }
278
279        self.to_release = Some(selected);
280    }
281}
282
283/// Scripting a top-level fold releases exactly **one named element per decision**
284/// (like a top-level `assume_ordering`), so intermediate fold states become observable
285/// exactly at the script's release points. The autonomous subset-and-permute path is
286/// never used: it is only sound when the fuzzer explores every subset split.
287impl<T> ScriptableHook for TopLevelFoldHook<T>
288where
289    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq,
290{
291    type Decision = TopLevelOrderingDecision<T>;
292    type Status = OrderingStatus;
293
294    fn is_honorable(&self, decision: &Self::Decision) -> Result<bool, String> {
295        let TopLevelOrderingDecision::Next(expected) = decision;
296        Ok(self.input.borrow().iter().any(|item| item == expected))
297    }
298
299    fn apply(&mut self, decision: Self::Decision) {
300        let TopLevelOrderingDecision::Next(expected) = decision;
301        let mut input = self.input.borrow_mut();
302        let index = input.iter().position(|item| item == &expected).unwrap();
303        self.to_release = Some(vec![input.remove(index).unwrap()]);
304    }
305
306    fn implicit(&mut self) {
307        // Implicit behavior exists for tick hooks whose tick is forced to run by *other*
308        // hooks; a top-level observation consists of exactly this hook, so it can never
309        // be forced to run without a scripted decision.
310        abort!("implicit decision invoked on a top-level fold hook");
311    }
312
313    fn status(&self) -> Self::Status {
314        OrderingStatus {
315            buffered: self.input.borrow().len(),
316        }
317    }
318
319    fn describe_pending(&self) -> Option<String> {
320        let input = self.input.borrow();
321        (!input.is_empty()).then(|| {
322            format!(
323                "{} buffered fold input(s): {:?}",
324                input.len(),
325                TruncatedVecDebug(RefCell::new(Some(input.iter())), 8, self.format_item_debug)
326            )
327        })
328    }
329}
330
331impl<T> ScriptableObservationHook for TopLevelFoldHook<T> where
332    T: serde::Serialize + serde::de::DeserializeOwned + PartialEq
333{
334}
335
336/// Keyed variant of [`TopLevelStreamOrderHook`]. Same one-at-a-time release
337/// strategy to simulate causal interleavings -- see the comment on
338/// [`TopLevelStreamOrderHook`] for the full explanation.
339pub struct TopLevelKeyedStreamOrderHook<K: Hash + Eq + Clone, V> {
340    pub input: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
341    pub to_release: Option<Vec<(K, V)>>,
342    pub output: Sender<(K, V)>,
343    pub location: HookLocationMeta,
344    pub format_item_debug: fn(&(K, V)) -> Option<String>,
345}
346
347impl<K: Hash + Eq + Clone, V> RuntimeHook for TopLevelKeyedStreamOrderHook<K, V> {
348    fn has_pending_input(&self) -> bool {
349        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
350        !self.input.borrow().values().all(|q| q.is_empty())
351    }
352
353    fn only_one_possible_decision(&self) -> bool {
354        // A sole buffered element (across all keys) has exactly one possible release.
355        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
356        let total: usize = self.input.borrow().values().map(|q| q.len()).sum();
357        total <= 1
358    }
359
360    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
361        if let Some(to_release) = self.to_release.take() {
362            if !to_release.is_empty()
363                && let Some(log_writer) = log_writer
364            {
365                let HookLocationMeta {
366                    location: batch_location,
367                    line,
368                    caret_indent,
369                } = self.location;
370                let note_str = format!(
371                    "^ observed non-deterministic order: {:?}",
372                    TruncatedVecDebug(
373                        RefCell::new(Some(to_release.iter())),
374                        8,
375                        self.format_item_debug
376                    )
377                );
378
379                let _ = writeln!(log_writer);
380                log_release(
381                    log_writer,
382                    batch_location,
383                    line,
384                    caret_indent,
385                    &note_str,
386                    colored::Color::Green,
387                );
388            }
389
390            for item in to_release {
391                self.output.try_send(item).unwrap();
392            }
393        } else {
394            panic!("No decision to release");
395        }
396    }
397
398    fn location_meta(&self) -> HookLocationMeta {
399        self.location
400    }
401}
402
403impl<K: Hash + Eq + Clone, V> ObservationHook for TopLevelKeyedStreamOrderHook<K, V> {
404    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
405        let mut current_input = self.input.borrow_mut();
406
407        // Collect non-empty keys with their queue lengths
408        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
409        let nonempty_keys: Vec<(K, usize)> = current_input
410            .iter()
411            .filter(|(_, q)| !q.is_empty())
412            .map(|(k, q)| (k.clone(), q.len()))
413            .collect();
414
415        // Pick which key to release from
416        let key_idx = (0..nonempty_keys.len()).generate(driver).unwrap();
417        let (key, queue_len) = &nonempty_keys[key_idx];
418
419        // Pick which item from that key's queue
420        let item_idx = (0..*queue_len).generate(driver).unwrap();
421        let item = current_input
422            .get_mut(key)
423            .unwrap()
424            .remove(item_idx)
425            .unwrap();
426
427        self.to_release = Some(vec![(key.clone(), item)]);
428    }
429}
430
431/// Top-level variant of [`PartiallyOrderedStreamHook`](super::PartiallyOrderedStreamHook). Same one-at-a-time release
432/// strategy as [`TopLevelKeyedStreamOrderHook`], but always takes from the FRONT
433/// of the chosen key's queue to preserve within-key order.
434pub struct TopLevelPartiallyOrderedStreamHook<K: Hash + Eq + Clone, V> {
435    pub input: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
436    pub to_release: Option<Vec<(K, V)>>,
437    pub output: Sender<(K, V)>,
438    pub location: HookLocationMeta,
439    pub format_item_debug: fn(&(K, V)) -> Option<String>,
440}
441
442impl<K: Hash + Eq + Clone, V> RuntimeHook for TopLevelPartiallyOrderedStreamHook<K, V> {
443    fn has_pending_input(&self) -> bool {
444        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
445        !self.input.borrow().values().all(|q| q.is_empty())
446    }
447
448    fn only_one_possible_decision(&self) -> bool {
449        // Within a key the front element is forced, so the only choice is which key
450        // releases next: a single non-empty key is fully forced regardless of depth.
451        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
452        let nonempty_keys = self
453            .input
454            .borrow()
455            .values()
456            .filter(|q| !q.is_empty())
457            .count();
458        nonempty_keys <= 1
459    }
460
461    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
462        if let Some(to_release) = self.to_release.take() {
463            if !to_release.is_empty()
464                && let Some(log_writer) = log_writer
465            {
466                let HookLocationMeta {
467                    location: batch_location,
468                    line,
469                    caret_indent,
470                } = self.location;
471                let note_str = format!(
472                    "^ observed partially-ordered interleaving: {:?}",
473                    TruncatedVecDebug(
474                        RefCell::new(Some(to_release.iter())),
475                        8,
476                        self.format_item_debug
477                    )
478                );
479
480                let _ = writeln!(log_writer);
481                log_release(
482                    log_writer,
483                    batch_location,
484                    line,
485                    caret_indent,
486                    &note_str,
487                    colored::Color::Green,
488                );
489            }
490
491            for item in to_release {
492                self.output.try_send(item).unwrap();
493            }
494        } else {
495            panic!("No decision to release");
496        }
497    }
498
499    fn location_meta(&self) -> HookLocationMeta {
500        self.location
501    }
502}
503
504impl<K: Hash + Eq + Clone, V> ObservationHook for TopLevelPartiallyOrderedStreamHook<K, V> {
505    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
506        let mut current_input = self.input.borrow_mut();
507
508        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
509        let nonempty_keys: Vec<K> = current_input
510            .iter()
511            .filter(|(_, q)| !q.is_empty())
512            .map(|(k, _)| k.clone())
513            .collect();
514
515        // Pick which key to release from
516        let key_idx = (0..nonempty_keys.len()).generate(driver).unwrap();
517        let key = &nonempty_keys[key_idx];
518
519        // Always take from the front to preserve within-key order
520        let item = current_input.get_mut(key).unwrap().pop_front().unwrap();
521
522        self.to_release = Some(vec![(key.clone(), item)]);
523    }
524}
525
526/// Top-level merge-ordered hook. Releases one element at a time, picking from
527/// the front of either the first or second input queue. This preserves per-input
528/// order while allowing feedback cycles to deliver elements between releases.
529pub struct TopLevelMergeOrderedHook<T> {
530    pub first: Rc<RefCell<VecDeque<T>>>,
531    pub second: Rc<RefCell<VecDeque<T>>>,
532    pub to_release: Option<Vec<T>>,
533    pub release_source: Option<&'static str>,
534    pub output: Sender<T>,
535    pub location: HookLocationMeta,
536    pub format_item_debug: fn(&T) -> Option<String>,
537}
538
539impl<T> RuntimeHook for TopLevelMergeOrderedHook<T> {
540    fn has_pending_input(&self) -> bool {
541        !self.first.borrow().is_empty() || !self.second.borrow().is_empty()
542    }
543
544    fn only_one_possible_decision(&self) -> bool {
545        // Each side's front element is forced, so the only choice is which side
546        // releases next: it exists only when both sides are non-empty.
547        self.first.borrow().is_empty() || self.second.borrow().is_empty()
548    }
549
550    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
551        if let Some(to_release) = self.to_release.take() {
552            let source = self.release_source.take();
553            if !to_release.is_empty()
554                && let Some(log_writer) = log_writer
555            {
556                let HookLocationMeta {
557                    location: batch_location,
558                    line,
559                    caret_indent,
560                } = self.location;
561                let source_label = source.unwrap_or("?");
562
563                let labeled_iter = to_release.iter().map(|item| (source_label, item));
564
565                let note_str = format!(
566                    "^ observed non-deterministic merge order: {:?}",
567                    TruncatedLabeledVecDebug(
568                        RefCell::new(Some(labeled_iter)),
569                        8,
570                        self.format_item_debug
571                    )
572                );
573
574                let _ = writeln!(log_writer);
575                log_release(
576                    log_writer,
577                    batch_location,
578                    line,
579                    caret_indent,
580                    &note_str,
581                    colored::Color::Green,
582                );
583            }
584
585            for item in to_release {
586                self.output.try_send(item).unwrap();
587            }
588        } else {
589            panic!("No decision to release");
590        }
591    }
592
593    fn location_meta(&self) -> HookLocationMeta {
594        self.location
595    }
596}
597
598impl<T> ObservationHook for TopLevelMergeOrderedHook<T> {
599    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
600        let first_empty = self.first.borrow().is_empty();
601        let second_empty = self.second.borrow().is_empty();
602
603        let (item, source) = if first_empty {
604            (self.second.borrow_mut().pop_front().unwrap(), "r")
605        } else if second_empty {
606            (self.first.borrow_mut().pop_front().unwrap(), "l")
607        } else {
608            let take_second: bool = produce().generate(driver).unwrap();
609            if take_second {
610                (self.second.borrow_mut().pop_front().unwrap(), "r")
611            } else {
612                (self.first.borrow_mut().pop_front().unwrap(), "l")
613            }
614        };
615
616        self.to_release = Some(vec![item]);
617        self.release_source = Some(source);
618    }
619}
620
621/// Keyed variant of [`TopLevelMergeOrderedHook`]. Releases one element at a
622/// time, picking from the front of some key's queue in either the first or the
623/// second input. This preserves per-input order within each key while allowing
624/// arbitrary interleaving both across the two inputs and across keys (which is
625/// unconstrained for keyed streams), and lets feedback cycles deliver elements
626/// between releases.
627pub struct TopLevelKeyedMergeOrderedHook<K: Hash + Eq + Clone, V> {
628    pub first: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
629    pub second: Rc<RefCell<FxHashMap<K, VecDeque<V>>>>,
630    pub to_release: Option<Vec<(K, V)>>,
631    pub release_source: Option<&'static str>,
632    pub output: Sender<(K, V)>,
633    pub location: HookLocationMeta,
634    pub format_item_debug: fn(&(K, V)) -> Option<String>,
635}
636
637impl<K: Hash + Eq + Clone, V> RuntimeHook for TopLevelKeyedMergeOrderedHook<K, V> {
638    fn has_pending_input(&self) -> bool {
639        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
640        let first_nonempty = !self.first.borrow().values().all(|q| q.is_empty());
641        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
642        let second_nonempty = !self.second.borrow().values().all(|q| q.is_empty());
643        first_nonempty || second_nonempty
644    }
645
646    fn only_one_possible_decision(&self) -> bool {
647        // Each (side, key) queue's front element is forced, so the only choice is which
648        // queue releases next.
649        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
650        let first_count = self
651            .first
652            .borrow()
653            .values()
654            .filter(|q| !q.is_empty())
655            .count();
656        #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
657        let second_count = self
658            .second
659            .borrow()
660            .values()
661            .filter(|q| !q.is_empty())
662            .count();
663        first_count + second_count <= 1
664    }
665
666    fn release_decision(&mut self, log_writer: Option<&mut dyn std::fmt::Write>) {
667        if let Some(to_release) = self.to_release.take() {
668            let source = self.release_source.take();
669            if !to_release.is_empty()
670                && let Some(log_writer) = log_writer
671            {
672                let HookLocationMeta {
673                    location: batch_location,
674                    line,
675                    caret_indent,
676                } = self.location;
677                let source_label = source.unwrap_or("?");
678
679                let labeled_iter = to_release.iter().map(|item| (source_label, item));
680
681                let note_str = format!(
682                    "^ observed non-deterministic merge order: {:?}",
683                    TruncatedLabeledVecDebug(
684                        RefCell::new(Some(labeled_iter)),
685                        8,
686                        self.format_item_debug
687                    )
688                );
689
690                let _ = writeln!(log_writer);
691                log_release(
692                    log_writer,
693                    batch_location,
694                    line,
695                    caret_indent,
696                    &note_str,
697                    colored::Color::Green,
698                );
699            }
700
701            for item in to_release {
702                self.output.try_send(item).unwrap();
703            }
704        } else {
705            panic!("No decision to release");
706        }
707    }
708
709    fn location_meta(&self) -> HookLocationMeta {
710        self.location
711    }
712}
713
714impl<K: Hash + Eq + Clone, V> ObservationHook for TopLevelKeyedMergeOrderedHook<K, V> {
715    fn autonomous_decision<'a>(&mut self, driver: &mut Borrowed<'a>) {
716        // Collect candidates: for each non-empty key queue in either input, we
717        // can release its front element. `false` = first input, `true` = second.
718        let mut candidates: Vec<(bool, K)> = Vec::new();
719        {
720            #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
721            for (k, q) in self.first.borrow().iter() {
722                if !q.is_empty() {
723                    candidates.push((false, k.clone()));
724                }
725            }
726            #[expect(clippy::disallowed_methods, reason = "FxHasher is deterministic")]
727            for (k, q) in self.second.borrow().iter() {
728                if !q.is_empty() {
729                    candidates.push((true, k.clone()));
730                }
731            }
732        }
733
734        let idx = (0..candidates.len()).generate(driver).unwrap();
735        let (take_second, key) = &candidates[idx];
736        let take_second = *take_second;
737
738        let item = if take_second {
739            self.second
740                .borrow_mut()
741                .get_mut(key)
742                .unwrap()
743                .pop_front()
744                .unwrap()
745        } else {
746            self.first
747                .borrow_mut()
748                .get_mut(key)
749                .unwrap()
750                .pop_front()
751                .unwrap()
752        };
753
754        self.to_release = Some(vec![(key.clone(), item)]);
755        self.release_source = Some(if take_second { "r" } else { "l" });
756    }
757}