1use 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
23pub 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 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 ¬e_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 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#[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 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
183pub 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 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 ¬e_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 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 *current_input = remaining;
268
269 {
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
284impl<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 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
337pub 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 #[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 ¬e_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 #[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 let key_idx = (0..nonempty_keys.len()).generate(driver).unwrap();
418 let (key, queue_len) = &nonempty_keys[key_idx];
419
420 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
432impl<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 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
486pub 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 #[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 ¬e_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 let key_idx = (0..nonempty_keys.len()).generate(driver).unwrap();
572 let key = &nonempty_keys[key_idx];
573
574 let item = current_input.get_mut(key).unwrap().pop_front().unwrap();
576
577 self.to_release = Some(vec![(key.clone(), item)]);
578 }
579}
580
581impl<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 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 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
639pub 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 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 ¬e_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#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
741pub enum MergeDecision<T, S = ()> {
742 First(T),
744 Second(T),
746 FirstNext(S),
748 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#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
770pub struct MergeStatus {
771 pub first_buffered: usize,
773 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 (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 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
850pub 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 #[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 ¬e_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 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 (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 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}