-- | Tests of the 'PersistentQueue'.
module Hydra.PersistentQueueSpec where

import Hydra.Prelude
import Test.Hydra.Prelude

import Control.Concurrent.Class.MonadSTM (check, newTVarIO, readTVarIO, writeTVar)
import Control.Monad.Class.MonadAsync (concurrently, wait, withAsync)
import Hydra.Logging (Envelope (message), nullTracer, traceInTVar)
import Hydra.Network.Etcd (EtcdLog (..), newPersistentQueue, nextPendingBatch, peekBatchPersistentQueue, peekPersistentQueue, popBatchPersistentQueue, writePersistentQueue)
import System.Directory (createDirectory, listDirectory, removeFile)
import System.FilePath ((</>))
import Test.QuickCheck (counterexample, generate, ioProperty)
import Test.QuickCheck.Instances.Natural ()

spec :: Spec
spec :: Spec
spec = do
  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"can be constructed" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    Natural
capacity <- Gen Natural -> IO Natural
forall a. Gen a -> IO a
generate Gen Natural
forall a. Arbitrary a => Gen a
arbitrary
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      IO (PersistentQueue IO Int) -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (PersistentQueue IO Int) -> IO ())
-> IO (PersistentQueue IO Int) -> IO ()
forall a b. (a -> b) -> a -> b
$ forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue @_ @Int Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
capacity

  String -> ([Int] -> Property) -> Spec
forall prop.
(HasCallStack, Testable prop) =>
String -> prop -> Spec
prop String
"is persistent with capacity" (([Int] -> Property) -> Spec) -> ([Int] -> Property) -> Spec
forall a b. (a -> b) -> a -> b
$ \([Int]
items :: [Int]) -> do
    let capacity :: Natural
capacity = Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int -> Natural) -> Int -> Natural
forall a b. (a -> b) -> a -> b
$ [Int] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [Int]
items
    String -> Property -> Property
forall prop. Testable prop => String -> prop -> Property
counterexample (String
"capacity: " String -> String -> String
forall a. Semigroup a => a -> a -> a
<> Natural -> String
forall b a. (Show a, IsString b) => a -> b
show Natural
capacity) (Property -> Property) -> Property -> Property
forall a b. (a -> b) -> a -> b
$
      IO () -> Property
forall prop. Testable prop => IO prop -> Property
ioProperty (IO () -> Property) -> IO () -> Property
forall a b. (a -> b) -> a -> b
$
        String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
          PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
capacity
          IO [()] -> IO ()
forall a. HasCallStack => IO a -> IO ()
shouldNotBlock_ (IO [()] -> IO ()) -> IO [()] -> IO ()
forall a b. (a -> b) -> a -> b
$ (Int -> IO ()) -> [Int] -> IO [()]
forall (t :: * -> *) (m :: * -> *) a b.
(Traversable t, Monad m) =>
(a -> m b) -> t a -> m (t b)
forall (m :: * -> *) a b. Monad m => (a -> m b) -> [a] -> m [b]
mapM (Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q) [Int]
items
          -- This is expected to block as we reached capacity
          Maybe ()
_ <- DiffTime -> IO () -> IO (Maybe ())
forall a. DiffTime -> IO a -> IO (Maybe a)
forall (m :: * -> *) a.
MonadTimer m =>
DiffTime -> m a -> m (Maybe a)
timeout DiffTime
0.01 (Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q Int
123)
          -- A new queue should be initialized with all the elements
          PersistentQueue IO Int
q2 <- IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int)
forall a. HasCallStack => IO a -> IO a
shouldNotBlock (IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int))
-> IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int)
forall a b. (a -> b) -> a -> b
$ forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue @_ @Int Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
capacity
          let expected :: Int
expected = Int -> (NonEmpty Int -> Int) -> Maybe (NonEmpty Int) -> Int
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Int
123 NonEmpty Int -> Int
forall (f :: * -> *) a. IsNonEmpty f a a "head" => f a -> a
head ([Int] -> Maybe (NonEmpty Int)
forall a. [a] -> Maybe (NonEmpty a)
nonEmpty [Int]
items)
          PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q2 IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
