1use std::any::Any;
8use std::cell::{Cell, RefCell};
9use std::collections::BTreeMap;
10use std::fmt;
11use std::rc::Rc;
12use std::sync::atomic::{AtomicU64, Ordering};
13
14use serde::{Deserialize, Serialize};
15use serde_json::{json, Number, Value};
16
17use crate::checkpoint::register_backend_state_contributor;
18use crate::graph::{Graph, GraphNodeOpts, RestoreFactoryMeta, TopologyGroupOptions};
19use crate::node::{Node, NodeOpts};
20use crate::operators::{init_node, Operator};
21use crate::protocol::{LockId, Message};
22
23type Disposer = Box<dyn FnOnce()>;
24type DisposerSlots = Rc<RefCell<Vec<Option<Disposer>>>>;
25type ViewDisposeAction = Rc<dyn Fn()>;
26type ViewDisposeHook = Box<dyn FnOnce()>;
27type PageView<T> = ReactiveView<LogChange<T>, Vec<T>>;
28type PageMemo<T> = Rc<RefCell<Vec<(String, PageView<T>)>>>;
29type IndexRangeView<K, S, V> = ReactiveView<IndexChange<K, S, V>, Vec<V>>;
30type IndexRangeMemo<K, S, V> = Rc<RefCell<Vec<((K, K), IndexRangeView<K, S, V>)>>>;
31type MapSelectView<K, V> = ReactiveView<MapChange<K, V>, BTreeMap<K, V>>;
32type MapSelectPredicate<K, V> = Rc<dyn Fn(&V, &K) -> bool>;
33type MapSelectMemo<K, V> = Rc<RefCell<Vec<(MapSelectPredicate<K, V>, MapSelectView<K, V>)>>>;
34
35static VIEW_PULL_SEQ: AtomicU64 = AtomicU64::new(0);
36
37#[derive(Clone, Copy)]
38struct ViewFactories {
39 group: &'static str,
40 delta: &'static str,
41 snapshot: &'static str,
42}
43
44#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
45#[serde(tag = "kind")]
46pub enum ListChange<T> {
48 #[serde(rename = "append")]
49 Append {
51 value: T,
53 },
54 #[serde(rename = "appendMany")]
55 AppendMany {
57 values: Vec<T>,
59 },
60 #[serde(rename = "insert")]
61 Insert {
63 index: usize,
65 value: T,
67 },
68 #[serde(rename = "insertMany")]
69 InsertMany {
71 index: usize,
73 values: Vec<T>,
75 },
76 #[serde(rename = "pop")]
77 Pop {
79 index: usize,
81 value: T,
83 },
84 #[serde(rename = "trimHead")]
85 TrimHead {
87 n: usize,
89 },
90 #[serde(rename = "clear")]
91 Clear {
93 count: usize,
95 },
96}
97
98#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
99#[serde(tag = "kind")]
100pub enum LogChange<T> {
102 #[serde(rename = "append")]
103 Append {
105 value: T,
107 },
108 #[serde(rename = "appendMany")]
109 AppendMany {
111 values: Vec<T>,
113 },
114 #[serde(rename = "trimHead")]
115 TrimHead {
117 n: usize,
119 },
120 #[serde(rename = "clear")]
121 Clear {
123 count: usize,
125 },
126}
127
128#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
129#[serde(tag = "kind")]
130pub enum IndexChange<K, S, V> {
132 #[serde(rename = "upsert")]
133 Upsert {
135 primary: K,
137 secondary: S,
139 value: V,
141 },
142 #[serde(rename = "delete")]
143 Delete {
145 primary: K,
147 },
148 #[serde(rename = "deleteMany")]
149 DeleteMany {
151 primaries: Vec<K>,
153 },
154 #[serde(rename = "clear")]
155 Clear {
157 count: usize,
159 },
160}
161
162#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
163#[serde(tag = "kind")]
164pub enum MapChange<K, V> {
166 #[serde(rename = "set")]
167 Set {
169 key: K,
171 value: V,
173 },
174 #[serde(rename = "delete")]
175 Delete {
177 key: K,
179 previous: V,
181 },
182 #[serde(rename = "clear")]
183 Clear {
185 count: usize,
187 },
188}
189
190#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
191pub struct IndexRow<K, S, V> {
193 pub primary: K,
195 pub secondary: S,
197 pub value: V,
199}
200
201#[derive(Clone, Default)]
202pub struct ReactiveListOptions {
204 pub name: Option<String>,
206 pub graph: Option<Graph>,
208 pub max_size: Option<usize>,
210}
211
212impl fmt::Debug for ReactiveListOptions {
213 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
214 f.debug_struct("ReactiveListOptions")
215 .field("name", &self.name)
216 .field("graph", &self.graph.as_ref().map(|g| g.name()))
217 .field("max_size", &self.max_size)
218 .finish()
219 }
220}
221
222impl ReactiveListOptions {
223 pub fn named(name: impl Into<String>) -> Self {
225 Self {
226 name: Some(name.into()),
227 ..Self::default()
228 }
229 }
230
231 pub fn graph(mut self, graph: Graph) -> Self {
233 self.graph = Some(graph);
234 self
235 }
236
237 pub fn max_size(mut self, max_size: usize) -> Self {
239 self.max_size = Some(max_size);
240 self
241 }
242}
243
244#[derive(Clone, Default)]
245pub struct ReactiveLogOptions {
247 pub name: Option<String>,
249 pub graph: Option<Graph>,
251 pub max_size: Option<usize>,
253}
254
255#[derive(Clone, Default)]
256pub struct ReactiveIndexOptions {
258 pub name: Option<String>,
260 pub graph: Option<Graph>,
262}
263
264#[derive(Clone, Default)]
265pub struct ReactiveMapOptions {
267 pub name: Option<String>,
269 pub graph: Option<Graph>,
271}
272
273impl fmt::Debug for ReactiveIndexOptions {
274 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
275 f.debug_struct("ReactiveIndexOptions")
276 .field("name", &self.name)
277 .field("graph", &self.graph.as_ref().map(|g| g.name()))
278 .finish()
279 }
280}
281
282impl fmt::Debug for ReactiveMapOptions {
283 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
284 f.debug_struct("ReactiveMapOptions")
285 .field("name", &self.name)
286 .field("graph", &self.graph.as_ref().map(|g| g.name()))
287 .finish()
288 }
289}
290
291impl ReactiveIndexOptions {
292 pub fn named(name: impl Into<String>) -> Self {
294 Self {
295 name: Some(name.into()),
296 ..Self::default()
297 }
298 }
299
300 pub fn graph(mut self, graph: Graph) -> Self {
302 self.graph = Some(graph);
303 self
304 }
305}
306
307impl ReactiveMapOptions {
308 pub fn named(name: impl Into<String>) -> Self {
310 Self {
311 name: Some(name.into()),
312 ..Self::default()
313 }
314 }
315
316 pub fn graph(mut self, graph: Graph) -> Self {
318 self.graph = Some(graph);
319 self
320 }
321}
322
323impl fmt::Debug for ReactiveLogOptions {
324 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
325 f.debug_struct("ReactiveLogOptions")
326 .field("name", &self.name)
327 .field("graph", &self.graph.as_ref().map(|g| g.name()))
328 .field("max_size", &self.max_size)
329 .finish()
330 }
331}
332
333impl ReactiveLogOptions {
334 pub fn named(name: impl Into<String>) -> Self {
336 Self {
337 name: Some(name.into()),
338 ..Self::default()
339 }
340 }
341
342 pub fn graph(mut self, graph: Graph) -> Self {
344 self.graph = Some(graph);
345 self
346 }
347
348 pub fn max_size(mut self, max_size: usize) -> Self {
350 self.max_size = Some(max_size);
351 self
352 }
353}
354
355#[derive(Clone)]
356struct ListBackend<T> {
357 buf: Rc<RefCell<Vec<T>>>,
358 version: Rc<Cell<u64>>,
359}
360
361impl<T: Clone> ListBackend<T> {
362 fn new(initial: Vec<T>, max_size: Option<usize>) -> Self {
363 let mut initial = initial;
364 trim_head_overflow(&mut initial, max_size);
365 Self {
366 buf: Rc::new(RefCell::new(initial)),
367 version: Rc::new(Cell::new(0)),
368 }
369 }
370
371 fn version(&self) -> u64 {
372 self.version.get()
373 }
374
375 fn bump(&self) {
376 self.version.set(self.version.get().wrapping_add(1));
377 }
378
379 fn snapshot(&self) -> Vec<T> {
380 self.buf.borrow().clone()
381 }
382
383 fn instance_token(&self) -> usize {
384 Rc::as_ptr(&self.buf) as usize
385 }
386
387 fn len(&self) -> usize {
388 self.buf.borrow().len()
389 }
390
391 fn at(&self, index: isize) -> Option<T> {
392 let buf = self.buf.borrow();
393 let i = normalize_read_index(index, buf.len())?;
394 buf.get(i).cloned()
395 }
396
397 fn append(&self, value: T) {
398 self.buf.borrow_mut().push(value);
399 self.bump();
400 }
401
402 fn append_many(&self, values: &[T]) {
403 if values.is_empty() {
404 return;
405 }
406 self.buf.borrow_mut().extend_from_slice(values);
407 self.bump();
408 }
409
410 fn insert(&self, index: usize, value: T) {
411 let mut buf = self.buf.borrow_mut();
412 assert!(
413 index <= buf.len(),
414 "insert: index {index} out of range [0, {}]",
415 buf.len()
416 );
417 buf.insert(index, value);
418 self.bump();
419 }
420
421 fn insert_many(&self, index: usize, values: &[T]) {
422 let mut buf = self.buf.borrow_mut();
423 assert!(
424 index <= buf.len(),
425 "insert_many: index {index} out of range [0, {}]",
426 buf.len()
427 );
428 if values.is_empty() {
429 return;
430 }
431 buf.splice(index..index, values.iter().cloned());
432 self.bump();
433 }
434
435 fn pop(&self, index: Option<isize>) -> (usize, T) {
436 let mut buf = self.buf.borrow_mut();
437 assert!(!buf.is_empty(), "pop from empty list");
438 let raw = index.unwrap_or(-1);
439 let i = normalize_read_index(raw, buf.len())
440 .unwrap_or_else(|| panic!("pop: index {raw} out of range"));
441 let value = buf.remove(i);
442 self.bump();
443 (i, value)
444 }
445
446 fn clear(&self) -> usize {
447 let mut buf = self.buf.borrow_mut();
448 let n = buf.len();
449 if n == 0 {
450 return 0;
451 }
452 buf.clear();
453 self.bump();
454 n
455 }
456
457 fn enforce_max_size(&self, max_size: Option<usize>) -> usize {
458 let mut buf = self.buf.borrow_mut();
459 let removed = trim_head_overflow(&mut buf, max_size);
460 if removed > 0 {
461 self.bump();
462 }
463 removed
464 }
465}
466
467#[derive(Clone)]
468struct LogBackend<T> {
469 buf: Rc<RefCell<Vec<T>>>,
470 version: Rc<Cell<u64>>,
471 max_size: Option<usize>,
472}
473
474#[derive(Clone)]
475struct IndexBackend<K, S, V> {
476 rows: Rc<RefCell<BTreeMap<K, (S, V)>>>,
477 version: Rc<Cell<u64>>,
478}
479
480#[derive(Clone)]
481struct MapBackend<K, V> {
482 rows: Rc<RefCell<BTreeMap<K, V>>>,
483 version: Rc<Cell<u64>>,
484}
485
486impl<K: Clone + Ord, S: Clone + Ord, V: Clone> IndexBackend<K, S, V> {
487 fn new(initial: Vec<IndexRow<K, S, V>>) -> Self {
488 let rows = initial
489 .into_iter()
490 .map(|row| (row.primary, (row.secondary, row.value)))
491 .collect();
492 Self {
493 rows: Rc::new(RefCell::new(rows)),
494 version: Rc::new(Cell::new(0)),
495 }
496 }
497
498 fn version(&self) -> u64 {
499 self.version.get()
500 }
501
502 fn bump(&self) {
503 self.version.set(self.version.get().wrapping_add(1));
504 }
505
506 fn instance_token(&self) -> usize {
507 Rc::as_ptr(&self.rows) as usize
508 }
509
510 fn len(&self) -> usize {
511 self.rows.borrow().len()
512 }
513
514 fn has(&self, primary: &K) -> bool {
515 self.rows.borrow().contains_key(primary)
516 }
517
518 fn get(&self, primary: &K) -> Option<V> {
519 self.rows
520 .borrow()
521 .get(primary)
522 .map(|(_, value)| value.clone())
523 }
524
525 fn snapshot(&self) -> Vec<IndexRow<K, S, V>> {
526 let mut rows = self
527 .rows
528 .borrow()
529 .iter()
530 .map(|(primary, (secondary, value))| IndexRow {
531 primary: primary.clone(),
532 secondary: secondary.clone(),
533 value: value.clone(),
534 })
535 .collect::<Vec<_>>();
536 rows.sort_by(|a, b| {
537 a.secondary
538 .cmp(&b.secondary)
539 .then_with(|| a.primary.cmp(&b.primary))
540 });
541 rows
542 }
543
544 fn range_by_primary(&self, start: &K, end: &K) -> Vec<V> {
545 if start >= end {
546 return Vec::new();
547 }
548 self.rows
549 .borrow()
550 .range(start.clone()..end.clone())
551 .map(|(_, (_, value))| value.clone())
552 .collect()
553 }
554
555 fn upsert(&self, primary: K, secondary: S, value: V) {
556 self.rows.borrow_mut().insert(primary, (secondary, value));
557 self.bump();
558 }
559
560 fn delete(&self, primary: &K) -> bool {
561 let removed = self.rows.borrow_mut().remove(primary).is_some();
562 if removed {
563 self.bump();
564 }
565 removed
566 }
567
568 fn delete_many(&self, primaries: &[K]) -> Vec<K> {
569 let mut rows = self.rows.borrow_mut();
570 let mut removed = Vec::new();
571 for primary in primaries {
572 if rows.remove(primary).is_some() {
573 removed.push(primary.clone());
574 }
575 }
576 if !removed.is_empty() {
577 drop(rows);
578 self.bump();
579 }
580 removed
581 }
582
583 fn clear(&self) -> usize {
584 let mut rows = self.rows.borrow_mut();
585 let count = rows.len();
586 if count == 0 {
587 return 0;
588 }
589 rows.clear();
590 self.bump();
591 count
592 }
593}
594
595impl<K: Clone + Ord, V: Clone> MapBackend<K, V> {
596 fn new(initial: Vec<(K, V)>) -> Self {
597 Self {
598 rows: Rc::new(RefCell::new(initial.into_iter().collect())),
599 version: Rc::new(Cell::new(0)),
600 }
601 }
602
603 fn version(&self) -> u64 {
604 self.version.get()
605 }
606
607 fn bump(&self) {
608 self.version.set(self.version.get().wrapping_add(1));
609 }
610
611 fn instance_token(&self) -> usize {
612 Rc::as_ptr(&self.rows) as usize
613 }
614
615 fn len(&self) -> usize {
616 self.rows.borrow().len()
617 }
618
619 fn has(&self, key: &K) -> bool {
620 self.rows.borrow().contains_key(key)
621 }
622
623 fn get(&self, key: &K) -> Option<V> {
624 self.rows.borrow().get(key).cloned()
625 }
626
627 fn snapshot(&self) -> BTreeMap<K, V> {
628 self.rows.borrow().clone()
629 }
630
631 fn entries_snapshot(&self) -> Vec<(K, V)> {
632 self.rows
633 .borrow()
634 .iter()
635 .map(|(key, value)| (key.clone(), value.clone()))
636 .collect()
637 }
638
639 fn selected_snapshot(&self, predicate: &MapSelectPredicate<K, V>) -> BTreeMap<K, V> {
640 self.rows
641 .borrow()
642 .iter()
643 .filter(|(key, value)| predicate(value, key))
644 .map(|(key, value)| (key.clone(), value.clone()))
645 .collect()
646 }
647
648 fn set(&self, key: K, value: V) {
649 self.rows.borrow_mut().insert(key, value);
650 self.bump();
651 }
652
653 fn delete(&self, key: &K) -> Option<V> {
654 let previous = self.rows.borrow_mut().remove(key);
655 if previous.is_some() {
656 self.bump();
657 }
658 previous
659 }
660
661 fn clear(&self) -> usize {
662 let mut rows = self.rows.borrow_mut();
663 let count = rows.len();
664 if count == 0 {
665 return 0;
666 }
667 rows.clear();
668 self.bump();
669 count
670 }
671}
672
673impl<T: Clone> LogBackend<T> {
674 fn new(initial: Vec<T>, max_size: Option<usize>) -> Self {
675 let mut initial = initial;
676 trim_head_overflow(&mut initial, max_size);
677 Self {
678 buf: Rc::new(RefCell::new(initial)),
679 version: Rc::new(Cell::new(0)),
680 max_size,
681 }
682 }
683
684 fn version(&self) -> u64 {
685 self.version.get()
686 }
687
688 fn bump(&self) {
689 self.version.set(self.version.get().wrapping_add(1));
690 }
691
692 fn snapshot(&self) -> Vec<T> {
693 self.buf.borrow().clone()
694 }
695
696 fn instance_token(&self) -> usize {
697 Rc::as_ptr(&self.buf) as usize
698 }
699
700 fn len(&self) -> usize {
701 self.buf.borrow().len()
702 }
703
704 fn at(&self, index: isize) -> Option<T> {
705 let buf = self.buf.borrow();
706 let i = normalize_read_index(index, buf.len())?;
707 buf.get(i).cloned()
708 }
709
710 fn append(&self, value: T) -> usize {
711 self.buf.borrow_mut().push(value);
712 self.bump();
713 self.enforce_max_size()
714 }
715
716 fn append_many(&self, values: &[T]) -> usize {
717 if values.is_empty() {
718 return 0;
719 }
720 self.buf.borrow_mut().extend_from_slice(values);
721 self.bump();
722 self.enforce_max_size()
723 }
724
725 fn clear(&self) -> usize {
726 let mut buf = self.buf.borrow_mut();
727 let n = buf.len();
728 if n == 0 {
729 return 0;
730 }
731 buf.clear();
732 self.bump();
733 n
734 }
735
736 fn trim_head(&self, n: usize) -> usize {
737 if n == 0 {
738 return 0;
739 }
740 let mut buf = self.buf.borrow_mut();
741 let removed = n.min(buf.len());
742 if removed == 0 {
743 return 0;
744 }
745 buf.drain(0..removed);
746 self.bump();
747 removed
748 }
749
750 fn enforce_max_size(&self) -> usize {
751 let mut buf = self.buf.borrow_mut();
752 let removed = trim_head_overflow(&mut buf, self.max_size);
753 if removed > 0 {
754 self.bump();
755 }
756 removed
757 }
758}
759
760#[derive(Clone)]
761pub struct ReactiveList<T> {
763 pub delta: Node<ListChange<T>>,
765 pub snapshot: Node<Vec<T>>,
767 pub pull_id: LockId,
769 backend: ListBackend<T>,
770 max_size: Option<usize>,
771 graph: Option<Graph>,
772 id_prefix: String,
773 bind_seq: Rc<Cell<usize>>,
774 disposers: DisposerSlots,
775}
776
777impl<T: Clone + 'static> ReactiveList<T> {
778 pub fn new(initial: Vec<T>, options: ReactiveListOptions) -> Self {
780 if matches!(options.max_size, Some(0)) {
781 panic!("reactive_list: max_size must be a positive integer");
782 }
783 let backend = ListBackend::new(initial, options.max_size);
784 let token = backend.instance_token();
785 let id_prefix = options
786 .name
787 .clone()
788 .unwrap_or_else(|| format!("reactiveList@{token:x}"));
789 let pull_id = LockId::new(format!(
790 "{id_prefix}.snapshot@{:x}",
791 backend.instance_token()
792 ));
793 let restorable =
794 options.graph.is_some() && options.name.is_some() && options.max_size.is_none();
795 let delta = make_delta_node::<T>(options.graph.as_ref(), &id_prefix, restorable);
796 if restorable {
797 let backend_for_checkpoint = backend.clone();
798 register_backend_state_json(&delta, move || {
799 vec_to_checkpoint_json(backend_for_checkpoint.snapshot())
800 });
801 }
802 let snapshot = make_snapshot_node(
803 options.graph.as_ref(),
804 &id_prefix,
805 &delta,
806 &backend,
807 pull_id.clone(),
808 restorable,
809 );
810 Self {
811 delta,
812 snapshot,
813 pull_id,
814 backend,
815 max_size: options.max_size,
816 graph: options.graph,
817 id_prefix,
818 bind_seq: Rc::new(Cell::new(0)),
819 disposers: Rc::new(RefCell::new(Vec::new())),
820 }
821 }
822
823 pub fn len(&self) -> usize {
825 self.backend.len()
826 }
827
828 pub fn is_empty(&self) -> bool {
830 self.len() == 0
831 }
832
833 pub fn at(&self, index: isize) -> Option<T> {
835 self.backend.at(index)
836 }
837
838 pub fn to_vec(&self) -> Vec<T> {
840 self.backend.snapshot()
841 }
842
843 pub fn append(&self, value: T) {
845 self.backend.append(value.clone());
846 self.emit(ListChange::Append { value });
847 self.enforce_capacity();
848 }
849
850 pub fn append_many(&self, values: Vec<T>) {
852 if values.is_empty() {
853 return;
854 }
855 self.backend.append_many(&values);
856 self.emit(ListChange::AppendMany { values });
857 self.enforce_capacity();
858 }
859
860 pub fn insert(&self, index: usize, value: T) {
862 self.backend.insert(index, value.clone());
863 self.emit(ListChange::Insert { index, value });
864 self.enforce_capacity();
865 }
866
867 pub fn insert_many(&self, index: usize, values: Vec<T>) {
869 self.backend.insert_many(index, &values);
870 if values.is_empty() {
871 return;
872 }
873 self.emit(ListChange::InsertMany { index, values });
874 self.enforce_capacity();
875 }
876
877 pub fn pop(&self, index: Option<isize>) -> T {
879 let (index, value) = self.backend.pop(index);
880 self.emit(ListChange::Pop {
881 index,
882 value: value.clone(),
883 });
884 value
885 }
886
887 pub fn clear(&self) {
889 let count = self.backend.clear();
890 if count > 0 {
891 self.emit(ListChange::Clear { count });
892 }
893 }
894
895 pub fn append_from(&self, src: &Node<T>) -> Disposer {
897 let graph = self.graph.as_ref().unwrap_or_else(|| {
898 panic!("reactive_list.append_from requires options.graph so the input fold is describe-visible (D61)")
899 });
900 let bind_idx = self.bind_seq.get();
901 self.bind_seq.set(bind_idx + 1);
902 let backend = self.backend.clone();
903 let delta = self.delta.clone();
904 let max_size = self.max_size;
905 let op = Operator::with_opts(
906 "reactiveList.bindSource",
907 NodeOpts {
908 partial: true,
909 ..NodeOpts::default()
910 },
911 move |ctx| {
912 for value in ctx.batch::<T>(0) {
913 let value = (*value).clone();
914 backend.append(value.clone());
915 delta.down(vec![Message::Data(Rc::new(ListChange::Append { value }))]);
916 emit_trimmed(&backend, &delta, max_size);
917 }
918 },
919 );
920 let mut opts = GraphNodeOpts::named(format!("{}.bind#{bind_idx}", self.id_prefix));
921 opts.meta
922 .insert("kind".to_owned(), "collection_bind_source".to_owned());
923 opts.meta
924 .insert("collection".to_owned(), "reactiveList".to_owned());
925 let src_core = src.erased();
926 let folder = graph.init_node::<ListChange<T>>(op, vec![src_core.clone()], opts);
927 let unsub = folder.subscribe(|_| {});
928 let slot = {
929 let mut disposers = self.disposers.borrow_mut();
930 let slot = disposers.len();
931 disposers.push(Some(unsub));
932 slot
933 };
934 let disposers = self.disposers.clone();
935 let folder_for_dispose = folder.clone();
936 Box::new(move || {
937 folder_for_dispose.unsubscribe_dep(src_core, |_| {});
938 if let Some(disposer) = disposers.borrow_mut().get_mut(slot).and_then(Option::take) {
939 disposer();
940 }
941 })
942 }
943
944 pub fn dispose(&self) {
946 for disposer in self.disposers.borrow_mut().iter_mut() {
947 if let Some(disposer) = disposer.take() {
948 disposer();
949 }
950 }
951 }
952
953 fn emit(&self, change: ListChange<T>) {
954 self.delta.down(vec![Message::Data(Rc::new(change))]);
955 }
956
957 fn enforce_capacity(&self) {
958 emit_trimmed(&self.backend, &self.delta, self.max_size);
959 }
960}
961
962#[derive(Clone)]
963pub struct ReactiveLog<T> {
965 pub delta: Node<LogChange<T>>,
967 pub snapshot: Node<Vec<T>>,
969 pub pull_id: LockId,
971 backend: LogBackend<T>,
972 graph: Option<Graph>,
973 id_prefix: String,
974 bind_seq: Rc<Cell<usize>>,
975 page_seq: Rc<Cell<usize>>,
976 page_memo: PageMemo<T>,
977 disposers: DisposerSlots,
978}
979
980pub struct ReactiveView<C, S> {
982 pub delta: Node<C>,
984 pub snapshot: Node<S>,
986 pub pull_id: LockId,
988 disposed: Rc<Cell<bool>>,
989 dispose_action: ViewDisposeAction,
990}
991
992#[derive(Clone)]
993pub struct ReactiveIndex<K, S, V> {
995 pub delta: Node<IndexChange<K, S, V>>,
997 pub snapshot: Node<Vec<IndexRow<K, S, V>>>,
999 pub pull_id: LockId,
1001 backend: IndexBackend<K, S, V>,
1002 graph: Option<Graph>,
1003 id_prefix: String,
1004 range_seq: Rc<Cell<usize>>,
1005 range_memo: IndexRangeMemo<K, S, V>,
1006}
1007
1008#[derive(Clone)]
1009pub struct ReactiveMap<K, V> {
1011 pub delta: Node<MapChange<K, V>>,
1013 pub snapshot: Node<BTreeMap<K, V>>,
1015 pub pull_id: LockId,
1017 backend: MapBackend<K, V>,
1018 graph: Option<Graph>,
1019 id_prefix: String,
1020 select_seq: Rc<Cell<usize>>,
1021 select_memo: MapSelectMemo<K, V>,
1022}
1023
1024impl<K, S, V> ReactiveIndex<K, S, V>
1025where
1026 K: Clone + Ord + fmt::Debug + 'static,
1027 S: Clone + Ord + 'static,
1028 V: Clone + 'static,
1029{
1030 pub fn new(initial: Vec<IndexRow<K, S, V>>, options: ReactiveIndexOptions) -> Self {
1032 let backend = IndexBackend::new(initial);
1033 let token = backend.instance_token();
1034 let id_prefix = options
1035 .name
1036 .clone()
1037 .unwrap_or_else(|| format!("reactiveIndex@{token:x}"));
1038 let pull_id = LockId::new(format!(
1039 "{id_prefix}.snapshot@{:x}",
1040 backend.instance_token()
1041 ));
1042 let restorable = options.graph.is_some() && options.name.is_some();
1043 let delta =
1044 make_index_delta_node::<K, S, V>(options.graph.as_ref(), &id_prefix, restorable);
1045 if restorable {
1046 let backend_for_checkpoint = backend.clone();
1047 register_backend_state_json(&delta, move || {
1048 index_rows_to_checkpoint_json(backend_for_checkpoint.snapshot())
1049 });
1050 }
1051 let snapshot = make_index_snapshot_node(
1052 options.graph.as_ref(),
1053 &id_prefix,
1054 &delta,
1055 &backend,
1056 pull_id.clone(),
1057 restorable,
1058 );
1059 Self {
1060 delta,
1061 snapshot,
1062 pull_id,
1063 backend,
1064 graph: options.graph,
1065 id_prefix,
1066 range_seq: Rc::new(Cell::new(0)),
1067 range_memo: Rc::new(RefCell::new(Vec::new())),
1068 }
1069 }
1070
1071 pub fn len(&self) -> usize {
1073 self.backend.len()
1074 }
1075
1076 pub fn is_empty(&self) -> bool {
1078 self.len() == 0
1079 }
1080
1081 pub fn has(&self, primary: &K) -> bool {
1083 self.backend.has(primary)
1084 }
1085
1086 pub fn get(&self, primary: &K) -> Option<V> {
1088 self.backend.get(primary)
1089 }
1090
1091 pub fn to_vec(&self) -> Vec<IndexRow<K, S, V>> {
1093 self.backend.snapshot()
1094 }
1095
1096 pub fn range_by_primary(&self, start: &K, end: &K) -> Vec<V> {
1098 self.backend.range_by_primary(start, end)
1099 }
1100
1101 pub fn upsert(&self, primary: K, secondary: S, value: V) {
1103 self.backend
1104 .upsert(primary.clone(), secondary.clone(), value.clone());
1105 self.delta
1106 .down(vec![Message::Data(Rc::new(IndexChange::Upsert {
1107 primary,
1108 secondary,
1109 value,
1110 }))]);
1111 }
1112
1113 pub fn delete(&self, primary: &K) {
1115 if self.backend.delete(primary) {
1116 self.delta
1117 .down(vec![Message::Data(Rc::new(IndexChange::Delete {
1118 primary: primary.clone(),
1119 }
1120 as IndexChange<K, S, V>))]);
1121 }
1122 }
1123
1124 pub fn delete_many(&self, primaries: Vec<K>) {
1126 let removed = self.backend.delete_many(&primaries);
1127 if !removed.is_empty() {
1128 self.delta
1129 .down(vec![Message::Data(Rc::new(
1130 IndexChange::DeleteMany { primaries: removed } as IndexChange<K, S, V>,
1131 ))]);
1132 }
1133 }
1134
1135 pub fn clear(&self) {
1137 let count = self.backend.clear();
1138 if count > 0 {
1139 self.delta.down(vec![Message::Data(Rc::new(
1140 IndexChange::Clear { count } as IndexChange<K, S, V>
1141 ))]);
1142 }
1143 }
1144
1145 pub fn range(&self, start: K, end: K) -> ReactiveView<IndexChange<K, S, V>, Vec<V>> {
1147 let key = (start.clone(), end.clone());
1148 if let Some((_, view)) = self
1149 .range_memo
1150 .borrow()
1151 .iter()
1152 .find(|(existing, _)| existing == &key)
1153 {
1154 return view.clone();
1155 }
1156 let name = self.graph.as_ref().map(|_| {
1157 let next = self.range_seq.get();
1158 self.range_seq.set(next + 1);
1159 format!("{}.range#{next}", self.id_prefix)
1160 });
1161 let backend = self.backend.clone();
1162 let materialize = move || backend.range_by_primary(&start, &end);
1163 let memo = self.range_memo.clone();
1164 let dispose_key = key.clone();
1165 let view = light_reactive_view::<IndexChange<K, S, V>, Vec<V>, _>(
1166 &self.delta,
1167 self.graph.as_ref(),
1168 ViewFactories {
1169 group: "reactiveIndex.range",
1170 delta: "reactiveIndex.range.delta",
1171 snapshot: "reactiveIndex.range.snapshot",
1172 },
1173 name,
1174 materialize,
1175 Some(Box::new(move || {
1176 memo.borrow_mut()
1177 .retain(|(existing, _)| existing != &dispose_key);
1178 })),
1179 );
1180 self.range_memo.borrow_mut().push((key, view.clone()));
1181 view
1182 }
1183
1184 pub fn dispose(&self) {
1186 let views = self
1187 .range_memo
1188 .borrow()
1189 .iter()
1190 .map(|(_, view)| view.clone())
1191 .collect::<Vec<_>>();
1192 for view in views {
1193 view.dispose();
1194 }
1195 self.range_memo.borrow_mut().clear();
1196 }
1197}
1198
1199impl<K, V> ReactiveMap<K, V>
1200where
1201 K: Clone + Ord + fmt::Debug + 'static,
1202 V: Clone + 'static,
1203{
1204 pub fn new(initial: Vec<(K, V)>, options: ReactiveMapOptions) -> Self {
1206 let backend = MapBackend::new(initial);
1207 let token = backend.instance_token();
1208 let id_prefix = options
1209 .name
1210 .clone()
1211 .unwrap_or_else(|| format!("reactiveMap@{token:x}"));
1212 let pull_id = LockId::new(format!(
1213 "{id_prefix}.snapshot@{:x}",
1214 backend.instance_token()
1215 ));
1216 let restorable = options.graph.is_some() && options.name.is_some();
1217 let delta = make_map_delta_node::<K, V>(options.graph.as_ref(), &id_prefix, restorable);
1218 if restorable {
1219 let backend_for_checkpoint = backend.clone();
1220 register_backend_state_json(&delta, move || {
1221 map_entries_to_checkpoint_json(backend_for_checkpoint.entries_snapshot())
1222 });
1223 }
1224 let snapshot = make_map_snapshot_node(
1225 options.graph.as_ref(),
1226 &id_prefix,
1227 &delta,
1228 &backend,
1229 pull_id.clone(),
1230 restorable,
1231 );
1232 Self {
1233 delta,
1234 snapshot,
1235 pull_id,
1236 backend,
1237 graph: options.graph,
1238 id_prefix,
1239 select_seq: Rc::new(Cell::new(0)),
1240 select_memo: Rc::new(RefCell::new(Vec::new())),
1241 }
1242 }
1243
1244 pub fn len(&self) -> usize {
1246 self.backend.len()
1247 }
1248
1249 pub fn is_empty(&self) -> bool {
1251 self.len() == 0
1252 }
1253
1254 pub fn has(&self, key: &K) -> bool {
1256 self.backend.has(key)
1257 }
1258
1259 pub fn get(&self, key: &K) -> Option<V> {
1261 self.backend.get(key)
1262 }
1263
1264 pub fn to_map(&self) -> BTreeMap<K, V> {
1266 self.backend.snapshot()
1267 }
1268
1269 pub fn set(&self, key: K, value: V) {
1271 self.backend.set(key.clone(), value.clone());
1272 self.delta
1273 .down(vec![Message::Data(Rc::new(MapChange::Set { key, value }))]);
1274 }
1275
1276 pub fn set_many(&self, entries: Vec<(K, V)>) {
1278 for (key, value) in entries {
1279 self.set(key, value);
1280 }
1281 }
1282
1283 pub fn delete(&self, key: &K) {
1285 if let Some(previous) = self.backend.delete(key) {
1286 self.delta
1287 .down(vec![Message::Data(Rc::new(MapChange::Delete {
1288 key: key.clone(),
1289 previous,
1290 }
1291 as MapChange<K, V>))]);
1292 }
1293 }
1294
1295 pub fn delete_many(&self, keys: Vec<K>) {
1297 for key in keys {
1298 self.delete(&key);
1299 }
1300 }
1301
1302 pub fn clear(&self) {
1304 let count = self.backend.clear();
1305 if count > 0 {
1306 self.delta.down(vec![Message::Data(Rc::new(
1307 MapChange::Clear { count } as MapChange<K, V>
1308 ))]);
1309 }
1310 }
1311
1312 pub fn select<F>(&self, predicate: F) -> ReactiveView<MapChange<K, V>, BTreeMap<K, V>>
1314 where
1315 F: Fn(&V, &K) -> bool + 'static,
1316 {
1317 let predicate: MapSelectPredicate<K, V> = Rc::new(predicate);
1318 self.select_by(predicate)
1319 }
1320
1321 pub fn select_by(
1323 &self,
1324 predicate: MapSelectPredicate<K, V>,
1325 ) -> ReactiveView<MapChange<K, V>, BTreeMap<K, V>> {
1326 if let Some((_, view)) = self
1327 .select_memo
1328 .borrow()
1329 .iter()
1330 .find(|(existing, _)| Rc::ptr_eq(existing, &predicate))
1331 {
1332 return view.clone();
1333 }
1334 let name = self.graph.as_ref().map(|_| {
1335 let next = self.select_seq.get();
1336 self.select_seq.set(next + 1);
1337 format!("{}.select#{next}", self.id_prefix)
1338 });
1339 let backend = self.backend.clone();
1340 let materialize_predicate = predicate.clone();
1341 let materialize = move || backend.selected_snapshot(&materialize_predicate);
1342 let memo = self.select_memo.clone();
1343 let dispose_predicate = predicate.clone();
1344 let view = light_reactive_view::<MapChange<K, V>, BTreeMap<K, V>, _>(
1345 &self.delta,
1346 self.graph.as_ref(),
1347 ViewFactories {
1348 group: "reactiveMap.select",
1349 delta: "reactiveMap.select.delta",
1350 snapshot: "reactiveMap.select.snapshot",
1351 },
1352 name,
1353 materialize,
1354 Some(Box::new(move || {
1355 memo.borrow_mut()
1356 .retain(|(existing, _)| !Rc::ptr_eq(existing, &dispose_predicate));
1357 })),
1358 );
1359 self.select_memo
1360 .borrow_mut()
1361 .push((predicate, view.clone()));
1362 view
1363 }
1364
1365 pub fn dispose(&self) {
1367 let views = self
1368 .select_memo
1369 .borrow()
1370 .iter()
1371 .map(|(_, view)| view.clone())
1372 .collect::<Vec<_>>();
1373 for view in views {
1374 view.dispose();
1375 }
1376 self.select_memo.borrow_mut().clear();
1377 }
1378}
1379
1380impl<C, S> Clone for ReactiveView<C, S> {
1381 fn clone(&self) -> Self {
1382 Self {
1383 delta: self.delta.clone(),
1384 snapshot: self.snapshot.clone(),
1385 pull_id: self.pull_id.clone(),
1386 disposed: self.disposed.clone(),
1387 dispose_action: self.dispose_action.clone(),
1388 }
1389 }
1390}
1391
1392impl<C, S> ReactiveView<C, S> {
1393 pub fn dispose(&self) {
1395 if self.disposed.get() {
1396 return;
1397 }
1398 (self.dispose_action)();
1399 self.disposed.set(true);
1400 }
1401}
1402
1403impl<T: Clone + 'static> ReactiveLog<T> {
1404 pub fn new(initial: Vec<T>, options: ReactiveLogOptions) -> Self {
1406 if matches!(options.max_size, Some(0)) {
1407 panic!("reactive_log: max_size must be a positive integer");
1408 }
1409 let backend = LogBackend::new(initial, options.max_size);
1410 let token = backend.instance_token();
1411 let id_prefix = options
1412 .name
1413 .clone()
1414 .unwrap_or_else(|| format!("reactiveLog@{token:x}"));
1415 let pull_id = LockId::new(format!(
1416 "{id_prefix}.snapshot@{:x}",
1417 backend.instance_token()
1418 ));
1419 let restorable = options.graph.is_some() && options.name.is_some();
1420 let delta = make_log_delta_node::<T>(
1421 options.graph.as_ref(),
1422 &id_prefix,
1423 restorable,
1424 options.max_size,
1425 );
1426 if restorable {
1427 let backend_for_checkpoint = backend.clone();
1428 register_backend_state_json(&delta, move || {
1429 vec_to_checkpoint_json(backend_for_checkpoint.snapshot())
1430 });
1431 }
1432 let snapshot = make_log_snapshot_node(
1433 options.graph.as_ref(),
1434 &id_prefix,
1435 &delta,
1436 &backend,
1437 pull_id.clone(),
1438 restorable,
1439 );
1440 Self {
1441 delta,
1442 snapshot,
1443 pull_id,
1444 backend,
1445 graph: options.graph,
1446 id_prefix,
1447 bind_seq: Rc::new(Cell::new(0)),
1448 page_seq: Rc::new(Cell::new(0)),
1449 page_memo: Rc::new(RefCell::new(Vec::new())),
1450 disposers: Rc::new(RefCell::new(Vec::new())),
1451 }
1452 }
1453
1454 pub fn len(&self) -> usize {
1456 self.backend.len()
1457 }
1458
1459 pub fn is_empty(&self) -> bool {
1461 self.len() == 0
1462 }
1463
1464 pub fn at(&self, index: isize) -> Option<T> {
1466 self.backend.at(index)
1467 }
1468
1469 pub fn to_vec(&self) -> Vec<T> {
1471 self.backend.snapshot()
1472 }
1473
1474 pub fn append(&self, value: T) {
1476 let trimmed = self.backend.append(value.clone());
1477 self.emit(LogChange::Append { value });
1478 self.emit_trimmed(trimmed);
1479 }
1480
1481 pub fn append_many(&self, values: Vec<T>) {
1483 if values.is_empty() {
1484 return;
1485 }
1486 let trimmed = self.backend.append_many(&values);
1487 self.emit(LogChange::AppendMany { values });
1488 self.emit_trimmed(trimmed);
1489 }
1490
1491 pub fn clear(&self) {
1493 let count = self.backend.clear();
1494 if count > 0 {
1495 self.emit(LogChange::Clear { count });
1496 }
1497 }
1498
1499 pub fn trim_head(&self, n: usize) {
1501 let removed = self.backend.trim_head(n);
1502 self.emit_trimmed(removed);
1503 }
1504
1505 pub fn tail(&self, n: usize) -> Node<Vec<T>> {
1507 let backend = self.backend.clone();
1508 Node::derived_opts(
1509 vec![self.delta.erased()],
1510 NodeOpts {
1511 partial: true,
1512 factory: Some("reactiveLog.tail".to_owned()),
1513 ..NodeOpts::default()
1514 },
1515 move |ctx| {
1516 let all = backend.snapshot();
1517 let start = all.len().saturating_sub(n);
1518 ctx.emit(all[start..].to_vec());
1519 },
1520 )
1521 }
1522
1523 pub fn slice(&self, start: usize, stop: Option<usize>) -> Node<Vec<T>> {
1525 let backend = self.backend.clone();
1526 Node::derived_opts(
1527 vec![self.delta.erased()],
1528 NodeOpts {
1529 partial: true,
1530 factory: Some("reactiveLog.slice".to_owned()),
1531 ..NodeOpts::default()
1532 },
1533 move |ctx| {
1534 let all = backend.snapshot();
1535 let end = stop.unwrap_or(all.len()).min(all.len());
1536 let start = start.min(end);
1537 ctx.emit(all[start..end].to_vec());
1538 },
1539 )
1540 }
1541
1542 pub fn page(&self, offset: usize, limit: usize) -> ReactiveView<LogChange<T>, Vec<T>> {
1544 let key = format!("{offset}:{limit}");
1545 if let Some((_, view)) = self
1546 .page_memo
1547 .borrow()
1548 .iter()
1549 .find(|(existing, _)| existing == &key)
1550 {
1551 return view.clone();
1552 }
1553 let name = self.graph.as_ref().map(|_| {
1554 let next = self.page_seq.get();
1555 self.page_seq.set(next + 1);
1556 format!("{}.page#{next}", self.id_prefix)
1557 });
1558 let backend = self.backend.clone();
1559 let materialize = move || {
1560 let all = backend.snapshot();
1561 let start = offset.min(all.len());
1562 let end = offset.saturating_add(limit).min(all.len());
1563 all[start..end].to_vec()
1564 };
1565 let memo = self.page_memo.clone();
1566 let dispose_key = key.clone();
1567 let view = light_reactive_view::<LogChange<T>, Vec<T>, _>(
1568 &self.delta,
1569 self.graph.as_ref(),
1570 ViewFactories {
1571 group: "reactiveLog.page",
1572 delta: "reactiveLog.page.delta",
1573 snapshot: "reactiveLog.page.snapshot",
1574 },
1575 name,
1576 materialize,
1577 Some(Box::new(move || {
1578 memo.borrow_mut()
1579 .retain(|(existing, _)| existing != &dispose_key);
1580 })),
1581 );
1582 self.page_memo.borrow_mut().push((key, view.clone()));
1583 view
1584 }
1585
1586 pub fn scan<A: Clone + 'static, F: Fn(A, &T) -> A + 'static>(
1588 &self,
1589 initial: A,
1590 step: F,
1591 ) -> Node<A> {
1592 let backend = self.backend.clone();
1593 Node::derived_opts(
1594 vec![self.delta.erased()],
1595 NodeOpts {
1596 partial: true,
1597 factory: Some("reactiveLog.scan".to_owned()),
1598 ..NodeOpts::default()
1599 },
1600 move |ctx| {
1601 let changes = ctx.batch::<LogChange<T>>(0);
1602 let appended = changes
1603 .iter()
1604 .map(|change| match change.as_ref() {
1605 LogChange::Append { .. } => 1,
1606 LogChange::AppendMany { values } => values.len(),
1607 LogChange::TrimHead { .. } | LogChange::Clear { .. } => 0,
1608 })
1609 .sum::<usize>();
1610 let reset_change = changes.iter().any(|change| {
1611 matches!(
1612 change.as_ref(),
1613 LogChange::TrimHead { .. } | LogChange::Clear { .. }
1614 )
1615 });
1616 let previous = ctx.state_get::<LogScanState<A>>();
1617 let mut state = previous
1618 .as_ref()
1619 .map(|s| LogScanState {
1620 acc: s.acc.clone(),
1621 processed: s.processed,
1622 })
1623 .unwrap_or_else(|| LogScanState {
1624 acc: initial.clone(),
1625 processed: 0,
1626 });
1627 let all = backend.snapshot();
1628 if reset_change || all.len() < state.processed {
1629 state.acc = initial.clone();
1630 state.processed = 0;
1631 }
1632 if appended > 0 && all.len() < state.processed.saturating_add(appended) {
1633 ctx.state_set(state);
1634 return;
1635 }
1636 for value in all.iter().skip(state.processed) {
1637 state.acc = step(state.acc, value);
1638 }
1639 state.processed = all.len();
1640 let out = state.acc.clone();
1641 ctx.state_set(state);
1642 ctx.emit(out);
1643 },
1644 )
1645 }
1646
1647 pub fn attach(&self, src: &Node<T>) -> Disposer {
1649 let graph = self.graph.as_ref().unwrap_or_else(|| {
1650 panic!("reactive_log.attach requires options.graph so the input fold is describe-visible (D61)")
1651 });
1652 let bind_idx = self.bind_seq.get();
1653 self.bind_seq.set(bind_idx + 1);
1654 let backend = self.backend.clone();
1655 let delta = self.delta.clone();
1656 let op = Operator::with_opts(
1657 "reactiveLog.bindSource",
1658 NodeOpts {
1659 partial: true,
1660 ..NodeOpts::default()
1661 },
1662 move |ctx| {
1663 for value in ctx.batch::<T>(0) {
1664 let value = (*value).clone();
1665 let trimmed = backend.append(value.clone());
1666 delta.down(vec![Message::Data(Rc::new(LogChange::Append { value }))]);
1667 emit_log_trimmed(&delta, trimmed);
1668 }
1669 },
1670 );
1671 let mut opts = GraphNodeOpts::named(format!("{}.bind#{bind_idx}", self.id_prefix));
1672 opts.meta
1673 .insert("kind".to_owned(), "collection_bind_source".to_owned());
1674 opts.meta
1675 .insert("collection".to_owned(), "reactiveLog".to_owned());
1676 let src_core = src.erased();
1677 let folder = graph.init_node::<LogChange<T>>(op, vec![src_core.clone()], opts);
1678 let unsub = folder.subscribe(|_| {});
1679 let slot = {
1680 let mut disposers = self.disposers.borrow_mut();
1681 let slot = disposers.len();
1682 disposers.push(Some(unsub));
1683 slot
1684 };
1685 let disposers = self.disposers.clone();
1686 let folder_for_dispose = folder.clone();
1687 Box::new(move || {
1688 folder_for_dispose.unsubscribe_dep(src_core, |_| {});
1689 if let Some(disposer) = disposers.borrow_mut().get_mut(slot).and_then(Option::take) {
1690 disposer();
1691 }
1692 })
1693 }
1694
1695 pub fn dispose(&self) {
1697 let views = self
1698 .page_memo
1699 .borrow()
1700 .iter()
1701 .map(|(_, view)| view.clone())
1702 .collect::<Vec<_>>();
1703 for view in views {
1704 view.dispose();
1705 }
1706 self.page_memo.borrow_mut().clear();
1707 for disposer in self.disposers.borrow_mut().iter_mut() {
1708 if let Some(disposer) = disposer.take() {
1709 disposer();
1710 }
1711 }
1712 }
1713
1714 fn emit(&self, change: LogChange<T>) {
1715 self.delta.down(vec![Message::Data(Rc::new(change))]);
1716 }
1717
1718 fn emit_trimmed(&self, n: usize) {
1719 emit_log_trimmed(&self.delta, n);
1720 }
1721}
1722
1723struct LogScanState<A> {
1724 acc: A,
1725 processed: usize,
1726}
1727
1728pub fn reactive_list<T: Clone + 'static>(
1730 initial: Vec<T>,
1731 options: ReactiveListOptions,
1732) -> ReactiveList<T> {
1733 ReactiveList::new(initial, options)
1734}
1735
1736pub fn reactive_log<T: Clone + 'static>(
1738 initial: Vec<T>,
1739 options: ReactiveLogOptions,
1740) -> ReactiveLog<T> {
1741 ReactiveLog::new(initial, options)
1742}
1743
1744pub fn restore_reactive_list<T: Clone + 'static>(
1746 state: crate::storage::ReactiveListRestoreState<T>,
1747 options: ReactiveListOptions,
1748) -> crate::storage::StorageResult<ReactiveList<T>> {
1749 if state.kind != crate::storage::ReactiveCollectionKind::ReactiveList {
1750 return Err(crate::storage::StorageError::backend(
1751 "restore_reactive_list: restore state kind must be reactiveList",
1752 ));
1753 }
1754 Ok(ReactiveList::new(state.state, options))
1755}
1756
1757pub fn restore_reactive_log<T: Clone + 'static>(
1759 state: crate::storage::ReactiveLogRestoreState<T>,
1760 options: ReactiveLogOptions,
1761) -> crate::storage::StorageResult<ReactiveLog<T>> {
1762 if state.kind != crate::storage::ReactiveCollectionKind::ReactiveLog {
1763 return Err(crate::storage::StorageError::backend(
1764 "restore_reactive_log: restore state kind must be reactiveLog",
1765 ));
1766 }
1767 Ok(ReactiveLog::new(state.state, options))
1768}
1769
1770pub fn restore_reactive_map<K, V>(
1772 state: crate::storage::ReactiveMapRestoreState<K, V>,
1773 options: ReactiveMapOptions,
1774) -> crate::storage::StorageResult<ReactiveMap<K, V>>
1775where
1776 K: Clone + Ord + fmt::Debug + 'static,
1777 V: Clone + 'static,
1778{
1779 if state.kind != crate::storage::ReactiveCollectionKind::ReactiveMap {
1780 return Err(crate::storage::StorageError::backend(
1781 "restore_reactive_map: restore state kind must be reactiveMap",
1782 ));
1783 }
1784 assert_unique_restore_map_keys(&state.state)?;
1785 Ok(ReactiveMap::new(state.state, options))
1786}
1787
1788pub fn restore_reactive_index<K, S, V>(
1790 state: crate::storage::ReactiveIndexRestoreState<K, S, V>,
1791 options: ReactiveIndexOptions,
1792) -> crate::storage::StorageResult<ReactiveIndex<K, S, V>>
1793where
1794 K: Clone + Ord + fmt::Debug + 'static,
1795 S: Clone + Ord + 'static,
1796 V: Clone + 'static,
1797{
1798 if state.kind != crate::storage::ReactiveCollectionKind::ReactiveIndex {
1799 return Err(crate::storage::StorageError::backend(
1800 "restore_reactive_index: restore state kind must be reactiveIndex",
1801 ));
1802 }
1803 assert_unique_restore_index_primaries(&state.state)?;
1804 Ok(ReactiveIndex::new(state.state, options))
1805}
1806
1807pub fn reactive_index<K, S, V>(
1809 initial: Vec<IndexRow<K, S, V>>,
1810 options: ReactiveIndexOptions,
1811) -> ReactiveIndex<K, S, V>
1812where
1813 K: Clone + Ord + fmt::Debug + 'static,
1814 S: Clone + Ord + 'static,
1815 V: Clone + 'static,
1816{
1817 ReactiveIndex::new(initial, options)
1818}
1819
1820pub fn reactive_map<K, V>(initial: Vec<(K, V)>, options: ReactiveMapOptions) -> ReactiveMap<K, V>
1822where
1823 K: Clone + Ord + fmt::Debug + 'static,
1824 V: Clone + 'static,
1825{
1826 ReactiveMap::new(initial, options)
1827}
1828
1829pub fn merge_reactive_logs<T: Clone + 'static>(logs: Vec<ReactiveLog<T>>) -> Node<LogChange<T>> {
1831 let deps = logs
1832 .iter()
1833 .map(|log| log.delta.erased())
1834 .collect::<Vec<_>>();
1835 Node::derived_opts(
1836 deps,
1837 NodeOpts {
1838 partial: true,
1839 factory: Some("mergeReactiveLogs".to_owned()),
1840 ..NodeOpts::default()
1841 },
1842 move |ctx| {
1843 for i in 0..logs.len() {
1844 for change in ctx.batch::<LogChange<T>>(i) {
1845 ctx.emit((*change).clone());
1846 }
1847 }
1848 },
1849 )
1850}
1851
1852pub fn scan_log<T: Clone + 'static, A: Clone + 'static, F: Fn(A, &T) -> A + 'static>(
1854 log: &ReactiveLog<T>,
1855 initial: A,
1856 step: F,
1857) -> Node<A> {
1858 log.scan(initial, step)
1859}
1860
1861fn make_delta_node<T: Clone + 'static>(
1862 graph: Option<&Graph>,
1863 id_prefix: &str,
1864 restorable: bool,
1865) -> Node<ListChange<T>> {
1866 match graph {
1867 Some(graph) => graph.empty_source(
1868 "reactiveList.delta",
1869 collection_node_opts(
1870 id_prefix,
1871 "delta",
1872 "collection_delta",
1873 "reactiveList.delta",
1874 restorable,
1875 ),
1876 ),
1877 None => Node::state_empty(),
1878 }
1879}
1880
1881fn make_log_delta_node<T: Clone + 'static>(
1882 graph: Option<&Graph>,
1883 id_prefix: &str,
1884 restorable: bool,
1885 max_size: Option<usize>,
1886) -> Node<LogChange<T>> {
1887 match graph {
1888 Some(graph) => graph.empty_source(
1889 "reactiveLog.delta",
1890 collection_node_opts_with_restore_config(
1891 id_prefix,
1892 "delta",
1893 "collection_delta",
1894 "reactiveLog.delta",
1895 restorable,
1896 max_size.map(|max_size| json!({ "maxSize": max_size })),
1897 ),
1898 ),
1899 None => Node::state_empty(),
1900 }
1901}
1902
1903fn make_index_delta_node<K: Clone + Ord + 'static, S: Clone + Ord + 'static, V: Clone + 'static>(
1904 graph: Option<&Graph>,
1905 id_prefix: &str,
1906 restorable: bool,
1907) -> Node<IndexChange<K, S, V>> {
1908 match graph {
1909 Some(graph) => graph.empty_source(
1910 "reactiveIndex.delta",
1911 collection_node_opts(
1912 id_prefix,
1913 "delta",
1914 "collection_delta",
1915 "reactiveIndex.delta",
1916 restorable,
1917 ),
1918 ),
1919 None => Node::state_empty(),
1920 }
1921}
1922
1923fn make_map_delta_node<K: Clone + Ord + 'static, V: Clone + 'static>(
1924 graph: Option<&Graph>,
1925 id_prefix: &str,
1926 restorable: bool,
1927) -> Node<MapChange<K, V>> {
1928 match graph {
1929 Some(graph) => graph.empty_source(
1930 "reactiveMap.delta",
1931 collection_node_opts(
1932 id_prefix,
1933 "delta",
1934 "collection_delta",
1935 "reactiveMap.delta",
1936 restorable,
1937 ),
1938 ),
1939 None => Node::state_empty(),
1940 }
1941}
1942
1943fn make_snapshot_node<T: Clone + 'static>(
1944 graph: Option<&Graph>,
1945 id_prefix: &str,
1946 delta: &Node<ListChange<T>>,
1947 backend: &ListBackend<T>,
1948 pull_id: LockId,
1949 restorable: bool,
1950) -> Node<Vec<T>> {
1951 let backend = backend.clone();
1952 let op = Operator::with_opts(
1953 "reactiveList.snapshot",
1954 NodeOpts {
1955 partial: true,
1956 pull_id: Some(pull_id),
1957 ..NodeOpts::default()
1958 },
1959 move |ctx| {
1960 let version = backend.version();
1961 let last = ctx.state_get::<u64>().map(|v| *v);
1962 if last == Some(version) {
1963 return;
1964 }
1965 ctx.state_set(version);
1966 ctx.emit(backend.snapshot());
1967 },
1968 );
1969 match graph {
1970 Some(graph) => graph.init_node(
1971 op,
1972 vec![delta.erased()],
1973 collection_node_opts(
1974 id_prefix,
1975 "snapshot",
1976 "collection_snapshot",
1977 "reactiveList.snapshot",
1978 restorable,
1979 ),
1980 ),
1981 None => init_node(op, vec![delta.erased()], NodeOpts::default()),
1982 }
1983}
1984
1985fn make_log_snapshot_node<T: Clone + 'static>(
1986 graph: Option<&Graph>,
1987 id_prefix: &str,
1988 delta: &Node<LogChange<T>>,
1989 backend: &LogBackend<T>,
1990 pull_id: LockId,
1991 restorable: bool,
1992) -> Node<Vec<T>> {
1993 let backend = backend.clone();
1994 let op = Operator::with_opts(
1995 "reactiveLog.snapshot",
1996 NodeOpts {
1997 partial: true,
1998 pull_id: Some(pull_id),
1999 ..NodeOpts::default()
2000 },
2001 move |ctx| {
2002 let version = backend.version();
2003 let last = ctx.state_get::<u64>().map(|v| *v);
2004 if last == Some(version) {
2005 return;
2006 }
2007 ctx.state_set(version);
2008 ctx.emit(backend.snapshot());
2009 },
2010 );
2011 match graph {
2012 Some(graph) => graph.init_node(
2013 op,
2014 vec![delta.erased()],
2015 collection_node_opts(
2016 id_prefix,
2017 "snapshot",
2018 "collection_snapshot",
2019 "reactiveLog.snapshot",
2020 restorable,
2021 ),
2022 ),
2023 None => init_node(op, vec![delta.erased()], NodeOpts::default()),
2024 }
2025}
2026
2027fn make_index_snapshot_node<
2028 K: Clone + Ord + 'static,
2029 S: Clone + Ord + 'static,
2030 V: Clone + 'static,
2031>(
2032 graph: Option<&Graph>,
2033 id_prefix: &str,
2034 delta: &Node<IndexChange<K, S, V>>,
2035 backend: &IndexBackend<K, S, V>,
2036 pull_id: LockId,
2037 restorable: bool,
2038) -> Node<Vec<IndexRow<K, S, V>>> {
2039 let backend = backend.clone();
2040 let op = Operator::with_opts(
2041 "reactiveIndex.snapshot",
2042 NodeOpts {
2043 partial: true,
2044 pull_id: Some(pull_id),
2045 ..NodeOpts::default()
2046 },
2047 move |ctx| {
2048 let version = backend.version();
2049 let last = ctx.state_get::<u64>().map(|v| *v);
2050 if last == Some(version) {
2051 return;
2052 }
2053 ctx.state_set(version);
2054 ctx.emit(backend.snapshot());
2055 },
2056 );
2057 match graph {
2058 Some(graph) => graph.init_node(
2059 op,
2060 vec![delta.erased()],
2061 collection_node_opts(
2062 id_prefix,
2063 "snapshot",
2064 "collection_snapshot",
2065 "reactiveIndex.snapshot",
2066 restorable,
2067 ),
2068 ),
2069 None => init_node(op, vec![delta.erased()], NodeOpts::default()),
2070 }
2071}
2072
2073fn make_map_snapshot_node<K: Clone + Ord + 'static, V: Clone + 'static>(
2074 graph: Option<&Graph>,
2075 id_prefix: &str,
2076 delta: &Node<MapChange<K, V>>,
2077 backend: &MapBackend<K, V>,
2078 pull_id: LockId,
2079 restorable: bool,
2080) -> Node<BTreeMap<K, V>> {
2081 let backend = backend.clone();
2082 let op = Operator::with_opts(
2083 "reactiveMap.snapshot",
2084 NodeOpts {
2085 partial: true,
2086 pull_id: Some(pull_id),
2087 ..NodeOpts::default()
2088 },
2089 move |ctx| {
2090 let version = backend.version();
2091 let last = ctx.state_get::<u64>().map(|v| *v);
2092 if last == Some(version) {
2093 return;
2094 }
2095 ctx.state_set(version);
2096 ctx.emit(backend.snapshot());
2097 },
2098 );
2099 match graph {
2100 Some(graph) => graph.init_node(
2101 op,
2102 vec![delta.erased()],
2103 collection_node_opts(
2104 id_prefix,
2105 "snapshot",
2106 "collection_snapshot",
2107 "reactiveMap.snapshot",
2108 restorable,
2109 ),
2110 ),
2111 None => init_node(op, vec![delta.erased()], NodeOpts::default()),
2112 }
2113}
2114
2115fn node_name(prefix: &str, suffix: &str) -> String {
2116 format!("{prefix}.{suffix}")
2117}
2118
2119fn collection_node_opts(
2120 id_prefix: &str,
2121 suffix: &str,
2122 kind: &str,
2123 restore_ref: &str,
2124 restorable: bool,
2125) -> GraphNodeOpts {
2126 collection_node_opts_with_restore_config(id_prefix, suffix, kind, restore_ref, restorable, None)
2127}
2128
2129fn collection_node_opts_with_restore_config(
2130 id_prefix: &str,
2131 suffix: &str,
2132 kind: &str,
2133 restore_ref: &str,
2134 restorable: bool,
2135 config: Option<Value>,
2136) -> GraphNodeOpts {
2137 let mut opts = GraphNodeOpts::named(node_name(id_prefix, suffix));
2138 opts.meta.insert("kind".to_owned(), kind.to_owned());
2139 if restorable {
2140 let restore = RestoreFactoryMeta::registry_ref(restore_ref);
2141 opts.restore = Some(match config {
2142 Some(config) => restore.with_config(config),
2143 None => restore,
2144 });
2145 }
2146 opts
2147}
2148
2149fn register_backend_state_json<T, F>(node: &Node<T>, snapshot: F)
2150where
2151 T: 'static,
2152 F: Fn() -> Result<Value, String> + 'static,
2153{
2154 register_backend_state_contributor(&node.erased(), Rc::new(move |_path| snapshot()));
2155}
2156
2157fn vec_to_checkpoint_json<T: Clone + 'static>(values: Vec<T>) -> Result<Value, String> {
2158 values
2159 .iter()
2160 .map(any_to_checkpoint_json)
2161 .collect::<Result<Vec<_>, _>>()
2162 .map(Value::Array)
2163}
2164
2165fn map_entries_to_checkpoint_json<K: Clone + 'static, V: Clone + 'static>(
2166 entries: Vec<(K, V)>,
2167) -> Result<Value, String> {
2168 entries
2169 .iter()
2170 .map(|(key, value)| {
2171 Ok(Value::Array(vec![
2172 any_to_checkpoint_json(key)?,
2173 any_to_checkpoint_json(value)?,
2174 ]))
2175 })
2176 .collect::<Result<Vec<_>, _>>()
2177 .map(Value::Array)
2178}
2179
2180fn index_rows_to_checkpoint_json<K: Clone + 'static, S: Clone + 'static, V: Clone + 'static>(
2181 rows: Vec<IndexRow<K, S, V>>,
2182) -> Result<Value, String> {
2183 rows.iter()
2184 .map(|row| {
2185 let mut out = serde_json::Map::new();
2186 out.insert("primary".to_owned(), any_to_checkpoint_json(&row.primary)?);
2187 out.insert(
2188 "secondary".to_owned(),
2189 any_to_checkpoint_json(&row.secondary)?,
2190 );
2191 out.insert("value".to_owned(), any_to_checkpoint_json(&row.value)?);
2192 Ok(Value::Object(out))
2193 })
2194 .collect::<Result<Vec<_>, _>>()
2195 .map(Value::Array)
2196}
2197
2198fn any_to_checkpoint_json<T: 'static>(value: &T) -> Result<Value, String> {
2199 let any = value as &dyn Any;
2200 if let Some(v) = any.downcast_ref::<Value>() {
2201 return Ok(v.clone());
2202 }
2203 if let Some(v) = any.downcast_ref::<String>() {
2204 return Ok(Value::String(v.clone()));
2205 }
2206 if let Some(v) = any.downcast_ref::<&'static str>() {
2207 return Ok(Value::String((*v).to_owned()));
2208 }
2209 if let Some(v) = any.downcast_ref::<bool>() {
2210 return Ok(Value::Bool(*v));
2211 }
2212 if let Some(v) = any.downcast_ref::<i32>() {
2213 return Ok(Value::Number(Number::from(*v)));
2214 }
2215 if let Some(v) = any.downcast_ref::<i64>() {
2216 return Ok(Value::Number(Number::from(*v)));
2217 }
2218 if let Some(v) = any.downcast_ref::<u32>() {
2219 return Ok(Value::Number(Number::from(*v)));
2220 }
2221 if let Some(v) = any.downcast_ref::<u64>() {
2222 return Ok(Value::Number(Number::from(*v)));
2223 }
2224 if let Some(v) = any.downcast_ref::<usize>() {
2225 return Ok(Value::Number(Number::from(*v as u64)));
2226 }
2227 if let Some(v) = any.downcast_ref::<f64>() {
2228 if let Some(n) = Number::from_f64(*v) {
2229 return Ok(Value::Number(n));
2230 }
2231 }
2232 Err("collection backend value is not strict JSON compatible (D160)".to_owned())
2233}
2234
2235fn assert_unique_restore_map_keys<K: 'static, V>(
2236 entries: &[(K, V)],
2237) -> crate::storage::StorageResult<()> {
2238 let mut seen = Vec::<Vec<u8>>::new();
2239 for (index, (key, _)) in entries.iter().enumerate() {
2240 let id = restore_json_identity(key)?;
2241 if seen.iter().any(|existing| existing == &id) {
2242 return Err(crate::storage::StorageError::backend(format!(
2243 "restore_reactive_map: entry {index} duplicates an earlier strict-JSON key"
2244 )));
2245 }
2246 seen.push(id);
2247 }
2248 Ok(())
2249}
2250
2251fn assert_unique_restore_index_primaries<K: 'static, S, V>(
2252 rows: &[IndexRow<K, S, V>],
2253) -> crate::storage::StorageResult<()> {
2254 let mut seen = Vec::<Vec<u8>>::new();
2255 for (index, row) in rows.iter().enumerate() {
2256 let id = restore_json_identity(&row.primary)?;
2257 if seen.iter().any(|existing| existing == &id) {
2258 return Err(crate::storage::StorageError::backend(format!(
2259 "restore_reactive_index: row {index} duplicates an earlier strict-JSON primary"
2260 )));
2261 }
2262 seen.push(id);
2263 }
2264 Ok(())
2265}
2266
2267fn restore_json_identity<T: 'static>(value: &T) -> crate::storage::StorageResult<Vec<u8>> {
2268 let json = any_to_checkpoint_json(value).map_err(crate::storage::StorageError::backend)?;
2269 crate::json::strict_canonical_json_bytes(&json)
2270 .map_err(|err| crate::storage::StorageError::backend(format!("restore_reactive_*: {err}")))
2271}
2272
2273fn light_reactive_view<C: Clone + 'static, S: Clone + 'static, F: Fn() -> S + 'static>(
2274 parent_delta: &Node<C>,
2275 graph: Option<&Graph>,
2276 factories: ViewFactories,
2277 name: Option<String>,
2278 materialize_snapshot: F,
2279 on_dispose: Option<ViewDisposeHook>,
2280) -> ReactiveView<C, S> {
2281 let group = graph.map(|graph| {
2282 graph.topology_group_opts(TopologyGroupOptions::named(
2283 name.clone().unwrap_or_else(|| factories.group.to_owned()),
2284 ))
2285 });
2286 let delta_op = Operator::with_opts(
2287 factories.delta,
2288 NodeOpts {
2289 partial: true,
2290 ..NodeOpts::default()
2291 },
2292 move |ctx| {
2293 for change in ctx.batch::<C>(0) {
2294 ctx.emit((*change).clone());
2295 }
2296 },
2297 );
2298 let delta = match group.as_ref() {
2299 Some(group) => {
2300 let mut opts =
2301 graph_node_opts_with_optional_name(name.as_ref().map(|n| format!("{n}.delta")));
2302 opts.meta
2303 .insert("kind".to_owned(), "collection_view_delta".to_owned());
2304 opts.meta
2305 .insert("factory".to_owned(), factories.group.to_owned());
2306 group.init_node(delta_op, vec![parent_delta.erased()], opts)
2307 }
2308 None => init_node(delta_op, vec![parent_delta.erased()], NodeOpts::default()),
2309 };
2310 let pull_id = LockId::new(match &name {
2311 Some(name) => format!("{name}.snapshot"),
2312 None => format!(
2313 "{}.snapshot#{}",
2314 factories.group,
2315 VIEW_PULL_SEQ.fetch_add(1, Ordering::Relaxed)
2316 ),
2317 });
2318 let snapshot_op = Operator::with_opts(
2319 factories.snapshot,
2320 NodeOpts {
2321 partial: true,
2322 pull_id: Some(pull_id.clone()),
2323 ..NodeOpts::default()
2324 },
2325 move |ctx| {
2326 ctx.emit(materialize_snapshot());
2327 },
2328 );
2329 let snapshot = match group.as_ref() {
2330 Some(group) => {
2331 let mut opts =
2332 graph_node_opts_with_optional_name(name.as_ref().map(|n| format!("{n}.snapshot")));
2333 opts.meta
2334 .insert("kind".to_owned(), "collection_view_snapshot".to_owned());
2335 opts.meta
2336 .insert("factory".to_owned(), factories.group.to_owned());
2337 group.init_node(snapshot_op, vec![delta.erased()], opts)
2338 }
2339 None => init_node(snapshot_op, vec![delta.erased()], NodeOpts::default()),
2340 };
2341 let delta_core = delta.erased();
2342 let snapshot_core = snapshot.erased();
2343 let group_for_dispose = group.clone();
2344 let on_dispose = Rc::new(RefCell::new(on_dispose));
2345 let dispose_action = Rc::new(move || {
2346 if let Some(group) = group_for_dispose.as_ref() {
2347 group.release_with_reason(factories.group);
2348 } else {
2349 assert!(
2350 snapshot_core.release_runtime_for_graph() && delta_core.release_runtime_for_graph(),
2351 "reactive view: cannot release runtime; view is not quiescent (D124)"
2352 );
2353 }
2354 if let Some(on_dispose) = on_dispose.borrow_mut().take() {
2355 on_dispose();
2356 }
2357 });
2358 ReactiveView {
2359 delta,
2360 snapshot,
2361 pull_id,
2362 disposed: Rc::new(Cell::new(false)),
2363 dispose_action,
2364 }
2365}
2366
2367fn graph_node_opts_with_optional_name(name: Option<String>) -> GraphNodeOpts {
2368 match name {
2369 Some(name) => GraphNodeOpts::named(name),
2370 None => GraphNodeOpts::default(),
2371 }
2372}
2373
2374fn emit_trimmed<T: Clone + 'static>(
2375 backend: &ListBackend<T>,
2376 delta: &Node<ListChange<T>>,
2377 max_size: Option<usize>,
2378) {
2379 let n = backend.enforce_max_size(max_size);
2380 if n > 0 {
2381 delta.down(vec![Message::Data(Rc::new(ListChange::<T>::TrimHead {
2382 n,
2383 }))]);
2384 }
2385}
2386
2387fn emit_log_trimmed<T: Clone + 'static>(delta: &Node<LogChange<T>>, n: usize) {
2388 if n > 0 {
2389 delta.down(vec![Message::Data(Rc::new(LogChange::<T>::TrimHead { n }))]);
2390 }
2391}
2392
2393fn normalize_read_index(index: isize, len: usize) -> Option<usize> {
2394 let len = isize::try_from(len).ok()?;
2395 let i = if index >= 0 { index } else { len + index };
2396 (i >= 0 && i < len).then_some(i as usize)
2397}
2398
2399fn trim_head_overflow<T>(buf: &mut Vec<T>, max_size: Option<usize>) -> usize {
2400 let Some(max_size) = max_size else {
2401 return 0;
2402 };
2403 assert!(max_size > 0, "max_size must be a positive integer");
2404 if buf.len() <= max_size {
2405 return 0;
2406 }
2407 let removed = buf.len() - max_size;
2408 buf.drain(0..removed);
2409 removed
2410}