module Hydra.Node.Outbox where
import Hydra.Prelude
import Control.Concurrent.Class.MonadSTM (
modifyTVar',
readTQueue,
tryReadTQueue,
writeTQueue,
writeTVar,
)
import Hydra.Network (StallReason (..))
data StallBounds = StallBounds
{ StallBounds -> DiffTime
noProgressFor :: DiffTime
, StallBounds -> Natural
maxPending :: Natural
}
data Outbox m = Outbox
{ forall (m :: * -> *). Outbox m -> m () -> m ()
submit :: m () -> m ()
, forall (m :: * -> *). Outbox m -> m (Maybe (StallReason, Natural))
outboxStalled :: m (Maybe (StallReason, Natural))
, forall (m :: * -> *). Outbox m -> STM m Natural
pendingActions :: STM m Natural
, forall (m :: * -> *). Outbox m -> m (Natural, DiffTime)
outboxBacklog :: m (Natural, DiffTime)
, forall (m :: * -> *). Outbox m -> m ()
runOutbox :: m ()
, forall (m :: * -> *). Outbox m -> m ()
drainOutbox :: m ()
, forall (m :: * -> *). Outbox m -> StallBounds
stallBounds :: StallBounds
}
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
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
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
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
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
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
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
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
}
serviceWindow :: DiffTime
serviceWindow :: DiffTime
serviceWindow = DiffTime
1