-- | A crash-safe, file-backed, bounded FIFO queue.
--
-- Each item is written to its own file (named by a monotonic index) before it
-- enters the in-memory 'TBQueue', and removed from disk only after it has been
-- popped, so a restart reloads exactly the items that were enqueued but not yet
-- consumed (at-least-once delivery). Items carry their CBOR serialization,
-- produced once on write and reused, and can be drained in byte- and
-- count-bounded batches with an in-flight pin for safe retries.
module Control.Concurrent.PersistentQueue (
  PersistentQueue,
  PersistentQueueLog (..),
  newPersistentQueue,
  writePersistentQueue,
  peekPersistentQueue,
  tryPeekPersistentQueue,
  peekBatchPersistentQueue,
  popBatchPersistentQueue,
  nextPendingBatch,
) where

import Cardano.Binary (FromCBOR, ToCBOR, decodeFull', serialize')
import Control.Concurrent.Class.Labelled (newLabelledTBQueueIO, newLabelledTVarIO)
import Control.Concurrent.Class.MonadSTM (
  MonadLabelledSTM,
  MonadSTM,
  TBQueue,
  TVar,
  atomically,
  isFullTBQueue,
  modifyTVar',
  peekTBQueue,
  readTBQueue,
  readTVar,
  readTVarIO,
  tryPeekTBQueue,
  tryReadTBQueue,
  unGetTBQueue,
  writeTBQueue,
  writeTVar,
 )
import Control.Exception (IOException, catch)
import Control.Monad (forM, forM_, unless, when)
import Control.Monad.Class.MonadThrow (MonadCatch, try)
import Control.Monad.IO.Class (MonadIO, liftIO)
import Control.Tracer (Tracer, traceWith)
import Data.ByteString (ByteString)
import Data.ByteString qualified as BS
import Data.List (sort)
import Data.Maybe (mapMaybe)
import Data.Text (Text)
import Data.Text qualified as Text
import Numeric.Natural (Natural)
import System.Directory (createDirectoryIfMissing, listDirectory, removeFile)
import System.FilePath ((</>))
import System.IO.Error (isDoesNotExistError)
import Text.Read (readMaybe)

-- | Events emitted while operating the queue.
data PersistentQueueLog
  = PersistentQueueLoadFailed {PersistentQueueLog -> Text
reason :: Text}
  | PersistentQueueFull
  | PersistentQueueDeleteFailed {PersistentQueueLog -> Natural
index :: Natural, reason :: Text}
  deriving stock (PersistentQueueLog -> PersistentQueueLog -> Bool
(PersistentQueueLog -> PersistentQueueLog -> Bool)
-> (PersistentQueueLog -> PersistentQueueLog -> Bool)
-> Eq PersistentQueueLog
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: PersistentQueueLog -> PersistentQueueLog -> Bool
== :: PersistentQueueLog -> PersistentQueueLog -> Bool
$c/= :: PersistentQueueLog -> PersistentQueueLog -> Bool
/= :: PersistentQueueLog -> PersistentQueueLog -> Bool
Eq, Int -> PersistentQueueLog -> ShowS
[PersistentQueueLog] -> ShowS
PersistentQueueLog -> String
(Int -> PersistentQueueLog -> ShowS)
-> (PersistentQueueLog -> String)
-> ([PersistentQueueLog] -> ShowS)
-> Show PersistentQueueLog
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> PersistentQueueLog -> ShowS
showsPrec :: Int -> PersistentQueueLog -> ShowS
$cshow :: PersistentQueueLog -> String
show :: PersistentQueueLog -> String
$cshowList :: [PersistentQueueLog] -> ShowS
showList :: [PersistentQueueLog] -> ShowS
Show)

-- | Queue elements carry the item's CBOR serialization, produced once on
-- write and reused for both the on-disk file and downstream consumers, so
-- consuming does not serialize twice.
data PersistentQueue m a = PersistentQueue
  { forall (m :: * -> *) a.
PersistentQueue m a -> TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
  , forall (m :: * -> *) a. PersistentQueue m a -> TVar m Natural
nextIx :: TVar m Natural
  , forall (m :: * -> *) a. PersistentQueue m a -> String
directory :: FilePath
  }

