{-# LANGUAGE DataKinds #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}
{-# LANGUAGE ScopedTypeVariables #-}

-- | Property tests for per-job concurrency limits: a deterministic exact-model
-- sweep over a job lifecycle (insert, claim, ack, retry, override, prune,
-- reconcile), a concurrent never-over-admit check, and a grouped drain-to-empty
-- check. Each model key is a seeded pool driven through a single suffix. Claims
-- are attributed.
module Arbiter.Test.ConcurrencyModel
  ( concurrencyModelSpec
  ) where

import Arbiter.Core.Concurrency.Schema (arbiterConcurrencyTable)
import Arbiter.Core.Concurrency.Stats (ConcurrencyPolicyUpdate (..))
import Arbiter.Core.HighLevel qualified as HL
import Arbiter.Core.Job.Schema (jobQueueTable)
import Arbiter.Core.Job.Types
  ( DedupKey (..)
  , JobRead
  , defaultGroupedJob
  , defaultJob
  , payload
  , setDedupKey
  , setMaxAttempts
  )
import Arbiter.Core.MonadArbiter (HasRegistry)
import Arbiter.Core.Operations qualified as Ops
import Control.Monad (foldM_, void)
import Data.Foldable (for_, traverse_)
import Data.Int (Int32)
import Data.List.NonEmpty qualified as NE
import Data.Map.Strict (Map)
import Data.Map.Strict qualified as Map
import Data.Maybe (fromMaybe)
import Data.Proxy (Proxy (..))
import Data.Set qualified as Set
import Data.String (fromString)
import Data.Text (Text)
import Data.Text qualified as T
import Data.UUID.Types qualified as UUID
import Database.PostgreSQL.Simple qualified as PG
import Hedgehog
import Hedgehog.Gen qualified as Gen
import Hedgehog.Range qualified as Range
import Test.Hspec
import UnliftIO.Async (mapConcurrently)

import Arbiter.Test.ConcurrencyLimit (CLPayload (..), CLReg, concurrencyTable)
import Arbiter.Test.Setup (execute_, seedConcurrencyPoolSQL)

-- A worker id that attributes claims.
worker :: UUID.UUID
worker :: UUID
worker = Word32 -> Word32 -> Word32 -> Word32 -> UUID
UUID.fromWords Word32
0 Word32
0 Word32
0 Word32
11

-- Pools with a fixed seeded limit. The model is keyed by pool prefix.
modelPools :: [(Text, Int)]
modelPools :: [(Text, Int)]
modelPools = [(Text
"mx", Int
1), (Text
"my", Int
2), (Text
"mz", Int
3)]

-- A single suffix per pool. Each pool drives one count-row key @prefix:k@.
poolSuffix :: Text
poolSuffix :: Text
poolSuffix = Text
"k"

storedKey :: Text -> Text
storedKey :: Text -> Text
storedKey Text
prefix = Text
prefix Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
":" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
poolSuffix

limitOf :: Text -> Int
limitOf :: Text -> Int
limitOf Text
prefix = Int -> Maybe Int -> Int
forall a. a -> Maybe a -> a
fromMaybe Int
1 (Text -> [(Text, Int)] -> Maybe Int
forall a b. Eq a => a -> [(a, b)] -> Maybe b
lookup Text
prefix [(Text, Int)]
modelPools)

-- A claim batch larger than any per-key gate.
claimBatch :: Int
claimBatch :: Int
claimBatch = Int
500

-- | The concurrency state-machine properties, run against any backend.
concurrencyModelSpec
  :: forall sm
   . (HasRegistry sm CLReg)
  => (forall a. sm a -> IO a)
  -> (forall a. (PG.Connection -> IO a) -> IO a)
  -> Text
  -> Spec
concurrencyModelSpec :: forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Spec
concurrencyModelSpec forall a. sm a -> IO a
run forall a. (Connection -> IO a) -> IO a
withConn Text
schema = do
  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"the lifecycle keeps in_flight exact for every key" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$
    Property -> IO Bool
forall (m :: * -> *). MonadIO m => Property -> m Bool
check ((forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
prop_model sm a -> IO a
forall a. sm a -> IO a
run (Connection -> IO a) -> IO a
forall a. (Connection -> IO a) -> IO a
withConn Text
schema) IO Bool -> (Bool -> 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
>>= (Bool -> Bool -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Bool
True)
  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"concurrent claimers never admit more than a key's limit" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$
    Property -> IO Bool
forall (m :: * -> *). MonadIO m => Property -> m Bool
check ((forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
prop_concurrent sm a -> IO a
forall a. sm a -> IO a
run (Connection -> IO a) -> IO a
forall a. (Connection -> IO a) -> IO a
withConn Text
schema) IO Bool -> (Bool -> 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
>>= (Bool -> Bool -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Bool
True)
  String -> IO () -> SpecWith (Arg (IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"grouped jobs always drain, whatever the batch size and key layout" (IO () -> SpecWith (Arg (IO ())))
-> IO () -> SpecWith (Arg (IO ()))
forall a b. (a -> b) -> a -> b
$
    Property -> IO Bool
forall (m :: * -> *). MonadIO m => Property -> m Bool
check ((forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
prop_groupedDrain sm a -> IO a
forall a. sm a -> IO a
run (Connection -> IO a) -> IO a
forall a. (Connection -> IO a) -> IO a
withConn Text
schema) IO Bool -> (Bool -> 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
>>= (Bool -> Bool -> IO ()
forall a. (HasCallStack, Show a, Eq a) => a -> a -> IO ()
`shouldBe` Bool
True)

-- Pure reference model.

data KS = KS
  { KS -> Int
ksLimit :: Int
  , KS -> Int
ksPending :: Int
  , KS -> Int
ksClaimed :: Int
  }

-- Count rows, keyed by pool prefix.
type MM = Map Text KS

-- Pool overrides on the policy, keyed by prefix.
type Ov = Map Text Int

eff :: Ov -> Text -> KS -> Int
eff :: Ov -> Text -> KS -> Int
eff Ov
overrides Text
prefix KS
state = Int -> Maybe Int -> Int
forall a. a -> Maybe a -> a
fromMaybe (KS -> Int
ksLimit KS
state) (Text -> Ov -> Maybe Int
forall k a. Ord k => k -> Map k a -> Maybe a
Map.lookup Text
prefix Ov
overrides)

ensure :: Text -> MM -> MM
ensure :: Text -> MM -> MM
ensure Text
prefix = (KS -> KS -> KS) -> Text -> KS -> MM -> MM
forall k a. Ord k => (a -> a -> a) -> k -> a -> Map k a -> Map k a
Map.insertWith (\KS
_ KS
old -> KS
old) Text
prefix (Int -> Int -> Int -> KS
KS (Text -> Int
limitOf Text
prefix) Int
0 Int
0)

data Op
  = OInsert Text Int
  | OClaim
  | OAck Text Int
  | ORetry Text Int
  | OOverride Text (Maybe Int)
  | OPrune
  | OReconcile
  | -- | Dedup-replace a fresh job from one pool's key onto another's. The third
    -- field is a dedup key that is unique per op.
    OMove Text Text Text
  deriving stock (Int -> Op -> ShowS
[Op] -> ShowS
Op -> String
(Int -> Op -> ShowS)
-> (Op -> String) -> ([Op] -> ShowS) -> Show Op
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> Op -> ShowS
showsPrec :: Int -> Op -> ShowS
$cshow :: Op -> String
show :: Op -> String
$cshowList :: [Op] -> ShowS
showList :: [Op] -> ShowS
Show)

genOps :: Gen [Op]
genOps :: Gen [Op]
genOps =
  Range Int -> GenT Identity Op -> Gen [Op]
forall (m :: * -> *) a. MonadGen m => Range Int -> m a -> m [a]
Gen.list (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
1 Int
25) (GenT Identity Op -> Gen [Op]) -> GenT Identity Op -> Gen [Op]
forall a b. (a -> b) -> a -> b
$
    [(Int, GenT Identity Op)] -> GenT Identity Op
forall (m :: * -> *) a.
(HasCallStack, MonadGen m) =>
[(Int, m a)] -> m a
Gen.frequency
      [ (Int
3, Text -> Int -> Op
OInsert (Text -> Int -> Op)
-> GenT Identity Text -> GenT Identity (Int -> Op)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> GenT Identity Text
genKey GenT Identity (Int -> Op) -> GenT Identity Int -> GenT Identity Op
forall a b.
GenT Identity (a -> b) -> GenT Identity a -> GenT Identity b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Range Int -> GenT Identity Int
forall (m :: * -> *). MonadGen m => Range Int -> m Int
Gen.int (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
1 Int
4))
      , (Int
4, Op -> GenT Identity Op
forall a. a -> GenT Identity a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Op
OClaim)
      , (Int
2, Text -> Int -> Op
OAck (Text -> Int -> Op)
-> GenT Identity Text -> GenT Identity (Int -> Op)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> GenT Identity Text
genKey GenT Identity (Int -> Op) -> GenT Identity Int -> GenT Identity Op
forall a b.
GenT Identity (a -> b) -> GenT Identity a -> GenT Identity b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Range Int -> GenT Identity Int
forall (m :: * -> *). MonadGen m => Range Int -> m Int
Gen.int (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
1 Int
4))
      , (Int
2, Text -> Int -> Op
ORetry (Text -> Int -> Op)
-> GenT Identity Text -> GenT Identity (Int -> Op)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> GenT Identity Text
genKey GenT Identity (Int -> Op) -> GenT Identity Int -> GenT Identity Op
forall a b.
GenT Identity (a -> b) -> GenT Identity a -> GenT Identity b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Range Int -> GenT Identity Int
forall (m :: * -> *). MonadGen m => Range Int -> m Int
Gen.int (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
1 Int
4))
      , (Int
1, Text -> Maybe Int -> Op
OOverride (Text -> Maybe Int -> Op)
-> GenT Identity Text -> GenT Identity (Maybe Int -> Op)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> GenT Identity Text
genKey GenT Identity (Maybe Int -> Op)
-> GenT Identity (Maybe Int) -> GenT Identity Op
forall a b.
GenT Identity (a -> b) -> GenT Identity a -> GenT Identity b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> GenT Identity Int -> GenT Identity (Maybe Int)
forall (m :: * -> *) a. MonadGen m => m a -> m (Maybe a)
Gen.maybe (Range Int -> GenT Identity Int
forall (m :: * -> *). MonadGen m => Range Int -> m Int
Gen.int (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
0 Int
4)))
      , (Int
1, Op -> GenT Identity Op
forall a. a -> GenT Identity a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Op
OPrune)
      , (Int
1, Op -> GenT Identity Op
forall a. a -> GenT Identity a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Op
OReconcile)
      , (Int
1, Text -> Text -> Text -> Op
OMove (Text -> Text -> Text -> Op)
-> GenT Identity Text -> GenT Identity (Text -> Text -> Op)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> GenT Identity Text
genKey GenT Identity (Text -> Text -> Op)
-> GenT Identity Text -> GenT Identity (Text -> Op)
forall a b.
GenT Identity (a -> b) -> GenT Identity a -> GenT Identity b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> GenT Identity Text
genKey GenT Identity (Text -> Op)
-> GenT Identity Text -> GenT Identity Op
forall a b.
GenT Identity (a -> b) -> GenT Identity a -> GenT Identity b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> GenT Identity Text
genDedup)
      ]
  where
    genKey :: GenT Identity Text
genKey = [Text] -> GenT Identity Text
forall (f :: * -> *) (m :: * -> *) a.
(HasCallStack, Foldable f, MonadGen m) =>
f a -> m a
Gen.element (((Text, Int) -> Text) -> [(Text, Int)] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map (Text, Int) -> Text
forall a b. (a, b) -> a
fst [(Text, Int)]
modelPools)
    genDedup :: GenT Identity Text
genDedup = Range Int -> GenT Identity Char -> GenT Identity Text
forall (m :: * -> *). MonadGen m => Range Int -> m Char -> m Text
Gen.text (Int -> Range Int
forall a. a -> Range a
Range.singleton Int
12) GenT Identity Char
forall (m :: * -> *). MonadGen m => m Char
Gen.alphaNum

-- Apply a non-claim op to the pure model. Claim is handled with the live result.
applyModel :: Op -> (MM, Ov) -> (MM, Ov)
applyModel :: Op -> (MM, Ov) -> (MM, Ov)
applyModel Op
operation (MM
model, Ov
overrides) = case Op
operation of
  OInsert Text
prefix Int
count -> ((KS -> KS) -> Text -> MM -> MM
forall k a. Ord k => (a -> a) -> k -> Map k a -> Map k a
Map.adjust (\KS
state -> KS
state {ksPending = ksPending state + count}) Text
prefix (Text -> MM -> MM
ensure Text
prefix MM
model), Ov
overrides)
  OAck Text
prefix Int
count -> ((KS -> KS) -> Text -> MM -> MM
forall k a. Ord k => (a -> a) -> k -> Map k a -> Map k a
Map.adjust (\KS
state -> KS
state {ksClaimed = ksClaimed state - min count (ksClaimed state)}) Text
prefix MM
model, Ov
overrides)
  ORetry Text
prefix Int
count ->
    ( (KS -> KS) -> Text -> MM -> MM
forall k a. Ord k => (a -> a) -> k -> Map k a -> Map k a
Map.adjust
        ( \KS
state ->
            let moved :: Int
moved = Int -> Int -> Int
forall a. Ord a => a -> a -> a
min Int
count (KS -> Int
ksClaimed KS
state)
             in KS
state {ksClaimed = ksClaimed state - moved, ksPending = ksPending state + moved}
        )
        Text
prefix
        MM
model
    , Ov
overrides
    )
  OOverride Text
prefix Maybe Int
newLimit -> (MM
model, (Ov -> Ov) -> (Int -> Ov -> Ov) -> Maybe Int -> Ov -> Ov
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (Text -> Ov -> Ov
forall k a. Ord k => k -> Map k a -> Map k a
Map.delete Text
prefix) (Text -> Int -> Ov -> Ov
forall k a. Ord k => k -> a -> Map k a -> Map k a
Map.insert Text
prefix) Maybe Int
newLimit Ov
overrides)
  Op
OPrune -> ((KS -> Bool) -> MM -> MM
forall a k. (a -> Bool) -> Map k a -> Map k a
Map.filter (\KS
state -> KS -> Int
ksPending KS
state Int -> Int -> Int
forall a. Num a => a -> a -> a
+ KS -> Int
ksClaimed KS
state Int -> Int -> Bool
forall a. Eq a => a -> a -> Bool
/= Int
0) MM
model, Ov
overrides)
  Op
OReconcile -> (MM
model, Ov
overrides)
  Op
OClaim -> (MM
model, Ov
overrides)
  -- The moved job is unclaimed. Pending gains one on the destination and the
  -- source keeps a drained count row.
  OMove Text
source Text
destination Text
_ ->
    ( (KS -> KS) -> Text -> MM -> MM
forall k a. Ord k => (a -> a) -> k -> Map k a -> Map k a
Map.adjust (\KS
state -> KS
state {ksPending = ksPending state + 1}) Text
destination (Text -> MM -> MM
ensure Text
destination (Text -> MM -> MM
ensure Text
source MM
model))
    , Ov
overrides
    )

-- Per-key admissions a claim grants.
claimDeltas :: Ov -> MM -> Map Text Int
claimDeltas :: Ov -> MM -> Ov
claimDeltas Ov
overrides = (Text -> KS -> Int) -> MM -> Ov
forall k a b. (k -> a -> b) -> Map k a -> Map k b
Map.mapWithKey (\Text
prefix KS
state -> Int -> Int -> Int
forall a. Ord a => a -> a -> a
min (KS -> Int
ksPending KS
state) (Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
0 (Ov -> Text -> KS -> Int
eff Ov
overrides Text
prefix KS
state Int -> Int -> Int
forall a. Num a => a -> a -> a
- KS -> Int
ksClaimed KS
state)))

prop_model
  :: (HasRegistry sm CLReg)
  => (forall a. sm a -> IO a)
  -> (forall a. (PG.Connection -> IO a) -> IO a)
  -> Text
  -> Property
prop_model :: forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
prop_model forall a. sm a -> IO a
run forall a. (Connection -> IO a) -> IO a
withConn Text
schema = TestLimit -> Property -> Property
withTests TestLimit
60 (Property -> Property) -> Property -> Property
forall a b. (a -> b) -> a -> b
$ HasCallStack => PropertyT IO () -> Property
PropertyT IO () -> Property
property (PropertyT IO () -> Property) -> PropertyT IO () -> Property
forall a b. (a -> b) -> a -> b
$ do
  ops <- Gen [Op] -> PropertyT IO [Op]
forall (m :: * -> *) a.
(Monad m, Show a, HasCallStack) =>
Gen a -> PropertyT m a
forAll Gen [Op]
genOps
  evalIO (resetState withConn schema)
  foldM_ (step run withConn schema) (Map.empty, Map.empty, Map.empty) ops

-- Live jobs currently held claimed, grouped by key.
type Held = Map Text [JobRead CLPayload]

step
  :: (HasRegistry sm CLReg)
  => (forall a. sm a -> IO a)
  -> (forall a. (PG.Connection -> IO a) -> IO a)
  -> Text
  -> (MM, Ov, Held)
  -> Op
  -> PropertyT IO (MM, Ov, Held)
step :: forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a)
-> Text
-> (MM, Ov, Held)
-> Op
-> PropertyT IO (MM, Ov, Held)
step forall a. sm a -> IO a
run forall a. (Connection -> IO a) -> IO a
withConn Text
schema (MM
model, Ov
overrides, Held
held) Op
operation = do
  -- Non-claim branches run their effect and defer the model update to applyModel.
  let done :: Held -> PropertyT IO (MM, Ov, Held)
done Held
held' = let (MM
nextModel, Ov
nextOverrides) = Op -> (MM, Ov) -> (MM, Ov)
applyModel Op
operation (MM
model, Ov
overrides) in (MM, Ov, Held) -> PropertyT IO (MM, Ov, Held)
forall a. a -> PropertyT IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (MM
nextModel, Ov
nextOverrides, Held
held')
  (model', overrides', held') <- case Op
operation of
    OInsert Text
prefix Int
count -> do
      let job :: JobWrite CLPayload
job = Maybe Int32 -> JobWrite CLPayload -> JobWrite CLPayload
forall payload. Maybe Int32 -> JobWrite payload -> JobWrite payload
setMaxAttempts (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
1000) (JobWrite CLPayload -> JobWrite CLPayload)
-> JobWrite CLPayload -> JobWrite CLPayload
forall a b. (a -> b) -> a -> b
$ CLPayload -> JobWrite CLPayload
forall payload. payload -> JobWrite payload
defaultJob (Text -> CLPayload
CLPayload Text
prefix)
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (IO [JobRead CLPayload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (sm [JobRead CLPayload] -> IO [JobRead CLPayload]
forall a. sm a -> IO a
run ([JobWrite CLPayload] -> sm [JobRead CLPayload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch (Int -> JobWrite CLPayload -> [JobWrite CLPayload]
forall a. Int -> a -> [a]
replicate Int
count JobWrite CLPayload
job)) :: IO [JobRead CLPayload]))
      Held -> PropertyT IO (MM, Ov, Held)
done Held
held
    Op
OClaim -> do
      claimed <- IO [JobRead CLPayload] -> PropertyT IO [JobRead CLPayload]
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (sm [JobRead CLPayload] -> IO [JobRead CLPayload]
forall a. sm a -> IO a
run (Int -> NominalDiffTime -> UUID -> sm [JobRead CLPayload]
forall payload (m :: * -> *).
QueueOperation m payload =>
Int -> NominalDiffTime -> UUID -> m [JobRead payload]
HL.claimNextVisibleJobsAs Int
claimBatch NominalDiffTime
60 UUID
worker) :: IO [JobRead CLPayload])
      let grouped = [JobRead CLPayload] -> Held
groupByKey [JobRead CLPayload]
claimed
          expected = Ov -> MM -> Ov
claimDeltas Ov
overrides MM
model
      -- The gate admitted exactly the model's per-key free slots.
      for_ (map fst modelPools) $ \Text
prefix ->
        Int -> Text -> Ov -> Int
forall k a. Ord k => a -> k -> Map k a -> a
Map.findWithDefault Int
0 Text
prefix (([JobRead CLPayload] -> Int) -> Held -> Ov
forall a b k. (a -> b) -> Map k a -> Map k b
Map.map [JobRead CLPayload] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length Held
grouped) Int -> Int -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, Eq a, Show a, HasCallStack) =>
a -> a -> m ()
=== Int -> Text -> Ov -> Int
forall k a. Ord k => a -> k -> Map k a -> a
Map.findWithDefault Int
0 Text
prefix Ov
expected
      let nextModel =
            (Text -> KS -> KS) -> MM -> MM
forall k a b. (k -> a -> b) -> Map k a -> Map k b
Map.mapWithKey
              ( \Text
prefix KS
state ->
                  let admitted :: Int
admitted = Int -> Text -> Ov -> Int
forall k a. Ord k => a -> k -> Map k a -> a
Map.findWithDefault Int
0 Text
prefix Ov
expected
                   in KS
state {ksClaimed = ksClaimed state + admitted, ksPending = ksPending state - admitted}
              )
              MM
model
      pure (nextModel, overrides, Map.unionWith (++) held grouped)
    OAck Text
prefix Int
count -> do
      let ([JobRead CLPayload]
toAck, [JobRead CLPayload]
rest) = Int
-> [JobRead CLPayload]
-> ([JobRead CLPayload], [JobRead CLPayload])
forall a. Int -> [a] -> ([a], [a])
splitAt Int
count ([JobRead CLPayload] -> Text -> Held -> [JobRead CLPayload]
forall k a. Ord k => a -> k -> Map k a -> a
Map.findWithDefault [] Text
prefix Held
held)
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (sm () -> IO ()
forall a. sm a -> IO a
run ((JobRead CLPayload -> sm Int64) -> [JobRead CLPayload] -> sm ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ JobRead CLPayload -> sm Int64
forall payload (m :: * -> *).
JobOperation m payload =>
JobRead payload -> m Int64
HL.ackJob [JobRead CLPayload]
toAck))
      Held -> PropertyT IO (MM, Ov, Held)
done (Text -> [JobRead CLPayload] -> Held -> Held
forall k a. Ord k => k -> a -> Map k a -> Map k a
Map.insert Text
prefix [JobRead CLPayload]
rest Held
held)
    ORetry Text
prefix Int
count -> do
      let ([JobRead CLPayload]
toRetry, [JobRead CLPayload]
rest) = Int
-> [JobRead CLPayload]
-> ([JobRead CLPayload], [JobRead CLPayload])
forall a. Int -> [a] -> ([a], [a])
splitAt Int
count ([JobRead CLPayload] -> Text -> Held -> [JobRead CLPayload]
forall k a. Ord k => a -> k -> Map k a -> a
Map.findWithDefault [] Text
prefix Held
held)
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (sm () -> IO ()
forall a. sm a -> IO a
run ((JobRead CLPayload -> sm Int64) -> [JobRead CLPayload] -> sm ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (NominalDiffTime -> Text -> JobRead CLPayload -> sm Int64
forall payload (m :: * -> *).
JobOperation m payload =>
NominalDiffTime -> Text -> JobRead payload -> m Int64
HL.updateJobForRetry NominalDiffTime
0 Text
"model retry") [JobRead CLPayload]
toRetry))
      Held -> PropertyT IO (MM, Ov, Held)
done (Text -> [JobRead CLPayload] -> Held -> Held
forall k a. Ord k => k -> a -> Map k a -> Map k a
Map.insert Text
prefix [JobRead CLPayload]
rest Held
held)
    OOverride Text
prefix Maybe Int
newLimit -> do
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO
        (IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (sm Int64 -> IO Int64
forall a. sm a -> IO a
run (Text -> ConcurrencyPolicyUpdate -> sm Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> ConcurrencyPolicyUpdate -> m Int64
HL.updateConcurrencyPolicyOverrides Text
prefix (Maybe (Maybe Int32) -> ConcurrencyPolicyUpdate
ConcurrencyPolicyUpdate (Maybe Int32 -> Maybe (Maybe Int32)
forall a. a -> Maybe a
Just (Int -> Int32
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int -> Int32) -> Maybe Int -> Maybe Int32
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int
newLimit))))))
      Held -> PropertyT IO (MM, Ov, Held)
done Held
held
    Op
OPrune -> do
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (sm Int64 -> IO Int64
forall a. sm a -> IO a
run sm Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m Int64
HL.pruneConcurrencyKeys))
      Held -> PropertyT IO (MM, Ov, Held)
done Held
held
    Op
OReconcile -> do
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (IO Int64 -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (sm Int64 -> IO Int64
forall a. sm a -> IO a
run sm Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m Int64
HL.reconcileConcurrencyCounts))
      Held -> PropertyT IO (MM, Ov, Held)
done Held
held
    OMove Text
source Text
destination Text
dedup -> do
      let mkJob :: Text -> JobWrite CLPayload
mkJob Text
prefix = 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
dedup)) (JobWrite CLPayload -> JobWrite CLPayload)
-> JobWrite CLPayload -> JobWrite CLPayload
forall a b. (a -> b) -> a -> b
$ Maybe Int32 -> JobWrite CLPayload -> JobWrite CLPayload
forall payload. Maybe Int32 -> JobWrite payload -> JobWrite payload
setMaxAttempts (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
1000) (JobWrite CLPayload -> JobWrite CLPayload)
-> JobWrite CLPayload -> JobWrite CLPayload
forall a b. (a -> b) -> a -> b
$ CLPayload -> JobWrite CLPayload
forall payload. payload -> JobWrite payload
defaultJob (Text -> CLPayload
CLPayload Text
prefix)
      -- Seed the source key, then dedup-replace onto the destination.
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (IO [JobRead CLPayload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (sm [JobRead CLPayload] -> IO [JobRead CLPayload]
forall a. sm a -> IO a
run ([JobWrite CLPayload] -> sm [JobRead CLPayload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch [Text -> JobWrite CLPayload
mkJob Text
source]) :: IO [JobRead CLPayload]))
      IO () -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, MonadIO m, HasCallStack) =>
IO a -> m a
evalIO (IO [JobRead CLPayload] -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (sm [JobRead CLPayload] -> IO [JobRead CLPayload]
forall a. sm a -> IO a
run ([JobWrite CLPayload] -> sm [JobRead CLPayload]
forall payload (m :: * -> *).
QueueOperation m payload =>
[JobWrite payload] -> m [JobRead payload]
HL.insertJobsBatch [Text -> JobWrite CLPayload
mkJob Text
destination]) :: IO [JobRead CLPayload]))
      Held -> PropertyT IO (MM, Ov, Held)
done Held
held
  -- A count row exists from the first insert until a prune drops it.
  stored <- evalIO (readCounts withConn schema)
  for_ (map fst modelPools) $ \Text
prefix ->
    case Text -> MM -> Maybe KS
forall k a. Ord k => k -> Map k a -> Maybe a
Map.lookup Text
prefix MM
model' of
      Just KS
state ->
        Text -> Map Text Int32 -> Maybe Int32
forall k a. Ord k => k -> Map k a -> Maybe a
Map.lookup (Text -> Text
storedKey Text
prefix) Map Text Int32
stored Maybe Int32 -> Maybe Int32 -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, Eq a, Show a, HasCallStack) =>
a -> a -> m ()
=== Int32 -> Maybe Int32
forall a. a -> Maybe a
Just (Int -> Int32
forall a b. (Integral a, Num b) => a -> b
fromIntegral (KS -> Int
ksClaimed KS
state))
      Maybe KS
Nothing -> Text -> Map Text Int32 -> Maybe Int32
forall k a. Ord k => k -> Map k a -> Maybe a
Map.lookup (Text -> Text
storedKey Text
prefix) Map Text Int32
stored Maybe Int32 -> Maybe Int32 -> PropertyT IO ()
forall (m :: * -> *) a.
(MonadTest m, Eq a, Show a, HasCallStack) =>
a -> a -> m ()
=== Maybe Int32
forall a. Maybe a
Nothing
  pure (model', overrides', held')

prop_concurrent
  :: (HasRegistry sm CLReg)
  => (forall a. sm a -> IO a)
  -> (forall a. (PG.Connection -> IO a) -> IO a)
  -> Text
  -> Property
prop_concurrent :: forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
prop_concurrent forall a. sm a -> IO a
run forall a. (Connection -> IO a) -> IO a
withConn Text
schema = TestLimit -> Property -> Property
withTests TestLimit
30 (Property -> Property) -> Property -> Property
forall a b. (a -> b) -> a -> b
$ HasCallStack => PropertyT IO () -> Property
PropertyT IO () -> Property
property (PropertyT IO () -> Property) -> PropertyT IO () -> Property
forall a b. (a -> b) -> a -> b
$ do
  lim <- GenT Identity Int -> PropertyT IO Int
forall (m :: * -> *) a.
(Monad m, Show a, HasCallStack) =>
Gen a -> PropertyT m a
forAll (Range Int -> GenT Identity Int
forall (m :: * -> *). MonadGen m => Range Int -> m Int
Gen.int (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
1 Int
3))
  extra <- forAll (Gen.int (Range.linear 1 6))
  claimers <- forAll (Gen.int (Range.linear 2 8))
  evalIO (resetState withConn schema)
  let prefix = Text
"mx"
      jobCount = Int
lim Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
extra
      job = Maybe Int32 -> JobWrite CLPayload -> JobWrite CLPayload
forall payload. Maybe Int32 -> JobWrite payload -> JobWrite payload
setMaxAttempts (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
1000) (JobWrite CLPayload -> JobWrite CLPayload)
-> JobWrite CLPayload -> JobWrite CLPayload
forall a b. (a -> b) -> a -> b
$ CLPayload -> JobWrite CLPayload
forall payload. payload -> JobWrite payload
defaultJob (Text -> CLPayload
CLPayload Text
prefix)
  evalIO (withConn $ \Connection
conn -> Connection -> Text -> Text -> Int -> IO ()
seedPool Connection
conn Text
schema Text
prefix Int
lim)
  evalIO (void (run (HL.insertJobsBatch (replicate jobCount job)) :: IO [JobRead CLPayload]))
  results <-
    evalIO
      (mapConcurrently (const (run (HL.claimNextVisibleJobsAs claimBatch 60 worker) :: IO [JobRead CLPayload])) [1 .. claimers])
  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)
  -- A key's cap holds under contention.
  assert (total <= lim)
  -- An uncontended follow-up claim fills the key to the cap. The contended total
  -- is not deterministic.
  more <- evalIO (run (HL.claimNextVisibleJobsAs claimBatch 60 worker) :: IO [JobRead CLPayload])
  let claimed = Int
total Int -> Int -> Int
forall a. Num a => a -> a -> a
+ [JobRead CLPayload] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [JobRead CLPayload]
more
  claimed === lim
  evalIO (void (run HL.reconcileConcurrencyCounts))
  stored <- evalIO (readCounts withConn schema)
  Map.lookup (storedKey prefix) stored === Just (fromIntegral claimed)

-- | One pool's key is pinned at its cap by a job that is never acked. Every group
-- with no job on that key still drains. Generated over batch size, per-poll slot
-- budget, group size, and key layout.
prop_groupedDrain
  :: (HasRegistry sm CLReg)
  => (forall a. sm a -> IO a)
  -> (forall a. (PG.Connection -> IO a) -> IO a)
  -> Text
  -> Property
prop_groupedDrain :: forall (sm :: * -> *).
HasRegistry sm CLReg =>
(forall a. sm a -> IO a)
-> (forall a. (Connection -> IO a) -> IO a) -> Text -> Property
prop_groupedDrain forall a. sm a -> IO a
run forall a. (Connection -> IO a) -> IO a
withConn Text
schema = TestLimit -> Property -> Property
withTests TestLimit
40 (Property -> Property) -> Property -> Property
forall a b. (a -> b) -> a -> b
$ HasCallStack => PropertyT IO () -> Property
PropertyT IO () -> Property
property (PropertyT IO () -> Property) -> PropertyT IO () -> Property
forall a b. (a -> b) -> a -> b
$ do
  perGroup <- GenT Identity Int -> PropertyT IO Int
forall (m :: * -> *) a.
(Monad m, Show a, HasCallStack) =>
Gen a -> PropertyT m a
forAll (Range Int -> GenT Identity Int
forall (m :: * -> *). MonadGen m => Range Int -> m Int
Gen.int (Int -> Int -> Range Int
forall a. Integral a => a -> a -> Range a
Range.linear Int
1 Int
3))
  batchSize <- forAll (Gen.int (Range.linear 1 3))
  -- Small enough for the slot budget to bind.
  maxBatches <- forAll (Gen.int (Range.linear 1 3))
  keys <- forAll (Gen.list (Range.linear 2 16) (Gen.element (map fst modelPools)))
  evalIO (resetState withConn schema)
  -- Pin the hot pool at its cap with a claim that is never acked.
  evalIO (void (run (HL.insertJobsBatch [defaultJob (CLPayload hotPool)]) :: IO [JobRead CLPayload]))
  pinned <- evalIO (run (HL.claimNextVisibleJobsAs claimBatch 600 worker) :: IO [JobRead CLPayload])
  length pinned === limitOf hotPool
  let groupOf Int
index = Text
"dg" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show (Int
index Int -> Int -> Int
forall a. Integral a => a -> a -> a
`div` Int
perGroup :: Int))
      tagged = [(Int -> Text
groupOf Int
index, Text
prefix) | (Int
index, Text
prefix) <- [Int] -> [Text] -> [(Int, Text)]
forall a b. [a] -> [b] -> [(a, b)]
zip [Int
0 :: Int ..] [Text]
keys]
      jobs = [Maybe Int32 -> JobWrite CLPayload -> JobWrite CLPayload
forall payload. Maybe Int32 -> JobWrite payload -> JobWrite payload
setMaxAttempts (Int32 -> Maybe Int32
forall a. a -> Maybe a
Just Int32
1000) (Text -> CLPayload -> JobWrite CLPayload
forall payload. Text -> payload -> JobWrite payload
defaultGroupedJob Text
groupKey (Text -> CLPayload
CLPayload Text
prefix)) | (Text
groupKey, Text
prefix) <- [(Text, Text)]
tagged]
      -- A group with no job on the pinned key is cold.
      cold =
        [ Text
groupKey
        | Text
groupKey <- Set Text -> [Text]
forall a. Set a -> [a]
Set.toList ([Text] -> Set Text
forall a. Ord a => [a] -> Set a
Set.fromList (((Text, Text) -> Text) -> [(Text, Text)] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map (Text, Text) -> Text
forall a b. (a, b) -> a
fst [(Text, Text)]
tagged))
        , ((Text, Text) -> Bool) -> [(Text, Text)] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
all ((Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
/= Text
hotPool) (Text -> Bool) -> ((Text, Text) -> Text) -> (Text, Text) -> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Text, Text) -> Text
forall a b. (a, b) -> b
snd) (((Text, Text) -> Bool) -> [(Text, Text)] -> [(Text, Text)]
forall a. (a -> Bool) -> [a] -> [a]
filter ((Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
groupKey) (Text -> Bool) -> ((Text, Text) -> Text) -> (Text, Text) -> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Text, Text) -> Text
forall a b. (a, b) -> a
fst) [(Text, Text)]
tagged)
        ]
      claimSql = Proxy CLPayload
-> Text
-> Text
-> Int
-> Int
-> NominalDiffTime
-> UUID
-> ClaimSql
forall payload (proxy :: * -> *).
JobPayload payload =>
proxy payload
-> Text
-> Text
-> Int
-> Int
-> NominalDiffTime
-> UUID
-> ClaimSql
Ops.mkClaimSql (Proxy CLPayload
forall {k} (t :: k). Proxy t
Proxy :: Proxy CLPayload) Text
schema Text
concurrencyTable Int
batchSize Int
0 NominalDiffTime
60 UUID
worker
      round_ =
        sm Int -> IO Int
forall a. sm a -> IO a
run
          ( do
              claimed <- ClaimSql -> Int -> sm [NonEmpty (JobRead CLPayload)]
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
ClaimSql -> Int -> m [NonEmpty (JobRead payload)]
Ops.claimJobsBatchedCached ClaimSql
claimSql Int
maxBatches
              let got = (NonEmpty (JobRead CLPayload) -> [JobRead CLPayload])
-> [NonEmpty (JobRead CLPayload)] -> [JobRead CLPayload]
forall (t :: * -> *) a b. Foldable t => (a -> [b]) -> t a -> [b]
concatMap NonEmpty (JobRead CLPayload) -> [JobRead CLPayload]
forall a. NonEmpty a -> [a]
NE.toList ([NonEmpty (JobRead CLPayload)]
claimed :: [NE.NonEmpty (JobRead CLPayload)])
              traverse_ HL.ackJob got
              pure (length got)
          )
      loop Int
remaining
        | Int
remaining Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
<= (Int
0 :: Int) = () -> IO ()
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
        | Bool
otherwise = IO Int
round_ IO Int -> (Int -> 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
>>= \Int
claimedCount -> if Int
claimedCount Int -> Int -> Bool
forall a. Eq a => a -> a -> Bool
== Int
0 then () -> IO ()
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure () else Int -> IO ()
loop (Int
remaining Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1)
  evalIO (void (run (HL.insertJobsBatch jobs) :: IO [JobRead CLPayload]))
  evalIO (loop (length jobs + 4))
  left <- evalIO (remainingGroups withConn schema)
  filter (`elem` cold) left === []

-- Helpers.

-- | The pool pinned at its cap by 'prop_groupedDrain'.
hotPool :: Text
hotPool :: Text
hotPool = Text
"mx"

-- | Group keys that still have rows.
remainingGroups :: (forall a. (PG.Connection -> IO a) -> IO a) -> Text -> IO [Text]
remainingGroups :: (forall a. (Connection -> IO a) -> IO a) -> Text -> IO [Text]
remainingGroups forall a. (Connection -> IO a) -> IO a
withConn Text
schema =
  (Connection -> IO [Text]) -> IO [Text]
forall a. (Connection -> IO a) -> IO a
withConn ((Connection -> IO [Text]) -> IO [Text])
-> (Connection -> IO [Text]) -> IO [Text]
forall a b. (a -> b) -> a -> b
$ \Connection
conn -> do
    rows <-
      Connection -> Query -> IO [Only Text]
forall r. FromRow r => Connection -> Query -> IO [r]
PG.query_
        Connection
conn
        (Text -> Query
stmt (Text
"SELECT DISTINCT group_key FROM " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> Text
jobQueueTable Text
schema Text
concurrencyTable Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" WHERE group_key IS NOT NULL"))
    pure (map PG.fromOnly rows)

groupByKey :: [JobRead CLPayload] -> Map Text [JobRead CLPayload]
groupByKey :: [JobRead CLPayload] -> Held
groupByKey = (Held -> JobRead CLPayload -> Held)
-> Held -> [JobRead CLPayload] -> Held
forall b a. (b -> a -> b) -> b -> [a] -> b
forall (t :: * -> *) b a.
Foldable t =>
(b -> a -> b) -> b -> t a -> b
foldl' (\Held
acc JobRead CLPayload
job -> ([JobRead CLPayload] -> [JobRead CLPayload] -> [JobRead CLPayload])
-> Text -> [JobRead CLPayload] -> Held -> Held
forall k a. Ord k => (a -> a -> a) -> k -> a -> Map k a -> Map k a
Map.insertWith [JobRead CLPayload] -> [JobRead CLPayload] -> [JobRead CLPayload]
forall a. [a] -> [a] -> [a]
(++) (JobRead CLPayload -> Text
forall {key} {q} {insertedAt} {adm}.
JobRecord CLPayload key q insertedAt adm -> Text
keyOf JobRead CLPayload
job) [JobRead CLPayload
job] Held
acc) Held
forall k a. Map k a
Map.empty
  where
    keyOf :: JobRecord CLPayload key q insertedAt adm -> Text
keyOf JobRecord CLPayload key q insertedAt adm
job = let CLPayload Text
prefix = JobRecord CLPayload key q insertedAt adm -> CLPayload
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> payload
payload JobRecord CLPayload key q insertedAt adm
job in Text
prefix

resetState :: (forall a. (PG.Connection -> IO a) -> IO a) -> Text -> IO ()
resetState :: (forall a. (Connection -> IO a) -> IO a) -> Text -> IO ()
resetState forall a. (Connection -> IO a) -> IO a
withConn Text
schema =
  (Connection -> IO ()) -> IO ()
forall a. (Connection -> IO a) -> IO a
withConn ((Connection -> IO ()) -> IO ()) -> (Connection -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \Connection
conn -> do
    Connection -> Text -> IO ()
execute_ Connection
conn (Text
"DELETE FROM " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> Text
jobQueueTable Text
schema Text
concurrencyTable)
    Connection -> Text -> IO ()
execute_ Connection
conn (Text
"TRUNCATE " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> Text
arbiterConcurrencyTable Text
schema)
    ((Text, Int) -> IO ()) -> [(Text, Int)] -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ ((Text -> Int -> IO ()) -> (Text, Int) -> IO ()
forall a b c. (a -> b -> c) -> (a, b) -> c
uncurry (Connection -> Text -> Text -> Int -> IO ()
seedPool Connection
conn Text
schema)) [(Text, Int)]
modelPools

seedPool :: PG.Connection -> Text -> Text -> Int -> IO ()
seedPool :: Connection -> Text -> Text -> Int -> IO ()
seedPool Connection
conn Text
schema Text
prefix Int
lim =
  (Text -> IO ()) -> [Text] -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (Connection -> Text -> IO ()
execute_ Connection
conn) (Text -> Text -> Int32 -> [Text]
seedConcurrencyPoolSQL Text
schema Text
prefix (Int -> Int32
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
lim))

readCounts :: (forall a. (PG.Connection -> IO a) -> IO a) -> Text -> IO (Map Text Int32)
readCounts :: (forall a. (Connection -> IO a) -> IO a)
-> Text -> IO (Map Text Int32)
readCounts forall a. (Connection -> IO a) -> IO a
withConn Text
schema =
  (Connection -> IO (Map Text Int32)) -> IO (Map Text Int32)
forall a. (Connection -> IO a) -> IO a
withConn ((Connection -> IO (Map Text Int32)) -> IO (Map Text Int32))
-> (Connection -> IO (Map Text Int32)) -> IO (Map Text Int32)
forall a b. (a -> b) -> a -> b
$ \Connection
conn -> do
    rows <-
      Connection -> Query -> IO [(Text, Int32)]
forall r. FromRow r => Connection -> Query -> IO [r]
PG.query_
        Connection
conn
        (Text -> Query
stmt (Text
"SELECT concurrency_key, in_flight FROM " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> Text
arbiterConcurrencyTable Text
schema))
    pure (Map.fromList rows)

stmt :: Text -> PG.Query
stmt :: Text -> Query
stmt = String -> Query
forall a. IsString a => String -> a
fromString (String -> Query) -> (Text -> String) -> Text -> Query
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> String
T.unpack