module Hydra.Node.InputQueue where
import Hydra.Prelude
import Control.Concurrent.Class.MonadSTM (
isEmptyTBQueue,
modifyTVar',
readTBQueue,
swapTVar,
writeTBQueue,
)
data InputQueue m e = InputQueue
{ forall (m :: * -> *) e. InputQueue m e -> e -> m ()
enqueue :: e -> m ()
, forall (m :: * -> *) e.
InputQueue m e -> DiffTime -> Queued e -> m ()
reenqueue :: DiffTime -> Queued e -> m ()
, forall (m :: * -> *) e. InputQueue m e -> Queued e -> m ()
park :: Queued e -> m ()
, forall (m :: * -> *) e. InputQueue m e -> m ()
releaseParked :: m ()
, forall (m :: * -> *) e. InputQueue m e -> m (Queued e)
dequeue :: m (Queued e)
, forall (m :: * -> *) e. InputQueue m e -> m Bool
isEmpty :: m Bool
}
data Queued a = Queued {forall a. Queued a -> Word64
queuedId :: Word64, forall a. Queued a -> a
queuedItem :: a}
createInputQueue ::
( MonadDelay m
, MonadAsync m
, MonadLabelledSTM m
) =>
m (InputQueue m e)
createInputQueue :: forall (m :: * -> *) e.
(MonadDelay m, MonadAsync m, MonadLabelledSTM m) =>
m (InputQueue m e)
createInputQueue = do
TVar m Integer
numThreads <- String -> Integer -> m (TVar m Integer)
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO String
"num-threads" (Integer
0 :: Integer)
TVar m Word64
nextId <- String -> Word64 -> m (TVar m Word64)
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO String
"nex-id" Word64
0
TVar m [Queued e]
parked <- String -> [Queued e] -> m (TVar m [Queued e])
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO String
"parked-queue" []
TBQueue m (Queued e)
q <- String -> Natural -> m (TBQueue m (Queued e))
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> Natural -> m (TBQueue m a)
newLabelledTBQueueIO String
"input-queue" Natural
100
InputQueue m e -> m (InputQueue m e)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure
InputQueue
{ $sel:enqueue:InputQueue :: e -> m ()
enqueue = \e
queuedItem ->
STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> STM m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
Word64
queuedId <- TVar m Word64 -> STM m Word64
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Word64
nextId
TBQueue m (Queued e) -> Queued e -> STM m ()
forall a. TBQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
writeTBQueue TBQueue m (Queued e)
q Queued{Word64
$sel:queuedId:Queued :: Word64
queuedId :: Word64
queuedId, e
$sel:queuedItem:Queued :: e
queuedItem :: e
queuedItem}
TVar m Word64 -> (Word64 -> Word64) -> STM m ()
forall a. TVar m a -> (a -> a) -> STM m ()
forall (m :: * -> *) a.
MonadSTM m =>
TVar m a -> (a -> a) -> STM m ()
modifyTVar' TVar m Word64
nextId Word64 -> Word64
forall a. Enum a => a -> a
succ
, $sel:reenqueue:InputQueue :: DiffTime -> Queued e -> m ()
reenqueue = \DiffTime
delay Queued e
e -> do
STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> STM m () -> m ()
forall a b. (a -> b) -> a -> b
$ TVar m Integer -> (Integer -> Integer) -> STM m ()
forall a. TVar m a -> (a -> a) -> STM m ()
forall (m :: * -> *) a.
MonadSTM m =>
TVar m a -> (a -> a) -> STM m ()
modifyTVar' TVar m Integer
numThreads Integer -> Integer
forall a. Enum a => a -> a
succ
m (Async m ()) -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m (Async m ()) -> m ())
-> (m () -> m (Async m ())) -> m () -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. String -> m () -> m (Async m ())
forall (m :: * -> *) a.
MonadAsync m =>
String -> m a -> m (Async m a)
asyncLabelled String
"input-queue-reenqueue" (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
DiffTime -> m ()
forall (m :: * -> *). MonadDelay m => DiffTime -> m ()
threadDelay DiffTime
delay
STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> STM m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
TVar m Integer -> (Integer -> Integer) -> STM m ()
forall a. TVar m a -> (a -> a) -> STM m ()
forall (m :: * -> *) a.
MonadSTM m =>
TVar m a -> (a -> a) -> STM m ()
modifyTVar' TVar m Integer
numThreads Integer -> Integer
forall a. Enum a => a -> a
pred
TBQueue m (Queued e) -> Queued e -> STM m ()
forall a. TBQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
writeTBQueue TBQueue m (Queued e)
q Queued e
e
, $sel:park:InputQueue :: Queued e -> m ()
park = \Queued e
e -> STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> STM m () -> m ()
forall a b. (a -> b) -> a -> b
$ TVar m [Queued e] -> ([Queued e] -> [Queued e]) -> STM m ()
forall a. TVar m a -> (a -> a) -> STM m ()
forall (m :: * -> *) a.
MonadSTM m =>
TVar m a -> (a -> a) -> STM m ()
modifyTVar' TVar m [Queued e]
parked ([Queued e] -> [Queued e] -> [Queued e]
forall a. Semigroup a => a -> a -> a
<> [Queued e
e])
, releaseParked :: m ()
releaseParked = do
[Queued e]
ps <- STM m [Queued e] -> m [Queued e]
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m [Queued e] -> m [Queued e])
-> STM m [Queued e] -> m [Queued e]
forall a b. (a -> b) -> a -> b
$ TVar m [Queued e] -> [Queued e] -> STM m [Queued e]
forall a. TVar m a -> a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m a
swapTVar TVar m [Queued e]
parked []
Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless ([Queued e] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [Queued e]
ps) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> STM m () -> m ()
forall a b. (a -> b) -> a -> b
$ TVar m Integer -> (Integer -> Integer) -> STM m ()
forall a. TVar m a -> (a -> a) -> STM m ()
forall (m :: * -> *) a.
MonadSTM m =>
TVar m a -> (a -> a) -> STM m ()
modifyTVar' TVar m Integer
numThreads Integer -> Integer
forall a. Enum a => a -> a
succ
m (Async m ()) -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m (Async m ()) -> m ())
-> (m () -> m (Async m ())) -> m () -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. String -> m () -> m (Async m ())
forall (m :: * -> *) a.
MonadAsync m =>
String -> m a -> m (Async m a)
asyncLabelled String
"input-queue-release-parked" (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
[Queued e] -> (Queued e -> m ()) -> m ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Queued e]
ps ((Queued e -> m ()) -> m ()) -> (Queued e -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> (Queued e -> STM m ()) -> Queued e -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. TBQueue m (Queued e) -> Queued e -> STM m ()
forall a. TBQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
writeTBQueue TBQueue m (Queued e)
q
STM m () -> m ()
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m () -> m ()) -> STM m () -> m ()
forall a b. (a -> b) -> a -> b
$ TVar m Integer -> (Integer -> Integer) -> STM m ()
forall a. TVar m a -> (a -> a) -> STM m ()
forall (m :: * -> *) a.
MonadSTM m =>
TVar m a -> (a -> a) -> STM m ()
modifyTVar' TVar m Integer
numThreads Integer -> Integer
forall a. Enum a => a -> a
pred
, $sel:dequeue:InputQueue :: m (Queued e)
dequeue =
STM m (Queued e) -> m (Queued e)
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m (Queued e) -> m (Queued e))
-> STM m (Queued e) -> m (Queued e)
forall a b. (a -> b) -> a -> b
$ TBQueue m (Queued e) -> STM m (Queued e)
forall a. TBQueue m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> STM m a
readTBQueue TBQueue m (Queued e)
q
, isEmpty :: m Bool
isEmpty = do
STM m Bool -> m Bool
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m Bool -> m Bool) -> STM m Bool -> m Bool
forall a b. (a -> b) -> a -> b
$ do
Integer
n <- TVar m Integer -> STM m Integer
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Integer
numThreads
Bool
isEmpty' <- TBQueue m (Queued e) -> STM m Bool
forall a. TBQueue m a -> STM m Bool
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> STM m Bool
isEmptyTBQueue TBQueue m (Queued e)
q
Bool -> STM m Bool
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Bool
isEmpty' Bool -> Bool -> Bool
&& Integer
n Integer -> Integer -> Bool
forall a. Eq a => a -> a -> Bool
== Integer
0)
}