packages feed

kiroku-store-0.10.0.0: src/Kiroku/Store/Settings.hs

{- | Interpreter-level hooks applied to 'EventData' before encoding on
the append path and to 'RecordedEvent' after decoding on the read and
subscription paths.

The hook seam lives inside 'Kiroku.Store.Effect.runStorePool' (and the
subscription publisher\/worker) rather than at the SQL encoder\/decoder
layer so the hook sees the typed value — payload and metadata as
'Data.Aeson.Value', the event type, ids — and can branch on event type
or mutate structured JSON. Plumbing it at the encoder layer would force
hooks to operate on opaque bytes.

Both fields default to 'Nothing'. With the defaults, the helpers below
retain the original event list/vector without traversal. Subscription
decoding adds one batch constructor, with no per-event wrappers.

A typical use case is enriching every appended event with an
OpenTelemetry trace context drawn from the calling thread:

@
storeSettings = 'defaultStoreSettings'
  { 'enrichEvent' = Just $ \\ed -> do
      ctx <- captureCurrentSpan        -- OpenTelemetry, OTLP, whatever
      pure (ed & #metadata %~ injectTraceContext ctx)
  , 'decodeHook' = Just $ \\re ->
      pure (Right (re & #metadata %~ Just . redactPII))
  }
@

Wire the resulting 'StoreSettings' into
'Kiroku.Store.Connection.ConnectionSettings' via its @storeSettings@
field; 'Kiroku.Store.Connection.withStore' copies it onto the
'Kiroku.Store.Connection.KirokuStore' handle for the interpreter to
reach.

Direct callers of 'Kiroku.Store.Transaction.appendToStreamTx' bypass
'runStorePool' and therefore the 'enrichEvent' hook. Use
'Kiroku.Store.Transaction.enrichEventsIO' to opt in to enrichment
manually before constructing the prepared event list.
-}
module Kiroku.Store.Settings (
    StoreSettings (..),
    defaultStoreSettings,
    enrichEvents,
    DecodeFailure (..),
    DecodedEvent (..),
    DecodedBatch (..),
    decodeEvents,
    decodeEvent,
    decodedEventRecorded,
    decodedBatchLength,
    decodedBatchLastEvent,
    filterDecodedBatch,
) where

import Data.Text (Text)
import Data.Vector (Vector)
import Data.Vector qualified as V
import GHC.Generics (Generic)
import Kiroku.Store.Types (EventData, EventId, RecordedEvent)

{- | A hook could not decode this event. Return it through 'Left'; throwing
from the hook remains a programming error. Default retry exhaustion is surfaced
by SubscriptionUndecodable, carrying this failure.
-}
data DecodeFailure = DecodeFailure
    { decodeFailureEventId :: !EventId
    , decodeFailureReason :: !Text
    }
    deriving stock (Eq, Show, Generic)

-- | A transformed event, or its raw value retained for retry and disposition.
data DecodedEvent
    = Decoded !RecordedEvent
    | Undecodable !RecordedEvent !DecodeFailure
    deriving stock (Eq, Show)

{- | No hook means no traversal or per-event wrappers. Hook results are shared
by all live publisher subscribers, whose retry/disposition is independent.
-}
data DecodedBatch
    = UnchangedBatch !(Vector RecordedEvent)
    | TransformedBatch !(Vector DecodedEvent)
    deriving stock (Eq, Show)

decodedEventRecorded :: DecodedEvent -> RecordedEvent
decodedEventRecorded = \case
    Decoded event -> event
    Undecodable event _ -> event

decodedBatchLength :: DecodedBatch -> Int
decodedBatchLength = \case
    UnchangedBatch events -> V.length events
    TransformedBatch events -> V.length events

decodedBatchLastEvent :: DecodedBatch -> RecordedEvent
decodedBatchLastEvent = \case
    UnchangedBatch events -> V.last events
    TransformedBatch events -> decodedEventRecorded (V.last events)

filterDecodedBatch :: (RecordedEvent -> Bool) -> DecodedBatch -> DecodedBatch
filterDecodedBatch predicate = \case
    UnchangedBatch events -> UnchangedBatch (V.filter predicate events)
    TransformedBatch events -> TransformedBatch (V.filter (predicate . decodedEventRecorded) events)

{- | Interpreter-level hooks for cross-cutting concerns at the
event-data boundary. All fields default to 'Nothing' (no-op).

* 'enrichEvent' fires on the append path before the SQL encoder runs,
  on the typed 'EventData' the caller supplied. Used to inject trace
  contexts, attach tenant ids, or encrypt payloads.

* 'decodeHook' fires on the read and subscription paths after the SQL
  decoder runs, on the typed 'RecordedEvent' about to be surfaced to
  the caller. Used to decrypt payloads, redact PII, or attach derived
  metadata.

When a field is 'Nothing', reads and enrichment return their input directly.
Subscription decoding retains the vector in one batch constructor without
traversal or per-event wrappers.
-}
data StoreSettings = StoreSettings
    { enrichEvent :: !(Maybe (EventData -> IO EventData))
    -- ^ Append-path hook. Runs once per appended event before encoding.
    , decodeHook :: !(Maybe (RecordedEvent -> IO (Either DecodeFailure RecordedEvent)))
    {- ^ Read- and subscription-path hook. Return 'Left' for an undecodable
    event: reads fail with a typed store error, subscriptions use their
    optional undecodable handler or retry and stop by default. Exceptions
    remain programming failures. Runs once per surfaced event; an undecodable
    event's retry re-applies the hook to its original value.
    -}
    }
    deriving stock (Generic)

-- | Defaults to both hooks being 'Nothing' — semantically a no-op.
defaultStoreSettings :: StoreSettings
defaultStoreSettings =
    StoreSettings
        { enrichEvent = Nothing
        , decodeHook = Nothing
        }

{- | Apply 'enrichEvent' to a list of events. When the hook is
'Nothing', returns the list unchanged with no traversal.
-}
enrichEvents :: StoreSettings -> [EventData] -> IO [EventData]
enrichEvents ss xs = case enrichEvent ss of
    Nothing -> pure xs
    Just f -> traverse f xs

{- | Apply 'decodeHook' to a vector of events. When the hook is
'Nothing', retains the vector unchanged without traversal or per-event wrappers.
-}
decodeEvents :: StoreSettings -> Vector RecordedEvent -> IO DecodedBatch
decodeEvents ss xs = case decodeHook ss of
    Nothing -> pure (UnchangedBatch xs)
    Just f -> TransformedBatch <$> V.mapM (applyDecode f) xs

-- | Re-apply the hook to one raw undecodable event on a subscriber retry.
decodeEvent :: StoreSettings -> RecordedEvent -> IO DecodedEvent
decodeEvent ss event = case decodeHook ss of
    Nothing -> pure (Decoded event)
    Just f -> applyDecode f event

applyDecode :: (RecordedEvent -> IO (Either DecodeFailure RecordedEvent)) -> RecordedEvent -> IO DecodedEvent
applyDecode f event = either (Undecodable event) Decoded <$> f event