diff --git a/streaming/Stream.ml b/streaming/Stream.ml index 7dc81e6..5b1e2ef 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 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 - self.stream (Sink { k with push }) + 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 } - - - (* 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..30cd7e1 100644 --- a/tests/Stream_tests.ml +++ b/tests/Stream_tests.ml @@ -315,6 +315,18 @@ 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 "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)) ))); + ]; + let t = T.test T.(list int) ~verbose in T.group "Stream.interpose" [ t "empty" ~expected:[]