-- | An ordered hand-off for effects that must not block the thread producing
-- them.
module Hydra.Node.Outbox where

import Hydra.Prelude

import Control.Concurrent.Class.MonadSTM (
  modifyTVar',
  readTQueue,
  tryReadTQueue,
  writeTQueue,
  writeTVar,
 )
import Hydra.Network (StallReason (..))

-- | Bounds beyond which an 'Outbox' reports itself stalled.
data StallBounds = StallBounds
  { StallBounds -> DiffTime
noProgressFor :: DiffTime
  -- ^ How long the outbox may fail to complete anything while it holds work.
  -- Also bounds how long the backlog may take to drain at the rate the
  -- consumer has recently been completing work, which catches a consumer
  -- that is moving but too slowly for what it holds, sooner than waiting for
  -- it to stop outright.
  , StallBounds -> Natural
maxPending :: Natural
  -- ^ How much work it may hold, whatever the rate it drains at: the memory
  -- backstop. Note it does not itself bound the queue: 'submit' never blocks,
  -- so only the producer can stop.
  --
  -- Both backlog limbs also fire on a producer outrunning a consumer that is
  -- still delivering, so a caller reporting them must describe the backlog
  -- rather than blame the consumer for being unreachable.
  }

-- | A single-consumer hand-off for actions that must not block the thread
-- submitting them.
--
-- The hydra node processes all its inputs on one thread and used to run every
-- effect inline on it, so an effect that blocked stopped it dequeuing
-- anything at all - including the chain observation and the client command it
-- needs to close and contest a head. See GHSA-3mmr-q43p-g6p2 and
-- 'Hydra.Node.withNetworkOutbox'.
--
-- 'submit' therefore only appends; 'runOutbox' is the one thing that performs
-- the actions, in submission order.
data Outbox m = Outbox
  { forall (m :: * -> *). Outbox m -> m () -> m ()
submit :: m () -> m ()
  -- ^ Append an action. Never blocks.
  , forall (m :: * -> *). Outbox m -> m (Maybe (StallReason, Natural))
outboxStalled :: m (Maybe (StallReason, Natural))
  -- ^ The cause, and the amount of queued work, once it exceeds the
  -- 'StallBounds'.
  --
  -- Progress-based first, on purpose: depth alone is a normal condition for
  -- a consumer that batches or rate-limits, and "non-empty for a while"
  -- would misfire on a busy producer whose queue simply never happens to be
  -- observed empty. What actually means "this is not getting through" is
  -- that nothing has completed - reported as 'NoProgress' even when the
  -- backlog also happens to be at 'maxPending'.
  --
  -- The backlog itself is judged by how long it would take to drain at the
  -- recent completion rate, so a burst that a fast consumer clears within
  -- 'noProgressFor' is not a stall however deep it gets, short of
  -- 'maxPending'. Both are reported as 'BacklogFull'.
  , forall (m :: * -> *). Outbox m -> STM m Natural
pendingActions :: STM m Natural
  -- ^ Submitted but not yet completed, including any action in flight.
  , forall (m :: * -> *). Outbox m -> m (Natural, DiffTime)
outboxBacklog :: m (Natural, DiffTime)
  -- ^ The same count, plus how long since the outbox last completed
  -- anything - zero while it holds nothing. Where 'outboxStalled' answers
  -- "should we cut the flow off", this is for reporting: it is the pair an
  -- operator needs to tell a backlog that is draining from one that is not,
  -- and to see which of the two 'StallBounds' limbs is being approached.
  , forall (m :: * -> *). Outbox m -> m ()
runOutbox :: m ()
  -- ^ Perform submitted actions forever, in submission order. Must run
  -- concurrently with the producer, or nothing is ever performed.
  , forall (m :: * -> *). Outbox m -> m ()
drainOutbox :: m ()
  -- ^ Perform everything queued right now and return. For producers driven
  -- step by step instead of alongside 'runOutbox'. Must not run concurrently
  -- with it: two consumers would interleave the actions.
  , forall (m :: * -> *). Outbox m -> StallBounds
stallBounds :: StallBounds
  -- ^ What this outbox was created with, so a watcher can pick a polling
  -- interval that matches.
  }

newOutbox :: (MonadLabelledSTM m, MonadMonotonicTime m) => StallBounds -> String -> m (Outbox m)
newOutbox :: forall (m :: * -> *).
(MonadLabelledSTM m, MonadMonotonicTime m) =>
StallBounds -> String -> m (Outbox m)
newOutbox bounds :: StallBounds
bounds@StallBounds{DiffTime
$sel:noProgressFor:StallBounds :: StallBounds -> DiffTime
noProgressFor :: DiffTime
noProgressFor, Natural
$sel:maxPending:StallBounds :: StallBounds -> Natural
maxPending :: Natural
maxPending} String
name = do
  -- Monotonic, not wall clock: an NTP step of more than 'noProgressFor' would
  -- otherwise fake a stall on a healthy node, or mask a real one.
  Time
