{-# LANGUAGE AllowAmbiguousTypes #-}
{-# LANGUAGE ConstraintKinds #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE FlexibleInstances #-}
{-# LANGUAGE MultiParamTypeClasses #-}
{-# LANGUAGE OverloadedStrings #-}

-- | Per-job concurrency limits. A payload's 'concurrencyFor' describes how to
-- select a pool and key. Static inspection finds all pools that the migration
-- must initialize. Each selected pool supplies the limit.
module Arbiter.Core.Concurrency.Spec
  ( -- * Core types
    ConcurrencyKey (..)
  , concurrencyKeyText
  , ConcurrencyPolicy (..)
  , concurrencyPool

    -- * Selecting a pool per job
  , HasConcurrency (..)
  , ConcurrencyFor
  , noConcurrency
  , concurrencyBy
  , globalConcurrency
  , concurrencyByCase
  , chooseWhen
  , runConcurrencyFor
  , collectPolicies

    -- * Registry reflection
  , RegistryConcurrencyPolicies
  , registryConcurrencyPolicies
  , registryConcurrencyTables
  ) where

import Data.Aeson (FromJSON (..), ToJSON (..))
import Data.Int (Int32)
import Data.Set (Set)
import Data.Text (Text)

import Arbiter.Core.Admission
  ( AdmissionPolicy (..)
  , CollectFor (..)
  , RegistryPolicies (..)
  , prefixedKeyParseJSON
  , prefixedKeyText
  , prefixedKeyToJSON
  , registryPolicies
  , registryPolicyTables
  , selectBy
  , selectNone
  )
import Arbiter.Core.Selector (Selector, chooseWhen, collectPolicies, runSelector, selectByCase)

-- | A resolved concurrency key with a pool prefix and per-key suffix. The
-- stored form is @prefix:suffix@. The separate prefix supports policy lookup.
data ConcurrencyKey = ConcurrencyKey
  { ConcurrencyKey -> Text
ckPrefix :: Text
  , ConcurrencyKey -> Text
ckSuffix :: Text
  }
  deriving stock (ConcurrencyKey -> ConcurrencyKey -> Bool
(ConcurrencyKey -> ConcurrencyKey -> Bool)
-> (ConcurrencyKey -> ConcurrencyKey -> Bool) -> Eq ConcurrencyKey
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: ConcurrencyKey -> ConcurrencyKey -> Bool
== :: ConcurrencyKey -> ConcurrencyKey -> Bool
$c/= :: ConcurrencyKey -> ConcurrencyKey -> Bool
/= :: ConcurrencyKey -> ConcurrencyKey -> Bool
Eq, Int -> ConcurrencyKey -> ShowS
[ConcurrencyKey] -> ShowS
ConcurrencyKey -> String
(Int -> ConcurrencyKey -> ShowS)
-> (ConcurrencyKey -> String)
-> ([ConcurrencyKey] -> ShowS)
-> Show ConcurrencyKey
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ConcurrencyKey -> ShowS
showsPrec :: Int -> ConcurrencyKey -> ShowS
$cshow :: ConcurrencyKey -> String
show :: ConcurrencyKey -> String
$cshowList :: [ConcurrencyKey] -> ShowS
showList :: [ConcurrencyKey] -> ShowS
Show)

instance ToJSON ConcurrencyKey where
  toJSON :: ConcurrencyKey -> Value
toJSON (ConcurrencyKey Text
prefix Text
suffix) = Text -> Text -> Value
prefixedKeyToJSON Text
prefix Text
suffix

instance FromJSON ConcurrencyKey where
  parseJSON :: Value -> Parser ConcurrencyKey
parseJSON = String
-> (Text -> Text -> ConcurrencyKey)
-> Value
-> Parser ConcurrencyKey
forall a. String -> (Text -> Text -> a) -> Value -> Parser a
prefixedKeyParseJSON String
"ConcurrencyKey" Text -> Text -> ConcurrencyKey
ConcurrencyKey

-- | The stored key text, @prefix:suffix@.
concurrencyKeyText :: ConcurrencyKey -> Text
concurrencyKeyText :: ConcurrencyKey -> Text
concurrencyKeyText (ConcurrencyKey Text
prefix Text
suffix) = Text -> Text -> Text
prefixedKeyText Text
prefix Text
suffix

-- | A concurrency pool. At most @cpLimit@ jobs share a key under @cpPrefix@. The
-- default is seeded. An operator override on the pool takes precedence.
data ConcurrencyPolicy = ConcurrencyPolicy
  { ConcurrencyPolicy -> Text
cpPrefix :: Text
  , ConcurrencyPolicy -> Int32
cpLimit :: Int32
  }
  deriving stock (ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
(ConcurrencyPolicy -> ConcurrencyPolicy -> Bool)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> Bool)
-> Eq ConcurrencyPolicy
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
== :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
$c/= :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
/= :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
Eq, Eq ConcurrencyPolicy
Eq ConcurrencyPolicy =>
(ConcurrencyPolicy -> ConcurrencyPolicy -> Ordering)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> Bool)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> Bool)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> Bool)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> Bool)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy)
-> (ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy)
-> Ord ConcurrencyPolicy
ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
ConcurrencyPolicy -> ConcurrencyPolicy -> Ordering
ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy
forall a.
Eq a =>
(a -> a -> Ordering)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> a)
-> (a -> a -> a)
-> Ord a
$ccompare :: ConcurrencyPolicy -> ConcurrencyPolicy -> Ordering
compare :: ConcurrencyPolicy -> ConcurrencyPolicy -> Ordering
$c< :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
< :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
$c<= :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
<= :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
$c> :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
> :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
$c>= :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
>= :: ConcurrencyPolicy -> ConcurrencyPolicy -> Bool
$cmax :: ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy
max :: ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy
$cmin :: ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy
min :: ConcurrencyPolicy -> ConcurrencyPolicy -> ConcurrencyPolicy
Ord, Int -> ConcurrencyPolicy -> ShowS
[ConcurrencyPolicy] -> ShowS
ConcurrencyPolicy -> String
(Int -> ConcurrencyPolicy -> ShowS)
-> (ConcurrencyPolicy -> String)
-> ([ConcurrencyPolicy] -> ShowS)
-> Show ConcurrencyPolicy
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ConcurrencyPolicy -> ShowS
showsPrec :: Int -> ConcurrencyPolicy -> ShowS
$cshow :: ConcurrencyPolicy -> String
show :: ConcurrencyPolicy -> String
$cshowList :: [ConcurrencyPolicy] -> ShowS
showList :: [ConcurrencyPolicy] -> ShowS
Show)