readFileBS :: MonadIO m => FilePath -> m ByteString
readFileBS :: forall (m :: * -> *). MonadIO m => String -> m ByteString
readFileBS = IO ByteString -> m ByteString
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO ByteString -> m ByteString)
-> (String -> IO ByteString) -> String -> m ByteString
forall b c a. (b -> c) -> (a -> b) -> a -> c
. String -> IO ByteString
BS.readFile

writeFileBS :: MonadIO m => FilePath -> ByteString -> m ()
writeFileBS :: forall (m :: * -> *). MonadIO m => String -> ByteString -> m ()
writeFileBS String
path = IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> (ByteString -> IO ()) -> ByteString -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. String -> ByteString -> IO ()
BS.writeFile String
path

-- | Create a new persistent queue at file path and given capacity.
newPersistentQueue ::
  (MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m, MonadFail m) =>
  Tracer IO PersistentQueueLog ->
  FilePath ->
  Natural ->
  m (PersistentQueue m a)
newPersistentQueue :: forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO PersistentQueueLog
-> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO PersistentQueueLog
tracer String
path Natural
capacity = do
  [Natural]
paths <- IO [Natural] -> m [Natural]
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO [Natural] -> m [Natural]) -> IO [Natural] -> m [Natural]
forall a b. (a -> b) -> a -> b
$ do
    Bool -> String -> IO ()
createDirectoryIfMissing Bool
True String
path
    [Natural] -> [Natural]
forall a. Ord a => [a] -> [a]
sort ([Natural] -> [Natural])
-> ([String] -> [Natural]) -> [String] -> [Natural]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (String -> Maybe Natural) -> [String] -> [Natural]
forall a b. (a -> Maybe b) -> [a] -> [b]
mapMaybe String -> Maybe Natural
forall a. Read a => String -> Maybe a
readMaybe ([String] -> [Natural]) -> IO [String] -> IO [Natural]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> String -> IO [String]
listDirectory String
path
  TBQueue m (Natural, a, ByteString)
queue <- String -> Natural -> m (TBQueue m (Natural, a, ByteString))
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> Natural -> m (TBQueue m a)
newLabelledTBQueueIO String
"persistent-queue" (Natural -> m (TBQueue m (Natural, a, ByteString)))
-> Natural -> m (TBQueue m (Natural, a, ByteString))
forall a b. (a -> b) -> a -> b
$ Natural -> Natural -> Natural
forall a. Ord a => a -> a -> a
max (Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int -> Natural) -> Int -> Natural
forall a b. (a -> b) -> a -> b
$ [Natural] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [Natural]
paths) Natural
capacity
  Natural
highestId <-
    m Natural -> m (Either IOException Natural)
forall e a. Exception e => m a -> m (Either e a)
forall (m :: * -> *) e a.
(MonadCatch m, Exception e) =>
m a -> m (Either e a)
try (TBQueue m (Natural, a, ByteString) -> [Natural] -> m Natural
forall {f :: * -> *} {b}.
(MonadIO f, FromCBOR b, MonadFail f, MonadSTM f) =>
TBQueue f (Natural, b, ByteString) -> [Natural] -> f Natural
loadExisting TBQueue m (Natural, a, ByteString)
queue [Natural]
paths) m (Either IOException Natural)
-> (Either IOException Natural -> m Natural) -> m Natural
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
      Left (IOException
e :: IOException) -> do
        IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> IO () -> m ()
forall a b. (a -> b) -> a -> b
$ do
          Tracer IO PersistentQueueLog -> PersistentQueueLog -> IO ()
forall (m :: * -> *) a. Tracer m a -> a -> m ()
traceWith Tracer IO PersistentQueueLog
tracer PersistentQueueLoadFailed{$sel:reason:PersistentQueueLoadFailed :: Text
reason = String -> Text
Text.pack (IOException -> String
forall a. Show a => a -> String
show IOException
e)}
          Bool -> String -> IO ()
createDirectoryIfMissing Bool
True String
path
        Natural -> m Natural
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Natural
0
      Right Natural
highest -> Natural -> m Natural
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Natural
highest
  TVar m Natural
