HORIZON HASKELLDocslts/ghc-9.10.x248f8f02026-10-05Search names, modules, packages, or :: a typeCtrl K

GHC 9.10.3 · lts/ghc-9.10.x · 248f8f0 · 2026-10-05

Modulestreamly-0.10.1Haskell2010

Streamly.Internal.Data.Stream.IsStream

Deprecated. Please use "Streamly.Internal.Data.Stream from streamly-core package", Streamly.Internal.Data.Stream.Concurrent, Streamly.Internal.Data.Stream.Exception.Lifted, & Streamly.Internal.Data.Stream.Time from streamly package instead.

This is an internal module which is a superset of the corresponding released module Streamly.Prelude. It contains some additional unreleased or experimental APIs.

  • 18 types
  • 3 classes
  • 427 values
  • Packagestreamly-0.10.1
  • Exports450
  • LanguageHaskell2010
  • LicenceBSD-3-Clause
  • SourceIsStream.hs
typetype Async = AsyncT IO
#

A demand driven left biased parallely composing IO stream of elements of type a. See AsyncT documentation for more details.

Since: 0.2.0 (Streamly)

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

Constructors

Instances10Functor, Foldable, Traversable, IsList, Read, Show, …
valuemkStream
  1. :: IsStream t
  2. => forall r. State StreamK m a -> (a -> t m a -> m r) -> (a -> m r) -> m r -> m r
  3. -> t m a
#

Build a stream from an SVar, a stop continuation, a singleton stream continuation and a yield continuation.

valuecons :: IsStream t => a -> t m a -> t m a
#

Construct a stream by adding a pure value at the head of an existing stream. For serial streams this is the same as (return a) `consM` r but more efficient. For concurrent streams this is not concurrent whereas consM is concurrent. For example:

> toList $ 1 `cons` 2 `cons` 3 `cons` nil
[1,2,3]
value(.:) :: IsStream t => a -> t m a -> t m a
#

Operator equivalent of cons.

> toList $ 1 .: 2 .: 3 .: nil
[1,2,3]
classclass (forall (m :: Type -> Type) a. MonadAsync m => Semigroup (t m a), forall (m :: Type -> Type) a. MonadAsync m => Monoid (t m a), forall (m :: Type -> Type). Monad m => Functor (t m), forall (m :: Type -> Type). MonadAsync m => Applicative (t m)) => IsStream (t :: (Type -> Type) -> Type -> Type) where
#

Class of types that can represent a stream of elements of some type a in some monad m.

Since: 0.2.0 (Streamly)

Methods

  • consM :: MonadAsync m => m a -> t m a -> t m ainfixr 5

    Constructs a stream by adding a monadic action at the head of an existing stream. For example:

    > toList $ getLine `consM` getLine `consM` nil
    hello
    world
    ["hello","world"]
    

    Concurrent (do not use fromParallel to construct infinite streams)

  • (|:) :: MonadAsync m => m a -> t m a -> t m ainfixr 5

    Operator equivalent of consM. We can read it as "parallel colon" to remember that | comes before :.

    > toList $ getLine |: getLine |: nil
    hello
    world
    ["hello","world"]
    
    let delay = threadDelay 1000000 >> print 1
    drain $ fromSerial  $ delay |: delay |: delay |: nil
    drain $ fromParallel $ delay |: delay |: delay |: nil
    

    Concurrent (do not use fromParallel to construct infinite streams)

Instances8IsStream, …
  • IsStream AheadTDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream AsyncTDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream WAsyncTDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream ParallelTDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream SerialTDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream WSerialTDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream ZipSerialMDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
  • IsStream ZipAsyncMDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Type
valueconcatMapWith
  1. :: IsStream t
  2. => t m b -> t m b -> t m b
  3. -> a -> t m b
  4. -> t m a
  5. -> t m b
#

concatMapWith mixer generator stream is a two dimensional looping combinator. The generator function is used to generate streams from the elements in the input stream and the mixer function is used to merge those streams.

Note we can merge streams concurrently by using a concurrent merge function.

Since: 0.7.0

Since: 0.8.0 (signature change)

valueconcatFoldableWith
  1. :: (IsStream t, Foldable f)
  2. => t m a -> t m a -> t m a
  3. -> f (t m a)
  4. -> t m a
#

A variant of fold that allows you to fold a Foldable container of streams using the specified stream sum operation.

concatFoldableWith async $ map return [1..3]

Equivalent to:

concatFoldableWith f = Prelude.foldr f D.nil
concatFoldableWith f = D.concatMapFoldableWith f id

Since: 0.8.0 (Renamed foldWith to concatFoldableWith)

Since: 0.1.0 (Streamly)

valueconcatMapFoldableWith
  1. :: (IsStream t, Foldable f)
  2. => t m b -> t m b -> t m b
  3. -> a -> t m b
  4. -> f a
  5. -> t m b
#

A variant of foldMap that allows you to map a monadic streaming action on a Foldable container and then fold it using the specified stream merge operation.

concatMapFoldableWith async return [1..3]

Equivalent to:

concatMapFoldableWith f g = Prelude.foldr (f . g) S.nil
concatMapFoldableWith f g xs = S.concatMapWith f g (S.fromFoldable xs)

Since: 0.8.0 (Renamed foldMapWith to concatMapFoldableWith)

Since: 0.1.0 (Streamly)

valueconcatForFoldableWith
  1. :: (IsStream t, Foldable f)
  2. => t m b -> t m b -> t m b
  3. -> f a
  4. -> a -> t m b
  5. -> t m b
#

Like concatMapFoldableWith but with the last two arguments reversed i.e. the monadic streaming function is the last argument.

Equivalent to:

concatForFoldableWith f xs g = Prelude.foldr (f . g) D.nil xs
concatForFoldableWith f = flip (D.concatMapFoldableWith f)

Since: 0.8.0 (Renamed forEachWith to concatForFoldableWith)

Since: 0.1.0 (Streamly)

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

For SerialT streams:

(<>) = Streamly.Prelude.serial                       -- Semigroup
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.serial -- Monad

A single Monad bind behaves like a for loop:

Example1 expression
:{IsStream.toList $ do     x <- IsStream.fromList [1,2] -- foreach x in stream     return x:}[1,2]

Nested monad binds behave like nested for loops:

Example1 expression
:{IsStream.toList $ do    x <- IsStream.fromList [1,2] -- foreach x in stream    y <- IsStream.fromList [3,4] -- foreach y in stream    return (x, y):}[(1,3),(1,4),(2,3),(2,4)]

Since: 0.2.0 (Streamly)

Instances22MonadTrans, IsStream, MonadReader, MonadState, Monad, Functor, …
newtypenewtype WSerialT (m :: Type -> Type) a
#

For WSerialT streams:

(<>) = Streamly.Prelude.wSerial                       -- Semigroup
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.wSerial -- Monad

Note that <> is associative only if we disregard the ordering of elements in the resulting stream.

A single Monad bind behaves like a for loop:

Example1 expression
:{IsStream.toList $ IsStream.fromWSerial $ do     x <- IsStream.fromList [1,2] -- foreach x in stream     return x:}[1,2]

Nested monad binds behave like interleaved nested for loops:

Example1 expression
:{IsStream.toList $ IsStream.fromWSerial $ do    x <- IsStream.fromList [1,2] -- foreach x in stream    y <- IsStream.fromList [3,4] -- foreach y in stream    return (x, y):}[(1,3),(2,3),(1,4),(2,4)]

It is a result of interleaving all the nested iterations corresponding to element 1 in the first stream with all the nested iterations of element 2:

Example2 expressions
import Streamly.Prelude (wSerial)IsStream.toList $ IsStream.fromList [(1,3),(1,4)] `IsStream.wSerial` IsStream.fromList [(2,3),(2,4)][(1,3),(2,3),(1,4),(2,4)]

The W in the name stands for wide or breadth wise scheduling in contrast to the depth wise scheduling behavior of SerialT.

Since: 0.2.0 (Streamly)

Instances22MonadTrans, IsStream, MonadReader, MonadState, Monad, Functor, …
newtypenewtype AheadT (m :: Type -> Type) a
#

For AheadT streams:

(<>) = Streamly.Prelude.ahead
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.ahead

A single Monad bind behaves like a for loop with iterations executed concurrently, ahead of time, producing side effects of iterations out of order, but results in order:

Example1 expression
:{Stream.toList $ Stream.fromAhead $ do     x <- Stream.fromList [2,1] -- foreach x in stream     Stream.fromEffect $ delay x:}1 sec2 sec[2,1]

Nested monad binds behave like nested for loops with nested iterations executed concurrently, ahead of time:

Example1 expression
:{Stream.toList $ Stream.fromAhead $ do    x <- Stream.fromList [1,2] -- foreach x in stream    y <- Stream.fromList [2,4] -- foreach y in stream    Stream.fromEffect $ delay (x + y):}3 sec4 sec5 sec6 sec[3,5,4,6]

The behavior can be explained as follows. All the iterations corresponding to the element 1 in the first stream constitute one output stream and all the iterations corresponding to 2 constitute another output stream and these two output streams are merged using ahead.

Since: 0.3.0 (Streamly)

Instances10IsStream, MonadReader, MonadState, Monad, Functor, Applicative, …
newtypenewtype AsyncT (m :: Type -> Type) a
#

For AsyncT streams:

(<>) = Streamly.Prelude.async
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.async

A single Monad bind behaves like a for loop with iterations of the loop executed concurrently a la the async combinator, producing results and side effects of iterations out of order:

Example1 expression
:{Stream.toList $ Stream.fromAsync $ do     x <- Stream.fromList [2,1] -- foreach x in stream     Stream.fromEffect $ delay x:}1 sec2 sec[1,2]

Nested monad binds behave like nested for loops with nested iterations executed concurrently, a la the async combinator:

Example1 expression
:{Stream.toList $ Stream.fromAsync $ do    x <- Stream.fromList [1,2] -- foreach x in stream    y <- Stream.fromList [2,4] -- foreach y in stream    Stream.fromEffect $ delay (x + y):}3 sec4 sec5 sec6 sec[3,4,5,6]

The behavior can be explained as follows. All the iterations corresponding to the element 1 in the first stream constitute one output stream and all the iterations corresponding to 2 constitute another output stream and these two output streams are merged using async.

Since: 0.1.0 (Streamly)

Instances10IsStream, MonadReader, MonadState, Monad, Functor, Applicative, …
newtypenewtype WAsyncT (m :: Type -> Type) a
#

For WAsyncT streams:

(<>) = Streamly.Prelude.wAsync
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.wAsync

A single Monad bind behaves like a for loop with iterations of the loop executed concurrently a la the wAsync combinator, producing results and side effects of iterations out of order:

Example1 expression
:{Stream.toList $ Stream.fromWAsync $ do     x <- Stream.fromList [2,1] -- foreach x in stream     Stream.fromEffect $ delay x:}1 sec2 sec[1,2]

Nested monad binds behave like nested for loops with nested iterations executed concurrently, a la the wAsync combinator:

Example1 expression
:{Stream.toList $ Stream.fromWAsync $ do    x <- Stream.fromList [1,2] -- foreach x in stream    y <- Stream.fromList [2,4] -- foreach y in stream    Stream.fromEffect $ delay (x + y):}3 sec4 sec5 sec6 sec[3,4,5,6]

The behavior can be explained as follows. All the iterations corresponding to the element 1 in the first stream constitute one WAsyncT output stream and all the iterations corresponding to 2 constitute another WAsyncT output stream and these two output streams are merged using wAsync.

The W in the name stands for wide or breadth wise scheduling in contrast to the depth wise scheduling behavior of AsyncT.

Since: 0.2.0 (Streamly)

Instances10IsStream, MonadReader, MonadState, Monad, Functor, Applicative, …
newtypenewtype ParallelT (m :: Type -> Type) a
#

For ParallelT streams:

(<>) = Streamly.Prelude.parallel
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.parallel

See Streamly.Prelude.AsyncT, ParallelT is similar except that all iterations are strictly concurrent while in AsyncT it depends on the consumer demand and available threads. See parallel for more details.

Since: 0.1.0 (Streamly)

Since: 0.7.0 (maxBuffer applies to ParallelT streams)

Instances10IsStream, MonadReader, MonadState, Monad, Functor, Applicative, …
newtypenewtype ZipSerialM (m :: Type -> Type) a
#

For ZipSerialM streams:

(<>) = Streamly.Prelude.serial
(*) = Streamly.Prelude.serial.zipWith id

Applicative evaluates the streams being zipped serially:

Example4 expressions
s1 = Stream.fromFoldable [1, 2]s2 = Stream.fromFoldable [3, 4]s3 = Stream.fromFoldable [5, 6]Stream.toList $ Stream.fromZipSerial $ (,,) <$> s1 <*> s2 <*> s3[(1,3,5),(2,4,6)]

Since: 0.2.0 (Streamly)

Instances16IsStream, Functor, Applicative, Foldable, Traversable, NFData1, …
newtypenewtype ZipAsyncM (m :: Type -> Type) a
#

For ZipAsyncM streams:

(<>) = Streamly.Prelude.serial
(*) = Streamly.Prelude.serial.zipAsyncWith id

Applicative evaluates the streams being zipped concurrently, the following would take half the time that it would take in serial zipping:

Example2 expressions
s = Stream.fromFoldableM $ Prelude.map delay [1, 1, 1]Stream.toList $ Stream.fromZipAsync $ (,) <$> s <*> s...[(1,1),(1,1),(1,1)]

Since: 0.2.0 (Streamly)

Instances5IsStream, Functor, Applicative, Semigroup, Monoid
typetype Serial = SerialT IO
#

A serial IO stream of elements of type a. See SerialT documentation for more details.

Since: 0.2.0 (Streamly)

typetype WSerial = WSerialT IO
#

An interleaving serial IO stream of elements of type a. See WSerialT documentation for more details.

Since: 0.2.0 (Streamly)

typetype Ahead = AheadT IO
#

A serial IO stream of elements of type a with concurrent lookahead. See AheadT documentation for more details.

Since: 0.3.0 (Streamly)

typetype WAsync = WAsyncT IO
#

A round robin parallely composing IO stream of elements of type a. See WAsyncT documentation for more details.

Since: 0.2.0 (Streamly)

typetype Parallel = ParallelT IO
#

A parallely composing IO stream of elements of type a. See ParallelT documentation for more details.

Since: 0.2.0 (Streamly)

typetype ZipSerial = ZipSerialM IO
#

An IO stream whose applicative instance zips streams serially.

Since: 0.2.0 (Streamly)

typetype ZipAsync = ZipAsyncM IO
#

An IO stream whose applicative instance zips streams wAsyncly.

Since: 0.2.0 (Streamly)

valueadapt :: (IsStream t1, IsStream t2) => t1 m a -> t2 m a
#

Adapt any specific stream type to any other specific stream type.

Since: 0.1.0 (Streamly)

valuefoldStreamShared
  1. :: IsStream t
  2. => State StreamK m a
  3. -> a -> t m a -> m r
  4. -> a -> m r
  5. -> m r
  6. -> t m a
  7. -> m r
#

Fold a stream by providing an SVar, a stop continuation, a singleton continuation and a yield continuation. The stream would share the current SVar passed via the State.

valuefoldStream
  1. :: IsStream t
  2. => State StreamK m a
  3. -> a -> t m a -> m r
  4. -> a -> m r
  5. -> m r
  6. -> t m a
  7. -> m r
#

Fold a stream by providing a State, stop continuation, a singleton continuation and a yield continuation. The stream will not use the SVar passed via State.

valuebindWith
  1. :: IsStream t
  2. => t m b -> t m b -> t m b
  3. -> t m a
  4. -> a -> t m b
  5. -> t m b
#
valuetoConsK
  1. :: IsStream t
  2. => m a -> t m a -> t m a
  3. -> m a
  4. -> StreamK m a
  5. -> StreamK m a
#

Adapt a polymorphic consM operation to a StreamK cons operation

valuefoldlx'
  1. :: (IsStream t, Monad m)
  2. => x -> a -> x
  3. -> x
  4. -> x -> b
  5. -> t m a
  6. -> m b
#

Strict left fold with an extraction function. Like the standard strict left fold, but applies a user supplied extraction function (the third argument) to the folded value at the end. This is designed to work with the foldl library. The suffix x is a mnemonic for extraction.

