Incrementalize sum and count aggregates

This commit is contained in:
sneeker committed 2017-04-20 08:51:00 +00:00
1 parent 19c440f03a
commit b0106d9a21
2 files changed
+127

No files matched your search

+126
View File
@@ -535,3 +535,129 @@ let filter_cases =
(Value.to_string applied); (Value.to_string applied);
check_equal_string "cached" "(collection (2 \"Bo\"))" (Value.to_string cached)) ); check_equal_string "cached" "(collection (2 \"Bo\"))" (Value.to_string cached)) );
] ]
let aggregate_cases =
let line price quantity =
Value.VRecord ("line", [ ("price", Value.VInt price); ("quantity", Value.VInt quantity) ])
in
let line_fixture query entries =
build_fixture
("type line = { price : int; quantity : int }\ninput lines : collection line\nquery q = " ^ query)
entries
in
[
( "sum adds insertions and subtracts removals",
fun () ->
let f = line_fixture "lines |> map (fun l -> l.price) |> sum" [ (1, line 10 1) ] in
(match step f [ Change.OpInsert (2, line 25 1) ] with
| Error message -> fail "insert" message
| Ok (applied, _, change) ->
check_equal_string "insert change" "+25" (Change.to_string change);
check_equal_string "applied" "35" (Value.to_string applied));
(match step f [ Change.OpRemove 1 ] with
| Error message -> fail "remove" message
| Ok (applied, _, change) ->
check_equal_string "remove change" "-10" (Change.to_string change);
check_equal_string "applied" "25" (Value.to_string applied)) );
( "a replacement applies the new contribution minus the old one",
fun () ->
let f = line_fixture "lines |> map (fun l -> l.price) |> sum" [ (1, line 10 1) ] in
(match step f [ Change.OpReplace (1, line 30 1) ] with
| Error message -> fail "replace" message
| Ok (applied, _, change) ->
check_equal_string "change" "+20" (Change.to_string change);
check_equal_string "applied" "30" (Value.to_string applied)) );
( "simultaneous operand changes include the cross term",
fun () ->
let f = line_fixture "lines |> map (fun l -> l.price * l.quantity) |> sum" [ (1, line 10 3) ] in
(match step f [ Change.OpReplace (1, line 20 5) ] with
| Error message -> fail "cross term" message
| Ok (applied, cached, change) ->
check_equal_string "change" "+70" (Change.to_string change);
check_equal_string "applied" "100" (Value.to_string applied);
check_equal_string "cached" "100" (Value.to_string cached);
check_equal_string "reference" (Value.to_string (reference_result f f.fx_state))
(Value.to_string applied)) );
( "the cross term is exact for negative operand changes",
fun () ->
let f = line_fixture "lines |> map (fun l -> l.price * l.quantity) |> sum" [ (1, line 10 3) ] in
(match step f [ Change.OpReplace (1, line 7 2) ] with
| Error message -> fail "negative" message
| Ok (applied, _, change) ->
check_equal_string "change" "-16" (Change.to_string change);
check_equal_string "applied" "14" (Value.to_string applied) );
(match step f [ Change.OpReplace (1, line (-4) 9) ] with
| Error message -> fail "mixed" message
| Ok (applied, _, change) ->
check_equal_string "change" "-50" (Change.to_string change);
check_equal_string "applied" "-36" (Value.to_string applied) ) );
( "a branch switch produces an additive change",
fun () ->
let f =
line_fixture "lines |> map (fun l -> if l.quantity > 0 then l.price else 0 - l.price) |> sum"
[ (1, line 10 3) ]
in
(match step f [ Change.OpReplace (1, line 10 (-3)) ] with
| Error message -> fail "branch" message
| Ok (applied, _, change) ->
check_equal_string "change" "-20" (Change.to_string change);
check_equal_string "applied" "-10" (Value.to_string applied)) );
( "division recomputes locally",
fun () ->
let f = line_fixture "lines |> map (fun l -> l.price / 10) |> sum" [ (1, line 100 1) ] in
(match step f [ Change.OpReplace (1, line 95 1) ] with
| Error message -> fail "division" message
| Ok (applied, _, change) ->
check_equal_string "change" "-1" (Change.to_string change);
check_equal_string "applied" "9" (Value.to_string applied)) );
( "division by zero during a replacement fails the batch",
fun () ->
let f = line_fixture "lines |> map (fun l -> l.price / l.quantity) |> sum" [ (1, line 100 2) ] in
let before = Value.to_string (Incremental.result f.fx_plan f.fx_state) in
(match step f [ Change.OpReplace (1, line 100 0) ] with
| Ok _ -> fail "division" "expected a failure"
| Error message -> check "mentions division" (String.length message > 0));
check_equal_string "state unchanged" before
(Value.to_string (Incremental.result f.fx_plan f.fx_state)) );
( "count only tracks membership",
fun () ->
let f = line_fixture "lines |> count" [ (1, line 10 1); (2, line 20 1) ] in
(match step f [ Change.OpReplace (2, line 99 9) ] with
| Error message -> fail "count replace" message
| Ok (_, cached, change) ->
check "no change" (Change.is_empty change);
check_equal_string "cached" "2" (Value.to_string cached));
(match step f [ Change.OpInsert (3, line 5 1); Change.OpRemove 1 ] with
| Error message -> fail "count insert remove" message
| Ok (applied, cached, change) ->
check_equal_string "change" "+0" (Change.to_string change);
check_equal_string "applied" "2" (Value.to_string applied);
check_equal_string "cached" "2" (Value.to_string cached)) );
( "aggregate updates track the reference over a run of batches",
fun () ->
let f =
line_fixture "lines |> filter (fun l -> l.quantity > 0) |> map (fun l -> l.price * l.quantity) |> sum"
[ (1, line 10 2); (2, line 5 0); (3, line 7 4) ]
in
let batches =
[
[ Change.OpReplace (1, line 11 3) ];
[ Change.OpInsert (4, line 2 2) ];
[ Change.OpReplace (2, line 9 1) ];
[ Change.OpRemove 3 ];
[ Change.OpReplace (4, line 0 5) ];
[ Change.OpRemove 1; Change.OpInsert (5, line 3 3) ];
]
in
List.iter
(fun ops ->
match step f ops with
| Error message -> fail "aggregate run" message
| Ok (applied, cached, _) ->
let expected = reference_result f f.fx_state in
if not (Value.equal applied expected) then
fail "applied matches the reference" (Printf.sprintf "ops=%d" (List.length ops));
if not (Value.equal cached expected) then
fail "cached matches the reference" (Printf.sprintf "ops=%d" (List.length ops)))
batches );
]
+1
View File
@@ -382,5 +382,6 @@ let () =
Test_harness.run_suite "graph" Test_incremental.graph_cases; Test_harness.run_suite "graph" Test_incremental.graph_cases;
Test_harness.run_suite "executor" Test_incremental.executor_cases; Test_harness.run_suite "executor" Test_incremental.executor_cases;
Test_harness.run_suite "filter" Test_incremental.filter_cases; Test_harness.run_suite "filter" Test_incremental.filter_cases;
Test_harness.run_suite "aggregates" Test_incremental.aggregate_cases;
Printf.printf "%d cases, %d failures\n" (Test_harness.case_count ()) (Test_harness.failure_count ()); Printf.printf "%d cases, %d failures\n" (Test_harness.case_count ()) (Test_harness.failure_count ());
exit (if Test_harness.failure_count () = 0 then 0 else 1) exit (if Test_harness.failure_count () = 0 then 0 else 1)