HORIZON HASKELLDocslts/ghc-9.10.xc74966e2026-09-27Search names, modules, packages, or :: a typeCtrl K

GHC 9.10.3 · lts/ghc-9.10.x · c74966e · 2026-09-27

Modulestreamly-core-0.2.2Haskell2010

Streamly.Internal.Data.Stream

Direct style re-implementation of CPS stream in Streamly.Internal.Data.StreamK. The symbol or suffix D in this module denotes the Direct style. GHC is able to INLINE and fuse direct style better, providing better performance than CPS implementation.

import qualified Streamly.Internal.Data.Stream as D
  • 9 types
  • 1 class
  • 345 values
datadata Step s a
#

A stream is a succession of Steps. A Yield produces a single value and the next state of the stream. Stop indicates there are no more values in the stream.

Constructors

Instances1Functor
  • Functor (Step s)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Step
datadata Stream (m :: Type -> Type) a
#

A stream consists of a step function that generates the next step given a current state, and the current state.

Constructors

Instances9Functor, Foldable, IsList, Eq, Ord, Read, …
valueconcat :: Monad m => Stream m (Stream m a) -> Stream m a
#

Flatten a stream of streams to a single stream.

Example1 expression
concat = Stream.concatMap id

Pre-release

valuefromPure :: Applicative m => a -> Stream m a
#

Create a singleton stream from a pure value.

Example3 expressions
fromPure a = a `Stream.cons` Stream.nilfromPure = purefromPure = Stream.fromEffect . pure
valuefoldr :: Monad m => (a -> b -> b) -> b -> Stream m a -> m b
#

Right fold, lazy for lazy monads and pure streams, and strict for strict monads.

Please avoid using this routine in strict monads like IO unless you need a strict right fold. This is provided only for use in lazy monads (e.g. Identity) or pure streams. Note that with this signature it is not possible to implement a lazy foldr when the monad m is strict. In that case it would be strict in its accumulator and therefore would necessarily consume all its input.

Example1 expression
foldr f z = Stream.foldrM (\a b -> f a <$> b) (return z)

Note: This is similar to Fold.foldr' (the right fold via left fold), but could be more efficient.

valuetoList :: Monad m => Stream m a -> m [a]
#

Definitions:

Example2 expressions
toList = Stream.foldr (:) []toList = Stream.fold Fold.toList

Convert a stream into a list in the underlying monad. The list can be consumed lazily in a lazy monad (e.g. Identity). In a strict monad (e.g. IO) the whole list is generated and buffered before it can be consumed.

Warning! working on large lists accumulated as buffers in memory could be very inefficient, consider using Streamly.Data.Array instead.

Note that this could a bit more efficient compared to Stream.fold Fold.toList, and it can fuse with pure list consumers.

valuefold :: Monad m => Fold m a b -> Stream m a -> m b
#

Fold a stream using the supplied left Fold and reducing the resulting expression strictly at each step. The behavior is similar to foldl'. A Fold can terminate early without consuming the full stream. See the documentation of individual Folds for termination behavior.

Definitions:

Example2 expressions
fold f = fmap fst . Stream.foldBreak ffold f = Stream.parse (Parser.fromFold f)

Example:

Example1 expression
Stream.fold Fold.sum (Stream.enumerateFromTo 1 100)5050
valuefoldBreak :: Monad m => Fold m a b -> Stream m a -> m (b, Stream m a)
#

Like fold but also returns the remaining stream. The resulting stream would be Stream.nil if the stream finished before the fold.

valuefoldEither
  1. :: Monad m
  2. => Fold m a b
  3. -> Stream m a
  4. -> m (Either (Fold m a b) (b, Stream m a))
#

Fold resulting in either breaking the stream or continuation of the fold. Instead of supplying the input stream in one go we can run the fold multiple times, each time supplying the next segment of the input stream. If the fold has not yet finished it returns a fold that can be run again otherwise it returns the fold result and the residual stream.

Internal

valuemapM :: Monad m => (a -> m b) -> Stream m a -> Stream m b
#
Example1 expression
mapM f = Stream.sequence . fmap f

Apply a monadic function to each element of the stream and replace it with the output of the resulting action.

Example2 expressions
s = Stream.fromList ["a", "b", "c"]Stream.fold Fold.drain $ Stream.mapM putStr sabc
valuezipWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m c
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.zipWith otherwise.

Stream a is evaluated first, followed by stream b, the resulting elements a and b are then zipped using the supplied zip function and the result c is yielded to the consumer.

If stream a or stream b ends, the zipped stream ends. If stream b ends first, the element a from previous evaluation of stream a is discarded.

Example3 expressions
s1 = Stream.fromList [1,2,3]s2 = Stream.fromList [4,5,6]Stream.fold Fold.toList $ Stream.zipWith (+) s1 s2[5,7,9]
valueconcatMap :: Monad m => (a -> Stream m b) -> Stream m a -> Stream m b
#

Map a stream producing function on each element of the stream and then flatten the results into a single stream.

Example3 expressions
concatMap f = Stream.concatMapM (return . f)concatMap f = Stream.concat . fmap fconcatMap f = Stream.unfoldMany (Unfold.lmap f Unfold.fromStream)

See unfoldMany for a fusible alternative.

valuenilM :: Applicative m => m b -> Stream m a
#

A stream that terminates without producing any output, but produces a side effect.

Example1 expression
Stream.fold Fold.toList (Stream.nilM (print "nil"))"nil"[]

Pre-release

valuefoldMany :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Apply a Fold repeatedly on a stream and emit the results in the output stream.

Definition:

Example1 expression
foldMany f = Stream.parseMany (Parser.fromFold f)

Example, empty stream:

Example3 expressions
f = Fold.take 2 Fold.sumfmany = Stream.fold Fold.toList . Stream.foldMany ffmany $ Stream.fromList [][]

Example, last fold empty:

Example1 expression
fmany $ Stream.fromList [1..4][3,7]

Example, last fold non-empty:

Example1 expression
fmany $ Stream.fromList [1..5][3,7,5]

Note that using a closed fold e.g. Fold.take 0, would result in an infinite stream on a non-empty input stream.

valuefromEffect :: Applicative m => m a -> Stream m a
#

Create a singleton stream from a monadic action.

Example2 expressions
fromEffect m = m `Stream.consM` Stream.nilfromEffect = Stream.sequence . Stream.fromPure
Example1 expression
Stream.fold Fold.drain $ Stream.fromEffect (putStrLn "hello")hello
valuegroupsOf :: Monad m => Int -> Fold m a b -> Stream m a -> Stream m b
#

Group the input stream into groups of n elements each and then fold each group using the provided fold function.

groupsOf n f = foldMany (FL.take n f)
Example1 expression
Stream.toList $ Stream.groupsOf 2 Fold.sum (Stream.enumerateFromTo 1 10)[3,7,11,15,19]

This can be considered as an n-fold version of take where we apply take repeatedly on the leftover stream until the stream exhausts.

valuedrain :: Monad m => Stream m a -> m ()
#

Definitions:

Example2 expressions
drain = Stream.fold Fold.draindrain = Stream.foldrM (\_ xs -> xs) (return ())

Run a stream, discarding the results.

valueunfold :: Applicative m => Unfold m a b -> a -> Stream m b
#

Convert an Unfold into a stream by supplying it an input seed.

Example2 expressions
s = Stream.unfold Unfold.replicateM (3, putStrLn "hello")Stream.fold Fold.drain shellohellohello
valueuncons :: Monad m => Stream m a -> m (Maybe (a, Stream m a))
#

Decompose a stream into its head and tail. If the stream is empty, returns Nothing. If the stream is non-empty, returns Just (a, ma), where a is the head of the stream and ma its tail.

Properties:

Example2 expressions
Nothing <- Stream.uncons Stream.nilJust ("a", t) <- Stream.uncons (Stream.cons "a" Stream.nil)

This can be used to consume the stream in an imperative manner one element at a time, as it just breaks down the stream into individual elements and we can loop over them as we deem fit. For example, this can be used to convert a streamly stream into other stream types.

All the folds in this module can be expressed in terms of uncons, however, this is generally less efficient than specific folds because it takes apart the stream one element at a time, therefore, does not take adavantage of stream fusion.

foldBreak is a more general way of consuming a stream piecemeal.

Example1 expression
:{uncons xs = do    r <- Stream.foldBreak Fold.one xs    return $ case r of        (Nothing, _) -> Nothing        (Just h, t) -> Just (h, t):}
valuefoldrM :: Monad m => (a -> m b -> m b) -> m b -> Stream m a -> m b
#

Right associative/lazy pull fold. foldrM build final stream constructs an output structure using the step function build. build is invoked with the next input element and the remaining (lazy) tail of the output structure. It builds a lazy output expression using the two. When the "tail structure" in the output expression is evaluated it calls build again thus lazily consuming the input stream until either the output expression built by build is free of the "tail" or the input is exhausted in which case final is used as the terminating case for the output structure. For more details see the description in the previous section.

Example, determine if any element is odd in a stream:

Example3 expressions
s = Stream.fromList (2:4:5:undefined)step x xs = if odd x then return True else xsStream.foldrM step (return False) sTrue
valuecrossWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m c
#

Definition:

Example1 expression
crossWith f m1 m2 = fmap f m1 `Stream.crossApply` m2

Note that the second stream is evaluated multiple times.

valueunfoldMany :: Monad m => Unfold m a b -> Stream m a -> Stream m b
#

unfoldMany unfold stream uses unfold to map the input stream elements to streams and then flattens the generated streams into a single output stream.

Like concatMap but uses an Unfold for stream generation. Unlike concatMap this can fuse the Unfold code with the inner loop and therefore provide many times better performance.

valueconcatEffect :: Monad m => m (Stream m a) -> Stream m a
#

Given a stream value in the underlying monad, lift and join the underlying monad with the stream monad.

Example2 expressions
concatEffect = Stream.concat . Stream.fromEffectconcatEffect eff = Stream.concatMapM (\() -> eff) (Stream.fromPure ())

See also: concat, sequence

valueconcatMapM :: Monad m => (a -> m (Stream m b)) -> Stream m a -> Stream m b
#

Map a stream producing monadic function on each element of the stream and then flatten the results into a single stream. Since the stream generation function is monadic, unlike concatMap, it can produce an effect at the beginning of each iteration of the inner loop.

See unfoldMany for a fusible alternative.

valuecross :: Monad m => Stream m a -> Stream m b -> Stream m (a, b)
#

Given a Stream m a and Stream m b generate a stream with all possible combinations of the tuple (a, b).

Definition:

Example1 expression
cross = Stream.crossWith (,)

The second stream is evaluated multiple times. If that is not desired it can be cached in an Array and then generated from the array before calling this function. Caching may also improve performance if the stream is expensive to evaluate.

See cross for a much faster fused alternative.

Time: O(m x n)

Pre-release

newtypenewtype CrossStream (m :: Type -> Type) a
#

A newtype wrapper for the Stream type with a cross product style monad instance.

A Monad bind behaves like a for loop:

Example1 expression
:{Stream.fold Fold.toList $ Stream.unCross $ do    x <- Stream.mkCross $ Stream.fromList [1,2]    -- Perform the following actions for each x in the stream    return x:}[1,2]

Nested monad binds behave like nested for loops:

