DFIR Operators
In our previous examples we made use of some of DFIR's operators. Here we document each operator in more detail. Many of these operators are based on the Rust equivalents for iterators; see the Rust documentation.
Maps:
_counter, enumerate, identity, inspect, map, resolve_futures, resolve_futures_blocking, resolve_futures_blocking_ordered, resolve_futures_ordered
Simple one-in-one-out operators.
Filters:
filter, filter_map
One-in zero-or-one-out operators.
Flattens:
flat_map, flat_map_stream_blocking, flatten, flatten_stream_blocking
One-in multiple-out operators.
Folds:
fold, fold_no_replay, reduce, reduce_no_replay, scan, scan_async_blocking
Operators which accumulate elements together.
Keyed Folds:
fold_keyed, reduce_keyed
Operators which accumulate elements together by key.
Lattice Folds:
lattice_fold, lattice_reduce
Folds based on lattice-merge.
Persistent Operators:
multiset_delta, defer_signal, persist, sort, sort_by_key, state, state_by, unique
Persistent (stateful) operators.
Multi-Input Operators:
anti_join, chain, chain_first_n, cross_join, cross_join_multiset, cross_singleton, difference, join, join_fused, join_fused_lhs, join_fused_rhs, join_multiset, join_multiset_half, lattice_bimorphism, union, zip, zip_longest
Operators with multiple inputs.
Multi-Output Operators:
demux_enum, partition, tee, unzip
Operators with multiple outputs.
Sources:
initialize, iter_ref, null, spin, source_file, source_interval, source_iter, source_json, source_stdin, source_stream, source_stream_serde
Operators which produce output elements (and consume no inputs).
Sinks:
dest_file, dest_sink, dest_sink_serde, for_each, null
Operators which consume input elements (and produce no outputs).
Control Flow Operators:
assert, assert_eq, defer_tick, defer_tick_lazy
Operators which affect control flow/scheduling.
Compiler Fusion Operators:
_lattice_fold_batch, _lattice_join_fused_join
Operators which are necessary to implement certain optimizations and rewrite rules.
Windowing Operators:
batch, batch_eager, batch_lazy
Operators for windowing loop inputs.
Un-Windowing Operators:
all_iterations
Operators for collecting loop outputs.
all_iterations
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> all_iterations() -> | exactly 1 | Streaming |
anti_join
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]anti_join() -> | exactly 1 | Streaming |
Input port names:
pos(streaming),neg(streaming)
assert
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> assert(A) | at least 0 and at most 1 | Streaming |
assert_eq
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> assert_eq(A) | at least 0 and at most 1 | Streaming |
batch
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> batch() -> | exactly 1 | Streaming |
batch_eager
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> batch_eager() -> | exactly 1 | Streaming |
batch_lazy
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> batch_lazy() -> | exactly 1 | Streaming |
chain
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]chain() -> | exactly 1 | Streaming |
chain_first_n
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| at least 2 | -> [<input_port>]chain_first_n(A) -> | exactly 1 | Streaming |
cross_join
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]cross_join() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
2 input streams of type S and T, 1 output stream of type (S, T)
Forms the cross-join (Cartesian product) of the items in the input streams, returning all tupled pairs.
source_iter(vec!["happy", "sad"]) -> [0]my_join;
source_iter(vec!["dog", "cat"]) -> [1]my_join;
my_join = cross_join() -> assert_eq([("happy", "dog"), ("sad", "dog"), ("happy", "cat"), ("sad", "cat")]);
cross_join can be provided with one or two generic lifetime persistence arguments
in the same way as join, see join's documentation for more info.
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<&str>();
let mut flow = dfir_rs::dfir_syntax! {
my_join = cross_join::<'tick>();
source_iter(["hello", "bye"]) -> [0]my_join;
source_stream(input_recv) -> [1]my_join;
my_join -> for_each(|(s, t)| println!("({}, {})", s, t));
};
input_send.send("oakland").unwrap();
flow.run_tick();
input_send.send("san francisco").unwrap();
flow.run_tick();
Prints only "(hello, oakland)" and "(bye, oakland)". The source_iter is only included in
the first tick, then forgotten, so when "san francisco" arrives on input [1] in the second tick,
there is nothing for it to match with from input [0], and therefore it does appear in the output.
cross_join_multiset
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]cross_join_multiset() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
cross_singleton
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]cross_singleton() -> | exactly 1 | Streaming |
Input port names:
input(streaming),single(streaming)
defer_signal
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]defer_signal() -> | exactly 1 | Streaming |
Input port names:
input(streaming),signal(streaming)
defer_tick
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> defer_tick() -> | exactly 1 | Blocking |
defer_tick_lazy
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> defer_tick_lazy() -> | exactly 1 | Blocking |
demux_enum
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> demux_enum() | any number of | Streaming |
Output port names: Variadic, as specified in arguments.
Generic Argument: A enum type which has
#[derive(DemuxEnum)]. Must match the items in the input stream.
Takes an input stream of enum instances and splits them into their variants.
#[derive(DemuxEnum)]
enum Shape {
Square(f64),
Rectangle(f64, f64),
Circle { r: f64 },
Triangle { w: f64, h: f64 }
}
let mut df = dfir_syntax! {
my_demux = source_iter([
Shape::Square(9.0),
Shape::Rectangle(10.0, 8.0),
Shape::Circle { r: 5.0 },
Shape::Triangle { w: 12.0, h: 13.0 },
]) -> demux_enum::<Shape>();
my_demux[Square] -> map(|s| s * s) -> out;
my_demux[Circle] -> map(|(r,)| std::f64::consts::PI * r * r) -> out;
my_demux[Rectangle] -> map(|(w, h)| w * h) -> out;
my_demux[Circle] -> map(|(w, h)| 0.5 * w * h) -> out;
out = union() -> for_each(|area| println!("Area: {}", area));
};
df.run_available();
dest_file
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> dest_file(A, B) | exactly 0 | Streaming |
0 input streams, 1 output stream
Arguments: (1) An
AsRef<Path>for a file to write to, and (2) a boolappend.
Consumes Strings by writing them as lines to a file. The file will be created if it doesn't
exist. Lines will be appended to the file if append is true, otherwise the file will be
truncated before lines are written.
Note this operator must be used within a Tokio runtime.
source_iter(1..=10) -> map(|n| format!("Line {}", n)) -> dest_file("dest.txt", false);
dest_sink
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> dest_sink(A) | exactly 0 | Streaming |
dest_sink_serde
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> dest_sink_serde(A) | exactly 0 | Streaming |
difference
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]difference() -> | exactly 1 | Streaming |
Input port names:
pos(streaming),neg(streaming)
enumerate
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> enumerate() -> | exactly 1 | Streaming |
1 input stream of type
T, 1 output stream of type(usize, T)
For each item passed in, enumerate it with its index: (0, x_0), (1, x_1), etc.
enumerate can also be provided with one generic lifetime persistence argument, either
'tick or 'static, to specify if indexing resets. If 'tick (the default) is specified, indexing will
restart at zero at the start of each tick. Otherwise 'static will never reset
and count monotonically upwards.
source_iter(vec!["hello", "world"])
-> enumerate()
-> assert_eq([(0, "hello"), (1, "world")]);
filter
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> filter(A) -> | exactly 1 | Streaming |
filter_map
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> filter_map(A) -> | exactly 1 | Streaming |
1 input stream, 1 output stream
An operator that both filters and maps. It yields only the items for which the supplied closure returns Some(value).
Note: The closure has access to the
contextobject.
source_iter(vec!["1", "hello", "world", "2"])
-> filter_map(|s| s.parse::<usize>().ok())
-> assert_eq([1, 2]);
flat_map
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> flat_map(A) -> | exactly 1 | Streaming |
1 input stream, 1 output stream
Arguments: A Rust closure that handles an iterator
For each item i passed in, treat i as an iterator and map the closure to that
iterator to produce items one by one. The type of the input items must be iterable.
Note: The closure has access to the
contextobject.
// should print out each character of each word on a separate line
source_iter(vec!["hello", "world"])
-> flat_map(|x| x.chars())
-> assert_eq(['h', 'e', 'l', 'l', 'o', 'w', 'o', 'r', 'l', 'd']);
flat_map_stream_blocking
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> flat_map_stream_blocking(A) -> | exactly 1 | Streaming |
1 input stream, 1 output stream
Arguments: A Rust closure that maps each item to a
Stream
For each item passed in, the closure is applied to produce a Stream, and the items
of that stream are emitted one by one. When the inner stream yields Pending, this
operator yields Pending as well.
source_iter(vec![1, 2, 3])
-> flat_map_stream_blocking(|x| futures::stream::iter(vec![x, x * 10]))
-> assert_eq([1, 10, 2, 20, 3, 30]);
flatten
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> flatten() -> | exactly 1 | Streaming |
flatten_stream_blocking
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> flatten_stream_blocking() -> | exactly 1 | Streaming |
fold
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> fold(A, B) | at least 0 and at most 1 | Streaming |
1 input stream, 1 output stream
Arguments: two arguments, both closures. The first closure is used to create the initial value for the accumulator, and the second is used to combine new items with the existing accumulator value. The second closure takes two two arguments: an
&mut Accumaccumulated value, and anItem.
Akin to Rust's built-in fold
operator, except that it takes the accumulator by &mut instead of by value. Folds every item
into an accumulator by applying a closure, returning the final result.
Note: The closures have access to the
contextobject.
fold can also be provided with one generic lifetime persistence argument, either
'tick or 'static, to specify how data persists. With 'tick, Items will only be collected
within the same tick. With 'static, the accumulated value will be remembered across ticks and
will be aggregated with items arriving in later ticks. When not explicitly specified
persistence defaults to 'tick.
// should print `Reassembled vector [1,2,3,4,5]`
source_iter([1,2,3,4,5])
-> fold::<'tick>(Vec::new, |accum: &mut Vec<_>, elem| {
accum.push(elem);
})
-> assert_eq([vec![1, 2, 3, 4, 5]]);
fold_keyed
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> fold_keyed(A, B) -> | exactly 1 | Streaming |
fold_no_replay
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> fold_no_replay(A, B) | at least 0 and at most 1 | Streaming |
for_each
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> for_each(A) | exactly 0 | Streaming |
identity
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> identity() -> | exactly 1 | Streaming |
initialize
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | initialize() -> | exactly 1 | Streaming |
0 input streams, 1 output stream
Arguments: None.
Emits a single unit () at the start of the first tick.
initialize()
-> assert_eq([()]);
inspect
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> inspect(A) | at least 0 and at most 1 | Streaming |
Arguments: A single closure
FnMut(&Item).
An operator which allows you to "inspect" each element of a stream without modifying it. The closure is called on a reference to each item. This is mainly useful for debugging as in the example below, and it is generally an anti-pattern to provide a closure with side effects.
Note: The closure has access to the
contextobject.
source_iter([1, 2, 3, 4])
-> inspect(|x| println!("{}", x))
-> assert_eq([1, 2, 3, 4]);
iter_ref
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | iter_ref(A) -> | exactly 1 | Streaming |
join
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]join() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
2 input streams of type
<(K, V1)>and<(K, V2)>, 1 output stream of type<(K, (V1, V2))>
Forms the equijoin of the tuples in the input streams by their first (key) attribute. Note that the result nests the 2nd input field (values) into a tuple in the 2nd output field.
source_iter(vec![("hello", "world"), ("stay", "gold"), ("hello", "world")]) -> [0]my_join;
source_iter(vec![("hello", "cleveland")]) -> [1]my_join;
my_join = join()
-> assert_eq([("hello", ("world", "cleveland"))]);
join can also be provided with one or two generic lifetime persistence arguments, either
'tick or 'static, to specify how join data persists. With 'tick, pairs will only be
joined with corresponding pairs within the same tick. With 'static, pairs will be remembered
across ticks and will be joined with pairs arriving in later ticks. When not explicitly
specified persistence defaults to `tick.
When two persistence arguments are supplied the first maps to port 0 and the second maps to
port 1.
When a single persistence argument is supplied, it is applied to both input ports.
When no persistence arguments are applied it defaults to 'tick for both.
The syntax is as follows:
join(); // Or
join::<'static>();
join::<'tick>();
join::<'static, 'tick>();
join::<'tick, 'static>();
// etc.
join is defined to treat its inputs as sets, meaning that it
eliminates duplicated values in its inputs. If you do not want
duplicates eliminated, use the join_multiset operator.
Examples
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<(&str, &str)>();
let mut flow = dfir_rs::dfir_syntax! {
source_iter([("hello", "world")]) -> [0]my_join;
source_stream(input_recv) -> [1]my_join;
my_join = join::<'tick>() -> for_each(|(k, (v1, v2))| println!("({}, ({}, {}))", k, v1, v2));
};
input_send.send(("hello", "oakland")).unwrap();
flow.run_tick();
input_send.send(("hello", "san francisco")).unwrap();
flow.run_tick();
Prints out "(hello, (world, oakland))" since source_iter([("hello", "world")]) is only
included in the first tick, then forgotten.
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<(&str, &str)>();
let mut flow = dfir_rs::dfir_syntax! {
source_iter([("hello", "world")]) -> [0]my_join;
source_stream(input_recv) -> [1]my_join;
my_join = join::<'static>() -> for_each(|(k, (v1, v2))| println!("({}, ({}, {}))", k, v1, v2));
};
input_send.send(("hello", "oakland")).unwrap();
flow.run_tick();
input_send.send(("hello", "san francisco")).unwrap();
flow.run_tick();
Prints out "(hello, (world, oakland))" and "(hello, (world, san francisco))" since the
inputs are peristed across ticks.
join_fused
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]join_fused(A, B) -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
join_fused_lhs
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]join_fused_lhs(A) -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
join_fused_rhs
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]join_fused_rhs(A) -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
join_multiset
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]join_multiset() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
join_multiset_half
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]join_multiset_half() -> | exactly 1 | Streaming |
Input port names:
build(streaming),probe(streaming)
lattice_bimorphism
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]lattice_bimorphism(A, B, C) -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
lattice_fold
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> lattice_fold(A) -> | exactly 1 | Streaming |
lattice_reduce
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> lattice_reduce() -> | exactly 1 | Streaming |
map
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> map(A) -> | exactly 1 | Streaming |
1 input stream, 1 output stream
Arguments: A Rust closure
For each item passed in, apply the closure to generate an item to emit.
If you do not want to modify the item stream and instead only want to view
each item use the inspect operator instead.
Note: The closure has access to the
contextobject.
source_iter(vec!["hello", "world"]) -> map(|x| x.to_uppercase())
-> assert_eq(["HELLO", "WORLD"]);
multiset_delta
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> multiset_delta() -> | exactly 1 | Streaming |
Multiset delta from the previous tick.
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<u32>();
let mut flow = dfir_rs::dfir_syntax! {
source_stream(input_recv)
-> multiset_delta()
-> for_each(|n| println!("{}", n));
};
input_send.send(3).unwrap();
input_send.send(4).unwrap();
input_send.send(3).unwrap();
flow.run_tick();
// 3, 4,
input_send.send(3).unwrap();
input_send.send(5).unwrap();
input_send.send(3).unwrap();
input_send.send(3).unwrap();
flow.run_tick();
// 5, 3
// First two "3"s are removed due to previous tick.
null
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| at least 0 and at most 1 | null() | at least 0 and at most 1 | Streaming |
partition
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> partition(A)[<output_port>] -> | at least 2 | Streaming |
Output port names: Variadic, as specified in arguments.
persist
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> persist() -> | exactly 1 | Streaming |
reduce
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> reduce(A) | at least 0 and at most 1 | Streaming |
reduce_keyed
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> reduce_keyed(A) -> | exactly 1 | Streaming |
1 input stream of type
(K, V), 1 output stream of type(K, V). The output will have one tuple for each distinctK, with an accumulated (reduced) value of typeV.
If you need the accumulated value to have a different type than the input, use fold_keyed.
Arguments: one Rust closures. The closure takes two arguments: an
&mut'accumulator', and an element. Accumulator should be updated based on the element.
A special case of reduce, in the spirit of SQL's GROUP BY and aggregation constructs. The input
is partitioned into groups by the first field, and for each group the values in the second
field are accumulated via the closures in the arguments.
Note: The closures have access to the
contextobject.
reduce_keyed can also be provided with one generic lifetime persistence argument, either
'tick or 'static, to specify how data persists. With 'tick, values will only be collected
within the same tick. With 'static, values will be remembered across ticks and will be
aggregated with pairs arriving in later ticks. When not explicitly specified persistence
defaults to 'tick.
reduce_keyed can also be provided with two type arguments, the key and value type. This is
required when using 'static persistence if the compiler cannot infer the types.
source_iter([("toy", 1), ("toy", 2), ("shoe", 11), ("shoe", 35), ("haberdashery", 7)])
-> reduce_keyed(|old: &mut u32, val: u32| *old += val)
-> assert_eq([("toy", 3), ("shoe", 46), ("haberdashery", 7)]);
Example using 'tick persistence and type arguments:
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<(&str, &str)>();
let mut flow = dfir_rs::dfir_syntax! {
source_stream(input_recv)
-> reduce_keyed::<'tick, &str>(|old: &mut _, val| *old = std::cmp::max(*old, val))
-> for_each(|(k, v)| println!("({:?}, {:?})", k, v));
};
input_send.send(("hello", "oakland")).unwrap();
input_send.send(("hello", "berkeley")).unwrap();
input_send.send(("hello", "san francisco")).unwrap();
flow.run_available();
// ("hello", "oakland, berkeley, san francisco, ")
input_send.send(("hello", "palo alto")).unwrap();
flow.run_available();
// ("hello", "palo alto, ")
reduce_no_replay
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> reduce_no_replay(A) | at least 0 and at most 1 | Streaming |
resolve_futures
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> resolve_futures() -> | exactly 1 | Streaming |
resolve_futures_blocking
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> resolve_futures_blocking() -> | exactly 1 | Streaming |
resolve_futures_blocking_ordered
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> resolve_futures_blocking_ordered() -> | exactly 1 | Streaming |
resolve_futures_ordered
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> resolve_futures_ordered() -> | exactly 1 | Streaming |
scan
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> scan(A, B) -> | exactly 1 | Streaming |
1 input stream, 1 output stream
Arguments: two arguments, both closures. The first closure is used to create the initial value for the accumulator, and the second is used to transform new items with the existing accumulator value. The second closure takes two arguments: an
&mut Accumaccumulated value, and anItem, and returns anOption<o>that will be emitted to the output stream if it'sSome, or terminate the stream if it'sNone.
Similar to Rust's standard library scan method. It applies a function to each element of the stream,
maintaining an internal state (accumulator) and emitting the values returned by the function.
The function can return None to terminate the stream early.
Note: The closures have access to the
contextobject.
scan can also be provided with one generic lifetime persistence argument, either
'tick or 'static, to specify how data persists. With 'tick, the accumulator will only be maintained
within the same tick. With 'static, the accumulated value will be remembered across ticks.
When not explicitly specified persistence defaults to 'tick.
// Running sum example
source_iter([1, 2, 3, 4])
-> scan::<'tick>(|| 0, |acc: &mut i32, x: i32| {
*acc += x;
Some(*acc)
})
-> assert_eq([1, 3, 6, 10]);
// Early termination example
source_iter([1, 2, 3, 4])
-> scan::<'tick>(|| 1, |state: &mut i32, x: i32| {
*state = *state * x;
if *state > 6 {
None
} else {
Some(-*state)
}
})
-> assert_eq([-1, -2, -6]);
scan_async_blocking
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> scan_async_blocking(A, B) -> | exactly 1 | Streaming |
sort
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> sort() -> | exactly 1 | Streaming |
sort_by_key
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> sort_by_key(A) -> | exactly 1 | Streaming |
source_file
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_file(A) -> | exactly 1 | Streaming |
source_interval
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_interval(A) -> | exactly 1 | Streaming |
source_iter
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_iter(A) -> | exactly 1 | Streaming |
source_json
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_json(A) -> | exactly 1 | Streaming |
source_stdin
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_stdin() -> | exactly 1 | Streaming |
source_stream
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_stream(A) -> | exactly 1 | Streaming |
0 input streams, 1 output stream
Arguments: The receive end of a tokio channel
Given a Stream
created in Rust code, source_stream
is passed the receive endpoint of the channel and emits each of the
elements it receives downstream.
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<&str>();
let mut flow = dfir_rs::dfir_syntax! {
source_stream(input_recv) -> map(|x| x.to_uppercase())
-> for_each(|x| println!("{}", x));
};
input_send.send("Hello").unwrap();
input_send.send("World").unwrap();
flow.run_available();
source_stream_serde
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | source_stream_serde(A) -> | exactly 1 | Streaming |
0 input streams, 1 output stream
Arguments:
Stream
Given a Stream
of (serialized payload, addr) pairs, deserializes the payload and emits each of the
elements it receives downstream.
async fn serde_in() {
let addr = dfir_rs::util::ipv4_resolve("localhost:9000".into()).unwrap();
let (outbound, inbound, _) = dfir_rs::util::bind_udp_bytes(addr).await;
let mut flow = dfir_rs::dfir_syntax! {
source_stream_serde(inbound) -> map(Result::unwrap) -> map(|(x, a): (String, std::net::SocketAddr)| x.to_uppercase())
-> for_each(|x| println!("{}", x));
};
flow.run_available();
}
spin
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 0 | spin() -> | exactly 1 | Streaming |
state
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> state()[<output_port>] -> | exactly 2 | Streaming |
Output port names:
items,state
state_by
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> state_by(A, B)[<output_port>] -> | exactly 2 | Streaming |
Output port names:
items,state
tee
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> tee()[<output_port>] -> | at least 2 | Streaming |
union
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| at least 2 | -> [<input_port>]union() -> | exactly 1 | Streaming |
unique
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> unique() -> | exactly 1 | Streaming |
Takes one stream as input and filters out any duplicate occurrences. The output contains all unique values from the input.
source_iter(vec![1, 1, 2, 3, 2, 1, 3])
-> unique()
-> assert_eq([1, 2, 3]);
unique can also be provided with one generic lifetime persistence argument, either
'tick or 'static, to specify how data persists. The default is 'tick.
With 'tick, uniqueness is only considered within the current tick, so across multiple ticks
duplicate values may be emitted.
With 'static, values will be remembered across ticks and no duplicates will ever be emitted.
let (input_send, input_recv) = dfir_rs::util::unbounded_channel::<usize>();
let mut flow = dfir_rs::dfir_syntax! {
source_stream(input_recv)
-> unique::<'tick>()
-> for_each(|n| println!("{}", n));
};
input_send.send(3).unwrap();
input_send.send(3).unwrap();
input_send.send(4).unwrap();
input_send.send(3).unwrap();
flow.run_available();
// 3, 4
input_send.send(3).unwrap();
input_send.send(5).unwrap();
flow.run_available();
// 3, 5
// Note: 3 is emitted again.
unzip
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> unzip()[<output_port>] -> | exactly 2 | Streaming |
Output port names:
0,1
zip
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]zip() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
2 input streams of type
V1andV2, 1 output stream of type(V1, V2)
Zips the streams together, forming paired tuples of the inputs. Note that zipping is done per-tick. If you do not
want to discard the excess, use zip_longest instead.
Takes in up to two generic lifetime persistence argument, one for each input. Within the lifetime, excess items
from one input or the other will be discarded. Using a 'static persistence lifetime may result in unbounded
buffering if the rates are mismatched.
source_iter(0..3) -> [0]my_zip;
source_iter(0..5) -> [1]my_zip;
my_zip = zip() -> assert_eq([(0, 0), (1, 1), (2, 2)]);
zip_longest
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]zip_longest() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)
_counter
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 1 | -> _counter(A, B) | at least 0 and at most 1 | Streaming |
_lattice_fold_batch
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]_lattice_fold_batch() -> | exactly 1 | Streaming |
Input port names:
input(streaming),signal(streaming)
2 input streams, 1 output stream, no arguments.
Batches streaming input and releases it downstream when a signal is delivered. This allows for buffering data and delivering it later while also folding it into a single lattice data structure.
This operator is similar to defer_signal in that it batches input and releases it when a signal is given. It is also similar to lattice_fold in that it folds the input into a single lattice.
So, _lattice_fold_batch is a combination of both defer_signal and lattice_fold. This operator is useful when trying to combine a sequence of defer_signal and lattice_fold operators without unnecessary memory consumption.
There are two inputs to _lattice_fold_batch, they are input and signal.
input is the input data flow. Data that is delivered on this input is collected in order inside of the _lattice_fold_batch operator.
When anything is sent to signal the collected data is released downstream. The entire signal input is consumed each tick, so sending 5 things on signal will not release inputs on the next 5 consecutive ticks.
use lattices::set_union::SetUnionHashSet;
use lattices::set_union::SetUnionSingletonSet;
source_iter([1, 2, 3])
-> map(SetUnionSingletonSet::new_from)
-> [input]batcher;
source_iter([()])
-> [signal]batcher;
batcher = _lattice_fold_batch::<SetUnionHashSet<usize>>()
-> assert_eq([SetUnionHashSet::new_from([1, 2, 3])]);
_lattice_join_fused_join
| Inputs | Syntax | Outputs | Flow |
|---|---|---|---|
| exactly 2 | -> [<input_port>]_lattice_join_fused_join() -> | exactly 1 | Streaming |
Input port names:
0(streaming),1(streaming)