{-# LANGUAGE TypeFamilies #-}

-- | Reading, combining, and storing the results of a rollup job's children.
module Arbiter.Worker.Results
  ( childResults
  , mergedChildResults
  , mergeChildResults
  , storeJobResult
  , storeEncodedResult
  , storeEncodedResults
  ) where

import Arbiter.Core.Job.Types (JobRead, parentId, primaryKey, queueName)
import Arbiter.Core.JobResult (EncodeJobResult, decodeJobResult, encodeJobResult)
import Arbiter.Core.MonadArbiter (MonadArbiter, ResultOf, getSchema)
import Arbiter.Core.Operations qualified as Ops
import Control.Monad (void)
import Data.Aeson (FromJSON, Value)
import Data.Either (partitionEithers)
import Data.Foldable (fold, foldMap')
import Data.Int (Int64)
import Data.Map.Strict (Map)
import Data.Map.Strict qualified as Map
import Data.Maybe (mapMaybe)
import Data.Text (Text)

-- | A rollup parent's immediate child results, keyed by child id, and its DLQ
-- errors, keyed by DLQ row id for 'Arbiter.Core.HighLevel.retryFromDLQ'. A
-- decode failure is returned as 'Left'.
childResults
  :: (FromJSON (ResultOf m payload), MonadArbiter m)
  => JobRead payload
  -> m (Map Int64 (Either Text (ResultOf m payload)), Map Int64 Text)
childResults :: forall (m :: * -> *) payload.
(FromJSON (ResultOf m payload), MonadArbiter m) =>
JobRead payload
-> m (Map Int64 (Either Text (ResultOf m payload)), Map Int64 Text)
childResults JobRead payload
job = do
  schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
  (results, failures, snapshot, dlqFailures) <-
    Ops.readChildResultsRaw schema (queueName job) (primaryKey job)
  let raw = Map Int64 Value
-> Map Int64 Text -> Maybe Value -> Map Int64 (Either Text Value)
Ops.mergeRawChildResults Map Int64 Value
results Map Int64 Text
failures Maybe Value
snapshot
  pure (Map.map (>>= decodeJobResult) raw, dlqFailures)

-- | 'childResults' with successfully decoded values combined through 'Monoid'.
mergedChildResults
  :: ( FromJSON (ResultOf m payload)
     , MonadArbiter m
     , Monoid (ResultOf m payload)
     )
  => JobRead payload
  -> m (ResultOf m payload, Map Int64 Text)
mergedChildResults :: forall (m :: * -> *) payload.
(FromJSON (ResultOf m payload), MonadArbiter m,
 Monoid (ResultOf m payload)) =>
JobRead payload -> m (ResultOf m payload, Map Int64 Text)
mergedChildResults JobRead payload
job = do
  (results, dlqFailures) <- JobRead payload
-> m (Map Int64 (Either Text (ResultOf m payload)), Map Int64 Text)
forall (m :: * -> *) payload.
(FromJSON (ResultOf m payload), MonadArbiter m) =>
JobRead payload
-> m (Map Int64 (Either Text (ResultOf m payload)), Map Int64 Text)
childResults JobRead payload
job
  pure (mergeChildResults results, dlqFailures)

-- | Combine successful child results, treating decode failures as 'mempty'.
mergeChildResults :: (Monoid a) => Map Int64 (Either Text a) -> a
mergeChildResults :: forall a. Monoid a => Map Int64 (Either Text a) -> a
mergeChildResults = (Either Text a -> a) -> Map Int64 (Either Text a) -> a
forall m a. Monoid m => (a -> m) -> Map Int64 a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap' Either Text a -> a
forall m. Monoid m => Either Text m -> m
forall (t :: * -> *) m. (Foldable t, Monoid m) => t m -> m
fold

-- | Store a job's result for its parent rollup, if it has one.
storeJobResult
  :: (EncodeJobResult result, MonadArbiter m)
  => Text
  -> JobRead payload
  -> result
  -> m ()
storeJobResult :: forall result (m :: * -> *) payload.
(EncodeJobResult result, MonadArbiter m) =>
Text -> JobRead payload -> result -> m ()
storeJobResult Text
schemaName JobRead payload
job = Text -> JobRead payload -> Maybe Value -> m ()
forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> JobRead payload -> Maybe Value -> m ()
storeEncodedResult Text
schemaName JobRead payload
job (Maybe Value -> m ()) -> (result -> Maybe Value) -> result -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. result -> Maybe Value
forall a. EncodeJobResult a => a -> Maybe Value
encodeJobResult

-- | 'storeJobResult' on an already-encoded result. 'Nothing' stores nothing.
storeEncodedResult
  :: (MonadArbiter m)
  => Text
  -> JobRead payload
  -> Maybe Value
  -> m ()
storeEncodedResult :: forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> JobRead payload -> Maybe Value -> m ()
storeEncodedResult Text
schemaName JobRead payload
job Maybe Value
mVal =
  case (JobRead payload -> Maybe Int64
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Maybe Int64
parentId JobRead payload
job, Maybe Value
mVal) of
    (Just Int64
pid, Just Value
val) ->
      m Int64 -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m Int64 -> m ()) -> m Int64 -> m ()
forall a b. (a -> b) -> a -> b
$ Text -> Text -> Int64 -> Int64 -> Value -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> Int64 -> Value -> m Int64
Ops.insertResult Text
schemaName (JobRead payload -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
queueName JobRead payload
job) Int64
pid (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job) Value
val
    (Maybe Int64
Nothing, Just Value
val)
      | JobRead payload -> Bool
forall payload. JobRead payload -> Bool
Ops.archivesOnAck JobRead payload
job ->
          m Int64 -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m Int64 -> m ()) -> m Int64 -> m ()
