MCPcopy Create free account
hub / github.com/comnik/declarative-dataflow / run_cases

Function run_cases

tests/or_test.rs:37–124  ·  view source on GitHub ↗
(mut cases: Vec<Case>)

Source from the content-addressed store, hash-verified

35}
36
37fn run_cases(mut cases: Vec<Case>) {
38 for case in cases.drain(..) {
39 timely::execute_directly(move |worker| {
40 let mut server = Server::<u64, u64>::new(Default::default());
41 let (send_results, results) = channel();
42
43 dbg!(case.description);
44
45 let mut deps = dependencies(&case);
46 let plan = case.plan.clone();
47
48 for tx in case.transactions.iter() {
49 for datum in tx {
50 deps.insert(datum.2.clone());
51 }
52 }
53
54 worker.dataflow::<u64, _, _>(|scope| {
55 for dep in deps.iter() {
56 let config = AttributeConfig {
57 trace_slack: Some(Time::TxId(1)),
58 query_support: QuerySupport::AdaptiveWCO,
59 index_direction: IndexDirection::Both,
60 ..Default::default()
61 };
62
63 server
64 .context
65 .internal
66 .create_transactable_attribute(dep, config, scope)
67 .unwrap();
68 }
69
70 server
71 .test_single(
72 scope,
73 Rule {
74 name: "query".to_string(),
75 plan,
76 },
77 )
78 .inner
79 .sink(Pipeline, "Results", move |input| {
80 input.for_each(|_time, data| {
81 for datum in data.iter() {
82 send_results.send(datum.clone()).unwrap()
83 }
84 });
85 });
86 });
87
88 let mut transactions = case.transactions.clone();
89 let mut next_tx = 0;
90
91 for (tx_id, tx_data) in transactions.drain(..).enumerate() {
92 next_tx += 1;
93
94 server.transact(tx_data, 0, 0).unwrap();

Callers 2

orFunction · 0.70
or_joinFunction · 0.70

Calls 7

test_singleMethod · 0.80
advance_domainMethod · 0.80
is_any_outdatedMethod · 0.80
dependenciesFunction · 0.70
sinkMethod · 0.45
transactMethod · 0.45

Tested by

no test coverage detected