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)
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)
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
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
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)
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)
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)
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
[(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') []
[(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
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)
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
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)}