valuerepeat :: (IsStream t, Monad m) => a -> t m a
#

Generate an infinite stream by repeating a pure value.

valuefromCallback :: MonadAsync m => ((a -> m ()) -> m ()) -> SerialT m a
#

Takes a callback setter function and provides it with a callback. The callback when invoked adds a value at the tail of the stream. Returns a stream of values generated by the callback.

Pre-release

valuecons :: IsStream t => a -> t m a -> t m a
#

Construct a stream by adding a pure value at the head of an existing stream. For serial streams this is the same as (return a) `consM` r but more efficient. For concurrent streams this is not concurrent whereas consM is concurrent. For example:

> toList $ 1 `cons` 2 `cons` 3 `cons` nil
[1,2,3]
value(.:) :: IsStream t => a -> t m a -> t m a
#

Operator equivalent of cons.

> toList $ 1 .: 2 .: 3 .: nil
[1,2,3]
methodconsM :: MonadAsync m => m a -> t m a -> t m a
#

Constructs a stream by adding a monadic action at the head of an existing stream. For example:

> toList $ getLine `consM` getLine `consM` nil
hello
world
["hello","world"]

Concurrent (do not use fromParallel to construct infinite streams)

method(|:) :: MonadAsync m => m a -> t m a -> t m a
#

Operator equivalent of consM. We can read it as "parallel colon" to remember that | comes before :.

> toList $ getLine |: getLine |: nil
hello
world
["hello","world"]
let delay = threadDelay 1000000 >> print 1
drain $ fromSerial  $ delay |: delay |: delay |: nil
drain $ fromParallel $ delay |: delay |: delay |: nil

Concurrent (do not use fromParallel to construct infinite streams)

valueunfold :: (IsStream t, Monad m) => Unfold m a b -> a -> t m b
#

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

Example1 expression
Stream.drain $ Stream.unfold Unfold.replicateM (3, putStrLn "hello")hellohellohello

Since: 0.7.0

valueunfoldr :: (Monad m, IsStream t) => (b -> Maybe (a, b)) -> b -> t m a
#
Example1 expression
:{unfoldr step s =    case step s of        Nothing -> Stream.nil        Just (a, b) -> a `Stream.cons` unfoldr step b:}

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
  1. :: (IsStream t, MonadAsync m)
  2. => b -> m (Maybe (a, b))
  3. -> b
  4. -> t 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]

When run concurrently, the next unfold step can run concurrently with the processing of the output of the previous step. Note that more than one step cannot run concurrently as the next step depends on the output of the previous step.

Example1 expression
:{let f b =        if b > 2        then return Nothing        else threadDelay 1000000 >> return (Just (b, b + 1))in Stream.toList $ Stream.delay 1 $ Stream.fromAsync $ Stream.unfoldrM f 0:}[0,1,2]

Concurrent

Since: 0.1.0

valuefromPure :: IsStream t => a -> t m a
#
fromPure a = a `cons` nil

Create a singleton stream from a pure value.

The following holds in monadic streams, but not in Zip streams:

fromPure = pure
fromPure = fromEffect . pure

In Zip applicative streams fromPure is not the same as pure because in that case pure is equivalent to repeat instead. fromPure and pure are equally efficient, in other cases fromPure may be slightly more efficient than the other equivalent definitions.

Since: 0.8.0 (Renamed yield to fromPure)

valuefromEffect :: (Monad m, IsStream t) => m a -> t m a
#
fromEffect m = m `consM` nil

Create a singleton stream from a monadic action.

> Stream.toList $ Stream.fromEffect getLine
hello
["hello"]

Since: 0.8.0 (Renamed yieldM to fromEffect)

valuerepeatM :: (IsStream t, MonadAsync m) => m a -> t m a
#
Example2 expressions
repeatM = fix . consMrepeatM = cycle1 . fromEffect

Generate a stream by repeatedly executing a monadic action forever.

Example1 expression
:{repeatAsync =       Stream.repeatM (threadDelay 1000000 >> print 1)     & Stream.take 10     & Stream.fromAsync     & Stream.drain:}

Concurrent, infinite (do not use with fromParallel)

valuereplicate :: (IsStream t, Monad m) => Int -> a -> t m a
#
Example1 expression
replicate n = Stream.take n . Stream.repeat

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

valuereplicateM :: (IsStream t, MonadAsync m) => Int -> m a -> t m a
#
Example1 expression
replicateM n = Stream.take n . Stream.repeatM

Generate a stream by performing a monadic action n times. Same as:

Example1 expression
pr n = threadDelay 1000000 >> print n

This runs serially and takes 3 seconds:

Example1 expression
Stream.drain $ Stream.fromSerial $ Stream.replicateM 3 $ pr 1111

This runs concurrently and takes just 1 second:

Example1 expression
Stream.drain $ Stream.fromAsync  $ Stream.replicateM 3 $ pr 1111

Concurrent

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 :: (IsStream t, Monad m) => a -> t 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.

    >>> 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.

    >>> Stream.toList $ Stream.take 4 $ Stream.enumerateFrom 1.1
    [1.1,2.1,3.1,4.1]
    
    
  • enumerateFromTo :: (IsStream t, Monad m) => a -> a -> t 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.

    >>> 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.

    >>> Stream.toList $ Stream.enumerateFromTo 1.1 4
    [1.1,2.1,3.1,4.1]
    
    >>> Stream.toList $ Stream.enumerateFromTo 1.1 4.6
    [1.1,2.1,3.1,4.1,5.1]
    
    
  • enumerateFromThen :: (IsStream t, Monad m) => a -> a -> t 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.

    >>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 2
    [0,2,4,6]
    
    >>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 (-2)
    [0,-2,-4,-6]
    
    
  • enumerateFromThenTo :: (IsStream t, Monad m) => a -> a -> a -> t 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.

    >>> Stream.toList $ Stream.enumerateFromThenTo 0 2 6
    [0,2,4,6]
    
    >>> Stream.toList $ Stream.enumerateFromThenTo 0 (-2) (-6)
    [0,-2,-4,-6]
    
    
Instances21Enumerable, …
  • Enumerable IntegerDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable NaturalDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int16Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int32Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int64Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int8Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word16Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word32Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word64Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word8Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable BoolDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable CharDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable DoubleDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable FloatDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable IntDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable OrderingDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable WordDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable ()Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Integral a => Enumerable (Ratio a)Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable a => Enumerable (Identity a)Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • HasResolution a => Enumerable (Fixed a)Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
