module Ecluse.Core.Queue.Buffer (
newEnqueueBuffer,
writeOrDrop,
reportWorthy,
) where
import Control.Concurrent.STM.TBQueue (TBQueue, isFullTBQueue, newTBQueueIO, readTBQueue, writeTBQueue)
import UnliftIO.Concurrent (threadDelay)
import UnliftIO.Exception (tryAny)
import Ecluse.Core.Fault (tfDetail)
import Ecluse.Core.Queue (MirrorJob, MirrorQueue (enqueue))
import Ecluse.Core.Supervision (BackoffSchedule (BackoffSchedule, bsBaseMicros, bsCapMicros), backoffMicros)
writeOrDrop :: TBQueue MirrorJob -> TVar Int -> MirrorJob -> STM (Maybe Int)
writeOrDrop :: TBQueue MirrorJob -> TVar Int -> MirrorJob -> STM (Maybe Int)
writeOrDrop TBQueue MirrorJob
queue TVar Int
dropCount MirrorJob
job = do
full <- TBQueue MirrorJob -> STM Bool
forall a. TBQueue a -> STM Bool
isFullTBQueue TBQueue MirrorJob
queue
if full
then Just <$> bumpCount dropCount
else writeTBQueue queue job $> Nothing
bumpCount :: TVar Int -> STM Int
bumpCount :: TVar Int -> STM Int
bumpCount TVar Int
counter = do
n <- (Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
1) (Int -> Int) -> STM Int -> STM Int
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TVar Int -> STM Int
forall a. TVar a -> STM a
readTVar TVar Int
counter
writeTVar counter n
pure n
reportWorthy :: Int -> Int -> Bool
reportWorthy :: Int -> Int -> Bool
reportWorthy Int
n Int
interval = Int
n Int -> Int -> Bool
forall a. Eq a => a -> a -> Bool
== Int
1 Bool -> Bool -> Bool
|| Int
n Int -> Int -> Int
forall a. Integral a => a -> a -> a
`mod` Int
interval Int -> Int -> Bool
forall a. Eq a => a -> a -> Bool
== Int
0
newEnqueueBuffer ::
Int ->
(Int -> IO ()) ->
(Int -> Text -> IO ()) ->
MirrorQueue ->
IO (MirrorQueue, IO ())
newEnqueueBuffer :: Int
-> (Int -> IO ())
-> (Int -> Text -> IO ())
-> MirrorQueue
-> IO (MirrorQueue, IO ())
newEnqueueBuffer Int
depth Int -> IO ()
onDrop Int -> Text -> IO ()
onDeliveryFailure MirrorQueue
backend = do
buffer <- Natural -> IO (TBQueue MirrorJob)
forall a. Natural -> IO (TBQueue a)
newTBQueueIO (Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
1 Int
depth))
dropCount <- newTVarIO (0 :: Int)
failureCount <- newTVarIO (0 :: Int)
let
handOff MirrorJob
job = do
dropped <- STM (Maybe Int) -> IO (Maybe Int)
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TBQueue MirrorJob -> TVar Int -> MirrorJob -> STM (Maybe Int)
writeOrDrop TBQueue MirrorJob
buffer TVar Int
dropCount MirrorJob
job)
whenJust dropped (void . tryAny . onDrop)
pure (Right ())
pure (backend{enqueue = handOff}, drainLoop buffer failureCount onDeliveryFailure backend)
drainLoop :: TBQueue MirrorJob -> TVar Int -> (Int -> Text -> IO ()) -> MirrorQueue -> IO ()
drainLoop :: TBQueue MirrorJob
-> TVar Int -> (Int -> Text -> IO ()) -> MirrorQueue -> IO ()
drainLoop TBQueue MirrorJob
buffer TVar Int
failureCount Int -> Text -> IO ()
onDeliveryFailure MirrorQueue
backend = Int -> IO ()
forall {b}. Int -> IO b
go Int
0
where
go :: Int -> IO b
go Int
consecutiveFailures = do
job <- STM MirrorJob -> IO MirrorJob
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TBQueue MirrorJob -> STM MirrorJob
forall a. TBQueue a -> STM a
readTBQueue TBQueue MirrorJob
buffer)
enqueue backend job >>= \case
Right () -> Int -> IO b
go Int
0
Left TransportFault
fault -> do
n <- STM Int -> IO Int
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TVar Int -> STM Int
bumpCount TVar Int
failureCount)
void (tryAny (onDeliveryFailure n (tfDetail fault)))
threadDelay (backoffMicros drainBackoff consecutiveFailures)
go (consecutiveFailures + 1)
drainBackoff :: BackoffSchedule
drainBackoff :: BackoffSchedule
drainBackoff = BackoffSchedule{bsBaseMicros :: Int
bsBaseMicros = Int
200_000, bsCapMicros :: Int
bsCapMicros = Int
30_000_000}