module Ecluse.Core.Queue.Memory (
MemoryQueueConfig (..),
defaultMemoryQueueConfig,
newBoundedInMemoryQueue,
memoryQueueDropReportInterval,
) where
import Control.Concurrent.STM.TBQueue (TBQueue, newTBQueueIO, readTBQueue, tryReadTBQueue)
import System.Timeout (timeout)
import Ecluse.Core.Queue (
DeadLetterTerminus (TerminusAbsent),
MirrorJob,
MirrorQueue (..),
QueueMessage (..),
defaultDeliveryBudget,
mkReceiptHandle,
)
import Ecluse.Core.Queue.Buffer (
reportWorthy,
writeOrDrop,
)
data MemoryQueueConfig = MemoryQueueConfig
{ MemoryQueueConfig -> Int
memQueueMaxDepth :: Int
, MemoryQueueConfig -> Int
memQueuePollWaitMicros :: Int
}
deriving stock (MemoryQueueConfig -> MemoryQueueConfig -> Bool
(MemoryQueueConfig -> MemoryQueueConfig -> Bool)
-> (MemoryQueueConfig -> MemoryQueueConfig -> Bool)
-> Eq MemoryQueueConfig
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
== :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
$c/= :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
/= :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
Eq, Int -> MemoryQueueConfig -> ShowS
[MemoryQueueConfig] -> ShowS
MemoryQueueConfig -> String
(Int -> MemoryQueueConfig -> ShowS)
-> (MemoryQueueConfig -> String)
-> ([MemoryQueueConfig] -> ShowS)
-> Show MemoryQueueConfig
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> MemoryQueueConfig -> ShowS
showsPrec :: Int -> MemoryQueueConfig -> ShowS
$cshow :: MemoryQueueConfig -> String
show :: MemoryQueueConfig -> String
$cshowList :: [MemoryQueueConfig] -> ShowS
showList :: [MemoryQueueConfig] -> ShowS
Show)
defaultMemoryQueueConfig :: Int -> MemoryQueueConfig
defaultMemoryQueueConfig :: Int -> MemoryQueueConfig
defaultMemoryQueueConfig Int
maxDepth =
MemoryQueueConfig
{ memQueueMaxDepth :: Int
memQueueMaxDepth = Int
maxDepth
, memQueuePollWaitMicros :: Int
memQueuePollWaitMicros = Int
20_000_000
}
memoryQueueBatchSize :: Int
memoryQueueBatchSize :: Int
memoryQueueBatchSize = Int
10
memoryQueueDropReportInterval :: Int
memoryQueueDropReportInterval :: Int
memoryQueueDropReportInterval = Int
1000
newBoundedInMemoryQueue ::
MemoryQueueConfig ->
(Int -> IO ()) ->
IO MirrorQueue
newBoundedInMemoryQueue :: MemoryQueueConfig -> (Int -> IO ()) -> IO MirrorQueue
newBoundedInMemoryQueue MemoryQueueConfig
cfg Int -> IO ()
onDrop = do
queue <- 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 (MemoryQueueConfig -> Int
memQueueMaxDepth MemoryQueueConfig
cfg)))
dropCount <- newTVarIO (0 :: Int)
nextReceipt <- newTVarIO (0 :: Word64)
pure
MirrorQueue
{ enqueue = \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
queue TVar Int
dropCount MirrorJob
job)
whenJust dropped (\Int
n -> Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int -> Int -> Bool
reportWorthy Int
n Int
memoryQueueDropReportInterval) (Int -> IO ()
onDrop Int
n))
pure (Right ())
,
receive = Right . fromMaybe [] <$> timeout (memQueuePollWaitMicros cfg) (atomically (receiveBatch queue nextReceipt))
,
ack = const (pure (Right ()))
, extendVisibility = \ReceiptHandle
_ Seconds
_ -> Either TransportFault () -> IO (Either TransportFault ())
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (() -> Either TransportFault ()
forall a b. b -> Either a b
Right ())
,
deadLetter = const (pure (Right ()))
,
deliveryBudget = defaultDeliveryBudget
, deadLetterTerminus = Right TerminusAbsent
}
receiveBatch :: TBQueue MirrorJob -> TVar Word64 -> STM [QueueMessage]
receiveBatch :: TBQueue MirrorJob -> TVar Word64 -> STM [QueueMessage]
receiveBatch TBQueue MirrorJob
queue TVar Word64
nextReceipt = do
headJob <- TBQueue MirrorJob -> STM MirrorJob
forall a. TBQueue a -> STM a
readTBQueue TBQueue MirrorJob
queue
rest <- drainUpTo (memoryQueueBatchSize - 1)
traverse assignReceipt (headJob : rest)
where
drainUpTo :: Int -> STM [MirrorJob]
drainUpTo :: Int -> STM [MirrorJob]
drainUpTo Int
budget
| Int
budget Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
<= Int
0 = [MirrorJob] -> STM [MirrorJob]
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure []
| Bool
otherwise =
TBQueue MirrorJob -> STM (Maybe MirrorJob)
forall a. TBQueue a -> STM (Maybe a)
tryReadTBQueue TBQueue MirrorJob
queue STM (Maybe MirrorJob)
-> (Maybe MirrorJob -> STM [MirrorJob]) -> STM [MirrorJob]
forall a b. STM a -> (a -> STM b) -> STM b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
Maybe MirrorJob
Nothing -> [MirrorJob] -> STM [MirrorJob]
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure []
Just MirrorJob
job -> (MirrorJob
job MirrorJob -> [MirrorJob] -> [MirrorJob]
forall a. a -> [a] -> [a]
:) ([MirrorJob] -> [MirrorJob]) -> STM [MirrorJob] -> STM [MirrorJob]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Int -> STM [MirrorJob]
drainUpTo (Int
budget Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1)
assignReceipt :: MirrorJob -> STM QueueMessage
assignReceipt :: MirrorJob -> STM QueueMessage
assignReceipt MirrorJob
job = do
n <- TVar Word64 -> STM Word64
forall a. TVar a -> STM a
readTVar TVar Word64
nextReceipt
writeTVar nextReceipt (n + 1)
pure QueueMessage{msgJob = job, msgReceipt = mkReceiptHandle (show n), msgReceiveCount = 1, msgLease = Nothing}