valueiterate :: (IsStream t, Monad m) => (a -> a) -> a -> t m a
#
Example1 expression
iterate f x = x `Stream.cons` iterate f x

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]
valueiterateM :: (IsStream t, MonadAsync m) => (a -> m a) -> m a -> t m a
#
Example1 expression
iterateM f m = m >>= \a -> return a `Stream.consM` iterateM f (f 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.

Example2 expressions
pr n = threadDelay 1000000 >> print n:{Stream.iterateM (\x -> pr x >> return (x + 1)) (return 0)    & Stream.take 3    & Stream.fromSerial    & Stream.toList:}01[0,1,2]

When run concurrently, the next iteration can run concurrently with the processing of the previous iteration. Note that more than one iteration cannot run concurrently as the next iteration depends on the output of the previous iteration.

Example1 expression
:{Stream.iterateM (\x -> pr x >> return (x + 1)) (return 0)    & Stream.delay 1    & Stream.take 3    & Stream.fromAsync    & Stream.toList:}01...

Concurrent

Since: 0.1.2

Since: 0.7.0 (signature change)

valuefromIndices :: (IsStream t, Monad m) => (Int -> a) -> t m a
#
Example2 expressions
fromIndices f = fmap f $ Stream.enumerateFrom 0fromIndices f = let g i = f i `Stream.cons` g (i + 1) in g 0

Generate an infinite stream, whose values are the output of a function f applied on the corresponding index. Index starts at 0.

Example1 expression
Stream.toList $ Stream.take 5 $ Stream.fromIndices id[0,1,2,3,4]
valuefromIndicesM :: (IsStream t, MonadAsync m) => (Int -> m a) -> t m a
#
Example2 expressions
fromIndicesM f = Stream.mapM f $ Stream.enumerateFrom 0fromIndicesM f = let g i = f i `Stream.consM` g (i + 1) in g 0

Generate an infinite stream, whose values are the output of a monadic function f applied on the corresponding index. Index starts at 0.

Concurrent

valuefromListM :: (MonadAsync m, IsStream t) => [m a] -> t m a
#
Example4 expressions
fromListM = Stream.fromFoldableMfromListM = Stream.sequence . Stream.fromListfromListM = Stream.mapM id . Stream.fromListfromListM = Prelude.foldr Stream.consM Stream.nil

Construct a stream from a list of monadic actions. This is more efficient than fromFoldableM for serial streams.

valuefromFoldable :: (IsStream t, Foldable f) => f a -> t m a
#
Example1 expression
fromFoldable = Prelude.foldr Stream.cons Stream.nil

Construct a stream from a Foldable containing pure values:

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

Construct a stream from a Foldable containing monadic actions.

Example2 expressions
pr n = threadDelay 1000000 >> print nStream.drain $ Stream.fromSerial $ Stream.fromFoldableM $ map pr [1,2,3]123
Example1 expression
Stream.drain $ Stream.fromAsync $ Stream.fromFoldableM $ map pr [1,2,3].........

Concurrent (do not use with fromParallel on infinite containers)

valueticks :: Rate -> t m ()
#

Generate ticks at the specified rate. The rate is adaptive, the tick generation speed can be increased or decreased at different times to achieve the specified rate. The specific behavior for different styles of Rate specifications is documented under Rate. The effective maximum rate achieved by a stream is governed by the processor speed.

Unimplemented

valuetimes :: (IsStream t, MonadAsync m) => t 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.

Example1 expression
Stream.mapM_ (\x -> print x >> threadDelay 1000000) $ 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

valueabsTimes :: (IsStream t, MonadAsync m, Functor (t m)) => t m AbsTime
#

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

Example1 expression
Stream.mapM_ print $ 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
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t 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.

Example1 expression
Stream.mapM_ print $ Stream.delayPre 1 $ Stream.take 3 $ 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 :: (IsStream t, MonadAsync m, Functor (t m)) => t m RelTime64
#

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

Example1 expression
Stream.mapM_ print $ 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

valuerelTimesWith
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t m RelTime64
#

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.

Example1 expression
Stream.mapM_ print $ 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

valuemfix :: (IsStream t, Monad m) => (m a -> t m a) -> t m a
#

We can define cyclic structures using let:

Example1 expression
let (a, b) = ([1, b], head a) in (a, b)([1,1],1)

The function fix defined as:

Example1 expression
fix f = let x = f x in x

ensures that the argument of a function and its output refer to the same lazy value x i.e. the same location in memory. Thus x can be defined in terms of itself, creating structures with cyclic references.

Example2 expressions
f ~(a, b) = ([1, b], head a)fix f([1,1],1)

mfix is essentially the same as fix but for monadic values.

Using mfix for streams we can construct a stream in which each element of the stream is defined in a cyclic fashion. The argument of the function being fixed represents the current element of the stream which is being returned by the stream monad. Thus, we can use the argument to construct itself.

Pre-release

valuetoList :: Monad m => SerialT m a -> m [a]
#
toList = Stream.foldr (:) []

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.Array instead.

valuefold :: Monad m => Fold m a b -> SerialT 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.

Example1 expression
Stream.fold Fold.sum (Stream.enumerateFromTo 1 100)5050

Folds never fail, therefore, they produce a default value even when no input is provided. It means we can always fold an empty stream and get a valid result. For example:

Example1 expression
Stream.fold Fold.sum Stream.nil0

However, foldMany on an empty stream results in an empty stream. Therefore, Stream.fold f is not the same as Stream.head . Stream.foldMany f.

fold f = Stream.parse (Parser.fromFold f)
valueuncons :: (IsStream t, Monad m) => SerialT m a -> m (Maybe (a, t 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.

This can be used to do pretty much anything in an imperative manner, 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.

valuetail :: (IsStream t, Monad m) => SerialT m a -> m (Maybe (t m a))
#
tail = fmap (fmap snd) . Stream.uncons

Extract all but the first element of the stream, if any.

valuefoldrM :: Monad m => (a -> m b -> m b) -> m b -> SerialT 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:

Example1 expression
Stream.foldrM (\x xs -> if odd x then return True else xs) (return False) $ Stream.fromList (2:4:5:undefined)True

Since: 0.7.0 (signature changed)

Since: 0.2.0 (signature changed)

Since: 0.1.0

valuefoldr :: Monad m => (a -> b -> b) -> b -> SerialT 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.

valuefoldl' :: Monad m => (b -> a -> b) -> b -> SerialT m a -> m b
#

Left associative/strict push fold. foldl' reduce initial stream invokes reduce with the accumulator and the next input in the input stream, using initial as the initial value of the current value of the accumulator. When the input is exhausted the current value of the accumulator is returned. Make sure to use a strict data structure for accumulator to not build unnecessary lazy expressions unless that's what you want. See the previous section for more details.

valuefoldl1' :: Monad m => (a -> a -> a) -> SerialT m a -> m (Maybe a)
#

Strict left fold, for non-empty streams, using first element as the starting value. Returns Nothing if the stream is empty.

valuefoldlM' :: Monad m => (b -> a -> m b) -> m b -> SerialT m a -> m b
#

Like foldl' but with a monadic step function.

Since: 0.2.0

Since: 0.8.0 (signature change)

valuedrain :: Monad m => SerialT m a -> m ()
#
drain = mapM_ (\_ -> return ())
drain = Stream.fold Fold.drain

Run a stream, discarding the results. By default it interprets the stream as SerialT, to run other types of streams use the type adapting combinators for example Stream.drain . fromAsync.

valuelast :: Monad m => SerialT m a -> m (Maybe a)
#

Extract the last element of the stream, if any.

last xs = xs !! (Stream.length xs - 1)
last = Stream.fold Fold.last
valuesum :: (Monad m, Num a) => SerialT m a -> m a
#

Determine the sum of all elements of a stream of numbers. Returns 0 when the stream is empty. Note that this is not numerically stable for floating point numbers.

sum = Stream.fold Fold.sum
valueproduct :: (Monad m, Num a) => SerialT m a -> m a
#

Determine the product of all elements of a stream of numbers. Returns 1 when the stream is empty.

product = Stream.fold Fold.product
valuemaximumBy :: Monad m => (a -> a -> Ordering) -> SerialT m a -> m (Maybe a)
#

Determine the maximum element in a stream using the supplied comparison function.

maximumBy = Stream.fold Fold.maximumBy
valueminimumBy :: Monad m => (a -> a -> Ordering) -> SerialT m a -> m (Maybe a)
#

Determine the minimum element in a stream using the supplied comparison function.

minimumBy = Stream.fold Fold.minimumBy
valuethe :: (Eq a, Monad m) => SerialT m a -> m (Maybe a)
#

Ensures that all the elements of the stream are identical and then returns that unique element.

valuedrainN :: Monad m => Int -> SerialT m a -> m ()
#
drainN n = Stream.drain . Stream.take n
drainN n = Stream.fold (Fold.take n Fold.drain)

Run maximum up to n iterations of a stream.

valuedrainWhile :: Monad m => (a -> Bool) -> SerialT m a -> m ()
#
drainWhile p = Stream.drain . Stream.takeWhile p

Run a stream as long as the predicate holds true.

valuehead :: Monad m => SerialT m a -> m (Maybe a)
#

Extract the first element of the stream, if any.

head = (!! 0)
head = Stream.fold Fold.one
valuefindM :: Monad m => (a -> m Bool) -> SerialT m a -> m (Maybe a)
#

Returns the first element that satisfies the given predicate.

findM = Stream.fold Fold.findM
valuefind :: Monad m => (a -> Bool) -> SerialT m a -> m (Maybe a)
#

Like findM but with a non-monadic predicate.

find p = findM (return . p)
find = Stream.fold Fold.find
valuelookup :: (Monad m, Eq a) => a -> SerialT m (a, b) -> m (Maybe b)
#

In a stream of (key-value) pairs (a, b), return the value b of the first pair where the key equals the given value a.

lookup = snd <$> Stream.find ((==) . fst)
lookup = Stream.fold Fold.lookup
valuefindIndex :: Monad m => (a -> Bool) -> SerialT m a -> m (Maybe Int)
#

Returns the first index that satisfies the given predicate.

findIndex = Stream.fold Fold.findIndex
valueelemIndex :: (Monad m, Eq a) => a -> SerialT m a -> m (Maybe Int)
#

Returns the first index where a given value is found in the stream.

elemIndex a = Stream.findIndex (== a)
valuenull :: Monad m => SerialT m a -> m Bool
#

Determine whether the stream is empty.

null = Stream.fold Fold.null
valueelem :: (Monad m, Eq a) => a -> SerialT m a -> m Bool
#

Determine whether an element is present in the stream.

elem = Stream.fold Fold.elem
valuenotElem :: (Monad m, Eq a) => a -> SerialT m a -> m Bool
#

Determine whether an element is not present in the stream.

notElem = Stream.fold Fold.length
valueall :: Monad m => (a -> Bool) -> SerialT m a -> m Bool
#

Determine whether all elements of a stream satisfy a predicate.

all = Stream.fold Fold.all
valueany :: Monad m => (a -> Bool) -> SerialT m a -> m Bool
#

Determine whether any of the elements of a stream satisfy a predicate.

any = Stream.fold Fold.any
valueand :: Monad m => SerialT m Bool -> m Bool
#

Determines if all elements of a boolean stream are True.

and = Stream.fold Fold.and
valueor :: Monad m => SerialT m Bool -> m Bool
#

Determines whether at least one element of a boolean stream is True.

or = Stream.fold Fold.or
value(|$.) :: (IsStream t, MonadAsync m) => (t m a -> m b) -> t m a -> m b
#

Parallel fold application operator; applies a fold function t m a -> m b to a stream t m a concurrently; The the input stream is evaluated asynchronously in an independent thread yielding elements to a buffer and the folding action runs in another thread consuming the input from the buffer.

If you read the signature as (t m a -> m b) -> (t m a -> m b) you can look at it as a transformation that converts a fold function to a buffered concurrent fold function.

The . at the end of the operator is a mnemonic for termination of the stream.

In the example below, each stage introduces a delay of 1 sec but output is printed every second because both stages are concurrent.

Example3 expressions
import Control.Concurrent (threadDelay)import Streamly.Prelude ((|$.)):{ Stream.foldlM' (\_ a -> threadDelay 1000000 >> print a) (return ())     |$. Stream.replicateM 3 (threadDelay 1000000 >> return 1):}111

Concurrent

Since: 0.3.0 (Streamly)

value(|&.) :: (IsStream t, MonadAsync m) => t m a -> (t m a -> m b) -> m b
#

Same as |$. but with arguments reversed.

(|&.) = flip (|$.)

Concurrent

Since: 0.3.0 (Streamly)

valueeqBy
  1. :: (IsStream t, Monad m)
  2. => a -> b -> Bool
  3. -> t m a
  4. -> t m b
  5. -> m Bool
#

Compare two streams for equality using an equality function.

valueisPrefixOf :: (Eq a, IsStream t, Monad m) => t m a -> t 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" :: SerialT IO Char)True
valueisSubsequenceOf :: (Eq a, IsStream t, Monad m) => t m a -> t 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" :: SerialT IO Char)True
valuestripPrefix
  1. :: (Eq a, IsStream t, Monad m)
  2. => t m a
  3. -> t m a
  4. -> m (Maybe (t m a))
#

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

See also "Streamly.Internal.Data.Stream.IsStream.Nesting.dropPrefix".

Space: O(1)

valuemapM_ :: Monad m => (a -> m b) -> SerialT m a -> m ()
#
mapM_ = Stream.drain . Stream.mapM

Apply a monadic action to each element of the stream and discard the output of the action. This is not really a pure transformation operation but a transformation followed by fold.

valuefoldx :: Monad m => (x -> a -> x) -> x -> (x -> b) -> SerialT m a -> m b
#

Strict left fold with an extraction function. Like the standard strict left fold, but applies a user supplied extraction function (the third argument) to the folded value at the end. This is designed to work with the foldl library. The suffix x is a mnemonic for extraction.

valuefoldxM
  1. :: Monad m
  2. => x -> a -> m x
  3. -> m x
  4. -> x -> m b
  5. -> SerialT m a
  6. -> m b
#

Like foldx, but with a monadic step function.

valuefoldr1 :: Monad m => (a -> a -> a) -> SerialT m a -> m (Maybe a)
#

Lazy right fold for non-empty streams, using first element as the starting value. Returns Nothing if the stream is empty.

valuerunStream :: Monad m => SerialT m a -> m ()
#

Run a stream, discarding the results. By default it interprets the stream as SerialT, to run other types of streams use the type adapting combinators for example runStream . fromAsync.

valuerunN :: Monad m => Int -> SerialT m a -> m ()
#
runN n = runStream . take n

Run maximum up to n iterations of a stream.

valuerunWhile :: Monad m => (a -> Bool) -> SerialT m a -> m ()
#
runWhile p = runStream . takeWhile p

Run a stream as long as the predicate holds true.

valueparse :: Monad m => Parser a m b -> SerialT m a -> m (Either ParseError b)
#

Parse a stream using the supplied Parser.

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:

fold f = Stream.parse (Parser.fromFold f)

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

Pre-release

valuefoldlS :: IsStream t => (t m b -> a -> t m b) -> t m b -> t m a -> t m b
#

Lazy left fold to a stream.

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

Lazy left fold to a transformer monad.

For example, to reverse a stream:

D.toList $ D.foldlT (flip D.cons) D.nil $ (D.fromList [1..5] :: SerialT IO Int)
valuemconcat :: (Monad m, Monoid a) => SerialT m a -> m a
#

Fold a stream of monoid elements by appending them.

mconcat = Stream.fold Fold.mconcat

Pre-release

valueheadElse :: Monad m => a -> SerialT m a -> m a
#

Extract the first element of the stream, if any, otherwise use the supplied default value. It can help avoid one branch in high performance code.

Pre-release

valuetoListRev :: Monad m => SerialT m a -> m [a]
#
toListRev = Stream.foldl' (flip (:)) []

Convert a stream into a list in reverse order in the underlying monad.

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

Pre-release

valuetoStreamRev :: Monad m => SerialT m a -> m (SerialT n a)
#

Convert a stream to a pure stream in reverse order.

toStreamRev = Stream.foldl' (flip Stream.cons) Stream.nil

Pre-release

valueisInfixOf
  1. :: (MonadIO m, Eq a, Enum a, Storable a, Unbox a)
  2. => SerialT m a
  3. -> SerialT 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.

Stream.isInfixOf (Stream.fromList "hello") (Stream.fromList "hello" :: SerialT IO Char)

True

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

Pre-release

Requires Storable constraint

valueisSuffixOf :: (Monad m, Eq a) => SerialT m a -> SerialT 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" :: SerialT 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. => SerialT m a
  3. -> SerialT m a
  4. -> m (Maybe (SerialT 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.

See also "Streamly.Internal.Data.Stream.IsStream.Nesting.dropSuffix".

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

Pre-release

valuemaxThreads :: IsStream t => Int -> t m a -> t m a
#

Specify the maximum number of threads that can be spawned concurrently for any concurrent combinator in a stream. A value of 0 resets the thread limit to default, a negative value means there is no limit. The default value is 1500. maxThreads does not affect ParallelT streams as they can use unbounded number of threads.

When the actions in a stream are IO bound, having blocking IO calls, this option can be used to control the maximum number of in-flight IO requests. When the actions are CPU bound this option can be used to control the amount of CPU used by the stream.

Since: 0.4.0 (Streamly)

valuemaxBuffer :: IsStream t => Int -> t m a -> t m a
#

Specify the maximum size of the buffer for storing the results from concurrent computations. If the buffer becomes full we stop spawning more concurrent tasks until there is space in the buffer. A value of 0 resets the buffer size to default, a negative value means there is no limit. The default value is 1500.

CAUTION! using an unbounded maxBuffer value (i.e. a negative value) coupled with an unbounded maxThreads value is a recipe for disaster in presence of infinite streams, or very large streams. Especially, it must not be used when pure is used in ZipAsyncM streams as pure in applicative zip streams generates an infinite stream causing unbounded concurrent generation with no limit on the buffer or threads.

Since: 0.4.0 (Streamly)

datadata Rate
#

Specifies the stream yield rate in yields per second (Hertz). We keep accumulating yield credits at rateGoal. At any point of time we allow only as many yields as we have accumulated as per rateGoal since the start of time. If the consumer or the producer is slower or faster, the actual rate may fall behind or exceed rateGoal. We try to recover the gap between the two by increasing or decreasing the pull rate from the producer. However, if the gap becomes more than rateBuffer we try to recover only as much as rateBuffer.

rateLow puts a bound on how low the instantaneous rate can go when recovering the rate gap. In other words, it determines the maximum yield latency. Similarly, rateHigh puts a bound on how high the instantaneous rate can go when recovering the rate gap. In other words, it determines the minimum yield latency. We reduce the latency by increasing concurrency, therefore we can say that it puts an upper bound on concurrency.

If the rateGoal is 0 or negative the stream never yields a value. If the rateBuffer is 0 or negative we do not attempt to recover.

Since: 0.5.0 (Streamly)

Constructors

valuerate :: IsStream t => Maybe Rate -> t m a -> t m a
#

Specify the pull rate of a stream. A Nothing value resets the rate to default which is unlimited. When the rate is specified, concurrent production may be ramped up or down automatically to achieve the specified yield rate. The specific behavior for different styles of Rate specifications is documented under Rate. The effective maximum production rate achieved by a stream is governed by:

  • The maxThreads limit

  • The maxBuffer limit

  • The maximum rate that the stream producer can achieve

  • The maximum rate that the stream consumer can achieve

Since: 0.5.0 (Streamly)

valueavgRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate (r/2) r (2*r) maxBound)

Specifies the average production rate of a stream in number of yields per second (i.e. Hertz). Concurrent production is ramped up or down automatically to achieve the specified average yield rate. The rate can go down to half of the specified rate on the lower side and double of the specified rate on the higher side.

Since: 0.5.0 (Streamly)

valueminRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate r r (2*r) maxBound)

Specifies the minimum rate at which the stream should yield values. As far as possible the yield rate would never be allowed to go below the specified rate, even though it may possibly go above it at times, the upper limit is double of the specified rate.

Since: 0.5.0 (Streamly)

valuemaxRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate (r/2) r r maxBound)

Specifies the maximum rate at which the stream should yield values. As far as possible the yield rate would never be allowed to go above the specified rate, even though it may possibly go below it at times, the lower limit is half of the specified rate. This can be useful in applications where certain resource usage must not be allowed to go beyond certain limits.

Since: 0.5.0 (Streamly)

valueconstRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate r r r 0)

Specifies a constant yield rate. If for some reason the actual rate goes above or below the specified rate we do not try to recover it by increasing or decreasing the rate in future. This can be useful in applications like graphics frame refresh where we need to maintain a constant refresh rate.

Since: 0.5.0 (Streamly)

valuemap :: (IsStream t, Monad m) => (a -> b) -> t m a -> t m b
#
map = fmap

Same as fmap.

> D.toList $ D.map (+1) $ D.fromList [1,2,3]
[2,3,4]
valuesequence :: (IsStream t, MonadAsync m) => t m (m a) -> t m a
#
sequence = mapM id

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

>>> drain $ Stream.sequence $ Stream.fromList [putStr "a", putStr "b", putStrLn "c"]
abc

>>> :{
drain $ Stream.replicateM 3 (return $ threadDelay 1000000 >> print 1)
 & (fromSerial . Stream.sequence)
:}
1
1
1

>>> :{
drain $ Stream.replicateM 3 (return $ threadDelay 1000000 >> print 1)
 & (fromAsync . Stream.sequence)
:}
1
1
1

Concurrent (do not use with fromParallel on infinite streams)

valuemapM :: (IsStream t, MonadAsync m) => (a -> m b) -> t m a -> t m b
#
mapM f = sequence . map f

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

>>> drain $ Stream.mapM putStr $ Stream.fromList ["a", "b", "c"]
abc

>>> :{
   drain $ Stream.replicateM 10 (return 1)
     & (fromSerial . Stream.mapM (x -> threadDelay 1000000 >> print x))
:}
1
...
1

> drain $ Stream.replicateM 10 (return 1)
 & (fromAsync . Stream.mapM (x -> threadDelay 1000000 >> print x))

Concurrent (do not use with fromParallel on infinite streams)

valuetrace :: (IsStream t, MonadAsync m) => (a -> m b) -> t m a -> t m a
#

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

Example1 expression
Stream.drain $ Stream.trace print (Stream.enumerateFromTo 1 2)12

Compare with tap.

valuetap :: (IsStream t, Monad m) => Fold m a b -> t m a -> t 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-----

Example1 expression
Stream.drain $ Stream.tap (Fold.drainBy print) (Stream.enumerateFromTo 1 2)12

Compare with trace.

valuedelay :: (IsStream t, MonadIO m) => Double -> t m a -> t m a
#

Introduce a delay of specified seconds before consuming an element of the stream except the first one.

Example1 expression
Stream.mapM_ print $ Stream.timestamped $ Stream.delay 1 $ Stream.enumerateFromTo 1 3(AbsTime (TimeSpec {sec = ..., nsec = ...}),1)(AbsTime (TimeSpec {sec = ..., nsec = ...}),2)(AbsTime (TimeSpec {sec = ..., nsec = ...}),3)
valuescanl' :: (IsStream t, Monad m) => (b -> a -> b) -> b -> t m a -> t m b
#

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

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

>>> 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':

>>> Stream.foldl' ((s, p) x -> (s + x, p * x)) (0,1) $ Stream.fromList 1,2,3,4

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:

>>> :{
  Stream.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.

Example3 expressions
scanl' step z = scan (Fold.foldl' step z)scanl' f z xs = scanlM' (\a b -> return (f a b)) (return z) xsscanl' f z xs = z `Stream.cons` postscanl' f z xs

See also: usingStateT

valuescanlM' :: (IsStream t, Monad m) => (b -> a -> m b) -> m b -> t m a -> t m b
#

Like scanl' but with a monadic step function and a monadic seed.

Since: 0.4.0

Since: 0.8.0 (signature change)

valuepostscanl' :: (IsStream t, Monad m) => (b -> a -> b) -> b -> t m a -> t m b
#

Like scanl' but does not stream the initial value of the accumulator.

Example3 expressions
postscanl' step z = postscan (Fold.foldl' step z)postscanl' f z = postscanlM' (\a b -> return (f a b)) (return z)postscanl' f z xs = Stream.drop 1 $ Stream.scanl' f z xs
valuepostscanlM'
  1. :: (IsStream t, Monad m)
  2. => b -> a -> m b
  3. -> m b
  4. -> t m a
  5. -> t m b
#

Like postscanl' but with a monadic step function and a monadic seed.

Example1 expression
postscanlM' f z xs = Stream.drop 1 $ Stream.scanlM' f z xs

Since: 0.7.0

Since: 0.8.0 (signature change)

valuescanl1' :: (IsStream t, Monad m) => (a -> a -> a) -> t m a -> t 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.

>>> Stream.toList $ Stream.scanl1' (+) $ fromList [1,2,3,4]
[1,3,6,10]

valuescan :: (IsStream t, Monad m) => Fold m a b -> t m a -> t m b
#

Scan a stream using the given monadic fold.

Example1 expression
Stream.toList $ Stream.takeWhile (< 10) $ Stream.scan Fold.sum (Stream.fromList [1..10])[0,1,3,6]
valuepostscan :: (IsStream t, Monad m) => Fold m a b -> t m a -> t 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:

Example3 expressions
import Data.Maybe (fromJust)let avg = Fold.teeWith (/) Fold.sum (fmap fromIntegral Fold.length):{ Stream.toList  $ Stream.map (fromJust . fst)  $ Stream.takeWhile (\(_,x) -> x <= 10)  $ Stream.postscan (Fold.tee Fold.last avg) (Stream.enumerateFromTo 1.0 100.0):}[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]
valuedeleteBy :: (IsStream t, Monad m) => (a -> a -> Bool) -> a -> t m a -> t m a
#

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

>>> Stream.toList $ Stream.deleteBy (==) 3 $ Stream.fromList [1,3,3,5]
[1,3,5]

valuefilter :: (IsStream t, Monad m) => (a -> Bool) -> t m a -> t m a
#

Include only those elements that pass a predicate.

valueuniq :: (Eq a, IsStream t, Monad m) => t m a -> t m a
#

Drop repeated elements that are adjacent to each other.

valuetake :: (IsStream t, Monad m) => Int -> t m a -> t m a
#

Take first n elements from the stream and discard the rest.

valuetakeWhile :: (IsStream t, Monad m) => (a -> Bool) -> t m a -> t m a
#

End the stream as soon as the predicate fails on an element.

valuedrop :: (IsStream t, Monad m) => Int -> t m a -> t m a
#

Discard first n elements from the stream and take the rest.

valuedropWhile :: (IsStream t, Monad m) => (a -> Bool) -> t m a -> t m a
#

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

valueinsertBy
  1. :: (IsStream t, Monad m)
  2. => a -> a -> Ordering
  3. -> a
  4. -> t m a
  5. -> t m a
#

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

insertBy cmp x = mergeBy cmp (fromPure x)
>>> Stream.toList $ Stream.insertBy compare 2 $ Stream.fromList [1,3,5]
[1,2,3,5]

valueintersperseM :: (IsStream t, MonadAsync m) => m a -> t m a -> t m a
#

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

Example1 expression
Stream.toList $ Stream.trace putChar $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.fromList "hello"h.,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.toList $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.trace putChar $ Stream.fromList "hello"he.l.l.o."h,e,l,l,o"
valueintersperse :: (IsStream t, MonadAsync m) => a -> t m a -> t m a
#

Insert a pure value between successive elements of a stream.

Example1 expression
Stream.toList $ Stream.intersperse ',' $ Stream.fromList "hello""h,e,l,l,o"
valuereverse :: (IsStream t, Monad m) => t m a -> t 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.

Example1 expression
reverse = Stream.foldlT (flip Stream.cons) Stream.nil

Since 0.7.0 (Monad m constraint)

Since: 0.1.1

valueindexed :: (IsStream t, Monad m) => t m a -> t m (Int, a)
#
indexed = Stream.postscanl' (\(i, _) x -> (i + 1, x)) (-1,undefined)
indexed = Stream.zipWith (,) (Stream.enumerateFrom 0)

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

Example1 expression
Stream.toList $ Stream.indexed $ Stream.fromList "hello"[(0,'h'),(1,'e'),(2,'l'),(3,'l'),(4,'o')]
valueindexedR :: (IsStream t, Monad m) => Int -> t m a -> t m (Int, a)
#
indexedR n = Stream.postscanl' (\(i, _) x -> (i - 1, x)) (n + 1,undefined)
indexedR n = Stream.zipWith (,) (Stream.enumerateFromThen n (n - 1))

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

Example1 expression
Stream.toList $ Stream.indexedR 10 $ Stream.fromList "hello"[(10,'h'),(9,'e'),(8,'l'),(7,'l'),(6,'o')]
valuefindIndices :: (IsStream t, Monad m) => (a -> Bool) -> t m a -> t m Int
#

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

findIndices = fold Fold.findIndices
valueelemIndices :: (IsStream t, Eq a, Monad m) => a -> t m a -> t m Int
#

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

elemIndices a = findIndices (== a)
valuemapMaybe :: (IsStream t, Monad m) => (a -> Maybe b) -> t m a -> t 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:

mapMaybe f = Stream.map fromJust . Stream.filter isJust . Stream.map f
value(|$) :: (IsStream t, MonadAsync m) => (t m a -> t m b) -> t m a -> t m b
#

Parallel transform application operator; applies a stream transformation function t m a -> t m b to a stream t m a concurrently; the input stream is evaluated asynchronously in an independent thread yielding elements to a buffer and the transformation function runs in another thread consuming the input from the buffer. |$ is just like regular function application operator $ except that it is concurrent.

If you read the signature as (t m a -> t m b) -> (t m a -> t m b) you can look at it as a transformation that converts a transform function to a buffered concurrent transform function.

The following code prints a value every second even though each stage adds a 1 second delay.

Example1 expression
:{Stream.drain $   Stream.mapM (\x -> threadDelay 1000000 >> print x)     |$ Stream.replicateM 3 (threadDelay 1000000 >> return 1):}111

Concurrent

Since: 0.3.0 (Streamly)

value(|&) :: (IsStream t, MonadAsync m) => t m a -> (t m a -> t m b) -> t m b
#

Same as |$ but with arguments reversed.

(|&) = flip (|$)

Concurrent

Since: 0.3.0 (Streamly)

valuemkAsync :: (IsStream t, MonadAsync m) => t m a -> t m a
#

Make the stream producer and consumer run concurrently by introducing a buffer between them. The producer thread evaluates the input stream until the buffer fills, it terminates if the buffer is full and a worker thread is kicked off again to evaluate the remaining stream when there is space in the buffer. The consumer consumes the stream lazily from the buffer.

Since: 0.2.0 (Streamly)

valuescanx
  1. :: (IsStream t, Monad m)
  2. => x -> a -> x
  3. -> x
  4. -> x -> b
  5. -> t m a
  6. -> t m b
#

Strict left scan with an extraction function. Like scanl', but applies a user supplied extraction function (the third argument) at each step. This is designed to work with the foldl library. The suffix x is a mnemonic for extraction.

Since 0.2.0

Since: 0.7.0 (Monad m constraint)

valuetapAsyncK :: (IsStream t, MonadAsync m) => (t m a -> m b) -> t m a -> t m a
#

Like tapAsyncF but uses a stream fold function instead of a Fold type.

Pre-release

valuetakeLastInterval :: Double -> t m a -> t m a
#

Take time interval i seconds at the end of the stream.

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

Unimplemented

valuedropLastInterval :: Int -> t m a -> t m a
#

Drop time interval i seconds at the end of the stream.

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

Unimplemented

valuescanlMAfter'
  1. :: (IsStream t, Monad m)
  2. => b -> a -> m b
  3. -> m b
  4. -> b -> m b
  5. -> t m a
  6. -> t 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

valuetimeIndexed
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => t m a
  3. -> t 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.mapM_ print $ Stream.timeIndexed $ Stream.delay 1 $ Stream.enumerateFromTo 1 3(RelTime64 (NanoSecond64 ...),1)(RelTime64 (NanoSecond64 ...),2)(RelTime64 (NanoSecond64 ...),3)

Pre-release

valuemkParallel :: (IsStream t, MonadAsync m) => t m a -> t m a
#

Make the stream producer and consumer run concurrently by introducing a buffer between them. The producer thread evaluates the input stream until the buffer fills, it blocks if the buffer is full until there is space in the buffer. The consumer consumes the stream lazily from the buffer.

mkParallel = IsStream.fromStreamD . mkParallelD . IsStream.toStreamD

Pre-release

valuefoldrS :: IsStream t => (a -> t m b -> t m b) -> t m b -> t m a -> t m b
#

Right fold to a streaming monad.

foldrS Stream.cons Stream.nil === id

foldrS can be used to perform stateless stream to stream transformations like map and filter in general. It can be coupled with a scan to perform stateful transformations. However, note that the custom map and filter routines can be much more efficient than this due to better stream fusion.

Example1 expression
Stream.toList $ Stream.foldrS Stream.cons Stream.nil $ Stream.fromList [1..5][1,2,3,4,5]

Find if any element in the stream is True:

Example1 expression
Stream.toList $ Stream.foldrS (\x xs -> if odd x then (Stream.fromPure True) else xs) (Stream.fromPure False) $ (Stream.fromList (2:4:5:undefined) :: Stream.SerialT IO Int)[True]

Map (+2) on odd elements and filter out the even elements:

Example1 expression
Stream.toList $ Stream.foldrS (\x xs -> if odd x then (x + 2) `Stream.cons` xs else xs) Stream.nil $ (Stream.fromList [1..5] :: Stream.SerialT IO Int)[3,5,7]

foldrM can also be represented in terms of foldrS, however, the former is much more efficient:

foldrM f z s = runIdentityT $ foldrS (\x xs -> lift $ f x (runIdentityT xs)) (lift z) s

Pre-release

valuefoldrT
  1. :: (IsStream t, Monad m, Monad (s m), MonadTrans s)
  2. => a -> s m b -> s m b
  3. -> s m b
  4. -> t m a
  5. -> s 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:

foldrS = foldrT
foldrM f z s = runIdentityT $ foldrT (\x xs -> lift $ f x (runIdentityT xs)) (lift z) s

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

Pre-release

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

A stateful mapM, equivalent to a left scan, more like mapAccumL. Hopefully, this is a better alternative to scan. Separation of state from the output makes it easier to think in terms of a shared state, and also makes it easier to keep the state fully strict and the output lazy.

See also: scanlM'

Pre-release

valuetrace_ :: (IsStream t, Monad m) => m b -> t m a -> t m a
#

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

>>> Stream.drain $ Stream.trace_ (print "got here") (Stream.enumerateFromTo 1 2)
"got here"
"got here"

Same as intersperseMPrefix_ but always serial.

See also: trace

Pre-release

valuetapOffsetEvery
  1. :: (IsStream t, Monad m)
  2. => Int
  3. -> Int
  4. -> Fold m a b
  5. -> t m a
  6. -> t m a
#

tapOffsetEvery offset n taps every nth element in the stream starting at offset. offset can be between 0 and n - 1. Offset 0 means start at the first element in the stream. If the offset is outside this range then offset mod n is used as offset.

Example1 expression
Stream.drain $ Stream.tapOffsetEvery 0 2 (Fold.rmapM print Fold.toList) $ Stream.enumerateFromTo 0 10[0,2,4,6,8,10]
valuetapAsync :: (IsStream t, MonadAsync m) => Fold m a b -> t m a -> t m a
#

Redirect a copy of the stream to a supplied fold and run it concurrently in an independent thread. The fold may buffer some elements. The buffer size is determined by the prevailing maxBuffer setting.

              Stream m a -> m b
                      |
-----stream m a ---------------stream m a-----

>>> Stream.drain $ Stream.tapAsync (Fold.drainBy print) (Stream.enumerateFromTo 1 2)
1
2

Exceptions from the concurrently running fold are propagated to the current computation. Note that, because of buffering in the fold, exceptions may be delayed and may not correspond to the current element being processed in the parent stream, but we guarantee that before the parent stream stops the tap finishes and all exceptions from it are drained.

Example1 expression
tapAsync f = Stream.tapAsyncK (Stream.fold f . Stream.adapt)

Compare with tap.

Pre-release

valuedistributeAsync_
  1. :: (Foldable f, IsStream t, MonadAsync m)
  2. => f (t m a -> m b)
  3. -> t m a
  4. -> t m a
#

Concurrently distribute a stream to a collection of fold functions, discarding the outputs of the folds.

> Stream.drain $ Stream.distributeAsync_ [Stream.mapM_ print, Stream.mapM_ print] (Stream.enumerateFromTo 1 2)
1
2
1
2

distributeAsync_ = flip (foldr tapAsync)

Pre-release

valuepollCounts
  1. :: (IsStream t, MonadAsync m)
  2. => a -> Bool
  3. -> t m Int -> m b
  4. -> t m a
  5. -> t m a
#

pollCounts predicate transform fold stream counts those elements in the stream that pass the predicate. The resulting count stream is sent to another thread which transforms it using transform and then folds it using fold. The thread is automatically cleaned up if the stream stops or aborts due to exception.

For example, to print the count of elements processed every second:

> Stream.drain $ Stream.pollCounts (const True) (Stream.rollingMap (-) . Stream.delayPost 1) (FLold.drainBy print)
          $ Stream.enumerateFrom 0

Note: This may not work correctly on 32-bit machines.

Pre-release

valuescanMany :: (IsStream t, Monad m) => Fold m a b -> t m a -> t m b
#

Like scan but restarts scanning afresh when the scanning fold terminates.

Pre-release

valueprescanl' :: (IsStream t, Monad m) => (b -> a -> b) -> b -> t m a -> t m b
#

Like scanl' but does not stream the final value of the accumulator.

Pre-release

valueprescanlM'
  1. :: (IsStream t, Monad m)
  2. => b -> a -> m b
  3. -> m b
  4. -> t m a
  5. -> t m b
#

Like prescanl' but with a monadic step function and a monadic seed.

Pre-release

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

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

filterWithIndex = with indexed filter
filterWithAbsTime = with timestamped filter
filterWithRelTime = with timeIndexed filter

Pre-release

valueuniqBy
  1. :: (IsStream t, Monad m, Functor (t m))
  2. => a -> a -> Bool
  3. -> t m a
  4. -> t m a
#

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

@uniq = uniqBy (==)

To strip duplicate path separators:

f x y = x == / && x == y
Stream.toList $ Stream.uniqBy f $ Stream.fromList "/a/b"
"ab"

Space: O(1)

See also: nubBy.

Pre-release

valuenubBy :: (a -> a -> Bool) -> t m a -> t m a
#

Drop repeated elements anywhere in the stream.

Caution: not scalable for infinite streams

See also: nubWindowBy

Unimplemented

valueprune :: (a -> Bool) -> t m a -> t m a
#

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

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

Space: O(1)

Unimplemented

valuerepeated :: t m a -> t m a
#

Emit only repeated elements, once.

Unimplemented

valuetakeLast :: Int -> t m a -> t m a
#

Take n elements at the end of the stream.

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

Unimplemented

valuetakeWhileLast :: (a -> Bool) -> t m a -> t 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 -> t m a -> t 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) -> t m a -> t 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 -> t m a -> t m a
#

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