forall a b. (a -> b) -> a -> b
$ Text -> Text -> Int64 -> Value -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> Value -> m Int64
Ops.updateArchiveResult Text
schemaName (JobRead payload -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
queueName JobRead payload
job) (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
job) Value
val
    (Maybe Int64, Maybe Value)
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()

-- | 'storeEncodedResult' over a batch from one queue. One statement stores the
-- child results and one the archived roots.
storeEncodedResults
  :: (MonadArbiter m)
  => Text
  -> [(JobRead payload, Maybe Value)]
  -> m ()
storeEncodedResults :: forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> [(JobRead payload, Maybe Value)] -> m ()
storeEncodedResults Text
_ [] = () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
storeEncodedResults Text
schemaName pairs :: [(JobRead payload, Maybe Value)]
pairs@((JobRead payload
firstJob, Maybe Value
_) : [(JobRead payload, Maybe Value)]
_) = do
  let ([(Int64, Int64, Value)]
childRows, [(Int64, Value)]
rootRows) = [Either (Int64, Int64, Value) (Int64, Value)]
-> ([(Int64, Int64, Value)], [(Int64, Value)])
forall a b. [Either a b] -> ([a], [b])
partitionEithers (((JobRead payload, Maybe Value)
 -> Maybe (Either (Int64, Int64, Value) (Int64, Value)))
-> [(JobRead payload, Maybe Value)]
-> [Either (Int64, Int64, Value) (Int64, Value)]
forall a b. (a -> Maybe b) -> [a] -> [b]
mapMaybe (JobRead payload, Maybe Value)
-> Maybe (Either (Int64, Int64, Value) (Int64, Value))
forall {payload} {b}.
(JobRecord payload Int64 Text UTCTime PayloadKeys, Maybe b)
-> Maybe (Either (Int64, Int64, b) (Int64, b))
resultRow [(JobRead payload, Maybe Value)]
pairs)
      queue :: Text
queue = JobRead payload -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
queueName JobRead payload
firstJob
  m Int64 -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m Int64 -> m ()) -> m Int64 -> m ()
forall a b. (a -> b) -> a -> b
$ Text -> Text -> [(Int64, Int64, Value)] -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [(Int64, Int64, Value)] -> m Int64
Ops.insertResultsBatch Text
schemaName Text
queue [(Int64, Int64, Value)]
childRows
  m Int64 -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m Int64 -> m ()) -> m Int64 -> m ()
forall a b. (a -> b) -> a -> b
$ Text -> Text -> [(Int64, Value)] -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [(Int64, Value)] -> m Int64
Ops.updateArchiveResultsBatch Text
schemaName Text
queue [(Int64, Value)]
rootRows
  where
    resultRow :: (JobRecord payload Int64 Text UTCTime PayloadKeys, Maybe b)
-> Maybe (Either (Int64, Int64, b) (Int64, b))
resultRow (JobRecord payload Int64 Text UTCTime PayloadKeys
job, Maybe b
mVal) = do
      val <- Maybe b
mVal
      case parentId job of
        Just Int64
pid -> Either (Int64, Int64, b) (Int64, b)
-> Maybe (Either (Int64, Int64, b) (Int64, b))
forall a. a -> Maybe a
Just ((Int64, Int64, b) -> Either (Int64, Int64, b) (Int64, b)
forall a b. a -> Either a b
Left (Int64
pid, JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRecord payload Int64 Text UTCTime PayloadKeys
job, b
val))
        Maybe Int64
Nothing
          | JobRecord payload Int64 Text UTCTime PayloadKeys -> Bool
forall payload. JobRead payload -> Bool
Ops.archivesOnAck JobRecord payload Int64 Text UTCTime PayloadKeys
job -> Either (Int64, Int64, b) (Int64, b)
-> Maybe (Either (Int64, Int64, b) (Int64, b))
forall a. a -> Maybe a
Just ((Int64, b) -> Either (Int64, Int64, b) (Int64, b)
forall a b. b -> Either a b
Right (JobRecord payload Int64 Text UTCTime PayloadKeys -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRecord payload Int64 Text UTCTime PayloadKeys
job, b
val))
          | Bool
otherwise -> Maybe (Either (Int64, Int64, b) (Int64, b))
forall a. Maybe a
Nothing