{-# LANGUAGE DataKinds #-}
{-# LANGUAGE DeriveAnyClass #-}
{-# LANGUAGE DeriveGeneric #-}
{-# LANGUAGE DerivingStrategies #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TypeApplications #-}

-- | Backend-parameterized integration tests for per-job concurrency limits over the shared
-- 'CLReg' registry. Each job's key comes from the 'HasConcurrency CLPayload' instance.
module Arbiter.Test.ConcurrencyLimit
  ( CLPayload (..)
  , CLReg
  , concurrencyTable
  , concurrencyLimitSpec
  ) where

import Arbiter.Core.Codec (Col (..), col, pval)
import Arbiter.Core.Concurrency.Schema (arbiterConcurrencyTable)
import Arbiter.Core.Concurrency.Spec
  ( ConcurrencyFor
  , ConcurrencyKey (..)
  , ConcurrencyPolicy (..)
  , HasConcurrency (..)
  , chooseWhen
  , collectPolicies
  , concurrencyBy
  , concurrencyByCase
  , concurrencyPool
  , registryConcurrencyPolicies
  , runConcurrencyFor
  )
import Arbiter.Core.Concurrency.Stats (ConcurrencyPolicyUpdate (..))
import Arbiter.Core.HighLevel qualified as HL
import Arbiter.Core.Job.DLQ qualified as DLQ
import Arbiter.Core.Job.Schema (jobQueueTable)
import Arbiter.Core.Job.Types
  ( DedupKey (..)
  , JobRead
  , JobWrite
  , defaultGroupedJob
  , defaultJob
  , payload
  , setDedupKey
  )
import Arbiter.Core.MonadArbiter (HasRegistry, getSchema, withDbTransaction)
import Arbiter.Core.MonadArbiter qualified as MA
import Arbiter.Core.QueueRegistry (Queue)
import Arbiter.Core.Sql.Concurrency qualified as Tmpl
import Control.Concurrent (threadDelay)
import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)
import Control.Monad (void)
import Control.Monad.IO.Class (liftIO)
import Data.Aeson (FromJSON, ToJSON)
import Data.Foldable (traverse_)
import Data.Int (Int32, Int64)
import Data.List.NonEmpty qualified as NE
import Data.Maybe (listToMaybe)
import Data.Set qualified as Set
import Data.Text (Text)
import Data.Text qualified as T
import Data.UUID.Types qualified as UUID
import GHC.Generics (Generic)
import System.Timeout (timeout)
import Test.Hspec
import UnliftIO.Async (async, mapConcurrently, wait)

import Arbiter.Test.Setup (drainWith, execQuery, execStatement, seedConcurrencyPoolSQL)

-- | A payload declaring one concurrency pool.
newtype CLPayload = CLPayload Text
  deriving stock (CLPayload -> CLPayload -> Bool
(CLPayload -> CLPayload -> Bool)
-> (CLPayload -> CLPayload -> Bool) -> Eq CLPayload
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: CLPayload -> CLPayload -> Bool
== :: CLPayload -> CLPayload -> Bool
$c/= :: CLPayload -> CLPayload -> Bool
/= :: CLPayload -> CLPayload -> Bool
Eq, (forall x. CLPayload -> Rep CLPayload x)
-> (forall x. Rep CLPayload x -> CLPayload) -> Generic CLPayload
forall x. Rep CLPayload x -> CLPayload
forall x. CLPayload -> Rep CLPayload x
forall a.
(forall x. a -> Rep a x) -> (forall x. Rep a x -> a) -> Generic a
$cfrom :: forall x. CLPayload -> Rep CLPayload x
from :: forall x. CLPayload -> Rep CLPayload x
$cto :: forall x. Rep CLPayload x -> CLPayload
to :: forall x. Rep CLPayload x -> CLPayload
Generic, Int -> CLPayload -> ShowS
[CLPayload] -> ShowS
CLPayload -> String
(Int -> CLPayload -> ShowS)
-> (CLPayload -> String)
-> ([CLPayload] -> ShowS)
-> Show CLPayload
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> CLPayload -> ShowS
showsPrec :: Int -> CLPayload -> ShowS
$cshow :: CLPayload -> String
show :: CLPayload -> String
$cshowList :: [CLPayload] -> ShowS
showList :: [CLPayload] -> ShowS
Show)
  deriving anyclass (Maybe CLPayload
Value -> Parser [CLPayload]
Value -> Parser CLPayload
(Value -> Parser CLPayload)
-> (Value -> Parser [CLPayload])
-> Maybe CLPayload
-> FromJSON CLPayload
forall a.
(Value -> Parser a)
-> (Value -> Parser [a]) -> Maybe a -> FromJSON a
$cparseJSON :: Value -> Parser CLPayload
parseJSON :: Value -> Parser CLPayload
$cparseJSONList :: Value -> Parser [CLPayload]
parseJSONList :: Value -> Parser [CLPayload]
$comittedField :: Maybe CLPayload
omittedField :: Maybe CLPayload
FromJSON, [CLPayload] -> Value
[CLPayload] -> Encoding
CLPayload -> Bool
CLPayload -> Value
CLPayload -> Encoding
(CLPayload -> Value)
-> (CLPayload -> Encoding)
-> ([CLPayload] -> Value)
-> ([CLPayload] -> Encoding)
-> (CLPayload -> Bool)
-> ToJSON CLPayload
forall a.
(a -> Value)
-> (a -> Encoding)
-> ([a] -> Value)
-> ([a] -> Encoding)
-> (a -> Bool)
-> ToJSON a
$ctoJSON :: CLPayload -> Value
toJSON :: CLPayload -> Value
$ctoEncoding :: CLPayload -> Encoding
toEncoding :: CLPayload -> Encoding
$ctoJSONList :: [CLPayload] -> Value
toJSONList :: [CLPayload] -> Value
$ctoEncodingList :: [CLPayload] -> Encoding
toEncodingList :: [CLPayload] -> Encoding
$comitField :: CLPayload -> Bool
omitField :: CLPayload -> Bool
ToJSON)

-- | A finite tag per payload. "mx", "my", and "mz" select those pools. Any other
-- text selects "declpool", keyed by the payload text.
data CLTag = CLDecl | CLMx | CLMy | CLMz
  deriving stock (CLTag
CLTag -> CLTag -> Bounded CLTag
forall a. a -> a -> Bounded a
$cminBound :: CLTag
minBound :: CLTag
$cmaxBound :: CLTag
maxBound :: CLTag
Bounded, Int -> CLTag
CLTag -> Int
CLTag -> [CLTag]
CLTag -> CLTag
CLTag -> CLTag -> [CLTag]
CLTag -> CLTag -> CLTag -> [CLTag]
(CLTag -> CLTag)
-> (CLTag -> CLTag)
-> (Int -> CLTag)
-> (CLTag -> Int)
-> (CLTag -> [CLTag])
-> (CLTag -> CLTag -> [CLTag])
-> (CLTag -> CLTag -> [CLTag])
-> (CLTag -> CLTag -> CLTag -> [CLTag])
-> Enum CLTag
forall a.
(a -> a)
-> (a -> a)
-> (Int -> a)
-> (a -> Int)
-> (a -> [a])
-> (a -> a -> [a])
-> (a -> a -> [a])
-> (a -> a -> a -> [a])
-> Enum a
$csucc :: CLTag -> CLTag
succ :: CLTag -> CLTag
$cpred :: CLTag -> CLTag
pred :: CLTag -> CLTag
$ctoEnum :: Int -> CLTag
toEnum :: Int -> CLTag
$cfromEnum :: CLTag -> Int
fromEnum :: CLTag -> Int
$cenumFrom :: CLTag -> [CLTag]
enumFrom :: CLTag -> [CLTag]
$cenumFromThen :: CLTag -> CLTag -> [CLTag]
enumFromThen :: CLTag -> CLTag -> [CLTag]
$cenumFromTo :: CLTag -> CLTag -> [CLTag]
enumFromTo :: CLTag -> CLTag -> [CLTag]
$cenumFromThenTo :: CLTag -> CLTag -> CLTag -> [CLTag]
enumFromThenTo :: CLTag -> CLTag -> CLTag -> [CLTag]
Enum, CLTag -> CLTag -> Bool
(CLTag -> CLTag -> Bool) -> (CLTag -> CLTag -> Bool) -> Eq CLTag
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: CLTag -> CLTag -> Bool
== :: CLTag -> CLTag -> Bool
$c/= :: CLTag -> CLTag -> Bool
/= :: CLTag -> CLTag -> Bool
Eq)

instance HasConcurrency CLPayload where
  concurrencyFor :: ConcurrencyFor CLPayload
concurrencyFor = (CLPayload -> CLTag)
-> (CLTag -> ConcurrencyFor CLPayload) -> ConcurrencyFor CLPayload
forall k payload.
(Bounded k, Enum k, Eq k) =>
(payload -> k)
-> (k -> ConcurrencyFor payload) -> ConcurrencyFor payload
concurrencyByCase CLPayload -> CLTag
tagOf CLTag -> ConcurrencyFor CLPayload
sel
    where
      tagOf :: CLPayload -> CLTag
tagOf (CLPayload Text
text)
        | Text
text Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
"mx" = CLTag
CLMx
        | Text
text Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
"my" = CLTag
CLMy
        | Text
text Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
"mz" = CLTag
CLMz
        | Bool
otherwise = CLTag
CLDecl
      sel :: CLTag -> ConcurrencyFor CLPayload
sel CLTag
CLDecl = ConcurrencyPolicy
-> (CLPayload -> Text) -> ConcurrencyFor CLPayload
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"declpool" Int32
2) (\(CLPayload Text
text) -> Text
text)
      sel CLTag
CLMx = ConcurrencyPolicy
-> (CLPayload -> Text) -> ConcurrencyFor CLPayload
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"mx" Int32
1) (Text -> CLPayload -> Text
forall a b. a -> b -> a
const Text
"k")
      sel CLTag
CLMy = ConcurrencyPolicy
-> (CLPayload -> Text) -> ConcurrencyFor CLPayload
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"my" Int32
2) (Text -> CLPayload -> Text
forall a b. a -> b -> a
const Text
"k")
      sel CLTag