> Stream.toList $ Stream.intersperseMWith 2 (return ',') $ Stream.fromList "hello"
"he,ll,o"

Unimplemented

valueintersperseMSuffix :: (IsStream t, Monad m) => m a -> t m a -> t m a
#

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

Example1 expression
Stream.toList $ Stream.trace putChar $ intersperseMSuffix (putChar '.' >> return ',') $ Stream.fromList "hello"h.,e.,l.,l.,o.,"h,e,l,l,o,"

Pre-release

valueintersperseMSuffixWith
  1. :: (IsStream t, Monad m)
  2. => Int
  3. -> m a
  4. -> t m a
  5. -> t m a
#

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

Example1 expression
Stream.toList $ Stream.intersperseMSuffixWith 2 (return ',') $ Stream.fromList "hello""he,ll,o,"

Pre-release

valueinterjectSuffix
  1. :: (IsStream t, MonadAsync m)
  2. => Double
  3. -> m a
  4. -> t m a
  5. -> t m a
#

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

> import Control.Concurrent (threadDelay)
> Stream.drain $ Stream.interjectSuffix 1 (putChar ',') $ Stream.mapM (x -> threadDelay 1000000 >> putChar x) $ Stream.fromList "hello"
h,e,l,l,o

Pre-release

valueintersperseM_ :: (IsStream t, Monad m) => m b -> t m a -> t m a
#

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

Example1 expression
Stream.drain $ Stream.trace putChar $ Stream.intersperseM_ (putChar '.') $ Stream.fromList "hello"h.e.l.l.o

Pre-release

valueintersperseMSuffix_ :: (IsStream t, Monad m) => m b -> t m a -> t m a
#

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

>>> Stream.mapM_ putChar $ Stream.intersperseMSuffix_ (threadDelay 1000000) $ Stream.fromList "hello"
hello

Pre-release

valuedelayPost :: (IsStream t, MonadIO m) => Double -> t m a -> t m a
#

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

Example1 expression
Stream.mapM_ print $ Stream.timestamped $ Stream.delayPost 1 $ Stream.enumerateFromTo 1 3(AbsTime (TimeSpec {sec = ..., nsec = ...}),1)(AbsTime (TimeSpec {sec = ..., nsec = ...}),2)(AbsTime (TimeSpec {sec = ..., nsec = ...}),3)

Pre-release

valueintersperseMPrefix_ :: (IsStream t, MonadAsync m) => m b -> t m a -> t m a
#

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

Example1 expression
Stream.toList $ Stream.trace putChar $ Stream.intersperseMPrefix_ (putChar '.' >> return ',') $ Stream.fromList "hello".h.e.l.l.o"hello"