Example1 expression
:{Stream.fold Fold.toList $ Stream.unCross $ do    x <- Stream.mkCross $ Stream.fromList [1,2]    y <- Stream.mkCross $ Stream.fromList [3,4]    -- Perform the following actions for each x, for each y    return (x, y):}[(1,3),(1,4),(2,3),(2,4)]
Instances14MonadTrans, Monad, Functor, Applicative, Foldable, MonadIO, …
valuefoldAddLazy :: Monad m => Fold m a b -> Stream m a -> Fold m a b
#

Append a stream to a fold lazily to build an accumulator incrementally.

Example, to continue folding a list of streams on the same sum fold:

Example3 expressions
streams = [Stream.fromList [1..5], Stream.fromList [6..10]]f = Prelude.foldl Stream.foldAddLazy Fold.sum streamsStream.fold f Stream.nil55
valuecrossApply :: Functor f => Stream f (a -> b) -> Stream f a -> Stream f b
#

Apply a stream of functions to a stream of values and flatten the results.

Note that the second stream is evaluated multiple times.

Example1 expression
crossApply = Stream.crossWith id
valueunfoldIterateDfs :: Monad m => Unfold m a a -> Stream m a -> Stream m a
#

Same as concatIterateDfs but more efficient due to stream fusion.

Example, list a directory tree using DFS:

Example3 expressions
f = Unfold.either Dir.eitherReaderPaths Unfold.nilinput = Stream.fromPure (Left ".")ls = Stream.unfoldIterateDfs f input

Pre-release

valueconcatIterateScan
  1. :: Monad m
  2. => b -> a -> m b
  3. -> b -> m (Maybe (b, Stream m a))
  4. -> b
  5. -> Stream m a
#

Generate a stream from an initial state, scan and concat the stream, generate a stream again from the final state of the previous scan and repeat the process.

valueconcatIterateDfs
  1. :: Monad m
  2. => a -> Maybe (Stream m a)
  3. -> Stream m a
  4. -> Stream m a
#

Traverse the stream in depth first style (DFS). Map each element in the input stream to a stream and flatten, recursively map the resulting elements as well to a stream and flatten until no more streams are generated.

Example, list a directory tree using DFS:

Example3 expressions
f = either (Just . Dir.readEitherPaths) (const Nothing)input = Stream.fromPure (Left ".")ls = Stream.concatIterateDfs f input

This is equivalent to using concatIterateWith StreamK.append.

Pre-release

valueconcatIterateBfs
  1. :: Monad m
  2. => a -> Maybe (Stream m a)
  3. -> Stream m a
  4. -> Stream m a
#

Similar to concatIterateDfs except that it traverses the stream in breadth first style (BFS). First, all the elements in the input stream are emitted, and then their traversals are emitted.

Example, list a directory tree using BFS:

Example3 expressions
f = either (Just . Dir.readEitherPaths) (const Nothing)input = Stream.fromPure (Left ".")ls = Stream.concatIterateBfs f input

Pre-release

valuefoldManyPost :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Like foldMany but evaluates the fold even if the fold did not receive any input, therefore, always results in a non-empty output even on an empty stream (default result of the fold).

Example, empty stream:

Example3 expressions
f = Fold.take 2 Fold.sumfmany = Stream.fold Fold.toList . Stream.foldManyPost ffmany $ Stream.fromList [][0]

Example, last fold empty:

Example1 expression
fmany $ Stream.fromList [1..4][3,7,0]

Example, last fold non-empty:

Example1 expression
fmany $ Stream.fromList [1..5][3,7,5]

Note that using a closed fold e.g. Fold.take 0, would result in an infinite stream without consuming the input.

Pre-release

valuereduceIterateBfs :: Monad m => (a -> a -> m a) -> Stream m a -> m (Maybe a)
#

Binary BFS style reduce, folds a level entirely using the supplied fold function, collecting the outputs as next level of the tree, then repeats the same process on the next level. The last elements of a previously folded level are folded first.

valuefoldIterateBfs :: Fold m a (Either a a) -> Stream m a -> m (Maybe a)
#

N-Ary BFS style iterative fold, if the input stream finished before the fold then it returns Left otherwise Right. If the fold returns Left we terminate.

Unimplemented

valueindexOnSuffix :: Monad m => (a -> Bool) -> Stream m a -> Stream m (Int, Int)
#

Like splitOnSuffix but generates a stream of (index, len) tuples marking the places where the predicate matches in the stream.

Pre-release

valueiterate :: Monad m => (a -> a) -> a -> Stream m a
#

Generate an infinite stream with x as the first element and each successive element derived by applying the function f on the previous element.

Example1 expression
Stream.toList $ Stream.take 5 $ Stream.iterate (+1) 1[1,2,3,4,5]
valuefromPure :: Applicative m => a -> Stream m a
#

Create a singleton stream from a pure value.

Example3 expressions
fromPure a = a `Stream.cons` Stream.nilfromPure = purefromPure = Stream.fromEffect . pure
valuenil :: Applicative m => Stream m a
#

A stream that terminates without producing any output or side effect.

Example1 expression
Stream.toList Stream.nil[]
valuerepeatM :: Monad m => m a -> Stream m a
#
Example1 expression
repeatM = Stream.sequence . Stream.repeat

Generate a stream by repeatedly executing a monadic action forever.

Example1 expression
:{repeatAction =       Stream.repeatM (threadDelay 1000000 >> print 1)     & Stream.take 10     & Stream.fold Fold.drain:}
valuereplicate :: Monad m => Int -> a -> Stream m a
#
Example2 expressions
replicate n = Stream.take n . Stream.repeatreplicate n x = Stream.replicateM n (pure x)

Generate a stream of length n by repeating a value n times.

valuereplicateM :: Monad m => Int -> m a -> Stream m a
#
Example1 expression
replicateM n = Stream.sequence . Stream.replicate n

Generate a stream by performing a monadic action n times.

valueiterateM :: Monad m => (a -> m a) -> m a -> Stream m a
#

Generate an infinite stream with the first element generated by the action m and each successive element derived by applying the monadic function f on the previous element.

Example1 expression
:{Stream.iterateM (\x -> print x >> return (x + 1)) (return 0)    & Stream.take 3    & Stream.toList:}01[0,1,2]
valuerepeat :: Monad m => a -> Stream m a
#

Generate an infinite stream by repeating a pure value.

Example1 expression
repeat x = Stream.repeatM (pure x)
valuenilM :: Applicative m => m b -> Stream m a
#

A stream that terminates without producing any output, but produces a side effect.

Example1 expression
Stream.fold Fold.toList (Stream.nilM (print "nil"))"nil"[]

Pre-release

valuefromEffect :: Applicative m => m a -> Stream m a
#

Create a singleton stream from a monadic action.

Example2 expressions
fromEffect m = m `Stream.consM` Stream.nilfromEffect = Stream.sequence . Stream.fromPure
Example1 expression
Stream.fold Fold.drain $ Stream.fromEffect (putStrLn "hello")hello
valuefromByteStr# :: Monad m => Addr# -> Stream m Word8
#

Read bytes from an immutable Addr# until a 0 byte is encountered, the 0 byte is not included in the stream.

Example2 expressions
:set -XMagicHashfromByteStr# addr = Stream.takeWhile (/= 0) $ Stream.fromPtr $ Ptr addr

Unsafe: The caller is responsible for safe addressing.

Note that this is completely safe when reading from Haskell string literals because they are guaranteed to be NULL terminated:

Example1 expression
Stream.toList $ Stream.fromByteStr# "\1\2\3\0"#[1,2,3]
valuecons :: Applicative m => a -> Stream m a -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of elements. Use the O(n) complexity StreamK.Streamly.Data.StreamK.cons unless you want to statically fuse just a few elements.

Fuse a pure value at the head of an existing stream::

Example2 expressions
s = 1 `Stream.cons` Stream.fromList [2,3]Stream.toList s[1,2,3]

Definition:

Example1 expression
cons x xs = return x `Stream.consM` xs
valueunfoldr :: Monad m => (s -> Maybe (a, s)) -> s -> Stream m a
#

Build a stream by unfolding a pure step function step starting from a seed s. The step function returns the next element in the stream and the next seed value. When it is done it returns Nothing and the stream ends. For example,

Example1 expression
:{let f b =        if b > 2        then Nothing        else Just (b, b + 1)in Stream.toList $ Stream.unfoldr f 0:}[0,1,2]
valueunfoldrM :: Monad m => (s -> m (Maybe (a, s))) -> s -> Stream m a
#

Build a stream by unfolding a monadic step function starting from a seed. The step function returns the next element in the stream and the next seed value. When it is done it returns Nothing and the stream ends. For example,

Example1 expression
:{let f b =        if b > 2        then return Nothing        else return (Just (b, b + 1))in Stream.toList $ Stream.unfoldrM f 0:}[0,1,2]
classclass Enum a => Enumerable a where
#

Types that can be enumerated as a stream. The operations in this type class are equivalent to those in the Enum type class, except that these generate a stream instead of a list. Use the functions in Streamly.Internal.Data.Stream.Enumeration module to define new instances.

Methods

  • enumerateFrom :: Monad m => a -> Stream m a

    enumerateFrom from generates a stream starting with the element from, enumerating up to maxBound when the type is Bounded or generating an infinite stream when the type is not Bounded.

    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFrom (0 :: Int)[0,1,2,3]

    For Fractional types, enumeration is numerically stable. However, no overflow or underflow checks are performed.

    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFrom 1.1[1.1,2.1,3.1,4.1]
  • enumerateFromTo :: Monad m => a -> a -> Stream m a

    Generate a finite stream starting with the element from, enumerating the type up to the value to. If to is smaller than from then an empty stream is returned.

    Example1 expression
    Stream.toList $ Stream.enumerateFromTo 0 4[0,1,2,3,4]

    For Fractional types, the last element is equal to the specified to value after rounding to the nearest integral value.

    Example1 expression
    Stream.toList $ Stream.enumerateFromTo 1.1 4[1.1,2.1,3.1,4.1]
    Example1 expression
    Stream.toList $ Stream.enumerateFromTo 1.1 4.6[1.1,2.1,3.1,4.1,5.1]
  • enumerateFromThen :: Monad m => a -> a -> Stream m a

    enumerateFromThen from then generates a stream whose first element is from, the second element is then and the successive elements are in increments of then - from. Enumeration can occur downwards or upwards depending on whether then comes before or after from. For Bounded types the stream ends when maxBound is reached, for unbounded types it keeps enumerating infinitely.

    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 2[0,2,4,6]
    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 (-2)[0,-2,-4,-6]
  • enumerateFromThenTo :: Monad m => a -> a -> a -> Stream m a

    enumerateFromThenTo from then to generates a finite stream whose first element is from, the second element is then and the successive elements are in increments of then - from up to to. Enumeration can occur downwards or upwards depending on whether then comes before or after from.

    Example1 expression
    Stream.toList $ Stream.enumerateFromThenTo 0 2 6[0,2,4,6]
    Example1 expression
    Stream.toList $ Stream.enumerateFromThenTo 0 (-2) (-6)[0,-2,-4,-6]