nextIx <- String -> Natural -> m (TVar m Natural)
forall (m :: * -> *) a.
MonadLabelledSTM m =>
String -> a -> m (TVar m a)
newLabelledTVarIO String
"persistent-next-ix" (Natural -> m (TVar m Natural)) -> Natural -> m (TVar m Natural)
forall a b. (a -> b) -> a -> b
$ Natural
highestId Natural -> Natural -> Natural
forall a. Num a => a -> a -> a
+ Natural
1
  PersistentQueue m a -> m (PersistentQueue m a)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure PersistentQueue{TBQueue m (Natural, a, ByteString)
$sel:queue:PersistentQueue :: TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
queue, TVar m Natural
$sel:nextIx:PersistentQueue :: TVar m Natural
nextIx :: TVar m Natural
nextIx, $sel:directory:PersistentQueue :: String
directory = String
path}
 where
  loadExisting :: TBQueue f (Natural, b, ByteString) -> [Natural] -> f Natural
loadExisting TBQueue f (Natural, b, ByteString)
queue = \case
    [] -> Natural -> f Natural
forall a. a -> f a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Natural
0
    [Natural]
idxs -> do
      [Natural] -> (Natural -> f ()) -> f ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Natural]
idxs ((Natural -> f ()) -> f ()) -> (Natural -> f ()) -> f ()
forall a b. (a -> b) -> a -> b
$ \(Natural
idx :: Natural) -> do
        ByteString
bs <- String -> f ByteString
forall (m :: * -> *). MonadIO m => String -> m ByteString
readFileBS (String
path String -> ShowS
</> Natural -> String
forall a. Show a => a -> String
show Natural
idx)
        case ByteString -> Either DecoderError b
forall a. FromCBOR a => ByteString -> Either DecoderError a
decodeFull' ByteString
bs of
          Left DecoderError
err ->
            String -> f ()
forall a. String -> f a
forall (m :: * -> *) a. MonadFail m => String -> m a
fail (String -> f ()) -> String -> f ()
forall a b. (a -> b) -> a -> b
$ String
"Failed to decode item: " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> DecoderError -> String
forall a. Show a => a -> String
show DecoderError
err
          Right b
item ->
            STM f () -> f ()
forall a. HasCallStack => STM f a -> f a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM f () -> f ()) -> STM f () -> f ()
forall a b. (a -> b) -> a -> b
$ TBQueue f (Natural, b, ByteString)
-> (Natural, b, ByteString) -> STM f ()
forall a. TBQueue f a -> a -> STM f ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
writeTBQueue TBQueue f (Natural, b, ByteString)
queue (Natural
idx, b
item, ByteString
bs)
      Natural -> f Natural
forall a. a -> f a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Natural -> f Natural) -> Natural -> f Natural
forall a b. (a -> b) -> a -> b
$ [Natural] -> Natural
forall a. HasCallStack => [a] -> a
last [Natural]
idxs

-- | Write a value to the queue, blocking if the queue is full.
writePersistentQueue :: (ToCBOR a, MonadSTM m, MonadIO m) => Tracer IO PersistentQueueLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue :: forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO PersistentQueueLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO PersistentQueueLog
tracer PersistentQueue{TBQueue m (Natural, a, ByteString)
$sel:queue:PersistentQueue :: forall (m :: * -> *) a.
PersistentQueue m a -> TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
queue, TVar m Natural
$sel:nextIx:PersistentQueue :: forall (m :: * -> *) a. PersistentQueue m a -> TVar m Natural
nextIx :: TVar m Natural
nextIx, String
$sel:directory:PersistentQueue :: forall (m :: * -> *) a. PersistentQueue m a -> String
directory :: String
directory} a
item = do
  Natural
next <- STM m Natural -> m Natural
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m Natural -> m Natural) -> STM m Natural -> m Natural
forall a b. (a -> b) -> a -> b
$ do
    Natural
next <- 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
nextIx
    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
nextIx (Natural -> Natural -> Natural
forall a. Num a => a -> a -> a
+ Natural
1)
    Natural -> STM m Natural
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Natural
next
  let !bytes :: ByteString
