{-# LANGUAGE DeriveAnyClass #-}
{-# LANGUAGE OverloadedStrings #-}

-- | Atomic parent-child job trees. Children run first. Their suspended finalizer
-- becomes claimable when no children remain in the main queue.
-- Child results are transient and are deleted when the finalizer is acked.
-- Finalizers must persist any results that need to outlive the tree.
module Arbiter.Core.JobTree
  ( -- * Tree type
    JobTree

    -- * Smart constructors
  , leaf
  , rollup

    -- * Operators
  , (<~~)

    -- * Interpreter
  , insertJobTree
  ) where

import Control.Exception (Exception)
import Control.Monad (when)
import Data.Aeson (Object, Value (..))
import Data.Int (Int64)
import Data.List.NonEmpty (NonEmpty (..))
import Data.List.NonEmpty qualified as NE
import Data.Text (Text)
import UnliftIO.Exception qualified as UE

import Arbiter.Core.Job.Types
  ( JobPayload
  , JobRead
  , JobWrite
  , primaryKey
  )
import Arbiter.Core.MonadArbiter (MonadArbiter (..))
import Arbiter.Core.Operations qualified as Ops
import Arbiter.Core.Trace (markSpanError, withPublishSpan)

-- | Aborts a tree insertion transaction. 'insertJobTree' catches it and returns @Left@.
newtype TreeInsertFailed = TreeInsertFailed Text
  deriving stock (Int -> TreeInsertFailed -> ShowS
[TreeInsertFailed] -> ShowS
TreeInsertFailed -> String
(Int -> TreeInsertFailed -> ShowS)
-> (TreeInsertFailed -> String)
-> ([TreeInsertFailed] -> ShowS)
-> Show TreeInsertFailed
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> TreeInsertFailed -> ShowS
showsPrec :: Int -> TreeInsertFailed -> ShowS
$cshow :: TreeInsertFailed -> String
show :: TreeInsertFailed -> String
$cshowList :: [TreeInsertFailed] -> ShowS
showList :: [TreeInsertFailed] -> ShowS
Show)
  deriving anyclass (Show TreeInsertFailed
Typeable TreeInsertFailed
(Typeable TreeInsertFailed, Show TreeInsertFailed) =>
(TreeInsertFailed -> SomeException)
-> (SomeException -> Maybe TreeInsertFailed)
-> (TreeInsertFailed -> String)
-> (TreeInsertFailed -> Bool)
-> Exception TreeInsertFailed
SomeException -> Maybe TreeInsertFailed
TreeInsertFailed -> Bool
TreeInsertFailed -> String
TreeInsertFailed -> SomeException
forall e.
(Typeable e, Show e) =>
(e -> SomeException)
-> (SomeException -> Maybe e)
-> (e -> String)
-> (e -> Bool)
-> Exception e
$ctoException :: TreeInsertFailed -> SomeException
toException :: TreeInsertFailed -> SomeException
$cfromException :: SomeException -> Maybe TreeInsertFailed
fromException :: SomeException -> Maybe TreeInsertFailed
$cdisplayException :: TreeInsertFailed -> String
displayException :: TreeInsertFailed -> String
$cbacktraceDesired :: TreeInsertFailed -> Bool
backtraceDesired :: TreeInsertFailed -> Bool
Exception)

-- | A tree of jobs. Leaves are single jobs. Finalizers are parents with
-- children that run immediately while the parent waits for completion.
data JobTree payload
  = -- | A single job with no children.
    Leaf (JobWrite payload)
  | -- | A finalizer job with children. The finalizer is suspended until
    -- all children complete, then it becomes claimable for a completion round.
    Finalizer (JobWrite payload) (NonEmpty (JobTree payload))

-- | A single job with no children.
leaf :: JobWrite payload -> JobTree payload
leaf :: forall payload. JobWrite payload -> JobTree payload
leaf = JobWrite payload -> JobTree payload
forall payload. JobWrite payload -> JobTree payload
Leaf

-- | A finalizer running once no child of it is left in the main queue. Nested rollups
-- do not merge on their own. An intermediate finalizer returns the merged value for
-- results to travel upward.
--
-- @
-- rollup (defaultJob root)
--   ( leaf (defaultJob leaf1)
--   :| [leaf (defaultJob leaf2)]
--   )
-- @
rollup :: JobWrite payload -> NonEmpty (JobTree payload) -> JobTree payload
rollup :: forall payload.
JobWrite payload -> NonEmpty (JobTree payload) -> JobTree payload
rollup = JobWrite payload -> NonEmpty (JobTree payload) -> JobTree payload
forall payload.
JobWrite payload -> NonEmpty (JobTree payload) -> JobTree payload
Finalizer

-- | The empty rollup snapshot @{}@ every finalizer is inserted with. A non-null snapshot
-- marks a job as a finalizer. The value is overwritten with the merged child results
-- before a DLQ move.
emptyState :: Value
emptyState :: Value
emptyState = Object -> Value
Object (Object
forall a. Monoid a => a
mempty :: Object)

-- | Infix 'rollup' for leaf-only children.
--
-- @
-- defaultJob reducer \<~~ (defaultJob mapper1 :| [defaultJob mapper2])
-- @
infixr 6 <~~

(<~~) :: JobWrite payload -> NonEmpty (JobWrite payload) -> JobTree payload
JobWrite payload
parent <~~ :: forall payload.
JobWrite payload -> NonEmpty (JobWrite payload) -> JobTree payload
<~~ NonEmpty (JobWrite payload)
children = JobWrite payload -> NonEmpty (JobTree payload) -> JobTree payload
forall payload.
JobWrite payload -> NonEmpty (JobTree payload) -> JobTree payload
Finalizer JobWrite payload
parent ((JobWrite payload -> JobTree payload)
-> NonEmpty (JobWrite payload) -> NonEmpty (JobTree payload)
forall a b. (a -> b) -> NonEmpty a -> NonEmpty b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap JobWrite payload -> JobTree payload
forall payload. JobWrite payload -> JobTree payload
Leaf NonEmpty (JobWrite payload)
children)

-- | Insert a 'JobTree' in one transaction, returning every inserted job root-first.
-- @Left@ on any failure, such as a dedup conflict, with nothing committed.
insertJobTree
  :: forall m payload
   . (JobPayload payload, MonadArbiter m)
  => Text
  -- ^ PostgreSQL schema name
  -> Text
  -- ^ Table name
  -> JobTree payload
  -> m (Either Text (NonEmpty (JobRead payload)))
insertJobTree :: forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> JobTree payload
-> m (Either Text (NonEmpty (JobRead payload)))
insertJobTree Text
schemaName Text
tableName JobTree payload
tree =
  Text
-> [JobWrite payload]
-> m (Either Text (NonEmpty (JobRead payload)))
-> m (Either Text (NonEmpty (JobRead payload)))
forall payload (m :: * -> *) a.
(HasKind payload, MonadUnliftIO m) =>
Text -> [JobWrite payload] -> m a -> m a
withPublishSpan Text
tableName (JobTree payload -> [JobWrite payload]
treeWrites JobTree payload
tree) (m (Either Text (NonEmpty (JobRead payload)))
 -> m (Either Text (NonEmpty (JobRead payload))))
-> m (Either Text (NonEmpty (JobRead payload)))
-> m (Either Text (NonEmpty (JobRead payload)))
forall a b. (a -> b) -> a -> b
$ do
    inserted <- m (NonEmpty (JobRead payload))
-> m (Either TreeInsertFailed (NonEmpty (JobRead payload)))
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> m (Either e a)
UE.try (m (NonEmpty (JobRead payload))
 -> m (Either TreeInsertFailed (NonEmpty (JobRead payload))))
-> m (NonEmpty (JobRead payload))
-> m (Either TreeInsertFailed (NonEmpty (JobRead payload)))
forall a b. (a -> b) -> a -> b
$ do
      stamp <- m (TraceStamp payload)
forall (m :: * -> *) payload. MonadIO m => m (TraceStamp payload)
Ops.traceStamp
      withDbTransaction $ go stamp Nothing (rootSuspended tree) tree
    either (\(TreeInsertFailed Text
msg) -> Text -> Either Text (NonEmpty (JobRead payload))
forall a b. a -> Either a b
Left Text
msg Either Text (NonEmpty (JobRead payload))
-> m () -> m (Either Text (NonEmpty (JobRead payload)))
forall a b. a -> m b -> m a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Text -> m ()
forall (m :: * -> *). MonadIO m => Text -> m ()
markSpanError Text
msg) (pure . Right) inserted
  where
    -- Finalizer roots are suspended (waiting for children to complete).
    rootSuspended :: JobTree payload -> Bool
    rootSuspended :: JobTree payload -> Bool
rootSuspended (Finalizer JobWrite payload
_ NonEmpty (JobTree payload)
_) = Bool
True
    rootSuspended JobTree payload
_ = Bool
False

    treeWrites :: JobTree payload -> [JobWrite payload]
    treeWrites :: JobTree payload -> [JobWrite payload]
treeWrites (Leaf JobWrite payload
job) = [JobWrite payload
job]
    treeWrites (Finalizer JobWrite payload
job NonEmpty (JobTree payload)
children) = JobWrite payload
job JobWrite payload -> [JobWrite payload] -> [JobWrite payload]
forall a. a -> [a] -> [a]
: (JobTree payload -> [JobWrite payload])
-> NonEmpty (JobTree payload) -> [JobWrite payload]
forall m a. Monoid m => (a -> m) -> NonEmpty a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap JobTree payload -> [JobWrite payload]
treeWrites NonEmpty (JobTree payload)
children

    go
      :: Ops.TraceStamp payload
      -> Maybe Int64
      -- \^ Parent primary key (Nothing for root)
      -> Bool
      -- \^ Whether this node should be inserted suspended
      -> JobTree payload
      -> m (NonEmpty (JobRead payload))
    go :: TraceStamp payload
-> Maybe Int64
-> Bool
-> JobTree payload
-> m (NonEmpty (JobRead payload))
go TraceStamp payload
stamp Maybe Int64
mParentId Bool
susp (Leaf JobWrite payload
jobW) = do
      mInserted <- Text
-> Text
-> TraceStamp payload
-> Maybe Int64
-> Maybe Value
-> Bool
-> JobWrite payload
-> m (Maybe (JobRead payload))
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> TraceStamp payload
-> Maybe Int64
-> Maybe Value
-> Bool
-> JobWrite payload
-> m (Maybe (JobRead payload))
Ops.insertJobTreeNodeStamped Text
schemaName Text
tableName TraceStamp payload
stamp Maybe Int64
mParentId Maybe Value
forall a. Maybe a
Nothing Bool
susp JobWrite payload
jobW
      case mInserted of
        Maybe (JobRead payload)
Nothing -> TreeInsertFailed -> m (NonEmpty (JobRead payload))
forall (m :: * -> *) e a. (MonadIO m, Exception e) => e -> m a
UE.throwIO (TreeInsertFailed -> m (NonEmpty (JobRead payload)))
-> TreeInsertFailed -> m (NonEmpty (JobRead payload))
forall a b. (a -> b) -> a -> b
$ Text -> TreeInsertFailed
TreeInsertFailed Text
"insertJobTree: job insert failed (dedup conflict)"
        Just JobRead payload
inserted -> NonEmpty (JobRead payload) -> m (NonEmpty (JobRead payload))
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (JobRead payload
inserted JobRead payload -> [JobRead payload] -> NonEmpty (JobRead payload)
forall a. a -> [a] -> NonEmpty a
:| [])
    go TraceStamp payload
stamp Maybe Int64
mParentId Bool
susp (Finalizer JobWrite payload
jobW NonEmpty (JobTree payload)
children) = do
      mInserted <- Text
-> Text
-> TraceStamp payload
-> Maybe Int64
-> Maybe Value
-> Bool
-> JobWrite payload
-> m (Maybe (JobRead payload))
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> TraceStamp payload
-> Maybe Int64
-> Maybe Value
-> Bool
-> JobWrite payload
-> m (Maybe (JobRead payload))
Ops.insertJobTreeNodeStamped Text
schemaName Text
tableName TraceStamp payload
stamp Maybe Int64
mParentId (Value -> Maybe Value
forall a. a -> Maybe a
Just Value
emptyState) Bool
susp JobWrite payload
jobW
      case mInserted of
        Maybe (JobRead payload)
Nothing -> TreeInsertFailed -> m (NonEmpty (JobRead payload))
forall (m :: * -> *) e a. (MonadIO m, Exception e) => e -> m a
UE.throwIO (TreeInsertFailed -> m (NonEmpty (JobRead payload)))
-> TreeInsertFailed -> m (NonEmpty (JobRead payload))
forall a b. (a -> b) -> a -> b
$ Text -> TreeInsertFailed
TreeInsertFailed Text
"insertJobTree: parent insert failed (dedup conflict)"
        Just JobRead payload
inserted -> do
          descendants <- Int64 -> [JobTree payload] -> m [JobRead payload]
insertChildren (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
primaryKey JobRead payload
inserted) (NonEmpty (JobTree payload) -> [JobTree payload]
forall a. NonEmpty a -> [a]
NE.toList NonEmpty (JobTree payload)
children)
          pure (inserted :| descendants)
      where
        -- Batch adjacent leaves. Nested trees keep their position.
        insertChildren :: Int64 -> [JobTree payload] -> m [JobRead payload]
insertChildren Int64
_ [] = [JobRead payload] -> m [JobRead payload]
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure []
        insertChildren Int64
parentPK children' :: [JobTree payload]
children'@(Leaf JobWrite payload
_ : [JobTree payload]
_) = do
          let ([JobTree payload]
leaves, [JobTree payload]
rest) = (JobTree payload -> Bool)
-> [JobTree payload] -> ([JobTree payload], [JobTree payload])
forall a. (a -> Bool) -> [a] -> ([a], [a])
span JobTree payload -> Bool
forall {payload}. JobTree payload -> Bool
isLeaf [JobTree payload]
children'
              leafWrites :: [JobWrite payload]
leafWrites = [JobWrite payload
job | Leaf JobWrite payload
job <- [JobTree payload]
leaves]
          leafJobs <- Text
-> Text
-> TraceStamp payload
-> Int64
-> [JobWrite payload]
-> m [JobRead payload]
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> TraceStamp payload
-> Int64
-> [JobWrite payload]
-> m [JobRead payload]
Ops.insertJobTreeLeavesStamped Text
schemaName Text
tableName TraceStamp payload
stamp Int64
parentPK [JobWrite payload]
leafWrites
          when (length leafJobs /= length leafWrites)
            $ UE.throwIO
            $ TreeInsertFailed "insertJobTree: leaf batch insert had dedup conflicts"
          (leafJobs <>) <$> insertChildren parentPK rest
        insertChildren Int64
parentPK (JobTree payload
subTree : [JobTree payload]
rest) = do
          subTreeJobs <- TraceStamp payload
-> Maybe Int64
-> Bool
-> JobTree payload
-> m (NonEmpty (JobRead payload))
go TraceStamp payload
stamp (Int64 -> Maybe Int64
forall a. a -> Maybe a
Just Int64
parentPK) Bool
True JobTree payload
subTree
          remaining <- insertChildren parentPK rest
          pure (NE.toList subTreeJobs <> remaining)

        isLeaf :: JobTree payload -> Bool
isLeaf (Leaf JobWrite payload
_) = Bool
True
        isLeaf JobTree payload
_ = Bool
False