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

Moduleconduit-1.3.6.1Haskell2010

Data.Conduit.Combinators

This module is meant as a replacement for Data.Conduit.List. That module follows a naming scheme which was originally inspired by its enumerator roots. This module is meant to introduce a naming scheme which encourages conduit best practices.

There are two versions of functions in this module. Those with a trailing E work in the individual elements of a chunk of data, e.g., the bytes of a ByteString, the Chars of a Text, or the Ints of a Vector Int. Those without a trailing E work on unchunked streams.

FIXME: discuss overall naming, usage of mono-traversable, etc

Mention take (Conduit) vs drop (Consumer)

  • 1 type
  • 150 values
  • Packageconduit-1.3.6.1
  • Exports151
  • LanguageHaskell2010
  • LicenceMIT
  • SourceCombinators.hs

Producers

0 declarations

Pure

valueyieldMany
  1. :: (Monad m, MonoFoldable mono)
  2. => mono
  3. -> ConduitT i (Element mono) m ()
#

Yield each of the values contained by the given MonoFoldable.

This will work on many data structures, including lists, ByteStrings, and Vectors.

Subject to fusion

valueunfold :: Monad m => (b -> Maybe (a, b)) -> b -> ConduitT i a m ()
#

Generate a producer from a seed value.

Subject to fusion

valueenumFromTo :: (Monad m, Enum a, Ord a) => a -> a -> ConduitT i a m ()
#

Enumerate from a value to a final value, inclusive, via succ.

This is generally more efficient than using Prelude's enumFromTo and combining with sourceList since this avoids any intermediate data structures.

Subject to fusion

valueiterate :: Monad m => (a -> a) -> a -> ConduitT i a m ()
#

Produces an infinite stream of repeated applications of f to x.

Subject to fusion

valuerepeat :: Monad m => a -> ConduitT i a m ()
#

Produce an infinite stream consisting entirely of the given value.

Subject to fusion

valuereplicate :: Monad m => Int -> a -> ConduitT i a m ()
#

Produce a finite stream consisting of n copies of the given value.

Subject to fusion

valuesourceLazy
  1. :: (Monad m, LazySequence lazy strict)
  2. => lazy
  3. -> ConduitT i strict m ()
#

Generate a producer by yielding each of the strict chunks in a LazySequence.

For more information, see toChunks.

Subject to fusion

Monadic

valuerepeatM :: Monad m => m a -> ConduitT i a m ()
#

Repeatedly run the given action and yield all values it produces.

Subject to fusion

valuerepeatWhileM :: Monad m => m a -> (a -> Bool) -> ConduitT i a m ()
#

Repeatedly run the given action and yield all values it produces, until the provided predicate returns False.

Subject to fusion

valuereplicateM :: Monad m => Int -> m a -> ConduitT i a m ()
#

Perform the given action n times, yielding each result.

Subject to fusion

I/O

valuesourceHandle :: MonadIO m => Handle -> ConduitT i ByteString m ()
#

Stream the contents of a Handle as binary data. Note that this function will not automatically close the Handle when processing completes, since it did not acquire the Handle in the first place.

Same as sourceHandle, but instead of allocating a new buffer for each incoming chunk of data, reuses the same buffer. Therefore, the ByteStrings yielded by this function are not referentially transparent between two different yields.

This function will be slightly more efficient than sourceHandle by avoiding allocations and reducing garbage collections, but should only be used if you can guarantee that you do not reuse a ByteString (or any slice thereof) between two calls to await.

Filesystem

Stream the contents of the given directory, without traversing deeply.

This function will return all of the contents of the directory, whether they be files, directories, etc.

Note that the generated filepaths will be the complete path, not just the filename. In other words, if you have a directory foo containing files bar and baz, and you use sourceDirectory on foo, the results will be foo/bar and foo/baz.

valuesourceDirectoryDeep
  1. :: MonadResource m
  2. => Bool

    Follow directory symlinks

  3. -> FilePath

    Root directory

  4. -> ConduitT i FilePath m ()
#

Deeply stream the contents of the given directory.

This works the same as sourceDirectory, but will not return directories at all. This function also takes an extra parameter to indicate whether symlinks will be followed.

Consumers

0 declarations

Pure

valuedrop :: Monad m => Int -> ConduitT a o m ()
#

Ignore a certain number of values in the stream.

Note: since this function doesn't produce anything, you probably want to use it with (>>) instead of directly plugging it into a pipeline:

