Modulestreamly-0.10.1Haskell2010
Streamly.Internal.Data.SVar
Deprecated. SVar is replaced by Channel.
- 18 types
- 73 values
- Packagestreamly-0.10.1
- Exports91
- LanguageHaskell2010
- LicenceBSD-3-Clause
- SourceType.hs
Constructors
SVarsvarStyle :: SVarStylesvarMrun :: RunInIO msvarStopStyle :: SVarStopStylesvarStopBy :: IORef ThreadIdoutputQueue :: IORef ([ChildEvent a], Int)outputDoorBell :: MVar ()readOutputQ :: m [ChildEvent a]postProcess :: m BooloutputQueueFromConsumer :: IORef ([ChildEvent a], Int)outputDoorBellFromConsumer :: MVar ()maxWorkerLimit :: LimitmaxBufferLimit :: LimitpushBufferSpace :: IORef CountpushBufferPolicy :: PushBufferPolicypushBufferMVar :: MVar ()remainingWork :: Maybe (IORef Count)yieldRateInfo :: Maybe YieldRateInfoenqueue :: (RunInIO m, t m a) -> IO ()isWorkDone :: IO BoolisQueueDone :: IO BoolneedDoorBell :: IORef BoolworkLoop :: Maybe WorkerInfo -> m ()workerThreads :: IORef (Set ThreadId)workerCount :: IORef IntaccountThread :: ThreadId -> m ()workerStopMVar :: MVar ()svarStats :: SVarStatssvarRef :: Maybe (IORef ())svarInspectMode :: BoolsvarCreator :: ThreadIdoutputHeap :: IORef (Heap (Entry Int (AheadHeapEntry t m a)), Maybe Int)aheadWorkQueue :: IORef ([t m a], Int)
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)
Instances9Bounded, Enum, Eq, Integral, Num, Ord, …
Bounded CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeEnum CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeEq CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeIntegral CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeNum CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeOrd CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeRead CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeReal CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeShow CountDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.Type
Instances2Show, Exception
Show ThreadAbortDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeException ThreadAbortDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.Type
Events that a child thread may send to a parent thread.
Constructors
Constructors
An SVar or a Stream Var is a conduit to the output from multiple streams running concurrently and asynchronously. An SVar can be thought of as an asynchronous IO handle. We can write any number of streams to an SVar in a non-blocking manner and then read them back at any time at any pace. The SVar would run the streams asynchronously and accumulate results. An SVar may not really execute the stream completely and accumulate all the results. However, it ensures that the reader can read the results at whatever paces it wants to read. The SVar monitors and adapts to the consumer's pace.
An SVar is a mini scheduler, it has an associated workLoop that holds the
stream tasks to be picked and run by a pool of worker threads. It has an
associated output queue where the output stream elements are placed by the
worker threads. A outputDoorBell is used by the worker threads to intimate the
consumer thread about availability of new results in the output queue. More
workers are added to the SVar by fromStreamVar on demand if the output
produced is not keeping pace with the consumer. On bounded SVars, workers
block on the output queue to provide throttling of the producer when the
consumer is not pulling fast enough. The number of workers may even get
reduced depending on the consuming pace.
New work is enqueued either at the time of creation of the SVar or as a
result of executing the parallel combinators i.e. <| and <|> when the
already enqueued computations get evaluated. See joinStreamVarAsync.
Constructors
Constructors
Instances1Show
Show LatencyRangeDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.Type
Constructors
YieldRateInfosvarLatencyTarget :: NanoSecond64svarLatencyRange :: LatencyRangesvarRateBuffer :: IntsvarGainedLostYields :: IORef CountsvarAllTimeLatency :: IORef (Count, AbsTime)workerBootstrapLatency :: Maybe NanoSecond64workerPollingInterval :: IORef CountworkerPendingLatency :: IORef (Count, Count, NanoSecond64)workerCollectedLatency :: IORef (Count, Count, NanoSecond64)workerMeasuredLatency :: IORef NanoSecond64
Adapt the stream state from one type to another.
Instances2Eq, Show
Eq SVarStopStyleDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.TypeShow SVarStopStyleDefined in streamly-core-0.2.2 · Streamly.Internal.Data.SVar.Type
Identify the type of the SVar. Two computations using the same style can be scheduled on the same SVar.
Sorting out-of-turn outputs in a heap for Ahead style streams
Constructors
AheadEntryNullAheadEntryPure aAheadEntryStream (RunInIO m, t m a)
Buffering policy for persistent push workers (in ParallelT). In a pull style SVar (in AsyncT, AheadT etc.), the consumer side dispatches workers on demand, workers terminate if the buffer is full or if the consumer is not cosuming fast enough. In a push style SVar, a worker is dispatched only once, workers are persistent and keep pushing work to the consumer via a bounded buffer. If the buffer becomes full the worker either blocks, or it can drop an item from the buffer to make space.
Pull style SVars are useful in lazy stream evaluation whereas push style SVars are useful in strict left Folds.
XXX Maybe we can separate the implementation in two different types instead of using a common SVar type.
This is a magic number and it is overloaded, and used at several places to achieve batching:
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.
Collected latencies are computed and transferred to measured latency after a minimum of this period.
Constructors
This function is used by the producer threads to queue output for the consumer thread to consume. Returns whether the queue has more space.
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.
In contrast to pushWorker which always happens only from the consumer thread, a pushWorkerPar can happen concurrently from multiple threads on the producer side. So we need to use a thread safe modification of workerThreads. Alternatively, we can use a CreateThread event to avoid using a CAS based modification.