now <- m Time
forall (m :: * -> *). MonadMonotonicTime m => m Time
getMonotonicTime
  TQueue m (m ())
queue <- String -> m (TQueue m (m ()))
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> m (TQueue m a)
newLabelledTQueueIO String
name
  TVar m Natural
count <- String -> Natural -> m (TVar m Natural)
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO (String
name String -> String -> String
forall a. Semigroup a => a -> a -> a
<> String
"-pending") Natural
0
  -- The last point at which the outbox was known to be keeping up: either an
  -- action completed, or the queue was empty and we started waiting. Both
  -- reset the clock, so neither an idle producer nor a busy one that keeps up
  -- is ever reported as stalled.
  TVar m Time
progressAt <- String -> Time -> m (TVar m Time)
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO (String
name String -> String -> String
forall a. Semigroup a => a -> a -> a
<> String
"-progress-at") Time
now
  -- Busy time and completions of the last full 'serviceWindow', and of the
  -- one in progress, from which 'outboxStalled' takes the mean time to
  -- complete one action. Each sample runs from 'progressAt', so it is the gap
  -- between consecutive completions while busy, and submission to completion
  -- otherwise; idle time is never counted. A mean over a window rather than a
  -- moving average, so one slow completion among many fast ones moves it only
  -- by its share of the window: with a deep backlog, a moving average turns a
  -- single pause into a refusal. The window in progress counts too, so the
  -- completions right after a slow one dilute it straight away.
  --
  -- A sample counts for at most 'serviceWindow', so one long pause, which is
  -- 'NoProgress' territory anyway, weighs no more than a window of its own.
  -- Both windows are cleared whenever the outbox empties: windows only roll
  -- over on busy time, so a backlog that drains quickly after a pause would
  -- otherwise leave that pause standing through any idle period and trip the
  -- next small burst.
  TVar m (DiffTime, Natural)
lastWindow <- String -> (DiffTime, Natural) -> m (TVar m (DiffTime, Natural))
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO (String
name String -> String -> String
forall a. Semigroup a => a -> a -> a
<> String
"-last-window") (DiffTime
0, Natural
0 :: Natural)
  TVar m (DiffTime, Natural)
window <- String -> (DiffTime, Natural) -> m (TVar m (DiffTime, Natural))
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO (String
name String -> String -> String
forall a. Semigroup a => a -> a -> a
<> String
"-window") (DiffTime
0, Natural
0 :: Natural)
  let
    -- NOTE: the count drops, and the clock resets, only once the action has
    -- completed. An action stuck inside the effect itself therefore still
    -- counts as pending and holds the clock, which is what makes a stall
    -- visible.
    perform :: m () -> m ()
perform m ()
action = do
      m ()
action
      Time
completedAt <- m Time
forall (m :: * -> *). MonadMonotonicTime m => m Time
getMonotonicTime
      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 Natural -> (Natural -> Natural) -> 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 Natural
count Natural -> Natural
forall a. Enum a => a -> a
pred
        Bool
idle <- (Natural -> Natural -> Bool
forall a. Eq a => a -> a -> Bool
== Natural
0) (Natural -> Bool) -> STM m Natural -> STM m Bool
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TVar m Natural -> STM m Natural
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Natural
count
        DiffTime
sample <- DiffTime -> DiffTime -> DiffTime
forall a. Ord a => a -> a -> a
min DiffTime
serviceWindow (DiffTime -> DiffTime) -> (Time -> DiffTime) -> Time -> DiffTime
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Time -> Time -> DiffTime
diffTime Time
completedAt (Time -> DiffTime) -> STM m Time -> STM m DiffTime
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TVar m Time -> STM m Time
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Time
progressAt
        (DiffTime
busy, Natural
completed) <- (DiffTime -> DiffTime)
-> (Natural -> Natural)
-> (DiffTime, Natural)
-> (DiffTime, Natural)
forall a b c d. (a -> b) -> (c -> d) -> (a, c) -> (b, d)
forall (p :: * -> * -> *) a b c d.
Bifunctor p =>
(a -> b) -> (c -> d) -> p a c -> p b d
bimap (DiffTime -> DiffTime -> DiffTime
forall a. Num a => a -> a -> a
+ DiffTime
sample) (Natural -> Natural -> Natural
forall a. Num a => a -> a -> a
+ Natural
1) ((DiffTime, Natural) -> (DiffTime, Natural))
-> STM m (DiffTime, Natural) -> STM m (DiffTime, Natural)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TVar m (DiffTime, Natural) -> STM m (DiffTime, Natural)
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m (DiffTime, Natural)
window
        if
          | Bool