Same as trace_ but may be concurrent.

Concurrent

Pre-release

valuedelayPre :: (IsStream t, MonadIO m) => Double -> t m a -> t m a
#

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

Example1 expression
Stream.mapM_ print $ Stream.timestamped $ Stream.delayPre 1 $ Stream.enumerateFromTo 1 3(AbsTime (TimeSpec {sec = ..., nsec = ...}),1)(AbsTime (TimeSpec {sec = ..., nsec = ...}),2)(AbsTime (TimeSpec {sec = ..., nsec = ...}),3)

Pre-release

valuereassembleBy :: Fold m a b -> (a -> a -> Int) -> t m a -> t 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

valuetimestampWith
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t m a
  4. -> t 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.mapM_ print $ 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. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t m a
  4. -> t 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.mapM_ print $ Stream.timeIndexWith 0.01 $ Stream.delay 1 $ Stream.enumerateFromTo 1 3(RelTime64 (NanoSecond64 ...),1)(RelTime64 (NanoSecond64 ...),2)(RelTime64 (NanoSecond64 ...),3)

Pre-release

valuerollingMap :: (IsStream t, Monad m) => (Maybe a -> a -> b) -> t m a -> t 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 :: (IsStream t, Monad m) => (a -> a -> b) -> t m a -> t 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

valueboth :: Functor (t m) => t m (Either a a) -> t m a
#

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

Pre-release

valuesampleOld :: Int -> t m a -> t m a
#

Evaluate the input stream continuously and keep only the oldest n elements in the buffer, discard the new ones when the buffer is full. When the output stream is evaluated it consumes the values from the buffer in a FIFO manner.

Unimplemented

valuesampleNew :: Int -> t m a -> t m a
#

Evaluate the input stream continuously and keep only the latest n elements in a ring buffer, keep discarding the older ones to make space for the new ones. When the output stream is evaluated it consumes the values from the buffer in a FIFO manner.

Unimplemented

valuesampleRate :: Double -> t m a -> t m a
#

Like sampleNew but samples at uniform intervals to match the consumer rate. Note that sampleNew leads to non-uniform sampling depending on the consumer pattern.

Unimplemented

valueinspectMode :: IsStream t => t m a -> t m a
#

Print debug information about an SVar when the stream ends

Pre-release

valueserial :: IsStream t => t m a -> t m a -> t m a
#

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

Example4 expressions
import Streamly.Prelude (serial)stream1 = Stream.fromList [1,2]stream2 = Stream.fromList [3,4]Stream.toList $ stream1 `serial` stream2[1,2,3,4]

This operation can be used to fold an infinite lazy container of streams.

Since: 0.2.0 (Streamly)

valuewSerial :: IsStream t => t m a -> t m a -> t m a
#

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.

Example4 expressions
import Streamly.Prelude (wSerial)stream1 = Stream.fromList [1,2]stream2 = Stream.fromList [3,4]Stream.toList $ Stream.fromWSerial $ stream1 `wSerial` stream2[1,3,2,4]

Note, for singleton streams wSerial and serial are identical.

Note that this operation cannot be used to fold a container of infinite streams but it can be used for very large streams as the state that it needs to maintain is proportional to the logarithm of the number of streams.

Since: 0.2.0 (Streamly)

valueahead :: (IsStream t, MonadAsync m) => t m a -> t m a -> t m a
#

Appends two streams, both the streams may be evaluated concurrently but the outputs are used in the same order as the corresponding actions in the original streams, side effects will happen in the order in which the streams are evaluated:

Example4 expressions
import Streamly.Prelude (ahead, SerialT)stream1 = Stream.fromEffect (delay 4) :: SerialT IO Intstream2 = Stream.fromEffect (delay 2) :: SerialT IO IntStream.toList $ stream1 `ahead` stream2 :: IO [Int]2 sec4 sec[4,2]

Multiple streams can be combined. With enough threads, all of them can be scheduled simultaneously:

Example2 expressions
stream3 = Stream.fromEffect (delay 1)Stream.toList $ stream1 `ahead` stream2 `ahead` stream31 sec2 sec4 sec[4,2,1]

With 2 threads, only two can be scheduled at a time, when one of those finishes, the third one gets scheduled:

Example1 expression
Stream.toList $ Stream.maxThreads 2 $ stream1 `ahead` stream2 `ahead` stream32 sec1 sec4 sec[4,2,1]

Only streams are scheduled for ahead evaluation, how actions within a stream are evaluated depends on the stream type. If it is a concurrent stream they will be evaluated concurrently. It may not make much sense combining serial streams using ahead.

ahead can be safely used to fold an infinite lazy container of streams.

Since: 0.3.0 (Streamly)

valueasync :: (IsStream t, MonadAsync m) => t m a -> t m a -> t m a
#

Merges two streams, both the streams may be evaluated concurrently, outputs from both are used as they arrive:

Example4 expressions
import Streamly.Prelude (async)stream1 = Stream.fromEffect (delay 4)stream2 = Stream.fromEffect (delay 2)Stream.toList $ stream1 `async` stream22 sec4 sec[2,4]

Multiple streams can be combined. With enough threads, all of them can be scheduled simultaneously:

Example2 expressions
stream3 = Stream.fromEffect (delay 1)Stream.toList $ stream1 `async` stream2 `async` stream3...[1,2,4]

With 2 threads, only two can be scheduled at a time, when one of those finishes, the third one gets scheduled:

Example1 expression
Stream.toList $ Stream.maxThreads 2 $ stream1 `async` stream2 `async` stream3...[2,1,4]

With a single thread, it becomes serial:

Example1 expression
Stream.toList $ Stream.maxThreads 1 $ stream1 `async` stream2 `async` stream3...[4,2,1]

Only streams are scheduled for async evaluation, how actions within a stream are evaluated depends on the stream type. If it is a concurrent stream they will be evaluated concurrently.

In the following example, both the streams are scheduled for concurrent evaluation but each individual stream is evaluated serially:

Example3 expressions
stream1 = Stream.fromListM $ Prelude.map delay [3,3] -- SerialT IO Intstream2 = Stream.fromListM $ Prelude.map delay [1,1] -- SerialT IO IntStream.toList $ stream1 `async` stream2 -- IO [Int]...[1,1,3,3]

If total threads are 2, the third stream is scheduled only after one of the first two has finished:

stream3 = Stream.fromListM $ Prelude.map delay [2,2] -- SerialT IO Int
Stream.toList $ Stream.maxThreads 2 $ stream1 `async` stream2 `async` stream3 -- IO [Int]

... [1,1,3,2,3,2]

Thus async goes deep in first few streams rather than going wide in all streams. It prefers to evaluate the leftmost streams as much as possible. Because of this behavior, async can be safely used to fold an infinite lazy container of streams.

Since: 0.2.0 (Streamly)

valuewAsync :: (IsStream t, MonadAsync m) => t m a -> t m a -> t m a
#

For singleton streams, wAsync is the same as async. See async for singleton stream behavior. For multi-element streams, while async is left biased i.e. it tries to evaluate the left side stream as much as possible, wAsync tries to schedule them both fairly. In other words, async goes deep while wAsync goes wide. However, outputs are always used as they arrive.

With a single thread, async starts behaving like serial while wAsync starts behaving like wSerial.

Example4 expressions
import Streamly.Prelude (async, wAsync)stream1 = Stream.fromList [1,2,3]stream2 = Stream.fromList [4,5,6]Stream.toList $ Stream.fromAsync $ Stream.maxThreads 1 $ stream1 `async` stream2[1,2,3,4,5,6]
Example1 expression
Stream.toList $ Stream.fromWAsync $ Stream.maxThreads 1 $ stream1 `wAsync` stream2[1,4,2,5,3,6]

With two threads available, and combining three streams:

Example2 expressions
stream3 = Stream.fromList [7,8,9]Stream.toList $ Stream.fromAsync $ Stream.maxThreads 2 $ stream1 `async` stream2 `async` stream3[1,2,3,4,5,6,7,8,9]
Example1 expression
Stream.toList $ Stream.fromWAsync $ Stream.maxThreads 2 $ stream1 `wAsync` stream2 `wAsync` stream3[1,4,2,7,5,3,8,6,9]

This operation cannot be used to fold an infinite lazy container of streams, because it schedules all the streams in a round robin manner.

Note that WSerialT and single threaded WAsyncT both interleave streams but the exact scheduling is slightly different in both cases.

Since: 0.2.0 (Streamly)

valueparallel :: (IsStream t, MonadAsync m) => t m a -> t m a -> t m a
#

Like async except that the execution is much more strict. There is no limit on the number of threads. While async may not schedule a stream if there is no demand from the consumer, parallel always evaluates both the streams immediately. The only limit that applies to parallel is Streamly.Prelude.maxBuffer. Evaluation may block if the output buffer becomes full.

Example3 expressions
import Streamly.Prelude (parallel)stream = Stream.fromEffect (delay 2) `parallel` Stream.fromEffect (delay 1)Stream.toList stream -- IO [Int]1 sec2 sec[1,2]

parallel guarantees that all the streams are scheduled for execution immediately, therefore, we could use things like starting timers inside the streams and relying on the fact that all timers were started at the same time.

Unlike async this operation cannot be used to fold an infinite lazy container of streams, because it schedules all the streams strictly concurrently.

Since: 0.2.0 (Streamly)

valuemergeBy :: IsStream t => (a -> a -> Ordering) -> t m a -> t m a -> t m a
#

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.

>>> Stream.toList $ Stream.mergeBy compare (Stream.fromList [1,3,5]) (Stream.fromList [2,4,6,8])
[1,2,3,4,5,6,8]

See also: mergeByMFused

valuemergeByM
  1. :: (IsStream t, Monad m)
  2. => a -> a -> m Ordering
  3. -> t m a
  4. -> t m a
  5. -> t m a
#

Like mergeBy but with a monadic comparison function.

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]

Merge two streams in a proportion of 2:1:

>>> :{
do
 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.toList $ Stream.mergeByM f (Stream.fromList [1,1,1,1,1,1]) (Stream.fromList [2,2,2])
 print xs
:}
[1,1,2,1,1,2,1,1,2]

See also: mergeByMFused

valuezipWith :: (IsStream t, Monad m) => (a -> b -> c) -> t m a -> t m b -> t m c
#

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.

> D.toList $ D.zipWith (+) (D.fromList [1,2,3]) (D.fromList [4,5,6])
[5,7,9]
valuezipAsyncWith
  1. :: (IsStream t, MonadAsync m)
  2. => a -> b -> c
  3. -> t m a
  4. -> t m b
  5. -> t m c
#

Like zipWith but zips concurrently i.e. both the streams being zipped are evaluated concurrently using the ParallelT concurrent evaluation style. The maximum number of elements of each stream evaluated in advance can be controlled by maxBuffer.

The stream ends if stream a or stream b ends. However, if stream b ends while we are still evaluating stream a and waiting for a result then stream will not end until after the evaluation of stream a finishes. This behavior can potentially be changed in future to end the stream immediately as soon as any of the stream end is detected.

valueintercalate :: (IsStream t, Monad m) => Unfold m b c -> b -> t m b -> t m c
#

intersperse followed by unfold and concat.

intercalate unf a str = unfoldMany unf $ intersperse a str
intersperse = intercalate (Unfold.function id)
unwords = intercalate Unfold.fromList " "
Example1 expression
Stream.toList $ Stream.intercalate Unfold.fromList " " $ Stream.fromList ["abc", "def", "ghi"]"abc def ghi"
valueintercalateSuffix
  1. :: (IsStream t, Monad m)
  2. => Unfold m b c
  3. -> b
  4. -> t m b
  5. -> t m c
#

intersperseMSuffix followed by unfold and concat.

intercalateSuffix unf a str = unfoldMany unf $ intersperseMSuffix a str
intersperseMSuffix = intercalateSuffix (Unfold.function id)
unlines = intercalateSuffix Unfold.fromList "\n"
Example1 expression
Stream.toList $ Stream.intercalateSuffix Unfold.fromList "\n" $ Stream.fromList ["abc", "def", "ghi"]"abc\ndef\nghi\n"
valueconcatMapWith
  1. :: IsStream t
  2. => t m b -> t m b -> t m b
  3. -> a -> t m b
  4. -> t m a
  5. -> t m b
#

concatMapWith mixer generator stream is a two dimensional looping combinator. The generator function is used to generate streams from the elements in the input stream and the mixer function is used to merge those streams.

Note we can merge streams concurrently by using a concurrent merge function.

Since: 0.7.0

Since: 0.8.0 (signature change)

valueconcatMap :: (IsStream t, Monad m) => (a -> t m b) -> t m a -> t 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.concatMapWith Stream.serial fconcatMap f = Stream.concat . Stream.map f
valueconcatMapM :: (IsStream t, Monad m) => (a -> m (t m b)) -> t m a -> t 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.

valueconcatFoldableWith
  1. :: (IsStream t, Foldable f)
  2. => t m a -> t m a -> t m a
  3. -> f (t m a)
  4. -> t m a
#

A variant of fold that allows you to fold a Foldable container of streams using the specified stream sum operation.

concatFoldableWith async $ map return [1..3]

Equivalent to:

concatFoldableWith f = Prelude.foldr f D.nil
concatFoldableWith f = D.concatMapFoldableWith f id

Since: 0.8.0 (Renamed foldWith to concatFoldableWith)

Since: 0.1.0 (Streamly)

valueconcatMapFoldableWith
  1. :: (IsStream t, Foldable f)
  2. => t m b -> t m b -> t m b
  3. -> a -> t m b
  4. -> f a
  5. -> t m b
#

A variant of foldMap that allows you to map a monadic streaming action on a Foldable container and then fold it using the specified stream merge operation.

concatMapFoldableWith async return [1..3]

Equivalent to:

concatMapFoldableWith f g = Prelude.foldr (f . g) S.nil
concatMapFoldableWith f g xs = S.concatMapWith f g (S.fromFoldable xs)

Since: 0.8.0 (Renamed foldMapWith to concatMapFoldableWith)

Since: 0.1.0 (Streamly)

valueconcatForFoldableWith
  1. :: (IsStream t, Foldable f)
  2. => t m b -> t m b -> t m b
  3. -> f a
  4. -> a -> t m b
  5. -> t m b
#

Like concatMapFoldableWith but with the last two arguments reversed i.e. the monadic streaming function is the last argument.

Equivalent to:

concatForFoldableWith f xs g = Prelude.foldr (f . g) D.nil xs
concatForFoldableWith f = flip (D.concatMapFoldableWith f)

Since: 0.8.0 (Renamed forEachWith to concatForFoldableWith)

Since: 0.1.0 (Streamly)

valuebindWith
  1. :: IsStream t
  2. => t m b -> t m b -> t m b
  3. -> t m a
  4. -> a -> t m b
  5. -> t m b
#
valueconcat :: (IsStream t, Monad m) => t m (t m a) -> t m a
#

Flatten a stream of streams to a single stream.

concat = concatMap id

Pre-release

valueconcatM :: (IsStream t, Monad m) => m (t m a) -> t m a
#

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

Example3 expressions
concatM = Stream.concat . Stream.fromEffectconcatM = Stream.concat . lift    -- requires (MonadTrans t)concatM = join . lift             -- requires (MonadTrans t, Monad (t m))

See also: concat, sequence

Internal

valueparallelFst :: (IsStream t, MonadAsync m) => t m a -> t m a -> t m a
#

Like parallel but stops the output as soon as the first stream stops.

Pre-release

valueappend :: (IsStream t, Monad m) => t m b -> t m b -> t m b
#

Append the outputs of two streams, yielding all the elements from the first stream and then yielding all the elements from the second stream.

IMPORTANT NOTE: This could be 100x faster than serial/<> for appending a few (say 100) streams because it can fuse via stream fusion. However, it does not scale for a large number of streams (say 1000s) and becomes qudartically slow. Therefore use this for custom appending of a few streams but use concatMap or 'concatMapWith serial' for appending n streams or infinite containers of streams.

Pre-release

valueinterleave :: (IsStream t, Monad m) => t m b -> t m b -> t m b
#

Interleaves the outputs of two streams, yielding elements from each stream alternately, starting from the first stream. If any of the streams finishes early the other stream continues alone until it too finishes.

