{-# LANGUAGE NumericUnderscores #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE TypeFamilies #-}
{-# OPTIONS_GHC -Wno-x-partial #-}

-- | Parameterized concurrency test suite that works with any MonadArbiter implementation.
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)

-- | Claim and ack jobs in a loop until no more are available.
-- Backs off with 10ms delay between empty claim attempts, gives up
-- after @maxStreak@ consecutive empties.
claimAckLoop
  :: IO [a]
  -> Int
  -- ^ Max consecutive empty claims before giving up
  -> (a -> IO ())
  -- ^ Per-job action
  -> 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

-- | Drain all remaining jobs with no streak limit.
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

-- | Find duplicate entries in a list. O(n log n).
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

-- | Install the gap-free HOL violation detector on a raw connection.
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)

-- | Count violations detected by the HOL detector trigger.
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

-- | Remove the HOL detector trigger, function, and violations table.
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)

-- | Parameterized concurrency test suite. These tests need a connection pool of
-- at least 10 connections.
concurrencySpec
  :: forall payload m env
   . ( JobPayload payload
     , KnownSymbol (TableForPayload payload (RegistryOf m))
     , MonadArbiter m
     )
  => (Text -> payload)
  -- ^ Constructor for a simple test message payload
  -> (forall a. env -> m a -> IO a)
  -- ^ Runner for monad actions
  -> 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
      -- Insert a job
      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"))

      -- Worker A claims the job with 5 second visibility
      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

      -- Worker A extends visibility to 60 seconds, as a heartbeat.
      rowsUpdated <- runM env (HL.setVisibilityTimeout 60 job)
      rowsUpdated `shouldBe` 1

      -- Worker B claims on a different pool connection and gets nothing.
      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
      -- Insert a job
      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"))

      -- Claim with a 2 second visibility.
      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

      -- Wait 1 second, then extend visibility by 10 more seconds
      threadDelay 1_000_000
      void $ runM env (HL.setVisibilityTimeout 10 job)

      -- Wait 2 more seconds, past the original timeout.
      threadDelay 2_000_000

      -- The extended visibility blocks a claim.
      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
      -- Insert and claim job
      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

      -- Set visibility to 0.
      void $ runM env (HL.setVisibilityTimeout 0 job)

      -- The job is claimable at once.
      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
      -- Insert 6 ungrouped jobs, remembering their ids.
      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)
      -- No job claimed twice in the racy pass.
      raced `shouldBe` nub raced
      -- Drain the jobs the race missed. A follow-up claim surfaces the unclaimed ones.
      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
      -- Every inserted job was claimed exactly once across the racing workers and the drain.
      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
      -- Each round seeds a group and races 10 claimers with 2 inserters on the
      -- same group. Per-group ordering allows at most 1 claim per round.
      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
        -- Insert 5 jobs in one group
        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, ()))

        -- Race the claimers and inserters on the same group.
        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 [])

        -- At most 1 claim per group. A round with no claim is normal when the
        -- insert trigger holds the groups row lock.
        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, ()))

        -- Ack the racy claims, then drain the rest, recording every processed id.
        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
      -- Every job inserted across all rounds was claimed and acked exactly once.
      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
      -- Insert a single job
      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"))

      -- Three workers race to claim the same job.
      [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]

      -- Exactly one worker claims it.
      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

-- | High-contention tests for concurrency defects. Many workers, large job
-- counts, and tight timing windows.
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

        -- Run multiple rounds of concurrent insert + claim
        replicateM_ numRounds $ do
          _ <-
            mapConcurrently
              id
              ( [ do
                    -- Inserter
                    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
                        -- Claimer
                        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 ()

        -- Drain remaining
        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

        -- No job is lost.
        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 -- 100ms visibility
        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"))

        -- Track which jobs were successfully acked
        ackedRef <- Set Int64 -> IO (IORef (Set Int64))
forall a. a -> IO (IORef a)
newIORef Set Int64
forall a. Set a
Set.empty

        replicateConcurrently_ numWorkers $
          -- A 50 empty-streak limit keeps workers retrying while jobs are in flight.
          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) -- Half the visibility time
              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
        -- Every job is acked exactly once.
        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

        -- Insert a fresh job per iteration.
        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)

        -- Each iteration has exactly 1 successful ack.
        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]

          -- Per-worker ids in each group are contiguous. A gap means another
          -- worker interleaved claims from the same group.
          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, ()))

          -- Drain remaining
          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
        -- Mix of batched and single mode workers on same table
        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
          -- Insert a job
          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"))

          -- Claim it
          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
              -- Many workers try to ack the same job concurrently
              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)

        -- Each iteration has exactly 1 successful ack.
        successCounts `shouldBe` replicate numIterations 1