CLMz = ConcurrencyPolicy
-> (CLPayload -> Text) -> ConcurrencyFor CLPayload
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"mz" Int32
3) (Text -> CLPayload -> Text
forall a b. a -> b -> a
const Text
"k")

-- | A one-queue registry over 'CLPayload'.
type CLReg = '[Queue "arbiter_concurrency_test" CLPayload]

-- | A second payload declaring a different pool, with a two-payload registry.
newtype CLPayload2 = CLPayload2 Text
  deriving stock (CLPayload2 -> CLPayload2 -> Bool
(CLPayload2 -> CLPayload2 -> Bool)
-> (CLPayload2 -> CLPayload2 -> Bool) -> Eq CLPayload2
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: CLPayload2 -> CLPayload2 -> Bool
== :: CLPayload2 -> CLPayload2 -> Bool
$c/= :: CLPayload2 -> CLPayload2 -> Bool
/= :: CLPayload2 -> CLPayload2 -> Bool
Eq, (forall x. CLPayload2 -> Rep CLPayload2 x)
-> (forall x. Rep CLPayload2 x -> CLPayload2) -> Generic CLPayload2
forall x. Rep CLPayload2 x -> CLPayload2
forall x. CLPayload2 -> Rep CLPayload2 x
forall a.
(forall x. a -> Rep a x) -> (forall x. Rep a x -> a) -> Generic a
$cfrom :: forall x. CLPayload2 -> Rep CLPayload2 x
from :: forall x. CLPayload2 -> Rep CLPayload2 x
$cto :: forall x. Rep CLPayload2 x -> CLPayload2
to :: forall x. Rep CLPayload2 x -> CLPayload2
Generic, Int -> CLPayload2 -> ShowS
[CLPayload2] -> ShowS
CLPayload2 -> String
(Int -> CLPayload2 -> ShowS)
-> (CLPayload2 -> String)
-> ([CLPayload2] -> ShowS)
-> Show CLPayload2
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> CLPayload2 -> ShowS
showsPrec :: Int -> CLPayload2 -> ShowS
$cshow :: CLPayload2 -> String
show :: CLPayload2 -> String
$cshowList :: [CLPayload2] -> ShowS
showList :: [CLPayload2] -> ShowS
Show)
  deriving anyclass (Maybe CLPayload2
Value -> Parser [CLPayload2]
Value -> Parser CLPayload2
(Value -> Parser CLPayload2)
-> (Value -> Parser [CLPayload2])
-> Maybe CLPayload2
-> FromJSON CLPayload2
forall a.
(Value -> Parser a)
-> (Value -> Parser [a]) -> Maybe a -> FromJSON a
$cparseJSON :: Value -> Parser CLPayload2
parseJSON :: Value -> Parser CLPayload2
$cparseJSONList :: Value -> Parser [CLPayload2]
parseJSONList :: Value -> Parser [CLPayload2]
$comittedField :: Maybe CLPayload2
omittedField :: Maybe CLPayload2
FromJSON, [CLPayload2] -> Value
[CLPayload2] -> Encoding
CLPayload2 -> Bool
CLPayload2 -> Value
CLPayload2 -> Encoding
(CLPayload2 -> Value)
-> (CLPayload2 -> Encoding)
-> ([CLPayload2] -> Value)
-> ([CLPayload2] -> Encoding)
-> (CLPayload2 -> Bool)
-> ToJSON CLPayload2
forall a.
(a -> Value)
-> (a -> Encoding)
-> ([a] -> Value)
-> ([a] -> Encoding)
-> (a -> Bool)
-> ToJSON a
$ctoJSON :: CLPayload2 -> Value
toJSON :: CLPayload2 -> Value
$ctoEncoding :: CLPayload2 -> Encoding
toEncoding :: CLPayload2 -> Encoding
$ctoJSONList :: [CLPayload2] -> Value
toJSONList :: [CLPayload2] -> Value
$ctoEncodingList :: [CLPayload2] -> Encoding
toEncodingList :: [CLPayload2] -> Encoding
$comitField :: CLPayload2 -> Bool
omitField :: CLPayload2 -> Bool
ToJSON)

