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

  • 12 types
  • 59 values
  • Packagestreamly-0.10.1
  • Exports71
  • LanguageHaskell2010
  • LicenceBSD-3-Clause
  • SourceDispatcher.hs

This is a magic number and it is overloaded, and used at several places to achieve batching:

  1. If we have to sleep to slowdown this is the minimum period that we accumulate before we sleep. Also, workers do not stop until this much sleep time is accumulated.

  2. Collected latencies are computed and transferred to measured latency after a minimum of this period.

valueallThreadsDone :: MonadIO m => IORef (Set ThreadId) -> m Bool
#

This is safe even if we are adding more threads concurrently because if a child thread is adding another thread then anyway workerThreads will not be empty.

datadata Config
#

An abstract type for specifying the configuration parameters of a Channel. Use Config -> Config modifier functions to modify the default configuration. See the individual modifier documentation for default values.

valuemaxThreads :: Int -> Config -> Config
#

Specify the maximum number of threads that can be spawned by the channel. A value of 0 resets the thread limit to default, a negative value means there is no limit. The default value is 1500.

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.

valuemaxBuffer :: Int -> Config -> Config
#

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.

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.

Constructors

valuerate :: Maybe Rate -> Config -> Config
#

Specify the stream evaluation rate of a channel.

A Nothing value means there is no smart rate control, concurrent execution blocks only if maxThreads or maxBuffer is reached, or there are no more concurrent tasks to execute. This is the default.

When rate (throughput) is specified, concurrent production may be ramped up or down automatically to achieve the specified stream throughput. The specific behavior for different styles of Rate specifications is documented under Rate. The effective maximum production rate achieved by a channel 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

Maximum production rate is given by:

rate = \frac{maxThreads}{latency}

If we know the average latency of the tasks we can set maxThreads accordingly.

valueavgRate :: Double -> Config -> Config
#

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.

valueminRate :: Double -> Config -> Config
#

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.

valuemaxRate :: Double -> Config -> Config
#

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.

valueconstRate :: Double -> Config -> Config
#

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.

datadata StopWhen
#

Specify when the Channel should stop.

Constructors

valueeager :: Bool -> Config -> Config
#

By default, processing of output from the worker threads is given priority over dispatching new workers. More workers are dispatched only when there is no output to process. When eager is set to True, workers are dispatched aggresively as long as there is more work to do irrespective of whether there is output pending to be processed by the stream consumer. However, dispatching may stop if maxThreads or maxBuffer is reached.

Note: This option has no effect when rate has been specified.

Note: Not supported with interleaved.

valueordered :: Bool -> Config -> Config
#

When enabled the streams may be evaluated cocnurrently but the results are produced in the same sequence as a serial evaluation would produce.

Note: Not supported with interleaved.

valueinterleaved :: Bool -> Config -> Config
#

Interleave the streams fairly instead of prioritizing the left stream. This schedules all streams in a round robin fashion over limited number of threads.

Note: Can only be used on finite number of streams.

Note: Not supported with ordered.

valueboundThreads :: Bool -> Config -> Config
#

Spawn bound threads (i.e., spawn threads using forkOS instead of forkIO). The default value is False.

Currently, this only takes effect only for concurrent folds.

valuedefaultConfig :: Config
#

The fields prefixed by an _ are not to be accessed or updated directly but via smart accessor APIs. Use get/set routines instead of directly accessing the Config fields

newtypenewtype Count
#

Constructors

Instances9Bounded, Enum, Eq, Integral, Num, Ord, …
  • Bounded CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Enum CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Eq CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Integral CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Num CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Ord CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Read CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Real CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Show CountDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
datadata Limit
#
Instances3Eq, Ord, Show
  • Eq LimitDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Ord LimitDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
  • Show LimitDefined in streamly-0.10.1 · Streamly.Internal.Data.Channel.Types
datadata ThreadAbort
#

Channel driver throws this exception to all active workers to clean up the channel.

Instances2Show, Exception
datadata YieldRateInfo
#

Rate control.

Constructors

  • YieldRateInfo
    • svarLatencyTarget :: NanoSecond64
    • svarLatencyRange :: LatencyRange
    • svarRateBuffer :: Int
    • svarGainedLostYields :: IORef Count
      LOCKING

      Unlocked access. Modified by the consumer thread and unsafely read by the worker threads

    • svarAllTimeLatency :: IORef (Count, AbsTime)

      Actual latency/througput as seen from the consumer side, we count the yields and the time it took to generates those yields. This is used to increase or decrease the number of workers needed to achieve the desired rate. The idle time of workers is adjusted in this, so that we only account for the rate when the consumer actually demands data. XXX interval latency is enough, we can move this under diagnostics build [LOCKING] Unlocked access. Modified by the consumer thread and unsafely read by the worker threads

    • workerBootstrapLatency :: Maybe NanoSecond64
    • workerPollingInterval :: IORef Count

      After how many yields the worker should update the latency information. If the latency is high, this count is kept lower and vice-versa. XXX If the latency suddenly becomes too high this count may remain too high for long time, in such cases the consumer can change it. 0 means no latency computation XXX this is derivable from workerMeasuredLatency, can be removed. [LOCKING] Unlocked access. Modified by the consumer thread and unsafely read by the worker threads

    • workerPendingLatency :: IORef (Count, Count, NanoSecond64)

      This is in progress latency stats maintained by the workers which we empty into workerCollectedLatency stats at certain intervals - whenever we process the stream elements yielded in this period. The first count is all yields, the second count is only those yields for which the latency was measured to be non-zero (note that if the timer resolution is low the measured latency may be zero e.g. on JS platform). [LOCKING] Locked access. Modified by the consumer thread as well as worker threads. Workers modify it periodically based on workerPollingInterval and not on every yield to reduce the locking overhead. (allYieldCount, yieldCount, timeTaken)

    • workerCollectedLatency :: IORef (Count, Count, NanoSecond64)

      This is the second level stat which is an accmulation from workerPendingLatency stats. We keep accumulating latencies in this bucket until we have stats for a sufficient period and then we reset it to start collecting for the next period and retain the computed average latency for the last period in workerMeasuredLatency. [LOCKING] Unlocked access. Modified by the consumer thread and unsafely read by the worker threads (allYieldCount, yieldCount, timeTaken)

    • workerMeasuredLatency :: IORef NanoSecond64

      Latency as measured by workers, aggregated for the last period. [LOCKING] Unlocked access. Modified by the consumer thread and unsafely read by the worker threads

A worker decrements the yield limit before it executes an action. However, the action may not result in an element being yielded, in that case we have to increment the yield limit.

Note that we need it to be an Int type so that we have the ability to undo a decrement that takes it below zero.

valuemagicMaxBuffer :: Word
#

A magical value for the buffer size arrived at by running the smallest possible task and measuring the optimal value of the buffer for that. This is obviously dependent on hardware, this figure is based on a 2.2GHz intel core-i7 processor.

valuewithDiagMVar :: Bool -> IO String -> String -> IO () -> IO ()
#

MVar diagnostics has some overhead - around 5% on AsyncT null benchmark, we can keep it on in production to debug problems quickly if and when they happen, but it may result in unexpected output when threads are left hanging until they are GCed because the consumer went away.