Example3 expressions
:set -XOverloadedStringsimport Data.Functor.Identity (Identity)Stream.interleave "ab" ",,,," :: Stream.SerialT Identity CharfromList "a,b,,,"
Example1 expression
Stream.interleave "abcd" ",," :: Stream.SerialT Identity CharfromList "a,b,cd"

interleave is dual to interleaveMin, it can be called interleaveMax.

Do not use at scale in concatMapWith.

Pre-release

valueinterleaveMin :: (IsStream t, Monad m) => t m b -> t m b -> t m b
#

Interleaves the outputs of two streams, yielding elements from each stream alternately, starting from the first stream. The output stops as soon as any of the two streams finishes, discarding the remaining part of the other stream. The last element of the resulting stream would be from the longer stream.

Example4 expressions
:set -XOverloadedStringsimport Data.Functor.Identity (Identity)Stream.interleaveMin "ab" ",,,," :: Stream.SerialT Identity CharfromList "a,b,"Stream.interleaveMin "abcd" ",," :: Stream.SerialT Identity CharfromList "a,b,c"

interleaveMin is dual to interleave.

Do not use at scale in concatMapWith.

Pre-release

valueinterleaveSuffix :: (IsStream t, Monad m) => t m b -> t m b -> t m b
#

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.interleaveSuffix "abc" ",,,," :: Stream.SerialT Identity CharfromList "a,b,c,"Stream.interleaveSuffix "abc" "," :: Stream.SerialT Identity CharfromList "a,bc"

interleaveSuffix is a dual of interleaveInfix.

Do not use at scale in concatMapWith.

Pre-release

valueinterleaveInfix :: (IsStream t, Monad m) => t m b -> t m b -> t m b
#

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.interleaveInfix "abc" ",,,," :: Stream.SerialT Identity CharfromList "a,b,c"Stream.interleaveInfix "abc" "," :: Stream.SerialT Identity CharfromList "a,bc"

interleaveInfix is a dual of interleaveSuffix.

Do not use at scale in concatMapWith.

Pre-release

valueroundrobin :: (IsStream t, Monad m) => t m b -> t m b -> t m b
#

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 chose 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 at scale in concatMapWith.

Pre-release

valuemerge :: (IsStream t, Ord a) => t m a -> t m a -> t m a
#

Same as mergeBy compare.

Example1 expression
Stream.toList $ Stream.merge (Stream.fromList [1,3,5]) (Stream.fromList [2,4,6,8])[1,2,3,4,5,6,8]

Internal

valuemergeByMFused
  1. :: (IsStream t, Monad m)
  2. => a -> a -> m Ordering
  3. -> t m a
  4. -> t m a
  5. -> t m a
#

Like mergeByM but much faster, works best when merging statically known number of streams. When merging more than two streams try to merge pairs and pair pf pairs in a tree like structure.mergeByM works better with variable number of streams being merged using concatPairsWith.

Internal

valuemergeMinBy :: (a -> a -> m Ordering) -> t m a -> t m a -> t m a
#

Like mergeByM but stops merging as soon as any of the two streams stops.

Unimplemented

valuemergeFstBy :: (a -> a -> m Ordering) -> t m a -> t m a -> t m a
#

Like mergeByM but stops merging as soon as the first stream stops.

Unimplemented

valueinterpose :: (IsStream t, Monad m) => c -> Unfold m b c -> t m b -> t m c
#

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

unwords = S.interpose ' '

Pre-release

valueinterposeSuffix
  1. :: (IsStream t, Monad m)
  2. => c
  3. -> Unfold m b c
  4. -> t m b
  5. -> t m c
#

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

unlines = S.interposeSuffix '\n'

Pre-release

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

Like concatMapWith but carries a state which can be used to share information across multiple steps of concat.

concatSmapMWith combine f initial = concatMapWith combine id . smapM f initial

Pre-release

valueconcatPairsWith
  1. :: IsStream t
  2. => t m b -> t m b -> t m b
  3. -> a -> t m b
  4. -> t m a
  5. -> t m b
#

Combine streams in pairs using a binary stream combinator, then combine the resulting streams in pairs recursively until we get to a single combined stream.

For example, you can sort a stream using merge sort like this:

Example1 expression
Stream.toList $ Stream.concatPairsWith (Stream.mergeBy compare) Stream.fromPure $ Stream.fromList [5,1,7,9,2][1,2,5,7,9]

Caution: the stream of streams must be finite

Pre-release

valueiterateMapWith
  1. :: IsStream t
  2. => t m a -> t m a -> t m a
  3. -> a -> t m a
  4. -> t m a
  5. -> t m a
#

Like iterateM but iterates after mapping a stream generator on the output.

Yield an input element in the output stream, map a stream generator on it and then do the same on the resulting stream. This can be used for a depth first traversal of a tree like structure.

Note that iterateM is a special case of iterateMapWith:

iterateM f = iterateMapWith serial (fromEffect . f) . fromEffect

It can be used to traverse a tree structure. For example, to list a directory tree:

Stream.iterateMapWith Stream.serial
    (either Dir.toEither (const nil))
    (fromPure (Left "tmp"))

Pre-release

valueiterateSmapMWith
  1. :: (IsStream t, Monad m)
  2. => t m a -> t m a -> t m a
  3. -> b -> a -> m (b, t m a)
  4. -> m b
  5. -> t m a
  6. -> t m a
#

Like iterateMap but carries a state in the stream generation function. This can be used to traverse graph like structures, we can remember the visited nodes in the state to avoid cycles.

Note that a combination of iterateMap and usingState can also be used to traverse graphs. However, this function provides a more localized state instead of using a global state.

See also: mfix

Pre-release

valueiterateMapLeftsWith
  1. :: (IsStream t, b ~ Either a c)
  2. => t m b -> t m b -> t m b
  3. -> a -> t m b
  4. -> t m b
  5. -> t m b
#

In an Either stream iterate on Lefts. This is a special case of iterateMapWith:

iterateMapLeftsWith combine f = iterateMapWith combine (either f (const nil))

To traverse a directory tree:

iterateMapLeftsWith serial Dir.toEither (fromPure (Left "tmp"))

Pre-release

valueiterateUnfold :: Unfold m a a -> t m a -> t m a
#

Same as iterateMapWith Stream.serial but more efficient due to stream fusion.

Unimplemented

valueintervalsOf
  1. :: (IsStream t, MonadAsync m)
  2. => Double
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

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

Example1 expression
Stream.toList $ Stream.take 5 $ Stream.intervalsOf 1 Fold.sum $ Stream.constRate 2 $ Stream.enumerateFrom 1[...,...,...,...,...]
valuefoldMany :: (IsStream t, Monad m) => Fold m a b -> t m a -> t m b
#

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

To sum every two contiguous elements in a stream:

Example2 expressions
f = Fold.take 2 Fold.sumStream.toList $ Stream.foldMany f $ Stream.fromList [1..10][3,7,11,15,19]

On an empty stream the output is empty:

Example1 expression
Stream.toList $ Stream.foldMany f $ Stream.fromList [][]

Note Stream.foldMany (Fold.take 0) would result in an infinite loop in a non-empty stream.

valuechunksOf :: (IsStream t, Monad m) => Int -> Fold m a b -> t m a -> t m b
#

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

Example1 expression
Stream.toList $ Stream.chunksOf 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.

chunksOf n f = foldMany (FL.take n f)
valuesplitOn
  1. :: (IsStream t, Monad m)
  2. => a -> Bool
  3. -> Fold m a b
  4. -> t m a
  5. -> t 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.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
valuesplitOnSuffix
  1. :: (IsStream t, Monad m)
  2. => a -> Bool
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

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

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

An empty stream results in an empty output stream:

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

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

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

A suffix is optional at the end of the stream:

Example1 expression
splitOnSuffix' (== '.') "a"["a"]
Example1 expression
splitOnSuffix' (== '.') ".a"["","a"]
Example1 expression
splitOnSuffix' (== '.') "a.b"["a","b"]
lines = splitOnSuffix (== '\n')

splitOnSuffix is an inverse of intercalateSuffix with a single element:

Stream.intercalateSuffix (Stream.fromPure '.') Unfold.fromList . Stream.splitOnSuffix (== '.') Fold.toList === id

Assuming the input stream does not contain the separator:

Stream.splitOnSuffix (== '.') Fold.toList . Stream.intercalateSuffix (Stream.fromPure '.') Unfold.fromList === id
valuesplitWithSuffix
  1. :: (IsStream t, Monad m)
  2. => a -> Bool
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Like splitOnSuffix but keeps the suffix attached to the resulting splits.

Example1 expression
splitWithSuffix' p xs = Stream.toList $ splitWithSuffix p Fold.toList (Stream.fromList xs)
Example1 expression
splitWithSuffix' (== '.') ""[]
Example1 expression
splitWithSuffix' (== '.') "."["."]
Example1 expression
splitWithSuffix' (== '.') "a"["a"]
Example1 expression
splitWithSuffix' (== '.') ".a"[".","a"]
Example1 expression
splitWithSuffix' (== '.') "a."["a."]
Example1 expression
splitWithSuffix' (== '.') "a.b"["a.","b"]
Example1 expression
splitWithSuffix' (== '.') "a.b."["a.","b."]
Example1 expression
splitWithSuffix' (== '.') "a..b.."["a.",".","b.","."]
valuewordsBy
  1. :: (IsStream t, Monad m)
  2. => a -> Bool
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

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

Example1 expression
wordsBy' p xs = Stream.toList $ Stream.wordsBy p Fold.toList (Stream.fromList xs)
Example1 expression
wordsBy' (== ',') ""[]
Example1 expression
wordsBy' (== ',') ","[]
Example1 expression
wordsBy' (== ',') ",a,,b,"["a","b"]
words = wordsBy isSpace
valuegroups :: (IsStream t, Monad m, Eq a) => Fold m a b -> t m a -> t m b
#
groups = groupsBy (==)
groups = groupsByRolling (==)

Groups contiguous spans of equal elements together in individual groups.

Example1 expression
Stream.toList $ Stream.groups Fold.toList $ Stream.fromList [1,1,2,2][[1,1],[2,2]]
valuegroupsBy
  1. :: (IsStream t, Monad m)
  2. => a -> a -> Bool
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

groupsBy cmp f $ S.fromList [a,b,c,...] assigns the element a to the first group, if b `cmp` a is True then b is also assigned to the same group. If c `cmp` a is True then c is also assigned to the same group and so on. When the comparison fails a new group is started. Each group is folded using the fold f and the result of the fold is emitted in the output stream.

Example1 expression
Stream.toList $ Stream.groupsBy (>) Fold.toList $ Stream.fromList [1,3,7,0,2,5][[1,3,7],[0,2,5]]
valuegroupsByRolling
  1. :: (IsStream t, Monad m)
  2. => a -> a -> Bool
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Unlike groupsBy this function performs a rolling comparison of two successive elements in the input stream. groupsByRolling cmp f $ S.fromList [a,b,c,...] assigns the element a to the first group, if a `cmp` b is True then b is also assigned to the same group. If b `cmp` c is True then c is also assigned to the same group and so on. When the comparison fails a new group is started. Each group is folded using the fold f.

Example1 expression
Stream.toList $ Stream.groupsByRolling (\a b -> a + 1 == b) Fold.toList $ Stream.fromList [1,2,3,7,8,9][[1,2,3],[7,8,9]]
valueclassifySessionsBy
  1. :: (IsStream t, MonadAsync m, Ord k)
  2. => Double

    timer tick in seconds

  3. -> Bool

    reset the timer when an event is received

  4. -> (Int -> m Bool)

    predicate to eject sessions based on session count

  5. -> Double

    session timeout in seconds

  6. -> Fold m a b

    Fold to be applied to session data

  7. -> t m (AbsTime, (k, a))

    timestamp, (session key, session data)

  8. -> t m (k, b)

    session key, fold result

#

classifySessionsBy tick keepalive predicate timeout fold stream classifies an input event stream consisting of (timestamp, (key, value)) into sessions based on the key, folding all the values corresponding to the same key into a session using the supplied fold.

When the fold terminates or a timeout occurs, a tuple consisting of the session key and the folded value is emitted in the output stream. The timeout is measured from the first event in the session. If the keepalive option is set to True the timeout is reset to 0 whenever an event is received.

The timestamp in the input stream is an absolute time from some epoch, characterizing the time when the input event was generated. The notion of current time is maintained by a monotonic event time clock using the timestamps seen in the input stream. The latest timestamp seen till now is used as the base for the current time. When no new events are seen, a timer is started with a clock resolution of tick seconds. This timer is used to detect session timeouts in the absence of new events.

To ensure an upper bound on the memory used the number of sessions can be limited to an upper bound. If the ejection predicate returns True, the oldest session is ejected before inserting a new session.

When the stream ends any buffered sessions are ejected immediately.

If a session key is received even after a session has finished, another session is created for that key.

Example1 expression
:{Stream.mapM_ print    $ Stream.classifySessionsBy 1 False (const (return False)) 3 (Fold.take 3 Fold.toList)    $ Stream.timestamped    $ Stream.delay 0.1    $ Stream.fromList ((,) <$> [1,2,3] <*> ['a','b','c']):}(1,"abc")(2,"abc")(3,"abc")

Pre-release

valuesplitOnSeq
  1. :: (IsStream t, MonadIO m, Storable a, Unbox a, Enum a, Eq a)
  2. => Array a
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Like splitOn but the separator is a sequence of elements instead of a single element.

For illustration, let's define a function that operates on pure lists:

Example1 expression
splitOnSeq' pat xs = Stream.toList $ Stream.splitOnSeq (Array.fromList pat) Fold.toList (Stream.fromList xs)
Example1 expression
splitOnSeq' "" "hello"["h","e","l","l","o"]
Example1 expression
splitOnSeq' "hello" ""[""]
Example1 expression
splitOnSeq' "hello" "hello"["",""]
Example1 expression
splitOnSeq' "x" "hello"["hello"]
Example1 expression
splitOnSeq' "h" "hello"["","ello"]
Example1 expression
splitOnSeq' "o" "hello"["hell",""]
Example1 expression
splitOnSeq' "e" "hello"["h","llo"]
Example1 expression
splitOnSeq' "l" "hello"["he","","o"]
Example1 expression
splitOnSeq' "ll" "hello"["he","o"]

splitOnSeq is an inverse of intercalate. The following law always holds:

intercalate . splitOnSeq == id

The following law holds when the separator is non-empty and contains none of the elements present in the input lists:

splitOnSeq . intercalate == id
Example1 expression
splitOnSeq pat f = Stream.foldManyPost (Fold.takeEndBySeq_ pat f)

Pre-release

valuedropPrefix :: t m a -> t m a -> t m a
#

Drop prefix from the input stream if present.

Space: O(1)

Unimplemented

valuedropInfix :: t m a -> t m a -> t 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 :: t m a -> t m a -> t 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

valuefoldManyPost :: (IsStream t, Monad m) => Fold m a b -> t m a -> t m b
#

Like foldMany but appends empty fold output if the fold and stream termination aligns:

Example4 expressions
f = Fold.take 2 Fold.sumStream.toList $ Stream.foldManyPost f $ Stream.fromList [][0]Stream.toList $ Stream.foldManyPost f $ Stream.fromList [1..9][3,7,11,15,9]Stream.toList $ Stream.foldManyPost f $ Stream.fromList [1..10][3,7,11,15,19,0]

Pre-release

valuefoldSequence :: t m (Fold m a b) -> t m a -> t m b
#

Apply a stream of folds to an input stream and emit the results in the output stream.

Unimplemented

valuefoldIterateM
  1. :: (IsStream t, Monad m)
  2. => b -> m (Fold m a b)
  3. -> m b
  4. -> t m a
  5. -> t 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.

>>> import Data.Monoid (Sum(..))
>>> f x = return (Fold.take 2 (Fold.sconcat x))
>>> s = Stream.map Sum $ Stream.fromList [1..10]
>>> Stream.toList $ Stream.map 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

valuearraysOf
  1. :: (IsStream t, MonadIO m, Unbox a)
  2. => Int
  3. -> t m a
  4. -> t m (Array a)
#

arraysOf n stream groups the elements in the input stream into arrays of n elements each.

Same as the following but may be more efficient:

arraysOf n = Stream.foldMany (A.writeN n)