instance HasConcurrency CLPayload2 where
  concurrencyFor :: ConcurrencyFor CLPayload2
concurrencyFor = ConcurrencyPolicy
-> (CLPayload2 -> Text) -> ConcurrencyFor CLPayload2
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"declpool2" Int32
5) (\(CLPayload2 Text
text) -> Text
text)

type CLReg2 = '[Queue "clq1" CLPayload, Queue "clq2" CLPayload2]

-- | Table name for 'CLReg', shared across backends.
concurrencyTable :: Text
concurrencyTable :: Text
concurrencyTable = Text
"arbiter_concurrency_test"

fullKey :: Text -> Text -> Text
fullKey :: Text -> Text -> Text
fullKey Text
pool Text
suffix = Text
"declpool:" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
pool Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
":" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
suffix

job :: Text -> Text -> JobWrite CLPayload
job :: Text -> Text -> JobWrite CLPayload
job Text
pool Text
suffix = CLPayload -> JobWrite CLPayload
forall payload. payload -> JobWrite payload
defaultJob (Text -> CLPayload
CLPayload (Text
pool Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
":" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
suffix))

groupedJob :: Text -> Text -> Text -> JobWrite CLPayload
groupedJob :: Text -> Text -> Text -> JobWrite CLPayload
groupedJob Text
groupKey Text
pool Text
suffix = Text -> CLPayload -> JobWrite CLPayload
forall payload. Text -> payload -> JobWrite payload
defaultGroupedJob Text
groupKey (Text -> CLPayload
CLPayload (Text
pool Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
":" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
suffix))

-- | The concurrency-limit suite, run against any backend.
concurrencyLimitSpec
  :: forall env m
   . (HasRegistry m CLReg)
  => (forall a. env -> m a -> IO a)
  -> SpecWith env
concurrencyLimitSpec :: forall env (m :: * -> *).
HasRegistry m CLReg =>
(forall a. env -> m a -> IO a) -> SpecWith env
concurrencyLimitSpec forall a. env -> m a -> IO a
runM = do
  let wid :: UUID
wid = UUID
UUID.nil
      tshow :: Int -> Text
tshow = String -> Text
T.pack (String -> Text) -> (Int -> String) -> Int -> Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Int -> String
forall a. Show a => a -> String
show :: Int -> Text
      enqueue :: env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [JobWrite CLPayload]
jobs = IO [JobRead CLPayload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (env -> m [JobRead CLPayload] -> IO [JobRead CLPayload]
forall a. env -> m a -> IO a
runM env
env ([JobWrite CLPayload] -> m [JobRead CLPayload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch [JobWrite CLPayload]
jobs) :: IO [JobRead CLPayload])
      -- Claims are attributed.
      claimAs :: env -> IO [JobRead CLPayload]
claimAs env
env = env -> m [JobRead CLPayload] -> IO [JobRead CLPayload]
forall a. env -> m a -> IO a
runM env
env (Int -> NominalDiffTime -> UUID -> m [JobRead CLPayload]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> UUID -> m [JobRead payload]
HL.claimNextVisibleJobsAs Int
100 NominalDiffTime
60 UUID
wid) :: IO [JobRead CLPayload]
      ackAll :: env -> [JobRead CLPayload] -> IO ()
ackAll env
env [JobRead CLPayload]
jobs = env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env ((JobRead CLPayload -> m Int64) -> [JobRead CLPayload] -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ JobRead CLPayload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob [JobRead CLPayload]
jobs)
      retryAll :: env -> [JobRead CLPayload] -> IO ()
retryAll env
env [JobRead CLPayload]
jobs = env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env ((JobRead CLPayload -> m Int64) -> [JobRead CLPayload] -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (NominalDiffTime -> Text -> JobRead CLPayload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
NominalDiffTime -> Text -> JobRead payload -> m Int64
HL.updateJobForRetry NominalDiffTime
60 Text
"boom") [JobRead CLPayload]
jobs)
      nackAll :: env -> [JobRead CLPayload] -> IO ()
nackAll env
env [JobRead CLPayload]
jobs = env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env ((JobRead CLPayload -> m Int64) -> [JobRead CLPayload] -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ JobRead CLPayload -> m Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.nackJob [JobRead CLPayload]
jobs)
      overridePool :: env -> Maybe Int32 -> IO ()
overridePool env
env Maybe Int32
lim = 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 (Text -> ConcurrencyPolicyUpdate -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> ConcurrencyPolicyUpdate -> m Int64
HL.updateConcurrencyPolicyOverrides Text
"declpool" (Maybe (Maybe Int32) -> ConcurrencyPolicyUpdate
ConcurrencyPolicyUpdate (Maybe Int32 -> Maybe (Maybe Int32)
forall a. a -> Maybe a
Just Maybe Int32
lim))) :: IO Int64)
      prune :: env -> IO Int64
prune env
env = env -> m Int64 -> IO Int64
forall a. env -> m a -> IO a
runM env
env m Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m Int64
HL.pruneConcurrencyKeys :: IO Int64
      reconcile :: env -> IO ()
reconcile env
env = 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 m Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m Int64
HL.reconcileConcurrencyCounts :: IO Int64)
      reconcileIfStale :: env -> IO ()