Example2 expressions
runConduit $ yieldMany [1..5] .| drop 2 .| sinkList[]runConduit $ yieldMany [1..5] .| (drop 2 >> sinkList)[3,4,5]
valuedropE :: (Monad m, IsSequence seq) => Index seq -> ConduitT seq o m ()
#

Drop a certain number of elements from a chunked stream.

Note: you likely want to use it with monadic composition. See the docs for drop.

valuedropWhile :: Monad m => (a -> Bool) -> ConduitT a o m ()
#

Drop all values which match the given predicate.

Note: you likely want to use it with monadic composition. See the docs for drop.

valuedropWhileE
  1. :: (Monad m, IsSequence seq)
  2. => Element seq -> Bool
  3. -> ConduitT seq o m ()
#

Drop all elements in the chunked stream which match the given predicate.

Note: you likely want to use it with monadic composition. See the docs for drop.

valuefoldl :: Monad m => (a -> b -> a) -> a -> ConduitT b o m a
#

A strict left fold.

Subject to fusion

valuefoldl1 :: Monad m => (a -> a -> a) -> ConduitT a o m (Maybe a)
#

A strict left fold with no starting value. Returns Nothing when the stream is empty.

Subject to fusion

valuefoldMap :: (Monad m, Monoid b) => (a -> b) -> ConduitT a o m b
#

Apply the provided mapping function and monoidal combine all values.

Subject to fusion

valueall :: Monad m => (a -> Bool) -> ConduitT a o m Bool
#

Check that all values in the stream return True.

Subject to shortcut logic: at the first False, consumption of the stream will stop.

Subject to fusion

valueallE
  1. :: (Monad m, MonoFoldable mono)
  2. => Element mono -> Bool
  3. -> ConduitT mono o m Bool
#

Check that all elements in the chunked stream return True.

Subject to shortcut logic: at the first False, consumption of the stream will stop.

Subject to fusion

valueany :: Monad m => (a -> Bool) -> ConduitT a o m Bool
#

Check that at least one value in the stream returns True.

Subject to shortcut logic: at the first True, consumption of the stream will stop.

Subject to fusion

valueanyE
  1. :: (Monad m, MonoFoldable mono)
  2. => Element mono -> Bool
  3. -> ConduitT mono o m Bool
#

Check that at least one element in the chunked stream returns True.

Subject to shortcut logic: at the first True, consumption of the stream will stop.

Subject to fusion

valueand :: Monad m => ConduitT Bool o m Bool
#

Are all values in the stream True?

Consumption stops once the first False is encountered.

Subject to fusion

valueor :: Monad m => ConduitT Bool o m Bool
#

Are any values in the stream True?

Consumption stops once the first True is encountered.

Subject to fusion

valueelem :: (Monad m, Eq a) => a -> ConduitT a o m Bool
#

Are any values in the stream equal to the given value?

Stops consuming as soon as a match is found.

Subject to fusion

valuenotElem :: (Monad m, Eq a) => a -> ConduitT a o m Bool
#

Are no values in the stream equal to the given value?

Stops consuming as soon as a match is found.

Subject to fusion

valuesinkLazy :: (Monad m, LazySequence lazy strict) => ConduitT strict o m lazy
#

Consume all incoming strict chunks into a lazy sequence. Note that the entirety of the sequence will be resident at memory.

This can be used to consume a stream of strict ByteStrings into a lazy ByteString, for example.

Subject to fusion

valuesinkList :: Monad m => ConduitT a o m [a]
#

Consume all values from the stream and return as a list. Note that this will pull all values into memory.

Subject to fusion

valuesinkVector :: (Vector v a, PrimMonad m) => ConduitT a o m (v a)
#

Sink incoming values into a vector, growing the vector as necessary to fit more elements.

Note that using this function is more memory efficient than sinkList and then converting to a Vector, as it avoids intermediate list constructors.

Subject to fusion

valuesinkVectorN
  1. :: (Vector v a, PrimMonad m)
  2. => Int

    maximum allowed size

  3. -> ConduitT a o m (v a)
#

Sink incoming values into a vector, up until size maxSize. Subsequent values will be left in the stream. If there are less than maxSize values present, returns a Vector of smaller size.

Note that using this function is more memory efficient than sinkList and then converting to a Vector, as it avoids intermediate list constructors.

Subject to fusion

Same as sinkBuilder, but afterwards convert the builder to its lazy representation.

Alternatively, this could be considered an alternative to sinkLazy, with the following differences:

  • This function will allow multiple input types, not just the strict version of the lazy structure.

  • Some buffer copying may occur in this version.