Pre-release

valuechunksOfTimeout
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Int
  3. -> Double
  4. -> Fold m a b
  5. -> t m a
  6. -> t m b
#

Like chunksOf but if the chunk is not completed within the specified time interval then emit whatever we have collected till now. The chunk timeout is reset whenever a chunk is emitted. The granularity of the clock is 100 ms.

Example2 expressions
s = Stream.delayPost 0.3 $ Stream.fromList [1..1000]f = Stream.mapM_ print $ Stream.chunksOfTimeout 5 1 Fold.toList s

Pre-release

valuesplitOnPrefix :: (a -> Bool) -> Fold m a b -> t m a -> t 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

valuesplitOnAny :: [Array a] -> Fold m a b -> t m a -> t m b
#

Split on any one of the given patterns.

Unimplemented

valuesplitBySeq
  1. :: (IsStream t, MonadAsync m, Storable a, Unbox a, Enum a, Eq a)
  2. => Array a
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Like splitOnSeq but splits the separator as well, as an infix token.

Example1 expression
splitOn'_ pat xs = Stream.toList $ Stream.splitBySeq (Array.fromList pat) Fold.toList (Stream.fromList xs)
Example1 expression
splitOn'_ "" "hello"["h","","e","","l","","l","","o"]
Example1 expression
splitOn'_ "hello" ""[""]
Example1 expression
splitOn'_ "hello" "hello"["","hello",""]
Example1 expression
splitOn'_ "x" "hello"["hello"]
Example1 expression
splitOn'_ "h" "hello"["","h","ello"]
Example1 expression
splitOn'_ "o" "hello"["hell","o",""]
Example1 expression
splitOn'_ "e" "hello"["h","e","llo"]
Example1 expression
splitOn'_ "l" "hello"["he","l","","l","o"]
Example1 expression
splitOn'_ "ll" "hello"["he","ll","o"]

Pre-release

valuesplitOnSuffixSeq
  1. :: (IsStream t, MonadIO m, Storable a, Unbox a, Enum a, Eq a)
  2. => Array a
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Like splitSuffixBy but the separator is a sequence of elements, instead of a predicate for a single element.

Example1 expression
splitOnSuffixSeq_ pat xs = Stream.toList $ Stream.splitOnSuffixSeq (Array.fromList pat) Fold.toList (Stream.fromList xs)
Example1 expression
splitOnSuffixSeq_ "." ""[]
Example1 expression
splitOnSuffixSeq_ "." "."[""]
Example1 expression
splitOnSuffixSeq_ "." "a"["a"]
Example1 expression
splitOnSuffixSeq_ "." ".a"["","a"]
Example1 expression
splitOnSuffixSeq_ "." "a."["a"]
Example1 expression
splitOnSuffixSeq_ "." "a.b"["a","b"]
Example1 expression
splitOnSuffixSeq_ "." "a.b."["a","b"]
Example1 expression
splitOnSuffixSeq_ "." "a..b.."["a","","b",""]
lines = splitOnSuffixSeq "\n"

splitOnSuffixSeq is an inverse of intercalateSuffix. The following law always holds:

intercalateSuffix . splitOnSuffixSeq == id

The following law holds when the separator is non-empty and contains none of the elements present in the input lists:

splitSuffixOn . intercalateSuffix == id
Example1 expression
splitOnSuffixSeq pat f = Stream.foldMany (Fold.takeEndBySeq_ pat f)

Pre-release

valuesplitWithSuffixSeq
  1. :: (IsStream t, MonadIO m, Storable a, Unbox a, Enum a, Eq a)
  2. => Array a
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Like splitOnSuffixSeq but keeps the suffix intact in the splits.

Example1 expression
splitWithSuffixSeq' pat xs = Stream.toList $ Stream.splitWithSuffixSeq (Array.fromList pat) Fold.toList (Stream.fromList xs)
Example1 expression
splitWithSuffixSeq' "." ""[]
Example1 expression
splitWithSuffixSeq' "." "."["."]
Example1 expression
splitWithSuffixSeq' "." "a"["a"]
Example1 expression
splitWithSuffixSeq' "." ".a"[".","a"]
Example1 expression
splitWithSuffixSeq' "." "a."["a."]
Example1 expression
splitWithSuffixSeq' "." "a.b"["a.","b"]
Example1 expression
splitWithSuffixSeq' "." "a.b."["a.","b."]
Example1 expression
splitWithSuffixSeq' "." "a..b.."["a.",".","b.","."]
Example1 expression
splitWithSuffixSeq pat f = Stream.foldMany (Fold.takeEndBySeq pat f)

Pre-release

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

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

This is the streaming equivalent of the many parse combinator.

Example1 expression
Stream.toList $ Stream.parseMany (Parser.takeBetween 0 2 Fold.sum) $ Stream.fromList [1..10][Right 3,Right 7,Right 11,Right 15,Right 19]
> Stream.toList $ Stream.parseMany (Parser.line Fold.toList) $ Stream.fromList "hello\nworld"
["hello\n","world"]

foldMany f = parseMany (fromFold f)

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

Pre-release

valueparseManyTill :: Parser a m b -> Parser a m x -> t m a -> t 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

valueparseSequence :: t m (Parser a m b) -> t m a -> t m b
#

Apply a stream of parsers to an input stream and emit the results in the output stream.

Unimplemented

valueparseIterate
  1. :: (IsStream t, Monad m)
  2. => b -> Parser a m b
  3. -> b
  4. -> t m a
  5. -> t 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.

Example2 expressions
import Data.Monoid (Sum(..))Stream.toList $ fmap getSum $ Stream.rights $ Stream.parseIterate (\b -> Parser.takeBetween 0 2 (Fold.sconcat b)) (Sum 0) $ fmap Sum $ Stream.fromList [1..10][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

valuesplitInnerBy
  1. :: (IsStream t, Monad m)
  2. => f a -> m (f a, Maybe (f a))
  3. -> f a -> f a -> m (f a)
  4. -> t m (f a)
  5. -> t m (f a)
#

splitInnerBy splitter joiner stream splits the inner containers f a of an input stream t m (f a) using the splitter function. Container elements f a are collected until a split occurs, then all the elements before the split are joined using the joiner function.

For example, if we have a stream of Array Word8, we may want to split the stream into arrays representing lines separated by 'n' byte such that the resulting stream after a split would be one array for each line.

CAUTION! This is not a true streaming function as the container size after the split and merge may not be bounded.

Pre-release

valueonException :: (IsStream t, MonadCatch m) => m b -> t m a -> t m a
#

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

Inhibits stream fusion

valuefinally :: (IsStream t, MonadAsync m, MonadCatch m) => m b -> t m a -> t m a
#

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

The semantics of running the action m b are similar to the cleanup action semantics described in bracket.

See also finally_

Inhibits stream fusion

valuebracket
  1. :: (IsStream t, MonadAsync m, MonadCatch m)
  2. => m b
  3. -> b -> m c
  4. -> b -> t m a
  5. -> t m a
#

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

b is usually a resource under the state of monad m, e.g. a file handle, that requires a cleanup after use. The cleanup action b -> m c, runs whenever the stream ends normally, due to a sync or async exception or if it gets garbage collected after a partial lazy evaluation.

bracket 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.

See also: bracket_

Inhibits stream fusion

valueafter :: (IsStream t, MonadRunInIO m) => m b -> t m a -> t m a
#

Run the action m b whenever the stream t m a stops normally, or if it is garbage collected after a partial lazy evaluation.

The semantics of the action m b are similar to the semantics of cleanup action in bracket.

See also after_

valuebefore :: (IsStream t, Monad m) => m b -> t m a -> t 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)
valuehandle
  1. :: (IsStream t, MonadCatch m, Exception e)
  2. => e -> t m a
  3. -> t m a
  4. -> t 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.

Inhibits stream fusion

valueretry
  1. :: (IsStream t, MonadCatch m, Exception e, Ord e)
  2. => Map e Int

    map from exception to retry count

  3. -> (e -> t m a)

    default handler for those exceptions that are not in the map

  4. -> t m a
  5. -> t m a
#

retry takes 3 arguments

  1. A map m whose keys are exceptions and values are the number of times to retry the action given that the exception occurs.

  2. A handler han that decides how to handle an exception when the exception cannot be retried.

  3. The stream itself that we want to run this mechanism on.

When evaluating a stream if an exception occurs,

  1. The stream evaluation aborts

  2. The exception is looked up in m

a. If the exception exists and the mapped value is > 0 then,

i. The value is decreased by 1.

ii. The stream is resumed from where the exception was called, retrying the action.

b. If the exception exists and the mapped value is == 0 then the stream evaluation stops.

c. If the exception does not exist then we handle the exception using han.

Internal

valueafter_ :: (IsStream t, Monad m) => m b -> t m a -> t 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:

after_ action xs = xs <> 'nilM' action

Pre-release

valuebracket_
  1. :: (IsStream t, MonadCatch m)
  2. => m b
  3. -> b -> m c
  4. -> b -> t m a
  5. -> t 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.

  • does not require a MonadAsync constraint.

  • has slightly better performance than bracket.

Inhibits stream fusion

Pre-release

valuebracket'
  1. :: (IsStream t, MonadAsync m, MonadCatch m)
  2. => m b
  3. -> b -> m c
  4. -> b -> m d
  5. -> b -> m e
  6. -> b -> t m a
  7. -> t m a
#

Like bracket but can use separate cleanup actions depending on the mode of termination. bracket' 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.

Pre-release

valuefinally_ :: (IsStream t, MonadCatch m) => m b -> t m a -> t m a
#

Like finally with following differences:

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

  • does not require a MonadAsync constraint.

  • has slightly better performance than finally.

Inhibits stream fusion

Pre-release

valueghandle
  1. :: (IsStream t, MonadCatch m, Exception e)
  2. => e -> t m a -> t m a
  3. -> t m a
  4. -> t 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

valuerunStateT :: Monad m => m s -> SerialT (StateT s m) a -> SerialT m (s, a)
#

Evaluate the inner monad of a stream as StateT and emit the resulting state and value pair after each step.

This is supported only for SerialT as concurrent state updation may not be safe.

valuehoist
  1. :: (Monad m, Monad n)
  2. => forall x. m x -> n x
  3. -> SerialT m a
  4. -> SerialT n a
#

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

Internal

valueevalStateT :: Monad m => m s -> SerialT (StateT s m) a -> SerialT m a
#

Evaluate the inner monad of a stream as StateT.

This is supported only for SerialT as concurrent state updation may not be safe.

evalStateT s = Stream.map snd . Stream.runStateT s

Internal

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

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

This is supported only for SerialT as concurrent state updation may not be safe.

usingStateT s f = evalStateT s . f . liftInner

See also: scanl'

Internal

valuesampleIntervalEnd
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t m a
  4. -> t m a
#

Continuously evaluate the input stream and sample the last event in time window of n seconds.

This is also known as throttle in some libraries.

sampleIntervalEnd n = Stream.catMaybes . Stream.intervalsOf n Fold.last

Pre-release

valuesampleBurstEnd
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t m a
  4. -> t m a
#

Sample one event at the end of each burst of events. A burst is a group of events close together in time, it ends when an event is spaced by more than the specified time interval (in seconds) from the previous event.

This is known as debounce in some libraries.

The clock granularity is 10 ms.

Pre-release

valuesampleFromThen
  1. :: (IsStream t, Monad m, Functor (t m))
  2. => Int
  3. -> Int
  4. -> t m a
  5. -> t m a
#

sampleFromthen offset stride samples the element at offset index and then every element at strides of stride.

Example1 expression
Stream.toList $ Stream.sampleFromThen 2 3 $ Stream.enumerateFromTo 0 10[2,5,8]

Pre-release

valuesortBy :: MonadCatch m => (a -> a -> Ordering) -> SerialT m a -> SerialT m a
#

Sort the input stream using a supplied comparison function.

O(n) space

Note: this is not the fastest possible implementation as of now.

Pre-release

valueintersectBy
  1. :: (IsStream t, Monad m)
  2. => a -> a -> Bool
  3. -> t m a
  4. -> t m a
  5. -> t m a
#

intersectBy is essentially a filtering operation that retains only those elements in the first stream that are present in the second stream.

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

intersectBy is similar to but not the same as joinInner:

Example1 expression
Stream.toList $ fmap fst $ Stream.joinInner (==) (Stream.fromList [1,2,2,4]) (Stream.fromList [2,1,1,3])[1,1,2,2]

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

valuedifferenceBy
  1. :: (IsStream t, Monad m)
  2. => a -> a -> Bool
  3. -> t m a
  4. -> t m a
  5. -> t m a
#

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

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

The following laws hold:

(s1 serial s2) `differenceBy eq` s1 === s2
(s1 wSerial s2) `differenceBy eq` s1 === s2

Same as the list Data.List.// operation.

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

valueunionBy
  1. :: (IsStream t, MonadAsync m, Semigroup (t m a))
  2. => a -> a -> Bool
  3. -> t m a
  4. -> t m a
  5. -> t m a
#

This is essentially an append operation that appends all the extra occurrences of elements from the second stream that are not already present in the first stream.

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

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

unionBy eq s1 s2 = s1 `serial` (s2 `differenceBy eq` s1)

Similar to joinOuter but not the same.

Space: O(n)

Time: O(m x n)

Pre-release

valuemergeUnionBy :: (a -> a -> Ordering) -> t m a -> t m a -> t m a
#

Like unionBy but works only on sorted streams.

Space: O(1)

Unimplemented

valuecrossJoin :: Monad (t m) => t m a -> t m b -> t m (a, b)
#

This is the same as Streamly.Internal.Data.Unfold.outerProduct but less efficient.

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

Time: O(m x n)

Pre-release

valuejoinInner
  1. :: (IsStream t, Monad m)
  2. => a -> b -> Bool
  3. -> t m a
  4. -> t m b
  5. -> t m (a, b)
#

For all elements in t m a, for all elements in t m b if a and b are equal by the given equality pedicate then return the tuple (a, b).

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

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

You should almost always use joinInnerMap instead of joinInner. joinInnerMap is an order of magnitude faster. joinInner may be used when the second stream is generated from a seed, therefore, need not be stored in memory and the amount of memory it takes is a concern.

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

Time: O(m x n)

Pre-release

valuejoinInnerMap
  1. :: (IsStream t, Monad m, Ord k)
  2. => t m (k, a)
  3. -> t m (k, b)
  4. -> t 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

valuejoinInnerMerge :: (a -> b -> Ordering) -> t m a -> t m b -> t m (a, b)
#

Like joinInner but works only on sorted streams.

Space: O(1)

Time: O(m + n)

Unimplemented

valuemergeLeftJoin :: (a -> b -> Ordering) -> t m a -> t m b -> t m (a, Maybe b)
#

Like joinLeft but works only on sorted streams.

Space: O(1)

Time: O(m + n)

Unimplemented

valuejoinLeftMap
  1. :: (IsStream t, Ord k, Monad m)
  2. => t m (k, a)
  3. -> t m (k, b)
  4. -> t m (k, a, Maybe b)
#

Like joinLeft but uses a hashmap for efficiency.

Space: O(n)

Time: O(m + n)

Pre-release

valuemergeOuterJoin
  1. :: a -> b -> Ordering
  2. -> t m a
  3. -> t m b
  4. -> t m (Maybe a, Maybe b)
#

Like joinOuter but works only on sorted streams.

Space: O(1)

Time: O(m + n)

Unimplemented

valuemaxThreads :: IsStream t => Int -> t m a -> t m a
#

Specify the maximum number of threads that can be spawned concurrently for any concurrent combinator in a stream. A value of 0 resets the thread limit to default, a negative value means there is no limit. The default value is 1500. maxThreads does not affect ParallelT streams as they can use unbounded number of threads.

When the actions in a stream are IO bound, having blocking IO calls, this option can be used to control the maximum number of in-flight IO requests. When the actions are CPU bound this option can be used to control the amount of CPU used by the stream.

Since: 0.4.0 (Streamly)

valuemaxBuffer :: IsStream t => Int -> t m a -> t m a
#

Specify the maximum size of the buffer for storing the results from concurrent computations. If the buffer becomes full we stop spawning more concurrent tasks until there is space in the buffer. A value of 0 resets the buffer size to default, a negative value means there is no limit. The default value is 1500.

