-- | 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 (isSuffixOf, 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, renameFile)
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) =>
  Tracer IO PersistentQueueLog ->
  FilePath ->
  Natural ->
  m (PersistentQueue m a)
newPersistentQueue :: forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m) =>
Tracer IO PersistentQueueLog
-> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO PersistentQueueLog
tracer String
path Natural
capacity = do
  ([Natural]
paths, [Natural]
quarantinedIxs) <- IO ([Natural], [Natural]) -> m ([Natural], [Natural])
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO ([Natural], [Natural]) -> m ([Natural], [Natural]))
-> IO ([Natural], [Natural]) -> m ([Natural], [Natural])
forall a b. (a -> b) -> a -> b
$ do
    Bool -> String -> IO ()
createDirectoryIfMissing Bool
True String
path
    [String]
entries <- String -> IO [String]
listDirectory String
path
    ([Natural], [Natural]) -> IO ([Natural], [Natural])
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([Natural] -> [Natural]
forall a. Ord a => [a] -> [a]
sort ([Natural] -> [Natural]) -> [Natural] -> [Natural]
forall a b. (a -> b) -> a -> b
$ (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]
entries, (String -> Maybe Natural) -> [String] -> [Natural]
forall a b. (a -> Maybe b) -> [a] -> [b]
mapMaybe String -> Maybe Natural
quarantinedIndex [String]
entries)
  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
  -- Load item by item: a single unreadable or torn file (nothing fsyncs these,
  -- so a crash can leave one) must cost only that item, not the whole pending
  -- queue.
  [Natural] -> (Natural -> m ()) -> m ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Natural]
paths ((Natural -> m ()) -> m ()) -> (Natural -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \(Natural
idx :: Natural) ->
    m ByteString -> m (Either IOException ByteString)
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 (String -> m ByteString
forall (m :: * -> *). MonadIO m => String -> m ByteString
readFileBS (String
path String -> ShowS
</> Natural -> String
forall a. Show a => a -> String
show Natural
idx)) m (Either IOException ByteString)
-> (Either IOException ByteString -> 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
      -- Reading can fail for reasons that are none of the item's fault and may
      -- not recur, so leave the file be: this run does not deliver it, a later
      -- one may. Two consequences worth knowing: an item recovered on a later
      -- start goes out after messages enqueued after it, so FIFO holds within
      -- a run rather than across restarts; and one that never becomes readable
      -- stays in the directory, as do quarantined items.
      Left (IOException
e :: IOException) -> Natural -> String -> m ()
forall {m :: * -> *} {a}.
(MonadIO m, Show a) =>
a -> String -> m ()
trace Natural
idx (IOException -> String
forall a. Show a => a -> String
show IOException
e)
      Right ByteString
bs -> case ByteString -> Either DecoderError a
forall a. FromCBOR a => ByteString -> Either DecoderError a
decodeFull' ByteString
bs of
        -- Set aside rather than deleted: undecodable is indistinguishable from
        -- "the encoding changed under us", so an upgrade must not be able to
        -- discard the pending queue. The suffix takes it out of the active
        -- set, since the listing above only accepts names that read as an
        -- index.
        Left DecoderError
err -> Natural -> String -> m ()
forall {m :: * -> *} {a}.
(MonadIO m, Show a) =>
a -> String -> m ()
trace Natural
idx (DecoderError -> String
forall a. Show a => a -> String
show DecoderError
err) m () -> m () -> m ()
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Natural -> m ()
forall {m :: * -> *} {a}. (MonadIO m, Show a) => a -> m ()
quarantine Natural
idx
        Right a
item -> 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
idx, a
item, ByteString
bs)
  -- Derived from what is on disk, not from what loaded: reusing an index
  -- would overwrite a file that is still pending. Quarantined indices count
  -- too, even though they are no longer live: 'quarantine' renames onto
  -- '<idx>.undecodable' unconditionally, so reissuing an already-quarantined
  -- index would let a later torn file silently clobber an earlier one's
  -- forensic copy.
  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] -> Natural
forall a. Ord a => [a] -> a
forall (t :: * -> *) a. (Foldable t, Ord a) => t a -> a
maximum (Natural
0 Natural -> [Natural] -> [Natural]
forall a. a -> [a] -> [a]
: [Natural]
paths [Natural] -> [Natural] -> [Natural]
forall a. Semigroup a => a -> a -> a
<> [Natural]
quarantinedIxs) 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
  trace :: a -> String -> m ()
trace a
idx String
reason =
    IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ())
-> (PersistentQueueLog -> IO ()) -> PersistentQueueLog -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Tracer IO PersistentQueueLog -> PersistentQueueLog -> IO ()
forall (m :: * -> *) a. Tracer m a -> a -> m ()
traceWith Tracer IO PersistentQueueLog
tracer (PersistentQueueLog -> m ()) -> PersistentQueueLog -> m ()
forall a b. (a -> b) -> a -> b
$
      PersistentQueueLoadFailed{$sel:reason:PersistentQueueLoadFailed :: Text
reason = String -> Text
Text.pack (String
"item " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> a -> String
forall a. Show a => a -> String
show a
idx String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
": " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
reason)}

  quarantinedIndex :: FilePath -> Maybe Natural
  quarantinedIndex :: String -> Maybe Natural
quarantinedIndex String
name
    | String
suffix String -> String -> Bool
forall a. Eq a => [a] -> [a] -> Bool
`isSuffixOf` String
name = String -> Maybe Natural
forall a. Read a => String -> Maybe a
readMaybe (Int -> ShowS
forall a. Int -> [a] -> [a]
take (String -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length String
name Int -> Int -> Int
forall a. Num a => a -> a -> a
- String -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length String
suffix) String
name)
    | Bool
otherwise = Maybe Natural
forall a. Maybe a
Nothing
   where
    suffix :: String
suffix = String
".undecodable"

  quarantine :: a -> m ()
quarantine a
idx =
    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 -> String -> IO ()
renameFile (String
path String -> ShowS
</> a -> String
forall a. Show a => a -> String
show a
idx) (String
path String -> ShowS
</> a -> String
forall a. Show a => a -> String
show a
idx String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
".undecodable")
        IO () -> (IOException -> IO ()) -> IO ()
forall e a. Exception e => IO a -> (e -> IO a) -> IO a
`catch` \(IOException
e :: IOException) ->
          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 (String
"could not set aside item " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> a -> String
forall a. Show a => a -> String
show a
idx String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
": " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> IOException -> String
forall a. Show a => a -> String
show IOException
e)}

-- | 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)}