-- | The general input queue from which the Hydra head is fed with inputs.
module Hydra.Node.InputQueue where

import Hydra.Prelude

import Control.Concurrent.Class.MonadSTM (
  isEmptyTBQueue,
  modifyTVar',
  readTBQueue,
  swapTVar,
  writeTBQueue,
 )

-- | The single, required queue in the system from which a hydra head is "fed".
-- NOTE(SN): handle pattern, but likely not required as there is no need for an
-- alternative implementation
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 ()
  -- ^ Hold an item aside without re-enqueuing it, until 'releaseParked' moves
  -- it back onto the queue. Used for inputs that can't be processed until the
  -- node is in sync. Parked items do not count towards 'isEmpty'.
  , forall (m :: * -> *) e. InputQueue m e -> m ()
releaseParked :: m ()
  -- ^ Move all parked items back onto the queue, preserving their order.
  , 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" []
  -- XXX: We bound the _input_ queue by the _logging_ queue size! This is a
  -- hack; but we do it because it seems that the logging queue blocking
  -- prevents further processing, _unless_ the input queue is also bounded.
  -- In truth it probably makes sense for this queue to be bounded anyway.
  -- See: <https://github.com/cardano-scaling/hydra/issues/2442>
  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 []
          -- Write the items back from a separate thread (as 'reenqueue' does),
          -- and one at a time, so a full queue can't deadlock the single
          -- consumer that calls this. 'numThreads' keeps them visible to
          -- 'isEmpty' until they are all back on the queue.
          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)
      }