keiro-0.6.0.0: src/Keiro/Outbox/Schema.hs
{-# LANGUAGE ApplicativeDo #-}
-- | Hasql-level storage for the durable integration-event outbox.
--
-- This module owns the SQL surface that publishers call into. Higher-level
-- helpers ('Keiro.Outbox') and transport adapters
-- ('Keiro.Outbox.Kafka') consume these primitives.
module Keiro.Outbox.Schema
( enqueueOutboxTx,
claimOutboxBatch,
requeueStuckOutbox,
markOutboxSent,
markOutboxSentBatch,
markOutboxFailedTx,
markOutboxSkippedTx,
lookupOutbox,
listOutbox,
countOutboxBacklog,
garbageCollectSent,
)
where
import Contravariant.Extras (contrazip2, contrazip3, contrazip5)
import Data.ByteString (ByteString)
import Data.Functor.Contravariant ((>$<))
import Data.Time.Clock (NominalDiffTime, addUTCTime)
import Data.UUID (UUID)
import Effectful (Eff, (:>))
import Hasql.Decoders qualified as D
import Hasql.Encoders qualified as E
import Hasql.Statement (Statement, preparable)
import Keiro.Integration.Event
( IntegrationEvent (..),
SchemaReference (..),
TraceContext (..),
contentTypeText,
parseContentType,
)
import Keiro.Outbox.Types
import Keiro.Prelude
import Kiroku.Store.Effect (Store)
import Kiroku.Store.Transaction (runTransaction)
import Kiroku.Store.Types (EventId (..), GlobalPosition (..))
import "hasql-transaction" Hasql.Transaction qualified as Tx
-- | Enqueue one integration event inside an existing transaction.
--
-- The @(source, message_id)@ unique constraint catches duplicate retries
-- from a saga/process-manager. Callers that mint a fresh @messageId@ per
-- attempt should also mint a fresh @outboxId@; callers that want
-- idempotent retries should reuse both.
--
-- The row's @created_at@ value is the PostgreSQL transaction-start time. The
-- per-key and per-source publisher policies therefore provide only best-effort
-- ordering when concurrent transactions enqueue the same key/source and commit
-- in the opposite order. Serialize those transactions when strict order matters;
-- the canonical producer subscription already does so.
enqueueOutboxTx :: OutboxMessage -> Tx.Transaction ()
enqueueOutboxTx message =
Tx.statement (toEncodedRow message) enqueueOutboxStmt
-- | Read a single outbox row by id. Used by tests and inspection tooling.
lookupOutbox :: (Store :> es) => OutboxId -> Eff es (Maybe OutboxRow)
lookupOutbox outboxId =
runTransaction $
Tx.statement (unOutboxId outboxId) lookupOutboxStmt
-- | List outbox rows for a source, ordered by @created_at@. Used by tests;
-- not intended for application traffic.
listOutbox :: (Store :> es) => Text -> Eff es [OutboxRow]
listOutbox source =
runTransaction $
Tx.statement source listOutboxStmt
-- | Count outbox rows awaiting publish (backlog gauge source).
--
-- Backlog = rows in a claimable, non-terminal state. Mirrors the claim
-- query's @status IN ('pending','failed')@ predicate so the gauge measures
-- exactly the rows a publisher still has to drain (rows held mid-pass in
-- @publishing@, and the terminal @sent@/@dead@ rows, are excluded).
countOutboxBacklog :: (Store :> es) => Eff es Int
countOutboxBacklog =
runTransaction (Tx.statement () countBacklogStmt)
-- | Delete @sent@ rows whose @published_at@ is older than @keepFor@ before
-- @now@.
--
-- Returns the number of rows deleted. @dead@ rows are never deleted: they are
-- operator action items proving an event was not published. The retention window
-- only bounds how long successful publish history remains queryable; consumer
-- dedupe lives in the inbox, not here.
garbageCollectSent ::
(Store :> es) =>
NominalDiffTime ->
UTCTime ->
Eff es Int
garbageCollectSent keepFor now = do
let cutoff = addUTCTime (negate keepFor) now
result <-
runTransaction $
Tx.statement cutoff gcSentStmt
pure (fromIntegral result)
-- | Claim up to @limit@ rows ready for publish.
--
-- Rows in @pending@ or @failed@ status whose @next_attempt_at@ has passed
-- become candidates. The selection is filtered by 'OrderingPolicy':
--
-- * 'PerKeyHeadOfLine' — a row is claimed only if every earlier
-- non-terminal row with the same @(source, message_key)@ is also claimed
-- by the same statement. Rows with @message_key IS NULL@ bypass the
-- per-key check.
-- * 'PerSourceStream' — a row is claimed only if every earlier
-- non-terminal row in the same @source@ is also claimed by the same
-- statement, regardless of key.
-- * 'StopTheLine' — same as 'PerKeyHeadOfLine' at claim time; the worker
-- halts on the first failure (decided at the worker level).
-- * 'BestEffort' — no head-of-line predicate.
--
-- The returned list preserves @(created_at, outbox_id)@ order. Per-key and
-- per-source subsequences are therefore gapless ordered runs. The @LIMIT@
-- applies to the locked candidate set before the post-filter; under
-- concurrent claimers or a limit cut through the middle of a run, a pass
-- can return fewer than @limit@ rows even when more rows are ready.
--
-- Claimed rows are transitioned to @publishing@ and have their
-- @attempt_count@ incremented atomically.
claimOutboxBatch ::
(Store :> es) =>
OrderingPolicy ->
Int ->
UTCTime ->
Eff es [OutboxRow]
claimOutboxBatch policy limit now =
runTransaction $
Tx.statement (fromIntegral limit, now) (claimStmt policy)
-- | Reclaim rows stranded in @publishing@ longer than @olderThan@.
--
-- Rows whose claim already consumed the attempt budget are dead-lettered; the
-- rest return to @failed@ so the regular claim query can retry them. Returns
-- @(requeued, deadLettered)@.
requeueStuckOutbox ::
(Store :> es) =>
Int ->
NominalDiffTime ->
UTCTime ->
Eff es (Int, Int)
requeueStuckOutbox maxAttempts olderThan now =
runTransaction $ do
let cutoff = addUTCTime (negate olderThan) now
dead <- Tx.statement (cutoff, fromIntegral maxAttempts, now) deadLetterStuckStmt
requeued <- Tx.statement (cutoff, fromIntegral maxAttempts, now) requeueStuckStmt
pure (fromIntegral requeued, fromIntegral dead)
-- | Mark a row as successfully published. Sets @published_at@ and clears
-- @last_error@. Returns 'False' if the row left @publishing@ before the mark,
-- for example because a stale-row sweeper or operator changed it while the
-- transport publish was in flight. The publish may still have happened; callers
-- must treat this as at-least-once delivery.
markOutboxSent :: (Store :> es) => OutboxId -> UTCTime -> Eff es Bool
markOutboxSent outboxId now =
runTransaction $
Tx.statement (unOutboxId outboxId, now) markSentStmt
-- | Mark many rows as successfully published in one statement.
--
-- Only rows still in @publishing@ transition. Returns how many rows changed;
-- callers treat a shortfall as benign because delivery is already at-least-once.
markOutboxSentBatch :: (Store :> es) => [OutboxId] -> UTCTime -> Eff es Int
markOutboxSentBatch [] _ = pure 0
markOutboxSentBatch outboxIds now =
fromIntegral
<$> runTransaction
( Tx.statement
(fmap unOutboxId outboxIds, now)
markSentBatchStmt
)
-- | Mark a row as failed and decide whether it is retryable or dead.
--
-- Reads the current @attempt_count@; if it is greater than or equal to
-- @maxAttempts@, transitions to 'OutboxDead'. Otherwise transitions to
-- 'OutboxFailed' and sets @next_attempt_at = now + delay@. Returns the
-- resulting status so the worker can update its summary counters.
--
-- Runs inside the caller's transaction to keep "read attempt count → write
-- status" atomic with respect to other workers.
--
-- Only rows still in @publishing@ are updated (matching 'markSentBatchStmt'
-- and 'markSkippedStmt'): if the row outlived 'publishingTimeout' and
-- 'outboxMaintenancePass' already requeued it — possibly handing it to
-- another worker — a late failure mark from the original worker must not
-- flip the re-claimed row mid-publish.
markOutboxFailedTx ::
OutboxId ->
Text ->
Int ->
NominalDiffTime ->
UTCTime ->
Tx.Transaction OutboxStatus
markOutboxFailedTx outboxId errMsg maxAttempts delay now = do
currentAttempt <- Tx.statement (unOutboxId outboxId) readAttemptCountStmt
let attempt = fromMaybe 0 currentAttempt
shouldDie = attempt >= maxAttempts
nextStatus = if shouldDie then OutboxDead else OutboxFailed
nextAttempt = addUTCTime delay now
Tx.statement
(unOutboxId outboxId, statusText nextStatus, errMsg, nextAttempt, now)
markFailedStmt
pure nextStatus
-- | Return a claimed row to @failed@ without consuming an attempt.
--
-- Used for rows skipped because an earlier row in the same ordered group failed
-- inside the same publish batch.
markOutboxSkippedTx :: OutboxId -> Text -> UTCTime -> Tx.Transaction ()
markOutboxSkippedTx outboxId errMsg now =
Tx.statement (unOutboxId outboxId, errMsg, now) markSkippedStmt
-- ---------------------------------------------------------------------------
-- Encoder support
-- ---------------------------------------------------------------------------
-- | Flattened row used as the encoder input. Field order is locked to
-- the INSERT statement; the encoder threads each field through a separate
-- 'E.Params' fragment joined by 'mconcat'.
data EncodedRow = EncodedRow
{ outboxId :: !UUID,
messageId :: !Text,
source :: !Text,
destination :: !Text,
messageKey :: !(Maybe Text),
eventType :: !Text,
schemaVersion :: !Int64,
contentType :: !Text,
schemaRegistry :: !(Maybe Text),
schemaSubject :: !(Maybe Text),
schemaVersionRef :: !(Maybe Int64),
schemaId :: !(Maybe Int64),
schemaFingerprint :: !(Maybe Text),
sourceEventId :: !(Maybe UUID),
sourceGlobalPosition :: !(Maybe Int64),
causationId :: !(Maybe UUID),
correlationId :: !(Maybe UUID),
traceparent :: !(Maybe Text),
tracestate :: !(Maybe Text),
payloadBytes :: !ByteString,
attributes :: !(Maybe Value),
occurredAt :: !UTCTime
}
deriving stock (Generic)
toEncodedRow :: OutboxMessage -> EncodedRow
toEncodedRow message =
let event = message ^. #event
mref = event ^. #schemaReference
mtrace = event ^. #traceContext
in EncodedRow
{ outboxId = unOutboxId (message ^. #outboxId),
messageId = event ^. #messageId,
source = event ^. #source,
destination = event ^. #destination,
messageKey = event ^. #key,
eventType = event ^. #eventType,
schemaVersion = fromIntegral (event ^. #schemaVersion),
contentType = contentTypeText (event ^. #contentType),
schemaRegistry = mref >>= (^. #registry),
schemaSubject = mref >>= (^. #subject),
schemaVersionRef = fmap fromIntegral (mref >>= (^. #version)),
schemaId = mref >>= (^. #schemaId),
schemaFingerprint = mref >>= (^. #fingerprint),
sourceEventId = fmap unEventId (event ^. #sourceEventId),
sourceGlobalPosition = fmap unGlobalPosition (event ^. #sourceGlobalPosition),
causationId = fmap unEventId (event ^. #causationId),
correlationId = fmap unEventId (event ^. #correlationId),
traceparent = fmap (^. #traceparent) mtrace,
tracestate = mtrace >>= (^. #tracestate),
payloadBytes = event ^. #payloadBytes,
attributes = event ^. #attributes,
occurredAt = event ^. #occurredAt
}
unEventId :: EventId -> UUID
unEventId (EventId u) = u
unGlobalPosition :: GlobalPosition -> Int64
unGlobalPosition (GlobalPosition i) = i
encodedRowEncoder :: E.Params EncodedRow
encodedRowEncoder =
mconcat
[ view #outboxId >$< E.param (E.nonNullable E.uuid),
view #messageId >$< E.param (E.nonNullable E.text),
view #source >$< E.param (E.nonNullable E.text),
view #destination >$< E.param (E.nonNullable E.text),
view #messageKey >$< E.param (E.nullable E.text),
view #eventType >$< E.param (E.nonNullable E.text),
view #schemaVersion >$< E.param (E.nonNullable E.int8),
view #contentType >$< E.param (E.nonNullable E.text),
view #schemaRegistry >$< E.param (E.nullable E.text),
view #schemaSubject >$< E.param (E.nullable E.text),
view #schemaVersionRef >$< E.param (E.nullable E.int8),
view #schemaId >$< E.param (E.nullable E.int8),
view #schemaFingerprint >$< E.param (E.nullable E.text),
view #sourceEventId >$< E.param (E.nullable E.uuid),
view #sourceGlobalPosition >$< E.param (E.nullable E.int8),
view #causationId >$< E.param (E.nullable E.uuid),
view #correlationId >$< E.param (E.nullable E.uuid),
view #traceparent >$< E.param (E.nullable E.text),
view #tracestate >$< E.param (E.nullable E.text),
view #payloadBytes >$< E.param (E.nonNullable E.bytea),
view #attributes >$< E.param (E.nullable E.jsonb),
view #occurredAt >$< E.param (E.nonNullable E.timestamptz)
]
-- ---------------------------------------------------------------------------
-- Statements
-- ---------------------------------------------------------------------------
enqueueOutboxStmt :: Statement EncodedRow ()
enqueueOutboxStmt =
preparable
"""
INSERT INTO keiro.keiro_outbox
( outbox_id
, message_id
, source
, destination
, message_key
, event_type
, schema_version
, content_type
, schema_registry
, schema_subject
, schema_version_ref
, schema_id
, schema_fingerprint
, source_event_id
, source_global_position
, causation_id
, correlation_id
, traceparent
, tracestate
, payload_bytes
, attributes
, occurred_at
)
VALUES
($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22)
ON CONFLICT (source, message_id) DO NOTHING
"""
encodedRowEncoder
D.noResult
claimStmt :: OrderingPolicy -> Statement (Int64, UTCTime) [OutboxRow]
claimStmt policy =
preparable
(claimSql (policyPredicates policy))
( contrazip2
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.timestamptz))
)
(D.rowList claimResultDecoder)
data PolicyPredicates = PolicyPredicates
{ preFilter :: !Text,
postFilter :: !Text
}
deriving stock (Generic)
policyPredicates :: OrderingPolicy -> PolicyPredicates
policyPredicates = \case
PerKeyHeadOfLine -> perKeyPredicates
PerSourceStream -> perSourcePredicates
StopTheLine -> perKeyPredicates
BestEffort -> PolicyPredicates {preFilter = "TRUE", postFilter = "TRUE"}
perKeyPredicates :: PolicyPredicates
perKeyPredicates =
PolicyPredicates
{ preFilter =
"""
( r.message_key IS NULL OR NOT EXISTS (
SELECT 1 FROM keiro.keiro_outbox earlier
WHERE earlier.source = r.source
AND earlier.message_key = r.message_key
AND (earlier.created_at, earlier.outbox_id) < (r.created_at, r.outbox_id)
AND earlier.status NOT IN ('sent', 'dead')
AND NOT (earlier.status IN ('pending', 'failed') AND earlier.next_attempt_at <= $2) ) )
""",
postFilter =
"""
( c.message_key IS NULL OR NOT EXISTS (
SELECT 1 FROM keiro.keiro_outbox earlier
WHERE earlier.source = c.source
AND earlier.message_key = c.message_key
AND (earlier.created_at, earlier.outbox_id) < (c.created_at, c.outbox_id)
AND earlier.status NOT IN ('sent', 'dead')
AND NOT EXISTS (
SELECT 1 FROM candidate c2
WHERE c2.outbox_id = earlier.outbox_id ) ) )
"""
}
perSourcePredicates :: PolicyPredicates
perSourcePredicates =
PolicyPredicates
{ preFilter =
"""
NOT EXISTS (
SELECT 1 FROM keiro.keiro_outbox earlier
WHERE earlier.source = r.source
AND (earlier.created_at, earlier.outbox_id) < (r.created_at, r.outbox_id)
AND earlier.status NOT IN ('sent', 'dead')
AND NOT (earlier.status IN ('pending', 'failed') AND earlier.next_attempt_at <= $2) )
""",
postFilter =
"""
NOT EXISTS (
SELECT 1 FROM keiro.keiro_outbox earlier
WHERE earlier.source = c.source
AND (earlier.created_at, earlier.outbox_id) < (c.created_at, c.outbox_id)
AND earlier.status NOT IN ('sent', 'dead')
AND NOT EXISTS (
SELECT 1 FROM candidate c2
WHERE c2.outbox_id = earlier.outbox_id ) )
"""
}
claimSql :: PolicyPredicates -> Text
claimSql PolicyPredicates {preFilter, postFilter} =
"""
WITH candidate AS (
SELECT r.outbox_id, r.source, r.message_key, r.created_at
FROM keiro.keiro_outbox r
WHERE r.status IN ('pending', 'failed')
AND r.next_attempt_at <= $2
AND (
"""
<> preFilter
<> """
)
ORDER BY r.created_at, r.outbox_id
LIMIT $1
FOR UPDATE SKIP LOCKED
),
ready AS (
SELECT c.outbox_id, c.created_at AS claim_created_at
FROM candidate c
WHERE (
"""
<> postFilter
<> """
)
),
updated AS (
UPDATE keiro.keiro_outbox kt
SET status = 'publishing', attempt_count = kt.attempt_count + 1, updated_at = $2
FROM ready
WHERE kt.outbox_id = ready.outbox_id
RETURNING ready.claim_created_at,
"""
<> rowColumns
<> """
)
SELECT
"""
<> unqualifiedRowColumns
<> """
FROM updated
ORDER BY claim_created_at, outbox_id
"""
rowColumns :: Text
rowColumns =
"""
kt.outbox_id, kt.message_id, kt.source, kt.destination, kt.message_key,
kt.event_type, kt.schema_version, kt.content_type, kt.schema_registry,
kt.schema_subject, kt.schema_version_ref, kt.schema_id, kt.schema_fingerprint,
kt.source_event_id, kt.source_global_position, kt.causation_id,
kt.correlation_id, kt.traceparent, kt.tracestate, kt.payload_bytes,
kt.attributes, kt.occurred_at, kt.status, kt.attempt_count,
kt.next_attempt_at, kt.last_error, kt.published_at, kt.created_at,
kt.updated_at
"""
unqualifiedRowColumns :: Text
unqualifiedRowColumns =
"""
outbox_id, message_id, source, destination, message_key,
event_type, schema_version, content_type, schema_registry,
schema_subject, schema_version_ref, schema_id, schema_fingerprint,
source_event_id, source_global_position, causation_id,
correlation_id, traceparent, tracestate, payload_bytes,
attributes, occurred_at, status, attempt_count,
next_attempt_at, last_error, published_at, created_at,
updated_at
"""
claimResultDecoder :: D.Row OutboxRow
claimResultDecoder = outboxRowDecoder
countBacklogStmt :: Statement () Int
countBacklogStmt =
preparable
"SELECT COUNT(*)::bigint FROM keiro.keiro_outbox WHERE status IN ('pending', 'failed')"
E.noParams
(fmap fromIntegral (D.singleRow (D.column (D.nonNullable D.int8))))
gcSentStmt :: Statement UTCTime Int64
gcSentStmt =
preparable
"""
WITH deleted AS (
DELETE FROM keiro.keiro_outbox
WHERE status = 'sent' AND published_at < $1
RETURNING 1
)
SELECT COALESCE(COUNT(*), 0)::bigint FROM deleted
"""
(E.param (E.nonNullable E.timestamptz))
(D.singleRow (D.column (D.nonNullable D.int8)))
readAttemptCountStmt :: Statement UUID (Maybe Int)
readAttemptCountStmt =
preparable
"SELECT attempt_count FROM keiro.keiro_outbox WHERE outbox_id = $1"
(E.param (E.nonNullable E.uuid))
(D.rowMaybe (fromIntegral <$> D.column (D.nonNullable D.int8)))
deadLetterStuckStmt :: Statement (UTCTime, Int64, UTCTime) Int64
deadLetterStuckStmt =
preparable
"""
UPDATE keiro.keiro_outbox
SET status = 'dead',
last_error = COALESCE(last_error, 'reclaimed: publisher crashed mid-publish'),
updated_at = $3
WHERE status = 'publishing'
AND updated_at <= $1
AND attempt_count >= $2
"""
( contrazip3
(E.param (E.nonNullable E.timestamptz))
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.timestamptz))
)
D.rowsAffected
requeueStuckStmt :: Statement (UTCTime, Int64, UTCTime) Int64
requeueStuckStmt =
preparable
"""
UPDATE keiro.keiro_outbox
SET status = 'failed',
updated_at = $3
WHERE status = 'publishing'
AND updated_at <= $1
AND attempt_count < $2
"""
( contrazip3
(E.param (E.nonNullable E.timestamptz))
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.timestamptz))
)
D.rowsAffected
markSentStmt :: Statement (UUID, UTCTime) Bool
markSentStmt =
preparable
"""
UPDATE keiro.keiro_outbox
SET status = 'sent',
published_at = $2,
last_error = NULL,
updated_at = $2
WHERE outbox_id = $1
AND status = 'publishing'
"""
( contrazip2
(E.param (E.nonNullable E.uuid))
(E.param (E.nonNullable E.timestamptz))
)
((> 0) <$> D.rowsAffected)
markSentBatchStmt :: Statement ([UUID], UTCTime) Int64
markSentBatchStmt =
preparable
"""
UPDATE keiro.keiro_outbox
SET status = 'sent',
published_at = $2,
last_error = NULL,
updated_at = $2
WHERE outbox_id = ANY($1)
AND status = 'publishing'
"""
( contrazip2
(E.param (E.nonNullable (E.foldableArray (E.nonNullable E.uuid))))
(E.param (E.nonNullable E.timestamptz))
)
D.rowsAffected
markFailedStmt :: Statement (UUID, Text, Text, UTCTime, UTCTime) ()
markFailedStmt =
preparable
"""
UPDATE keiro.keiro_outbox
SET status = $2,
last_error = $3,
next_attempt_at = $4,
updated_at = $5
WHERE outbox_id = $1
AND status = 'publishing'
"""
( contrazip5
(E.param (E.nonNullable E.uuid))
(E.param (E.nonNullable E.text))
(E.param (E.nonNullable E.text))
(E.param (E.nonNullable E.timestamptz))
(E.param (E.nonNullable E.timestamptz))
)
D.noResult
markSkippedStmt :: Statement (UUID, Text, UTCTime) ()
markSkippedStmt =
preparable
"""
UPDATE keiro.keiro_outbox
SET status = 'failed',
attempt_count = GREATEST(attempt_count - 1, 0),
last_error = $2,
next_attempt_at = $3,
updated_at = $3
WHERE outbox_id = $1
AND status = 'publishing'
"""
( contrazip3
(E.param (E.nonNullable E.uuid))
(E.param (E.nonNullable E.text))
(E.param (E.nonNullable E.timestamptz))
)
D.noResult
lookupOutboxStmt :: Statement UUID (Maybe OutboxRow)
lookupOutboxStmt =
preparable
(selectAllSql <> " WHERE outbox_id = $1")
(E.param (E.nonNullable E.uuid))
(D.rowMaybe outboxRowDecoder)
listOutboxStmt :: Statement Text [OutboxRow]
listOutboxStmt =
preparable
(selectAllSql <> " WHERE source = $1 ORDER BY created_at, outbox_id")
(E.param (E.nonNullable E.text))
(D.rowList outboxRowDecoder)
selectAllSql :: Text
selectAllSql =
"""
SELECT outbox_id, message_id, source, destination, message_key, event_type,
schema_version, content_type, schema_registry, schema_subject,
schema_version_ref, schema_id, schema_fingerprint, source_event_id,
source_global_position, causation_id, correlation_id, traceparent,
tracestate, payload_bytes, attributes, occurred_at, status,
attempt_count, next_attempt_at, last_error, published_at, created_at,
updated_at
FROM keiro.keiro_outbox
"""
outboxRowDecoder :: D.Row OutboxRow
outboxRowDecoder = fmap assembleRow rawRowDecoder
data RawRow = RawRow
{ outboxId :: !OutboxId,
messageId :: !Text,
source :: !Text,
destination :: !Text,
key :: !(Maybe Text),
eventType :: !Text,
schemaVersion :: !Int,
contentType :: !Text,
schemaRegistry :: !(Maybe Text),
schemaSubject :: !(Maybe Text),
schemaVersionRef :: !(Maybe Int),
schemaId :: !(Maybe Int64),
schemaFingerprint :: !(Maybe Text),
sourceEventId :: !(Maybe EventId),
sourceGlobalPosition :: !(Maybe GlobalPosition),
causationId :: !(Maybe EventId),
correlationId :: !(Maybe EventId),
traceparent :: !(Maybe Text),
tracestate :: !(Maybe Text),
payloadBytes :: !ByteString,
attributes :: !(Maybe Value),
occurredAt :: !UTCTime,
status :: !OutboxStatus,
attemptCount :: !Int,
nextAttemptAt :: !UTCTime,
lastError :: !(Maybe Text),
publishedAt :: !(Maybe UTCTime),
createdAt :: !UTCTime,
updatedAt :: !UTCTime
}
deriving stock (Generic)
rawRowDecoder :: D.Row RawRow
rawRowDecoder =
RawRow
<$> (OutboxId <$> D.column (D.nonNullable D.uuid))
<*> D.column (D.nonNullable D.text)
<*> D.column (D.nonNullable D.text)
<*> D.column (D.nonNullable D.text)
<*> D.column (D.nullable D.text)
<*> D.column (D.nonNullable D.text)
<*> (fromIntegral <$> D.column (D.nonNullable D.int8))
<*> D.column (D.nonNullable D.text)
<*> D.column (D.nullable D.text)
<*> D.column (D.nullable D.text)
<*> (fmap fromIntegral <$> D.column (D.nullable D.int8))
<*> D.column (D.nullable D.int8)
<*> D.column (D.nullable D.text)
<*> (fmap EventId <$> D.column (D.nullable D.uuid))
<*> (fmap GlobalPosition <$> D.column (D.nullable D.int8))
<*> (fmap EventId <$> D.column (D.nullable D.uuid))
<*> (fmap EventId <$> D.column (D.nullable D.uuid))
<*> D.column (D.nullable D.text)
<*> D.column (D.nullable D.text)
<*> D.column (D.nonNullable D.bytea)
<*> D.column (D.nullable D.jsonb)
<*> D.column (D.nonNullable D.timestamptz)
<*> D.column (D.nonNullable (D.refine parseStatus D.text))
<*> (fromIntegral <$> D.column (D.nonNullable D.int8))
<*> D.column (D.nonNullable D.timestamptz)
<*> D.column (D.nullable D.text)
<*> D.column (D.nullable D.timestamptz)
<*> D.column (D.nonNullable D.timestamptz)
<*> D.column (D.nonNullable D.timestamptz)
assembleRow :: RawRow -> OutboxRow
assembleRow raw =
let traceContext = case raw ^. #traceparent of
Nothing -> Nothing
Just tp -> Just (TraceContext tp (raw ^. #tracestate))
schemaReference =
case ( raw ^. #schemaRegistry,
raw ^. #schemaSubject,
raw ^. #schemaVersionRef,
raw ^. #schemaId,
raw ^. #schemaFingerprint
) of
(Nothing, Nothing, Nothing, Nothing, Nothing) -> Nothing
_ ->
Just
( SchemaReference
(raw ^. #schemaRegistry)
(raw ^. #schemaSubject)
(raw ^. #schemaVersionRef)
(raw ^. #schemaId)
(raw ^. #schemaFingerprint)
)
event =
IntegrationEvent
{ messageId = raw ^. #messageId,
source = raw ^. #source,
destination = raw ^. #destination,
key = raw ^. #key,
eventType = raw ^. #eventType,
schemaVersion = raw ^. #schemaVersion,
contentType = parseContentType (raw ^. #contentType),
schemaReference,
sourceEventId = raw ^. #sourceEventId,
sourceGlobalPosition = raw ^. #sourceGlobalPosition,
payloadBytes = raw ^. #payloadBytes,
occurredAt = raw ^. #occurredAt,
causationId = raw ^. #causationId,
correlationId = raw ^. #correlationId,
traceContext,
attributes = raw ^. #attributes
}
in OutboxRow
{ outboxId = raw ^. #outboxId,
event,
status = raw ^. #status,
attemptCount = raw ^. #attemptCount,
nextAttemptAt = raw ^. #nextAttemptAt,
lastError = raw ^. #lastError,
publishedAt = raw ^. #publishedAt,
createdAt = raw ^. #createdAt,
updatedAt = raw ^. #updatedAt
}