{-# LANGUAGE DataKinds #-}
{-# LANGUAGE DeriveAnyClass #-}
{-# LANGUAGE DeriveGeneric #-}
{-# LANGUAGE DerivingStrategies #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TypeApplications #-}
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)
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)
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")
type CLReg = '[Queue "arbiter_concurrency_test" CLPayload]
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]
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))
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])
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 ()
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 :: 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]
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)
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
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
_ ->
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
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
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
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
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
pruned <- timeout 5000000 (prune env)
pruned `shouldBe` Just 1
inFlight env (fullKey "dk" "a") `shouldReturn` Nothing
putMVar pruneDone ()
wait enqueuer
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
ackAll env claimed
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
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
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
truncateCounts env
inFlight env (fullKey "stale" "a") `shouldReturn` Nothing
blocked <- claimAs env
length blocked `shouldBe` 0
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
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
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
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)
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)
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
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` []
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
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
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
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
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"]