idle -> TVar m (DiffTime, Natural) -> (DiffTime, Natural) -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m (DiffTime, Natural)
lastWindow (DiffTime
0, Natural
0) STM m () -> STM m () -> STM m ()
forall a b. STM m a -> STM m b -> STM m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> TVar m (DiffTime, Natural) -> (DiffTime, Natural) -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m (DiffTime, Natural)
window (DiffTime
0, Natural
0)
          | DiffTime
busy DiffTime -> DiffTime -> Bool
forall a. Ord a => a -> a -> Bool
>= DiffTime
serviceWindow -> TVar m (DiffTime, Natural) -> (DiffTime, Natural) -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m (DiffTime, Natural)
lastWindow (DiffTime
busy, Natural
completed) STM m () -> STM m () -> STM m ()
forall a b. STM m a -> STM m b -> STM m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> TVar m (DiffTime, Natural) -> (DiffTime, Natural) -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m (DiffTime, Natural)
window (DiffTime
0, Natural
0)
          | Bool
otherwise -> TVar m (DiffTime, Natural) -> (DiffTime, Natural) -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m (DiffTime, Natural)
window (DiffTime
busy, Natural
completed)
        TVar m Time -> Time -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m Time
progressAt Time
completedAt
  Outbox m -> m (Outbox m)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure
    Outbox
      { submit :: m () -> m ()
submit = \m ()
action -> do
          Time
submittedAt <- m Time
forall (m :: * -> *). MonadMonotonicTime m => m Time
getMonotonicTime
          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
            Natural
n <- TVar m Natural -> STM m Natural
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Natural
count
            -- NOTE: the clock is read before this transaction, so a
            -- completion committing in between would otherwise be rolled back
            -- to an older timestamp and fake a stall. Only ever move forward.
            Bool -> STM m () -> STM m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Natural
n Natural -> Natural -> Bool
forall a. Eq a => a -> a -> Bool
== Natural
0) (STM m () -> STM m ()) -> STM m () -> STM m ()
forall a b. (a -> b) -> a -> b
$ TVar m Time -> (Time -> Time) -> 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 Time
progressAt (Time -> Time -> Time
forall a. Ord a => a -> a -> a
max Time
submittedAt)
            TVar m Natural -> Natural -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m Natural
count (Natural
n Natural -> Natural -> Natural
forall a. Num a => a -> a -> a
+ Natural
1)
            TQueue m (m ()) -> m () -> STM m ()
forall a. TQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TQueue m a -> a -> STM m ()
writeTQueue TQueue m (m ())
queue m ()
action
      , $sel:outboxStalled:Outbox :: m (Maybe (StallReason, Natural))
outboxStalled = do
          -- NOTE: our own clock, deliberately. Callers may be working off
          -- another notion of time, such as chain time.
          Time
asOf <- m Time
forall (m :: * -> *). MonadMonotonicTime m => m Time
getMonotonicTime
          STM m (Maybe (StallReason, Natural))
-> m (Maybe (StallReason, Natural))
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m (Maybe (StallReason, Natural))
 -> m (Maybe (StallReason, Natural)))
-> STM m (Maybe (StallReason, Natural))
-> m (Maybe (StallReason, Natural))
forall a b. (a -> b) -> a -> b
$ do
            Natural
n <- TVar m Natural -> STM m Natural
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Natural
count
            Time
since <- TVar m Time -> STM m Time
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Time
progressAt
            (DiffTime
lastBusy, Natural
lastCompleted) <- TVar m (DiffTime, Natural) -> STM m (DiffTime, Natural)
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m (DiffTime, Natural)
lastWindow
            (DiffTime
busy, Natural
completed) <- TVar m (DiffTime, Natural) -> STM m (DiffTime, Natural)
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m (DiffTime, Natural)
window
            -- Zero until a full window has elapsed, so the first few
            -- completions after startup cannot trip it on their own.
            let perAction :: DiffTime
perAction
                  | Natural
lastCompleted Natural -> Natural -> Bool
forall a. Eq a => a -> a -> Bool
== Natural
0 = DiffTime
0
                  | Bool
otherwise = (DiffTime
lastBusy DiffTime -> DiffTime -> DiffTime
forall a. Num a => a -> a -> a
+ DiffTime
busy) DiffTime -> DiffTime -> DiffTime
forall a. Fractional a => a -> a -> a
/ Natural -> DiffTime
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Natural
lastCompleted Natural -> Natural -> Natural
forall a. Num a => a -> a -> a
+ Natural
completed)
            Maybe (StallReason, Natural)