Subject to fusion

valuesinkNull :: Monad m => ConduitT a o m ()
#

Consume and discard all remaining values in the stream.

Subject to fusion

valueheadDef :: Monad m => a -> ConduitT a o m a
#

Same as head, but returns a default value if none are available from the stream.

valuelast :: Monad m => ConduitT a o m (Maybe a)
#

Retrieve the last value in the stream, if present.

Subject to fusion

valuelastDef :: Monad m => a -> ConduitT a o m a
#

Same as last, but returns a default value if none are available from the stream.

valuelength :: (Monad m, Num len) => ConduitT a o m len
#

Count how many values are in the stream.

Subject to fusion

valuelengthIf :: (Monad m, Num len) => (a -> Bool) -> ConduitT a o m len
#

Count how many values in the stream pass the given predicate.

Subject to fusion

valuenull :: Monad m => ConduitT a o m Bool
#

True if there are no values in the stream.

This function does not modify the stream.

valuenullE :: (Monad m, MonoFoldable mono) => ConduitT mono o m Bool
#

True if there are no elements in the chunked stream.

This function may remove empty leading chunks from the stream, but otherwise will not modify it.

valuesum :: (Monad m, Num a) => ConduitT a o m a
#

Get the sum of all values in the stream.

Subject to fusion

Monadic

valuemapM_ :: Monad m => (a -> m ()) -> ConduitT a o m ()
#

Apply the action to all values in the stream.

Note: if you want to pass the values instead of consuming them, use iterM instead.

Subject to fusion

valuemapM_E
  1. :: (Monad m, MonoFoldable mono)
  2. => Element mono -> m ()
  3. -> ConduitT mono o m ()
#

Apply the action to all elements in the chunked stream.

Note: the same caveat as with mapM_ applies. If you don't want to consume the values, you can use iterM:

iterM (omapM_ f)

Subject to fusion

valuefoldM :: Monad m => (a -> b -> m a) -> a -> ConduitT b o m a
#

A monadic strict left fold.

Subject to fusion

valuefoldMapM :: (Monad m, Monoid w) => (a -> m w) -> ConduitT a o m w
#

Apply the provided monadic mapping function and monoidal combine all values.

Subject to fusion

I/O

Cautious version of sinkFile. The idea here is to stream the values to a temporary file in the same directory of the destination file, and only on successfully writing the entire file, moves it atomically to the destination path.

In the event of an exception occurring, the temporary file will be deleted and no move will be made. If the application shuts down without running exception handling (such as machine failure or a SIGKILL), the temporary file will remain and the destination file will be untouched.

valuesinkTempFile
  1. :: MonadResource m
  2. => FilePath

    temp directory

  3. -> String

    filename pattern

  4. -> ConduitM ByteString o m FilePath
#

Stream data into a temporary file in the given directory with the given filename pattern, and return the temporary filename. The temporary file will be automatically deleted when exiting the active ResourceT block, if it still exists.

valuesinkHandleBuilder :: MonadIO m => Handle -> ConduitM Builder o m ()
#

Stream incoming builders, executing them directly on the buffer of the given Handle. Note that this function does not automatically close the Handle when processing completes. Pass flush to flush the buffer.

Transformers

0 declarations

Pure

valuemap :: Monad m => (a -> b) -> ConduitT a b m ()
#

Apply a transformation to all values in a stream.

Subject to fusion

valuemapE :: (Monad m, Functor f) => (a -> b) -> ConduitT (f a) (f b) m ()
#

Apply a transformation to all elements in a chunked stream.

Subject to fusion

valueomapE
  1. :: (Monad m, MonoFunctor mono)
  2. => Element mono -> Element mono
  3. -> ConduitT mono mono m ()
#

Apply a monomorphic transformation to all elements in a chunked stream.

Unlike mapE, this will work on types like ByteString and Text which are MonoFunctor but not Functor.

Subject to fusion

valueconcatMap
  1. :: (Monad m, MonoFoldable mono)
  2. => a -> mono
  3. -> ConduitT a (Element mono) m ()
#

Apply the function to each value in the stream, resulting in a foldable value (e.g., a list). Then yield each of the individual values in that foldable value separately.

Generalizes concatMap, mapMaybe, and mapFoldable.

Subject to fusion

valueconcatMapE
  1. :: (Monad m, MonoFoldable mono, Monoid w)
  2. => Element mono -> w
  3. -> ConduitT mono w m ()
#

Apply the function to each element in the chunked stream, resulting in a foldable value (e.g., a list). Then yield each of the individual values in that foldable value separately.

