From 7ca400d5b7d140854b8146404479f44855c41c02 Mon Sep 17 00:00:00 2001 From: Philippe Veber Date: Wed, 6 Jan 2021 22:51:19 +0100 Subject: [PATCH 1/2] added Stream.group --- streaming/Stream.ml | 21 ++++++++++++--------- streaming/Streaming.mli | 6 +++--- tests/Stream_tests.ml | 10 ++++++++++ 3 files changed, 25 insertions(+), 12 deletions(-) diff --git a/streaming/Stream.ml b/streaming/Stream.ml index 7dc81e6..93da674 100644 --- a/streaming/Stream.ml +++ b/streaming/Stream.ml @@ -396,20 +396,23 @@ let split ~by:pred self = in { stream } - - -let group ?equal:(_ =Pervasives.(=)) self = +let group ~break self = let stream (Sink k) = - let push r x = - k.push r x + let init () = (k.init (), None, empty) in + let push (r, maybe_previous, acc) x = + match maybe_previous with + | None -> (r, Some x, acc ++ single x) + | Some previous -> + if break previous x then + (k.push r acc, Some x, single x) + else (r, Some x, acc ++ single x) in - self.stream (Sink { k with push }) + let stop (r, _, acc) = k.stop (k.push r acc) in + let full (r, _, _) = k.full r in + self.stream (Sink { init; push; full; stop }) in { stream } - - - (* IO *) let of_file path = diff --git a/streaming/Streaming.mli b/streaming/Streaming.mli index 1151ee4..1262d97 100644 --- a/streaming/Streaming.mli +++ b/streaming/Streaming.mli @@ -1120,9 +1120,9 @@ module Stream : sig val partition : int -> 'a t -> 'a t t (** [partition n] partitions the stream into sub-streams of size [n]. *) - (* TODO *) - (* val group : by *) - + val group : break:('a -> 'a -> bool) -> 'a t -> 'a t t + (** [group ~break stream] splits [stream] each time two consecutive + elements [x] and [y] verify [break x y]. *) (* TODO: Add variants for splitting once. Consider renaming: divide. *) (* split, partition, divide, etc is too confusing. *) diff --git a/tests/Stream_tests.ml b/tests/Stream_tests.ml index cf0af9f..56806b7 100644 --- a/tests/Stream_tests.ml +++ b/tests/Stream_tests.ml @@ -315,6 +315,16 @@ let () = ~actual:S.(to_list (map to_list (partition 0 (0-<5)))); ]; + let t = T.test T.(list (list int)) ~verbose in + T.group "Stream.group" [ + t "empty" ~expected:[[]] + ~actual:S.(to_list (map to_list (group ~break:( <> ) empty))); + t "of_list" ~expected:[[1;1];[2];[3;3]] + ~actual:S.(to_list (map to_list (group ~break:( <> ) (of_list [1;1;2;3;3])))); + t "repeat,concat" ~expected:[[1;1];[2;2;2]] + ~actual:S.(to_list (map to_list (group ~break:( <> ) (concat (repeat ~times:2 1) (repeat ~times:3 2)) ))); + ]; + let t = T.test T.(list int) ~verbose in T.group "Stream.interpose" [ t "empty" ~expected:[] From 025757b2e796d0a5d64131e997f40da624851bed Mon Sep 17 00:00:00 2001 From: Philippe Veber Date: Wed, 6 Jan 2021 23:14:06 +0100 Subject: [PATCH 2/2] Stream.group: optimization --- streaming/Stream.ml | 20 ++++++++++---------- tests/Stream_tests.ml | 2 ++ 2 files changed, 12 insertions(+), 10 deletions(-) diff --git a/streaming/Stream.ml b/streaming/Stream.ml index 93da674..5b1e2ef 100644 --- a/streaming/Stream.ml +++ b/streaming/Stream.ml @@ -398,17 +398,17 @@ let split ~by:pred self = let group ~break self = let stream (Sink k) = - let init () = (k.init (), None, empty) in - let push (r, maybe_previous, acc) x = - match maybe_previous with - | None -> (r, Some x, acc ++ single x) - | Some previous -> - if break previous x then - (k.push r acc, Some x, single x) - else (r, Some x, acc ++ single x) + let send r xs = k.push r (of_list (List.rev xs)) in + let init () = (k.init (), []) in + let push (r, acc) x = + match acc with + | [] -> (r, [x]) + | h :: _ -> + if break h x then (send r acc, [x]) + else (r, x :: acc) in - let stop (r, _, acc) = k.stop (k.push r acc) in - let full (r, _, _) = k.full r in + let stop (r, acc) = k.stop (send r acc) in + let full (r, _) = k.full r in self.stream (Sink { init; push; full; stop }) in { stream } diff --git a/tests/Stream_tests.ml b/tests/Stream_tests.ml index 56806b7..30cd7e1 100644 --- a/tests/Stream_tests.ml +++ b/tests/Stream_tests.ml @@ -321,6 +321,8 @@ let () = ~actual:S.(to_list (map to_list (group ~break:( <> ) empty))); t "of_list" ~expected:[[1;1];[2];[3;3]] ~actual:S.(to_list (map to_list (group ~break:( <> ) (of_list [1;1;2;3;3])))); + t "of_list, mod" ~expected:[[1;3];[2];[1;3]] + ~actual:S.(to_list (map to_list (group ~break:(fun x y -> x mod 2 <> y mod 2) (of_list [1;3;2;1;3])))); t "repeat,concat" ~expected:[[1;1];[2;2;2]] ~actual:S.(to_list (map to_list (group ~break:( <> ) (concat (repeat ~times:2 1) (repeat ~times:3 2)) ))); ];