{-# LANGUAGE NumericUnderscores #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE TypeFamilies #-}
{-# OPTIONS_GHC -Wno-x-partial #-}
module Arbiter.Test.Concurrency
( concurrencySpec
, raceConditionSpec
, installHolDetector
, countHolViolations
, removeHolDetector
) where
import Arbiter.Core.HighLevel qualified as HL
import Arbiter.Core.Job.Types
import Arbiter.Core.MonadArbiter (MonadArbiter, RegistryOf)
import Arbiter.Core.QueueRegistry (TableForPayload)
import Control.Concurrent (threadDelay)
import Control.Monad (forM, forM_, replicateM, replicateM_, void, when)
import Data.Foldable (traverse_)
import Data.IORef (atomicModifyIORef', newIORef, readIORef)
import Data.Int (Int64)
import Data.List (nub, sort)
import Data.List.NonEmpty qualified as NE
import Data.Set qualified as Set
import Data.String (fromString)
import Data.Text (Text)
import Data.Text qualified as T
import Database.PostgreSQL.Simple qualified as PG
import GHC.TypeLits (KnownSymbol)
import Test.Hspec
import UnliftIO.Async (mapConcurrently, replicateConcurrently_)
import Arbiter.Test.StateMachine (holInstallSql, holRemoveSql, holViolTbl)
claimAckLoop
:: IO [a]
-> Int
-> (a -> IO ())
-> IO ()
claimAckLoop :: forall a. IO [a] -> Int -> (a -> IO ()) -> IO ()
claimAckLoop IO [a]
claim Int
maxStreak a -> IO ()
onJob = Int -> IO ()
go Int
0
where
go :: Int -> IO ()
go Int
streak = do
jobs <- IO [a]
claim
if null jobs
then when (streak < maxStreak) $ threadDelay 10_000 >> go (streak + 1)
else forM_ jobs onJob >> go 0
drainAll :: IO [a] -> (a -> IO ()) -> IO ()
drainAll :: forall a. IO [a] -> (a -> IO ()) -> IO ()
drainAll IO [a]
claim a -> IO ()
onJob = do
jobs <- IO [a]
claim
if null jobs then pure () else forM_ jobs onJob >> drainAll claim onJob
findDuplicates :: (Ord a) => [a] -> [a]
findDuplicates :: forall a. Ord a => [a] -> [a]
findDuplicates = Set a -> Set a -> [a] -> [a]
forall {a}. Ord a => Set a -> Set a -> [a] -> [a]
go Set a
forall a. Set a
Set.empty Set a
forall a. Set a
Set.empty
where
go :: Set a -> Set a -> [a] -> [a]
go Set a
_ Set a
dups [] = Set a -> [a]
forall a. Set a -> [a]
Set.toList Set a
dups
go Set a
seen Set a
dups (a
item : [a]
rest)
| a -> Set a -> Bool
forall a. Ord a => a -> Set a -> Bool
Set.member a
item Set a
seen = Set a -> Set a -> [a] -> [a]
go Set a
seen (a -> Set a -> Set a
forall a. Ord a => a -> Set a -> Set a
Set.insert a
item Set a
dups) [a]
rest
| Bool
otherwise = Set a -> Set a -> [a] -> [a]
go (a -> Set a -> Set a
forall a. Ord a => a -> Set a -> Set a
Set.insert a
item Set a
seen) Set a
dups [a]
rest
installHolDetector :: PG.Connection -> Text -> Text -> IO ()
installHolDetector :: Connection -> Text -> Text -> IO ()
installHolDetector Connection
conn Text
schemaName Text
tableName =
(Text -> IO ()) -> [Text] -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> (Text -> IO Int64) -> Text -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Connection -> Query -> IO Int64
PG.execute_ Connection
conn (Query -> IO Int64) -> (Text -> Query) -> Text -> IO Int64
forall b c a. (b -> c) -> (a -> b) -> a -> c
. String -> Query
forall a. IsString a => String -> a
fromString (String -> Query) -> (Text -> String) -> Text -> Query
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> String
T.unpack) (Text -> Text -> [Text]
holInstallSql Text
schemaName Text
tableName)
countHolViolations :: PG.Connection -> Text -> Text -> IO Int64
countHolViolations :: Connection -> Text -> Text -> IO Int64
countHolViolations Connection
conn Text
schemaName Text
tableName = do
[PG.Only count] <-
Connection -> Query -> IO [Only Int64]
forall r. FromRow r => Connection -> Query -> IO [r]
PG.query_ Connection
conn
(Query -> IO [Only Int64]) -> Query -> IO [Only Int64]
forall a b. (a -> b) -> a -> b
$ String -> Query
forall a. IsString a => String -> a
fromString
(String -> Query) -> String -> Query
forall a b. (a -> b) -> a -> b
$ Text -> String
T.unpack
(Text -> String) -> Text -> String
forall a b. (a -> b) -> a -> b
$ Text
"SELECT count(*)::bigint FROM " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> Text
holViolTbl Text
schemaName Text
tableName
pure count
removeHolDetector :: PG.Connection -> Text -> Text -> IO ()
removeHolDetector :: Connection -> Text -> Text -> IO ()
removeHolDetector Connection
conn Text
schemaName Text
tableName =
(Text -> IO ()) -> [Text] -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> (Text -> IO Int64) -> Text -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Connection -> Query -> IO Int64
PG.execute_ Connection
conn (Query -> IO Int64) -> (Text -> Query) -> Text -> IO Int64
forall b c a. (b -> c) -> (a -> b) -> a -> c
. String -> Query
forall a. IsString a => String -> a
fromString (String -> Query) -> (Text -> String) -> Text -> Query
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> String
T.unpack) (Text -> Text -> [Text]
holRemoveSql Text
schemaName Text
tableName)
concurrencySpec
:: forall payload m env
. ( JobPayload payload
, KnownSymbol (TableForPayload payload (RegistryOf m))
, MonadArbiter m
)
=> (Text -> payload)
-> (forall a. env -> m a -> IO a)
-> SpecWith env
concurrencySpec :: forall payload (m :: * -> *) env.
(JobPayload payload,
KnownSymbol (TableForPayload payload (RegistryOf m)),
MonadArbiter m) =>
(Text -> payload) -> (forall a. env -> m a -> IO a) -> SpecWith env
concurrencySpec Text -> payload
mkMessage forall a. env -> m a -> IO a
runM = do
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Heartbeat Visibility Across Connections" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"setVisibilityTimeout updates are immediately visible to concurrent workers" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ())
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall a b. (a -> b) -> a -> b
$ env
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys)))
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a b. (a -> b) -> a -> b
$ JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"heartbeat-visibility-test") (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"Heartbeat"))
claimed1 <- env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (Int
-> NominalDiffTime
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
1 NominalDiffTime
5) :: IO [JobRead payload]
length claimed1 `shouldBe` 1
let job = [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> JobRecord payload Int64 Text UTCTime PayloadKeys
forall a. HasCallStack => [a] -> a
head [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimed1
rowsUpdated <- runM env (HL.setVisibilityTimeout 60 job)
rowsUpdated `shouldBe` 1
claimed2 <- runM env (HL.claimNextVisibleJobs 1 5) :: IO [JobRead payload]
length claimed2 `shouldBe` 0
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"heartbeat prevents reclaim after original visibility timeout expires" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ())
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall a b. (a -> b) -> a -> b
$ env
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys)))
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a b. (a -> b) -> a -> b
$ JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"heartbeat-timeout-test") (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"Timeout"))
claimed1 <- env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (Int
-> NominalDiffTime
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
1 NominalDiffTime
2) :: IO [JobRead payload]
length claimed1 `shouldBe` 1
let job = [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> JobRecord payload Int64 Text UTCTime PayloadKeys
forall a. HasCallStack => [a] -> a
head [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimed1
threadDelay 1_000_000
void $ runM env (HL.setVisibilityTimeout 10 job)
threadDelay 2_000_000
claimed2 <- runM env (HL.claimNextVisibleJobs 1 2) :: IO [JobRead payload]
length claimed2 `shouldBe` 0
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"visibility timeout of 0 makes job immediately claimable" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ())
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall a b. (a -> b) -> a -> b
$ env
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys)))
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a b. (a -> b) -> a -> b
$ JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"visibility-zero-test") (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"Zero"))
claimed1 <- env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (Int
-> NominalDiffTime
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
1 NominalDiffTime
60) :: IO [JobRead payload]
length claimed1 `shouldBe` 1
let job = [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> JobRecord payload Int64 Text UTCTime PayloadKeys
forall a. HasCallStack => [a] -> a
head [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimed1
void $ runM env (HL.setVisibilityTimeout 0 job)
claimed2 <- runM env (HL.claimNextVisibleJobs 1 60) :: IO [JobRead payload]
length claimed2 `shouldBe` 1
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Concurrent Job Claims" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrent workers claim disjoint ungrouped jobs" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
inserted <- env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys])
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload]
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch (Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
6 (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"Concurrent"))
let insertedIds = [Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ((JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64)
-> [JobRecord payload Int64 Text UTCTime PayloadKeys] -> [Int64]
forall a b. (a -> b) -> [a] -> [b]
map JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey [JobRecord payload Int64 Text UTCTime PayloadKeys]
inserted)
[claimsA, claimsB, claimsC] <-
mapConcurrently
(\Int
_ -> env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (Int
-> NominalDiffTime
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
2 NominalDiffTime
60) :: IO [JobRead payload])
[1 :: Int, 2, 3]
let raced = (JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64)
-> [JobRecord payload Int64 Text UTCTime PayloadKeys] -> [Int64]
forall a b. (a -> b) -> [a] -> [b]
map JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey ([JobRecord payload Int64 Text UTCTime PayloadKeys]
claimsA [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. Semigroup a => a -> a -> a
<> [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimsB [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. Semigroup a => a -> a -> a
<> [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimsC)
raced `shouldBe` nub raced
drainedRef <- newIORef ([] :: [Int64])
drainAll (runM env (HL.claimNextVisibleJobs 6 60) :: IO [JobRead payload]) $ \JobRecord payload Int64 Text UTCTime PayloadKeys
job ->
IORef [Int64] -> ([Int64] -> ([Int64], ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef [Int64]
drainedRef (\[Int64]
acc -> (JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRecord payload Int64 Text UTCTime PayloadKeys
job Int64 -> [Int64] -> [Int64]
forall a. a -> [a] -> [a]
: [Int64]
acc, ()))
drained <- readIORef drainedRef
let allClaimed = [Int64]
raced [Int64] -> [Int64] -> [Int64]
forall a. Semigroup a => a -> a -> a
<> [Int64]
drained
allClaimed `shouldBe` nub allClaimed
Set.fromList allClaimed `shouldBe` insertedIds
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrent workers respect per-group ordering for grouped jobs" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numIterations :: Int
numIterations = Int
500
numWorkers :: Int
numWorkers = Int
10
seedPerRound :: Int
seedPerRound = Int
5
insertersPerRound :: Int
insertersPerRound = Int
2
violationRef <- Int -> IO (IORef Int)
forall a. a -> IO (IORef a)
newIORef (Int
0 :: Int)
insertedRef <- newIORef (0 :: Int)
processedRef <- newIORef (Set.empty :: Set.Set Int64)
forM_ [(1 :: Int) .. numIterations] $ \Int
_ -> do
IO [JobRecord payload Int64 Text UTCTime PayloadKeys] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void
(IO [JobRecord payload Int64 Text UTCTime PayloadKeys] -> IO ())
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys] -> IO ()
forall a b. (a -> b) -> a -> b
$ env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env
(m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys])
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload]
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch (Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
seedPerRound (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"hol-stress") (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"Grouped"))
IORef Int -> (Int -> (Int, ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef Int
insertedRef (\Int
count -> (Int
count Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
seedPerRound Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
insertersPerRound, ()))
results <-
(IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys])
-> [IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
-> IO [[JobRecord payload Int64 Text UTCTime PayloadKeys]]
forall (m :: * -> *) (t :: * -> *) a b.
(MonadUnliftIO m, Traversable t) =>
(a -> m b) -> t a -> m (t b)
mapConcurrently
IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. a -> a
id
([IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
-> IO [[JobRecord payload Int64 Text UTCTime PayloadKeys]])
-> [IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
-> IO [[JobRecord payload Int64 Text UTCTime PayloadKeys]]
forall a b. (a -> b) -> a -> b
$ (Int -> IO [JobRecord payload Int64 Text UTCTime PayloadKeys])
-> [Int] -> [IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
forall a b. (a -> b) -> [a] -> [b]
map
(\Int
_ -> env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (Int
-> NominalDiffTime
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
1 NominalDiffTime
60) :: IO [JobRead payload])
[(Int
1 :: Int) .. Int
numWorkers]
[IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
-> [IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
-> [IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
forall a. Semigroup a => a -> a -> a
<> Int
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> [IO [JobRecord payload Int64 Text UTCTime PayloadKeys]]
forall a. Int -> a -> [a]
replicate
Int
insertersPerRound
(IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (env
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys)))
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a b. (a -> b) -> a -> b
$ JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys)))
-> JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a b. (a -> b) -> a -> b
$ Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"hol-stress") (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"concurrent")) IO ()
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure [])
let totalClaimed = [Int] -> Int
forall a. Num a => [a] -> a
forall (t :: * -> *) a. (Foldable t, Num a) => t a -> a
sum (([JobRecord payload Int64 Text UTCTime PayloadKeys] -> Int)
-> [[JobRecord payload Int64 Text UTCTime PayloadKeys]] -> [Int]
forall a b. (a -> b) -> [a] -> [b]
map [JobRecord payload Int64 Text UTCTime PayloadKeys] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [[JobRecord payload Int64 Text UTCTime PayloadKeys]]
results)
when (totalClaimed > 1) $
atomicModifyIORef' violationRef (\Int
count -> (Int
count Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
1, ()))
let record JobRecord payload Int64 Text UTCTime PayloadKeys
job = IORef (Set Int64) -> (Set Int64 -> (Set Int64, ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef (Set Int64)
processedRef (\Set Int64
acc -> (Int64 -> Set Int64 -> Set Int64
forall a. Ord a => a -> Set a -> Set a
Set.insert (JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRecord payload Int64 Text UTCTime PayloadKeys
job) Set Int64
acc, ()))
forM_ (concat results) $ \JobRecord payload Int64 Text UTCTime PayloadKeys
job -> JobRecord payload Int64 Text UTCTime PayloadKeys -> IO ()
record JobRecord payload Int64 Text UTCTime PayloadKeys
job IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRecord payload Int64 Text UTCTime PayloadKeys -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRecord payload Int64 Text UTCTime PayloadKeys
job))
drainAll (runM env (HL.claimNextVisibleJobs 100 60) :: IO [JobRead payload]) $
\JobRecord payload Int64 Text UTCTime PayloadKeys
job -> JobRecord payload Int64 Text UTCTime PayloadKeys -> IO ()
record JobRecord payload Int64 Text UTCTime PayloadKeys
job IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRecord payload Int64 Text UTCTime PayloadKeys -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRecord payload Int64 Text UTCTime PayloadKeys
job))
violations <- readIORef violationRef
violations `shouldBe` 0
inserted <- readIORef insertedRef
processed <- readIORef processedRef
Set.size processed `shouldBe` inserted
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"FOR UPDATE SKIP LOCKED prevents concurrent claims of same job" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ())
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO ()
forall a b. (a -> b) -> a -> b
$ env
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys)))
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
-> IO (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall a b. (a -> b) -> a -> b
$ JobWrite payload
-> m (Maybe (JobRecord payload Int64 Text UTCTime PayloadKeys))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"skip-locked-test") (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"Single"))
[claimsA, claimsB, claimsC] <-
(Int -> IO [JobRecord payload Int64 Text UTCTime PayloadKeys])
-> [Int] -> IO [[JobRecord payload Int64 Text UTCTime PayloadKeys]]
forall (m :: * -> *) (t :: * -> *) a b.
(MonadUnliftIO m, Traversable t) =>
(a -> m b) -> t a -> m (t b)
mapConcurrently
(\Int
_ -> env
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> IO [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall a. env -> m a -> IO a
runM env
env (Int
-> NominalDiffTime
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
2 NominalDiffTime
60) :: IO [JobRead payload])
[Int
1 :: Int, Int
2, Int
3]
let totalClaimed = [JobRecord payload Int64 Text UTCTime PayloadKeys] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimsA Int -> Int -> Int
forall a. Num a => a -> a -> a
+ [JobRecord payload Int64 Text UTCTime PayloadKeys] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimsB Int -> Int -> Int
forall a. Num a => a -> a -> a
+ [JobRecord payload Int64 Text UTCTime PayloadKeys] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [JobRecord payload Int64 Text UTCTime PayloadKeys]
claimsC
totalClaimed `shouldBe` 1
raceConditionSpec
:: forall payload m env
. ( JobPayload payload
, KnownSymbol (TableForPayload payload (RegistryOf m))
, MonadArbiter m
)
=> (Text -> payload)
-> (forall a. env -> m a -> IO a)
-> SpecWith env
raceConditionSpec :: forall payload (m :: * -> *) env.
(JobPayload payload,
KnownSymbol (TableForPayload payload (RegistryOf m)),
MonadArbiter m) =>
(Text -> payload) -> (forall a. env -> m a -> IO a) -> SpecWith env
raceConditionSpec Text -> payload
mkMessage forall a. env -> m a -> IO a
runM = do
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Race Condition Stress Tests" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"High-Contention Claim Stress" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"no duplicate claims under extreme contention (ungrouped)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numJobs :: Int
numJobs = Int
500
numWorkers :: Int
numWorkers = Int
40
IO [JobRead payload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO [JobRead payload] -> IO ()) -> IO [JobRead payload] -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env (m [JobRead payload] -> IO [JobRead payload])
-> m [JobRead payload] -> IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload] -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch (Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
numJobs (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"stress"))
claimedRef <- [Int64] -> IO (IORef [Int64])
forall a. a -> IO (IORef a)
newIORef ([] :: [Int64])
replicateConcurrently_ numWorkers $
claimAckLoop (runM env (HL.claimNextVisibleJobs 3 60) :: IO [JobRead payload]) 20 $ \JobRead payload
job -> do
IORef [Int64] -> ([Int64] -> ([Int64], ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef [Int64]
claimedRef (\[Int64]
acc -> (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job Int64 -> [Int64] -> [Int64]
forall a. a -> [a] -> [a]
: [Int64]
acc, ()))
IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> IO Int64 -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
allClaimed <- readIORef claimedRef
findDuplicates allClaimed `shouldBe` []
length allClaimed `shouldBe` numJobs
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"no duplicate claims under extreme contention (many groups)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numGroups :: Int
numGroups = Int
50
jobsPerGroup :: Int
jobsPerGroup = Int
10
numWorkers :: Int
numWorkers = Int
40
[Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1 .. Int
numGroups] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \(Int
groupIndex :: Int) -> do
let groupName :: Text
groupName = Text
"stress-group-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
groupIndex)
IO [JobRead payload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void
(IO [JobRead payload] -> IO ()) -> IO [JobRead payload] -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env
(m [JobRead payload] -> IO [JobRead payload])
-> m [JobRead payload] -> IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload] -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch
(Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
jobsPerGroup (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
groupName) (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"grouped"))
claimedRef <- [Int64] -> IO (IORef [Int64])
forall a. a -> IO (IORef a)
newIORef ([] :: [Int64])
replicateConcurrently_ numWorkers $
claimAckLoop (runM env (HL.claimNextVisibleJobs 5 60) :: IO [JobRead payload]) 20 $ \JobRead payload
job -> do
IORef [Int64] -> ([Int64] -> ([Int64], ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef [Int64]
claimedRef (\[Int64]
acc -> (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job Int64 -> [Int64] -> [Int64]
forall a. a -> [a] -> [a]
: [Int64]
acc, ()))
IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> IO Int64 -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
allClaimed <- readIORef claimedRef
let totalJobs = Int
numGroups Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
jobsPerGroup
findDuplicates allClaimed `shouldBe` []
length allClaimed `shouldBe` totalJobs
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Insert + Claim Race" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"no lost jobs under concurrent insert/claim pressure" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numRounds :: Int
numRounds = Int
100
jobsPerRound :: Int
jobsPerRound = Int
20
numClaimers :: Int
numClaimers = Int
10
insertedRef <- Set Int64 -> IO (IORef (Set Int64))
forall a. a -> IO (IORef a)
newIORef Set Int64
forall a. Set a
Set.empty
claimedRef <- newIORef Set.empty
replicateM_ numRounds $ do
_ <-
mapConcurrently
id
( [ do
jobs <- runM env $ HL.insertJobsBatch (replicate jobsPerRound $ defaultJob (mkMessage "race"))
let ids = [Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ([Int64] -> Set Int64) -> [Int64] -> Set Int64
forall a b. (a -> b) -> a -> b
$ (JobRead payload -> Int64) -> [JobRead payload] -> [Int64]
forall a b. (a -> b) -> [a] -> [b]
map JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey [JobRead payload]
jobs
atomicModifyIORef' insertedRef (\Set Int64
acc -> (Set Int64 -> Set Int64 -> Set Int64
forall a. Ord a => Set a -> Set a -> Set a
Set.union Set Int64
acc Set Int64
ids, ()))
]
<> replicate
numClaimers
( do
claimed <- runM env (HL.claimNextVisibleJobs 10 60) :: IO [JobRead payload]
let ids = [Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ([Int64] -> Set Int64) -> [Int64] -> Set Int64
forall a b. (a -> b) -> a -> b
$ (JobRead payload -> Int64) -> [JobRead payload] -> [Int64]
forall a b. (a -> b) -> [a] -> [b]
map JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey [JobRead payload]
claimed
atomicModifyIORef' claimedRef (\Set Int64
acc -> (Set Int64 -> Set Int64 -> Set Int64
forall a. Ord a => Set a -> Set a -> Set a
Set.union Set Int64
acc Set Int64
ids, ()))
forM_ claimed $ \JobRead payload
job -> env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
)
)
pure ()
drainAll (runM env (HL.claimNextVisibleJobs 100 60) :: IO [JobRead payload]) $ \JobRead payload
job -> do
IORef (Set Int64) -> (Set Int64 -> (Set Int64, ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef (Set Int64)
claimedRef (\Set Int64
acc -> (Int64 -> Set Int64 -> Set Int64
forall a. Ord a => a -> Set a -> Set a
Set.insert (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job) Set Int64
acc, ()))
IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> IO Int64 -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
inserted <- readIORef insertedRef
claimed <- readIORef claimedRef
Set.size inserted `shouldBe` (numRounds * jobsPerRound)
claimed `shouldBe` inserted
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Visibility Timeout Races" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"no duplicate processing with aggressive visibility expiration" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numJobs :: Int
numJobs = Int
50
numWorkers :: Int
numWorkers = Int
20
visibilityMs :: Int
visibilityMs = Int
100
IO [JobRead payload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO [JobRead payload] -> IO ()) -> IO [JobRead payload] -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env (m [JobRead payload] -> IO [JobRead payload])
-> m [JobRead payload] -> IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload] -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch (Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
numJobs (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"timeout-race"))
ackedRef <- Set Int64 -> IO (IORef (Set Int64))
forall a. a -> IO (IORef a)
newIORef Set Int64
forall a. Set a
Set.empty
replicateConcurrently_ numWorkers $
claimAckLoop (runM env (HL.claimNextVisibleJobs 2 (fromIntegral visibilityMs / 1000)) :: IO [JobRead payload]) 50 $
\JobRead payload
job -> do
Int -> IO ()
threadDelay (Int
visibilityMs Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
500)
rows <- env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
when (rows > 0) $
atomicModifyIORef' ackedRef (\Set Int64
acc -> (Int64 -> Set Int64 -> Set Int64
forall a. Ord a => a -> Set a -> Set a
Set.insert (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job) Set Int64
acc, ()))
acked <- readIORef ackedRef
Set.size acked `shouldBe` numJobs
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"exactly one ack succeeds per job under reclaim pressure" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numAttempts :: Int
numAttempts = Int
50
numWorkersPerAttempt :: Int
numWorkersPerAttempt = Int
5
successCounts <- Int -> IO Int -> IO [Int]
forall (m :: * -> *) a. Applicative m => Int -> m a -> m [a]
replicateM Int
numAttempts (IO Int -> IO [Int]) -> IO Int -> IO [Int]
forall a b. (a -> b) -> a -> b
$ do
IO (Maybe (JobRead payload)) -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Maybe (JobRead payload)) -> IO ())
-> IO (Maybe (JobRead payload)) -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m (Maybe (JobRead payload)) -> IO (Maybe (JobRead payload))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRead payload)) -> IO (Maybe (JobRead payload)))
-> m (Maybe (JobRead payload)) -> IO (Maybe (JobRead payload))
forall a b. (a -> b) -> a -> b
$ JobWrite payload -> m (Maybe (JobRead payload))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"single-race"))
results <-
(IO Int -> IO Int) -> [IO Int] -> IO [Int]
forall (m :: * -> *) (t :: * -> *) a b.
(MonadUnliftIO m, Traversable t) =>
(a -> m b) -> t a -> m (t b)
mapConcurrently
IO Int -> IO Int
forall a. a -> a
id
([IO Int] -> IO [Int]) -> [IO Int] -> IO [Int]
forall a b. (a -> b) -> a -> b
$ Int -> IO Int -> [IO Int]
forall a. Int -> a -> [a]
replicate Int
numWorkersPerAttempt
(IO Int -> [IO Int]) -> IO Int -> [IO Int]
forall a b. (a -> b) -> a -> b
$ do
claimed <- env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env (Int -> NominalDiffTime -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
1 NominalDiffTime
0.1) :: IO [JobRead payload]
case claimed of
[] -> Int -> IO Int
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Int
0
(JobRead payload
job : [JobRead payload]
_) -> do
Int -> IO ()
threadDelay Int
50_000
rows <- env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
pure (fromIntegral rows :: Int)
pure (sum results)
successCounts `shouldBe` replicate numAttempts 1
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Batched Mode Stress" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"no duplicate claims in batched mode under heavy contention" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numGroups :: Int
numGroups = Int
30
jobsPerGroup :: Int
jobsPerGroup = Int
20
numWorkers :: Int
numWorkers = Int
25
[Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1 .. Int
numGroups] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \(Int
groupIndex :: Int) -> do
let groupName :: Text
groupName = Text
"batch-stress-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
groupIndex)
IO [JobRead payload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void
(IO [JobRead payload] -> IO ()) -> IO [JobRead payload] -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env
(m [JobRead payload] -> IO [JobRead payload])
-> m [JobRead payload] -> IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload] -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch
(Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
jobsPerGroup (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
groupName) (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"batched"))
claimedRef <- [Int64] -> IO (IORef [Int64])
forall a. a -> IO (IORef a)
newIORef ([] :: [Int64])
replicateConcurrently_ numWorkers $
claimAckLoop (concatMap NE.toList <$> runM env (HL.claimNextVisibleJobsBatched 5 4 60) :: IO [JobRead payload]) 20 $ \JobRead payload
job -> do
IORef [Int64] -> ([Int64] -> ([Int64], ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef [Int64]
claimedRef (\[Int64]
acc -> (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job Int64 -> [Int64] -> [Int64]
forall a. a -> [a] -> [a]
: [Int64]
acc, ()))
IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> IO Int64 -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
allClaimed <- readIORef claimedRef
let totalJobs = Int
numGroups Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
jobsPerGroup
findDuplicates allClaimed `shouldBe` []
length allClaimed `shouldBe` totalJobs
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"batched claims from same group are contiguous (no interleaving)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numGroups :: Int
numGroups = Int
5
jobsPerGroup :: Int
jobsPerGroup = Int
50
numWorkers :: Int
numWorkers = Int
20
numRounds :: Int
numRounds = Int
30
gapCount <- Int -> IO (IORef Int)
forall a. a -> IO (IORef a)
newIORef (Int
0 :: Int)
replicateM_ numRounds $ do
forM_ [1 .. numGroups] $ \(Int
groupIndex :: Int) ->
IO [JobRead payload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void
(IO [JobRead payload] -> IO ()) -> IO [JobRead payload] -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env
(m [JobRead payload] -> IO [JobRead payload])
-> m [JobRead payload] -> IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload] -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch
(Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
jobsPerGroup (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just (Text
"split-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
groupIndex))) (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"split"))
resultsRef <- newIORef ([] :: [(Int, [(Maybe Text, Int64)])])
_ <-
mapConcurrently
( \Int
workerId -> do
claimed <- (NonEmpty (JobRead payload) -> [JobRead payload])
-> [NonEmpty (JobRead payload)] -> [JobRead payload]
forall (t :: * -> *) a b. Foldable t => (a -> [b]) -> t a -> [b]
concatMap NonEmpty (JobRead payload) -> [JobRead payload]
forall a. NonEmpty a -> [a]
NE.toList ([NonEmpty (JobRead payload)] -> [JobRead payload])
-> IO [NonEmpty (JobRead payload)] -> IO [JobRead payload]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> env
-> m [NonEmpty (JobRead payload)]
-> IO [NonEmpty (JobRead payload)]
forall a. env -> m a -> IO a
runM env
env (Int -> Int -> NominalDiffTime -> m [NonEmpty (JobRead payload)]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> Int -> NominalDiffTime -> m [NonEmpty (JobRead payload)]
HL.claimNextVisibleJobsBatched Int
10 Int
3 NominalDiffTime
60) :: IO [JobRead payload]
atomicModifyIORef' resultsRef (\[(Int, [(Maybe Text, Int64)])]
acc -> ((Int
workerId, (JobRead payload -> (Maybe Text, Int64))
-> [JobRead payload] -> [(Maybe Text, Int64)]
forall a b. (a -> b) -> [a] -> [b]
map (\JobRead payload
job -> (JobRead payload -> Maybe Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> Maybe Text
groupKey JobRead payload
job, JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job)) [JobRead payload]
claimed) (Int, [(Maybe Text, Int64)])
-> [(Int, [(Maybe Text, Int64)])] -> [(Int, [(Maybe Text, Int64)])]
forall a. a -> [a] -> [a]
: [(Int, [(Maybe Text, Int64)])]
acc, ()))
forM_ claimed $ \JobRead payload
job -> env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
)
[1 .. numWorkers]
results <- readIORef resultsRef
let perGroup Text
wantedGroup =
[ [Int64] -> [Int64]
forall a. Ord a => [a] -> [a]
sort [Int64]
ids
| (Int
_, [(Maybe Text, Int64)]
jobs) <- [(Int, [(Maybe Text, Int64)])]
results
, let ids :: [Int64]
ids = [Int64
jobId | (Just Text
groupName, Int64
jobId) <- [(Maybe Text, Int64)]
jobs, Text
groupName Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
wantedGroup]
, Bool -> Bool
not ([Int64] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [Int64]
ids)
]
allGroupKeys = [Text] -> [Text]
forall a. Eq a => [a] -> [a]
nub [Text
groupName | (Int
_, [(Maybe Text, Int64)]
jobs) <- [(Int, [(Maybe Text, Int64)])]
results, (Just Text
groupName, Int64
_) <- [(Maybe Text, Int64)]
jobs]
hasGaps [a]
ids = ((a, a) -> Bool) -> [(a, a)] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
any (\(a
prev, a
next) -> a
next a -> a -> a
forall a. Num a => a -> a -> a
- a
prev a -> a -> Bool
forall a. Eq a => a -> a -> Bool
/= a
1) ([(a, a)] -> Bool) -> [(a, a)] -> Bool
forall a b. (a -> b) -> a -> b
$ [a] -> [a] -> [(a, a)]
forall a b. [a] -> [b] -> [(a, b)]
zip [a]
ids (Int -> [a] -> [a]
forall a. Int -> [a] -> [a]
drop Int
1 [a]
ids)
gappedGroups =
[ Text
groupName
| Text
groupName <- [Text]
allGroupKeys
, ([Int64] -> Bool) -> [[Int64]] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
any [Int64] -> Bool
forall {a}. (Eq a, Num a) => [a] -> Bool
hasGaps (Text -> [[Int64]]
perGroup Text
groupName)
]
when (not (null gappedGroups)) $
atomicModifyIORef' gapCount (\Int
count -> (Int
count Int -> Int -> Int
forall a. Num a => a -> a -> a
+ [Text] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [Text]
gappedGroups, ()))
drainAll (concatMap NE.toList <$> runM env (HL.claimNextVisibleJobsBatched 100 20 60) :: IO [JobRead payload]) $
\JobRead payload
job -> IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> IO Int64 -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
readIORef gapCount >>= (`shouldBe` 0)
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"batched and single mode don't interfere" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numGroups :: Int
numGroups = Int
20
jobsPerGroup :: Int
jobsPerGroup = Int
15
numBatchedWorkers :: Int
numBatchedWorkers = Int
10
numSingleWorkers :: Int
numSingleWorkers = Int
10
[Int] -> (Int -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [Int
1 .. Int
numGroups] ((Int -> IO ()) -> IO ()) -> (Int -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \(Int
groupIndex :: Int) -> do
let groupName :: Text
groupName = Text
"mixed-mode-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
groupIndex)
IO [JobRead payload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void
(IO [JobRead payload] -> IO ()) -> IO [JobRead payload] -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env
(m [JobRead payload] -> IO [JobRead payload])
-> m [JobRead payload] -> IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ [JobWrite payload] -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch
(Int -> JobWrite payload -> [JobWrite payload]
forall a. Int -> a -> [a]
replicate Int
jobsPerGroup (JobWrite payload -> [JobWrite payload])
-> JobWrite payload -> [JobWrite payload]
forall a b. (a -> b) -> a -> b
$ Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
groupName) (JobWrite payload -> JobWrite payload)
-> JobWrite payload -> JobWrite payload
forall a b. (a -> b) -> a -> b
$ payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"mixed"))
claimedRef <- [Int64] -> IO (IORef [Int64])
forall a. a -> IO (IORef a)
newIORef ([] :: [Int64])
let onJob JobRead payload
job = do
IORef [Int64] -> ([Int64] -> ([Int64], ())) -> IO ()
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef [Int64]
claimedRef (\[Int64]
acc -> (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job Int64 -> [Int64] -> [Int64]
forall a. a -> [a] -> [a]
: [Int64]
acc, ()))
IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Int64 -> IO ()) -> IO Int64 -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
batchedWorker =
IO [JobRead payload] -> Int -> (JobRead payload -> IO ()) -> IO ()
forall a. IO [a] -> Int -> (a -> IO ()) -> IO ()
claimAckLoop ((NonEmpty (JobRead payload) -> [JobRead payload])
-> [NonEmpty (JobRead payload)] -> [JobRead payload]
forall (t :: * -> *) a b. Foldable t => (a -> [b]) -> t a -> [b]
concatMap NonEmpty (JobRead payload) -> [JobRead payload]
forall a. NonEmpty a -> [a]
NE.toList ([NonEmpty (JobRead payload)] -> [JobRead payload])
-> IO [NonEmpty (JobRead payload)] -> IO [JobRead payload]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> env
-> m [NonEmpty (JobRead payload)]
-> IO [NonEmpty (JobRead payload)]
forall a. env -> m a -> IO a
runM env
env (Int -> Int -> NominalDiffTime -> m [NonEmpty (JobRead payload)]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> Int -> NominalDiffTime -> m [NonEmpty (JobRead payload)]
HL.claimNextVisibleJobsBatched Int
3 Int
3 NominalDiffTime
60) :: IO [JobRead payload]) Int
50 JobRead payload -> IO ()
onJob
singleWorker =
IO [JobRead payload] -> Int -> (JobRead payload -> IO ()) -> IO ()
forall a. IO [a] -> Int -> (a -> IO ()) -> IO ()
claimAckLoop (env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env (Int -> NominalDiffTime -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
2 NominalDiffTime
60) :: IO [JobRead payload]) Int
50 JobRead payload -> IO ()
onJob
_ <-
mapConcurrently
id
(replicate numBatchedWorkers batchedWorker <> replicate numSingleWorkers singleWorker)
allClaimed <- readIORef claimedRef
let totalJobs = Int
numGroups Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
jobsPerGroup
findDuplicates allClaimed `shouldBe` []
length allClaimed `shouldBe` totalJobs
String -> SpecWith env -> SpecWith env
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"Ack Race Conditions" (SpecWith env -> SpecWith env) -> SpecWith env -> SpecWith env
forall a b. (a -> b) -> a -> b
$ do
String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrent ack attempts on same job - exactly one succeeds" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
let numIterations :: Int
numIterations = Int
100
numRacingWorkers :: Int
numRacingWorkers = Int
10
successCounts <- [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
1 .. Int
numIterations] ((Int -> IO Int) -> IO [Int]) -> (Int -> IO Int) -> IO [Int]
forall a b. (a -> b) -> a -> b
$ \(Int
_ :: Int) -> do
IO (Maybe (JobRead payload)) -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Maybe (JobRead payload)) -> IO ())
-> IO (Maybe (JobRead payload)) -> IO ()
forall a b. (a -> b) -> a -> b
$ env -> m (Maybe (JobRead payload)) -> IO (Maybe (JobRead payload))
forall a. env -> m a -> IO a
runM env
env (m (Maybe (JobRead payload)) -> IO (Maybe (JobRead payload)))
-> m (Maybe (JobRead payload)) -> IO (Maybe (JobRead payload))
forall a b. (a -> b) -> a -> b
$ JobWrite payload -> m (Maybe (JobRead payload))
forall payload (m :: * -> *).
QueueOperation m payload =>
JobWrite payload -> m (Maybe (JobRead payload))
HL.insertJob (payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob (Text -> payload
mkMessage Text
"ack-race"))
claimed <- env -> m [JobRead payload] -> IO [JobRead payload]
forall a. env -> m a -> IO a
runM env
env (Int -> NominalDiffTime -> m [JobRead payload]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> m [JobRead payload]
HL.claimNextVisibleJobs Int
1 NominalDiffTime
60) :: IO [JobRead payload]
case claimed of
[] -> Int -> IO Int
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Int
0
(JobRead payload
job : [JobRead payload]
_) -> do
results <-
(IO Int -> IO Int) -> [IO Int] -> IO [Int]
forall (m :: * -> *) (t :: * -> *) a b.
(MonadUnliftIO m, Traversable t) =>
(a -> m b) -> t a -> m (t b)
mapConcurrently
IO Int -> IO Int
forall a. a -> a
id
([IO Int] -> IO [Int]) -> [IO Int] -> IO [Int]
forall a b. (a -> b) -> a -> b
$ Int -> IO Int -> [IO Int]
forall a. Int -> a -> [a]
replicate Int
numRacingWorkers
(IO Int -> [IO Int]) -> IO Int -> [IO Int]
forall a b. (a -> b) -> a -> b
$ do
rows <- env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env (JobRead payload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob JobRead payload
job)
pure (fromIntegral rows :: Int)
pure (sum results)
successCounts `shouldBe` replicate numIterations 1