Generalizes concatMap, mapMaybe, and mapFoldable.

Subject to fusion

valuetake :: Monad m => Int -> ConduitT a a m ()
#

Stream up to n number of values downstream.

Note that, if downstream terminates early, not all values will be consumed. If you want to force exactly the given number of values to be consumed, see takeExactly.

Subject to fusion

valuetakeE :: (Monad m, IsSequence seq) => Index seq -> ConduitT seq seq m ()
#

Stream up to n number of elements downstream in a chunked stream.

Note that, if downstream terminates early, not all values will be consumed. If you want to force exactly the given number of values to be consumed, see takeExactlyE.

valuetakeWhile :: Monad m => (a -> Bool) -> ConduitT a a m ()
#

Stream all values downstream that match the given predicate.

Same caveats regarding downstream termination apply as with take.

valuetakeExactly :: Monad m => Int -> ConduitT a b m r -> ConduitT a b m r
#

Consume precisely the given number of values and feed them downstream.

This function is in contrast to take, which will only consume up to the given number of values, and will terminate early if downstream terminates early. This function will discard any additional values in the stream if they are unconsumed.

Note that this function takes a downstream ConduitT as a parameter, as opposed to working with normal fusion. For more information, see http://www.yesodweb.com/blog/2013/10/core-flaw-pipes-conduit, the section titled "pipes and conduit: isolate".

valueconcat :: (Monad m, MonoFoldable mono) => ConduitT mono (Element mono) m ()
#

Flatten out a stream by yielding the values contained in an incoming MonoFoldable as individually yielded values.

Subject to fusion

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

Keep only values in the stream passing a given predicate.

Subject to fusion

valueconduitVector
  1. :: (Vector v a, PrimMonad m)
  2. => Int

    maximum allowed size

  3. -> ConduitT a (v a) m ()
#

Break up a stream of values into vectors of size n. The final vector may be smaller than n if the total number of values is not a strict multiple of n. No empty vectors will be yielded.

valuemapAccumWhile
  1. :: Monad m
  2. => a -> s -> Either s (s, b)
  3. -> s
  4. -> ConduitT a b m s
#

mapWhile with a break condition dependent on a strict accumulator. Equivalently, mapAccum as long as the result is Right. Instead of producing a leftover, the breaking input determines the resulting accumulator via Left.

Subject to fusion

valueintersperse :: Monad m => a -> ConduitT a a m ()
#

Insert the given value between each two values in the stream.

Subject to fusion

valueslidingWindow
  1. :: (Monad m, IsSequence seq, Element seq ~ a)
  2. => Int
  3. -> ConduitT a seq m ()
#

Sliding window of values 1,2,3,4,5 with window size 2 gives [1,2],[2,3],[3,4],[4,5]

Best used with structures that support O(1) snoc.

Subject to fusion

Monadic

valuemapM :: Monad m => (a -> m b) -> ConduitT a b m ()
#

Apply a monadic transformation to all values in a stream.

If you do not need the transformed values, and instead just want the monadic side-effects of running the action, see mapM_.

Subject to fusion

valuemapME :: (Monad m, Traversable f) => (a -> m b) -> ConduitT (f a) (f b) m ()
#

Apply a monadic transformation to all elements in a chunked stream.

Subject to fusion

valueomapME
  1. :: (Monad m, MonoTraversable mono)
  2. => Element mono -> m (Element mono)
  3. -> ConduitT mono mono m ()
#

Apply a monadic monomorphic transformation to all elements in a chunked stream.

Unlike mapME, this will work on types like ByteString and Text which are MonoFunctor but not Functor.

Subject to fusion

valueconcatMapM
  1. :: (Monad m, MonoFoldable mono)
  2. => a -> m mono
  3. -> ConduitT a (Element mono) m ()
#

Apply the monadic function to each value in the stream, resulting in a foldable value (e.g., a list). Then yield each of the individual values in that foldable value separately.

Generalizes concatMapM, mapMaybeM, and mapFoldableM.

Subject to fusion

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

Keep only values in the stream passing a given monadic predicate.

Subject to fusion

valueiterM :: Monad m => (a -> m ()) -> ConduitT a a m ()
#

Apply a monadic action on all values in a stream.

This Conduit can be used to perform a monadic side-effect for every value, whilst passing the value through the Conduit as-is.

iterM f = mapM (\a -> f a >>= \() -> return a)

Subject to fusion

