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

Function run_cases

tests/pull_test.rs:23–115  ·  view source on GitHub ↗
(mut cases: Vec<Case>)

Source from the content-addressed store, hash-verified

21}
22
23fn run_cases(mut cases: Vec<Case>) {
24 for case in cases.drain(..) {
25 timely::execute_directly(move |worker| {
26 let mut server = Server::<u64, u64>::new(Default::default());
27 let (send_results, results) = channel();
28
29 dbg!(case.description);
30
31 let mut deps = case.plan.dependencies();
32 let plan = case.plan.clone();
33
34 dbg!(&plan);
35
36 for tx in case.transactions.iter() {
37 for datum in tx {
38 deps.attributes.insert(datum.2.clone());
39 }
40 }
41
42 worker.dataflow::<u64, _, _>(|scope| {
43 for dep in deps.attributes.iter() {
44 let config = AttributeConfig {
45 trace_slack: Some(Time::TxId(1)),
46 // @TODO Forward delta should be enough eventually
47 query_support: QuerySupport::AdaptiveWCO,
48 index_direction: IndexDirection::Both,
49 // query_support: QuerySupport::Delta,
50 // index_direction: IndexDirection::Forward,
51 ..Default::default()
52 };
53
54 server
55 .context
56 .internal
57 .create_transactable_attribute(dep, config, scope)
58 .unwrap();
59 }
60
61 server
62 .test_single(
63 scope,
64 Rule {
65 name: "query".to_string(),
66 plan,
67 },
68 )
69 .inner
70 .sink(Pipeline, "Results", move |input| {
71 input.for_each(|_time, data| {
72 for datum in data.iter() {
73 send_results.send(datum.clone()).unwrap()
74 }
75 });
76 });
77 });
78
79 let mut transactions = case.transactions.clone();
80 let mut next_tx = 0;

Callers 2

pull_levelFunction · 0.70
graph_qlFunction · 0.70

Calls 7

test_singleMethod · 0.80
advance_domainMethod · 0.80
is_any_outdatedMethod · 0.80
dependenciesMethod · 0.45
sinkMethod · 0.45
transactMethod · 0.45

Tested by

no test coverage detected