module Ecluse.Core.Server.Admission.Weighted (
WeightedAdmission,
newWeightedAdmission,
withWeightedAdmission,
AdmissionObservers (..),
admissionWaitMicros,
) where
import Control.Concurrent.STM (retry)
import GHC.Conc (registerDelay)
import UnliftIO (MonadUnliftIO)
import UnliftIO.Exception qualified as UE
data WeightedAdmission = WeightedAdmission
{ WeightedAdmission -> TVar Int
waAvailable :: TVar Int
, WeightedAdmission -> TVar Int
waWaiting :: TVar Int
, WeightedAdmission -> Int
waWaitingRoom :: Int
, WeightedAdmission -> Int
waWaitMicros :: Int
}
data AdmissionObservers = AdmissionObservers
{ AdmissionObservers -> IO ()
onQueued :: IO ()
, AdmissionObservers -> IO ()
onShed :: IO ()
, AdmissionObservers -> Int -> IO ()
onInFlightDelta :: Int -> IO ()
}
admissionWaitMicros :: Int
admissionWaitMicros :: Int
admissionWaitMicros = Int
1_000_000
newWeightedAdmission :: Int -> Int -> Int -> IO WeightedAdmission
newWeightedAdmission :: Int -> Int -> Int -> IO WeightedAdmission
newWeightedAdmission Int
capacity Int
room Int
waitMicros = do
available <- Int -> IO (TVar Int)
forall (m :: * -> *) a. MonadIO m => a -> m (TVar a)
newTVarIO Int
capacity
waiting <- newTVarIO 0
pure
WeightedAdmission
{ waAvailable = available
, waWaiting = waiting
, waWaitingRoom = max 0 room
, waWaitMicros = max 0 waitMicros
}
data Gate = Admitted | Queued | Refused
doorDecision :: WeightedAdmission -> Int -> STM Gate
doorDecision :: WeightedAdmission -> Int -> STM Gate
doorDecision WeightedAdmission
wa Int
weight = do
available <- TVar Int -> STM Int
forall a. TVar a -> STM a
readTVar (WeightedAdmission -> TVar Int
waAvailable WeightedAdmission
wa)
waiting <- readTVar (waWaiting wa)
if available >= weight && waiting == 0
then writeTVar (waAvailable wa) (available - weight) $> Admitted
else
if waiting >= waWaitingRoom wa
then pure Refused
else writeTVar (waWaiting wa) (waiting + 1) $> Queued
acquireOrExpire :: WeightedAdmission -> Int -> TVar Bool -> STM Bool
acquireOrExpire :: WeightedAdmission -> Int -> TVar Bool -> STM Bool
acquireOrExpire WeightedAdmission
wa Int
weight TVar Bool
deadline = do
available <- TVar Int -> STM Int
forall a. TVar a -> STM a
readTVar (WeightedAdmission -> TVar Int
waAvailable WeightedAdmission
wa)
if available >= weight
then writeTVar (waAvailable wa) (available - weight) $> True
else do
expired <- readTVar deadline
if expired then pure False else retry
{-# INLINE withWeightedAdmission #-}
withWeightedAdmission ::
(MonadUnliftIO m) =>
AdmissionObservers ->
WeightedAdmission ->
Int ->
m a ->
m (Maybe a)
withWeightedAdmission :: forall (m :: * -> *) a.
MonadUnliftIO m =>
AdmissionObservers
-> WeightedAdmission -> Int -> m a -> m (Maybe a)
withWeightedAdmission AdmissionObservers
obs WeightedAdmission
wa Int
weight m a
action =
((forall a. m a -> m a) -> m (Maybe a)) -> m (Maybe a)
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> m a) -> m b) -> m b
UE.mask (((forall a. m a -> m a) -> m (Maybe a)) -> m (Maybe a))
-> ((forall a. m a -> m a) -> m (Maybe a)) -> m (Maybe a)
forall a b. (a -> b) -> a -> b
$ \forall a. m a -> m a
restore -> do
gate <- STM Gate -> m Gate
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (WeightedAdmission -> Int -> STM Gate
doorDecision WeightedAdmission
wa Int
weight)
case gate of
Gate
Refused -> AdmissionObservers -> m (Maybe a)
forall (m :: * -> *) a.
MonadIO m =>
AdmissionObservers -> m (Maybe a)
shedRecording AdmissionObservers
obs
Gate
Admitted -> AdmissionObservers
-> WeightedAdmission
-> Int
-> IO ()
-> (m a -> m a)
-> m a
-> m (Maybe a)
forall (m :: * -> *) a.
MonadUnliftIO m =>
AdmissionObservers
-> WeightedAdmission
-> Int
-> IO ()
-> (m a -> m a)
-> m a
-> m (Maybe a)
admittedRun AdmissionObservers
obs WeightedAdmission
wa Int
weight (() -> IO ()
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()) m a -> m a
forall a. m a -> m a
restore m a
action
Gate
Queued -> AdmissionObservers
-> WeightedAdmission -> Int -> (m a -> m a) -> m a -> m (Maybe a)
forall (m :: * -> *) a.
MonadUnliftIO m =>
AdmissionObservers
-> WeightedAdmission -> Int -> (m a -> m a) -> m a -> m (Maybe a)
queuedWait AdmissionObservers
obs WeightedAdmission
wa Int
weight m a -> m a
forall a. m a -> m a
restore m a
action
{-# INLINE shedRecording #-}
shedRecording :: (MonadIO m) => AdmissionObservers -> m (Maybe a)
shedRecording :: forall (m :: * -> *) a.
MonadIO m =>
AdmissionObservers -> m (Maybe a)
shedRecording AdmissionObservers
obs = IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (AdmissionObservers -> IO ()
onShed AdmissionObservers
obs) m () -> Maybe a -> m (Maybe a)
forall (f :: * -> *) a b. Functor f => f a -> b -> f b
$> Maybe a
forall a. Maybe a
Nothing
{-# INLINE queuedWait #-}
queuedWait ::
(MonadUnliftIO m) =>
AdmissionObservers ->
WeightedAdmission ->
Int ->
(m a -> m a) ->
m a ->
m (Maybe a)
queuedWait :: forall (m :: * -> *) a.
MonadUnliftIO m =>
AdmissionObservers
-> WeightedAdmission -> Int -> (m a -> m a) -> m a -> m (Maybe a)
queuedWait AdmissionObservers
obs WeightedAdmission
wa Int
weight m a -> m a
restore m a
action = do
deadline <- IO (TVar Bool) -> m (TVar Bool)
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Int -> IO (TVar Bool)
registerDelay (WeightedAdmission -> Int
waWaitMicros WeightedAdmission
wa))
acquired <-
atomically (acquireOrExpire wa weight deadline)
`UE.finally` atomically (modifyTVar' (waWaiting wa) (subtract 1))
if acquired
then admittedRun obs wa weight (onQueued obs) restore action
else shedRecording obs
{-# INLINE admittedRun #-}
admittedRun ::
(MonadUnliftIO m) =>
AdmissionObservers ->
WeightedAdmission ->
Int ->
IO () ->
(m a -> m a) ->
m a ->
m (Maybe a)
admittedRun :: forall (m :: * -> *) a.
MonadUnliftIO m =>
AdmissionObservers
-> WeightedAdmission
-> Int
-> IO ()
-> (m a -> m a)
-> m a
-> m (Maybe a)
admittedRun AdmissionObservers
obs WeightedAdmission
wa Int
weight IO ()
afterArm m a -> m a
restore m a
action =
a -> Maybe a
forall a. a -> Maybe a
Just
(a -> Maybe a) -> m a -> m (Maybe a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ( (IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (AdmissionObservers -> Int -> IO ()
onInFlightDelta AdmissionObservers
obs Int
weight IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> IO ()
afterArm) m () -> m a -> m a
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> m a -> m a
restore m a
action)
m a -> m () -> m a
forall (m :: * -> *) a b. MonadUnliftIO m => m a -> m b -> m a
`UE.finally` AdmissionObservers -> WeightedAdmission -> Int -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
AdmissionObservers -> WeightedAdmission -> Int -> m ()
releaseWeight AdmissionObservers
obs WeightedAdmission
wa Int
weight
)
{-# INLINE releaseWeight #-}
releaseWeight :: (MonadUnliftIO m) => AdmissionObservers -> WeightedAdmission -> Int -> m ()
releaseWeight :: forall (m :: * -> *).
MonadUnliftIO m =>
AdmissionObservers -> WeightedAdmission -> Int -> m ()
releaseWeight AdmissionObservers
obs WeightedAdmission
wa Int
weight =
IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (AdmissionObservers -> Int -> IO ()
onInFlightDelta AdmissionObservers
obs (Int -> Int
forall a. Num a => a -> a
negate Int
weight))
m () -> m () -> m ()
forall (m :: * -> *) a b. MonadUnliftIO m => m a -> m b -> m a
`UE.finally` STM () -> m ()
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TVar Int -> (Int -> Int) -> STM ()
forall a. TVar a -> (a -> a) -> STM ()
modifyTVar' (WeightedAdmission -> TVar Int
waAvailable WeightedAdmission
wa) (Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
weight))