Skip to main content

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​

InputsSyntaxOutputsFlow
exactly 1-> all_iterations() ->exactly 1Streaming

anti_join​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]anti_join() ->exactly 1Streaming

Input port names: pos (streaming), neg (streaming)

assert​

InputsSyntaxOutputsFlow
exactly 1-> assert(A)at least 0 and at most 1Streaming

assert_eq​

InputsSyntaxOutputsFlow
exactly 1-> assert_eq(A)at least 0 and at most 1Streaming

batch​

InputsSyntaxOutputsFlow
exactly 1-> batch() ->exactly 1Streaming

batch_eager​

InputsSyntaxOutputsFlow
exactly 1-> batch_eager() ->exactly 1Streaming

batch_lazy​

InputsSyntaxOutputsFlow
exactly 1-> batch_lazy() ->exactly 1Streaming

chain​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]chain() ->exactly 1Streaming

chain_first_n​

InputsSyntaxOutputsFlow
at least 2-> [<input_port>]chain_first_n(A) ->exactly 1Streaming

cross_join​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]cross_join() ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]cross_join_multiset() ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

cross_singleton​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]cross_singleton() ->exactly 1Streaming

Input port names: input (streaming), single (streaming)

defer_signal​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]defer_signal() ->exactly 1Streaming

Input port names: input (streaming), signal (streaming)

defer_tick​

InputsSyntaxOutputsFlow
exactly 1-> defer_tick() ->exactly 1Blocking

defer_tick_lazy​

InputsSyntaxOutputsFlow
exactly 1-> defer_tick_lazy() ->exactly 1Blocking

demux_enum​

InputsSyntaxOutputsFlow
exactly 1-> demux_enum()any number ofStreaming

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​

InputsSyntaxOutputsFlow
exactly 1-> dest_file(A, B)exactly 0Streaming

0 input streams, 1 output stream

Arguments: (1) An AsRef<Path> for a file to write to, and (2) a bool append.

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​

InputsSyntaxOutputsFlow
exactly 1-> dest_sink(A)exactly 0Streaming

dest_sink_serde​

InputsSyntaxOutputsFlow
exactly 1-> dest_sink_serde(A)exactly 0Streaming

difference​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]difference() ->exactly 1Streaming

Input port names: pos (streaming), neg (streaming)

enumerate​

InputsSyntaxOutputsFlow
exactly 1-> enumerate() ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 1-> filter(A) ->exactly 1Streaming

filter_map​

InputsSyntaxOutputsFlow
exactly 1-> filter_map(A) ->exactly 1Streaming

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 context object.

source_iter(vec!["1", "hello", "world", "2"])
-> filter_map(|s| s.parse::<usize>().ok())
-> assert_eq([1, 2]);

flat_map​

InputsSyntaxOutputsFlow
exactly 1-> flat_map(A) ->exactly 1Streaming

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 context object.

// 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​

InputsSyntaxOutputsFlow
exactly 1-> flat_map_stream_blocking(A) ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 1-> flatten() ->exactly 1Streaming

flatten_stream_blocking​

InputsSyntaxOutputsFlow
exactly 1-> flatten_stream_blocking() ->exactly 1Streaming

fold​

InputsSyntaxOutputsFlow
exactly 1-> fold(A, B)at least 0 and at most 1Streaming

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 Accum accumulated value, and an Item.

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 context object.

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​

InputsSyntaxOutputsFlow
exactly 1-> fold_keyed(A, B) ->exactly 1Streaming

fold_no_replay​

InputsSyntaxOutputsFlow
exactly 1-> fold_no_replay(A, B)at least 0 and at most 1Streaming

for_each​

InputsSyntaxOutputsFlow
exactly 1-> for_each(A)exactly 0Streaming

identity​

InputsSyntaxOutputsFlow
exactly 1-> identity() ->exactly 1Streaming

initialize​

InputsSyntaxOutputsFlow
exactly 0initialize() ->exactly 1Streaming

0 input streams, 1 output stream

Arguments: None.

Emits a single unit () at the start of the first tick.

initialize()
-> assert_eq([()]);

inspect​

InputsSyntaxOutputsFlow
exactly 1-> inspect(A)at least 0 and at most 1Streaming

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 context object.

source_iter([1, 2, 3, 4])
-> inspect(|x| println!("{}", x))
-> assert_eq([1, 2, 3, 4]);

iter_ref​

InputsSyntaxOutputsFlow
exactly 0iter_ref(A) ->exactly 1Streaming

join​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]join() ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]join_fused(A, B) ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

join_fused_lhs​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]join_fused_lhs(A) ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

join_fused_rhs​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]join_fused_rhs(A) ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

join_multiset​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]join_multiset() ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

join_multiset_half​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]join_multiset_half() ->exactly 1Streaming

Input port names: build (streaming), probe (streaming)

