Skip to main content

zng_var/var_impl/
merge_var.rs

1//! Variable that merges multiple others.
2
3///<span data-del-macro-root></span> Initializes a new [`Var<T>`](crate::Var) with value made
4/// by merging multiple other variables.
5///
6/// # Arguments
7///
8/// All arguments are separated by comma like a function call.
9///
10/// * `var0..N`: A list of *vars*, minimal 2.
11/// * `merge`: A new closure that produces a new value from references to all variable values. `FnMut(&var0_T, ..) -> merge_T`
12///
13/// Note that the *vars* can be of any of `Var<T>`, `ContextVar<T>` or `ResponseVar<T>`, that is, already constructed
14/// var types, not all types that convert into var.
15///
16/// # Contextualized
17///
18/// The merge var is contextualized when needed, meaning if any input is [contextual] at the moment the var is created it
19/// is also contextual. The full output type of this macro is a `Var<O>`, the `O` type is defined by the output of the merge closure.
20///
21/// [contextual]: crate::VarCapability::CONTEXT
22///
23/// # Examples
24///
25/// ```
26/// # use zng_var::*;
27/// # use zng_txt::*;
28/// # macro_rules! Text { ($($tt:tt)*) => { () } }
29/// let var0: Var<Txt> = var_from("Hello");
30/// let var1: Var<Txt> = var_from("World");
31///
32/// let greeting_text = Text!(merge_var!(var0, var1, |a, b| formatx!("{a} {b}!")));
33/// ```
34#[macro_export]
35macro_rules! merge_var {
36    ($($tt:tt)+) => {
37        $crate::__merge_var! {
38            $crate,
39            $($tt)+
40        }
41    };
42}
43
44use std::{any::TypeId, fmt::Write as _, marker::PhantomData, sync::Arc};
45
46use parking_lot::Mutex;
47use zng_clone_move::clmv;
48use zng_txt::Txt;
49#[doc(hidden)]
50pub use zng_var_proc_macros::merge_var as __merge_var;
51
52use crate::{
53    AnyVar, AnyVarModify, AnyVarValue, BoxAnyVarValue, ContextVar, Response, ResponseVar, Var, VarImpl, VarInstanceTag, VarModify,
54    VarValue, WeakAnyVar, any_contextual_var, any_var,
55};
56
57#[doc(hidden)]
58pub fn merge_var_input<I: VarValue>(input: impl MergeInput<I>) -> AnyVar {
59    input.into_merge_input().into()
60}
61
62#[doc(hidden)]
63pub fn merge_var_with(var: &AnyVar, visitor: &mut dyn FnMut(&dyn AnyVarValue)) {
64    var.0.with(visitor);
65}
66
67#[doc(hidden)]
68pub fn merge_var_output<O: VarValue>(output: O) -> BoxAnyVarValue {
69    BoxAnyVarValue::new(output)
70}
71
72#[doc(hidden)]
73pub fn merge_var<O: VarValue>(inputs: Box<[AnyVar]>, merge: impl FnMut(&[AnyVar]) -> BoxAnyVarValue + Send + 'static) -> Var<O> {
74    Var::new_any(merge_var_impl(inputs, Arc::new(Mutex::new(merge)), TypeId::of::<O>()))
75}
76
77#[doc(hidden)]
78#[diagnostic::on_unimplemented(note = "merge_var! and expr_var! inputs can be: Var<T>, ContextVar<T> or ResponseVar<T>")]
79pub trait MergeInput<T: VarValue> {
80    fn into_merge_input(self) -> Var<T>;
81}
82impl<T: VarValue> MergeInput<T> for Var<T> {
83    fn into_merge_input(self) -> Var<T> {
84        self
85    }
86}
87impl<T: VarValue> MergeInput<T> for ContextVar<T> {
88    fn into_merge_input(self) -> Var<T> {
89        self.into()
90    }
91}
92impl<T: VarValue> MergeInput<Response<T>> for ResponseVar<T> {
93    fn into_merge_input(self) -> Var<Response<T>> {
94        self.into()
95    }
96}
97
98fn merge_var_impl(inputs: Box<[AnyVar]>, merge: MergeFn, value_type: TypeId) -> AnyVar {
99    if inputs.iter().any(|i| i.capabilities().is_contextual()) {
100        return any_contextual_var(
101            move || {
102                let mut inputs = inputs.clone();
103                for v in inputs.iter_mut() {
104                    if v.capabilities().is_contextual() {
105                        *v = v.current_context();
106                    }
107                }
108                merge_var_tail(inputs, merge.clone())
109            },
110            value_type,
111        );
112    }
113    merge_var_tail(inputs, merge)
114}
115fn merge_var_tail(inputs: Box<[AnyVar]>, merge: MergeFn) -> AnyVar {
116    let output = any_var(merge.lock()(&inputs));
117    struct InputData {
118        inputs: Box<[AnyVar]>,
119        merge: MergeFn,
120        output_wk: WeakAnyVar,
121    }
122    let input_data = Arc::new(InputData {
123        inputs,
124        merge,
125        output_wk: output.downgrade(),
126    });
127    for input in &input_data.inputs {
128        let input_data_wk = Arc::downgrade(&input_data);
129        input
130            .hook(move |a| {
131                // modify on each input update, if multiple inputs update on the same cycle
132                // modify multiple times anyway, because services may be responding to the
133                // *partial* merge state as it happens.
134                if let Some(input_data) = input_data_wk.upgrade()
135                    && let Some(output) = input_data.output_wk.upgrade()
136                {
137                    let update = a.update();
138                    output.modify(clmv!(input_data_wk, |output| if let Some(input_data) = input_data_wk.upgrade() {
139                        let new_value = input_data.merge.lock()(&input_data.inputs);
140                        output.set(new_value);
141                        if update {
142                            output.update();
143                        }
144                    }));
145                    true
146                } else {
147                    false
148                }
149            })
150            .perm();
151    }
152
153    output.hold(input_data).perm();
154
155    output.read_only()
156}
157
158fn merge_var_bidi_impl(inputs: Box<[AnyVar]>, merge: MergeFn, map_back: MapBackFn, value_type: TypeId) -> AnyVar {
159    if inputs.iter().any(|i| i.capabilities().is_contextual()) {
160        return any_contextual_var(
161            move || {
162                let mut inputs = inputs.clone();
163                for v in inputs.iter_mut() {
164                    if v.capabilities().is_contextual() {
165                        *v = v.current_context();
166                    }
167                }
168                merge_var_bidi_tail(inputs, merge.clone(), map_back.clone())
169            },
170            value_type,
171        );
172    }
173    merge_var_bidi_tail(inputs, merge, map_back)
174}
175fn merge_var_bidi_tail(inputs: Box<[AnyVar]>, merge: MergeFn, map_back: MapBackFn) -> AnyVar {
176    let output = any_var(merge.lock()(&inputs));
177    struct InputData {
178        inputs: Box<[AnyVar]>,
179        merge: MergeFn,
180        map_back: MapBackFn,
181        output_wk: WeakAnyVar,
182    }
183    let input_data = Arc::new(InputData {
184        inputs,
185        merge,
186        map_back,
187        output_wk: output.downgrade(),
188    });
189    for input in &input_data.inputs {
190        let input_data_wk = Arc::downgrade(&input_data);
191        input
192            .hook(move |a| {
193                // modify on each input update, if multiple inputs update on the same cycle
194                // modify multiple times anyway, because services may be responding to the
195                // *partial* merge state as it happens.
196                if let Some(input_data) = input_data_wk.upgrade()
197                    && let Some(output) = input_data.output_wk.upgrade()
198                {
199                    let update = a.update();
200                    output.modify(clmv!(input_data_wk, |output| if let Some(input_data) = input_data_wk.upgrade() {
201                        let new_value = input_data.merge.lock()(&input_data.inputs);
202                        output.set(new_value);
203                        if update {
204                            output.update();
205                        }
206                    }));
207                    true
208                } else {
209                    false
210                }
211            })
212            .perm();
213    }
214
215    output
216        .hook(move |a| {
217            let mut map_back = input_data.map_back.lock();
218            for (i, input) in input_data.inputs.iter().enumerate() {
219                if input.capabilities().can_modify() {
220                    input.set(map_back(a.value(), i));
221                }
222            }
223            true
224        })
225        .perm();
226
227    output
228}
229
230fn merge_var_bidi_modify_impl(inputs: Box<[AnyVar]>, merge: MergeFn, modify_back: ModifyBackFn, value_type: TypeId) -> AnyVar {
231    if inputs.iter().any(|i| i.capabilities().is_contextual()) {
232        return any_contextual_var(
233            move || {
234                let mut inputs = inputs.clone();
235                for v in inputs.iter_mut() {
236                    if v.capabilities().is_contextual() {
237                        *v = v.current_context();
238                    }
239                }
240                merge_var_bidi_modify_tail(inputs, merge.clone(), modify_back.clone())
241            },
242            value_type,
243        );
244    }
245    merge_var_bidi_modify_tail(inputs, merge, modify_back)
246}
247fn merge_var_bidi_modify_tail(inputs: Box<[AnyVar]>, merge: MergeFn, modify_back: ModifyBackFn) -> AnyVar {
248    let output = any_var(merge.lock()(&inputs));
249    struct InputData {
250        inputs: Box<[AnyVar]>,
251        merge: MergeFn,
252        modify_back: ModifyBackFn,
253        output_wk: WeakAnyVar,
254    }
255    let input_data = Arc::new(InputData {
256        inputs,
257        merge,
258        modify_back,
259        output_wk: output.downgrade(),
260    });
261    #[derive(Debug, PartialEq, Clone, Copy)]
262    struct InputToOutputTag(VarInstanceTag);
263    #[derive(Debug, PartialEq, Clone, Copy)]
264    struct OutputToInputsTag(VarInstanceTag);
265    let input_to_output_tag = InputToOutputTag(output.var_instance_tag());
266    let output_to_inputs_tag = OutputToInputsTag(output.var_instance_tag());
267    for input in &input_data.inputs {
268        let input_data_wk = Arc::downgrade(&input_data);
269        input
270            .hook(move |a| {
271                // modify on each input update, if multiple inputs update on the same cycle
272                // modify multiple times anyway, because services may be responding to the
273                // *partial* merge state as it happens.
274                if let Some(input_data) = input_data_wk.upgrade()
275                    && let Some(output) = input_data.output_wk.upgrade()
276                {
277                    if a.contains_tag(&output_to_inputs_tag) {
278                        return true;
279                    }
280
281                    let update = a.update();
282                    output.modify(clmv!(input_data_wk, |output| if let Some(input_data) = input_data_wk.upgrade() {
283                        let new_value = input_data.merge.lock()(&input_data.inputs);
284                        let changed = output.set(new_value);
285                        if update {
286                            output.update();
287                        }
288                        if changed || update {
289                            output.push_tag(input_to_output_tag);
290                        }
291                    }));
292                    true
293                } else {
294                    false
295                }
296            })
297            .perm();
298    }
299
300    output
301        .hook(move |a| {
302            if a.contains_tag(&input_to_output_tag) {
303                return true;
304            }
305            for (i, input) in input_data.inputs.iter().enumerate() {
306                if input.capabilities().can_modify() {
307                    let output_wk = input_data.output_wk.clone();
308                    let modify_back = input_data.modify_back.clone();
309                    let update = a.update();
310                    input.modify(move |m| {
311                        if let Some(output) = output_wk.upgrade() {
312                            let has_updated = m.check_update(|m| {
313                                output.with(|o| {
314                                    modify_back.lock()(o, i, m);
315                                });
316                                if update {
317                                    m.update();
318                                }
319                            });
320                            if has_updated {
321                                m.push_tag(output_to_inputs_tag);
322                            }
323                        }
324                    });
325                }
326            }
327            true
328        })
329        .perm();
330
331    output
332}
333
334type MergeFn = Arc<Mutex<dyn FnMut(&[AnyVar]) -> BoxAnyVarValue + Send + 'static>>;
335type MapBackFn = Arc<Mutex<dyn FnMut(&dyn AnyVarValue, usize) -> BoxAnyVarValue + Send + 'static>>;
336type ModifyBackFn = Arc<Mutex<dyn FnMut(&dyn AnyVarValue, usize, &mut AnyVarModify) + Send + 'static>>;
337
338/// Build a [`merge_var!`] from any number of input vars of any type.
339#[derive(Clone)]
340pub struct AnyMergeVarBuilder {
341    inputs: Vec<AnyVar>,
342}
343impl Default for AnyMergeVarBuilder {
344    fn default() -> Self {
345        Self::new()
346    }
347}
348impl AnyMergeVarBuilder {
349    /// new empty.
350    pub fn new() -> Self {
351        Self { inputs: vec![] }
352    }
353
354    /// New with pre-allocated inputs.
355    pub fn with_capacity(capacity: usize) -> Self {
356        Self {
357            inputs: Vec::with_capacity(capacity),
358        }
359    }
360
361    /// Push an input.
362    pub fn push(&mut self, input: AnyVar) {
363        self.inputs.push(input);
364    }
365
366    fn read_only_inputs(self) -> Box<[AnyVar]> {
367        let mut inputs = self.inputs;
368        for input in &mut inputs {
369            if !input.capabilities().is_always_read_only() {
370                let v = input.read_only();
371                *input = v;
372            }
373        }
374        inputs.into_boxed_slice()
375    }
376
377    /// Build a read-only merge var.
378    pub fn build_any(self, merge: impl FnMut(&[AnyVar]) -> BoxAnyVarValue + Send + 'static, output_type: TypeId) -> AnyVar {
379        merge_var_impl(self.read_only_inputs(), Arc::new(Mutex::new(merge)), output_type)
380    }
381
382    /// Build a read-only strongly typed merge var.
383    pub fn build<O: VarValue>(self, mut merge: impl FnMut(&[AnyVar]) -> O + Send + 'static) -> Var<O> {
384        Var::new_any(self.build_any(move |inputs| BoxAnyVarValue::new(merge(inputs)), TypeId::of::<O>()))
385    }
386
387    /// Build a read-write merge var.
388    pub fn build_bidi_any(
389        self,
390        merge: impl FnMut(&[AnyVar]) -> BoxAnyVarValue + Send + 'static,
391        map_back: impl FnMut(&dyn AnyVarValue, usize) -> BoxAnyVarValue + Send + 'static,
392        output_type: TypeId,
393    ) -> AnyVar {
394        merge_var_bidi_impl(
395            self.inputs.into_boxed_slice(),
396            Arc::new(Mutex::new(merge)),
397            Arc::new(Mutex::new(map_back)),
398            output_type,
399        )
400    }
401
402    /// Build a read-write merge var that modifies each input back.
403    pub fn build_bidi_any_modify(
404        self,
405        merge: impl FnMut(&[AnyVar]) -> BoxAnyVarValue + Send + 'static,
406        modify_back: impl FnMut(&dyn AnyVarValue, usize, &mut AnyVarModify) + Send + 'static,
407        output_type: TypeId,
408    ) -> AnyVar {
409        merge_var_bidi_modify_impl(
410            self.inputs.into_boxed_slice(),
411            Arc::new(Mutex::new(merge)),
412            Arc::new(Mutex::new(modify_back)),
413            output_type,
414        )
415    }
416
417    /// Build a read-write strongly typed merge var.
418    pub fn build_bidi<O: VarValue>(
419        self,
420        mut merge: impl FnMut(&[AnyVar]) -> O + Send + 'static,
421        mut map_back: impl FnMut(&O, usize) -> BoxAnyVarValue + Send + 'static,
422    ) -> Var<O> {
423        Var::new_any(self.build_bidi_any(
424            move |inputs| BoxAnyVarValue::new(merge(inputs)),
425            move |output, input_idx| map_back(output.downcast_ref::<O>().unwrap(), input_idx),
426            TypeId::of::<O>(),
427        ))
428    }
429
430    /// Build a read-write strongly typed merge var that modifies each input back.
431    pub fn build_bidi_modify<O: VarValue>(
432        self,
433        mut merge: impl FnMut(&[AnyVar]) -> O + Send + 'static,
434        mut modify_back: impl FnMut(&O, usize, &mut AnyVarModify) + Send + 'static,
435    ) -> Var<O> {
436        Var::new_any(self.build_bidi_any_modify(
437            move |inputs| BoxAnyVarValue::new(merge(inputs)),
438            move |output, input_idx, m| modify_back(output.downcast_ref::<O>().unwrap(), input_idx, m),
439            TypeId::of::<O>(),
440        ))
441    }
442
443    /// Convert into a [`MergeVarBuilder<I>`].
444    ///
445    /// # Panics
446    ///
447    /// Panics if any input is not of type `I`.
448    pub fn into_typed<I: VarValue>(self) -> MergeVarBuilder<I> {
449        let i_id = TypeId::of::<I>();
450        for input in &self.inputs {
451            assert_eq!(i_id, input.value_type());
452        }
453        MergeVarBuilder {
454            builder: self,
455            _type: PhantomData,
456        }
457    }
458}
459
460/// Build a [`merge_var!`] from any number of input vars of the same type `I`.
461#[derive(Clone)]
462pub struct MergeVarBuilder<I: VarValue> {
463    builder: AnyMergeVarBuilder,
464    _type: PhantomData<fn() -> I>,
465}
466impl<I: VarValue> MergeVarBuilder<I> {
467    /// New empty.
468    pub fn new() -> Self {
469        Self {
470            builder: AnyMergeVarBuilder::new(),
471            _type: PhantomData,
472        }
473    }
474
475    /// New with pre-allocated inputs.
476    pub fn with_capacity(capacity: usize) -> Self {
477        Self {
478            builder: AnyMergeVarBuilder::with_capacity(capacity),
479            _type: PhantomData,
480        }
481    }
482
483    /// Push an input.
484    pub fn push(&mut self, input: impl MergeInput<I>) {
485        self.builder.push(input.into_merge_input().into())
486    }
487
488    /// Build a red-only merge var.
489    pub fn build<O: VarValue>(self, mut merge: impl FnMut(MergeVarInputs<I>) -> O + Send + 'static) -> Var<O> {
490        self.builder.build(move |inputs| {
491            merge(MergeVarInputs {
492                inputs,
493                _input_type: PhantomData,
494            })
495        })
496    }
497
498    /// Build a read-write merge var.
499    pub fn build_bidi<O: VarValue>(
500        self,
501        mut merge: impl FnMut(MergeVarInputs<I>) -> O + Send + 'static,
502        mut map_back: impl FnMut(&O, usize) -> I + Send + 'static,
503    ) -> Var<O> {
504        self.builder.build_bidi(
505            move |inputs| {
506                merge(MergeVarInputs {
507                    inputs,
508                    _input_type: PhantomData,
509                })
510            },
511            move |output, input_idx| BoxAnyVarValue::new(map_back(output, input_idx)),
512        )
513    }
514
515    /// Build a read-write merge var that modifies each input back.
516    pub fn build_bidi_modify<O: VarValue>(
517        self,
518        mut merge: impl FnMut(MergeVarInputs<I>) -> O + Send + 'static,
519        mut modify_back: impl FnMut(&O, usize, VarModify<I>) + Send + 'static,
520    ) -> Var<O> {
521        self.builder.build_bidi_modify(
522            move |inputs| {
523                merge(MergeVarInputs {
524                    inputs,
525                    _input_type: PhantomData,
526                })
527            },
528            move |output, input_idx, m| modify_back(output, input_idx, m.downcast().unwrap()),
529        )
530    }
531}
532impl<I: VarValue> Default for MergeVarBuilder<I> {
533    fn default() -> Self {
534        Self::new()
535    }
536}
537impl<T: VarValue + AsRef<str>> MergeVarBuilder<T> {
538    /// Build to a var that joins texts placing a `separator` between each.
539    pub fn join_txt(self, separator: impl Into<Txt>) -> Var<Txt> {
540        self.join_txt_impl(separator.into())
541    }
542    fn join_txt_impl(self, separator: Txt) -> Var<Txt> {
543        self.build(move |t| {
544            let mut s = String::new();
545            let mut sep = "";
546            t.with_each(|_, t| {
547                write!(&mut s, "{sep}{}", t.as_ref()).unwrap();
548                sep = &separator;
549            });
550            s.into()
551        })
552    }
553}
554
555/// Strongly typed input vars for [`MergeVarBuilder`]
556pub struct MergeVarInputs<'a, I: VarValue> {
557    inputs: &'a [AnyVar],
558    _input_type: PhantomData<fn() -> &'a I>,
559}
560impl<'a, I: VarValue> MergeVarInputs<'a, I> {
561    /// Number of inputs.
562    pub fn len(&self) -> usize {
563        self.inputs.len()
564    }
565
566    /// If has no inputs.
567    pub fn is_empty(&self) -> bool {
568        self.inputs.is_empty()
569    }
570
571    /// Visit the current value of the `index` input.
572    pub fn with<R>(&self, index: usize, visit: impl FnOnce(&I) -> R) -> R {
573        self.inputs[index].with(|v| visit(v.downcast_ref().unwrap()))
574    }
575
576    /// Clone the current value of the `index` input.
577    pub fn get(&self, index: usize) -> I {
578        self.with(index, |v| v.clone())
579    }
580
581    /// Clone the `index` input var.
582    pub fn input(&self, index: usize) -> Var<I> {
583        self.inputs[index].read_only().downcast().unwrap()
584    }
585
586    /// Iterate over clones of the current value of each input.
587    pub fn iter(&self) -> std::iter::Map<std::slice::Iter<'a, AnyVar>, fn(&AnyVar) -> I> {
588        Self {
589            inputs: self.inputs,
590            _input_type: self._input_type,
591        }
592        .into_iter()
593    }
594
595    /// Visit the current value of each input.
596    pub fn with_each(&self, mut visit: impl FnMut(usize, &I)) {
597        for (i, var) in self.inputs.iter().enumerate() {
598            var.with(|v| visit(i, v.downcast_ref().unwrap()))
599        }
600    }
601}
602impl<'a, I: VarValue> std::iter::IntoIterator for MergeVarInputs<'a, I> {
603    type Item = I;
604
605    type IntoIter = std::iter::Map<std::slice::Iter<'a, AnyVar>, fn(&AnyVar) -> I>;
606
607    fn into_iter(self) -> Self::IntoIter {
608        self.inputs.iter().map(downcast)
609    }
610}
611fn downcast<I: VarValue>(v: &AnyVar) -> I {
612    v.with(|v| v.downcast_ref::<I>().unwrap().clone())
613}