expected

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"pop unblocks a blocked writer when at capacity" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      TVar [Envelope EtcdLog]
traces <- [Envelope EtcdLog] -> IO (TVar IO [Envelope EtcdLog])
forall a. a -> IO (TVar IO a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TVar m a)
newTVarIO []
      let tracer :: Tracer IO EtcdLog
tracer = TVar IO [Envelope EtcdLog] -> Text -> Tracer IO EtcdLog
forall (m :: * -> *) msg.
(MonadFork m, MonadTime m, MonadSTM m) =>
TVar m [Envelope msg] -> Text -> Tracer m msg
traceInTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces Text
"PersistentQueueSpec"
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
tracer String
dir Natural
1
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q (Int
1 :: Int)
      IO () -> (Async IO () -> IO ()) -> IO ()
forall a b. IO a -> (Async IO a -> IO b) -> IO b
forall (m :: * -> *) a b.
MonadAsync m =>
m a -> (Async m a -> m b) -> m b
withAsync (Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q Int
2) ((Async IO () -> IO ()) -> IO ())
-> (Async IO () -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \Async IO ()
writer -> do
        -- The writer traces 'PersistentQueueFull' right before blocking on
        -- the full queue, so this doubles as the signal it reached that point.
        IO () -> IO ()
forall a. HasCallStack => IO a -> IO ()
shouldNotBlock_ (IO () -> IO ()) -> (STM () -> IO ()) -> STM () -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. STM () -> IO ()
STM IO () -> IO ()
forall a. HasCallStack => STM IO a -> IO a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          [Envelope EtcdLog]
entries <- TVar IO [Envelope EtcdLog] -> STM IO [Envelope EtcdLog]
forall a. TVar IO a -> STM IO a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces
          Bool -> STM IO ()
forall (m :: * -> *). MonadSTM m => Bool -> STM m ()
check (Bool -> STM IO ()) -> Bool -> STM IO ()
forall a b. (a -> b) -> a -> b
$ EtcdLog
PersistentQueueFull EtcdLog -> [EtcdLog] -> Bool
forall (f :: * -> *) a.
(Foldable f, DisallowElem f, Eq a) =>
a -> f a -> Bool
`elem` (Envelope EtcdLog -> EtcdLog
forall a. Envelope a -> a
message (Envelope EtcdLog -> EtcdLog) -> [Envelope EtcdLog] -> [EtcdLog]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [Envelope EtcdLog]
entries)
        Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q [(Int
1, ByteString
"")]
        IO () -> IO ()
forall a. HasCallStack => IO a -> IO ()
shouldNotBlock_ (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ Async IO () -> IO ()
forall a. Async IO a -> IO a
forall (m :: * -> *) a. MonadAsync m => Async m a -> m a
wait Async IO ()
writer
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
2

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"does not deadlock when writing beyond capacity with a concurrent consumer" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      let items :: [Int]
items = [Int
1 .. Int
100 :: Int]
      (()
_, [Int]
received) <-
        IO ((), [Int]) -> IO ((), [Int])
forall a. HasCallStack => IO a -> IO a
shouldNotBlock (IO ((), [Int]) -> IO ((), [Int]))
-> IO ((), [Int]) -> IO ((), [Int])
forall a b. (a -> b) -> a -> b
$
          IO () -> IO [Int] -> IO ((), [Int])
forall a b. IO a -> IO b -> IO (a, b)
forall (m :: * -> *) a b. MonadAsync m => m a -> m b -> m (a, b)
concurrently
            ([Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int]
items ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q)
            ( [Int] -> (Int -> IO Int) -> IO [Int]
forall (t :: * -> *) (m :: * -> *) a b.
(Traversable t, Monad m) =>
t a -> (a -> m b) -> m (t b)
forM [Int]
items ((Int -> IO Int) -> IO [Int]) -> (Int -> IO Int) -> IO [Int]
forall a b. (a -> b) -> a -> b
$ \Int
_ -> do
                Int
item <- PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q
                Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q [(Int
item, ByteString
"")]
                Int -> IO Int
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Int
item
            )
      [Int]
received [Int] -> [Int] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [Int]
items
      String -> IO [String]
listDirectory String
dir IO [String] -> [String] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` []

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"pop removes the head item from memory and disk" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      [Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1, Int
2, Int
3 :: Int] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
1
      Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q [(Int
1, ByteString
"")]
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
2
      -- A popped item must not be reloaded on restart (it was already
      -- broadcast and would be duplicated)
      PersistentQueue IO Int
q2 <- IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int)
forall a. HasCallStack => IO a -> IO a
shouldNotBlock (IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int))
-> IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int)
forall a b. (a -> b) -> a -> b
$ forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue @_ @Int Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q2 IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
2

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"pop tolerates a missing backing file" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q (Int
1 :: Int)
      [String]
files <- String -> IO [String]
listDirectory String
dir
      [String] -> (String -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [String]
files ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
f -> String -> IO ()
removeFile (String
dir String -> String -> String
</> String
f)
      Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q [(Int
1, ByteString
"")]

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"pop survives and traces a failing backing file deletion" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      TVar [Envelope EtcdLog]
traces <- [Envelope EtcdLog] -> IO (TVar IO [Envelope EtcdLog])
forall a. a -> IO (TVar IO a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TVar m a)
newTVarIO []
      let tracer :: Tracer IO EtcdLog
tracer = TVar IO [Envelope EtcdLog] -> Text -> Tracer IO EtcdLog
forall (m :: * -> *) msg.
(MonadFork m, MonadTime m, MonadSTM m) =>
TVar m [Envelope msg] -> Text -> Tracer m msg
traceInTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces Text
"PersistentQueueSpec"
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
tracer String
dir Natural
10
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q (Int
1 :: Int)
      -- Replace the backing file with a same-named directory so removeFile
      -- fails with EISDIR instead of ENOENT (works even when running as root)
      [String
file] <- String -> IO [String]
listDirectory String
dir
      String -> IO ()
removeFile (String
dir String -> String -> String
</> String
file)
      String -> IO ()
createDirectory (String
dir String -> String -> String
</> String
file)
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q Int
2
      Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q [(Int
1, ByteString
"")]
      [Envelope EtcdLog]
entries <- TVar IO [Envelope EtcdLog] -> IO [Envelope EtcdLog]
forall a. TVar IO a -> IO a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> m a
readTVarIO TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces
      (Envelope EtcdLog -> EtcdLog
forall a. Envelope a -> a
message (Envelope EtcdLog -> EtcdLog) -> [Envelope EtcdLog] -> [EtcdLog]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [Envelope EtcdLog]
entries) [EtcdLog] -> ([EtcdLog] -> Bool) -> IO ()
forall a. (HasCallStack, Show a) => a -> (a -> Bool) -> IO ()
`shouldSatisfy` (EtcdLog -> Bool) -> [EtcdLog] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
any EtcdLog -> Bool
isDeleteFailed
      -- The consumer keeps making progress
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
2

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"traces PersistentQueueLoadFailed on corrupt items and starts empty" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      String -> ByteString -> IO ()
forall (m :: * -> *). MonadIO m => String -> ByteString -> m ()
writeFileBS (String
dir String -> String -> String
</> String
"1") ByteString
"not-valid-cbor"
      TVar [Envelope EtcdLog]
traces <- [Envelope EtcdLog] -> IO (TVar IO [Envelope EtcdLog])
forall a. a -> IO (TVar IO a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TVar m a)
newTVarIO []
      let tracer :: Tracer IO EtcdLog
tracer = TVar IO [Envelope EtcdLog] -> Text -> Tracer IO EtcdLog
forall (m :: * -> *) msg.
(MonadFork m, MonadTime m, MonadSTM m) =>
TVar m [Envelope msg] -> Text -> Tracer m msg
traceInTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces Text
"PersistentQueueSpec"
      PersistentQueue IO Int
q <- forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue @_ @Int Tracer IO EtcdLog
tracer String
dir Natural
10
      [Envelope EtcdLog]
entries <- TVar IO [Envelope EtcdLog] -> IO [Envelope EtcdLog]
forall a. TVar IO a -> IO a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> m a
readTVarIO TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces
      (Envelope EtcdLog -> EtcdLog
forall a. Envelope a -> a
message (Envelope EtcdLog -> EtcdLog) -> [Envelope EtcdLog] -> [EtcdLog]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [Envelope EtcdLog]
entries) [EtcdLog] -> ([EtcdLog] -> Bool) -> IO ()
forall a. (HasCallStack, Show a) => a -> (a -> Bool) -> IO ()
`shouldSatisfy` (EtcdLog -> Bool) -> [EtcdLog] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
any EtcdLog -> Bool
isLoadFailed
      -- The queue starts empty but stays functional
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q Int
42
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
42

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"peekBatch returns the FIFO prefix and popBatch removes exactly it" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      [Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1 .. Int
5 :: Int] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q
      [(Int, ByteString)]
batch <- PersistentQueue IO Int -> Int -> Int -> IO [(Int, ByteString)]
forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
peekBatchPersistentQueue PersistentQueue IO Int
q Int
3 Int
forall a. Bounded a => a
maxBound
      ((Int, ByteString) -> Int
forall a b. (a, b) -> a
fst ((Int, ByteString) -> Int) -> [(Int, ByteString)] -> [Int]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [(Int, ByteString)]
batch) [Int] -> [Int] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [Int
1, Int
2, Int
3]
      -- Peeking does not consume
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
1
      Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q [(Int, ByteString)]
batch
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
4
      -- Popped items must not be reloaded on restart
      PersistentQueue IO Int
q2 <- IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int)
forall a. HasCallStack => IO a -> IO a
shouldNotBlock (IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int))
-> IO (PersistentQueue IO Int) -> IO (PersistentQueue IO Int)
forall a b. (a -> b) -> a -> b
$ forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue @_ @Int Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q2 IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
4

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"popBatch unblocks a blocked writer when at capacity" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      TVar [Envelope EtcdLog]
traces <- [Envelope EtcdLog] -> IO (TVar IO [Envelope EtcdLog])
forall a. a -> IO (TVar IO a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TVar m a)
newTVarIO []
      let tracer :: Tracer IO EtcdLog
tracer = TVar IO [Envelope EtcdLog] -> Text -> Tracer IO EtcdLog
forall (m :: * -> *) msg.
(MonadFork m, MonadTime m, MonadSTM m) =>
TVar m [Envelope msg] -> Text -> Tracer m msg
traceInTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces Text
"PersistentQueueSpec"
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
tracer String
dir Natural
2
      [Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1, Int
2 :: Int] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q
      IO () -> (Async IO () -> IO ()) -> IO ()
forall a b. IO a -> (Async IO a -> IO b) -> IO b
forall (m :: * -> *) a b.
MonadAsync m =>
m a -> (Async m a -> m b) -> m b
withAsync (Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q Int
3) ((Async IO () -> IO ()) -> IO ())
-> (Async IO () -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \Async IO ()
writer -> do
        IO () -> IO ()
forall a. HasCallStack => IO a -> IO ()
shouldNotBlock_ (IO () -> IO ()) -> (STM () -> IO ()) -> STM () -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. STM () -> IO ()
STM IO () -> IO ()
forall a. HasCallStack => STM IO a -> IO a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          [Envelope EtcdLog]
entries <- TVar IO [Envelope EtcdLog] -> STM IO [Envelope EtcdLog]
forall a. TVar IO a -> STM IO a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> STM m a
readTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces
          Bool -> STM IO ()
forall (m :: * -> *). MonadSTM m => Bool -> STM m ()
check (Bool -> STM IO ()) -> Bool -> STM IO ()
forall a b. (a -> b) -> a -> b
$ EtcdLog
PersistentQueueFull EtcdLog -> [EtcdLog] -> Bool
forall (f :: * -> *) a.
(Foldable f, DisallowElem f, Eq a) =>
a -> f a -> Bool
`elem` (Envelope EtcdLog -> EtcdLog
forall a. Envelope a -> a
message (Envelope EtcdLog -> EtcdLog) -> [Envelope EtcdLog] -> [EtcdLog]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [Envelope EtcdLog]
entries)
        [(Int, ByteString)]
batch <- PersistentQueue IO Int -> Int -> Int -> IO [(Int, ByteString)]
forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
peekBatchPersistentQueue PersistentQueue IO Int
q Int
2 Int
forall a. Bounded a => a
maxBound
        Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q [(Int, ByteString)]
batch
        IO () -> IO ()
forall a. HasCallStack => IO a -> IO ()
shouldNotBlock_ (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ Async IO () -> IO ()
forall a. Async IO a -> IO a
forall (m :: * -> *) a. MonadAsync m => Async m a -> m a
wait Async IO ()
writer
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
3

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"popBatch tolerates missing backing files and traces failing deletions" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      TVar [Envelope EtcdLog]
traces <- [Envelope EtcdLog] -> IO (TVar IO [Envelope EtcdLog])
forall a. a -> IO (TVar IO a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TVar m a)
newTVarIO []
      let tracer :: Tracer IO EtcdLog
tracer = TVar IO [Envelope EtcdLog] -> Text -> Tracer IO EtcdLog
forall (m :: * -> *) msg.
(MonadFork m, MonadTime m, MonadSTM m) =>
TVar m [Envelope msg] -> Text -> Tracer m msg
traceInTVar TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces Text
"PersistentQueueSpec"
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
tracer String
dir Natural
10
      [Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1, Int
2 :: Int] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q
      [String]
files <- [String] -> [String]
forall a. Ord a => [a] -> [a]
sort ([String] -> [String]) -> IO [String] -> IO [String]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> String -> IO [String]
listDirectory String
dir
      case [String]
files of
        [String
f1, String
f2] -> do
          -- One file vanished (ENOENT, tolerated silently), the other
          -- becomes a directory so deletion fails with EISDIR (traced)
          String -> IO ()
removeFile (String
dir String -> String -> String
</> String
f1)
          String -> IO ()
removeFile (String
dir String -> String -> String
</> String
f2)
          String -> IO ()
createDirectory (String
dir String -> String -> String
</> String
f2)
        [String]
_ -> String -> IO ()
forall (m :: * -> *) a.
(HasCallStack, MonadThrow m) =>
String -> m a
failure (String -> IO ()) -> String -> IO ()
forall a b. (a -> b) -> a -> b
$ String
"expected two backing files, got: " String -> String -> String
forall a. Semigroup a => a -> a -> a
<> [String] -> String
forall b a. (Show a, IsString b) => a -> b
show [String]
files
      [(Int, ByteString)]
batch <- PersistentQueue IO Int -> Int -> Int -> IO [(Int, ByteString)]
forall (m :: * -> *) a.
MonadSTM m =>
PersistentQueue m a -> Int -> Int -> m [(a, ByteString)]
peekBatchPersistentQueue PersistentQueue IO Int
q Int
2 Int
forall a. Bounded a => a
maxBound
      Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q [(Int, ByteString)]
batch
      [Envelope EtcdLog]
entries <- TVar IO [Envelope EtcdLog] -> IO [Envelope EtcdLog]
forall a. TVar IO a -> IO a
forall (m :: * -> *) a. MonadSTM m => TVar m a -> m a
readTVarIO TVar [Envelope EtcdLog]
TVar IO [Envelope EtcdLog]
traces
      (Envelope EtcdLog -> EtcdLog
forall a. Envelope a -> a
message (Envelope EtcdLog -> EtcdLog) -> [Envelope EtcdLog] -> [EtcdLog]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [Envelope EtcdLog]
entries) [EtcdLog] -> ([EtcdLog] -> Bool) -> IO ()
forall a. (HasCallStack, Show a) => a -> (a -> Bool) -> IO ()
`shouldSatisfy` (EtcdLog -> Bool) -> [EtcdLog] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
any EtcdLog -> Bool
isDeleteFailed
      -- The queue keeps operating
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
tracer PersistentQueue IO Int
q (Int
3 :: Int)
      PersistentQueue IO Int -> IO Int
forall (m :: * -> *) a. MonadSTM m => PersistentQueue m a -> m a
peekPersistentQueue PersistentQueue IO Int
q IO Int -> Int -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int
3

  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"keeps the in-flight batch pinned across retries" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$ do
    String -> (String -> IO ()) -> IO ()
forall (m :: * -> *) r.
(MonadIO m, MonadMask m) =>
String -> (String -> m r) -> m r
withTempDir String
"persistent-queue" ((String -> IO ()) -> IO ()) -> (String -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \String
dir -> do
      PersistentQueue IO Int
q <- Tracer IO EtcdLog
-> String -> Natural -> IO (PersistentQueue IO Int)
forall (m :: * -> *) a.
(MonadLabelledSTM m, MonadIO m, FromCBOR a, MonadCatch m,
 MonadFail m) =>
Tracer IO EtcdLog -> String -> Natural -> m (PersistentQueue m a)
newPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer String
dir Natural
10
      TVar (Maybe [(Int, ByteString)])
inFlightVar <- Maybe [(Int, ByteString)]
-> IO (TVar IO (Maybe [(Int, ByteString)]))
forall a. a -> IO (TVar IO a)
forall (m :: * -> *) a. MonadSTM m => a -> m (TVar m a)
newTVarIO Maybe [(Int, ByteString)]
forall a. Maybe a
Nothing
      [Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1, Int
2 :: Int] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q
      Just [(Int, ByteString)]
batch <- TVar IO (Maybe [(Int, ByteString)])
-> PersistentQueue IO Int
-> Int
-> Int
-> IO (Maybe [(Int, ByteString)])
forall (m :: * -> *) a.
MonadSTM m =>
TVar m (Maybe [(a, ByteString)])
-> PersistentQueue m a -> Int -> Int -> m (Maybe [(a, ByteString)])
nextPendingBatch TVar (Maybe [(Int, ByteString)])
TVar IO (Maybe [(Int, ByteString)])
inFlightVar PersistentQueue IO Int
q Int
50 Int
forall a. Bounded a => a
maxBound
      ((Int, ByteString) -> Int
forall a b. (a, b) -> a
fst ((Int, ByteString) -> Int) -> [(Int, ByteString)] -> [Int]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [(Int, ByteString)]
batch) [Int] -> [Int] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [Int
1, Int
2]
      -- Messages arriving while a send is in flight (e.g. during a transient
      -- gRPC retry) must not grow the retried batch: the compare-fail dedup
      -- in putMessage would declare a grown batch delivered and its
      -- never-sent tail would be popped and lost.
      Tracer IO EtcdLog -> PersistentQueue IO Int -> Int -> IO ()
forall a (m :: * -> *).
(ToCBOR a, MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog -> PersistentQueue m a -> a -> m ()
writePersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q Int
3
      Just [(Int, ByteString)]
retried <- TVar IO (Maybe [(Int, ByteString)])
-> PersistentQueue IO Int
-> Int
-> Int
-> IO (Maybe [(Int, ByteString)])
forall (m :: * -> *) a.
MonadSTM m =>
TVar m (Maybe [(a, ByteString)])
-> PersistentQueue m a -> Int -> Int -> m (Maybe [(a, ByteString)])
nextPendingBatch TVar (Maybe [(Int, ByteString)])
TVar IO (Maybe [(Int, ByteString)])
inFlightVar PersistentQueue IO Int
q Int
50 Int
forall a. Bounded a => a
maxBound
      ((Int, ByteString) -> Int
forall a b. (a, b) -> a
fst ((Int, ByteString) -> Int) -> [(Int, ByteString)] -> [Int]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [(Int, ByteString)]
retried) [Int] -> [Int] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [Int
1, Int
2]
      -- After a successful send the pin is cleared and the next batch picks
      -- up the remaining messages.
      Tracer IO EtcdLog
-> PersistentQueue IO Int -> [(Int, ByteString)] -> IO ()
forall (m :: * -> *) a.
(MonadSTM m, MonadIO m) =>
Tracer IO EtcdLog
-> PersistentQueue m a -> [(a, ByteString)] -> m ()
popBatchPersistentQueue Tracer IO EtcdLog
forall (m :: * -> *) a. Applicative m => Tracer m a
nullTracer PersistentQueue IO Int
q [(Int, ByteString)]
retried
      STM IO () -> IO ()
forall a. HasCallStack => STM IO a -> IO a
forall (m :: * -> *) a.
(MonadSTM m, HasCallStack) =>
STM m a -> m a
atomically (STM IO () -> IO ()) -> STM IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ TVar IO (Maybe [(Int, ByteString)])
-> Maybe [(Int, ByteString)] -> STM IO ()
forall a. TVar IO a -> a -> STM IO ()
forall (m :: * -> *) a. MonadSTM m => TVar m a -> a -> STM m ()
writeTVar TVar (Maybe [(Int, ByteString)])
TVar IO (Maybe [(Int, ByteString)])
inFlightVar Maybe [(Int, ByteString)]
forall a. Maybe a
Nothing
      Just [(Int, ByteString)]
next <- TVar IO (Maybe [(Int, ByteString)])
-> PersistentQueue IO Int
-> Int
-> Int
-> IO (Maybe [(Int, ByteString)])
forall (m :: * -> *) a.
MonadSTM m =>
TVar m (Maybe [(a, ByteString)])
-> PersistentQueue m a -> Int -> Int -> m (Maybe [(a, ByteString)])
nextPendingBatch TVar (Maybe [(Int, ByteString)])
TVar IO (Maybe [(Int, ByteString)])
inFlightVar PersistentQueue IO Int
q Int
50 Int
forall a. Bounded a => a
maxBound
      ((Int, ByteString) -> Int
forall a b. (a, b) -> a
fst ((Int, ByteString) -> Int) -> [(Int, ByteString)] -> [Int]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [(Int, ByteString)]
next) [Int] -> [Int] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [Int
3]

isDeleteFailed :: EtcdLog -> Bool
isDeleteFailed :: EtcdLog -> Bool
isDeleteFailed = \case
  PersistentQueueDeleteFailed{} -> Bool
True
  EtcdLog
_ -> Bool
False

isLoadFailed :: EtcdLog -> Bool
isLoadFailed :: EtcdLog -> Bool
isLoadFailed = \case
  PersistentQueueLoadFailed{} -> Bool
True
  EtcdLog
_ -> Bool
False

shouldNotBlock :: HasCallStack => IO a -> IO a
shouldNotBlock :: forall a. HasCallStack => IO a -> IO a
shouldNotBlock IO a
action = do
  -- Generous enough to survive heavy parallel load on CI; the goal is to
  -- detect a hung action, not benchmark responsiveness.
  DiffTime -> IO a -> IO (Maybe a)
forall a. DiffTime -> IO a -> IO (Maybe a)
forall (m :: * -> *) a.
MonadTimer m =>
DiffTime -> m a -> m (Maybe a)
timeout DiffTime
5 IO a
action IO (Maybe a) -> (Maybe a -> IO a) -> IO a
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
    Maybe a
Nothing -> String -> IO a
forall (m :: * -> *) a.
(HasCallStack, MonadThrow m) =>
String -> m a
failure String
"blocked unexpectedly"
    Just a
a -> a -> IO a
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure a
a

shouldNotBlock_ :: HasCallStack => IO a -> IO ()
shouldNotBlock_ :: forall a. HasCallStack => IO a -> IO ()
shouldNotBlock_ = IO () -> IO ()
forall a. HasCallStack => IO a -> IO a
shouldNotBlock (IO () -> IO ()) -> (IO a -> IO ()) -> IO a -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO a -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void