{-# LANGUAGE DeriveDataTypeable #-} ----------------------------------------------------------------------- -- | -- Module : Control/Concurrent/Queue/Blocking.hs -- Copyright : (c) Ivan Tomac 2012 -- License : BSD3 -- -- Maintainer : ivan `dot` tomac `at` google `dot` com -- Stability : experimental -- -- Single writer, multiple reader queue, blocking and non-blocking. -- Slightly less efficient than the strictly non-blocking version. -- -- A queue is represented by a single 'IORef' containing a read -- barrier. -- Writing to the 'Queue' is guarded by an 'MVar'. -- When the 'Queue' is initialized, the barrier is set to @Left v@ -- where @v@ is an empty 'MVar'. -- Blocking read reads the 'MVar' while the non-blocking first looks -- at whether the barrier is wrapped with a 'Left' (empty) or 'Right' -- (full) constructor. -- -- 'Queue' 's elements can safely be read multiple times from multiple -- threads. -- ----------------------------------------------------------------------- module Control.Concurrent.Queue.Blocking ( -- * The 'Queue' type Queue (..) -- * Operations , newQ , newQ' , readQ , readBlockingQ , unsafeWriteQ ) where import Control.Applicative import Control.Concurrent.MVar import Data.Function import Data.IORef import Data.Typeable newtype Queue a = Q (IORef (Barrier a)) deriving (Eq, Typeable) type Barrier a = Either (QVar a) (QVar a) type QVar a = MVar (a, Queue a) -- | Returns a new 'Queue' and a function to write to it. newQ :: IO (a -> IO (), Queue a) newQ = newQ' >>= liftA2 fmap result newMVar where result q v = (modifyMVar_ v . flip unsafeWriteQ, q) -- | Returns only the 'Queue', without the writing function. newQ' :: IO (Queue a) newQ' = Q <$> (newEmptyMVar >>= newIORef . Left) -- | Returns the next value in the 'Queue' if one is available, along -- with the rest of the 'Queue'. -- If the 'Queue' is empty, it returns 'Nothing'. readQ :: Queue a -> IO (Maybe (a, Queue a)) readQ (Q ref) = readIORef ref >>= either (const $ return Nothing) (fmap Just . readMVar) -- | Returns the next value in the 'Queue', along with the new head of -- the 'Queue'. -- If the 'Queue' is empty, it waits for a value to be written. readBlockingQ :: Queue a -> IO (a, Queue a) readBlockingQ (Q ref) = readIORef ref >>= readMVar . either id id -- | Writes a value to the 'Queue'. 'unsafeWriteQ' is not thread-safe. -- It is meant to be used for building higher level concurrency -- primitives and is not intended to be used directly. unsafeWriteQ :: Queue a -> a -> IO (Queue a) unsafeWriteQ (Q ref) x = newQ' >>= liftA2 (>>) write return where write q = do barrier <- either id id <$> readIORef ref putMVar barrier (x, q) writeIORef ref (Right barrier)