1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
use std::sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
};
use super::{CellValue, Gettable, Watchable};
use crate::{
cell::{Cell, CellImmutable, CellMutable},
signal::Signal,
};
pub trait MergeMapExt<T>: Watchable<T> {
#[track_caller]
fn merge_map<U, F>(&self, f: F) -> Cell<U, CellImmutable>
where
T: CellValue,
U: CellValue,
F: Fn(&T) -> Cell<U, CellImmutable> + Send + Sync + 'static,
Self: Clone + Send + Sync + 'static,
{
let first_inner = f(&self.get());
let cell = Cell::<U, CellMutable>::new(first_inner.get());
let cell = if let Some(name) = self.name() {
cell.with_name(format!("{}::merge_map", name))
} else {
cell
};
// Complete when outer completes AND all inner cells complete
// Track: outer_complete flag + count of active (non-complete) inner cells
let outer_complete = Arc::new(AtomicBool::new(false));
let active_inners = Arc::new(AtomicUsize::new(1)); // Start with 1 for first_inner
// One-shot guard so the terminal `Complete` fires exactly once. The
// outer's completion and the final inner's completion are distinct
// same-height cells that can run concurrently under the scheduler's
// wave-parallel drain, and both can observe the "outer done AND no active
// inners" condition — a non-atomic check-then-notify that would emit
// `Complete` twice (confirmed by repro). Whoever wins this swap emits;
// the other skips. (A lost completion is impossible — at least one side
// observes the final state — so exactly-once holds.)
let completed_emitted = Arc::new(AtomicBool::new(false));
// Subscribe to first inner
let weak = cell.downgrade();
let oc = outer_complete.clone();
let ai = active_inners.clone();
let ce = completed_emitted.clone();
let first_inner_guard = first_inner.subscribe(move |signal| {
if let Some(c) = weak.upgrade() {
match signal {
Signal::Value(_) => c.notify(signal.clone()),
Signal::Complete => {
let remaining = ai.fetch_sub(1, Ordering::SeqCst) - 1;
if remaining == 0
&& oc.load(Ordering::SeqCst)
&& !ce.swap(true, Ordering::SeqCst)
{
c.notify(Signal::Complete);
}
}
Signal::Error(e) => c.notify(Signal::Error(e.clone())),
}
}
});
cell.own(first_inner_guard);
// When outer changes, subscribe to new inner (without unsubscribing from previous)
// Note: merge_map accumulates subscriptions by design - each inner cell stays subscribed
let weak_outer = cell.downgrade();
let f = Arc::new(f);
let first = Arc::new(AtomicBool::new(true));
let oc2 = outer_complete.clone();
let ai2 = active_inners.clone();
let ce2 = completed_emitted.clone();
let outer_guard = self.subscribe(move |signal| {
match signal {
Signal::Value(outer_value) => {
if first.swap(false, Ordering::SeqCst) {
return;
}
let Some(c) = weak_outer.upgrade() else {
return;
};
// Increment active count before creating inner
ai2.fetch_add(1, Ordering::SeqCst);
let inner = f(outer_value.as_ref());
// Subscribe to new inner - these subscriptions accumulate
let weak_inner = weak_outer.clone();
let oc_inner = oc2.clone();
let ai_inner = ai2.clone();
let ce_inner = ce2.clone();
let inner_guard = inner.subscribe(move |signal| {
if let Some(c) = weak_inner.upgrade() {
match signal {
Signal::Value(_) => c.notify(signal.clone()),
Signal::Complete => {
let remaining = ai_inner.fetch_sub(1, Ordering::SeqCst) - 1;
if remaining == 0
&& oc_inner.load(Ordering::SeqCst)
&& !ce_inner.swap(true, Ordering::SeqCst)
{
c.notify(Signal::Complete);
}
}
Signal::Error(e) => c.notify(Signal::Error(e.clone())),
}
}
});
c.own(inner_guard);
}
Signal::Complete => {
outer_complete.store(true, Ordering::SeqCst);
if active_inners.load(Ordering::SeqCst) == 0
&& let Some(c) = weak_outer.upgrade()
&& !ce2.swap(true, Ordering::SeqCst)
{
c.notify(Signal::Complete);
}
}
Signal::Error(e) => {
if let Some(c) = weak_outer.upgrade() {
c.notify(Signal::Error(e.clone()));
}
}
}
});
cell.own(outer_guard);
cell.lock()
}
}
impl<T, W: Watchable<T>> MergeMapExt<T> for W {}
#[cfg(test)]
mod tests {
use super::*;
use crate::{MapExt, MaterializeDefinite};
#[test]
fn test_merge_map_merges() {
let source = Cell::new(1u64);
let merged = source.merge_map(|v| {
// Must return CellImmutable, use map to create one
Cell::new(*v).map(|x| x * 10).materialize()
});
assert_eq!(merged.get(), 10);
}
}