lattice_bimorphism​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]lattice_bimorphism(A, B, C) ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

lattice_fold​

InputsSyntaxOutputsFlow
exactly 1-> lattice_fold(A) ->exactly 1Streaming

lattice_reduce​

InputsSyntaxOutputsFlow
exactly 1-> lattice_reduce() ->exactly 1Streaming

map​

InputsSyntaxOutputsFlow
exactly 1-> map(A) ->exactly 1Streaming

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 context object.

source_iter(vec!["hello", "world"]) -> map(|x| x.to_uppercase())
-> assert_eq(["HELLO", "WORLD"]);

multiset_delta​

InputsSyntaxOutputsFlow
exactly 1-> multiset_delta() ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
at least 0 and at most 1null()at least 0 and at most 1Streaming

partition​

InputsSyntaxOutputsFlow
exactly 1-> partition(A)[<output_port>] ->at least 2Streaming

Output port names: Variadic, as specified in arguments.

persist​

InputsSyntaxOutputsFlow
exactly 1-> persist() ->exactly 1Streaming

reduce​

InputsSyntaxOutputsFlow
exactly 1-> reduce(A)at least 0 and at most 1Streaming

reduce_keyed​

InputsSyntaxOutputsFlow
exactly 1-> reduce_keyed(A) ->exactly 1Streaming

1 input stream of type (K, V), 1 output stream of type (K, V). The output will have one tuple for each distinct K, with an accumulated (reduced) value of type V.

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 context object.

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​

InputsSyntaxOutputsFlow
exactly 1-> reduce_no_replay(A)at least 0 and at most 1Streaming

resolve_futures​

InputsSyntaxOutputsFlow
exactly 1-> resolve_futures() ->exactly 1Streaming

resolve_futures_blocking​

InputsSyntaxOutputsFlow
exactly 1-> resolve_futures_blocking() ->exactly 1Streaming

resolve_futures_blocking_ordered​

InputsSyntaxOutputsFlow
exactly 1-> resolve_futures_blocking_ordered() ->exactly 1Streaming

resolve_futures_ordered​

InputsSyntaxOutputsFlow
exactly 1-> resolve_futures_ordered() ->exactly 1Streaming

scan​

InputsSyntaxOutputsFlow
exactly 1-> scan(A, B) ->exactly 1Streaming

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 Accum accumulated value, and an Item, and returns an Option<o> that will be emitted to the output stream if it's Some, or terminate the stream if it's None.

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 context object.

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​

InputsSyntaxOutputsFlow
exactly 1-> scan_async_blocking(A, B) ->exactly 1Streaming

sort​

InputsSyntaxOutputsFlow
exactly 1-> sort() ->exactly 1Streaming

sort_by_key​

InputsSyntaxOutputsFlow
exactly 1-> sort_by_key(A) ->exactly 1Streaming

source_file​

InputsSyntaxOutputsFlow
exactly 0source_file(A) ->exactly 1Streaming

source_interval​

InputsSyntaxOutputsFlow
exactly 0source_interval(A) ->exactly 1Streaming

source_iter​

InputsSyntaxOutputsFlow
exactly 0source_iter(A) ->exactly 1Streaming

source_json​

InputsSyntaxOutputsFlow
exactly 0source_json(A) ->exactly 1Streaming

source_stdin​

InputsSyntaxOutputsFlow
exactly 0source_stdin() ->exactly 1Streaming

source_stream​

InputsSyntaxOutputsFlow
exactly 0source_stream(A) ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 0source_stream_serde(A) ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 0spin() ->exactly 1Streaming

state​

InputsSyntaxOutputsFlow
exactly 1-> state()[<output_port>] ->exactly 2Streaming

Output port names: items, state

state_by​

InputsSyntaxOutputsFlow
exactly 1-> state_by(A, B)[<output_port>] ->exactly 2Streaming

Output port names: items, state

tee​

InputsSyntaxOutputsFlow
exactly 1-> tee()[<output_port>] ->at least 2Streaming

union​

InputsSyntaxOutputsFlow
at least 2-> [<input_port>]union() ->exactly 1Streaming

unique​

InputsSyntaxOutputsFlow
exactly 1-> unique() ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 1-> unzip()[<output_port>] ->exactly 2Streaming

Output port names: 0, 1

zip​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]zip() ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

2 input streams of type V1 and V2, 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​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]zip_longest() ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)

_counter​

InputsSyntaxOutputsFlow
exactly 1-> _counter(A, B)at least 0 and at most 1Streaming

_lattice_fold_batch​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]_lattice_fold_batch() ->exactly 1Streaming

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​

InputsSyntaxOutputsFlow
exactly 2-> [<input_port>]_lattice_join_fused_join() ->exactly 1Streaming

Input port names: 0 (streaming), 1 (streaming)