CAUTION! using an unbounded maxBuffer value (i.e. a negative value) coupled with an unbounded maxThreads value is a recipe for disaster in presence of infinite streams, or very large streams. Especially, it must not be used when pure is used in ZipAsyncM streams as pure in applicative zip streams generates an infinite stream causing unbounded concurrent generation with no limit on the buffer or threads.

Since: 0.4.0 (Streamly)

valuerate :: IsStream t => Maybe Rate -> t m a -> t m a
#

Specify the pull rate of a stream. A Nothing value resets the rate to default which is unlimited. When the rate is specified, concurrent production may be ramped up or down automatically to achieve the specified yield rate. The specific behavior for different styles of Rate specifications is documented under Rate. The effective maximum production rate achieved by a stream is governed by:

  • The maxThreads limit

  • The maxBuffer limit

  • The maximum rate that the stream producer can achieve

  • The maximum rate that the stream consumer can achieve

Since: 0.5.0 (Streamly)

valueavgRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate (r/2) r (2*r) maxBound)

Specifies the average production rate of a stream in number of yields per second (i.e. Hertz). Concurrent production is ramped up or down automatically to achieve the specified average yield rate. The rate can go down to half of the specified rate on the lower side and double of the specified rate on the higher side.

Since: 0.5.0 (Streamly)

valueminRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate r r (2*r) maxBound)

Specifies the minimum rate at which the stream should yield values. As far as possible the yield rate would never be allowed to go below the specified rate, even though it may possibly go above it at times, the upper limit is double of the specified rate.

Since: 0.5.0 (Streamly)

valuemaxRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate (r/2) r r maxBound)

Specifies the maximum rate at which the stream should yield values. As far as possible the yield rate would never be allowed to go above the specified rate, even though it may possibly go below it at times, the lower limit is half of the specified rate. This can be useful in applications where certain resource usage must not be allowed to go beyond certain limits.

Since: 0.5.0 (Streamly)

valueconstRate :: IsStream t => Double -> t m a -> t m a
#

Same as rate (Just $ Rate r r r 0)

Specifies a constant yield rate. If for some reason the actual rate goes above or below the specified rate we do not try to recover it by increasing or decreasing the rate in future. This can be useful in applications like graphics frame refresh where we need to maintain a constant refresh rate.

Since: 0.5.0 (Streamly)

valueinspectMode :: IsStream t => t m a -> t m a
#

Print debug information about an SVar when the stream ends

Pre-release

valuefromPure :: IsStream t => a -> t m a
#
fromPure a = a `cons` nil

Create a singleton stream from a pure value.

The following holds in monadic streams, but not in Zip streams:

fromPure = pure
fromPure = fromEffect . pure

In Zip applicative streams fromPure is not the same as pure because in that case pure is equivalent to repeat instead. fromPure and pure are equally efficient, in other cases fromPure may be slightly more efficient than the other equivalent definitions.

Since: 0.8.0 (Renamed yield to fromPure)

valuefromEffect :: (Monad m, IsStream t) => m a -> t m a
#
fromEffect m = m `consM` nil

Create a singleton stream from a monadic action.

> Stream.toList $ Stream.fromEffect getLine
hello
["hello"]

Since: 0.8.0 (Renamed yieldM to fromEffect)

valuerepeatM :: (IsStream t, MonadAsync m) => m a -> t m a
#
Example2 expressions
repeatM = fix . consMrepeatM = cycle1 . fromEffect

Generate a stream by repeatedly executing a monadic action forever.

Example1 expression
:{repeatAsync =       Stream.repeatM (threadDelay 1000000 >> print 1)     & Stream.take 10     & Stream.fromAsync     & Stream.drain:}

Concurrent, infinite (do not use with fromParallel)

valuefold :: Monad m => Fold m a b -> SerialT 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.

Example1 expression
Stream.fold Fold.sum (Stream.enumerateFromTo 1 100)5050

Folds never fail, therefore, they produce a default value even when no input is provided. It means we can always fold an empty stream and get a valid result. For example:

Example1 expression
Stream.fold Fold.sum Stream.nil0

However, foldMany on an empty stream results in an empty stream. Therefore, Stream.fold f is not the same as Stream.head . Stream.foldMany f.

fold f = Stream.parse (Parser.fromFold f)
valuemap :: (IsStream t, Monad m) => (a -> b) -> t m a -> t m b
#
map = fmap

Same as fmap.

> D.toList $ D.map (+1) $ D.fromList [1,2,3]
[2,3,4]
valuepostscanlM'
  1. :: (IsStream t, Monad m)
  2. => b -> a -> m b
  3. -> m b
  4. -> t m a
  5. -> t m b
#

Like postscanl' but with a monadic step function and a monadic seed.

Example1 expression
postscanlM' f z xs = Stream.drop 1 $ Stream.scanlM' f z xs

Since: 0.7.0

Since: 0.8.0 (signature change)

valuetake :: (IsStream t, Monad m) => Int -> t m a -> t m a
#

Take first n elements from the stream and discard the rest.

valuetakeWhile :: (IsStream t, Monad m) => (a -> Bool) -> t m a -> t m a
#

End the stream as soon as the predicate fails on an element.

valuedrop :: (IsStream t, Monad m) => Int -> t m a -> t m a
#

Discard first n elements from the stream and take the rest.

valueintersperseM :: (IsStream t, MonadAsync m) => m a -> t m a -> t m a
#

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

Example1 expression
Stream.toList $ Stream.trace putChar $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.fromList "hello"h.,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.toList $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.trace putChar $ Stream.fromList "hello"he.l.l.o."h,e,l,l,o"
valuereverse :: (IsStream t, Monad m) => t m a -> t 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.

Example1 expression
reverse = Stream.foldlT (flip Stream.cons) Stream.nil

Since 0.7.0 (Monad m constraint)

Since: 0.1.1

valuefindIndices :: (IsStream t, Monad m) => (a -> Bool) -> t m a -> t m Int
#

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

findIndices = fold Fold.findIndices
valuemkAsync :: (IsStream t, MonadAsync m) => t m a -> t m a
#

Make the stream producer and consumer run concurrently by introducing a buffer between them. The producer thread evaluates the input stream until the buffer fills, it terminates if the buffer is full and a worker thread is kicked off again to evaluate the remaining stream when there is space in the buffer. The consumer consumes the stream lazily from the buffer.

Since: 0.2.0 (Streamly)

valuezipWith :: (IsStream t, Monad m) => (a -> b -> c) -> t m a -> t m b -> t m c
#

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.

> D.toList $ D.zipWith (+) (D.fromList [1,2,3]) (D.fromList [4,5,6])
[5,7,9]
valueconcatMap :: (IsStream t, Monad m) => (a -> t m b) -> t m a -> t 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.concatMapWith Stream.serial fconcatMap f = Stream.concat . Stream.map f
valueconcatMapM :: (IsStream t, Monad m) => (a -> m (t m b)) -> t m a -> t 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.

valuescanlMAfter'
  1. :: (IsStream t, Monad m)
  2. => b -> a -> m b
  3. -> m b
  4. -> b -> m b
  5. -> t m a
  6. -> t 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

valueabsTimesWith
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t 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.

Example1 expression
Stream.mapM_ print $ Stream.delayPre 1 $ Stream.take 3 $ 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

valuerelTimesWith
  1. :: (IsStream t, MonadAsync m, Functor (t m))
  2. => Double
  3. -> t m RelTime64
#

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.

Example1 expression
Stream.mapM_ print $ 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

valueconcatM :: (IsStream t, Monad m) => m (t m a) -> t m a
#

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

Example3 expressions
concatM = Stream.concat . Stream.fromEffectconcatM = Stream.concat . lift    -- requires (MonadTrans t)concatM = join . lift             -- requires (MonadTrans t, Monad (t m))

See also: concat, sequence

Internal

valuetimesWith
  1. :: (IsStream t, MonadAsync m)
  2. => Double
  3. -> t 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)import Streamly.Internal.Data.Stream.IsStream.Common as Stream (timesWith)Stream.mapM_ (\x -> print x >> threadDelay 1000000) $ 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

valuesplitOnSeq
  1. :: (IsStream t, MonadIO m, Storable a, Unbox a, Enum a, Eq a)
  2. => Array a
  3. -> Fold m a b
  4. -> t m a
  5. -> t m b
#

Like splitOn but the separator is a sequence of elements instead of a single element.

For illustration, let's define a function that operates on pure lists:

Example1 expression
splitOnSeq' pat xs = Stream.toList $ Stream.splitOnSeq (Array.fromList pat) Fold.toList (Stream.fromList xs)
Example1 expression
splitOnSeq' "" "hello"["h","e","l","l","o"]
Example1 expression
splitOnSeq' "hello" ""[""]
Example1 expression
splitOnSeq' "hello" "hello"["",""]
Example1 expression
splitOnSeq' "x" "hello"["hello"]
Example1 expression
splitOnSeq' "h" "hello"["","ello"]
Example1 expression
splitOnSeq' "o" "hello"["hell",""]
Example1 expression
splitOnSeq' "e" "hello"["h","llo"]
Example1 expression
splitOnSeq' "l" "hello"["he","","o"]
Example1 expression
splitOnSeq' "ll" "hello"["he","o"]

splitOnSeq is an inverse of intercalate. The following law always holds:

intercalate . splitOnSeq == id

The following law holds when the separator is non-empty and contains none of the elements present in the input lists:

splitOnSeq . intercalate == id
Example1 expression
splitOnSeq pat f = Stream.foldManyPost (Fold.takeEndBySeq_ pat f)

Pre-release

valuemkParallel :: (IsStream t, MonadAsync m) => t m a -> t m a
#

Make the stream producer and consumer run concurrently by introducing a buffer between them. The producer thread evaluates the input stream until the buffer fills, it blocks if the buffer is full until there is space in the buffer. The consumer consumes the stream lazily from the buffer.

mkParallel = IsStream.fromStreamD . mkParallelD . IsStream.toStreamD

Pre-release

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

A stateful mapM, equivalent to a left scan, more like mapAccumL. Hopefully, this is a better alternative to scan. Separation of state from the output makes it easier to think in terms of a shared state, and also makes it easier to keep the state fully strict and the output lazy.

See also: scanlM'

Pre-release

valueinterjectSuffix
  1. :: (IsStream t, MonadAsync m)
  2. => Double
  3. -> m a
  4. -> t m a
  5. -> t m a
#

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

> import Control.Concurrent (threadDelay)
> Stream.drain $ Stream.interjectSuffix 1 (putChar ',') $ Stream.mapM (x -> threadDelay 1000000 >> putChar x) $ Stream.fromList "hello"
h,e,l,l,o

Pre-release

valueparallelFst :: (IsStream t, MonadAsync m) => t m a -> t m a -> t m a
#

Like parallel but stops the output as soon as the first stream stops.

Pre-release

valuefoldManyPost :: (IsStream t, Monad m) => Fold m a b -> t m a -> t m b
#

Like foldMany but appends empty fold output if the fold and stream termination aligns:

Example4 expressions
f = Fold.take 2 Fold.sumStream.toList $ Stream.foldManyPost f $ Stream.fromList [][0]Stream.toList $ Stream.foldManyPost f $ Stream.fromList [1..9][3,7,11,15,9]Stream.toList $ Stream.foldManyPost f $ Stream.fromList [1..10][3,7,11,15,19,0]

Pre-release

valuefoldContinue :: Monad m => Fold m a b -> SerialT m a -> Fold m a b
#

We can create higher order folds using foldContinue. We can fold a number of streams to a given fold efficiently with full stream fusion. For example, to fold a list of streams on the same sum fold:

concatFold = Prelude.foldl Stream.foldContinue Fold.sum
fold f = Fold.extractM . Stream.foldContinue f

Internal

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 :: (IsStream t, Monad m) => a -> t 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.

    >>> 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.

    >>> Stream.toList $ Stream.take 4 $ Stream.enumerateFrom 1.1
    [1.1,2.1,3.1,4.1]
    
    
  • enumerateFromTo :: (IsStream t, Monad m) => a -> a -> t 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.

    >>> 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.

    >>> Stream.toList $ Stream.enumerateFromTo 1.1 4
    [1.1,2.1,3.1,4.1]
    
    >>> Stream.toList $ Stream.enumerateFromTo 1.1 4.6
    [1.1,2.1,3.1,4.1,5.1]
    
    
  • enumerateFromThen :: (IsStream t, Monad m) => a -> a -> t 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.

    >>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 2
    [0,2,4,6]
    
    >>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 (-2)
    [0,-2,-4,-6]
    
    
  • enumerateFromThenTo :: (IsStream t, Monad m) => a -> a -> a -> t 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.

    >>> Stream.toList $ Stream.enumerateFromThenTo 0 2 6
    [0,2,4,6]
    
    >>> Stream.toList $ Stream.enumerateFromThenTo 0 (-2) (-6)
    [0,-2,-4,-6]
    
    
Instances21Enumerable, …
  • Enumerable IntegerDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable NaturalDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int16Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int32Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int64Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Int8Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word16Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word32Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word64Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable Word8Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable BoolDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable CharDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable DoubleDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable FloatDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable IntDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable OrderingDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable WordDefined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable ()Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Integral a => Enumerable (Ratio a)Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • Enumerable a => Enumerable (Identity a)Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
  • HasResolution a => Enumerable (Fixed a)Defined in streamly-0.10.1 · Streamly.Internal.Data.Stream.IsStream.Enumeration
valueenumerateFromIntegral
  1. :: (IsStream t, Monad m, Integral a, Bounded a)
  2. => a
  3. -> t 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.

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromIntegral (0 :: Int)
[0,1,2,3]

valueenumerateFromThenIntegral
  1. :: (IsStream t, Monad m, Integral a, Bounded a)
  2. => a
  3. -> a
  4. -> t 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.

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenIntegral (0 :: Int) 2
[0,2,4,6]

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenIntegral (0 :: Int) (-2)
[0,-2,-4,-6]

valueenumerateFromToIntegral
  1. :: (IsStream t, Monad m, Integral a)
  2. => a
  3. -> a
  4. -> t 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.

>>> Stream.toList $ Stream.enumerateFromToIntegral 0 4
[0,1,2,3,4]

valueenumerateFromThenToIntegral
  1. :: (IsStream t, Monad m, Integral a)
  2. => a
  3. -> a
  4. -> a
  5. -> t 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.

>>> Stream.toList $ Stream.enumerateFromThenToIntegral 0 2 6
[0,2,4,6]

>>> Stream.toList $ Stream.enumerateFromThenToIntegral 0 (-2) (-6)
[0,-2,-4,-6]

valueenumerateFromStepIntegral
  1. :: (IsStream t, Monad m, Integral a)
  2. => a
  3. -> a
  4. -> t 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.

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromStepIntegral 0 2
[0,2,4,6]

>>> Stream.toList $ Stream.take 3 $ Stream.enumerateFromStepIntegral 0 (-2)
[0,-2,-4]

valueenumerateFromFractional :: (IsStream t, Monad m, Fractional a) => a -> t 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:

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromFractional 1.1
[1.1,2.1,3.1,4.1]

valueenumerateFromToFractional
  1. :: (IsStream t, Monad m, Fractional a, Ord a)
  2. => a
  3. -> a
  4. -> t 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:

>>> Stream.toList $ Stream.enumerateFromToFractional 1.1 4
[1.1,2.1,3.1,4.1]

>>> 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. :: (IsStream t, Monad m, Fractional a)
  2. => a
  3. -> a
  4. -> t 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:

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenFractional 1.1 2.1
[1.1,2.1,3.1,4.1]

>>> Stream.toList $ Stream.take 4 $ Stream.enumerateFromThenFractional 1.1 (-2.1)
[1.1,-2.1,-5.300000000000001,-8.500000000000002]

valueenumerateFromThenToFractional
  1. :: (IsStream t, Monad m, Fractional a, Ord a)
  2. => a
  3. -> a
  4. -> a
  5. -> t 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:

>>> Stream.toList $ Stream.enumerateFromThenToFractional 0.1 2 6
[0.1,2.0,3.9,5.799999999999999]

>>> Stream.toList $ Stream.enumerateFromThenToFractional 0.1 (-2) (-6)
[0.1,-2.0,-4.1000000000000005,-6.200000000000001]