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