-- | High-level worker and queue registry operations.
module Arbiter.Core.HighLevel.Runtime
  ( registerWorker
  , heartbeatWorker
  , setWorkerPaused
  , markWorkerShuttingDown
  , deregisterWorker
  , listWorkers
  , sweepStaleWorkers
  , ensureQueue
  , setQueuePaused
  , getQueue
  , listQueues
  ) where

import Data.Aeson (Value)
import Data.Int (Int32, Int64)
import Data.Text (Text)
import Data.Time (NominalDiffTime)
import Data.UUID.Types (UUID)

import Arbiter.Core.MonadArbiter (MonadArbiter, getSchema)
import Arbiter.Core.Operations qualified as Ops
import Arbiter.Core.Queues (QueueRow)
import Arbiter.Core.Worker (WorkerRow)

-- | Register or refresh a worker and return its effective pause state.
registerWorker
  :: (MonadArbiter m)
  => UUID
  -> Text
  -> Maybe Text
  -> Maybe Int32
  -> NominalDiffTime
  -> Maybe Value
  -> m (Maybe Bool)
registerWorker :: forall (m :: * -> *).
MonadArbiter m =>
UUID
-> Text
-> Maybe Text
-> Maybe Int32
-> NominalDiffTime
-> Maybe Value
-> m (Maybe Bool)
registerWorker UUID
workerId Text
queue Maybe Text
host Maybe Int32
threads NominalDiffTime
staleThreshold Maybe Value
metadata =
  m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m (Maybe Bool)) -> m (Maybe Bool)
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text
-> UUID
-> Text
-> Maybe Text
-> Maybe Int32
-> NominalDiffTime
-> Maybe Value
-> m (Maybe Bool)
forall (m :: * -> *).
MonadArbiter m =>
Text
-> UUID
-> Text
-> Maybe Text
-> Maybe Int32
-> NominalDiffTime
-> Maybe Value
-> m (Maybe Bool)
Ops.registerWorker Text
schema UUID
workerId Text
queue Maybe Text
host Maybe Int32
threads NominalDiffTime
staleThreshold Maybe Value
metadata

-- | Record a heartbeat and return the worker's effective pause state.
heartbeatWorker :: (MonadArbiter m) => UUID -> m (Maybe Bool)
heartbeatWorker :: forall (m :: * -> *). MonadArbiter m => UUID -> m (Maybe Bool)
heartbeatWorker UUID
workerId = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m (Maybe Bool)) -> m (Maybe Bool)
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> UUID -> m (Maybe Bool)
forall (m :: * -> *).
MonadArbiter m =>
Text -> UUID -> m (Maybe Bool)
Ops.heartbeatWorker Text
schema UUID
workerId

-- | Set a worker's pause flag.
setWorkerPaused :: (MonadArbiter m) => UUID -> Bool -> m Int64
setWorkerPaused :: forall (m :: * -> *). MonadArbiter m => UUID -> Bool -> m Int64
setWorkerPaused UUID
workerId Bool
paused = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m Int64) -> m Int64
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> UUID -> Bool -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> UUID -> Bool -> m Int64
Ops.setWorkerPaused Text
schema UUID
workerId Bool
paused

-- | Mark a worker as gracefully draining.
markWorkerShuttingDown :: (MonadArbiter m) => UUID -> m Int64
markWorkerShuttingDown :: forall (m :: * -> *). MonadArbiter m => UUID -> m Int64
markWorkerShuttingDown UUID
workerId = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m Int64) -> m Int64
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> UUID -> m Int64
forall (m :: * -> *). MonadArbiter m => Text -> UUID -> m Int64
Ops.markWorkerShuttingDown Text
schema UUID
workerId

-- | Remove a worker registry row.
deregisterWorker :: (MonadArbiter m) => UUID -> m Int64
deregisterWorker :: forall (m :: * -> *). MonadArbiter m => UUID -> m Int64
deregisterWorker UUID
workerId = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m Int64) -> m Int64
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> UUID -> m Int64
forall (m :: * -> *). MonadArbiter m => Text -> UUID -> m Int64
Ops.deregisterWorker Text
schema UUID
workerId

-- | List workers, optionally filtered by queue and heartbeat age.
listWorkers :: (MonadArbiter m) => Maybe Text -> Maybe NominalDiffTime -> m [WorkerRow]
listWorkers :: forall (m :: * -> *).
MonadArbiter m =>
Maybe Text -> Maybe NominalDiffTime -> m [WorkerRow]
listWorkers Maybe Text
queue Maybe NominalDiffTime
liveSecs = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m [WorkerRow]) -> m [WorkerRow]
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> Maybe Text -> Maybe NominalDiffTime -> m [WorkerRow]
forall (m :: * -> *).
MonadArbiter m =>
Text -> Maybe Text -> Maybe NominalDiffTime -> m [WorkerRow]
Ops.listWorkers Text
schema Maybe Text
queue Maybe NominalDiffTime
liveSecs

-- | Delete workers older than their recorded stale threshold.
sweepStaleWorkers :: (MonadArbiter m) => m Int64
sweepStaleWorkers :: forall (m :: * -> *). MonadArbiter m => m Int64
sweepStaleWorkers = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m Int64) -> m Int64
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= Text -> m Int64
forall (m :: * -> *). MonadArbiter m => Text -> m Int64
Ops.sweepStaleWorkers

-- | Ensure a queue registry row exists.
ensureQueue :: (MonadArbiter m) => Text -> m Int64
ensureQueue :: forall (m :: * -> *). MonadArbiter m => Text -> m Int64
ensureQueue Text
queue = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m Int64) -> m Int64
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> Text -> m Int64
forall (m :: * -> *). MonadArbiter m => Text -> Text -> m Int64
Ops.ensureQueue Text
schema Text
queue

-- | Set a queue's pause flag and notify its workers.
setQueuePaused :: (MonadArbiter m) => Text -> Bool -> m Int64
setQueuePaused :: forall (m :: * -> *). MonadArbiter m => Text -> Bool -> m Int64
setQueuePaused Text
queue Bool
paused = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m Int64) -> m Int64
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> Text -> Bool -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Bool -> m Int64
Ops.setQueuePaused Text
schema Text
queue Bool
paused

-- | Get one queue registry row.
getQueue :: (MonadArbiter m) => Text -> m (Maybe QueueRow)
getQueue :: forall (m :: * -> *). MonadArbiter m => Text -> m (Maybe QueueRow)
getQueue Text
queue = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m (Maybe QueueRow)) -> m (Maybe QueueRow)
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Text
schema -> Text -> Text -> m (Maybe QueueRow)
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> m (Maybe QueueRow)
Ops.getQueue Text
schema Text
queue

-- | List all queues registered in the schema.
listQueues :: (MonadArbiter m) => m [QueueRow]
listQueues :: forall (m :: * -> *). MonadArbiter m => m [QueueRow]
listQueues = m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema m Text -> (Text -> m [QueueRow]) -> m [QueueRow]
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= Text -> m [QueueRow]
forall (m :: * -> *). MonadArbiter m => Text -> m [QueueRow]
Ops.listQueues