{-# LANGUAGE DeriveFunctor #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Operations.Gates
( runGated
, runGatedBounded
, runGatedShared
, runGatedState
, runGatedStateBounded
, setLocalStatementTimeout
, gateNameFor
, micros
, Shared (..)
) where
import Control.Exception qualified as E
import Control.Monad (void)
import Data.Aeson (FromJSON, ToJSON, Value, parseJSON, toJSON)
import Data.Aeson.Types (parseEither, parseMaybe)
import Data.Bifunctor (second)
import Data.List (sort)
import Data.Maybe (fromMaybe, listToMaybe)
import Data.Text (Text)
import Data.Text qualified as T
import Data.Time (NominalDiffTime)
import UnliftIO (tryAny, withRunInIO)
import UnliftIO qualified as UIO
import Arbiter.Core.Job.Schema (SchemaName)
import Arbiter.Core.MonadArbiter (MonadArbiter, withDbTransaction)
import Arbiter.Core.MonadArbiter qualified as MA
import Arbiter.Core.Sql.Gates qualified as Sql
import Arbiter.Core.Sql.QQ qualified as QQ
setLocalStatementTimeout :: (MonadArbiter m) => NominalDiffTime -> m ()
setLocalStatementTimeout :: forall (m :: * -> *). MonadArbiter m => NominalDiffTime -> m ()
setLocalStatementTimeout NominalDiffTime
limit =
let millis :: Int
millis = Double -> Int
forall b. Integral b => Double -> b
forall a b. (RealFrac a, Integral b) => a -> b
ceiling (NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac NominalDiffTime
limit Double -> Double -> Double
forall a. Num a => a -> a -> a
* Double
1000 :: Double) :: Int
millisText :: Text
millisText = String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
millis)
in m [Text] -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m [Text] -> m ()) -> m [Text] -> m ()
forall a b. (a -> b) -> a -> b
$
Query Text -> m [Text]
forall a. Query a -> m [a]
forall (m :: * -> *) a. MonadArbiter m => Query a -> m [a]
MA.executeQuery
[QQ.sql|SELECT set_config('statement_timeout', '${millisText}', true) AS @{set_config :: CText}|]
runGatedBounded :: (MonadArbiter m) => SchemaName -> Text -> NominalDiffTime -> NominalDiffTime -> m a -> m (Maybe a)
runGatedBounded :: forall (m :: * -> *) a.
MonadArbiter m =>
Text
-> Text -> NominalDiffTime -> NominalDiffTime -> m a -> m (Maybe a)
runGatedBounded Text
schemaName Text
task NominalDiffTime
interval NominalDiffTime
limit m a
work =
Text -> Text -> NominalDiffTime -> m a -> m (Maybe a)
forall (m :: * -> *) a.
MonadArbiter m =>
Text -> Text -> NominalDiffTime -> m a -> m (Maybe a)
runGated Text
schemaName Text
task NominalDiffTime
interval (NominalDiffTime -> m ()
forall (m :: * -> *). MonadArbiter m => NominalDiffTime -> m ()
setLocalStatementTimeout NominalDiffTime
limit m () -> m a -> m a
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> m a
work)
runGated
:: (MonadArbiter m)
=> SchemaName
-> Text
-> NominalDiffTime
-> m a
-> m (Maybe a)
runGated :: forall (m :: * -> *) a.
MonadArbiter m =>
Text -> Text -> NominalDiffTime -> m a -> m (Maybe a)
runGated Text
schemaName Text
task NominalDiffTime
interval m a
work =
Text
-> Text
-> NominalDiffTime
-> (Maybe Value -> m (a, Maybe Value))
-> m (Maybe a)
forall (m :: * -> *) a.
MonadArbiter m =>
Text
-> Text
-> NominalDiffTime
-> (Maybe Value -> m (a, Maybe Value))
-> m (Maybe a)
runGatedInner Text
schemaName Text
task NominalDiffTime
interval (m (a, Maybe Value) -> Maybe Value -> m (a, Maybe Value)
forall a b. a -> b -> a
const ((,Maybe Value
forall a. Maybe a
Nothing) (a -> (a, Maybe Value)) -> m a -> m (a, Maybe Value)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> m a
work))
runGatedState
:: (FromJSON s, MonadArbiter m, ToJSON s)
=> SchemaName
-> Text
-> NominalDiffTime
-> (Maybe s -> m (a, s))
-> m (Maybe a)
runGatedState :: forall s (m :: * -> *) a.
(FromJSON s, MonadArbiter m, ToJSON s) =>
Text
-> Text -> NominalDiffTime -> (Maybe s -> m (a, s)) -> m (Maybe a)
runGatedState Text
schemaName Text
task NominalDiffTime
interval Maybe s -> m (a, s)
work =
Text
-> Text
-> NominalDiffTime
-> (Maybe Value -> m (a, Maybe Value))
-> m (Maybe a)
forall (m :: * -> *) a.
MonadArbiter m =>
Text
-> Text
-> NominalDiffTime
-> (Maybe Value -> m (a, Maybe Value))
-> m (Maybe a)
runGatedInner Text
schemaName Text
task NominalDiffTime
interval (((a, s) -> (a, Maybe Value)) -> m (a, s) -> m (a, Maybe Value)
forall a b. (a -> b) -> m a -> m b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap ((s -> Maybe Value) -> (a, s) -> (a, Maybe Value)
forall b c a. (b -> c) -> (a, b) -> (a, c)
forall (p :: * -> * -> *) b c a.
Bifunctor p =>
(b -> c) -> p a b -> p a c
second (Value -> Maybe Value
forall a. a -> Maybe a
Just (Value -> Maybe Value) -> (s -> Value) -> s -> Maybe Value
forall b c a. (b -> c) -> (a -> b) -> a -> c
. s -> Value
forall a. ToJSON a => a -> Value
toJSON)) (m (a, s) -> m (a, Maybe Value))
-> (Maybe Value -> m (a, s)) -> Maybe Value -> m (a, Maybe Value)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Maybe s -> m (a, s)
work (Maybe s -> m (a, s))
-> (Maybe Value -> Maybe s) -> Maybe Value -> m (a, s)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Maybe Value -> (Value -> Maybe s) -> Maybe s
forall a b. Maybe a -> (a -> Maybe b) -> Maybe b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= (Value -> Parser s) -> Value -> Maybe s
forall a b. (a -> Parser b) -> a -> Maybe b
parseMaybe Value -> Parser s
forall a. FromJSON a => Value -> Parser a
parseJSON))
runGatedStateBounded
:: (FromJSON s, MonadArbiter m, ToJSON s)
=> SchemaName
-> Text
-> NominalDiffTime
-> NominalDiffTime
-> (Maybe s -> m (a, s))
-> m (Maybe a)
runGatedStateBounded :: forall s (m :: * -> *) a.
(FromJSON s, MonadArbiter m, ToJSON s) =>
Text
-> Text
-> NominalDiffTime
-> NominalDiffTime
-> (Maybe s -> m (a, s))
-> m (Maybe a)
runGatedStateBounded Text
schemaName Text
task NominalDiffTime
interval NominalDiffTime
limit Maybe s -> m (a, s)
work =
Text
-> Text -> NominalDiffTime -> (Maybe s -> m (a, s)) -> m (Maybe a)
forall s (m :: * -> *) a.
(FromJSON s, MonadArbiter m, ToJSON s) =>
Text
-> Text -> NominalDiffTime -> (Maybe s -> m (a, s)) -> m (Maybe a)
runGatedState Text
schemaName Text
task NominalDiffTime
interval (\Maybe s
state -> NominalDiffTime -> m ()
forall (m :: * -> *). MonadArbiter m => NominalDiffTime -> m ()
setLocalStatementTimeout NominalDiffTime
limit m () -> m (a, s) -> m (a, s)
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Maybe s -> m (a, s)
work Maybe s
state)
runGatedInner
:: (MonadArbiter m)
=> SchemaName
-> Text
-> NominalDiffTime
-> (Maybe Value -> m (a, Maybe Value))
-> m (Maybe a)
runGatedInner :: forall (m :: * -> *) a.
MonadArbiter m =>
Text
-> Text
-> NominalDiffTime
-> (Maybe Value -> m (a, Maybe Value))
-> m (Maybe a)
runGatedInner Text
schemaName Text
task NominalDiffTime
interval Maybe Value -> m (a, Maybe Value)
work = do
_ <-
Query () -> m Int64
forall a. Query a -> m Int64
forall (m :: * -> *) a. MonadArbiter m => Query a -> m Int64
MA.executeStatement
(Text -> Text -> Query ()
Sql.ensureGateRowSQL Text
schemaName Text
task)
gateOpen <- checkGateOuter
if not gateOpen
then pure Nothing
else withDbTransaction $ tryClaimGate >>= traverse ran
where
intervalSecs :: Double
intervalSecs = NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac NominalDiffTime
interval :: Double
checkGateOuter :: m Bool
checkGateOuter = do
rows <- Query Bool -> m [Bool]
forall a. Query a -> m [a]
forall (m :: * -> *) a. MonadArbiter m => Query a -> m [a]
MA.executeQuery (Text -> Double -> Text -> Query Bool
Sql.checkGateSQL Text
schemaName Double
intervalSecs Text
task)
pure $ fromMaybe True (listToMaybe rows)
tryClaimGate :: m (Maybe (Maybe Value))
tryClaimGate = [Maybe Value] -> Maybe (Maybe Value)
forall a. [a] -> Maybe a
listToMaybe ([Maybe Value] -> Maybe (Maybe Value))
-> m [Maybe Value] -> m (Maybe (Maybe Value))
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Query (Maybe Value) -> m [Maybe Value]
forall a. Query a -> m [a]
forall (m :: * -> *) a. MonadArbiter m => Query a -> m [a]
MA.executeQuery (Text -> Text -> Double -> Query (Maybe Value)
Sql.tryClaimGateSQL Text
schemaName Text
task Double
intervalSecs)
ran :: Maybe Value -> m a
ran Maybe Value
state = do
(result, next) <- Maybe Value -> m (a, Maybe Value)
work Maybe Value
state
result <$ MA.executeStatement (maybe (Sql.bumpGateSQL schemaName task) (Sql.bumpGateStateSQL schemaName task) next)
gateNameFor :: (MonadArbiter m) => Text -> [Text] -> m Text
gateNameFor :: forall (m :: * -> *). MonadArbiter m => Text -> [Text] -> m Text
gateNameFor Text
prefix [Text]
parts
| Text -> Int
T.length Text
joined Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
<= Int
maxGateNameLength = Text -> m Text
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Text
prefix Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
":" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
joined)
| Bool
otherwise = do
rows <- Query Text -> m [Text]
forall a. Query a -> m [a]
forall (m :: * -> *) a. MonadArbiter m => Query a -> m [a]
MA.executeQuery (Text -> Query Text
Sql.gateNameDigestSQL Text
joined)
pure (prefix <> ":#" <> fromMaybe joined (listToMaybe rows))
where
joined :: Text
joined = Text -> [Text] -> Text
T.intercalate Text
"," ([Text] -> [Text]
forall a. Ord a => [a] -> [a]
sort [Text]
parts)
maxGateNameLength :: Int
maxGateNameLength :: Int
maxGateNameLength = Int
200
data Shared a
=
Ran a
|
Published Double a
|
Unreadable Text
deriving stock (Shared a -> Shared a -> Bool
(Shared a -> Shared a -> Bool)
-> (Shared a -> Shared a -> Bool) -> Eq (Shared a)
forall a. Eq a => Shared a -> Shared a -> Bool
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: forall a. Eq a => Shared a -> Shared a -> Bool
== :: Shared a -> Shared a -> Bool
$c/= :: forall a. Eq a => Shared a -> Shared a -> Bool
/= :: Shared a -> Shared a -> Bool
Eq, (forall a b. (a -> b) -> Shared a -> Shared b)
-> (forall a b. a -> Shared b -> Shared a) -> Functor Shared
forall a b. a -> Shared b -> Shared a
forall a b. (a -> b) -> Shared a -> Shared b
forall (f :: * -> *).
(forall a b. (a -> b) -> f a -> f b)
-> (forall a b. a -> f b -> f a) -> Functor f
$cfmap :: forall a b. (a -> b) -> Shared a -> Shared b
fmap :: forall a b. (a -> b) -> Shared a -> Shared b
$c<$ :: forall a b. a -> Shared b -> Shared a
<$ :: forall a b. a -> Shared b -> Shared a
Functor, Int -> Shared a -> ShowS
[Shared a] -> ShowS
Shared a -> String
(Int -> Shared a -> ShowS)
-> (Shared a -> String) -> ([Shared a] -> ShowS) -> Show (Shared a)
forall a. Show a => Int -> Shared a -> ShowS
forall a. Show a => [Shared a] -> ShowS
forall a. Show a => Shared a -> String
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: forall a. Show a => Int -> Shared a -> ShowS
showsPrec :: Int -> Shared a -> ShowS
$cshow :: forall a. Show a => Shared a -> String
show :: Shared a -> String
$cshowList :: forall a. Show a => [Shared a] -> ShowS
showList :: [Shared a] -> ShowS
Show)
runGatedShared
:: (FromJSON a, MonadArbiter m, ToJSON a)
=> SchemaName
-> Text
-> NominalDiffTime
-> NominalDiffTime
-> m a
-> m (Maybe (Shared a))
runGatedShared :: forall a (m :: * -> *).
(FromJSON a, MonadArbiter m, ToJSON a) =>
Text
-> Text
-> NominalDiffTime
-> NominalDiffTime
-> m a
-> m (Maybe (Shared a))
runGatedShared Text
schemaName Text
task NominalDiffTime
interval NominalDiffTime
maxAge m a
work =
Query (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
-> m [(Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)]
forall a. Query a -> m [a]
forall (m :: * -> *) a. MonadArbiter m => Query a -> m [a]
MA.executeQuery Query (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
claimOrRead m [(Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)]
-> ([(Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)]
-> m (Maybe (Shared a)))
-> m (Maybe (Shared a))
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= m (Maybe (Shared a))
-> ((Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
-> m (Maybe (Shared a)))
-> Maybe (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
-> m (Maybe (Shared a))
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (Maybe (Shared a) -> m (Maybe (Shared a))
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe (Shared a)
forall a. Maybe a
Nothing) (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
-> m (Maybe (Shared a))
shared (Maybe (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
-> m (Maybe (Shared a)))
-> ([(Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)]
-> Maybe (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double))
-> [(Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)]
-> m (Maybe (Shared a))
forall b c a. (b -> c) -> (a -> b) -> a -> c
. [(Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)]
-> Maybe (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
forall a. [a] -> Maybe a
listToMaybe
where
claimOrRead :: Query (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
claimOrRead =
Text
-> Text
-> Double
-> Double
-> Query (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
Sql.claimOrReadGateSQL Text
schemaName Text
task (NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac NominalDiffTime
interval) (NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac NominalDiffTime
maxAge)
shared :: (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
-> m (Maybe (Shared a))
shared (Maybe UTCTime
mClaimedAt, Maybe UTCTime
mPrevious, Maybe Value
mPayload, Maybe Double
mAge) = case (Maybe UTCTime
mClaimedAt, Maybe UTCTime
mPrevious) of
(Just UTCTime
claimedAt, Just UTCTime
previous) -> Shared a -> Maybe (Shared a)
forall a. a -> Maybe a
Just (Shared a -> Maybe (Shared a))
-> (a -> Shared a) -> a -> Maybe (Shared a)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. a -> Shared a
forall a. a -> Shared a
Ran (a -> Maybe (Shared a)) -> m a -> m (Maybe (Shared a))
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> UTCTime -> UTCTime -> m a
publish UTCTime
claimedAt UTCTime
previous
(Maybe UTCTime, Maybe UTCTime)
_ -> Maybe (Shared a) -> m (Maybe (Shared a))
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Value -> Double -> Shared a
decoded (Value -> Double -> Shared a)
-> Maybe Value -> Maybe (Double -> Shared a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Value
mPayload Maybe (Double -> Shared a) -> Maybe Double -> Maybe (Shared a)
forall a b. Maybe (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Maybe Double
mAge)
publish :: UTCTime -> UTCTime -> m a
publish UTCTime
claimedAt UTCTime
previous = ((forall a. m a -> IO a) -> IO a) -> m a
forall b. ((forall a. m a -> IO a) -> IO b) -> m b
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> IO a) -> IO b) -> m b
withRunInIO (((forall a. m a -> IO a) -> IO a) -> m a)
-> ((forall a. m a -> IO a) -> IO a) -> m a
forall a b. (a -> b) -> a -> b
$ \forall a. m a -> IO a
run ->
m a -> IO a
forall a. m a -> IO a
run
(m a
work m a -> (a -> m a) -> m a
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \a
result -> a
result a -> m Int64 -> m a
forall a b. a -> m b -> m a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Query () -> m Int64
forall a. Query a -> m Int64
forall (m :: * -> *) a. MonadArbiter m => Query a -> m Int64
MA.executeStatement (Text -> Value -> UTCTime -> Text -> Query ()
Sql.setGateMetadataSQL Text
schemaName (a -> Value
forall a. ToJSON a => a -> Value
toJSON a
result) UTCTime
claimedAt Text
task))
IO a -> IO () -> IO a
forall a b. IO a -> IO b -> IO a
`E.onException` m () -> IO ()
forall a. m a -> IO a
run (UTCTime -> UTCTime -> m ()
reopen UTCTime
claimedAt UTCTime
previous)
reopen :: UTCTime -> UTCTime -> m ()
reopen UTCTime
claimedAt UTCTime
previous =
m (Either SomeException (Maybe Int64)) -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void
(m (Maybe Int64) -> m (Either SomeException (Maybe Int64))
forall (m :: * -> *) a.
MonadUnliftIO m =>
m a -> m (Either SomeException a)
tryAny (Int -> m Int64 -> m (Maybe Int64)
forall (m :: * -> *) a.
MonadUnliftIO m =>
Int -> m a -> m (Maybe a)
UIO.timeout (NominalDiffTime -> Int
micros NominalDiffTime
interval) (Query () -> m Int64
forall a. Query a -> m Int64
forall (m :: * -> *) a. MonadArbiter m => Query a -> m Int64
MA.executeStatement (Text -> Text -> UTCTime -> UTCTime -> Query ()
Sql.releaseGateSQL Text
schemaName Text
task UTCTime
claimedAt UTCTime
previous))))
decoded :: Value -> Double -> Shared a
decoded Value
value Double
age = (String -> Shared a)
-> (a -> Shared a) -> Either String a -> Shared a
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either (Text -> Shared a
forall a. Text -> Shared a
Unreadable (Text -> Shared a) -> (String -> Text) -> String -> Shared a
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (\String
err -> Text
task Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" gate payload: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> String -> Text
T.pack String
err)) (Double -> a -> Shared a
forall a. Double -> a -> Shared a
Published Double
age) ((Value -> Parser a) -> Value -> Either String a
forall a b. (a -> Parser b) -> a -> Either String b
parseEither Value -> Parser a
forall a. FromJSON a => Value -> Parser a
parseJSON Value
value)
micros :: NominalDiffTime -> Int
micros :: NominalDiffTime -> Int
micros NominalDiffTime
seconds = Double -> Int
forall b. Integral b => Double -> b
forall a b. (RealFrac a, Integral b) => a -> b
round (NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac NominalDiffTime
seconds Double -> Double -> Double
forall a. Num a => a -> a -> a
* Double
1_000_000 :: Double)