{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}

-- | Cross-process coordination for database gauge scans.
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)

-- | A prepared refresh operation and its timing values. A refresh that says nothing
-- about the database yields @Right Nothing@.
data RefreshSource = RefreshSource
  { RefreshSource
-> IO (Either SomeException (Maybe (Shared Snapshot)))
runRefresh :: IO (Either SomeException (Maybe (Shared Snapshot)))
  , RefreshSource -> NominalDiffTime
minimumDelay :: NominalDiffTime
  }

-- | The gate interval over the refresh interval. Under the loop's period, so a
-- replica can win the gate on consecutive ticks.
gateIntervalFactor :: NominalDiffTime
gateIntervalFactor :: NominalDiffTime
gateIntervalFactor = NominalDiffTime
0.9

-- | Prepare a coordinated refresh operation for one queue set.
prepareRefreshSource
  :: (MonadArbiter m)
  => LogConfig
  -> (forall a. m a -> IO a)
  -> SchemaName
  -> [(TableName, [Text])]
  -> NominalDiffTime
  -> NominalDiffTime
  -- ^ How old a reading may be before it is stale.
  -> IO Bool
  -- ^ Whether the local cache is stale.
  -> 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

    -- Fixed for the registration. A long name takes a query of its own.
    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

    -- A scan that outlives the freshness window is abandoned when a fresh reading
    -- stands in for it. Only unblocks a driver that yields to the runtime.
    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"

    -- An unreadable payload falls back to a local scan when the cache is stale.
    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)