-> STM m (Maybe (StallReason, Natural))
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Maybe (StallReason, Natural)
 -> STM m (Maybe (StallReason, Natural)))
-> Maybe (StallReason, Natural)
-> STM m (Maybe (StallReason, Natural))
forall a b. (a -> b) -> a -> b
$
              if
                | Natural
n Natural -> Natural -> Bool
forall a. Eq a => a -> a -> Bool
== Natural
0 -> Maybe (StallReason, Natural)
forall a. Maybe a
Nothing
                | Time
asOf Time -> Time -> DiffTime
`diffTime` Time
since DiffTime -> DiffTime -> Bool
forall a. Ord a => a -> a -> Bool
> DiffTime
noProgressFor -> (StallReason, Natural) -> Maybe (StallReason, Natural)
forall a. a -> Maybe a
Just (StallReason
NoProgress, Natural
n)
                | Natural
n Natural -> Natural -> Bool
forall a. Ord a => a -> a -> Bool
>= Natural
maxPending -> (StallReason, Natural) -> Maybe (StallReason, Natural)
forall a. a -> Maybe a
Just (StallReason
BacklogFull, Natural
n)
                | Natural -> DiffTime
forall a b. (Integral a, Num b) => a -> b
fromIntegral Natural
n DiffTime -> DiffTime -> DiffTime
forall a. Num a => a -> a -> a
* DiffTime
perAction DiffTime -> DiffTime -> Bool
forall a. Ord a => a -> a -> Bool
> DiffTime
noProgressFor -> (StallReason, Natural) -> Maybe (StallReason, Natural)
forall a. a -> Maybe a
Just (StallReason
BacklogFull, Natural
n)
                | Bool
otherwise -> Maybe (StallReason, Natural)
forall a. Maybe a
Nothing
      , $sel:outboxBacklog:Outbox :: m (Natural, DiffTime)
outboxBacklog = do
          Time
asOf <- m Time
forall (m :: * -> *). MonadMonotonicTime m => m Time
getMonotonicTime
          STM m (Natural, DiffTime) -> m (Natural, DiffTime)
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m (Natural, DiffTime) -> m (Natural, DiffTime))
-> STM m (Natural, DiffTime) -> m (Natural, DiffTime)
forall a b. (a -> b) -> a -> b
$ do
            Natural
n <- TVar m Natural -> STM m Natural
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Natural
count
            Time
since <- TVar m Time -> STM m Time
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Time
progressAt
            (Natural, DiffTime) -> STM m (Natural, DiffTime)
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Natural
n, if Natural
n Natural -> Natural -> Bool
forall a. Eq a => a -> a -> Bool
== Natural
0 then DiffTime
0 else Time
asOf Time -> Time -> DiffTime
`diffTime` Time
since)
      , $sel:pendingActions:Outbox :: STM m Natural
pendingActions = TVar m Natural -> STM m Natural
forall a. TVar m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar m Natural
count
      , runOutbox :: m ()
runOutbox = m () -> m ()
forall (f :: * -> *) a b. Applicative f => f a -> f b
forever (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ STM m (m ()) -> m (m ())
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (TQueue m (m ()) -> STM m (m ())
forall a. TQueue m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TQueue m a -> STM m a
readTQueue TQueue m (m ())
queue) m (m ()) -> (m () -> m ()) -> m ()
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= m () -> m ()
perform
      , $sel:drainOutbox:Outbox :: m ()
drainOutbox =
          let go :: m ()
go =
                STM m (Maybe (m ())) -> m (Maybe (m ()))
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (TQueue m (m ()) -> STM m (Maybe (m ()))
forall a. TQueue m a -> STM m (Maybe a)
forall (m :: * -> *) a. MonadSTM m => TQueue m a -> STM m (Maybe a)
tryReadTQueue TQueue m (m ())
queue) m (Maybe (m ())) -> (Maybe (m ()) -> m ()) -> m ()
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
                  Maybe (m ())
Nothing -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
                  Just m ()
action -> m () -> m ()
perform m ()
action m () -> m () -> m ()
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> m ()
go
           in m ()
go
      , $sel:stallBounds:Outbox :: StallBounds
stallBounds = StallBounds
bounds
      }

-- | How much busy time the service time estimate averages over. Long enough
-- that a pause of a few hundred milliseconds barely moves it, short against
-- any 'noProgressFor' worth configuring.
serviceWindow :: DiffTime
serviceWindow :: DiffTime
serviceWindow = DiffTime
1