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

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

Modulestreamly-core-0.2.2Haskell2010

Streamly.Internal.Data.StreamK

  • 4 types
  • 142 values

Setup

0 declarations

To execute the code examples provided in this module in ghci, please run the following commands first.

Example4 expressions
:mimport Control.Concurrent (threadDelay)import Data.Function (fix, (&))import Data.Semigroup (cycle1)
Example1 expression
effect n = print n >> return n
Example6 expressions
import Streamly.Data.StreamK (StreamK)import qualified Streamly.Data.Fold as Foldimport qualified Streamly.Data.Parser as Parserimport qualified Streamly.Data.Stream as Streamimport qualified Streamly.Data.StreamK as StreamKimport qualified Streamly.FileSystem.Dir as Dir

For APIs that have not been released yet.

Example2 expressions
import qualified Streamly.Internal.Data.StreamK as StreamKimport qualified Streamly.Internal.FileSystem.Dir as Dir

The stream type

86 declarations
valuefoldl' :: Monad m => (b -> a -> b) -> b -> StreamK m a -> m b
#

Strict left associative fold.

typetype Stream = StreamK
#

Deprecated. Please use StreamK instead.

Continuation Passing Style (CPS) version of Streamly.Data.Stream.Stream. Unlike Streamly.Data.Stream.Stream, StreamK can be composed recursively without affecting performance.

Semigroup instance appends two streams:

Example1 expression
(<>) = Stream.append
valuefoldr :: Monad m => (a -> b -> b) -> b -> StreamK m a -> m b
#

Lazy right associative fold.

valuenil :: StreamK m a
#

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

Example1 expression
Stream.fold Fold.toList (StreamK.toStream StreamK.nil)[]
newtypenewtype StreamK (m :: Type -> Type) a
#

Constructors

Instances10Functor, Foldable, Traversable, IsList, Read, Show, …
valuerepeat :: a -> StreamK m a
#

Generate an infinite stream by repeating a pure value.

Pre-release

valueinit :: Applicative m => StreamK m a -> m (Maybe (StreamK m a))
#

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

Note: This will end up buffering the entire stream.

Pre-release

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

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

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

Pre-release

valueconsM :: Monad m => m a -> StreamK m a -> StreamK m a
#

A right associative prepend operation to add an effectful value at the head of an existing stream::

Example2 expressions
s = putStrLn "hello" `StreamK.consM` putStrLn "world" `StreamK.consM` StreamK.nilStream.fold Fold.drain (StreamK.toStream s)helloworld

It can be used efficiently with foldr:

Example1 expression
fromFoldableM = Prelude.foldr StreamK.consM StreamK.nil

Same as the following but more efficient:

Example1 expression
consM x xs = StreamK.fromEffect x `StreamK.append` xs
valuecons :: a -> StreamK m a -> StreamK m a
#

A right associative prepend operation to add a pure value at the head of an existing stream::

Example2 expressions
s = 1 `StreamK.cons` 2 `StreamK.cons` 3 `StreamK.cons` StreamK.nilStream.fold Fold.toList (StreamK.toStream s)[1,2,3]

It can be used efficiently with foldr:

Example1 expression
fromFoldable = Prelude.foldr StreamK.cons StreamK.nil

Same as the following but more efficient:

Example1 expression
cons x xs = return x `StreamK.consM` xs
valueunfoldr :: (b -> Maybe (a, b)) -> b -> StreamK m a
#
Example1 expression
:{unfoldr step s =    case step s of        Nothing -> StreamK.nil        Just (a, b) -> a `StreamK.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 StreamK.toList $ StreamK.unfoldr f 0:}[0,1,2]
valueunfoldrM :: Monad m => (b -> m (Maybe (a, b))) -> b -> StreamK 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 StreamK.toList $ StreamK.unfoldrM f 0:}[0,1,2]
valuefoldrM :: (a -> m b -> m b) -> m b -> StreamK m a -> m b
#

Lazy right fold with a monadic step function.

valueinterleave :: StreamK m a -> StreamK m a -> StreamK m a
#

Note: When joining many streams in a left associative manner earlier streams will get exponential priority than the ones joining later. Because of exponentially high weighting of left streams it can be used with concatMapWith even on a large number of streams.

valuecrossWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> StreamK m a
  4. -> StreamK m b
  5. -> StreamK m c
#

Definition:

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

Note that the second stream is evaluated multiple times.

valuefromFoldable :: Foldable f => f a -> StreamK m a
#
Example1 expression
fromFoldable = Prelude.foldr StreamK.cons StreamK.nil

Construct a stream from a Foldable containing pure values:

valuemergeMapWith
  1. :: StreamK m b -> StreamK m b -> StreamK m b
  2. -> a -> StreamK m b
  3. -> StreamK m a
  4. -> StreamK m b
#

Combine streams in pairs using a binary combinator, the resulting streams are then combined again in pairs recursively until we get to a single combined stream. The composition would thus form a binary tree.

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

Example4 expressions
s = StreamK.fromStream $ Stream.fromList [5,1,7,9,2]generate = StreamK.fromPurecombine = StreamK.mergeBy compareStream.fold Fold.toList $ StreamK.toStream $ StreamK.mergeMapWith combine generate s[1,2,5,7,9]

Note that if the stream length is not a power of 2, the binary tree composed by mergeMapWith would not be balanced, which may or may not be important depending on what you are trying to achieve.

Caution: the stream of streams must be finite

Pre-release

valuefoldlMx'
  1. :: Monad m
  2. => x -> a -> m x
  3. -> m x
  4. -> x -> m b
  5. -> StreamK m a
  6. -> m b
#

Like foldx, but with a monadic step function.

valuefoldlx' :: Monad m => (x -> a -> x) -> x -> (x -> b) -> StreamK 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.

Note that the accumulator is always evaluated including the initial value.

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

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

Definition:

Example1 expression
cross = StreamK.crossWith (,)

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

See cross for a much faster fused alternative.

Time: O(m x n)

Pre-release

valuefoldrS
  1. :: a -> StreamK m b -> StreamK m b
  2. -> StreamK m b
  3. -> StreamK m a
  4. -> StreamK m b
#

Right fold to a streaming monad.

foldrS StreamK.cons StreamK.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.

Example2 expressions
input = StreamK.fromStream $ Stream.fromList [1..5]Stream.fold Fold.toList $ StreamK.toStream $ StreamK.foldrS StreamK.cons StreamK.nil input[1,2,3,4,5]

Find if any element in the stream is True:

Example3 expressions
step x xs = if odd x then StreamK.fromPure True else xsinput = StreamK.fromStream (Stream.fromList (2:4:5:undefined)) :: StreamK IO IntStream.fold Fold.toList $ StreamK.toStream $ StreamK.foldrS step (StreamK.fromPure False) input[True]

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

Example3 expressions
step x xs = if odd x then (x + 2) `StreamK.cons` xs else xsinput = StreamK.fromStream (Stream.fromList [1..5]) :: StreamK IO IntStream.fold Fold.toList $ StreamK.toStream $ StreamK.foldrS step StreamK.nil input[3,5,7]

Pre-release

valuecrossApply :: StreamK m (a -> b) -> StreamK m a -> StreamK m b
#

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

Note that the second stream is evaluated multiple times.

Definition:

Example2 expressions
crossApply = StreamK.crossApplyWith StreamK.appendcrossApply = Stream.crossWith id
valuebuild :: (forall b. (a -> b -> b) -> b -> b) -> StreamK m a
#
newtypenewtype CrossStreamK (m :: Type -> Type) a
#

A newtype wrapper for the StreamK type adding a cross product style monad instance.

A Monad bind behaves like a for loop:

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

Nested monad binds behave like nested for loops:

Example1 expression
:{Stream.fold Fold.toList $ StreamK.toStream $ StreamK.unCross $ do    x <- StreamK.mkCross $ StreamK.fromStream $ Stream.fromList [1,2]    y <- StreamK.mkCross $ StreamK.fromStream $ Stream.fromList [3,4]    -- Perform the following actions for each x, for each y    return (x, y):}[(1,3),(1,4),(2,3),(2,4)]
Instances15MonadTrans, Monad, Functor, Applicative, Foldable, Traversable, …
valuefoldStream
  1. :: State StreamK m a
  2. -> a -> StreamK m a -> m r
  3. -> a -> m r
  4. -> m r
  5. -> StreamK m a
  6. -> 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.

valuefoldStreamShared
  1. :: State StreamK m a
  2. -> a -> StreamK m a -> m r
  3. -> a -> m r
  4. -> m r
  5. -> StreamK m a
  6. -> 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.

valuefromStopK :: StopK m -> StreamK m a
#

Make an empty stream from a stop function.

valuefromYieldK :: YieldK m a -> StreamK m a
#

Make a singleton stream from a callback function. The callback function calls the one-shot yield continuation to yield an element.

valueconsK :: YieldK m a -> StreamK m a -> StreamK m a
#

Add a yield function at the head of the stream.

value(.:) :: a -> StreamK m a -> StreamK m a
#

Operator equivalent of cons.

> toList $ 1 .: 2 .: 3 .: nil
[1,2,3]
valuerepeatMWith :: (m a -> t m a -> t m a) -> m a -> t m a
#

Like repeatM but takes a stream cons operation to combine the actions in a stream specific manner. A serial cons would repeat the values serially while an async cons would repeat concurrently.

Pre-release

valuemfix :: Monad m => (m a -> StreamK m a) -> StreamK 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.

In the following example, the argument action of the function f represents the tuple (x,y) returned by it in a given iteration. We define the first element of the tuple in terms of the second.

Example1 expression
import System.IO.Unsafe (unsafeInterleaveIO)
Example1 expression
:{main = Stream.fold (Fold.drainMapM print) $ StreamK.toStream $ StreamK.mfix f    where    f action = StreamK.unCross $ do        let incr n act = fmap ((+n) . snd) $ unsafeInterleaveIO act        x <- StreamK.mkCross $ StreamK.fromStream $ Stream.sequence $ Stream.fromList [incr 1 action, incr 2 action]        y <- StreamK.mkCross $ StreamK.fromStream $ Stream.fromList [4,5]        return (x, y):}

Note: you cannot achieve this by just changing the order of the monad statements because that would change the order in which the stream elements are generated.

Note that the function f must be lazy in its argument, that's why we use unsafeInterleaveIO on action because IO monad is strict.

Pre-release

valueconcatIterateWith
  1. :: StreamK m a -> StreamK m a -> StreamK m a
  2. -> a -> StreamK m a
  3. -> StreamK m a
  4. -> StreamK m a
#

Yield an input element in the output stream, map a stream generator on it and repeat the process on the resulting stream. Resulting streams are flattened using the concatMapWith combinator. This can be used for a depth first style (DFS) traversal of a tree like structure.

Example, list a directory tree using DFS:

Example3 expressions
f = StreamK.fromStream . either Dir.readEitherPaths (const Stream.nil)input = StreamK.fromPure (Left ".")ls = StreamK.concatIterateWith StreamK.append f input

Note that iterateM is a special case of concatIterateWith:

Example1 expression
iterateM f = StreamK.concatIterateWith StreamK.append (StreamK.fromEffect . f) . StreamK.fromEffect

Pre-release

valueconcatIterateLeftsWith
  1. :: b ~ Either a c
  2. => StreamK m b -> StreamK m b -> StreamK m b
  3. -> a -> StreamK m b
  4. -> StreamK m b
  5. -> StreamK m b
#

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

Example1 expression
concatIterateLeftsWith combine f = StreamK.concatIterateWith combine (either f (const StreamK.nil))

To traverse a directory tree:

Example2 expressions
input = StreamK.fromPure (Left ".")ls = StreamK.concatIterateLeftsWith StreamK.append (StreamK.fromStream . Dir.readEither) input

Pre-release

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

valuemergeIterateWith
  1. :: StreamK m a -> StreamK m a -> StreamK m a
  2. -> a -> StreamK m a
  3. -> StreamK m a
  4. -> StreamK m a
#

Like concatIterateWith but uses the pairwise flattening combinator mergeMapWith for flattening the resulting streams. This can be used for a balanced traversal of a tree like structure.

Example, list a directory tree using balanced traversal:

Example3 expressions
f = StreamK.fromStream . either Dir.readEitherPaths (const Stream.nil)input = StreamK.fromPure (Left ".")ls = StreamK.mergeIterateWith StreamK.interleave f input

Pre-release

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

Constructors

Instances10Functor, Foldable, Traversable, IsList, Read, Show, …
valuefromStream :: Monad m => Stream m a -> StreamK m a
#

Convert a fused Stream to StreamK.

For example:

Example3 expressions
s1 = StreamK.fromStream $ Stream.fromList [1,2]s2 = StreamK.fromStream $ Stream.fromList [3,4]Stream.fold Fold.toList $ StreamK.toStream $ s1 `StreamK.append` s2[1,2,3,4]

Specialized Generation

valuerepeatM :: Monad m => m a -> StreamK m a
#
Example3 expressions
repeatM = StreamK.sequence . StreamK.repeatrepeatM = fix . StreamK.consMrepeatM = cycle1 . StreamK.fromEffect

Generate a stream by repeatedly executing a monadic action forever.

Example1 expression
:{repeatAction =       StreamK.repeatM (threadDelay 1000000 >> print 1)     & StreamK.take 10     & StreamK.fold Fold.drain:}
valueiterate :: (a -> a) -> a -> StreamK m a
#
Example1 expression
iterate f x = x `StreamK.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
StreamK.toList $ StreamK.take 5 $ StreamK.iterate (+1) 1[1,2,3,4,5]
valueiterateM :: Monad m => (a -> m a) -> m a -> StreamK m a
#
Example1 expression
iterateM f m = m >>= \a -> return a `StreamK.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.

Example1 expression
:{StreamK.iterateM (\x -> print x >> return (x + 1)) (return 0)    & StreamK.take 3    & StreamK.toList:}01[0,1,2]

Elimination

0 declarations

General Folds

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

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

Definitions:

Example2 expressions
fold f = fmap fst . StreamK.foldBreak ffold f = StreamK.parseD (Parser.fromFold f)

Example:

Example1 expression
StreamK.fold Fold.sum $ StreamK.fromStream $ Stream.enumerateFromTo 1 1005050
valuefoldBreak :: Monad m => Fold m a b -> StreamK m a -> m (b, StreamK m a)
#

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

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

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

Internal

valuefoldConcat
  1. :: Monad m
  2. => Producer m a b
  3. -> Fold m b c
  4. -> StreamK m a
  5. -> m (c, StreamK m a)
#

Generate streams from individual elements of a stream and fold the concatenation of those streams using the supplied fold. Return the result of the fold and residual stream.

For example, this can be used to efficiently fold an Array Word8 stream using Word8 folds.

Internal

Specialized Folds

Map and Fold

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

Apply a monadic action to each element of the stream and discard the output of the action.

Conversions

Transformation

0 declarations

By folding (scans)

Filtering

Mapping

Inserting

Deleting

Reordering

valuesortBy :: Monad m => (a -> a -> Ordering) -> StreamK m a -> StreamK m a
#

Sort the input stream using a supplied comparison function.

Sorting can be achieved by simply:

Example1 expression
sortBy cmp = StreamK.mergeMapWith (StreamK.mergeBy cmp) StreamK.fromPure

However, this combinator uses a parser to first split the input stream into down and up sorted segments and then merges them to optimize sorting when pre-sorted sequences exist in the input stream.

O(n) space

Map and Filter

Zipping

Merging

Transformation comprehensions

Exceptions

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

Like Streamly.Data.Stream.Streamly.Data.Stream.handle but with one significant difference, this function observes exceptions from the consumer of the stream as well.

You can also convert StreamK to Stream and use exception handling from Stream module:

Example1 expression
handle f s = StreamK.fromStream $ Stream.handle (\e -> StreamK.toStream (f e)) (StreamK.toStream s)

Resource Management

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

Like Streamly.Data.Stream.Streamly.Data.Stream.bracketIO but with one significant difference, this function observes exceptions from the consumer of the stream as well. Therefore, it cleans up the resource promptly when the consumer encounters an exception.

You can also convert StreamK to Stream and use resource handling from Stream module:

Example1 expression
bracketIO bef aft bet = StreamK.fromStream $ Stream.bracketIO bef aft (StreamK.toStream . bet)