valuescanlM :: Monad m => (a -> b -> m a) -> a -> ConduitT b a m ()
#

Analog of scanl for lists, monadic.

Subject to fusion

Textual

valueline
  1. :: (Monad m, IsSequence seq, Element seq ~ Char)
  2. => ConduitT seq o m r
  3. -> ConduitT seq o m r
#

Stream in the entirety of a single line.

Like takeExactly, this will consume the entirety of the line regardless of the behavior of the inner Conduit.

valuelinesUnbounded :: (Monad m, IsSequence seq, Element seq ~ Char) => ConduitT seq seq m ()
#

Convert a stream of arbitrarily-chunked textual data into a stream of data where each chunk represents a single line. Note that, if you have unknown or untrusted input, this function is unsafe, since it would allow an attacker to form lines of massive length and exhaust memory.

Subject to fusion

valuesplitOnUnboundedE
  1. :: (Monad m, IsSequence seq)
  2. => Element seq -> Bool
  3. -> ConduitT seq seq m ()
#

Split a stream of arbitrarily-chunked data, based on a predicate on elements. Elements that satisfy the predicate will cause chunks to be split, and aren't included in these output chunks. Note that, if you have unknown or untrusted input, this function is unsafe, since it would allow an attacker to form chunks of massive length and exhaust memory.

Builders

Incrementally execute builders on the given buffer and pass on the filled chunks as bytestrings. Note that, if the given buffer is too small for the execution of a build step, a larger one will be allocated.

WARNING: This conduit yields bytestrings that are NOT referentially transparent. Their content will be overwritten as soon as control is returned from the inner sink!

typetype BufferAllocStrategy = (IO Buffer, Int -> Buffer -> IO (IO Buffer))
#

A buffer allocation strategy (buf0, nextBuf) specifies the initial buffer to use and how to compute a new buffer nextBuf minSize buf with at least size minSize from a filled buffer buf. The double nesting of the IO monad helps to ensure that the reference to the filled buffer buf is lost as soon as possible, but the new buffer doesn't have to be allocated too early.

The simplest buffer allocation strategy: whenever a buffer is requested, allocate a new one that is big enough for the next build step to execute.

NOTE that this allocation strategy may spill quite some memory upon direct insertion of a bytestring by the builder. Thats no problem for garbage collection, but it may lead to unreasonably high memory consumption in special circumstances.

An unsafe, but possibly more efficient buffer allocation strategy: reuse the buffer, if it is big enough for the next build step to execute.

Special

4 declarations
valuevectorBuilder
  1. :: (PrimMonad m, PrimMonad n, Vector v e, PrimState m ~ PrimState n)
  2. => Int

    size

  3. -> ((e -> n ()) -> ConduitT i Void m r)
  4. -> ConduitT i (v e) m r
#

Generally speaking, yielding values from inside a Conduit requires some allocation for constructors. This can introduce an overhead, similar to the overhead needed to represent a list of values instead of a vector. This overhead is even more severe when talking about unboxed values.

This combinator allows you to overcome this overhead, and efficiently fill up vectors. It takes two parameters. The first is the size of each mutable vector to be allocated. The second is a function. The function takes an argument which will yield the next value into a mutable vector.

Under the surface, this function uses a number of tricks to get high performance. For more information on both usage and implementation, please see: https://www.schoolofhaskell.com/user/snoyberg/library-documentation/vectorbuilder

valuemapAccumS
  1. :: Monad m
  2. => a -> s -> ConduitT b Void m s
  3. -> s
  4. -> ConduitT () b m ()
  5. -> ConduitT a Void m s
#

Consume a source with a strict accumulator, in a way piecewise defined by a controlling stream. The latter will be evaluated until it terminates.

Example2 expressions
let f a s = liftM (:s) $ mapC (*a) =$ CL.take areverse $ runIdentity $ yieldMany [0..3] $$ mapAccumS f [] (yieldMany [1..])[[],[1],[4,6],[12,15,18]] :: [[Int]]
valuepeekForever :: Monad m => ConduitT i o m () -> ConduitT i o m ()
#

Run a consuming conduit repeatedly, only stopping when there is no more data available from upstream.

valuepeekForeverE
  1. :: (Monad m, MonoFoldable i)
  2. => ConduitT i o m ()
  3. -> ConduitT i o m ()
#

Run a consuming conduit repeatedly, only stopping when there is no more data available from upstream.

In contrast to peekForever, this function will ignore empty chunks of data. So for example, if a stream of data contains an empty ByteString, it is still treated as empty, and the consuming function is not called.