reconcileIfStale env
env = 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 m Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m Int64
HL.reconcileConcurrencyCountsIfStale) :: IO ()
      -- Empty the UNLOGGED count table the way a Postgres crash recovery would.
      truncateCounts :: env -> IO ()
truncateCounts env
env =
        env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env (m () -> IO ()) -> m () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          void $ execStatement ("DELETE FROM " <> arbiterConcurrencyTable schema) []
      -- Seed the declpool default limit and clear any override.
      seed :: env -> Int -> IO ()
seed env
env (Int
lim :: Int) =
        env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env (m () -> IO ()) -> m () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          traverse_
            (\Text
statement -> m Int64 -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (Text -> Params -> m Int64
forall (m :: * -> *). MonadArbiter m => Text -> Params -> m Int64
execStatement Text
statement []))
            (seedConcurrencyPoolSQL schema "declpool" (fromIntegral lim))
      corrupt :: env -> Text -> Int -> IO ()
corrupt env
env Text
key Int
count =
        env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env (m () -> IO ()) -> m () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          void $
            execStatement
              ("UPDATE " <> arbiterConcurrencyTable schema <> " SET in_flight = " <> tshow count <> " WHERE concurrency_key = ?")
              [pval CText key]
      -- Mark a grouped job as already failed.
      markAttempted :: env -> Text -> IO ()
markAttempted env
env Text
key =
        env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env (m () -> IO ()) -> m () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          void $
            execStatement
              ( "UPDATE "
                  <> jobQueueTable schema concurrencyTable
                  <> " SET attempts = 1 WHERE concurrency_key = ? AND group_key IS NOT NULL"
              )
              [pval CText key]
      deleteRow :: env -> Text -> IO ()
deleteRow env
env Text
key =
        env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env (m () -> IO ()) -> m () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
          schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          void $
            execStatement
              ("DELETE FROM " <> arbiterConcurrencyTable schema <> " WHERE concurrency_key = ?")
              [pval CText key]
      inFlight :: env -> Text -> IO (Maybe Int32)
inFlight env
env Text
key =
        env -> m (Maybe Int32) -> IO (Maybe Int32)
forall a. env -> m a -> IO a
runM env
env (m (Maybe Int32) -> IO (Maybe Int32))
-> m (Maybe Int32) -> IO (Maybe Int32)
forall a b. (a -> b) -> a -> b
$ do
          schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          rows <-
            execQuery
              ("SELECT in_flight FROM " <> arbiterConcurrencyTable schema <> " WHERE concurrency_key = ?")
              [pval CText key]
              (col "in_flight" CInt4)
          pure (listToMaybe rows :: Maybe Int32)
      -- Backdate every job's visibility. Claimed jobs time out and become reclaimable.
      timeOut :: env -> Int -> IO ()
timeOut env
env Int
secs = env -> m () -> IO ()
forall a. env -> m a -> IO a
runM env
env (m () -> IO ()) -> m () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
        schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
        void $
          execStatement
            ( "UPDATE "
                <> jobQueueTable schema concurrencyTable
                <> " SET not_visible_until = not_visible_until - "
                <> tshow secs
                <> " * interval '1 second'"
            )
            []

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"admits up to the pool limit and blocks the rest" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"cap" Text
"a"))
    first <- env -> IO [JobRead CLPayload]
claimAs env
env
    length first `shouldBe` 3
    again <- claimAs env
    length again `shouldBe` 0

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"an undeclared pool runs uncapped (fail open)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- Claim well above any seeded limit.
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
20 (Text -> Text -> JobWrite CLPayload
job Text
"undeclared" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 20

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"caps by a HasConcurrency instance's declared pool (no manual key)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"declpool" Text
"tx"))
    first <- env -> IO [JobRead CLPayload]
claimAs env
env
    length first `shouldBe` 2
    inFlight env (fullKey "declpool" "tx") `shouldReturn` Just 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"the registry collects the declared pool for migration seeding" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
_ ->
    forall (registry :: JobPayloadRegistry).
RegistryConcurrencyPolicies registry =>
Set ConcurrencyPolicy
registryConcurrencyPolicies @CLReg
      Set ConcurrencyPolicy -> Set ConcurrencyPolicy -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [ConcurrencyPolicy] -> Set ConcurrencyPolicy
forall a. Ord a => [a] -> Set a
Set.fromList
        [ Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"declpool" Int32
2
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"mx" Int32
1
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"my" Int32
2
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"mz" Int32
3
        ]

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"the registry unions declared pools across all payloads" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
_ ->
    -- Two payloads, two pools.
    forall (registry :: JobPayloadRegistry).
