keiro-dsl-0.17.0.0: test/conformance-process-timers/Generated/ProcessTimers/IncidentTimers/Process.hs
{-# LANGUAGE DeriveAnyClass #-}
{-# OPTIONS_GHC -Wno-missing-signatures #-}
-- @generated by keiro-dsl 0.17.0.0 (language keiro-dsl 6) from process reaction IncidentTimers; do not edit.
module Generated.ProcessTimers.IncidentTimers.Process
( IncidentTimersInput (..)
, incidentTimersProcessName
, incidentTimersCategory
, incidentTimersProcessWorkerOptions
, incidentTimersReactionVersion
, incidentTimersReactionFingerprint
, incidentTimersReact
, incidentTimersProcessManager
, incidentTimersRunProcessWorker
, EscalationPayload (..)
, incidentTimersEscalationTimerRequest
, ReminderPayload (..)
, incidentTimersReminderTimerRequest
, incidentTimersTimerWorkerOptions
, incidentTimersFireTimer
) where
import Data.Text (Text)
import Data.Aeson (FromJSON, ToJSON)
import Data.Aeson qualified as Aeson
import Data.ByteString qualified as BS
import Data.Text qualified as T
import Data.Text.Encoding (encodeUtf8)
import Data.Time (UTCTime, addUTCTime)
import Data.UUID (UUID)
import Data.UUID.V5 qualified as UUID.V5
import GHC.Generics (Generic)
import Numeric.Natural (Natural)
import Generated.ProcessTimers.IncidentTimers.Input (IncidentTimersInput (..))
import ProcessTimers.IncidentTimers.ProcessHoles (decodeIncidentTimersInput)
import Generated.ProcessTimers.IncidentSaga.Domain qualified as Saga
import Generated.ProcessTimers.IncidentSaga.EventStream (incidentSagaEventStream, IncidentSagaEventStreamDef)
import Generated.ProcessTimers.Incident.Domain qualified as Target
import Generated.ProcessTimers.Incident.EventStream (incidentCategory, incidentEventStream)
import Generated.ProcessTimers.Nominals qualified as N
import Keiro.Command (DomainCommandHandler (..), SilentDomainDecision (..))
import Keiro.Command (CommandError (..), RunCommandOptions (..), runCommand)
import Keiro.ProcessManager (confirmBenignDuplicate)
import Keiro.ProcessManager.Reaction qualified as Reaction
import Keiro.Stream qualified as Stream
import Keiro.DeterministicId (identitySeedBytes)
import Keiro.Timer (TimerId (..), TimerRequest (..), TimerWorkerOptions (..))
import Kiroku.Store.Types (EventId (..))
import Keiro.ProcessManager (PoisonPolicy (..), RejectedCommandPolicy (..), WorkerOptions (..))
import Shibuya.Core.Ack (RetryDelay (..))
data EscalationPayload = EscalationPayload
{ kind :: !Text
, incidentId :: !N.IncidentId
, detail :: !Text
}
deriving stock (Generic, Eq, Show)
deriving anyclass (FromJSON, ToJSON)
incidentTimersEscalationTimerRequest :: Text -> UTCTime -> N.IncidentId -> Text -> TimerRequest
incidentTimersEscalationTimerRequest correlationId fireAtTime payloadIncidentId payloadDetail =
TimerRequest
{ timerId = TimerId (reactionIdentity "incident-escalation-timer:" correlationId)
, processManagerName = incidentTimersProcessName
, correlationId = correlationId
, fireAt = fireAtTime
, payload = Aeson.toJSON (EscalationPayload { kind = "escalation", incidentId = payloadIncidentId, detail = payloadDetail })
}
data ReminderPayload = ReminderPayload
{ kind :: !Text
, incidentId :: !N.IncidentId
}
deriving stock (Generic, Eq, Show)
deriving anyclass (FromJSON, ToJSON)
incidentTimersReminderTimerRequest :: Text -> UTCTime -> N.IncidentId -> TimerRequest
incidentTimersReminderTimerRequest correlationId fireAtTime payloadIncidentId =
TimerRequest
{ timerId = TimerId (reactionIdentity "incident-reminder-timer:" correlationId)
, processManagerName = incidentTimersProcessName
, correlationId = correlationId
, fireAt = fireAtTime
, payload = Aeson.toJSON (ReminderPayload { kind = "reminder", incidentId = payloadIncidentId })
}
incidentTimersProcessName :: Text
incidentTimersProcessName = "incident-timers"
incidentTimersCategory :: Stream.StreamCategory IncidentSagaEventStreamDef
incidentTimersCategory = Stream.categoryUnsafe "incidentTimers"
incidentTimersProcessWorkerOptions :: WorkerOptions es msg
incidentTimersProcessWorkerOptions =
WorkerOptions
{ poisonPolicy = PoisonHalt,
rejectedCommandPolicy = RejectedHalt,
transientRetryDelay = RetryDelay 5, -- matches defaultWorkerOptions; runtime tuning
metrics = Nothing -- runtime configuration; install at call site
}
incidentTimersTimerWorkerOptions :: TimerWorkerOptions
incidentTimersTimerWorkerOptions =
TimerWorkerOptions
{ maxAttempts = Just 5
, requeueStuckAfter = Just 300
}
-- Operator dead-letter guidance: "incident timer exceeded ceiling"
incidentTimersReactionVersion :: Natural
incidentTimersReactionVersion = 1
incidentTimersReactionFingerprint :: Text
incidentTimersReactionFingerprint = "5a6afd7c9691c5716a9c7b3140b5c6c54ff8689aa6e0fbc9acf576ce56aa4904"
incidentTimersCorrelate :: IncidentTimersInput -> Text
incidentTimersCorrelate input = case input of
IncidentReported { incidentId } -> (N.incidentIdText incidentId)
ResponderAcked { incidentId } -> (N.incidentIdText incidentId)
incidentTimersReact :: IncidentTimersInput -> Reaction.ReactionPlan Saga.IncidentSagaCommand Target.IncidentCommand
incidentTimersReact input = case input of
IncidentReported { incidentId, raisedAt, detail }
-> Reaction.AdvanceReaction { command = Saga.RecordIncident (Saga.RecordIncidentData { Saga.incidentId = incidentId }), followUps = [Reaction.FollowSchedule Reaction.Rearm (incidentTimersEscalationTimerRequest (N.incidentIdText incidentId) (addUTCTime 300 raisedAt) incidentId detail), Reaction.FollowSchedule Reaction.Once (incidentTimersReminderTimerRequest (N.incidentIdText incidentId) (addUTCTime 3600 raisedAt) incidentId)], onAccepted = [] }
ResponderAcked { incidentId }
-> Reaction.NoAdvance [Reaction.FollowCancel (TimerId (reactionIdentity "incident-reminder-timer:" (N.incidentIdText incidentId)))]
incidentTimersProcessManager =
Reaction.ReactiveProcessManager
{ name = incidentTimersProcessName
, correlate = incidentTimersCorrelate
, sagaHandler = DomainCommandHandler { eventStream = incidentSagaEventStream, classifySilent = \_ -> SilentNoOp () }
, streamFor = Stream.entityStream incidentTimersCategory
, targetEventStream = incidentEventStream
, targetProjections = const []
, react = incidentTimersReact
}
incidentTimersRunProcessWorker options adapter =
Reaction.runReactiveProcessManagerWorkerWith
incidentTimersProcessWorkerOptions
options
incidentTimersProcessManager
adapter
(\event -> case decodeIncidentTimersInput event of Nothing -> Nothing; Just input -> Just (event, input))
incidentTimersFireTimer options timer
| timer.timerId == TimerId (reactionIdentity "incident-escalation-timer:" timer.correlationId) = incidentTimersEscalationFire options timer
| timer.timerId == TimerId (reactionIdentity "incident-reminder-timer:" timer.correlationId) = incidentTimersReminderFire options timer
| otherwise = pure Nothing
incidentTimersEscalationFire options timer
| timer.processManagerName /= incidentTimersProcessName = pure Nothing
| otherwise =
case (Aeson.fromJSON timer.payload :: Aeson.Result EscalationPayload) of
Aeson.Error _ -> pure Nothing
Aeson.Success decoded -> do
let firedId = EventId (reactionIdentity "incident-escalation-fired:" timer.correlationId)
target = Stream.entityStream incidentCategory timer.correlationId
result <-
runCommand
(options { eventIds = [firedId] })
incidentEventStream
target
(Target.EscalateIncident (Target.EscalateIncidentData { Target.incidentId = decoded.incidentId }))
case result of
Right {} -> pure (Just firedId)
Left err -> do
benign <- confirmBenignDuplicate (Stream.streamName target) firedId err
pure $ if benign then Just firedId else case err of
CommandRejected -> Just firedId
CommandAmbiguous _ -> Nothing
_ -> Nothing
incidentTimersReminderFire options timer
| timer.processManagerName /= incidentTimersProcessName = pure Nothing
| otherwise =
case (Aeson.fromJSON timer.payload :: Aeson.Result ReminderPayload) of
Aeson.Error _ -> pure Nothing
Aeson.Success decoded -> do
let firedId = EventId (reactionIdentity "incident-reminder-fired:" timer.correlationId)
target = Stream.entityStream incidentCategory timer.correlationId
result <-
runCommand
(options { eventIds = [firedId] })
incidentEventStream
target
(Target.RemindIncident (Target.RemindIncidentData { Target.incidentId = decoded.incidentId }))
case result of
Right {} -> pure (Just firedId)
Left err -> do
benign <- confirmBenignDuplicate (Stream.streamName target) firedId err
pure $ if benign then Just firedId else case err of
CommandRejected -> Just firedId
CommandAmbiguous _ -> Nothing
_ -> Nothing
reactionIdentity :: Text -> Text -> UUID
reactionIdentity prefix correlation =
UUID.V5.generateNamed UUID.V5.namespaceURL (identitySeedBytes (T.concat (map field [prefix, correlation])))
where
field value = T.pack (show (BS.length (encodeUtf8 value))) <> ":" <> value