{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}
module Arbiter.Otel.Gauges.Coordination
( RefreshSource (..)
, prepareRefreshSource
) where
import Arbiter.Core.Job.Schema (SchemaName, TableName)
import Arbiter.Core.MonadArbiter (MonadArbiter)
import Arbiter.Core.Operations (Shared (..), gateNameFor, micros, runGatedShared)
import Arbiter.Worker (LogConfig, LogLevel (Warning), tryLog)
import Control.Exception (SomeException)
import Data.IORef (newIORef, readIORef, writeIORef)
import Data.Text (Text)
import Data.Time (NominalDiffTime)
import System.Timeout (timeout)
import UnliftIO.Exception (tryAny)
import Arbiter.Otel.Gauges.Cache (Snapshot)
import Arbiter.Otel.Gauges.Scan (scanSnapshot)
data RefreshSource = RefreshSource
{ RefreshSource
-> IO (Either SomeException (Maybe (Shared Snapshot)))
runRefresh :: IO (Either SomeException (Maybe (Shared Snapshot)))
, RefreshSource -> NominalDiffTime
minimumDelay :: NominalDiffTime
}
gateIntervalFactor :: NominalDiffTime
gateIntervalFactor :: NominalDiffTime
gateIntervalFactor = NominalDiffTime
0.9
prepareRefreshSource
:: (MonadArbiter m)
=> LogConfig
-> (forall a. m a -> IO a)
-> SchemaName
-> [(TableName, [Text])]
-> NominalDiffTime
-> NominalDiffTime
-> IO Bool
-> IO RefreshSource
prepareRefreshSource :: forall (m :: * -> *).
MonadArbiter m =>
LogConfig
-> (forall a. m a -> IO a)
-> SchemaName
-> [(SchemaName, [SchemaName])]
-> NominalDiffTime
-> NominalDiffTime
-> IO Bool
-> IO RefreshSource
prepareRefreshSource LogConfig
logCfg forall a. m a -> IO a
runDb SchemaName
schema [(SchemaName, [SchemaName])]
queueKinds NominalDiffTime
refreshInterval NominalDiffTime
freshnessWindow IO Bool
cacheIsStale = do
gateCache <- Maybe SchemaName -> IO (IORef (Maybe SchemaName))
forall a. a -> IO (IORef a)
newIORef Maybe SchemaName
forall a. Maybe a
Nothing
pure
RefreshSource
{ runRefresh = refresh gateCache
, minimumDelay = gateInterval
}
where
gateInterval :: NominalDiffTime
gateInterval = NominalDiffTime
gateIntervalFactor NominalDiffTime -> NominalDiffTime -> NominalDiffTime
forall a. Num a => a -> a -> a
* NominalDiffTime
refreshInterval
queueTables :: [SchemaName]
queueTables = ((SchemaName, [SchemaName]) -> SchemaName)
-> [(SchemaName, [SchemaName])] -> [SchemaName]
forall a b. (a -> b) -> [a] -> [b]
map (SchemaName, [SchemaName]) -> SchemaName
forall a b. (a, b) -> a
fst [(SchemaName, [SchemaName])]
queueKinds
scan :: m Snapshot
scan = NominalDiffTime
-> SchemaName -> [(SchemaName, [SchemaName])] -> m Snapshot
forall (m :: * -> *).
MonadArbiter m =>
NominalDiffTime
-> SchemaName -> [(SchemaName, [SchemaName])] -> m Snapshot
scanSnapshot NominalDiffTime
freshnessWindow SchemaName
schema [(SchemaName, [SchemaName])]
queueKinds
refresh :: IORef (Maybe SchemaName)
-> IO (Either SomeException (Maybe (Shared Snapshot)))
refresh IORef (Maybe SchemaName)
gateCache = IO SchemaName -> IO (Either SomeException SchemaName)
forall (m :: * -> *) a.
MonadUnliftIO m =>
m a -> m (Either SomeException a)
tryAny (IORef (Maybe SchemaName) -> IO SchemaName
resolveGate IORef (Maybe SchemaName)
gateCache) IO (Either SomeException SchemaName)
-> (Either SomeException SchemaName
-> IO (Either SomeException (Maybe (Shared Snapshot))))
-> IO (Either SomeException (Maybe (Shared Snapshot)))
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= (SomeException
-> IO (Either SomeException (Maybe (Shared Snapshot))))
-> (SchemaName
-> IO (Either SomeException (Maybe (Shared Snapshot))))
-> Either SomeException SchemaName
-> IO (Either SomeException (Maybe (Shared Snapshot)))
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either (Either SomeException (Maybe (Shared Snapshot))
-> IO (Either SomeException (Maybe (Shared Snapshot)))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Either SomeException (Maybe (Shared Snapshot))
-> IO (Either SomeException (Maybe (Shared Snapshot))))
-> (SomeException
-> Either SomeException (Maybe (Shared Snapshot)))
-> SomeException
-> IO (Either SomeException (Maybe (Shared Snapshot)))
forall b c a. (b -> c) -> (a -> b) -> a -> c
. SomeException -> Either SomeException (Maybe (Shared Snapshot))
forall a b. a -> Either a b
Left) SchemaName -> IO (Either SomeException (Maybe (Shared Snapshot)))
sharedScan
sharedScan :: SchemaName -> IO (Either SomeException (Maybe (Shared Snapshot)))
sharedScan SchemaName
gate =
IO (Maybe (Shared Snapshot))
-> IO (Either SomeException (Maybe (Shared Snapshot)))
bounded (m (Maybe (Shared Snapshot)) -> IO (Maybe (Shared Snapshot))
forall a. m a -> IO a
runDb (SchemaName
-> SchemaName
-> NominalDiffTime
-> NominalDiffTime
-> m Snapshot
-> m (Maybe (Shared Snapshot))
forall a (m :: * -> *).
(FromJSON a, MonadArbiter m, ToJSON a) =>
SchemaName
-> SchemaName
-> NominalDiffTime
-> NominalDiffTime
-> m a
-> m (Maybe (Shared a))
runGatedShared SchemaName
schema SchemaName
gate NominalDiffTime
gateInterval NominalDiffTime
freshnessWindow m Snapshot
scan)) IO (Either SomeException (Maybe (Shared Snapshot)))
-> (Either SomeException (Maybe (Shared Snapshot))
-> IO (Either SomeException (Maybe (Shared Snapshot))))
-> IO (Either SomeException (Maybe (Shared Snapshot)))
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
Right (Just (Unreadable SchemaName
why)) -> SchemaName -> IO (Either SomeException (Maybe (Shared Snapshot)))
localScan SchemaName
why
Either SomeException (Maybe (Shared Snapshot))
outcome -> Either SomeException (Maybe (Shared Snapshot))
-> IO (Either SomeException (Maybe (Shared Snapshot)))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Either SomeException (Maybe (Shared Snapshot))
outcome
resolveGate :: IORef (Maybe SchemaName) -> IO SchemaName
resolveGate IORef (Maybe SchemaName)
gateCache =
IORef (Maybe SchemaName) -> IO (Maybe SchemaName)
forall a. IORef a -> IO a
readIORef IORef (Maybe SchemaName)
gateCache
IO (Maybe SchemaName)
-> (Maybe SchemaName -> IO SchemaName) -> IO SchemaName
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= IO SchemaName
-> (SchemaName -> IO SchemaName)
-> Maybe SchemaName
-> IO SchemaName
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (m SchemaName -> IO SchemaName
forall a. m a -> IO a
runDb (SchemaName -> [SchemaName] -> m SchemaName
forall (m :: * -> *).
MonadArbiter m =>
SchemaName -> [SchemaName] -> m SchemaName
gateNameFor SchemaName
"refresh-gauges" [SchemaName]
queueTables) IO SchemaName -> (SchemaName -> IO SchemaName) -> IO SchemaName
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \SchemaName
gate -> SchemaName
gate SchemaName -> IO () -> IO SchemaName
forall a b. a -> IO b -> IO a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ IORef (Maybe SchemaName) -> Maybe SchemaName -> IO ()
forall a. IORef a -> a -> IO ()
writeIORef IORef (Maybe SchemaName)
gateCache (SchemaName -> Maybe SchemaName
forall a. a -> Maybe a
Just SchemaName
gate)) SchemaName -> IO SchemaName
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure
bounded :: IO (Maybe (Shared Snapshot)) -> IO (Either SomeException (Maybe (Shared Snapshot)))
bounded :: IO (Maybe (Shared Snapshot))
-> IO (Either SomeException (Maybe (Shared Snapshot)))
bounded IO (Maybe (Shared Snapshot))
act = do
stale <- IO Bool
cacheIsStale
tryAny (if stale then Just <$> act else timeout (micros freshnessWindow) act)
>>= traverse (maybe (Nothing <$ tryLog logCfg Warning abandoned) pure)
abandoned :: SchemaName
abandoned = SchemaName
"Gauge scan outlived the freshness window, abandoned"
localScan :: SchemaName -> IO (Either SomeException (Maybe (Shared Snapshot)))
localScan SchemaName
why = do
LogConfig -> LogLevel -> SchemaName -> IO ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> SchemaName -> m ()
tryLog LogConfig
logCfg LogLevel
Warning (SchemaName
"Shared gauge payload unreadable, scanning locally: " SchemaName -> SchemaName -> SchemaName
forall a. Semigroup a => a -> a -> a
<> SchemaName
why)
stale <- IO Bool
cacheIsStale
if stale then tryAny (Just . Ran <$> runDb scan) else pure (Right Nothing)