RegistryConcurrencyPolicies registry =>
Set ConcurrencyPolicy
registryConcurrencyPolicies @CLReg2
      Set ConcurrencyPolicy -> Set ConcurrencyPolicy -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` [ConcurrencyPolicy] -> Set ConcurrencyPolicy
forall a. Ord a => [a] -> Set a
Set.fromList
        [ Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"declpool" Int32
2
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"mx" Int32
1
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"my" Int32
2
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"mz" Int32
3
        , Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
"declpool2" Int32
5
        ]

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrencyPool floors a non-positive limit at 1 (the policies table requires a positive default)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
_ -> do
    ConcurrencyPolicy -> Int32
cpLimit (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"p" Int32
0) Int32 -> Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Int32
1
    ConcurrencyPolicy -> Int32
cpLimit (Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"p" (-Int32
5)) Int32 -> Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Int32
1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"chooseWhen collects both concurrency branches and runs the chosen one" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
_ -> do
    let policyA :: ConcurrencyPolicy
policyA = Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"ca" Int32
1
        policyB :: ConcurrencyPolicy
policyB = Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
"cb" Int32
2
        sel :: ConcurrencyFor Bool
        sel :: ConcurrencyFor Bool
sel = (Bool -> Bool)
-> ConcurrencyFor Bool
-> ConcurrencyFor Bool
-> ConcurrencyFor Bool
forall payload policy a.
(payload -> Bool)
-> Selector policy payload a
-> Selector policy payload a
-> Selector policy payload a
chooseWhen Bool -> Bool
forall a. a -> a
id (ConcurrencyPolicy -> (Bool -> Text) -> ConcurrencyFor Bool
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy ConcurrencyPolicy
policyA (Text -> Bool -> Text
forall a b. a -> b -> a
const Text
"x")) (ConcurrencyPolicy -> (Bool -> Text) -> ConcurrencyFor Bool
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy ConcurrencyPolicy
policyB (Text -> Bool -> Text
forall a b. a -> b -> a
const Text
"y"))
    Set ConcurrencyPolicy -> [ConcurrencyPolicy]
forall a. Set a -> [a]
Set.toList (ConcurrencyFor Bool -> Set ConcurrencyPolicy
forall policy payload a.
Ord policy =>
Selector policy payload a -> Set policy
collectPolicies ConcurrencyFor Bool
sel) [ConcurrencyPolicy] -> [ConcurrencyPolicy] -> IO ()
forall a. (HasCallStack, Show a, Eq a) => [a] -> [a] -> IO ()
`shouldMatchList` [ConcurrencyPolicy
policyA, ConcurrencyPolicy
policyB]
    (ConcurrencyKey -> Text
ckPrefix (ConcurrencyKey -> Text) -> Maybe ConcurrencyKey -> Maybe Text
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Bool -> ConcurrencyFor Bool -> Maybe ConcurrencyKey
forall payload.
payload -> ConcurrencyFor payload -> Maybe ConcurrencyKey
runConcurrencyFor Bool
True ConcurrencyFor Bool
sel) Maybe Text -> Maybe Text -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"ca"
    (ConcurrencyKey -> Text
ckPrefix (ConcurrencyKey -> Text) -> Maybe ConcurrencyKey -> Maybe Text
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Bool -> ConcurrencyFor Bool -> Maybe ConcurrencyKey
forall payload.
payload -> ConcurrencyFor payload -> Maybe ConcurrencyKey
runConcurrencyFor Bool
False ConcurrencyFor Bool
sel) Maybe Text -> Maybe Text -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Text -> Maybe Text
forall a. a -> Maybe a
Just Text
"cb"

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"frees a slot on ack" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"freeack" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    ackAll env (take 1 claimed)
    refilled <- claimAs env
    length refilled `shouldBe` 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"frees a slot when a job goes back for retry" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"freeretry" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    retryAll env (take 2 claimed)
    refilled <- claimAs env
    length refilled `shouldBe` 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"frees a slot on nack" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"freenack" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    nackAll env (take 2 claimed)
    refilled <- claimAs env
    length refilled `shouldBe` 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"reclaims a timed-out job without exceeding the limit" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"reclaim" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    timeOut env 120
    reclaimed <- claimAs env
    length reclaimed `shouldBe` 3

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"keeps separate keys under one pool independent" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"iso" Text
"x") [JobWrite CLPayload]
-> [JobWrite CLPayload] -> [JobWrite CLPayload]
forall a. Semigroup a => a -> a -> a
<> Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"iso" Text
"y"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 4

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a pool override lowers the cap live" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
5 (Text -> Text -> JobWrite CLPayload
job Text
"ovlo" Text
"a"))
    env -> Maybe Int32 -> IO ()
overridePool env
env (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
1)
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a pool override raises the cap live, across every key under the prefix" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
6 (Text -> Text -> JobWrite CLPayload
job Text
"ovhi" Text
"x") [JobWrite CLPayload]
-> [JobWrite CLPayload] -> [JobWrite CLPayload]
forall a. Semigroup a => a -> a -> a
<> Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
6 (Text -> Text -> JobWrite CLPayload
job Text
"ovhi" Text
"y"))
    env -> Maybe Int32 -> IO ()
overridePool env
env (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
5)
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    -- Both keys under the pool admit up to 5.
    length claimed `shouldBe` 10

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"an override of 0 pauses the pool until cleared" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
3 (Text -> Text -> JobWrite CLPayload
job Text
"pause" Text
"a"))
    env -> Maybe Int32 -> IO ()
overridePool env
env (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
0)
    paused <- env -> IO [JobRead CLPayload]
claimAs env
env
    length paused `shouldBe` 0
    overridePool env Nothing
    resumed <- claimAs env
    length resumed `shouldBe` 3

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"an empty policy patch leaves the override unchanged, an explicit null clears it" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
6 (Text -> Text -> JobWrite CLPayload
job Text
"patch" Text
"a"))
    env -> Maybe Int32 -> IO ()
overridePool env
env (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
1)
    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 (Text -> ConcurrencyPolicyUpdate -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> ConcurrencyPolicyUpdate -> m Int64
HL.updateConcurrencyPolicyOverrides Text
"declpool" (Maybe (Maybe Int32) -> ConcurrencyPolicyUpdate
ConcurrencyPolicyUpdate Maybe (Maybe Int32)
forall a. Maybe a
Nothing)) :: IO Int64)
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 1
    overridePool env Nothing
    more <- claimAs env
    length more `shouldBe` 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"in_flight tracks claims and acks" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
4
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
4 (Text -> Text -> JobWrite CLPayload
job Text
"track" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 4
    inFlight env (fullKey "track" "a") `shouldReturn` Just 4
    ackAll env (take 3 claimed)
    inFlight env (fullKey "track" "a") `shouldReturn` Just 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"floors in_flight at zero when a decrement would underflow a drifted count" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- Drift can leave a count row below its live claimed jobs. Acking one clamps
    -- the count at 0.
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"floor" Text
"a"]
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 1
    -- One claimed job. Force the count to read 0.
    corrupt env (fullKey "floor" "a") 0
    ackAll env claimed
    inFlight env (fullKey "floor" "a") `shouldReturn` Just 0

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"prunes drained key rows" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
2 (Text -> Text -> JobWrite CLPayload
job Text
"drain" Text
"x") [JobWrite CLPayload]
-> [JobWrite CLPayload] -> [JobWrite CLPayload]
forall a. Semigroup a => a -> a -> a
<> Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
2 (Text -> Text -> JobWrite CLPayload
job Text
"drain" Text
"y"))
    env -> IO [JobRead CLPayload]
claimAs env
env IO [JobRead CLPayload] -> ([JobRead CLPayload] -> IO ()) -> IO ()
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= env -> [JobRead CLPayload] -> IO ()
ackAll env
env
    env -> IO Int64
prune env
env IO Int64 -> Int64 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int64
2
    env -> Text -> IO (Maybe Int32)
inFlight env
env (Text -> Text -> Text
fullKey Text
"drain" Text
"x") IO (Maybe Int32) -> Maybe Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Maybe Int32
forall a. Maybe a
Nothing
    env -> Text -> IO (Maybe Int32)
inFlight env
env (Text -> Text -> Text
fullKey Text
"drain" Text
"y") IO (Maybe Int32) -> Maybe Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Maybe Int32
forall a. Maybe a
Nothing

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"prune skips a key held by an open enqueue transaction and the job still runs" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"skip" Text
"a"]
    env -> IO [JobRead CLPayload]
claimAs env
env IO [JobRead CLPayload] -> ([JobRead CLPayload] -> IO ()) -> IO ()
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= env -> [JobRead CLPayload] -> IO ()
ackAll env
env
    env -> Text -> IO (Maybe Int32)
inFlight env
env (Text -> Text -> Text
fullKey Text
"skip" Text
"a") IO (Maybe Int32) -> Maybe Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
0
    entered <- IO (MVar ())
forall a. IO (MVar a)
newEmptyMVar
    release <- newEmptyMVar
    enqueuer <- async $ runM env $ withDbTransaction $ do
      _ <- HL.insertJobsBatch [job "skip" "a"] :: m [JobRead CLPayload]
      liftIO $ putMVar entered ()
      liftIO $ takeMVar release
    takeMVar entered
    -- The open enqueue holds the key's shared advisory lock. Prune returns without it.
    pruned <- timeout 5000000 (prune env)
    pruned `shouldBe` Just 0
    inFlight env (fullKey "skip" "a") `shouldReturn` Just 0
    putMVar release ()
    wait enqueuer
    claimed <- claimAs env
    length claimed `shouldBe` 1
    ackAll env claimed
    prune env `shouldReturn` 1
    inFlight env (fullKey "skip" "a") `shouldReturn` Nothing

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"prune racing a multi-statement descending-key enqueue transaction does not deadlock" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"dk" Text
"a", Text -> Text -> JobWrite CLPayload
job Text
"dk" Text
"b"]
    env -> IO [JobRead CLPayload]
claimAs env
env IO [JobRead CLPayload] -> ([JobRead CLPayload] -> IO ()) -> IO ()
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= env -> [JobRead CLPayload] -> IO ()
ackAll env
env
    entered <- IO (MVar ())
forall a. IO (MVar a)
newEmptyMVar
    pruneDone <- newEmptyMVar
    enqueuer <- async $ runM env $ withDbTransaction $ do
      _ <- HL.insertJobsBatch [job "dk" "b"] :: m [JobRead CLPayload]
      liftIO $ putMVar entered ()
      liftIO $ takeMVar pruneDone
      _ <- HL.insertJobsBatch [job "dk" "a"] :: m [JobRead CLPayload]
      pure ()
    takeMVar entered
    -- Prune sees both keys dead. It takes a and skips the held b.
    pruned <- timeout 5000000 (prune env)
    pruned `shouldBe` Just 1
    inFlight env (fullKey "dk" "a") `shouldReturn` Nothing
    putMVar pruneDone ()
    wait enqueuer
    -- The second insert reseeds the pruned key. Both jobs stay claimable.
    claimed <- claimAs env
    length claimed `shouldBe` 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"preserves the concurrency cap through a DLQ round-trip" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
1
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"dlq" Text
"a"]
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    runM env (traverse_ (void . HL.moveToDLQ "boom") claimed)
    dlqs <- runM env (HL.listDLQJobs 10 0) :: IO [DLQ.DLQJob CLPayload]
    runM
      env
      (traverse_ (\DLQJob CLPayload
dlqJob -> m (Maybe (JobRead CLPayload)) -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (Int64 -> m (Maybe (JobRead CLPayload))
forall payload (m :: * -> *).
QueueOperation m payload =>
Int64 -> m (Maybe (JobRead payload))
HL.retryFromDLQ (DLQJob CLPayload -> Int64
forall payload. DLQJob payload -> Int64
DLQ.dlqPrimaryKey DLQJob CLPayload
dlqJob) :: m (Maybe (JobRead CLPayload)))) dlqs)
    reclaimed <- claimAs env
    length reclaimed `shouldBe` 1
    inFlight env (fullKey "dlq" "a") `shouldReturn` Just 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"reconcile repairs a drifted in_flight count" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
5
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
3 (Text -> Text -> JobWrite CLPayload
job Text
"recon" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    corrupt env (fullKey "recon" "a") 99
    inFlight env (fullKey "recon" "a") `shouldReturn` Just 99
    reconcile env
    inFlight env (fullKey "recon" "a") `shouldReturn` Just 3

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"reconcile does not crash on a drained-but-unpruned key" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
3
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
2 (Text -> Text -> JobWrite CLPayload
job Text
"draincon" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 2
    -- Ack deletes the jobs. The count row drains and stays unpruned.
    ackAll env claimed
    -- Reconcile repairs a key with no live jobs.
    reconcile env
    inFlight env (fullKey "draincon" "a") `shouldReturn` Just 0

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"reconcile recreates a missing count row from live jobs" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
5
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
3 (Text -> Text -> JobWrite CLPayload
job Text
"reconm" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    deleteRow env (fullKey "reconm" "a")
    inFlight env (fullKey "reconm" "a") `shouldReturn` Nothing
    reconcile env
    inFlight env (fullKey "reconm" "a") `shouldReturn` Just 3

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"reconcile against an in-flight claim does not overwrite the count below the truth" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- Reconcile runs while the second job is claimed in an open transaction. It counts that job.
    env -> Int -> IO ()
seed env
env Int
5
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"recrace" Text
"a"]
    firstClaim <- env -> IO [JobRead CLPayload]
claimAs env
env
    length firstClaim `shouldBe` 1
    enqueue env [job "recrace" "a"]
    claimed <- newEmptyMVar
    release <- newEmptyMVar
    claimer <- async $ runM env $ withDbTransaction $ do
      claimedJobs <- HL.claimNextVisibleJobsAs 100 60 wid :: m [JobRead CLPayload]
      liftIO $ putMVar claimed (length claimedJobs)
      liftIO $ takeMVar release
    takeMVar claimed `shouldReturn` 1
    reconciler <- async (reconcile env)
    threadDelay 200000
    putMVar release ()
    wait claimer
    wait reconciler
    inFlight env (fullKey "recrace" "a") `shouldReturn` Just 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"reconcile does not overwrite a claim on a key seeded after its lock pass" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- The key is born between the lock pass and the recount. The recount leaves
    -- it to its triggers.
    env -> Int -> IO ()
seed env
env Int
5
    lockTaken <- IO (MVar ())
forall a. IO (MVar a)
newEmptyMVar
    resume <- newEmptyMVar
    reconciler <- async $ runM env $ withDbTransaction $ do
      schema <- getSchema
      held <- MA.executeQuery (Tmpl.lockConcurrencyCountsSQL schema)
      liftIO $ putMVar lockTaken ()
      liftIO $ takeMVar resume
      void $
        MA.executeQuery
          (Tmpl.reconcileConcurrencyCountsSQL schema [concurrencyTable] held)
    takeMVar lockTaken
    enqueue env [job "lateseed" "a"]
    claimed <- newEmptyMVar
    release <- newEmptyMVar
    claimer <- async $ runM env $ withDbTransaction $ do
      claimedJobs <- HL.claimNextVisibleJobsAs 100 60 wid :: m [JobRead CLPayload]
      liftIO $ putMVar claimed (length claimedJobs)
      liftIO $ takeMVar release
    takeMVar claimed `shouldReturn` 1
    putMVar resume ()
    threadDelay 200000
    putMVar release ()
    wait claimer
    wait reconciler
    inFlight env (fullKey "lateseed" "a") `shouldReturn` Just 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a stale rebuild restores counts a crash truncated, so keyed jobs stay claimable" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
4 (Text -> Text -> JobWrite CLPayload
job Text
"stale" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 2
    -- A crash truncates the UNLOGGED count table. Keyed jobs are now unclaimable.
    truncateCounts env
    inFlight env (fullKey "stale" "a") `shouldReturn` Nothing
    blocked <- claimAs env
    length blocked `shouldBe` 0
    -- The startup stale rebuild repopulates the count from the live claimed jobs.
    reconcileIfStale env
    inFlight env (fullKey "stale" "a") `shouldReturn` Just 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a stale rebuild still fires when a post-crash enqueue re-seeds one key" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
4 (Text -> Text -> JobWrite CLPayload
job Text
"stale" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 2
    truncateCounts env
    -- A fresh enqueue seeds its own count row. The truncation is still detected.
    enqueue env [job "stale" "b"]
    reconcileIfStale env
    inFlight env (fullKey "stale" "a") `shouldReturn` Just 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a stale rebuild is a no-op on a healthy count table (does not clobber drift it should not touch)" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
5
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
3 (Text -> Text -> JobWrite CLPayload
job Text
"healthy" Text
"a"))
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 3
    -- The table is non-empty. The stale check skips the full reconcile.
    corrupt env (fullKey "healthy" "a") 99
    reconcileIfStale env
    inFlight env (fullKey "healthy" "a") `shouldReturn` Just 99

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"moves the count row when a dedup replace changes the concurrency key" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
1
    let mkJob :: Text -> JobWrite CLPayload
mkJob Text
suffix = Maybe DedupKey -> JobWrite CLPayload -> JobWrite CLPayload
forall payload.
Maybe DedupKey -> JobWrite payload -> JobWrite payload
setDedupKey (DedupKey -> Maybe DedupKey
forall a. a -> Maybe a
Just (Text -> DedupKey
ReplaceDuplicate Text
"movekey")) (JobWrite CLPayload -> JobWrite CLPayload)
-> JobWrite CLPayload -> JobWrite CLPayload
forall a b. (a -> b) -> a -> b
$ Text -> Text -> JobWrite CLPayload
job Text
"ck" Text
suffix
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> JobWrite CLPayload
mkJob Text
"old"]
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> JobWrite CLPayload
mkJob Text
"new"]
    env -> Text -> IO (Maybe Int32)
inFlight env
env (Text -> Text -> Text
fullKey Text
"ck" Text
"old") IO (Maybe Int32) -> Maybe Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
0
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 1
    inFlight env (fullKey "ck" "new") `shouldReturn` Just 1
    void (prune env)
    inFlight env (fullKey "ck" "old") `shouldReturn` Nothing

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"caps across groups sharing one key" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> Text -> JobWrite CLPayload
groupedJob (Text
"g-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
tshow Int
index) Text
"shared" Text
"c" | Int
index <- [Int
1 .. Int
5]]
    claimed <- env -> IO [JobRead CLPayload]
claimAs env
env
    length claimed `shouldBe` 2

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"stalls a concurrency-blocked grouped head over a free-key sibling in batched mode" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- With the head's key at its cap, a batched claim defers the whole group.
    env -> Int -> IO ()
seed env
env Int
1
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"bh" Text
"head"]
    filled <- env -> IO [JobRead CLPayload]
claimAs env
env
    length filled `shouldBe` 1
    enqueue env [groupedJob "bgrp" "bh" "head", groupedJob "bgrp" "bh" "sib"]
    stalled <- runM env (HL.claimNextVisibleJobsBatched 5 100 60) :: IO [NE.NonEmpty (JobRead CLPayload)]
    concatMap NE.toList stalled `shouldSatisfy` null

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrent claimers never exceed the limit, and the key fills to the cap" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    env -> Int -> IO ()
seed env
env Int
5
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
40 (Text -> Text -> JobWrite CLPayload
job Text
"race" Text
"a"))
    results <- (Int -> IO [JobRead CLPayload])
-> [Int] -> IO [[JobRead CLPayload]]
forall (m :: * -> *) (t :: * -> *) a b.
(MonadUnliftIO m, Traversable t) =>
(a -> m b) -> t a -> m (t b)
mapConcurrently (IO [JobRead CLPayload] -> Int -> IO [JobRead CLPayload]
forall a b. a -> b -> a
const (env -> IO [JobRead CLPayload]
claimAs env
env)) [Int
1 .. Int
8 :: Int]
    let total = [Int] -> Int
forall a. Num a => [a] -> a
forall (t :: * -> *) a. (Foldable t, Num a) => t a -> a
sum (([JobRead CLPayload] -> Int) -> [[JobRead CLPayload]] -> [Int]
forall a b. (a -> b) -> [a] -> [b]
map [JobRead CLPayload] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [[JobRead CLPayload]]
results)
    -- At least one claimer makes progress. None pushes the key over the cap.
    total `shouldSatisfy` (\Int
count -> Int
count Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
1 Bool -> Bool -> Bool
&& Int
count Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
<= Int
5)
    -- An uncontended follow-up fills the key to the cap. The contended total is
    -- not deterministic.
    more <- claimAs env
    (total + length more) `shouldBe` 5
    inFlight env (fullKey "race" "a") `shouldReturn` Just 5

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrent claimers across many keys never exceed any cap" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- Many distinct keys at once.
    env -> Int -> IO ()
seed env
env Int
2
    let keys :: [Text]
keys = [Int -> Text
tshow Int
index | Int
index <- [Int
1 .. Int
12]]
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env ([[JobWrite CLPayload]] -> [JobWrite CLPayload]
forall (t :: * -> *) a. Foldable t => t [a] -> [a]
concat [Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
6 (Text -> Text -> JobWrite CLPayload
job Text
"declpool" Text
suffix) | Text
suffix <- [Text]
keys])
    _ <- (Int -> IO [JobRead CLPayload])
-> [Int] -> IO [[JobRead CLPayload]]
forall (m :: * -> *) (t :: * -> *) a b.
(MonadUnliftIO m, Traversable t) =>
(a -> m b) -> t a -> m (t b)
mapConcurrently (IO [JobRead CLPayload] -> Int -> IO [JobRead CLPayload]
forall a b. a -> b -> a
const (env -> IO [JobRead CLPayload]
claimAs env
env)) [Int
1 .. Int
8 :: Int]
    overCap <- runM env $ do
      schema <- getSchema
      execQuery
        ( "SELECT concurrency_key FROM "
            <> jobQueueTable schema concurrencyTable
            <> " WHERE concurrency_key IS NOT NULL GROUP BY concurrency_key"
            <> " HAVING count(*) FILTER (WHERE claimed_by IS NOT NULL) > 2"
        )
        []
        (col "concurrency_key" CText)
    (overCap :: [Text]) `shouldBe` []
    -- An uncontended drain fills every key to the cap.
    void (drainWith (claimAs env))
    traverse_ (\Text
suffix -> env -> Text -> IO (Maybe Int32)
inFlight env
env (Text -> Text -> Text
fullKey Text
"declpool" Text
suffix) IO (Maybe Int32) -> Maybe Int32 -> IO ()
forall a. (HasCallStack, Show a, Eq a) => IO a -> a -> IO ()
`shouldReturn` Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
2) keys

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a full concurrency key does not starve admissible ungrouped jobs behind it" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- Fill the hot key to its cap, then flood it past the bounded candidate window.
    -- The cold job on another key still claims.
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
150 (Text -> Text -> JobWrite CLPayload
job Text
"declpool" Text
"hot"))
    filled <- env -> IO [JobRead CLPayload]
claimAs env
env
    length filled `shouldBe` 2
    enqueue env [job "declpool" "cold"]
    _ <- claimAs env
    inFlight env (fullKey "declpool" "cold") `shouldReturn` Just 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"a full concurrency key does not starve admissible grouped jobs behind it" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- Each blocked group takes a slot in the bounded window. A flood of groups on
    -- one full key leaves room for a cold group behind them.
    env -> Int -> IO ()
seed env
env Int
2
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> Text -> JobWrite CLPayload
groupedJob (Text
"hotg-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
tshow Int
index) Text
"declpool" Text
"ghot" | Int
index <- [Int
1 .. Int
150]]
    filled <- env -> IO [JobRead CLPayload]
claimAs env
env
    length filled `shouldBe` 2
    enqueue env [groupedJob "coldg" "declpool" "gcold"]
    _ <- claimAs env
    inFlight env (fullKey "declpool" "gcold") `shouldReturn` Just 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"gates a group on the row it would claim" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- A failed job keeps the head of its group's line. The head gate judges that
    -- row. Here the fresh low-id sibling sits on a full key.
    env -> Int -> IO ()
seed env
env Int
1
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"declpool" Text
"gfresh"]
    filled <- env -> IO [JobRead CLPayload]
claimAs env
env
    length filled `shouldBe` 1
    enqueue env [groupedJob "rg" "declpool" "gfresh"]
    enqueue env [groupedJob "rg" "declpool" "gretry"]
    markAttempted env (fullKey "declpool" "gretry")
    _ <- claimAs env
    inFlight env (fullKey "declpool" "gretry") `shouldReturn` Just 1

  String -> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"keeps a group whose batch still has a claimable row under a blocked retry" ((env -> IO ()) -> SpecWith (Arg (env -> IO ())))
-> (env -> IO ()) -> SpecWith (Arg (env -> IO ()))
forall a b. (a -> b) -> a -> b
$ \env
env -> do
    -- The batch is taken attempts-first and cut by (priority, id). It yields a
    -- claim when its lowest-id row is admissible. The gate judges that row.
    env -> Int -> IO ()
seed env
env Int
1
    env -> [JobWrite CLPayload] -> IO ()
enqueue env
env [Text -> Text -> JobWrite CLPayload
job Text
"declpool" Text
"bhot"]
    filled <- env -> IO [JobRead CLPayload]
claimAs env
env
    length filled `shouldBe` 1
    enqueue env [groupedJob "bg" "declpool" "bfree"]
    enqueue env [groupedJob "bg" "declpool" "bhot"]
    markAttempted env (fullKey "declpool" "bhot")
    claimed <- runM env (HL.claimNextVisibleJobsBatched 2 100 60) :: IO [NE.NonEmpty (JobRead CLPayload)]
    map payload (concatMap NE.toList claimed) `shouldBe` [CLPayload "declpool:bfree"]