instance AdmissionPolicy ConcurrencyPolicy where
  policyPrefixOf :: ConcurrencyPolicy -> Text
policyPrefixOf = ConcurrencyPolicy -> Text
cpPrefix

-- | A pool named @prefix@ admitting at most @limit@ concurrent jobs per key. The cap is
-- floored at 1. Pause a pool via an override. The prefix must not contain @:@, the key
-- separator. The migration enforces this.
concurrencyPool :: Text -> Int32 -> ConcurrencyPolicy
concurrencyPool :: Text -> Int32 -> ConcurrencyPolicy
concurrencyPool Text
prefix Int32
limit = Text -> Int32 -> ConcurrencyPolicy
ConcurrencyPolicy Text
prefix (Int32 -> Int32 -> Int32
forall a. Ord a => a -> a -> a
max Int32
1 Int32
limit)

-- | A selective description of the concurrency key for a payload. Evaluation
-- returns the job key. Static inspection returns the reachable pools.
type ConcurrencyFor payload = Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey)

-- | This payload is unbounded.
noConcurrency :: ConcurrencyFor payload
noConcurrency :: forall payload. ConcurrencyFor payload
noConcurrency = Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey)
forall p payload key. Selector p payload (Maybe key)
selectNone

-- | Cap by a fixed pool, keyed by a per-job suffix (e.g. a tenant id).
concurrencyBy :: ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy :: forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy = (Text -> Text -> ConcurrencyKey)
-> ConcurrencyPolicy
-> (payload -> Text)
-> Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey)
forall p key payload.
AdmissionPolicy p =>
(Text -> Text -> key)
-> p -> (payload -> Text) -> Selector p payload (Maybe key)
selectBy Text -> Text -> ConcurrencyKey
ConcurrencyKey

