Skip to content

Commit eeb4824

Browse files
committed
refactor(dfir_lang): remove stratum, add push codegen, test sort, sort_by_key
PR: #2968
1 parent 3a91784 commit eeb4824

17 files changed

Lines changed: 156 additions & 219 deletions

dfir_lang/src/graph/ops/sort.rs

Lines changed: 18 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
use quote::quote_spanned;
22

33
use super::{
4-
DelayType, OperatorCategory, OperatorConstraints, OperatorWriteOutput, RANGE_0, RANGE_1,
4+
OperatorCategory, OperatorConstraints, OperatorWriteOutput, RANGE_0, RANGE_1,
55
WriteContextArgs,
66
};
77

@@ -12,9 +12,7 @@ use super::{
1212
/// -> sort()
1313
/// -> assert_eq([1, 2, 3]);
1414
/// ```
15-
///
16-
/// `sort` is blocking. Only the values collected within a single tick will be sorted and
17-
/// emitted.
15+
/// Within a tick, only the values received within that tick will be sorted and emitted.
1816
pub const SORT: OperatorConstraints = OperatorConstraints {
1917
name: "sort",
2018
categories: &[OperatorCategory::Persistence],
@@ -29,27 +27,32 @@ pub const SORT: OperatorConstraints = OperatorConstraints {
2927
flo_type: None,
3028
ports_inn: None,
3129
ports_out: None,
32-
input_delaytype_fn: |_| Some(DelayType::Stratum),
30+
input_delaytype_fn: |_| None,
3331
write_fn: |&WriteContextArgs {
3432
root,
3533
op_span,
3634
work_fn_async,
3735
ident,
3836
inputs,
37+
outputs,
3938
is_pull,
4039
..
4140
},
4241
_| {
43-
assert!(is_pull);
44-
45-
let input = &inputs[0];
46-
let write_iterator = quote_spanned! {op_span=>
47-
// TODO(mingwei): unnecessary extra handoff into_iter() then collect().
48-
let #ident = {
49-
let mut tmp = #work_fn_async(#root::dfir_pipes::pull::Pull::collect::<::std::vec::Vec<_>>(#input)).await;
50-
<[_]>::sort_unstable(&mut tmp);
51-
#root::dfir_pipes::pull::iter(tmp)
52-
};
42+
let write_iterator = if is_pull {
43+
let input = &inputs[0];
44+
quote_spanned! {op_span=>
45+
let #ident = {
46+
let mut tmp = #work_fn_async(#root::dfir_pipes::pull::Pull::collect::<::std::vec::Vec<_>>(#input)).await;
47+
<[_]>::sort_unstable(&mut tmp);
48+
#root::dfir_pipes::pull::iter(tmp)
49+
};
50+
}
51+
} else {
52+
let output = &outputs[0];
53+
quote_spanned! {op_span=>
54+
let #ident = #root::dfir_pipes::push::Sort::new(#output);
55+
}
5356
};
5457
Ok(OperatorWriteOutput {
5558
write_iterator,

dfir_lang/src/graph/ops/sort_by_key.rs

Lines changed: 30 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
use quote::quote_spanned;
22

33
use super::{
4-
DelayType, OperatorCategory, OperatorConstraints, OperatorWriteOutput, RANGE_0, RANGE_1,
4+
OperatorCategory, OperatorConstraints, OperatorWriteOutput, RANGE_0, RANGE_1,
55
WriteContextArgs,
66
};
77

@@ -29,27 +29,46 @@ pub const SORT_BY_KEY: OperatorConstraints = OperatorConstraints {
2929
flo_type: None,
3030
ports_inn: None,
3131
ports_out: None,
32-
input_delaytype_fn: |_| Some(DelayType::Stratum),
32+
input_delaytype_fn: |_| None,
3333
write_fn: |&WriteContextArgs {
3434
root,
3535
op_span,
3636
work_fn_async,
3737
ident,
3838
inputs,
39+
outputs,
3940
is_pull,
4041
arguments,
4142
..
4243
},
4344
_| {
44-
assert!(is_pull);
45-
let input = &inputs[0];
46-
let write_iterator = quote_spanned! {op_span=>
47-
// TODO(mingwei): unnecessary extra handoff into_iter() then collect().
48-
let #ident = {
49-
let mut tmp = #work_fn_async(#root::dfir_pipes::pull::Pull::collect::<::std::vec::Vec<_>>(#input)).await;
50-
#root::util::sort_unstable_by_key_hrtb(&mut tmp, #arguments);
51-
#root::dfir_pipes::pull::iter(tmp)
52-
};
45+
let write_iterator = if is_pull {
46+
let input = &inputs[0];
47+
quote_spanned! {op_span=>
48+
let #ident = {
49+
let mut tmp = #work_fn_async(#root::dfir_pipes::pull::Pull::collect::<::std::vec::Vec<_>>(#input)).await;
50+
#root::util::sort_unstable_by_key_hrtb(&mut tmp, #arguments);
51+
#root::dfir_pipes::pull::iter(tmp)
52+
};
53+
}
54+
} else {
55+
let output = &outputs[0];
56+
quote_spanned! {op_span=>
57+
let #ident = #root::dfir_pipes::push::fold(
58+
::std::vec::Vec::new(),
59+
|__buf: &mut ::std::vec::Vec<_>, __item| {
60+
__buf.push(__item);
61+
},
62+
#root::dfir_pipes::push::flat_map(
63+
|__buf: ::std::vec::Vec<_>| {
64+
let mut __buf = __buf;
65+
#root::util::sort_unstable_by_key_hrtb(&mut __buf, #arguments);
66+
__buf
67+
},
68+
#output,
69+
),
70+
);
71+
}
5372
};
5473
Ok(OperatorWriteOutput {
5574
write_iterator,

dfir_rs/tests/snapshots/surface_cross_singleton__union_defer_tick@graphvis_dot.snap

Lines changed: 23 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -24,28 +24,26 @@ digraph {
2424
n17v1 [label="(n17v1) handoff", shape=parallelogram, fillcolor="#ddddff"]
2525
n18v1 [label="(n18v1) handoff", shape=parallelogram, fillcolor="#ddddff"]
2626
n19v1 [label="(n19v1) handoff", shape=parallelogram, fillcolor="#ddddff"]
27-
n20v1 [label="(n20v1) handoff", shape=parallelogram, fillcolor="#ddddff"]
2827
n2v1 -> n3v1
29-
n1v1 -> n15v1
30-
n3v1 -> n16v1
28+
n1v1 -> n2v1
29+
n3v1 -> n15v1
3130
n4v1 -> n7v1 [label="0"]
32-
n14v1 -> n17v1
31+
n14v1 -> n16v1
3332
n5v1 -> n6v1
3433
n6v1 -> n7v1 [label="1"]
35-
n7v1 -> n18v1
34+
n7v1 -> n17v1
3635
n8v1 -> n9v1
3736
n9v1 -> n10v1
3837
n9v1 -> n11v1
39-
n3v1 -> n19v1
40-
n11v1 -> n20v1
38+
n3v1 -> n18v1
39+
n11v1 -> n19v1
4140
n13v1 -> n14v1
4241
n12v1 -> n13v1
43-
n15v1 -> n2v1 [color=red]
44-
n16v1 -> n8v1 [label="input"]
45-
n17v1 -> n4v1 [color=red]
46-
n18v1 -> n8v1 [label="single", color=red]
47-
n19v1 -> n12v1 [label="input"]
48-
n20v1 -> n12v1 [label="single", color=red]
42+
n15v1 -> n8v1 [label="input"]
43+
n16v1 -> n4v1 [color=red]
44+
n17v1 -> n8v1 [label="single", color=red]
45+
n18v1 -> n12v1 [label="input"]
46+
n19v1 -> n12v1 [label="single", color=red]
4947
subgraph sg_1v1 {
5048
cluster=true
5149
fillcolor="#dddddd"
@@ -55,68 +53,58 @@ digraph {
5553
cluster=true
5654
label="var teed_in"
5755
n1v1
58-
}
59-
}
60-
subgraph sg_2v1 {
61-
cluster=true
62-
fillcolor="#dddddd"
63-
style=filled
64-
label = "sg_2v1"
65-
subgraph sg_2v1_var_teed_in {
66-
cluster=true
67-
label="var teed_in"
6856
n2v1
6957
n3v1
7058
}
7159
}
72-
subgraph sg_3v1 {
60+
subgraph sg_2v1 {
7361
cluster=true
7462
fillcolor="#dddddd"
7563
style=filled
76-
label = "sg_3v1"
64+
label = "sg_2v1"
7765
n4v1
78-
subgraph sg_3v1_var_persisted_stream {
66+
subgraph sg_2v1_var_persisted_stream {
7967
cluster=true
8068
label="var persisted_stream"
8169
n5v1
8270
n6v1
8371
}
84-
subgraph sg_3v1_var_unioned_stream {
72+
subgraph sg_2v1_var_unioned_stream {
8573
cluster=true
8674
label="var unioned_stream"
8775
n7v1
8876
}
8977
}
90-
subgraph sg_4v1 {
78+
subgraph sg_3v1 {
9179
cluster=true
9280
fillcolor="#dddddd"
9381
style=filled
94-
label = "sg_4v1"
82+
label = "sg_3v1"
9583
n10v1
96-
subgraph sg_4v1_var_folded_thing {
84+
subgraph sg_3v1_var_folded_thing {
9785
cluster=true
9886
label="var folded_thing"
9987
n11v1
10088
}
101-
subgraph sg_4v1_var_join {
89+
subgraph sg_3v1_var_join {
10290
cluster=true
10391
label="var join"
10492
n8v1
10593
n9v1
10694
}
10795
}
108-
subgraph sg_5v1 {
96+
subgraph sg_4v1 {
10997
cluster=true
11098
fillcolor="#dddddd"
11199
style=filled
112-
label = "sg_5v1"
113-
subgraph sg_5v1_var_deferred_stream {
100+
label = "sg_4v1"
101+
subgraph sg_4v1_var_deferred_stream {
114102
cluster=true
115103
label="var deferred_stream"
116104
n13v1
117105
n14v1
118106
}
119-
subgraph sg_5v1_var_joined_folded {
107+
subgraph sg_4v1_var_joined_folded {
120108
cluster=true
121109
label="var joined_folded"
122110
n12v1

dfir_rs/tests/snapshots/surface_cross_singleton__union_defer_tick@graphvis_mermaid.snap

Lines changed: 20 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -27,65 +27,59 @@ linkStyle default stroke:#aaa
2727
17v1["(17v1) <code>handoff</code>"]:::otherClass
2828
18v1["(18v1) <code>handoff</code>"]:::otherClass
2929
19v1["(19v1) <code>handoff</code>"]:::otherClass
30-
20v1["(20v1) <code>handoff</code>"]:::otherClass
3130
2v1-->3v1
32-
1v1-->15v1
33-
3v1-->16v1
31+
1v1-->2v1
32+
3v1-->15v1
3433
4v1-->|0|7v1
35-
14v1-->17v1
34+
14v1-->16v1
3635
5v1-->6v1
3736
6v1-->|1|7v1
38-
7v1-->18v1
37+
7v1-->17v1
3938
8v1-->9v1
4039
9v1-->10v1
4140
9v1-->11v1
42-
3v1-->19v1
43-
11v1-->20v1
41+
3v1-->18v1
42+
11v1-->19v1
4443
13v1-->14v1
4544
12v1-->13v1
46-
15v1--x2v1; linkStyle 15 stroke:red
47-
16v1-->|input|8v1
48-
17v1--o4v1; linkStyle 17 stroke:red
49-
18v1--x|single|8v1; linkStyle 18 stroke:red
50-
19v1-->|input|12v1
51-
20v1--x|single|12v1; linkStyle 20 stroke:red
45+
15v1-->|input|8v1
46+
16v1--o4v1; linkStyle 16 stroke:red
47+
17v1--x|single|8v1; linkStyle 17 stroke:red
48+
18v1-->|input|12v1
49+
19v1--x|single|12v1; linkStyle 19 stroke:red
5250
subgraph sg_1v1 ["sg_1v1"]
5351
subgraph sg_1v1_var_teed_in ["var <tt>teed_in</tt>"]
5452
1v1
55-
end
56-
end
57-
subgraph sg_2v1 ["sg_2v1"]
58-
subgraph sg_2v1_var_teed_in ["var <tt>teed_in</tt>"]
5953
2v1
6054
3v1
6155
end
6256
end
63-
subgraph sg_3v1 ["sg_3v1"]
57+
subgraph sg_2v1 ["sg_2v1"]
6458
4v1
65-
subgraph sg_3v1_var_persisted_stream ["var <tt>persisted_stream</tt>"]
59+
subgraph sg_2v1_var_persisted_stream ["var <tt>persisted_stream</tt>"]
6660
5v1
6761
6v1
6862
end
69-
subgraph sg_3v1_var_unioned_stream ["var <tt>unioned_stream</tt>"]
63+
subgraph sg_2v1_var_unioned_stream ["var <tt>unioned_stream</tt>"]
7064
7v1
7165
end
7266
end
73-
subgraph sg_4v1 ["sg_4v1"]
67+
subgraph sg_3v1 ["sg_3v1"]
7468
10v1
75-
subgraph sg_4v1_var_folded_thing ["var <tt>folded_thing</tt>"]
69+
subgraph sg_3v1_var_folded_thing ["var <tt>folded_thing</tt>"]
7670
11v1
7771
end
78-
subgraph sg_4v1_var_join ["var <tt>join</tt>"]
72+
subgraph sg_3v1_var_join ["var <tt>join</tt>"]
7973
8v1
8074
9v1
8175
end
8276
end
83-
subgraph sg_5v1 ["sg_5v1"]
84-
subgraph sg_5v1_var_deferred_stream ["var <tt>deferred_stream</tt>"]
77+
subgraph sg_4v1 ["sg_4v1"]
78+
subgraph sg_4v1_var_deferred_stream ["var <tt>deferred_stream</tt>"]
8579
13v1
8680
14v1
8781
end
88-
subgraph sg_5v1_var_joined_folded ["var <tt>joined_folded</tt>"]
82+
subgraph sg_4v1_var_joined_folded ["var <tt>joined_folded</tt>"]
8983
12v1
9084
end
9185
end

dfir_rs/tests/snapshots/surface_difference__diff_multiset_static@graphvis_dot.snap

Lines changed: 5 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -13,15 +13,13 @@ digraph {
1313
n6v1 [label="(n6v1) tee()", shape=house, fillcolor="#ffff88"]
1414
n7v1 [label="(n7v1) for_each(|x| println!(\"neg: {:?}\", x))", shape=house, fillcolor="#ffff88"]
1515
n8v1 [label="(n8v1) handoff", shape=parallelogram, fillcolor="#ddddff"]
16-
n9v1 [label="(n9v1) handoff", shape=parallelogram, fillcolor="#ddddff"]
1716
n2v1 -> n3v1
18-
n1v1 -> n8v1
17+
n1v1 -> n2v1
1918
n4v1 -> n1v1 [label="pos"]
2019
n5v1 -> n6v1
21-
n6v1 -> n9v1
20+
n6v1 -> n8v1
2221
n6v1 -> n7v1
23-
n8v1 -> n2v1 [color=red]
24-
n9v1 -> n1v1 [label="neg", color=red]
22+
n8v1 -> n1v1 [label="neg", color=red]
2523
subgraph sg_1v1 {
2624
cluster=true
2725
fillcolor="#dddddd"
@@ -44,23 +42,13 @@ digraph {
4442
cluster=true
4543
label="var diff"
4644
n1v1
45+
n2v1
46+
n3v1
4747
}
4848
subgraph sg_2v1_var_poss {
4949
cluster=true
5050
label="var poss"
5151
n4v1
5252
}
5353
}
54-
subgraph sg_3v1 {
55-
cluster=true
56-
fillcolor="#dddddd"
57-
style=filled
58-
label = "sg_3v1"
59-
subgraph sg_3v1_var_diff {
60-
cluster=true
61-
label="var diff"
62-
n2v1
63-
n3v1
64-
}
65-
}
6654
}

0 commit comments

Comments
 (0)