keiro-0.4.0.1: src/Keiro/Snapshot/Schema.hs
{- | The @keiro_snapshots@ table: persistence for aggregate snapshots.
One row per stream holds the latest snapshot of its folded state as JSONB,
tagged with the 'stateCodecVersion', 'regfileShapeHash', and 'stateShapeHash'
that produced it. 'lookupSnapshot' fetches the newest row matching all three
discriminators (so incompatible snapshots are simply not found). Within one
discriminator tuple, 'writeSnapshotRow' keeps only the highest stream version,
so a late or out-of-order write cannot regress the snapshot.
A write with a /different/ discriminator deliberately replaces the row even
at a lower stream version. This lets a rolled-back deployment
reclaim the single snapshot slot instead of being locked out by a newer codec
forever. During a mixed-version deployment, however, writers with incompatible
codecs can thrash that row and each side will miss the other's snapshot. The
cost is repeated full replay, never incorrect state: the event log remains the
source of truth.
This module is the storage layer beneath "Keiro.Snapshot"; callers normally
go through 'Keiro.Snapshot.hydrateWithSnapshot' and
'Keiro.Snapshot.writeSnapshot' rather than these statements directly.
-}
module Keiro.Snapshot.Schema (
-- * Rows
SnapshotRow (..),
SnapshotWrite (..),
-- * Storage
lookupSnapshot,
writeSnapshotRow,
)
where
import Contravariant.Extras (contrazip4, contrazip6)
import Effectful (Eff, (:>))
import Hasql.Decoders qualified as D
import Hasql.Encoders qualified as E
import Hasql.Statement (Statement, preparable)
import Keiro.Prelude
import Kiroku.Store.Effect (Store)
import Kiroku.Store.Transaction (runTransaction)
import Kiroku.Store.Types (StreamId (..), StreamVersion (..))
import "hasql-transaction" Hasql.Transaction qualified as Tx
import Prelude qualified
{- | A snapshot row as read back from @keiro_snapshots@: the stored 'state'
JSON, the 'streamVersion' it captures, the 'stateCodecVersion' and
'regfileShapeHash' and 'stateShapeHash' that gate compatibility, and the create/update
timestamps.
-}
data SnapshotRow = SnapshotRow
{ streamId :: !StreamId
, streamVersion :: !StreamVersion
, state :: !Value
, stateCodecVersion :: !Int
, regfileShapeHash :: !Text
, stateShapeHash :: !Text
, createdAt :: !UTCTime
, updatedAt :: !UTCTime
}
deriving stock (Generic, Eq, Show)
{- | The fields needed to write a snapshot — 'SnapshotRow' minus the
database-managed timestamps.
-}
data SnapshotWrite = SnapshotWrite
{ streamId :: !StreamId
, streamVersion :: !StreamVersion
, state :: !Value
, stateCodecVersion :: !Int
, regfileShapeHash :: !Text
, stateShapeHash :: !Text
}
deriving stock (Generic, Eq, Show)
{- | Fetch the latest snapshot for a stream that matches all three codec
discriminators. Returns 'Nothing' when no compatible snapshot exists, so an
incompatible one is treated as absent.
-}
lookupSnapshot ::
(Store :> es) =>
StreamId ->
Int ->
Text ->
Text ->
Eff es (Maybe SnapshotRow)
lookupSnapshot streamId version shapeHash stateShapeHash =
runTransaction
$ Tx.statement
(streamIdToInt streamId, Prelude.fromIntegral version, shapeHash, stateShapeHash)
lookupSnapshotStmt
{- | Upsert a snapshot row for its stream. For the same discriminator tuple,
the write only takes effect when its 'streamVersion' is at least the stored
one. Any incompatible discriminator replaces the row even at a lower version
so codec rollback can make progress; see the module header for the
mixed-deployment performance caveat.
-}
writeSnapshotRow ::
(Store :> es) =>
SnapshotWrite ->
Eff es ()
writeSnapshotRow snapshot =
runTransaction
$ Tx.statement (snapshotWriteParams snapshot) writeSnapshotStmt
lookupSnapshotStmt :: Statement (Int64, Int64, Text, Text) (Maybe SnapshotRow)
lookupSnapshotStmt =
preparable
"""
SELECT stream_id, stream_version, state, state_codec_version, regfile_shape_hash, state_shape_hash, created_at, updated_at
FROM keiro.keiro_snapshots
WHERE stream_id = $1
AND state_codec_version = $2
AND regfile_shape_hash = $3
AND state_shape_hash = $4
ORDER BY stream_version DESC
LIMIT 1
"""
( contrazip4
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.text))
(E.param (E.nonNullable E.text))
)
(D.rowMaybe snapshotRowDecoder)
writeSnapshotStmt :: Statement (Int64, Int64, Value, Int64, Text, Text) ()
writeSnapshotStmt =
preparable
"""
INSERT INTO keiro.keiro_snapshots
(stream_id, stream_version, state, state_codec_version, regfile_shape_hash, state_shape_hash)
VALUES
($1, $2, $3, $4, $5, $6)
ON CONFLICT (stream_id) DO UPDATE
SET stream_version = EXCLUDED.stream_version,
state = EXCLUDED.state,
state_codec_version = EXCLUDED.state_codec_version,
regfile_shape_hash = EXCLUDED.regfile_shape_hash,
state_shape_hash = EXCLUDED.state_shape_hash,
updated_at = now()
WHERE keiro_snapshots.stream_version <= EXCLUDED.stream_version
OR keiro_snapshots.state_codec_version <> EXCLUDED.state_codec_version
OR keiro_snapshots.regfile_shape_hash <> EXCLUDED.regfile_shape_hash
OR keiro_snapshots.state_shape_hash <> EXCLUDED.state_shape_hash
"""
( contrazip6
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.jsonb))
(E.param (E.nonNullable E.int8))
(E.param (E.nonNullable E.text))
(E.param (E.nonNullable E.text))
)
D.noResult
snapshotRowDecoder :: D.Row SnapshotRow
snapshotRowDecoder =
SnapshotRow
<$> (StreamId <$> D.column (D.nonNullable D.int8))
<*> (StreamVersion <$> D.column (D.nonNullable D.int8))
<*> D.column (D.nonNullable D.jsonb)
<*> (Prelude.fromIntegral <$> D.column (D.nonNullable D.int8))
<*> D.column (D.nonNullable D.text)
<*> D.column (D.nonNullable D.text)
<*> D.column (D.nonNullable D.timestamptz)
<*> D.column (D.nonNullable D.timestamptz)
streamIdToInt :: StreamId -> Int64
streamIdToInt (StreamId value) = value
streamVersionToInt :: StreamVersion -> Int64
streamVersionToInt (StreamVersion value) = value
snapshotWriteParams :: SnapshotWrite -> (Int64, Int64, Value, Int64, Text, Text)
snapshotWriteParams snapshot =
( streamIdToInt (snapshot ^. #streamId)
, streamVersionToInt (snapshot ^. #streamVersion)
, snapshot ^. #state
, Prelude.fromIntegral (snapshot ^. #stateCodecVersion)
, snapshot ^. #regfileShapeHash
, snapshot ^. #stateShapeHash
)