-- | Cap by a fixed pool under one shared key (a single global pool).
globalConcurrency :: ConcurrencyPolicy -> Text -> ConcurrencyFor payload
globalConcurrency :: forall payload. ConcurrencyPolicy -> Text -> ConcurrencyFor payload
globalConcurrency ConcurrencyPolicy
pol Text
suffix = ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
forall payload.
ConcurrencyPolicy -> (payload -> Text) -> ConcurrencyFor payload
concurrencyBy ConcurrencyPolicy
pol (Text -> payload -> Text
forall a b. a -> b -> a
const Text
suffix)

-- | Concurrency 'selectByCase'. See 'selectByCase' for the totality requirement on @k@.
concurrencyByCase
  :: (Bounded k, Enum k, Eq k) => (payload -> k) -> (k -> ConcurrencyFor payload) -> ConcurrencyFor payload
concurrencyByCase :: forall k payload.
(Bounded k, Enum k, Eq k) =>
(payload -> k)
-> (k -> ConcurrencyFor payload) -> ConcurrencyFor payload
concurrencyByCase = (payload -> k)
-> (k -> Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey))
-> Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey)
forall k payload policy a.
(Bounded k, Enum k, Eq k) =>
(payload -> k)
-> (k -> Selector policy payload a) -> Selector policy payload a
selectByCase

-- | Run a selector against a concrete job to get its key.
runConcurrencyFor :: payload -> ConcurrencyFor payload -> Maybe ConcurrencyKey
runConcurrencyFor :: forall payload.
payload -> ConcurrencyFor payload -> Maybe ConcurrencyKey
runConcurrencyFor = payload
-> Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey)
-> Maybe ConcurrencyKey
forall policy payload a. payload -> Selector policy payload a -> a
runSelector

-- | A payload's per-job pool selection. Defaults to unbounded. Only capped
-- payloads need an instance.
class HasConcurrency payload where
  concurrencyFor :: ConcurrencyFor payload
  concurrencyFor = ConcurrencyFor payload
forall payload. ConcurrencyFor payload
noConcurrency

instance {-# OVERLAPPABLE #-} HasConcurrency payload

instance (HasConcurrency payload) => CollectFor payload ConcurrencyPolicy where
  collectFor :: Set ConcurrencyPolicy
collectFor = Selector ConcurrencyPolicy payload (Maybe ConcurrencyKey)
-> Set ConcurrencyPolicy
forall policy payload a.
Ord policy =>
Selector policy payload a -> Set policy
collectPolicies (forall payload. HasConcurrency payload => ConcurrencyFor payload
concurrencyFor @payload)

-- | Collect every pool declared across a registry's payloads, by statically
-- inspecting each payload's 'concurrencyFor'. The migration seeds these.
type RegistryConcurrencyPolicies registry = RegistryPolicies registry ConcurrencyPolicy

-- | Every distinct pool declared across the registry's payloads.
registryConcurrencyPolicies :: forall registry. (RegistryConcurrencyPolicies registry) => Set ConcurrencyPolicy
registryConcurrencyPolicies :: forall (registry :: JobPayloadRegistry).
RegistryConcurrencyPolicies registry =>
Set ConcurrencyPolicy
registryConcurrencyPolicies = forall (registry :: JobPayloadRegistry) p.
(Ord p, RegistryPolicies registry p) =>
Set p
registryPolicies @registry @ConcurrencyPolicy

-- | Each registry table paired with whether its payload declares any pool.
registryConcurrencyTables :: forall registry. (RegistryConcurrencyPolicies registry) => [(Text, Bool)]
registryConcurrencyTables :: forall (registry :: JobPayloadRegistry).
RegistryConcurrencyPolicies registry =>
[(Text, Bool)]
registryConcurrencyTables = forall (registry :: JobPayloadRegistry) p.
RegistryPolicies registry p =>
[(Text, Bool)]
registryPolicyTables @registry @ConcurrencyPolicy