Instances21Enumerable, …
  • Enumerable IntegerDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable NaturalDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int16Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int32Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int64Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int8Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word16Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word32Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word64Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word8Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable BoolDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable CharDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable DoubleDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable FloatDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable IntDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable OrderingDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable WordDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable ()Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Integral a => Enumerable (Ratio a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable a => Enumerable (Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • HasResolution a => Enumerable (Fixed a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
valueunfold :: Applicative m => Unfold m a b -> a -> Stream m b
#

Convert an Unfold into a stream by supplying it an input seed.

Example2 expressions
s = Stream.unfold Unfold.replicateM (3, putStrLn "hello")Stream.fold Fold.drain shellohellohello
valuefromFoldable :: (Monad m, Foldable f) => f a -> Stream m a
#
Example1 expression
fromFoldable = Prelude.foldr Stream.cons Stream.nil

Construct a stream from a Foldable containing pure values:

/WARNING: O(n^2), suitable only for a small number of elements in the stream/

valueenumerateFromStepNum :: (Monad m, Num a) => a -> a -> Stream m a
#

For floating point numbers if the increment is less than the precision then it just gets lost. Therefore we cannot always increment it correctly by just repeated addition. 9007199254740992 + 1 + 1 :: Double => 9.007199254740992e15 9007199254740992 + 2 :: Double => 9.007199254740994e15

Instead we accumulate the increment counter and compute the increment every time before adding it to the starting number.

This works for Integrals as well as floating point numbers, but enumerateFromStepIntegral is faster for integrals.

valueenumerateFromIntegral :: (Monad m, Integral a, Bounded a) => a -> Stream m a
#

Enumerate an Integral type. enumerateFromIntegral from generates a stream whose first element is from and the successive elements are in increments of 1. The stream is bounded by the size of the Integral type.

Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromIntegral (0 :: Int)[0,1,2,3]
valueenumerateFromThenIntegral
  1. :: (Monad m, Integral a, Bounded a)
  2. => a
  3. -> a
  4. -> Stream m a
#

Enumerate an Integral type in steps. enumerateFromThenIntegral from then generates a stream whose first element is from, the second element is then and the successive elements are in increments of then - from. The stream is bounded by the size of the Integral type.

Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenIntegral (0 :: Int) 2[0,2,4,6]
Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenIntegral (0 :: Int) (-2)[0,-2,-4,-6]
valueenumerateFromToIntegral :: (Monad m, Integral a) => a -> a -> Stream m a
#

Enumerate an Integral type up to a given limit. enumerateFromToIntegral from to generates a finite stream whose first element is from and successive elements are in increments of 1 up to to.

Example1 expression
Stream.toList $ Stream.enumerateFromToIntegral 0 4[0,1,2,3,4]
valueenumerateFromThenToIntegral
  1. :: (Monad m, Integral a)
  2. => a
  3. -> a
  4. -> a
  5. -> Stream m a
#

Enumerate an Integral type in steps up to a given limit. enumerateFromThenToIntegral from then to generates a finite stream whose first element is from, the second element is then and the successive elements are in increments of then - from up to to.

Example1 expression
Stream.toList $ Stream.enumerateFromThenToIntegral 0 2 6[0,2,4,6]
Example1 expression
Stream.toList $ Stream.enumerateFromThenToIntegral 0 (-2) (-6)[0,-2,-4,-6]
valueenumerateFromStepIntegral :: (Integral a, Monad m) => a -> a -> Stream m a
#

enumerateFromStepIntegral from step generates an infinite stream whose first element is from and the successive elements are in increments of step.

CAUTION: This function is not safe for finite integral types. It does not check for overflow, underflow or bounds.

Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromStepIntegral 0 2[0,2,4,6]
Example1 expression
Stream.toList $ Stream.take 3 $ Stream.enumerateFromStepIntegral 0 (-2)[0,-2,-4]
valueenumerateFromFractional :: (Monad m, Fractional a) => a -> Stream m a
#

Numerically stable enumeration from a Fractional number in steps of size 1. enumerateFromFractional from generates a stream whose first element is from and the successive elements are in increments of 1. No overflow or underflow checks are performed.

This is the equivalent to enumFrom for Fractional types. For example:

Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromFractional 1.1[1.1,2.1,3.1,4.1]
valueenumerateFromToFractional
  1. :: (Monad m, Fractional a, Ord a)
  2. => a
  3. -> a
  4. -> Stream m a
#

Numerically stable enumeration from a Fractional number to a given limit. enumerateFromToFractional from to generates a finite stream whose first element is from and successive elements are in increments of 1 up to to.

This is the equivalent of enumFromTo for Fractional types. For example:

Example1 expression
Stream.toList $ Stream.enumerateFromToFractional 1.1 4[1.1,2.1,3.1,4.1]
Example1 expression
Stream.toList $ Stream.enumerateFromToFractional 1.1 4.6[1.1,2.1,3.1,4.1,5.1]

Notice that the last element is equal to the specified to value after rounding to the nearest integer.

valueenumerateFromThenFractional
  1. :: (Monad m, Fractional a)
  2. => a
  3. -> a
  4. -> Stream m a
#

Numerically stable enumeration from a Fractional number in steps. enumerateFromThenFractional from then generates a stream whose first element is from, the second element is then and the successive elements are in increments of then - from. No overflow or underflow checks are performed.

This is the equivalent of enumFromThen for Fractional types. For example:

Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenFractional 1.1 2.1[1.1,2.1,3.1,4.1]
Example1 expression
Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenFractional 1.1 (-2.1)[1.1,-2.1,-5.300000000000001,-8.500000000000002]
valueenumerateFromThenToFractional
  1. :: (Monad m, Fractional a, Ord a)
  2. => a
  3. -> a
  4. -> a
  5. -> Stream m a
#

Numerically stable enumeration from a Fractional number in steps up to a given limit. enumerateFromThenToFractional from then to generates a finite stream whose first element is from, the second element is then and the successive elements are in increments of then - from up to to.

This is the equivalent of enumFromThenTo for Fractional types. For example:

Example1 expression
Stream.toList $ Stream.enumerateFromThenToFractional 0.1 2 6[0.1,2.0,3.9,5.799999999999999]
Example1 expression
Stream.toList $ Stream.enumerateFromThenToFractional 0.1 (-2) (-6)[0.1,-2.0,-4.1000000000000005,-6.200000000000001]
valuetimes :: MonadIO m => Stream m (AbsTime, RelTime64)
#

times returns a stream of time value tuples with clock of 10 ms granularity. The first component of the tuple is an absolute time reference (epoch) denoting the start of the stream and the second component is a time relative to the reference.

Example2 expressions
f = Fold.drainMapM (\x -> print x >> threadDelay 1000000)Stream.fold f $ Stream.take 3 $ Stream.times(AbsTime (TimeSpec {sec = ..., nsec = ...}),RelTime64 (NanoSecond64 ...))(AbsTime (TimeSpec {sec = ..., nsec = ...}),RelTime64 (NanoSecond64 ...))(AbsTime (TimeSpec {sec = ..., nsec = ...}),RelTime64 (NanoSecond64 ...))

Note: This API is not safe on 32-bit machines.

Pre-release

valuetimesWith :: MonadIO m => Double -> Stream m (AbsTime, RelTime64)
#

timesWith g returns a stream of time value tuples. The first component of the tuple is an absolute time reference (epoch) denoting the start of the stream and the second component is a time relative to the reference.

The argument g specifies the granularity of the relative time in seconds. A lower granularity clock gives higher precision but is more expensive in terms of CPU usage. Any granularity lower than 1 ms is treated as 1 ms.

Example3 expressions
import Control.Concurrent (threadDelay)f = Fold.drainMapM (\x -> print x >> threadDelay 1000000)Stream.fold f $ Stream.take 3 $ Stream.timesWith 0.01(AbsTime (TimeSpec {sec = ..., nsec = ...}),RelTime64 (NanoSecond64 ...))(AbsTime (TimeSpec {sec = ..., nsec = ...}),RelTime64 (NanoSecond64 ...))(AbsTime (TimeSpec {sec = ..., nsec = ...}),RelTime64 (NanoSecond64 ...))

Note: This API is not safe on 32-bit machines.

Pre-release

valueabsTimes :: MonadIO m => Stream m AbsTime
#

absTimes returns a stream of absolute timestamps using a clock of 10 ms granularity.

Example2 expressions
f = Fold.drainMapM printStream.fold f $ Stream.delayPre 1 $ Stream.take 3 $ Stream.absTimesAbsTime (TimeSpec {sec = ..., nsec = ...})AbsTime (TimeSpec {sec = ..., nsec = ...})AbsTime (TimeSpec {sec = ..., nsec = ...})

Note: This API is not safe on 32-bit machines.

Pre-release

valueabsTimesWith :: MonadIO m => Double -> Stream m AbsTime
#

absTimesWith g returns a stream of absolute timestamps using a clock of granularity g specified in seconds. A low granularity clock is more expensive in terms of CPU usage. Any granularity lower than 1 ms is treated as 1 ms.

Example2 expressions
f = Fold.drainMapM printStream.fold f $ Stream.delayPre 1 $ Stream.take 3 $ Stream.absTimesWith 0.01AbsTime (TimeSpec {sec = ..., nsec = ...})AbsTime (TimeSpec {sec = ..., nsec = ...})AbsTime (TimeSpec {sec = ..., nsec = ...})

Note: This API is not safe on 32-bit machines.

Pre-release

valuerelTimes :: MonadIO m => Stream m RelTime64
#

relTimes returns a stream of relative time values starting from 0, using a clock of granularity 10 ms.

Example2 expressions
f = Fold.drainMapM printStream.fold f $ Stream.delayPre 1 $ Stream.take 3 $ Stream.relTimesRelTime64 (NanoSecond64 ...)RelTime64 (NanoSecond64 ...)RelTime64 (NanoSecond64 ...)

Note: This API is not safe on 32-bit machines.

Pre-release

relTimesWith g returns a stream of relative time values starting from 0, using a clock of granularity g specified in seconds. A low granularity clock is more expensive in terms of CPU usage. Any granularity lower than 1 ms is treated as 1 ms.

Example2 expressions
f = Fold.drainMapM printStream.fold f $ Stream.delayPre 1 $ Stream.take 3 $ Stream.relTimesWith 0.01RelTime64 (NanoSecond64 ...)RelTime64 (NanoSecond64 ...)RelTime64 (NanoSecond64 ...)

Note: This API is not safe on 32-bit machines.

Pre-release

valuedurations :: Double -> t m RelTime64
#

durations g returns a stream of relative time values measuring the time elapsed since the immediate predecessor element of the stream was generated. The first element of the stream is always 0. durations uses a clock of granularity g specified in seconds. A low granularity clock is more expensive in terms of CPU usage. The minimum granularity is 1 millisecond. Durations lower than 1 ms will be 0.

Note: This API is not safe on 32-bit machines.

Unimplemented

valuetimeout :: AbsTime -> t m ()
#

Generate a singleton event at or after the specified absolute time. Note that this is different from a threadDelay, a threadDelay starts from the time when the action is evaluated, whereas if we use AbsTime based timeout it will immediately expire if the action is evaluated too late.

Unimplemented

valuefromFoldableM :: (Monad m, Foldable f) => f (m a) -> Stream m a
#
Example1 expression
fromFoldableM = Prelude.foldr Stream.consM Stream.nil

Construct a stream from a Foldable containing pure values:

/WARNING: O(n^2), suitable only for a small number of elements in the stream/

valuefromPtrN :: (Monad m, Storable a) => Int -> Ptr a -> Stream m a
#

Take n Storable elements starting from an immutable Ptr onwards.

Example1 expression
fromPtrN n = Stream.take n . Stream.fromPtr

Unsafe: The caller is responsible for safe addressing.

Pre-release

valuefoldr :: Monad m => (a -> b -> b) -> b -> Stream m a -> m b
#

Right fold, lazy for lazy monads and pure streams, and strict for strict monads.

Please avoid using this routine in strict monads like IO unless you need a strict right fold. This is provided only for use in lazy monads (e.g. Identity) or pure streams. Note that with this signature it is not possible to implement a lazy foldr when the monad m is strict. In that case it would be strict in its accumulator and therefore would necessarily consume all its input.

Example1 expression
foldr f z = Stream.foldrM (\a b -> f a <$> b) (return z)

Note: This is similar to Fold.foldr' (the right fold via left fold), but could be more efficient.

valuetoList :: Monad m => Stream m a -> m [a]
#

Definitions:

Example2 expressions
toList = Stream.foldr (:) []toList = Stream.fold Fold.toList

Convert a stream into a list in the underlying monad. The list can be consumed lazily in a lazy monad (e.g. Identity). In a strict monad (e.g. IO) the whole list is generated and buffered before it can be consumed.

Warning! working on large lists accumulated as buffers in memory could be very inefficient, consider using Streamly.Data.Array instead.

Note that this could a bit more efficient compared to Stream.fold Fold.toList, and it can fuse with pure list consumers.

valuefold :: Monad m => Fold m a b -> Stream m a -> m b
#

Fold a stream using the supplied left Fold and reducing the resulting expression strictly at each step. The behavior is similar to foldl'. A Fold can terminate early without consuming the full stream. See the documentation of individual Folds for termination behavior.

Definitions:

Example2 expressions
fold f = fmap fst . Stream.foldBreak ffold f = Stream.parse (Parser.fromFold f)

Example:

Example1 expression
Stream.fold Fold.sum (Stream.enumerateFromTo 1 100)5050
valueparse :: Monad m => Parser a m b -> Stream m a -> m (Either ParseError b)
#

Parse a stream using the supplied Parser.

Parsers (See Streamly.Internal.Data.Parser) are more powerful folds that add backtracking and error functionality to terminating folds. Unlike folds, parsers may not always result in a valid output, they may result in an error. For example:

Example1 expression
Stream.parse (Parser.takeEQ 1 Fold.drain) Stream.nilLeft (ParseError "takeEQ: Expecting exactly 1 elements, input terminated on 0")

Note: parse p is not the same as head . parseMany p on an empty stream.

valuemapM_ :: Monad m => (a -> m b) -> Stream m a -> m ()
#

Execute a monadic action for each element of the Stream

valuedrain :: Monad m => Stream m a -> m ()
#

Definitions:

Example2 expressions
drain = Stream.fold Fold.draindrain = Stream.foldrM (\_ xs -> xs) (return ())

Run a stream, discarding the results.

valueuncons :: Monad m => Stream m a -> m (Maybe (a, Stream m a))
#

Decompose a stream into its head and tail. If the stream is empty, returns Nothing. If the stream is non-empty, returns Just (a, ma), where a is the head of the stream and ma its tail.

Properties:

Example2 expressions
Nothing <- Stream.uncons Stream.nilJust ("a", t) <- Stream.uncons (Stream.cons "a" Stream.nil)

This can be used to consume the stream in an imperative manner one element at a time, as it just breaks down the stream into individual elements and we can loop over them as we deem fit. For example, this can be used to convert a streamly stream into other stream types.

All the folds in this module can be expressed in terms of uncons, however, this is generally less efficient than specific folds because it takes apart the stream one element at a time, therefore, does not take adavantage of stream fusion.

foldBreak is a more general way of consuming a stream piecemeal.

Example1 expression
:{uncons xs = do    r <- Stream.foldBreak Fold.one xs    return $ case r of        (Nothing, _) -> Nothing        (Just h, t) -> Just (h, t):}
valuefoldrM :: Monad m => (a -> m b -> m b) -> m b -> Stream m a -> m b
#

Right associative/lazy pull fold. foldrM build final stream constructs an output structure using the step function build. build is invoked with the next input element and the remaining (lazy) tail of the output structure. It builds a lazy output expression using the two. When the "tail structure" in the output expression is evaluated it calls build again thus lazily consuming the input stream until either the output expression built by build is free of the "tail" or the input is exhausted in which case final is used as the terminating case for the output structure. For more details see the description in the previous section.

Example, determine if any element is odd in a stream:

Example3 expressions
s = Stream.fromList (2:4:5:undefined)step x xs = if odd x then return True else xsStream.foldrM step (return False) sTrue
valueisPrefixOf :: (Monad m, Eq a) => Stream m a -> Stream m a -> m Bool
#

Returns True if the first stream is the same as or a prefix of the second. A stream is a prefix of itself.

Example1 expression
Stream.isPrefixOf (Stream.fromList "hello") (Stream.fromList "hello" :: Stream IO Char)True
valueisSubsequenceOf :: (Monad m, Eq a) => Stream m a -> Stream m a -> m Bool
#

Returns True if all the elements of the first stream occur, in order, in the second stream. The elements do not have to occur consecutively. A stream is a subsequence of itself.

Example1 expression
Stream.isSubsequenceOf (Stream.fromList "hlo") (Stream.fromList "hello" :: Stream IO Char)True
valuestripPrefix
  1. :: (Monad m, Eq a)
  2. => Stream m a
  3. -> Stream m a
  4. -> m (Maybe (Stream m a))
#

stripPrefix prefix input strips the prefix stream from the input stream if it is a prefix of input. Returns Nothing if the input does not start with the given prefix, stripped input otherwise. Returns Just nil when the prefix is the same as the input stream.

Space: O(1)

valueisInfixOf
  1. :: (MonadIO m, Eq a, Enum a, Storable a, Unbox a)
  2. => Stream m a
  3. -> Stream m a
  4. -> m Bool
#

Returns True if the first stream is an infix of the second. A stream is considered an infix of itself.

Example2 expressions
s = Stream.fromList "hello" :: Stream IO CharStream.isInfixOf s sTrue

Space: O(n) worst case where n is the length of the infix.

Pre-release

Requires Storable constraint

valueisSuffixOf :: (Monad m, Eq a) => Stream m a -> Stream m a -> m Bool
#

Returns True if the first stream is a suffix of the second. A stream is considered a suffix of itself.

Example1 expression
Stream.isSuffixOf (Stream.fromList "hello") (Stream.fromList "hello" :: Stream IO Char)True

Space: O(n), buffers entire input stream and the suffix.

Pre-release

Suboptimal - Help wanted.

valuestripSuffix
  1. :: (Monad m, Eq a)
  2. => Stream m a
  3. -> Stream m a
  4. -> m (Maybe (Stream m a))
#

Drops the given suffix from a stream. Returns Nothing if the stream does not end with the given suffix. Returns Just nil when the suffix is the same as the stream.

It may be more efficient to convert the stream to an Array and use stripSuffix on that especially if the elements have a Storable or Prim instance.

Space: O(n), buffers the entire input stream as well as the suffix

Pre-release

valuehandle
  1. :: (MonadCatch m, Exception e)
  2. => e -> m (Stream m a)
  3. -> Stream m a
  4. -> Stream m a
#

When evaluating a stream if an exception occurs, stream evaluation aborts and the specified exception handler is run with the exception as argument. The exception is caught and handled unless the handler decides to rethrow it. Note that exception handling is not applied to the stream returned by the exception handler.

Observes exceptions only in the stream generation, and not in stream consumers.

Inhibits stream fusion

valuebracketIO
  1. :: (MonadIO m, MonadCatch m)
  2. => IO b
  3. -> b -> IO c
  4. -> b -> Stream m a
  5. -> Stream m a
#

Run the alloc action IO b with async exceptions disabled but keeping blocking operations interruptible (see mask). Use the output b of the IO action as input to the function b -> Stream m a to generate an output stream.

b is usually a resource under the IO monad, e.g. a file handle, that requires a cleanup after use. The cleanup action b -> IO c, runs whenever (1) the stream ends normally, (2) due to a sync or async exception or, (3) if it gets garbage collected after a partial lazy evaluation. The exception is not caught, it is rethrown.

bracketIO only guarantees that the cleanup action runs, and it runs with async exceptions enabled. The action must ensure that it can successfully cleanup the resource in the face of sync or async exceptions.

When the stream ends normally or on a sync exception, cleanup action runs immediately in the current thread context, whereas in other cases it runs in the GC context, therefore, cleanup may be delayed until the GC gets to run. An example where GC based cleanup happens is when a stream is being folded but the fold terminates without draining the entire stream or if the consumer of the stream encounters an exception.

Observes exceptions only in the stream generation, and not in stream consumers.

See also: bracketUnsafe

Inhibits stream fusion

valuegbracket_
  1. :: Monad m
  2. => m c

    before

  3. -> (c -> m d)

    after, on normal stop

  4. -> (c -> e -> Stream m b -> m (Stream m b))

    on exception

  5. -> (forall s. m s -> m (Either e s))

    try (exception handling)

  6. -> (c -> Stream m b)

    stream generator

  7. -> Stream m b
#

Like gbracket but with following differences:

  • alloc action m c runs with async exceptions enabled

  • cleanup action c -> m d won't run if the stream is garbage collected after partial evaluation.

Inhibits stream fusion

Pre-release

valuebefore :: Monad m => m b -> Stream m a -> Stream m a
#

Run the action m b before the stream yields its first element.

Same as the following but more efficient due to fusion:

Example2 expressions
before action xs = Stream.nilM action <> xsbefore action xs = Stream.concatMap (const xs) (Stream.fromEffect action)
valueafterIO :: MonadIO m => IO b -> Stream m a -> Stream m a
#

Run the action IO b whenever the stream is evaluated to completion, or if it is garbage collected after a partial lazy evaluation.

The semantics of the action IO b are similar to the semantics of cleanup action in bracketIO.

See also afterUnsafe

valuefinallyIO :: (MonadIO m, MonadCatch m) => IO b -> Stream m a -> Stream m a
#

Run the action IO b whenever the stream stream stops normally, aborts due to an exception or if it is garbage collected after a partial lazy evaluation.

The semantics of running the action IO b are similar to the cleanup action semantics described in bracketIO.

Example1 expression
finallyIO release = Stream.bracketIO (return ()) (const release)

See also finallyUnsafe

Inhibits stream fusion

valueonException :: MonadCatch m => m b -> Stream m a -> Stream m a
#

Run the action m b if the stream evaluation is aborted due to an exception. The exception is not caught, simply rethrown.

Observes exceptions only in the stream generation, and not in stream consumers.

Inhibits stream fusion

valuebracketIO3
  1. :: (MonadIO m, MonadCatch m)
  2. => IO b
  3. -> b -> IO c
  4. -> b -> IO d
  5. -> b -> IO e
  6. -> b -> Stream m a
  7. -> Stream m a
#

Like bracketIO but can use 3 separate cleanup actions depending on the mode of termination:

  1. When the stream stops normally

  2. When the stream is garbage collected

  3. When the stream encounters an exception

bracketIO3 before onStop onGC onException action runs action using the result of before. If the stream stops, onStop action is executed, if the stream is abandoned onGC is executed, if the stream encounters an exception onException is executed.

The exception is not caught, it is rethrown.

Inhibits stream fusion

Pre-release

valuegbracket
  1. :: MonadIO m
  2. => IO c

    before

  3. -> (c -> IO d1)

    on normal stop

  4. -> (c -> e -> Stream m b -> IO (Stream m b))

    on exception

  5. -> (c -> IO d2)

    on GC without normal stop or exception

  6. -> (forall s. m s -> m (Either e s))

    try (exception handling)

  7. -> (c -> Stream m b)

    stream generator

  8. -> Stream m b
#

Run the alloc action m c with async exceptions disabled but keeping blocking operations interruptible (see mask). Use the output c as input to c -> Stream m b to generate an output stream. When generating the stream use the supplied try operation forall s. m s -> m (Either e s) to catch synchronous exceptions. If an exception occurs run the exception handler c -> e -> Stream m b -> m (Stream m b). Note that gbracket does not rethrow the exception, it has to be done by the exception handler if desired.

The cleanup action c -> m d, runs whenever the stream ends normally, due to a sync or async exception or if it gets garbage collected after a partial lazy evaluation. See bracket for the semantics of the cleanup action.

gbracket can express all other exception handling combinators.

Inhibits stream fusion

Pre-release

valueafterUnsafe :: Monad m => m b -> Stream m a -> Stream m a
#

Like after, with following differences:

  • action m b won't run if the stream is garbage collected after partial evaluation.

  • Monad m does not require any other constraints.

  • has slightly better performance than after.

Same as the following, but with stream fusion:

Example1 expression
afterUnsafe action xs = xs <> Stream.nilM action

Pre-release

valuebracketUnsafe
  1. :: MonadCatch m
  2. => m b
  3. -> b -> m c
  4. -> b -> Stream m a
  5. -> Stream m a
#

Like bracket but with following differences:

  • alloc action m b runs with async exceptions enabled

  • cleanup action b -> m c won't run if the stream is garbage collected after partial evaluation.

  • has slightly better performance than bracketIO.

Inhibits stream fusion

Pre-release

valuefinallyUnsafe :: MonadCatch m => m b -> Stream m a -> Stream m a
#

Like finally with following differences:

  • action m b won't run if the stream is garbage collected after partial evaluation.

  • has slightly better performance than finallyIO.

Inhibits stream fusion

Pre-release

valueghandle
  1. :: (MonadCatch m, Exception e)
  2. => e -> Stream m a -> m (Stream m a)
  3. -> Stream m a
  4. -> Stream m a
#

Like handle but the exception handler is also provided with the stream that generated the exception as input. The exception handler can thus re-evaluate the stream to retry the action that failed. The exception handler can again call ghandle on it to retry the action multiple times.

This is highly experimental. In a stream of actions we can map the stream with a retry combinator to retry each action on failure.

Inhibits stream fusion

Pre-release

valuemorphInner :: Monad n => (forall x. m x -> n x) -> Stream m a -> Stream n a
#

Transform the inner monad of a stream using a natural transformation.

Example, generalize the inner monad from Identity to any other:

Example1 expression
generalizeInner = Stream.morphInner (return . runIdentity)

Also known as hoist.

valueliftInnerWith
  1. :: Monad (t m)
  2. => forall b. m b -> t m b
  3. -> Stream m a
  4. -> Stream (t m) a
#

Lift the inner monad m of a stream Stream m a to t m using the supplied lift function.

valuerunInnerWith
  1. :: Monad m
  2. => forall b. t m b -> m b
  3. -> Stream (t m) a
  4. -> Stream m a
#

Evaluate the inner monad of a stream using the supplied runner function.

valuerunInnerWithState
  1. :: Monad m
  2. => forall b. s -> t m b -> m (b, s)
  3. -> m s
  4. -> Stream (t m) a
  5. -> Stream m (s, a)
#

Evaluate the inner monad of a stream using the supplied stateful runner function and the initial state. The state returned by an invocation of the runner is supplied as input state to the next invocation.

valuefoldrT
  1. :: (Monad m, Monad (t m), MonadTrans t)
  2. => a -> t m b -> t m b
  3. -> t m b
  4. -> Stream m a
  5. -> t m b
#

Right fold to a transformer monad. This is the most general right fold function. foldrS is a special case of foldrT, however foldrS implementation can be more efficient:

Example1 expression
foldrS = Stream.foldrT
Example2 expressions
step f x xs = lift $ f x (runIdentityT xs)foldrM f z s = runIdentityT $ Stream.foldrT (step f) (lift z) s

foldrT can be used to translate streamly streams to other transformer monads e.g. to a different streaming type.

Pre-release

valueusingStateT
  1. :: Monad m
  2. => m s
  3. -> Stream (StateT s m) a -> Stream (StateT s m) a
  4. -> Stream m a
  5. -> Stream m a
#

Run a stateful (StateT) stream transformation using a given state.

Example1 expression
usingStateT s f = Stream.evalStateT s . f . Stream.liftInner

See also: scan

valueappend :: Monad m => Stream m a -> Stream m a -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.append otherwise.

Fuses two streams sequentially, yielding all elements from the first stream, and then all elements from the second stream.

Example3 expressions
s1 = Stream.fromList [1,2]s2 = Stream.fromList [3,4]Stream.fold Fold.toList $ s1 `Stream.append` s2[1,2,3,4]
valuezipWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m c
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.zipWith otherwise.

Stream a is evaluated first, followed by stream b, the resulting elements a and b are then zipped using the supplied zip function and the result c is yielded to the consumer.

If stream a or stream b ends, the zipped stream ends. If stream b ends first, the element a from previous evaluation of stream a is discarded.

Example3 expressions
s1 = Stream.fromList [1,2,3]s2 = Stream.fromList [4,5,6]Stream.fold Fold.toList $ Stream.zipWith (+) s1 s2[5,7,9]
valuemergeBy
  1. :: Monad m
  2. => a -> a -> Ordering
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.mergeBy otherwise.

Merge two streams using a comparison function. The head elements of both the streams are compared and the smaller of the two elements is emitted, if both elements are equal then the element from the first stream is used first.

If the streams are sorted in ascending order, the resulting stream would also remain sorted in ascending order.

Example3 expressions
s1 = Stream.fromList [1,3,5]s2 = Stream.fromList [2,4,6,8]Stream.fold Fold.toList $ Stream.mergeBy compare s1 s2[1,2,3,4,5,6,8]
valuemergeByM
  1. :: Monad m
  2. => a -> a -> m Ordering
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

Like mergeBy but with a monadic comparison function.

Example, to merge two streams randomly:

> randomly _ _ = randomIO >>= x -> return $ if x then LT else GT
> Stream.toList $ Stream.mergeByM randomly (Stream.fromList [1,1,1,1]) (Stream.fromList [2,2,2,2])
[2,1,2,2,2,1,1,1]

Example, merge two streams in a proportion of 2:1:

Example1 expression
:{do let s1 = Stream.fromList [1,1,1,1,1,1]     s2 = Stream.fromList [2,2,2] let proportionately m n = do      ref <- newIORef $ cycle $ Prelude.concat [Prelude.replicate m LT, Prelude.replicate n GT]      return $ \_ _ -> do         r <- readIORef ref         writeIORef ref $ Prelude.tail r         return $ Prelude.head r f <- proportionately 2 1 xs <- Stream.fold Fold.toList $ Stream.mergeByM f s1 s2 print xs:}[1,1,2,1,1,2,1,1,2]
valueconcatMap :: Monad m => (a -> Stream m b) -> Stream m a -> Stream m b
#

Map a stream producing function on each element of the stream and then flatten the results into a single stream.

Example3 expressions
concatMap f = Stream.concatMapM (return . f)concatMap f = Stream.concat . fmap fconcatMap f = Stream.unfoldMany (Unfold.lmap f Unfold.fromStream)

See unfoldMany for a fusible alternative.

valuefoldMany :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Apply a Fold repeatedly on a stream and emit the results in the output stream.

Definition:

Example1 expression
foldMany f = Stream.parseMany (Parser.fromFold f)

Example, empty stream:

Example3 expressions
f = Fold.take 2 Fold.sumfmany = Stream.fold Fold.toList . Stream.foldMany ffmany $ Stream.fromList [][]

Example, last fold empty:

Example1 expression
fmany $ Stream.fromList [1..4][3,7]

Example, last fold non-empty:

Example1 expression
fmany $ Stream.fromList [1..5][3,7,5]

Note that using a closed fold e.g. Fold.take 0, would result in an infinite stream on a non-empty input stream.

valueroundRobin :: Monad m => Stream m a -> Stream m a -> Stream m a
#

Schedule the execution of two streams in a fair round-robin manner, executing each stream once, alternately. Execution of a stream may not necessarily result in an output, a stream may choose to Skip producing an element until later giving the other stream a chance to run. Therefore, this combinator fairly interleaves the execution of two streams rather than fairly interleaving the output of the two streams. This can be useful in co-operative multitasking without using explicit threads. This can be used as an alternative to async.

Do not use dynamically.

Pre-release

valueinterpose :: Monad m => c -> Unfold m b c -> Stream m b -> Stream m c
#

Unfold the elements of a stream, intersperse the given element between the unfolded streams and then concat them into a single stream.

Example1 expression
unwords = Stream.interpose ' '

Pre-release

valueinterposeSuffix :: Monad m => c -> Unfold m b c -> Stream m b -> Stream m c
#

Unfold the elements of a stream, append the given element after each unfolded stream and then concat them into a single stream.

Example1 expression
unlines = Stream.interposeSuffix '\n'

Pre-release

valueintercalateSuffix
  1. :: Monad m
  2. => Unfold m b c
  3. -> b
  4. -> Stream m b
  5. -> Stream m c
#

intersperseMSuffix followed by unfold and concat.

Example3 expressions
intercalateSuffix u a = Stream.unfoldMany u . Stream.intersperseMSuffix aintersperseMSuffix = Stream.intercalateSuffix Unfold.identityunlines = Stream.intercalateSuffix Unfold.fromList "\n"
Example2 expressions
input = Stream.fromList ["abc", "def", "ghi"]Stream.fold Fold.toList $ Stream.intercalateSuffix Unfold.fromList "\n" input"abc\ndef\nghi\n"
valuegroupsOf :: Monad m => Int -> Fold m a b -> Stream m a -> Stream m b
#

Group the input stream into groups of n elements each and then fold each group using the provided fold function.

groupsOf n f = foldMany (FL.take n f)
Example1 expression
Stream.toList $ Stream.groupsOf 2 Fold.sum (Stream.enumerateFromTo 1 10)[3,7,11,15,19]

This can be considered as an n-fold version of take where we apply take repeatedly on the leftover stream until the stream exhausts.

valueinterleave :: Monad m => Stream m a -> Stream m a -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.interleave otherwise.

Interleaves two streams, yielding one element from each stream alternately. When one stream stops the rest of the other stream is used in the output stream.

valueunfoldMany :: Monad m => Unfold m a b -> Stream m a -> Stream m b
#

unfoldMany unfold stream uses unfold to map the input stream elements to streams and then flattens the generated streams into a single output stream.

Like concatMap but uses an Unfold for stream generation. Unlike concatMap this can fuse the Unfold code with the inner loop and therefore provide many times better performance.

valueintercalate :: Monad m => Unfold m b c -> b -> Stream m b -> Stream m c
#

intersperse followed by unfold and concat.

Example3 expressions
intercalate u a = Stream.unfoldMany u . Stream.intersperse aintersperse = Stream.intercalate Unfold.identityunwords = Stream.intercalate Unfold.fromList " "
Example2 expressions
input = Stream.fromList ["abc", "def", "ghi"]Stream.fold Fold.toList $ Stream.intercalate Unfold.fromList " " input"abc def ghi"
valueconcatMapM :: Monad m => (a -> m (Stream m b)) -> Stream m a -> Stream m b
#

Map a stream producing monadic function on each element of the stream and then flatten the results into a single stream. Since the stream generation function is monadic, unlike concatMap, it can produce an effect at the beginning of each iteration of the inner loop.

See unfoldMany for a fusible alternative.

valueparseMany
  1. :: Monad m
  2. => Parser a m b
  3. -> Stream m a
  4. -> Stream m (Either ParseError b)
#

Apply a Parser repeatedly on a stream and emit the parsed values in the output stream.

Example:

Example3 expressions
s = Stream.fromList [1..10]parser = Parser.takeBetween 0 2 Fold.sumStream.fold Fold.toList $ Stream.parseMany parser s[Right 3,Right 7,Right 11,Right 15,Right 19]

This is the streaming equivalent of the Streamly.Data.Parser.many parse combinator.

Known Issues: When the parser fails there is no way to get the remaining stream.

valuewordsBy :: Monad m => (a -> Bool) -> Fold m a b -> Stream m a -> Stream m b
#

Split the stream after stripping leading, trailing, and repeated separators as per the fold supplied. Therefore, ".a..b." with . as the separator would be parsed as ["a","b"]. In other words, its like parsing words from whitespace separated text.

valueinterleaveFst :: Monad m => Stream m a -> Stream m a -> Stream m a
#

Interleaves the outputs of two streams, yielding elements from each stream alternately, starting from the first stream and ending at the first stream. If the second stream is longer than the first, elements from the second stream are infixed with elements from the first stream. If the first stream is longer then it continues yielding elements even after the second stream has finished.

Example4 expressions
:set -XOverloadedStringsimport Data.Functor.Identity (Identity)Stream.interleaveFst "abc" ",,,," :: Stream Identity CharfromList "a,b,c"Stream.interleaveFst "abc" "," :: Stream Identity CharfromList "a,bc"

interleaveFst is a dual of interleaveFstSuffix.

Do not use dynamically.

Pre-release

valueinterleaveFstSuffix :: Monad m => Stream m a -> Stream m a -> Stream m a
#

Interleaves the outputs of two streams, yielding elements from each stream alternately, starting from the first stream. As soon as the first stream finishes, the output stops, discarding the remaining part of the second stream. In this case, the last element in the resulting stream would be from the second stream. If the second stream finishes early then the first stream still continues to yield elements until it finishes.

Example4 expressions
:set -XOverloadedStringsimport Data.Functor.Identity (Identity)Stream.interleaveFstSuffix "abc" ",,,," :: Stream Identity CharfromList "a,b,c,"Stream.interleaveFstSuffix "abc" "," :: Stream Identity CharfromList "a,bc"

interleaveFstSuffix is a dual of interleaveFst.

Do not use dynamically.

Pre-release

valueunfoldInterleave :: Monad m => Unfold m a b -> Stream m a -> Stream m b
#

This does not pair streams like mergeMapWith, instead, it goes through each stream one by one and yields one element from each stream. After it goes to the last stream it reverses the traversal to come back to the first stream yielding elements from each stream on its way back to the first stream and so on.

Example3 expressions
lists = Stream.fromList [[1,1],[2,2],[3,3],[4,4],[5,5]]interleaved = Stream.unfoldInterleave Unfold.fromList listsStream.fold Fold.toList interleaved[1,2,3,4,5,5,4,3,2,1]

Note that this is order of magnitude more efficient than "mergeMapWith interleave" because of fusion.

valueunfoldRoundRobin :: Monad m => Unfold m a b -> Stream m a -> Stream m b
#

unfoldInterleave switches to the next stream whenever a value from a stream is yielded, it does not switch on a Skip. So if a stream keeps skipping for long time other streams won't get a chance to run. unfoldRoundRobin switches on Skip as well. So it basically schedules each stream fairly irrespective of whether it produces a value or not.

valuefoldIterateM
  1. :: Monad m
  2. => b -> m (Fold m a b)
  3. -> m b
  4. -> Stream m a
  5. -> Stream m b
#

Iterate a fold generator on a stream. The initial value b is used to generate the first fold, the fold is applied on the stream and the result of the fold is used to generate the next fold and so on.

Example4 expressions
import Data.Monoid (Sum(..))f x = return (Fold.take 2 (Fold.sconcat x))s = fmap Sum $ Stream.fromList [1..10]Stream.fold Fold.toList $ fmap getSum $ Stream.foldIterateM f (pure 0) s[3,10,21,36,55,55]

This is the streaming equivalent of monad like sequenced application of folds where next fold is dependent on the previous fold.

Pre-release

valueparseManyTill :: Parser a m b -> Parser a m x -> Stream m a -> Stream m b
#

parseManyTill collect test stream tries the parser test on the input, if test fails it backtracks and tries collect, after collect succeeds test is tried again and so on. The parser stops when test succeeds. The output of test is discarded and the output of collect is emitted in the output stream. The parser fails if collect fails.

Unimplemented

valueparseIterate
  1. :: Monad m
  2. => b -> Parser a m b
  3. -> b
  4. -> Stream m a
  5. -> Stream m (Either ParseError b)
#

Iterate a parser generating function on a stream. The initial value b is used to generate the first parser, the parser is applied on the stream and the result is used to generate the next parser and so on.

Example3 expressions
import Data.Monoid (Sum(..))s = Stream.fromList [1..10]Stream.fold Fold.toList $ fmap getSum $ Stream.catRights $ Stream.parseIterate (\b -> Parser.takeBetween 0 2 (Fold.sconcat b)) (Sum 0) $ fmap Sum s[3,10,21,36,55,55]

This is the streaming equivalent of monad like sequenced application of parsers where next parser is dependent on the previous parser.

Pre-release

valuegroupsBy
  1. :: Monad m
  2. => a -> a -> Bool
  3. -> Fold m a b
  4. -> Stream m a
  5. -> Stream m b
#

Deprecated. Please use groupsWhile instead. Please note the change in the argument order of the comparison function.

valuegroupsWhile
  1. :: Monad m
  2. => a -> a -> Bool
  3. -> Fold m a b
  4. -> Stream m a
  5. -> Stream m b
#

The argument order of the comparison function in groupsWhile is different than that of groupsBy.

In groupsBy the comparison function takes the next element as the first argument and the previous element as the second argument. In groupsWhile the first argument is the previous element and second argument is the next element.

valuesplitOnPrefix :: (a -> Bool) -> Fold m a b -> Stream m a -> Stream m b
#

Split on a prefixed separator element, dropping the separator. The supplied Fold is applied on the split segments.

> splitOnPrefix' p xs = Stream.toList $ Stream.splitOnPrefix p (Fold.toList) (Stream.fromList xs)
> splitOnPrefix' (== .) ".a.b"
["a","b"]

An empty stream results in an empty output stream: > splitOnPrefix' (== .) "" []

An empty segment consisting of only a prefix is folded to the default output of the fold:

> splitOnPrefix' (== .) "."
[""]

> splitOnPrefix' (== .) ".a.b."
["a","b",""]

> splitOnPrefix' (== .) ".a..b"
["a","","b"]

A prefix is optional at the beginning of the stream:

> splitOnPrefix' (== .) "a"
["a"]

> splitOnPrefix' (== .) "a.b"
["a","b"]

splitOnPrefix is an inverse of intercalatePrefix with a single element:

Stream.intercalatePrefix (Stream.fromPure '.') Unfold.fromList . Stream.splitOnPrefix (== '.') Fold.toList === id

Assuming the input stream does not contain the separator:

Stream.splitOnPrefix (== '.') Fold.toList . Stream.intercalatePrefix (Stream.fromPure '.') Unfold.fromList === id

Unimplemented

valuedropInfix :: Stream m a -> Stream m a -> Stream m a
#

Drop all matching infix from the input stream if present. Infix stream may be consumed multiple times.

Space: O(n) where n is the length of the infix.

Unimplemented

valuedropSuffix :: Stream m a -> Stream m a -> Stream m a
#

Drop suffix from the input stream if present. Suffix stream may be consumed multiple times.

Space: O(n) where n is the length of the suffix.

Unimplemented

valuefindIndices :: Monad m => (a -> Bool) -> Stream m a -> Stream m Int
#

Find all the indices where the element in the stream satisfies the given predicate.

Example1 expression
findIndices p = Stream.scanMaybe (Fold.findIndices p)
valuescanl' :: Monad m => (b -> a -> b) -> b -> Stream m a -> Stream m b
#

Strict left scan. Like map, scanl' too is a one to one transformation, however it adds an extra element.

Example1 expression
Stream.toList $ Stream.scanl' (+) 0 $ Stream.fromList [1,2,3,4][0,1,3,6,10]
Example1 expression
Stream.toList $ Stream.scanl' (flip (:)) [] $ Stream.fromList [1,2,3,4][[],[1],[2,1],[3,2,1],[4,3,2,1]]

The output of scanl' is the initial value of the accumulator followed by all the intermediate steps and the final result of foldl'.

By streaming the accumulated state after each fold step, we can share the state across multiple stages of stream composition. Each stage can modify or extend the state, do some processing with it and emit it for the next stage, thus modularizing the stream processing. This can be useful in stateful or event-driven programming.

Consider the following monolithic example, computing the sum and the product of the elements in a stream in one go using a foldl':

Example1 expression
Stream.fold (Fold.foldl' (\(s, p) x -> (s + x, p * x)) (0,1)) $ Stream.fromList [1,2,3,4](10,24)

Using scanl' we can make it modular by computing the sum in the first stage and passing it down to the next stage for computing the product:

Example1 expression
:{  Stream.fold (Fold.foldl' (\(_, p) (s, x) -> (s, p * x)) (0,1))  $ Stream.scanl' (\(s, _) x -> (s + x, x)) (0,1)  $ Stream.fromList [1,2,3,4]:}(10,24)

IMPORTANT: scanl' evaluates the accumulator to WHNF. To avoid building lazy expressions inside the accumulator, it is recommended that a strict data structure is used for accumulator.

Example2 expressions
scanl' step z = Stream.scan (Fold.foldl' step z)scanl' f z xs = Stream.scanlM' (\a b -> return (f a b)) (return z) xs

See also: usingStateT

valuefilter :: Monad m => (a -> Bool) -> Stream m a -> Stream m a
#

Include only those elements that pass a predicate.

Example3 expressions
filter p = Stream.filterM (return . p)filter p = Stream.mapMaybe (\x -> if p x then Just x else Nothing)filter p = Stream.scanMaybe (Fold.filtering p)
valuedropWhile :: Monad m => (a -> Bool) -> Stream m a -> Stream m a
#

Drop elements in the stream as long as the predicate succeeds and then take the rest of the stream.

valuemapM :: Monad m => (a -> m b) -> Stream m a -> Stream m b
#
Example1 expression
mapM f = Stream.sequence . fmap f

Apply a monadic function to each element of the stream and replace it with the output of the resulting action.

Example2 expressions
s = Stream.fromList ["a", "b", "c"]Stream.fold Fold.drain $ Stream.mapM putStr sabc
valuesequence :: Monad m => Stream m (m a) -> Stream m a
#
Example1 expression
sequence = Stream.mapM id

Replace the elements of a stream of monadic actions with the outputs of those actions.

Example2 expressions
s = Stream.fromList [putStr "a", putStr "b", putStrLn "c"]Stream.fold Fold.drain $ Stream.sequence sabc
valueintersperseM :: Monad m => m a -> Stream m a -> Stream m a
#

Insert an effect and its output before consuming an element of a stream except the first one.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.trace putChar $ Stream.intersperseM (putChar '.' >> return ',') inputh.,e.,l.,l.,o"h,e,l,l,o"

Be careful about the order of effects. In the above example we used trace after the intersperse, if we use it before the intersperse the output would be he.l.l.o."h,e,l,l,o".

Example1 expression
Stream.fold Fold.toList $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.trace putChar inputhe.l.l.o."h,e,l,l,o"
valueintersperse :: Monad m => a -> Stream m a -> Stream m a
#

Insert a pure value between successive elements of a stream.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.intersperse ',' input"h,e,l,l,o"
valueinsertBy :: Monad m => (a -> a -> Ordering) -> a -> Stream m a -> Stream m a
#

insertBy cmp elem stream inserts elem before the first element in stream that is less than elem when compared using cmp.

Example1 expression
insertBy cmp x = Stream.mergeBy cmp (Stream.fromPure x)
Example2 expressions
input = Stream.fromList [1,3,5]Stream.fold Fold.toList $ Stream.insertBy compare 2 input[1,2,3,5]
valuedeleteBy :: Monad m => (a -> a -> Bool) -> a -> Stream m a -> Stream m a
#

Deletes the first occurrence of the element in the stream that satisfies the given equality predicate.

Example2 expressions
input = Stream.fromList [1,3,3,5]Stream.fold Fold.toList $ Stream.deleteBy (==) 3 input[1,3,5]
valuemapMaybe :: Monad m => (a -> Maybe b) -> Stream m a -> Stream m b
#

Map a Maybe returning function to a stream, filter out the Nothing elements, and return a stream of values extracted from Just.

Equivalent to:

Example1 expression
mapMaybe f = Stream.catMaybes . fmap f
valuereverse :: Monad m => Stream m a -> Stream m a
#

Returns the elements of the stream in reverse order. The stream must be finite. Note that this necessarily buffers the entire stream in memory.

Definition:

Example1 expression
reverse m = Stream.concatEffect $ Stream.fold Fold.toListRev m >>= return . Stream.fromList
valuepostscan :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Postscan a stream using the given monadic fold.

The following example extracts the input stream up to a point where the running average of elements is no more than 10:

Example4 expressions
import Data.Maybe (fromJust)let avg = Fold.teeWith (/) Fold.sum (fmap fromIntegral Fold.length)s = Stream.enumerateFromTo 1.0 100.0:{ Stream.fold Fold.toList  $ fmap (fromJust . fst)  $ Stream.takeWhile (\(_,x) -> x <= 10)  $ Stream.postscan (Fold.tee Fold.latest avg) s:}[1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0,12.0,13.0,14.0,15.0,16.0,17.0,18.0,19.0]
valuescan :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Strict left scan. Scan a stream using the given monadic fold.

Example2 expressions
s = Stream.fromList [1..10]Stream.fold Fold.toList $ Stream.takeWhile (< 10) $ Stream.scan Fold.sum s[0,1,3,6]

See also: usingStateT

valuefilterM :: Monad m => (a -> m Bool) -> Stream m a -> Stream m a
#

Same as filter but with a monadic predicate.

Example2 expressions
f p x = p x >>= \r -> return $ if r then Just x else NothingfilterM p = Stream.mapMaybeM (f p)
valueindexed :: Monad m => Stream m a -> Stream m (Int, a)
#
Example4 expressions
f = Fold.foldl' (\(i, _) x -> (i + 1, x)) (-1,undefined)indexed = Stream.postscan findexed = Stream.zipWith (,) (Stream.enumerateFrom 0)indexedR n = fmap (\(i, a) -> (n - i, a)) . indexed

Pair each element in a stream with its index, starting from index 0.

Example1 expression
Stream.fold Fold.toList $ Stream.indexed $ Stream.fromList "hello"[(0,'h'),(1,'e'),(2,'l'),(3,'l'),(4,'o')]
valuesplitOn :: Monad m => (a -> Bool) -> Fold m a b -> Stream m a -> Stream m b
#

Split on an infixed separator element, dropping the separator. The supplied Fold is applied on the split segments. Splits the stream on separator elements determined by the supplied predicate, separator is considered as infixed between two segments:

Example2 expressions
splitOn' p xs = Stream.fold Fold.toList $ Stream.splitOn p Fold.toList (Stream.fromList xs)splitOn' (== '.') "a.b"["a","b"]

An empty stream is folded to the default value of the fold:

Example1 expression
splitOn' (== '.') ""[""]

If one or both sides of the separator are missing then the empty segment on that side is folded to the default output of the fold:

Example1 expression
splitOn' (== '.') "."["",""]
Example1 expression
splitOn' (== '.') ".a"["","a"]
Example1 expression
splitOn' (== '.') "a."["a",""]
Example1 expression
splitOn' (== '.') "a..b"["a","","b"]

splitOn is an inverse of intercalating single element:

Stream.intercalate (Stream.fromPure '.') Unfold.fromList . Stream.splitOn (== '.') Fold.toList === id

Assuming the input stream does not contain the separator:

Stream.splitOn (== '.') Fold.toList . Stream.intercalate (Stream.fromPure '.') Unfold.fromList === id
valueintersperseMSuffix :: Monad m => m a -> Stream m a -> Stream m a
#

Insert an effect and its output after consuming an element of a stream.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.trace putChar $ Stream.intersperseMSuffix (putChar '.' >> return ',') inputh.,e.,l.,l.,o.,"h,e,l,l,o,"

Pre-release

valueelemIndices :: (Monad m, Eq a) => a -> Stream m a -> Stream m Int
#

Find all the indices where the value of the element in the stream is equal to the given value.

Example1 expression
elemIndices a = Stream.findIndices (== a)
valueuniqBy :: Monad m => (a -> a -> Bool) -> Stream m a -> Stream m a
#

Drop repeated elements that are adjacent to each other using the supplied comparison function.

Example1 expression
uniq = Stream.uniqBy (==)

To strip duplicate path separators:

Example3 expressions
input = Stream.fromList "//a//b"f x y = x == '/' && y == '/'Stream.fold Fold.toList $ Stream.uniqBy f input"/a/b"

Space: O(1)

Pre-release

valuescanMaybe :: Monad m => Fold m a (Maybe b) -> Stream m a -> Stream m b
#

Use a filtering fold on a stream.

Example1 expression
scanMaybe f = Stream.catMaybes . Stream.postscan f
valuecatMaybes :: Monad m => Stream m (Maybe a) -> Stream m a
#

In a stream of Maybes, discard Nothings and unwrap Justs.

Example2 expressions
catMaybes = Stream.mapMaybe idcatMaybes = fmap fromJust . Stream.filter isJust

Pre-release

valuecatEithers :: Monad m => Stream m (Either a a) -> Stream m a
#

Remove the either wrapper and flatten both lefts and as well as rights in the output stream.

Example1 expression
catEithers = fmap (either id id)

Pre-release

valuetrace :: Monad m => (a -> m b) -> Stream m a -> Stream m a
#

Apply a monadic function to each element flowing through the stream and discard the results.

Example2 expressions
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.trace print s12

Compare with tap.

valuetap :: Monad m => Fold m a b -> Stream m a -> Stream m a
#

Tap the data flowing through a stream into a Fold. For example, you may add a tap to log the contents flowing through the stream. The fold is used only for effects, its result is discarded.

                  Fold m a b
                      |
-----stream m a ---------------stream m a-----

Example2 expressions
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.tap (Fold.drainMapM print) s12

Compare with trace.

valuedelay :: MonadIO m => Double -> Stream m a -> Stream m a
#

Introduce a delay of specified seconds between elements of the stream.

Definition:

Example2 expressions
sleep n = liftIO $ threadDelay $ round $ n * 1000000delay = Stream.intersperseM_ . sleep

Example:

Example2 expressions
input = Stream.enumerateFromTo 1 3Stream.fold (Fold.drainMapM print) $ Stream.delay 1 input123
valueintersperseM_ :: Monad m => m b -> Stream m a -> Stream m a
#

Insert a side effect before consuming an element of a stream except the first one.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.drain $ Stream.trace putChar $ Stream.intersperseM_ (putChar '.') inputh.e.l.l.o

Pre-release

valuemapMaybeM :: Monad m => (a -> m (Maybe b)) -> Stream m a -> Stream m b
#

Like mapMaybe but maps a monadic function.

Equivalent to:

Example1 expression
mapMaybeM f = Stream.catMaybes . Stream.mapM f
Example1 expression
mapM f = Stream.mapMaybeM (\x -> Just <$> f x)
valueuniq :: (Eq a, Monad m) => Stream m a -> Stream m a
#

Drop repeated elements that are adjacent to each other.

Example1 expression
uniq = Stream.uniqBy (==)
valueprune :: (a -> Bool) -> Stream m a -> Stream m a
#

Strip all leading and trailing occurrences of an element passing a predicate and make all other consecutive occurrences uniq.

> prune p = Stream.dropWhileAround p $ Stream.uniqBy (x y -> p x && p y)
> Stream.prune isSpace (Stream.fromList "  hello      world!   ")
"hello world!"

Space: O(1)

Unimplemented

valuewith
  1. :: Monad m
  2. => Stream m a -> Stream m (s, a)
  3. -> ((s, a) -> b) -> Stream m (s, a) -> Stream m (s, a)
  4. -> (s, a) -> b
  5. -> Stream m a
  6. -> Stream m a
#

Modify a Stream m a -> Stream m a stream transformation that accepts a predicate (a -> b) to accept ((s, a) -> b) instead, provided a transformation Stream m a -> Stream m (s, a). Convenient to filter with index or time.

Example1 expression
filterWithIndex = Stream.with Stream.indexed Stream.filter

Pre-release

valuetrace_ :: Monad m => m b -> Stream m a -> Stream m a
#

Perform a side effect before yielding each element of the stream and discard the results.

Example2 expressions
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.trace_ (print "got here") s"got here""got here"

Same as intersperseMPrefix_ but always serial.

See also: trace

Pre-release

valuescanlMAfter'
  1. :: Monad m
  2. => b -> a -> m b
  3. -> m b
  4. -> b -> m b
  5. -> Stream m a
  6. -> Stream m b
#

scanlMAfter' accumulate initial done stream is like scanlM' except that it provides an additional done function to be applied on the accumulator when the stream stops. The result of done is also emitted in the stream.

This function can be used to allocate a resource in the beginning of the scan and release it when the stream ends or to flush the internal state of the scan at the end.

Pre-release

valuescanl1' :: Monad m => (a -> a -> a) -> Stream m a -> Stream m a
#

Like scanl' but for a non-empty stream. The first element of the stream is used as the initial value of the accumulator. Does nothing if the stream is empty.

Example1 expression
Stream.toList $ Stream.scanl1' (+) $ Stream.fromList [1,2,3,4][1,3,6,10]
valuetakeWhileLast :: (a -> Bool) -> Stream m a -> Stream m a
#

Take all consecutive elements at the end of the stream for which the predicate is true.

O(n) space, where n is the number elements taken.

Unimplemented

valuedropLast :: Int -> Stream m a -> Stream m a
#

Drop n elements at the end of the stream.

O(n) space, where n is the number elements dropped.

Unimplemented

valuedropWhileLast :: (a -> Bool) -> Stream m a -> Stream m a
#

Drop all consecutive elements at the end of the stream for which the predicate is true.

O(n) space, where n is the number elements dropped.

Unimplemented

valueintersperseMWith :: Int -> m a -> Stream m a -> Stream m a
#

Intersperse a monadic action into the input stream after every n elements.

> input = Stream.fromList "hello"
> Stream.fold Fold.toList $ Stream.intersperseMWith 2 (return ',') input

"he,ll,o"

Unimplemented

valueintersperseMSuffixWith :: Monad m => Int -> m a -> Stream m a -> Stream m a
#

Like intersperseMSuffix but intersperses an effectful action into the input stream after every n elements and after the last element.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.intersperseMSuffixWith 2 (return ',') input"he,ll,o,"

Pre-release

valueintersperseMSuffix_ :: Monad m => m b -> Stream m a -> Stream m a
#

Insert a side effect after consuming an element of a stream.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.intersperseMSuffix_ (threadDelay 1000000) input"hello"

Pre-release

valueintersperseMPrefix_ :: Monad m => m b -> Stream m a -> Stream m a
#

Insert a side effect before consuming an element of a stream.

Definition:

Example1 expression
intersperseMPrefix_ m = Stream.mapM (\x -> void m >> return x)
Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.trace putChar $ Stream.intersperseMPrefix_ (putChar '.' >> return ',') input.h.e.l.l.o"hello"

Same as trace_.

Pre-release

valuedelayPre :: MonadIO m => Double -> Stream m a -> Stream m a
#

Introduce a delay of specified seconds before consuming an element of a stream.

Definition:

Example2 expressions
sleep n = liftIO $ threadDelay $ round $ n * 1000000delayPre = Stream.intersperseMPrefix_. sleep

Example:

Example2 expressions
input = Stream.enumerateFromTo 1 3Stream.fold (Fold.drainMapM print) $ Stream.delayPre 1 input123

Pre-release

valuedelayPost :: MonadIO m => Double -> Stream m a -> Stream m a
#

Introduce a delay of specified seconds after consuming an element of a stream.

Definition:

Example2 expressions
sleep n = liftIO $ threadDelay $ round $ n * 1000000delayPost = Stream.intersperseMSuffix_ . sleep

Example:

Example2 expressions
input = Stream.enumerateFromTo 1 3Stream.fold (Fold.drainMapM print) $ Stream.delayPost 1 input123

Pre-release

valuereassembleBy :: Fold m a b -> (a -> a -> Int) -> Stream m a -> Stream m b
#

Buffer until the next element in sequence arrives. The function argument determines the difference in sequence numbers. This could be useful in implementing sequenced streams, for example, TCP reassembly.

Unimplemented

valueindexedR :: Monad m => Int -> Stream m a -> Stream m (Int, a)
#
Example2 expressions
f n = Fold.foldl' (\(i, _) x -> (i - 1, x)) (n + 1,undefined)indexedR n = Stream.postscan (f n)
Example2 expressions
s n = Stream.enumerateFromThen n (n - 1)indexedR n = Stream.zipWith (,) (s n)

Pair each element in a stream with its index, starting from the given index n and counting down.

Example1 expression
Stream.fold Fold.toList $ Stream.indexedR 10 $ Stream.fromList "hello"[(10,'h'),(9,'e'),(8,'l'),(7,'l'),(6,'o')]
valuetimestampWith :: MonadIO m => Double -> Stream m a -> Stream m (AbsTime, a)
#

Pair each element in a stream with an absolute timestamp, using a clock of specified granularity. The timestamp is generated just before the element is consumed.

Example1 expression
Stream.fold Fold.toList $ Stream.timestampWith 0.01 $ Stream.delay 1 $ Stream.enumerateFromTo 1 3[(AbsTime (TimeSpec {sec = ..., nsec = ...}),1),(AbsTime (TimeSpec {sec = ..., nsec = ...}),2),(AbsTime (TimeSpec {sec = ..., nsec = ...}),3)]

Pre-release

valuetimeIndexWith
  1. :: MonadIO m
  2. => Double
  3. -> Stream m a
  4. -> Stream m (RelTime64, a)
#

Pair each element in a stream with relative times starting from 0, using a clock with the specified granularity. The time is measured just before the element is consumed.

Example1 expression
Stream.fold Fold.toList $ Stream.timeIndexWith 0.01 $ Stream.delay 1 $ Stream.enumerateFromTo 1 3[(RelTime64 (NanoSecond64 ...),1),(RelTime64 (NanoSecond64 ...),2),(RelTime64 (NanoSecond64 ...),3)]

Pre-release

valuetimeIndexed :: MonadIO m => Stream m a -> Stream m (RelTime64, a)
#

Pair each element in a stream with relative times starting from 0, using a 10 ms granularity clock. The time is measured just before the element is consumed.

Example1 expression
Stream.fold Fold.toList $ Stream.timeIndexed $ Stream.delay 1 $ Stream.enumerateFromTo 1 3[(RelTime64 (NanoSecond64 ...),1),(RelTime64 (NanoSecond64 ...),2),(RelTime64 (NanoSecond64 ...),3)]

Pre-release

valuerollingMap :: Monad m => (Maybe a -> a -> b) -> Stream m a -> Stream m b
#

Apply a function on every two successive elements of a stream. The first argument of the map function is the previous element and the second argument is the current element. When the current element is the first element, the previous element is Nothing.

Pre-release

valuerollingMap2 :: Monad m => (a -> a -> b) -> Stream m a -> Stream m b
#

Like rollingMap but requires at least two elements in the stream, returns an empty stream otherwise.

This is the stream equivalent of the list idiom zipWith f xs (tail xs).

Pre-release

valuejoinInnerGeneric
  1. :: Monad m
  2. => a -> b -> Bool
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m (a, b)
#

Like cross but emits only those tuples where a == b using the supplied equality predicate.

Definition:

Example1 expression
joinInnerGeneric eq s1 s2 = Stream.filter (\(a, b) -> a `eq` b) $ Stream.cross s1 s2

You should almost always prefer joinInnerOrd over joinInnerGeneric if possible. joinInnerOrd is an order of magnitude faster but may take more space for caching the second stream.

See joinInnerGeneric for a much faster fused alternative.

Time: O(m x n)

Pre-release

valuestrideFromThen :: Monad m => Int -> Int -> Stream m a -> Stream m a
#

strideFromthen offset stride takes the element at offset index and then every element at strides of stride.

Example1 expression
Stream.fold Fold.toList $ Stream.strideFromThen 2 3 $ Stream.enumerateFromTo 0 10[2,5,8]
valuefilterInStreamGenericBy
  1. :: Monad m
  2. => a -> a -> Bool
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

filterInStreamGenericBy retains only those elements in the second stream that are present in the first stream.

Example1 expression
Stream.fold Fold.toList $ Stream.filterInStreamGenericBy (==) (Stream.fromList [1,2,2,4]) (Stream.fromList [2,1,1,3])[2,1,1]
Example1 expression
Stream.fold Fold.toList $ Stream.filterInStreamGenericBy (==) (Stream.fromList [2,1,1,3]) (Stream.fromList [1,2,2,4])[1,2,2]

Similar to the list intersectBy operation but with the stream argument order flipped.

The first stream must be finite and must not block. Second stream is processed only after the first stream is fully realized.

Space: O(n) where n is the number of elements in the second stream.

Time: O(m x n) where m is the number of elements in the first stream and n is the number of elements in the second stream.

Pre-release

valuedeleteInStreamGenericBy
  1. :: Monad m
  2. => a -> a -> Bool
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

Delete all elements of the first stream from the seconds stream. If an element occurs multiple times in the first stream as many occurrences of it are deleted from the second stream.

Example1 expression
Stream.fold Fold.toList $ Stream.deleteInStreamGenericBy (==) (Stream.fromList [1,2,3]) (Stream.fromList [1,2,2])[2]

The following laws hold:

deleteInStreamGenericBy (==) s1 (s1 `append` s2) === s2
deleteInStreamGenericBy (==) s1 (s1 `interleave` s2) === s2

Same as the list Data.List.// operation but with argument order flipped.

The first stream must be finite and must not block. Second stream is processed only after the first stream is fully realized.

Space: O(m) where m is the number of elements in the first stream.

Time: O(m x n) where m is the number of elements in the first stream and n is the number of elements in the second stream.

Pre-release

valueunionWithStreamGenericBy
  1. :: MonadIO m
  2. => a -> a -> Bool
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

This essentially appends to the second stream all the occurrences of elements in the first stream that are not already present in the second stream.

Equivalent to the following except that s2 is evaluated only once:

Example1 expression
unionWithStreamGenericBy eq s1 s2 = s2 `Stream.append` (Stream.deleteInStreamGenericBy eq s2 s1)

Example:

Example1 expression
Stream.fold Fold.toList $ Stream.unionWithStreamGenericBy (==) (Stream.fromList [1,1,2,3]) (Stream.fromList [1,2,2,4])[1,2,2,4,3]

Space: O(n)

Time: O(m x n)

Pre-release

valuenub :: (Monad m, Ord a) => Stream m a -> Stream m a
#

The memory used is proportional to the number of unique elements in the stream. If we want to limit the memory we can just use "take" to limit the uniq elements in the stream.

valuejoinLeftGeneric
  1. :: Monad m
  2. => a -> b -> Bool
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m (a, Maybe b)
#

Like joinInner but emit (a, Just b), and additionally, for those a's that are not equal to any b emit (a, Nothing).

The second stream is evaluated multiple times. If the stream is a consume-once stream then the caller should cache it in an Array before calling this function. Caching may also improve performance if the stream is expensive to evaluate.

Example1 expression
joinRightGeneric eq = flip (Stream.joinLeftGeneric eq)

Space: O(n) assuming the second stream is cached in memory.

Time: O(m x n)

Unimplemented

valuejoinOuterGeneric
  1. :: MonadIO m
  2. => a -> b -> Bool
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m (Maybe a, Maybe b)
#

Like joinLeft but emits a (Just a, Just b). Like joinLeft, for those a's that are not equal to any b emit (Just a, Nothing), but additionally, for those b's that are not equal to any a emit (Nothing, Just b).

For space efficiency use the smaller stream as the second stream.

Space: O(n)

Time: O(m x n)

Pre-release

valuejoinInner
  1. :: (Monad m, Ord k)
  2. => Stream m (k, a)
  3. -> Stream m (k, b)
  4. -> Stream m (k, a, b)
#

Like joinInner but uses a Map for efficiency.

If the input streams have duplicate keys, the behavior is undefined.

For space efficiency use the smaller stream as the second stream.

Space: O(n)

Time: O(m + n)

Pre-release