bytes = a -> ByteString
forall a. ToCBOR a => a -> ByteString
serialize' a
item
  String -> ByteString -> m ()
forall (m :: * -> *). MonadIO m => String -> ByteString -> m ()
writeFileBS (String
directory String -> ShowS
</> Natural -> String
forall a. Show a => a -> String
show Natural
next) ByteString
bytes
  Bool
full <- STM m Bool -> m Bool
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m Bool -> m Bool) -> STM m Bool -> m Bool
forall a b. (a -> b) -> a -> b
$ TBQueue m (Natural, a, ByteString) -> STM m Bool
forall a. TBQueue m a -> STM m Bool
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> STM m Bool
isFullTBQueue TBQueue m (Natural, a, ByteString)
queue
  Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when Bool
full (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> IO () -> m ()
forall a b. (a -> b) -> a -> b
$ Tracer IO PersistentQueueLog -> PersistentQueueLog -> IO ()
forall (m :: * -> *) a. Tracer m a -> a -> m ()
traceWith Tracer IO PersistentQueueLog
tracer PersistentQueueLog
PersistentQueueFull
  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
$ TBQueue m (Natural, a, ByteString)
-> (Natural, a, ByteString) -> STM m ()
forall a. TBQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
writeTBQueue TBQueue m (Natural, a, ByteString)
queue (Natural
next, a
item, ByteString
bytes)

-- | Get the next value from the queue without removing it, blocking if the
-- queue is empty.
peekPersistentQueue :: MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue :: forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue{TBQueue m (Natural, a, ByteString)
$sel:queue:PersistentQueue :: forall (m :: * -> *) a.
PersistentQueue m a -> TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
queue} = do
  (\(Natural
_, a
item, ByteString
_) -> a
item) ((Natural, a, ByteString) -> a)
-> m (Natural, a, ByteString) -> m a
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> STM m (Natural, a, ByteString) -> m (Natural, a, ByteString)
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (TBQueue m (Natural, a, ByteString)
-> STM m (Natural, a, ByteString)
forall a. TBQueue m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> STM m a
peekTBQueue TBQueue m (Natural, a, ByteString)
queue)

-- | Like 'peekPersistentQueue', but returns 'Nothing' instead of blocking
-- when the queue is empty.
tryPeekPersistentQueue :: MonadSTM m => PersistentQueue m a -> m (Maybe a)
tryPeekPersistentQueue :: forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> m (Maybe a)
tryPeekPersistentQueue PersistentQueue{TBQueue m (Natural, a, ByteString)
$sel:queue:PersistentQueue :: forall (m :: * -> *) a.
PersistentQueue m a -> TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
queue} = do
  ((Natural, a, ByteString) -> a)
-> Maybe (Natural, a, ByteString) -> Maybe a
forall a b. (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (\(Natural
_, a
item, ByteString
_) -> a
item) (Maybe (Natural, a, ByteString) -> Maybe a)
-> m (Maybe (Natural, a, ByteString)) -> m (Maybe a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> STM m (Maybe (Natural, a, ByteString))
-> m (Maybe (Natural, a, ByteString))
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (TBQueue m (Natural, a, ByteString)
-> STM m (Maybe (Natural, a, ByteString))
forall a. TBQueue m a -> STM m (Maybe a)
forall (m :: * -> *) a.
MonadSTM m =>
TBQueue m a -> STM m (Maybe a)
tryPeekTBQueue TBQueue m (Natural, a, ByteString)
queue)

-- | Get all pending values and their serializations, up to the given count
-- and total byte limits, blocking until at least one is available. Values
-- are not removed; use 'popBatchPersistentQueue' after they were sent.
peekBatchPersistentQueue :: MonadSTM m => PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
peekBatchPersistentQueue :: forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
peekBatchPersistentQueue PersistentQueue{TBQueue m (Natural, a, ByteString)
$sel:queue:PersistentQueue :: forall (m :: * -> *) a.
PersistentQueue m a -> TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
queue} Int
maxCount Int
maxBytes = STM m [(a, ByteString)] -> m [(a, ByteString)]
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m [(a, ByteString)] -> m [(a, ByteString)])
-> STM m [(a, ByteString)] -> m [(a, ByteString)]
forall a b. (a -> b) -> a -> b
$ do
  (Natural, a, ByteString)
first' <- TBQueue m (Natural, a, ByteString)
-> STM m (Natural, a, ByteString)
forall a. TBQueue m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> STM m a
readTBQueue TBQueue m (Natural, a, ByteString)
queue
  -- Collected in reverse consumption order
  [(Natural, a, ByteString)]
rest <- Int
-> [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
forall {m :: * -> *}.
(TBQueue m ~ TBQueue m, MonadSTM m) =>
Int
-> [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
go ((Natural, a, ByteString) -> Int
forall {a} {b}. (a, b, ByteString) -> Int
remainingAfter (Natural, a, ByteString)
first') []
  -- Restore everything we consumed: 'unGetTBQueue' pushes to the front, so
  -- restoring newest-first re-establishes the original queue order.
  [(Natural, a, ByteString)]
-> ((Natural, a, ByteString) -> STM m ()) -> STM m ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ ([(Natural, a, ByteString)]
rest [(Natural, a, ByteString)]
-> [(Natural, a, ByteString)] -> [(Natural, a, ByteString)]
forall a. Semigroup a => a -> a -> a
<> [(Natural, a, ByteString)
first']) (((Natural, a, ByteString) -> STM m ()) -> STM m ())
-> ((Natural, a, ByteString) -> STM m ()) -> STM m ()
forall a b. (a -> b) -> a -> b
$ TBQueue m (Natural, a, ByteString)
-> (Natural, a, ByteString) -> STM m ()
forall a. TBQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
unGetTBQueue TBQueue m (Natural, a, ByteString)
queue
  [(a, ByteString)] -> STM m [(a, ByteString)]
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([(a, ByteString)] -> STM m [(a, ByteString)])
-> [(a, ByteString)] -> STM m [(a, ByteString)]
forall a b. (a -> b) -> a -> b
$ (\(Natural
_, a
item, ByteString
bytes) -> (a
item, ByteString
bytes)) ((Natural, a, ByteString) -> (a, ByteString))
-> [(Natural, a, ByteString)] -> [(a, ByteString)]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ((Natural, a, ByteString)
first' (Natural, a, ByteString)
-> [(Natural, a, ByteString)] -> [(Natural, a, ByteString)]
forall a. a -> [a] -> [a]
: [(Natural, a, ByteString)] -> [(Natural, a, ByteString)]
forall a. [a] -> [a]
reverse [(Natural, a, ByteString)]
rest)
 where
  go :: Int
-> [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
go Int
budget [(Natural, a, ByteString)]
acc
    | [(Natural, a, ByteString)] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [(Natural, a, ByteString)]
acc Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
maxCount Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1 = [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure [(Natural, a, ByteString)]
acc
    | Bool
otherwise =
        TBQueue m (Natural, a, ByteString)
-> STM m (Maybe (Natural, a, ByteString))
forall a. TBQueue m a -> STM m (Maybe a)
forall (m :: * -> *) a.
MonadSTM m =>
TBQueue m a -> STM m (Maybe a)
tryReadTBQueue TBQueue m (Natural, a, ByteString)
TBQueue m (Natural, a, ByteString)
queue STM m (Maybe (Natural, a, ByteString))
-> (Maybe (Natural, a, ByteString)
    -> STM m [(Natural, a, ByteString)])
-> STM m [(Natural, a, ByteString)]
forall a b. STM m a -> (a -> STM m b) -> STM m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
          Maybe (Natural, a, ByteString)
Nothing -> [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure [(Natural, a, ByteString)]
acc
          Just next :: (Natural, a, ByteString)
next@(Natural
_, a
_, ByteString
bytes)
            | ByteString -> Int
BS.length ByteString
bytes Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
budget -> do
                TBQueue m (Natural, a, ByteString)
-> (Natural, a, ByteString) -> STM m ()
forall a. TBQueue m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> a -> STM m ()
unGetTBQueue TBQueue m (Natural, a, ByteString)
TBQueue m (Natural, a, ByteString)
queue (Natural, a, ByteString)
next
                [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
forall a. a -> STM m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure [(Natural, a, ByteString)]
acc
            | Bool
otherwise -> Int
-> [(Natural, a, ByteString)] -> STM m [(Natural, a, ByteString)]
go (Int
budget Int -> Int -> Int
forall a. Num a => a -> a -> a
- ByteString -> Int
BS.length ByteString
bytes) ((Natural, a, ByteString)
next (Natural, a, ByteString)
-> [(Natural, a, ByteString)] -> [(Natural, a, ByteString)]
forall a. a -> [a] -> [a]
: [(Natural, a, ByteString)]
acc)

  remainingAfter :: (a, b, ByteString) -> Int
remainingAfter (a
_, b
_, ByteString
bytes) = Int
maxBytes Int -> Int -> Int
forall a. Num a => a -> a -> a
- ByteString -> Int
BS.length ByteString
bytes

-- | Get the batch to broadcast next: the batch already in flight if there is
-- one, otherwise a fresh one peeked from the queue (and recorded as in
-- flight). Returns 'Nothing' when nothing is pending. The caller must clear
-- the in-flight var after popping a successfully sent batch.
--
-- Pinning the in-flight batch across transient retries matters for consumers
-- with an at-least-once send that can commit while the caller sees a
-- transient failure: the retry must send (and afterwards pop) exactly the
-- content of the committed attempt. Re-peeking on retry could pick up
-- messages enqueued in the meantime, and a dedup that declared the grown
-- batch delivered would pop and lose the never-sent tail.
nextPendingBatch ::
  MonadSTM m =>
  TVar m (Maybe [(a, ByteString)]) ->
  PersistentQueue m a ->
  Int ->
  Int ->
  m (Maybe [(a, ByteString)])
nextPendingBatch :: forall (m :: * -> *) a.
MonadSTM m =>
TVar m (Maybe [(a, ByteString)])
-> PersistentQueue m a -> Int -> Int -> m (Maybe [(a, ByteString)])
nextPendingBatch TVar m (Maybe [(a, ByteString)])
inFlightVar PersistentQueue m a
queue Int
maxCount Int
maxBytes =
  TVar m (Maybe [(a, ByteString)]) -> m (Maybe [(a, ByteString)])
forall a. TVar m a -> m a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> m a
readTVarIO TVar m (Maybe [(a, ByteString)])
inFlightVar m (Maybe [(a, ByteString)])
-> (Maybe [(a, ByteString)] -> m (Maybe [(a, ByteString)]))
-> m (Maybe [(a, ByteString)])
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
    Just [(a, ByteString)]
batch -> Maybe [(a, ByteString)] -> m (Maybe [(a, ByteString)])
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([(a, ByteString)] -> Maybe [(a, ByteString)]
forall a. a -> Maybe a
Just [(a, ByteString)]
batch)
    Maybe [(a, ByteString)]
Nothing ->
      PersistentQueue m a -> m (Maybe a)
forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> m (Maybe a)
tryPeekPersistentQueue PersistentQueue m a
queue m (Maybe a)
-> (Maybe a -> m (Maybe [(a, ByteString)]))
-> m (Maybe [(a, ByteString)])
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 a
Nothing -> Maybe [(a, ByteString)] -> m (Maybe [(a, ByteString)])
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe [(a, ByteString)]
forall a. Maybe a
Nothing
        Just a
_ -> do
          [(a, ByteString)]
batch <- PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
peekBatchPersistentQueue PersistentQueue m a
queue Int
maxCount Int
maxBytes
          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
$ TVar m (Maybe [(a, ByteString)])
-> Maybe [(a, ByteString)] -> STM m ()
forall a. TVar m a -> a -> STM m ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar m (Maybe [(a, ByteString)])
inFlightVar ([(a, ByteString)] -> Maybe [(a, ByteString)]
forall a. a -> Maybe a
Just [(a, ByteString)]
batch)
          Maybe [(a, ByteString)] -> m (Maybe [(a, ByteString)])
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([(a, ByteString)] -> Maybe [(a, ByteString)]
forall a. a -> Maybe a
Just [(a, ByteString)]
batch)

-- | Remove a batch previously returned by 'peekBatchPersistentQueue'. Pops
-- unconditionally, one item per batch entry: the caller is expected to be the
-- sole consumer, so the queue head still holds exactly the peeked items (an
-- item-matching guard could only silently no-op and wedge the queue).
popBatchPersistentQueue :: (MonadSTM m, MonadIO m) => Tracer IO PersistentQueueLog -> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue :: forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO PersistentQueueLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO PersistentQueueLog
tracer PersistentQueue{TBQueue m (Natural, a, ByteString)
$sel:queue:PersistentQueue :: forall (m :: * -> *) a.
PersistentQueue m a -> TBQueue m (Natural, a, ByteString)
queue :: TBQueue m (Natural, a, ByteString)
queue, String
$sel:directory:PersistentQueue :: forall (m :: * -> *) a. PersistentQueue m a -> String
directory :: String
directory} [(a, ByteString)]
batch = do
  [Natural]
indices <- STM m [Natural] -> m [Natural]
forall a. HasCallStack => STM m a -> m a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM m [Natural] -> m [Natural]) -> STM m [Natural] -> m [Natural]
forall a b. (a -> b) -> a -> b
$ [(a, ByteString)]
-> ((a, ByteString) -> STM m Natural) -> STM m [Natural]
forall (t :: * -> *) (m :: * -> *) a b.
(Traversable t, Monad m) =>
t a -> (a -> m b) -> m (t b)
forM [(a, ByteString)]
batch (((a, ByteString) -> STM m Natural) -> STM m [Natural])
-> ((a, ByteString) -> STM m Natural) -> STM m [Natural]
forall a b. (a -> b) -> a -> b
$ \(a, ByteString)
_ -> (\(Natural
ix, a
_, ByteString
_) -> Natural
ix) ((Natural, a, ByteString) -> Natural)
-> STM m (Natural, a, ByteString) -> STM m Natural
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TBQueue m (Natural, a, ByteString)
-> STM m (Natural, a, ByteString)
forall a. TBQueue m a -> STM m a
forall (m :: * -> *) a. MonadSTM m => TBQueue m a -> STM m a
readTBQueue TBQueue m (Natural, a, ByteString)
queue
  [Natural] -> (Natural -> m ()) -> m ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Natural]
indices ((Natural -> m ()) -> m ()) -> (Natural -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ Tracer IO PersistentQueueLog -> String -> Natural -> m ()
forall (m :: * -> *).
MonadIO m =>
Tracer IO PersistentQueueLog -> String -> Natural -> m ()
removeQueueFile Tracer IO PersistentQueueLog
tracer String
directory

-- | Delete the backing file of a popped queue item. Failing to delete is
-- traced but not fatal: the item was already consumed, so a leftover file
-- only means it may be re-delivered after a restart (at-least-once delivery,
-- same as the crash-recovery path).
removeQueueFile :: MonadIO m => Tracer IO PersistentQueueLog -> FilePath -> Natural -> m ()
removeQueueFile :: forall (m :: * -> *).
MonadIO m =>
Tracer IO PersistentQueueLog -> String -> Natural -> m ()
removeQueueFile Tracer IO PersistentQueueLog
tracer String
directory Natural
ix =
  IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> IO () -> m ()
forall a b. (a -> b) -> a -> b
$
    String -> IO ()
removeFile (String
directory String -> ShowS
</> Natural -> String
forall a. Show a => a -> String
show Natural
ix) IO () -> (IOException -> IO ()) -> IO ()
forall e a. Exception e => IO a -> (e -> IO a) -> IO a
`catch` \IOException
e ->
      Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless (IOException -> Bool
isDoesNotExistError IOException
e) (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$
        Tracer IO PersistentQueueLog -> PersistentQueueLog -> IO ()
forall (m :: * -> *) a. Tracer m a -> a -> m ()
traceWith Tracer IO PersistentQueueLog
tracer PersistentQueueDeleteFailed{$sel:index:PersistentQueueLoadFailed :: Natural
index = Natural
ix, $sel:reason:PersistentQueueLoadFailed :: Text
reason = String -> Text
Text.pack (IOException -> String
forall a. Show a => a -> String
show IOException
e)}