From 5c39cb5683fc50e251ffa374ee9a344483880baf Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 26 Jan 2024 10:50:52 +0200 Subject: [PATCH 01/15] WIP: add UtxoWithSpent Signed-off-by: Ana Pantilie --- .../marconi-cardano-indexers.cabal | 2 + .../src/Marconi/Cardano/Indexers/Utxo.hs | 4 +- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 169 ++++++++++++++++++ 3 files changed, 172 insertions(+), 3 deletions(-) create mode 100644 marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs diff --git a/marconi-cardano-indexers/marconi-cardano-indexers.cabal b/marconi-cardano-indexers/marconi-cardano-indexers.cabal index 5a5226b3a2..c6df40ed7e 100644 --- a/marconi-cardano-indexers/marconi-cardano-indexers.cabal +++ b/marconi-cardano-indexers/marconi-cardano-indexers.cabal @@ -69,6 +69,7 @@ library Marconi.Cardano.Indexers.SyncHelper Marconi.Cardano.Indexers.Utxo Marconi.Cardano.Indexers.UtxoQuery + Marconi.Cardano.Indexers.UtxoWithSpent -------------------- -- Local components @@ -112,6 +113,7 @@ library , sqlite-simple , text , time + , transformers , vector-map test-suite marconi-cardano-indexers-test diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index bb37ccc780..b0a9c009a9 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -1,5 +1,4 @@ {-# LANGUAGE AllowAmbiguousTypes #-} -{-# LANGUAGE DeriveGeneric #-} {-# LANGUAGE DerivingStrategies #-} {-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE FlexibleInstances #-} @@ -76,7 +75,6 @@ import Database.SQLite.Simple qualified as SQL import Database.SQLite.Simple.QQ (sql) import Database.SQLite.Simple.ToField (ToField (toField)) import Database.SQLite.Simple.ToRow (ToRow (toRow)) -import GHC.Generics (Generic) import Marconi.Cardano.Core.Indexer.Worker ( StandardSQLiteIndexer, StandardWorker, @@ -106,7 +104,7 @@ data Utxo = Utxo , _inlineScript :: !(Maybe C.ScriptInAnyLang) , _inlineScriptHash :: !(Maybe C.ScriptHash) } - deriving (Show, Eq, Generic) + deriving (Show, Eq) -- | An alias for a non-empty list of @Utxo@, it's the event potentially produced on each block type UtxoEvent = NonEmpty Utxo diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs new file mode 100644 index 0000000000..5413f4c1a0 --- /dev/null +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -0,0 +1,169 @@ +{-# LANGUAGE FlexibleContexts #-} +{-# LANGUAGE NamedFieldPuns #-} +{-# LANGUAGE QuasiQuotes #-} +{-# LANGUAGE TemplateHaskell #-} +{-# OPTIONS_GHC -Wno-orphans #-} + +module Marconi.Cardano.Indexers.UtxoWithSpent ( + +) where + +import Cardano.Api qualified as C +import Control.Lens ((^.)) +import Control.Lens qualified as Lens +import Control.Monad.Error.Class (MonadError) +import Control.Monad.Except (MonadIO) +import Control.Monad.Trans.Maybe (MaybeT (MaybeT, runMaybeT)) +import Data.Aeson.TH qualified as Aeson +import Data.List.NonEmpty (NonEmpty) +import Data.List.NonEmpty qualified as NonEmpty +import Database.SQLite.Simple (FromRow (fromRow), ToRow (toRow)) +import Database.SQLite.Simple qualified as SQL +import Database.SQLite.Simple.QQ (sql) +import Database.SQLite.Simple.ToField (ToField (toField)) +import Marconi.Cardano.Core.Indexer.Worker (StandardSQLiteIndexer) +import Marconi.Cardano.Core.Orphans () +import Marconi.Cardano.Core.Types (TxIndexInBlock) +import Marconi.Cardano.Indexers.SyncHelper qualified as Sync +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) +import Marconi.Core.Indexer.SQLiteIndexer qualified as Core +import Marconi.Core.Type qualified as Core + +{- | Indexer representation of a UTxO together with the information + about where it was spent. If this information is missing, then it isn't + spent and we consider it active. +-} +data UtxoWithSpent = UtxoWithSpent + { _txIn :: !C.TxIn + -- ^ the tx id and tx index at which the tx out is created + , _address :: !C.AddressAny + -- ^ the address of the tx out + , _value :: !C.Value + -- ^ the value of the tx out + , _datumHash :: !(Maybe (C.Hash C.ScriptData)) + -- ^ the datum hash of the tx out + , _inlineScript :: !(Maybe C.ScriptInAnyLang) + -- ^ the inline script of the tx out + , _inlineScriptHash :: !(Maybe C.ScriptHash) + -- ^ the inline script hash of the tx out + , _txIndex :: TxIndexInBlock + -- ^ the index at which the tx is present in the block + , _spentAt :: !(Maybe C.TxIn) + -- ^ the tx id and index at which the tx out is spent + } + deriving (Show, Eq) + +Aeson.deriveJSON Aeson.defaultOptions{Aeson.fieldLabelModifier = tail} ''UtxoWithSpent + +Lens.makeLenses ''UtxoWithSpent + +instance SQL.ToRow (Core.Timed C.ChainPoint UtxoWithSpent) where + toRow u = + let (C.TxIn txid txix) = u ^. Core.event . txIn + spentAtFields = + case u ^. Core.event . spentAt of + (Just (C.TxIn txIdSpent txIxSpent)) -> [toField txIdSpent, toField txIxSpent] + Nothing -> [] + in toRow + [ toField $ u ^. Core.event . address + , toField $ u ^. Core.event . txIndex + , toField txid + , toField txix + , toField $ u ^. Core.event . datumHash + , toField $ u ^. Core.event . value + , toField $ u ^. Core.event . inlineScript + , toField $ u ^. Core.event . inlineScriptHash + ] + <> toRow spentAtFields + <> toRow (u ^. Core.point) + +instance FromRow (Core.Timed C.ChainPoint UtxoWithSpent) where + fromRow = do + utxo <- fromRow + point <- fromRow + pure $ Core.Timed point utxo + +instance FromRow UtxoWithSpent where + fromRow = do + _address <- SQL.field + _txIndex <- SQL.field + txId <- SQL.field + txIx <- SQL.field + _datumHash <- SQL.field + _value <- SQL.field + _inlineScript <- SQL.field + _inlineScriptHash <- SQL.field + _spentAt <- runMaybeT parseSpent + pure $ + UtxoWithSpent + { _address + , _txIndex + , _txIn = C.TxIn txId txIx + , _datumHash + , _value + , _inlineScript + , _inlineScriptHash + , _spentAt + } + where + parseSpent = MaybeT $ do + txIdSpent <- SQL.field + txIxSpent <- SQL.field + pure . pure $ C.TxIn txIdSpent txIxSpent + +-- | An alias for a non-empty list of @Utxo@, it's the event potentially produced on each block +type UtxoWithSpentEvent = NonEmpty UtxoWithSpent + +type instance Core.Point UtxoWithSpentEvent = C.ChainPoint +type UtxoWithSpentIndexer = Core.SQLiteIndexer UtxoWithSpentEvent +type StandardUtxoWithSpentIndexer m = StandardSQLiteIndexer m UtxoWithSpentEvent + +-- | Make a SQLiteIndexer for UtxoWithSpent +mkUtxoWithSpentIndexer + :: (MonadIO m, MonadError Core.IndexerError m) + => SQLiteDBLocation + -- ^ SQL connection to database + -> m UtxoWithSpentIndexer +mkUtxoWithSpentIndexer path = do + let createUtxoWithSpent = + [sql|CREATE TABLE IF NOT EXISTS utxo_with_spent + ( address BLOB NOT NULL + , txIndex INT NOT NULL + , txId TEXT NOT NULL + , txIx INT NOT NULL + , datumHash BLOB + , value BLOB + , inlineScript BLOB + , inlineScriptHash BLOB + , slotNo INT NOT NULL + , blockHeaderHash BLOB NOT NULL + , txIdSpent TEXT + , txIxSpent INT + )|] + -- TODO: fix + utxoInsertQuery :: SQL.Query + utxoInsertQuery = + [sql|INSERT INTO utxo ( + address, + txIndex, + txId, + txIx, + datumHash, + value, + inlineScript, + inlineScriptHash, + slotNo, + blockHeaderHash, + txIdSpent, + txIxSpent + ) VALUES + (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + createUtxoTables = [createUtxoWithSpent] + insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) utxoInsertQuery] + + Sync.mkSyncedSqliteIndexer + path + createUtxoTables + [insertEvent] + -- TODO: fix + [Core.SQLRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo] From dc051ffefd0f26b55058bd2465a4faf4ac3c9c3c Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 26 Jan 2024 13:47:48 +0200 Subject: [PATCH 02/15] Generalise RollbackPlan Signed-off-by: Ana Pantilie --- .../tutorials/BasicApp.hs | 16 +++--- .../src/Marconi/Cardano/Indexers/BlockInfo.hs | 2 +- .../src/Marconi/Cardano/Indexers/Datum.hs | 2 +- .../Marconi/Cardano/Indexers/EpochNonce.hs | 2 +- .../src/Marconi/Cardano/Indexers/EpochSDD.hs | 2 +- .../Cardano/Indexers/MintTokenEvent.hs | 4 +- .../src/Marconi/Cardano/Indexers/Spent.hs | 2 +- .../src/Marconi/Cardano/Indexers/Utxo.hs | 2 +- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 2 +- marconi-core/src/Marconi/Core.hs | 4 +- .../src/Marconi/Core/Indexer/SQLiteIndexer.hs | 54 ++++++++++--------- marconi-core/test/Marconi/CoreSpec.hs | 2 +- .../Marconi/Starter/Indexers/AddressCount.hs | 2 +- 13 files changed, 52 insertions(+), 44 deletions(-) diff --git a/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs b/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs index 3d813623d6..e79bf62642 100644 --- a/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs +++ b/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs @@ -341,13 +341,15 @@ mkBlockInfoSqliteIndexer dbPath = do ] -- Requests launched when a rollback occurs [ Core.SQLRollbackPlan - -- Name of the SQLite table - "block_info_table" - -- Field name of the SQLite table which we will use for handling - -- rollbacks (which deletes any database after that point) - "slotNo" - -- Translate the point of the 'Core.Timed point event' to a 'SlotNo' - C.chainPointToSlotNo + ( Core.defaultRollbackPlan + -- Name of the SQLite table + "block_info_table" + -- Field name of the SQLite table which we will use for handling + -- rollbacks (which deletes any database after that point) + "slotNo" + -- Translate the point of the 'Core.Timed point event' to a 'SlotNo' + C.chainPointToSlotNo + ) ] where dbCreation = diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs index fbd261daed..a19844f27f 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs @@ -133,7 +133,7 @@ mkBlockInfoIndexer path = do id createBlockInfoTable blockInfoInsertQuery - (Core.SQLRollbackPlan "blockInfo" "slotNo" C.chainPointToSlotNo) + (Core.SQLRollbackPlan (Core.defaultRollbackPlan "blockInfo" "slotNo" C.chainPointToSlotNo)) catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () catchupConfigEventHook indexerName stdoutTrace dbPath Core.Synced = do diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs index 6b8a9b84fc..df8c6d5609 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs @@ -125,7 +125,7 @@ mkDatumIndexer path = do path createDatumTables [[Core.SQLInsertPlan (traverse NonEmpty.toList) datumInsertQuery]] - [Core.SQLRollbackPlan "datum" "slotNo" C.chainPointToSlotNo] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "datum" "slotNo" C.chainPointToSlotNo)] -- | A worker with catchup for a 'DatumIndexer' datumWorker diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs index 153cbc84ed..392eff4ec8 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs @@ -112,7 +112,7 @@ mkEpochNonceIndexer path = do path [createNonce] [insertEvent] - [Core.SQLRollbackPlan "epoch_nonce" "slotNo" C.chainPointToSlotNo] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "epoch_nonce" "slotNo" C.chainPointToSlotNo)] newtype EpochNonceWorkerConfig input = EpochNonceWorkerConfig { epochNonceWorkerConfigExtractor :: input -> C.EpochNo diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs index 2245dab475..fd9017ba08 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs @@ -119,7 +119,7 @@ mkEpochSDDIndexer path = do path [createSDD] [insertEvent] - [Core.SQLRollbackPlan "epoch_sdd" "slotNo" C.chainPointToSlotNo] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "epoch_sdd" "slotNo" C.chainPointToSlotNo)] newtype EpochSDDWorkerConfig input = EpochSDDWorkerConfig { epochSDDWorkerConfigExtractor :: input -> C.EpochNo diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs index 509c60ceb6..9b3aebb941 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs @@ -445,7 +445,9 @@ mkMintTokenIndexer dbPath = do dbPath createMintPolicyEventTables [mintInsertPlans] - [Core.SQLRollbackPlan "minting_policy_events" "slotNo" C.chainPointToSlotNo] + [ Core.SQLRollbackPlan + (Core.defaultRollbackPlan "minting_policy_events" "slotNo" C.chainPointToSlotNo) + ] catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () catchupConfigEventHook indexerName stdoutTrace dbPath Core.Synced = do diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index e8aeeac5e4..98899c8be0 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -134,7 +134,7 @@ mkSpentIndexer path = do path createSpentTables [spentInsert] - [Core.SQLRollbackPlan "spent" "slotNo" C.chainPointToSlotNo] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "spent" "slotNo" C.chainPointToSlotNo)] catchupConfigEventHook :: Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () catchupConfigEventHook stdoutTrace dbPath Core.Synced = do diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index b0a9c009a9..3a4a6616f2 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -167,7 +167,7 @@ mkUtxoIndexer path = do path createUtxoTables [insertEvent] - [Core.SQLRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo)] catchupConfigEventHook :: Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () catchupConfigEventHook stdoutTrace dbPath Core.Synced = do diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 5413f4c1a0..7676c21e96 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -166,4 +166,4 @@ mkUtxoWithSpentIndexer path = do createUtxoTables [insertEvent] -- TODO: fix - [Core.SQLRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo)] diff --git a/marconi-core/src/Marconi/Core.hs b/marconi-core/src/Marconi/Core.hs index 31e941bf53..7e0157e69f 100644 --- a/marconi-core/src/Marconi/Core.hs +++ b/marconi-core/src/Marconi/Core.hs @@ -283,7 +283,8 @@ module Marconi.Core ( parseDBLocation, connection, SQLInsertPlan (SQLInsertPlan, planExtractor, planInsert), - SQLRollbackPlan (SQLRollbackPlan, tableName, pointName, pointExtractor), + SQLRollbackPlan (SQLRollbackPlan, planRollback), + defaultRollbackPlan, -- **** Reexport from SQLite ToRow (..), @@ -620,6 +621,7 @@ import Marconi.Core.Indexer.SQLiteIndexer ( ToRow (..), connection, dbLastSync, + defaultRollbackPlan, handleSQLErrors, inMemoryDB, mkSingleInsertSqliteIndexer, diff --git a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs index 4d47a461d1..479eb77a99 100644 --- a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs +++ b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs @@ -33,7 +33,8 @@ module Marconi.Core.Indexer.SQLiteIndexer ( handleSQLErrors, dbLastSync, SQLInsertPlan (SQLInsertPlan, planInsert, planExtractor), - SQLRollbackPlan (SQLRollbackPlan, tableName, pointName, pointExtractor), + SQLRollbackPlan (SQLRollbackPlan, planRollback), + defaultRollbackPlan, -- * Reexport from SQLite SQL.ToRow (..), @@ -102,17 +103,8 @@ newtype InsertPointQuery = InsertPointQuery {getInsertPointQuery :: SQL.Query} {- | A 'SQLRollbackPlan' provides a piece of information on how to perform a rollback on the data inserted in the database. -} -data SQLRollbackPlan point = forall a. - (SQL.ToField a) => - SQLRollbackPlan - { tableName :: String - -- ^ the table to rollback - , pointName :: String - -- ^ The name of the point field in the table - , pointExtractor :: point -> Maybe a - -- ^ How we transform the data to the point field. Returning 'Nothing' essentially means that we - -- delete all information from the database. Returning 'Just a' means that we will delete all - -- rows with a point higher than 'point'. +newtype SQLRollbackPlan point = SQLRollbackPlan + { planRollback :: point -> SQL.Connection -> IO () } -- | A newtype to set the last stable point of an indexer. @@ -298,6 +290,29 @@ indexEvents evts@(e : _) indexer = do runIndexQueries (indexer ^. connection) evts (indexer ^. insertPlan) setDbLastSync (e ^. point) indexer +defaultRollbackPlan + :: (SQL.ToField a) + => String + -> String + -> (point -> Maybe a) + -> point + -> SQL.Connection + -> IO () +defaultRollbackPlan tableName pointName extractor p c = + let deleteAllQuery tName = "DELETE FROM " <> tName + deleteAll = SQL.execute_ c . deleteAllQuery . SQL.Query . Text.pack + deleteUntilQuery tName pName = + deleteAllQuery tName <> " WHERE " <> pName <> " > :point" + deleteUntil :: (SQL.ToField a) => String -> String -> a -> IO () + deleteUntil tName pName pt = + SQL.executeNamed + c + (deleteUntilQuery (SQL.Query $ Text.pack tName) (SQL.Query $ Text.pack pName)) + [":point" SQL.:= pt] + in case extractor p of + Nothing -> deleteAll tableName + Just pt -> deleteUntil tableName pointName pt + runLastStablePointQuery :: (MonadError IndexerError m, MonadIO m, SQL.FromRow r) => SQL.Connection @@ -317,20 +332,7 @@ instance rollback p indexer = do let c = indexer ^. connection - deleteAllQuery tName = "DELETE FROM " <> tName - deleteAll = SQL.execute_ c . deleteAllQuery . SQL.Query . Text.pack - deleteUntilQuery tName pName = - deleteAllQuery tName <> " WHERE " <> pName <> " > :point" - deleteUntil :: (SQL.ToField a) => String -> String -> a -> IO () - deleteUntil tName pName pt = - SQL.executeNamed - c - (deleteUntilQuery (SQL.Query $ Text.pack tName) (SQL.Query $ Text.pack pName)) - [":point" SQL.:= pt] - rollbackTable (SQLRollbackPlan tableName pointName extractor) = - case extractor p of - Nothing -> deleteAll tableName - Just pt -> deleteUntil tableName pointName pt + rollbackTable (SQLRollbackPlan planRollback) = planRollback p c liftIO $ SQL.withTransaction c $ traverse_ rollbackTable (indexer ^. rollbackPlan) diff --git a/marconi-core/test/Marconi/CoreSpec.hs b/marconi-core/test/Marconi/CoreSpec.hs index 710e225aa0..39103e2b6d 100644 --- a/marconi-core/test/Marconi/CoreSpec.hs +++ b/marconi-core/test/Marconi/CoreSpec.hs @@ -679,7 +679,7 @@ sqliteModelIndexerWithFile filepath = do "INSERT INTO index_model VALUES (?, ?, ?)" ] ] - [Core.SQLRollbackPlan "index_model" "pointSlotNo" extractor] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "index_model" "pointSlotNo" extractor)] (Core.SetLastStablePointQuery "INSERT OR REPLACE INTO sync (pointSlotNo, pointHash) VALUES (?,?)") (Core.GetLastStablePointQuery "SELECT pointSlotNo, pointHash FROM sync") diff --git a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs index 5dd8bd368d..81c58ab2d0 100644 --- a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs +++ b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs @@ -199,7 +199,7 @@ mkAddressCountSqliteIndexer dbPath = do [ SQLInsertPlan eventToRows addressCountInsertQuery ] ] -- requests launched when an event is stored - [SQLRollbackPlan "address_count" "slotNo" C.chainPointToSlotNo] + [SQLRollbackPlan (Core.defaultRollbackPlan "address_count" "slotNo" C.chainPointToSlotNo)] where dbCreation = [sql|CREATE TABLE IF NOT EXISTS address_count From 016943337bf47006778957fe12f61b73c13c4134 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Mon, 29 Jan 2024 16:45:24 +0200 Subject: [PATCH 03/15] Add indexer to CatchupConfig Signed-off-by: Ana Pantilie --- .../Marconi/Cardano/ChainIndex/Indexers.hs | 17 +++++++--- .../Marconi/Cardano/ChainIndex/Indexers.hs | 3 +- .../Marconi/Cardano/Core/Indexer/Worker.hs | 16 ++++++---- .../src/Marconi/Cardano/Indexers/BlockInfo.hs | 4 +-- .../src/Marconi/Cardano/Indexers/Datum.hs | 5 ++- .../Marconi/Cardano/Indexers/EpochNonce.hs | 4 +-- .../src/Marconi/Cardano/Indexers/EpochSDD.hs | 4 +-- .../Cardano/Indexers/MintTokenEvent.hs | 4 +-- .../Cardano/Indexers/SnapshotBlockEvent.hs | 4 +-- .../src/Marconi/Cardano/Indexers/Spent.hs | 4 +-- .../src/Marconi/Cardano/Indexers/Utxo.hs | 4 +-- .../test-lib/Test/Integration.hs | 2 +- .../Marconi/Core/Transformer/WithCatchup.hs | 31 ++++++++++--------- .../src/Marconi/Sidechain/Experimental/Env.hs | 2 +- .../Sidechain/Experimental/Indexers.hs | 7 +++-- 15 files changed, 66 insertions(+), 45 deletions(-) diff --git a/marconi-cardano-chain-index/src/Marconi/Cardano/ChainIndex/Indexers.hs b/marconi-cardano-chain-index/src/Marconi/Cardano/ChainIndex/Indexers.hs index 4f9c20b70d..c85437cc9f 100644 --- a/marconi-cardano-chain-index/src/Marconi/Cardano/ChainIndex/Indexers.hs +++ b/marconi-cardano-chain-index/src/Marconi/Cardano/ChainIndex/Indexers.hs @@ -3,6 +3,7 @@ {-# LANGUAGE FlexibleInstances #-} {-# LANGUAGE LambdaCase #-} {-# LANGUAGE OverloadedStrings #-} +{-# LANGUAGE RankNTypes #-} {-# LANGUAGE TemplateHaskell #-} {-# LANGUAGE ViewPatterns #-} @@ -111,7 +112,7 @@ and expose a single coordinator to operate them -} buildIndexers :: SecurityParam - -> Core.CatchupConfig + -> (forall indexer event. Core.CatchupConfig indexer event) -- (Core.WithTransform Core.SQLiteIndexer (NonEmpty Datum.DatumInfo)) (WithDistance (Maybe (NonEmpty Datum.DatumInfo))) -> Utxo.UtxoIndexerConfig -> MintTokenEvent.MintTokenEventConfig -> ExtLedgerStateCoordinator.ExtLedgerStateWorkerConfig EpochEvent (WithDistance BlockEvent) @@ -249,6 +250,8 @@ epochNonceBuilder :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) => SecurityParam -> Core.CatchupConfig + (Core.WithTransform indexer Nonce.EpochNonce) + (WithDistance (Maybe Nonce.EpochNonce)) -> BM.Trace m Text -> FilePath -> n @@ -278,6 +281,8 @@ epochSDDBuilder :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) => SecurityParam -> Core.CatchupConfig + (Core.WithTransform indexer (NonEmpty SDD.EpochSDD)) + (WithDistance (Maybe (NonEmpty SDD.EpochSDD))) -> BM.Trace m Text -> FilePath -> n @@ -307,6 +312,8 @@ snapshotBlockEventBuilder :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) => SecurityParam -> Core.CatchupConfig + (Core.WithTransform indexer SnapshotBlockEvent) + (WithDistance (Maybe SnapshotBlockEvent)) -> BM.Trace m Text -> FilePath -> BlockRange @@ -346,7 +353,7 @@ extractSnapshotBlockEvent = -- | Builds the coordinators for each sub-chain serializers. buildIndexersForSnapshot :: SecurityParam - -> Core.CatchupConfig + -> (forall indexer event. Core.CatchupConfig indexer event) -> ExtLedgerStateCoordinator.ExtLedgerStateWorkerConfig (ExtLedgerStateEvent, WithDistance BlockEvent) (WithDistance BlockEvent) @@ -419,6 +426,8 @@ snapshotExtLedgerStateEventBuilder :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) => SecurityParam -> Core.CatchupConfig + (Core.WithTransform indexer ExtLedgerStateEvent) + (WithDistance (Maybe ExtLedgerStateEvent)) -> BM.Trace m Text -> FilePath -> BlockRange @@ -449,9 +458,9 @@ snapshotExtLedgerStateEventBuilder securityParam catchupConfig textLogger path b path snapshotExtLedgerStateEventWorker - :: forall input m n + :: forall indexer input m n . (MonadIO m, MonadError Core.IndexerError m, MonadIO n) - => StandardWorkerConfig n input ExtLedgerStateEvent + => StandardWorkerConfig n indexer input ExtLedgerStateEvent -> SnapshotWorkerConfig input -> FilePath -> m diff --git a/marconi-cardano-chain-index/test-lib/Test/Marconi/Cardano/ChainIndex/Indexers.hs b/marconi-cardano-chain-index/test-lib/Test/Marconi/Cardano/ChainIndex/Indexers.hs index 1946034da9..7175ba265a 100644 --- a/marconi-cardano-chain-index/test-lib/Test/Marconi/Cardano/ChainIndex/Indexers.hs +++ b/marconi-cardano-chain-index/test-lib/Test/Marconi/Cardano/ChainIndex/Indexers.hs @@ -1,4 +1,5 @@ {-# LANGUAGE OverloadedStrings #-} +{-# LANGUAGE RankNTypes #-} {-# LANGUAGE TemplateHaskell #-} {- | Generators and helpers for testing @Marconi.Cardano.Indexers@, namely the @@ -142,7 +143,7 @@ we cannot generate explicitly. -} buildIndexers :: SecurityParam - -> Core.CatchupConfig + -> (forall indexer event. Core.CatchupConfig indexer event) -> Utxo.UtxoIndexerConfig -> MintTokenEvent.MintTokenEventConfig -> ExtLedgerStateCoordinator.ExtLedgerStateWorkerConfig EpochEvent (WithDistance BlockEvent) diff --git a/marconi-cardano-core/src/Marconi/Cardano/Core/Indexer/Worker.hs b/marconi-cardano-core/src/Marconi/Cardano/Core/Indexer/Worker.hs index 0e2a5409ca..a4b0393ac5 100644 --- a/marconi-cardano-core/src/Marconi/Cardano/Core/Indexer/Worker.hs +++ b/marconi-cardano-core/src/Marconi/Cardano/Core/Indexer/Worker.hs @@ -26,6 +26,7 @@ import Marconi.Cardano.Core.Extract.WithDistance (WithDistance) import Marconi.Cardano.Core.Extract.WithDistance qualified as Distance import Marconi.Cardano.Core.Orphans () import Marconi.Cardano.Core.Types (SecurityParam) +import Marconi.Core (WithTransform) import Marconi.Core qualified as Core -- | An alias for an indexer with catchup and transformation to perform filtering @@ -38,10 +39,13 @@ type StandardIndexer m indexer event = -- | An alias for an SQLiteWorker with catchup and transformation to perform filtering type StandardSQLiteIndexer m event = StandardIndexer m Core.SQLiteIndexer event -data StandardWorkerConfig m input event = StandardWorkerConfig +type StandardCatchupConfig indexer event = + Core.CatchupConfig (WithTransform indexer event) (WithDistance (Maybe event)) + +data StandardWorkerConfig m indexer input event = StandardWorkerConfig { workerName :: Text , securityParamConfig :: !SecurityParam - , catchupConfig :: Core.CatchupConfig + , catchupConfig :: StandardCatchupConfig indexer event , eventExtractor :: input -> m (Maybe event) , logger :: Trace m (Core.IndexerEvent C.ChainPoint) } @@ -57,7 +61,7 @@ mkStandardIndexer :: ( MonadIO m , Core.Point event ~ C.ChainPoint ) - => StandardWorkerConfig m a b + => StandardWorkerConfig m indexer a event -> indexer event -> StandardIndexer m indexer event mkStandardIndexer config indexer = @@ -75,7 +79,7 @@ mkStandardWorker , Core.IsSync n event indexer , Core.Point event ~ C.ChainPoint ) - => StandardWorkerConfig m input event + => StandardWorkerConfig m indexer input event -> indexer event -> n (StandardWorker m input event indexer) mkStandardWorker config = mkStandardWorkerWithFilter config Just @@ -85,7 +89,7 @@ mkStandardIndexerWithFilter :: ( MonadIO m , Core.Point event ~ C.ChainPoint ) - => StandardWorkerConfig m a b + => StandardWorkerConfig m indexer a event -> (event -> Maybe event) -> indexer event -> StandardIndexer m indexer event @@ -105,7 +109,7 @@ mkStandardWorkerWithFilter , Ord (Core.Point event) , Core.Point event ~ C.ChainPoint ) - => StandardWorkerConfig m input event + => StandardWorkerConfig m indexer input event -> (event -> Maybe event) -> indexer event -> n (StandardWorker m input event indexer) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs index a19844f27f..8a1b09fe6d 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs @@ -151,7 +151,7 @@ blockInfoWorker , MonadError Core.IndexerError n , MonadIO m ) - => StandardWorkerConfig m input BlockInfo + => StandardWorkerConfig m Core.SQLiteIndexer input BlockInfo -- ^ General indexer configuration -> SQLiteDBLocation -- ^ SQLite database location @@ -166,7 +166,7 @@ creating 'StandardWorkerConfig', including a preprocessor. blockInfoBuilder :: (MonadIO n, MonadError Core.IndexerError n) => SecurityParam - -> Core.CatchupConfig + -> Core.CatchupConfig indexer event -> BM.Trace IO Text -> FilePath -> n (StandardWorker IO BlockEvent BlockInfo Core.SQLiteIndexer) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs index df8c6d5609..31237671d1 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs @@ -52,6 +52,7 @@ import Database.SQLite.Simple (NamedParam ((:=))) import Database.SQLite.Simple qualified as SQL import Database.SQLite.Simple.QQ (sql) import GHC.Generics (Generic) +import Marconi.Cardano.Core.Extract.WithDistance (WithDistance) import Marconi.Cardano.Core.Indexer.Worker ( StandardSQLiteIndexer, StandardWorker, @@ -130,7 +131,7 @@ mkDatumIndexer path = do -- | A worker with catchup for a 'DatumIndexer' datumWorker :: (MonadIO m, MonadIO n, MonadError Core.IndexerError n) - => StandardWorkerConfig m input DatumEvent + => StandardWorkerConfig m Core.SQLiteIndexer input DatumEvent -- ^ General configuration of the indexer (mostly for logging purpose) -> SQLiteDBLocation -- ^ SQLite database location @@ -146,6 +147,8 @@ datumBuilder :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) => SecurityParam -> Core.CatchupConfig + (Core.WithTransform Core.SQLiteIndexer (NonEmpty DatumInfo)) + (WithDistance (Maybe (NonEmpty DatumInfo))) -> Trace m Text -> FilePath -> n (StandardWorker m [AnyTxBody] DatumEvent Core.SQLiteIndexer) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs index 392eff4ec8..322ef88d0c 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs @@ -119,9 +119,9 @@ newtype EpochNonceWorkerConfig input = EpochNonceWorkerConfig } epochNonceWorker - :: forall input m n + :: forall indexer input m n . (MonadIO m, MonadError Core.IndexerError m, MonadIO n) - => StandardWorkerConfig n input EpochNonce + => StandardWorkerConfig n indexer input EpochNonce -> EpochNonceWorkerConfig input -> SQLiteDBLocation -> m (Core.WorkerIndexer n input EpochNonce (Core.WithTrace n Core.SQLiteIndexer)) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs index fd9017ba08..ff78e27cc2 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs @@ -126,9 +126,9 @@ newtype EpochSDDWorkerConfig input = EpochSDDWorkerConfig } epochSDDWorker - :: forall input m n + :: forall indexer input m n . (MonadIO m, MonadError Core.IndexerError m, MonadIO n) - => StandardWorkerConfig n input (NonEmpty EpochSDD) + => StandardWorkerConfig n indexer input (NonEmpty EpochSDD) -> EpochSDDWorkerConfig input -> SQLiteDBLocation -> m (Core.WorkerIndexer n input (NonEmpty EpochSDD) (Core.WithTrace n Core.SQLiteIndexer)) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs index 9b3aebb941..e867420d52 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs @@ -296,7 +296,7 @@ type instance Core.Point MintTokenBlockEvents = C.ChainPoint -- | Create a worker for the MintTokenEvent indexer mintTokenEventWorker :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) - => StandardWorkerConfig m input MintTokenBlockEvents + => StandardWorkerConfig m Core.SQLiteIndexer input MintTokenBlockEvents -- ^ General configuration of the indexer (mostly for logging purpose) -> MintTokenEventConfig -- ^ Specific configuration of the indexer (mostly for logging purpose and filtering for target @@ -315,7 +315,7 @@ creating 'StandardWorkerConfig', including a preprocessor. mintTokenEventBuilder :: (MonadIO n, MonadError Core.IndexerError n) => SecurityParam - -> Core.CatchupConfig + -> Core.CatchupConfig indexer event -> MintTokenEventConfig -> BM.Trace IO Text -> FilePath diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SnapshotBlockEvent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SnapshotBlockEvent.hs index c15e76acf1..2b28e811fc 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SnapshotBlockEvent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SnapshotBlockEvent.hs @@ -240,9 +240,9 @@ deserialiseMetadata _other = Nothing any blocks which are not in the given 'BlockRange'. -} snapshotBlockEventWorker - :: forall input m n + :: forall indexer input m n . (MonadIO m, MonadError Core.IndexerError m, MonadIO n) - => StandardWorkerConfig n input SnapshotBlockEvent + => StandardWorkerConfig n indexer input SnapshotBlockEvent -> SnapshotWorkerConfig input -> FilePath -> m diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index 98899c8be0..597f360307 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -163,7 +163,7 @@ catchupConfigEventHook stdoutTrace dbPath Core.Synced = do -- | A minimal worker for the UTXO indexer, with catchup and filtering. spentWorker :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) - => StandardWorkerConfig m input SpentInfoEvent + => StandardWorkerConfig m Core.SQLiteIndexer input SpentInfoEvent -- ^ General configuration of a worker -> SQLiteDBLocation -- ^ SQLite database location @@ -178,7 +178,7 @@ creating 'StandardWorkerConfig', including a preprocessor. spentBuilder :: (MonadIO n, MonadError Core.IndexerError n) => SecurityParam - -> Core.CatchupConfig + -> Core.CatchupConfig indexer event -> BM.Trace IO Text -> FilePath -> n (StandardWorker IO [AnyTxBody] SpentInfoEvent Core.SQLiteIndexer) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index 3a4a6616f2..4e7212344a 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -189,7 +189,7 @@ catchupConfigEventHook stdoutTrace dbPath Core.Synced = do -- | A minimal worker for the UTXO indexer, with catchup and filtering. utxoWorker :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) - => StandardWorkerConfig m input UtxoEvent + => StandardWorkerConfig m Core.SQLiteIndexer input UtxoEvent -- ^ General configuration of the indexer (mostly for logging purpose) -> UtxoIndexerConfig -- ^ Specific configuration of the indexer (mostly for logging purpose) @@ -230,7 +230,7 @@ creating 'StandardWorkerConfig', including a preprocessor. utxoBuilder :: (MonadIO n, MonadError Core.IndexerError n) => SecurityParam - -> Core.CatchupConfig + -> Core.CatchupConfig indexer event -> UtxoIndexerConfig -> BM.Trace IO Text -> FilePath diff --git a/marconi-cardano-indexers/test-lib/Test/Integration.hs b/marconi-cardano-indexers/test-lib/Test/Integration.hs index 63342a7bc7..4e91dc7f43 100644 --- a/marconi-cardano-indexers/test-lib/Test/Integration.hs +++ b/marconi-cardano-indexers/test-lib/Test/Integration.hs @@ -208,7 +208,7 @@ quantityToZero = C.valueFromList . map (fmap (const 0)) . C.valueToList {- | @Core.'CatchupConfig'@ with values suitable for end-to-end tests. The current values are taken from those hard-coded in the marconi-cardano-chain-index application. -} -mkEndToEndCatchupConfig :: Core.CatchupConfig +mkEndToEndCatchupConfig :: Core.CatchupConfig indexer event mkEndToEndCatchupConfig = Core.mkCatchupConfig 5000 100 {- | Build a @Runner.'RunIndexerConfig'@ with values suitable for end-to-end tests with diff --git a/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs b/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs index 0107083e76..98a91c9221 100644 --- a/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs +++ b/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs @@ -25,6 +25,7 @@ import Control.Lens qualified as Lens import Control.Lens.Operators ((%~), (+~), (.~), (?~), (^.)) import Control.Monad.IO.Class (MonadIO, liftIO) import Data.Function ((&)) +import Data.Kind (Type) import Data.Word (Word64) import Marconi.Core.Class ( Closeable, @@ -53,7 +54,7 @@ import Marconi.Core.Type (Point, Timed (Timed), point) {- | The visible part of the catchup configuration, it allows you to configure the size of the batch and to control when the batch mecanism stops -} -data CatchupConfig = CatchupConfig +data CatchupConfig (indexer :: Type -> Type) (event :: Type) = CatchupConfig { _configCatchupBatchSize :: Word64 -- ^ Maximal number of events in one batch , _configCatchupBypassDistance :: Word64 @@ -64,15 +65,15 @@ data CatchupConfig = CatchupConfig data CatchupEvent = Synced -mkCatchupConfig :: Word64 -> Word64 -> CatchupConfig +mkCatchupConfig :: Word64 -> Word64 -> CatchupConfig indexer event mkCatchupConfig batchSize bypassDistance = CatchupConfig batchSize bypassDistance Nothing Lens.makeLenses ''CatchupConfig -data CatchupContext event = CatchupContext +data CatchupContext (indexer :: Type -> Type) (event :: Type) = CatchupContext { _contextDistanceComputation :: Point event -> event -> Word64 -- ^ How we compute distance to tip - , _contextCatchupConfig :: CatchupConfig + , _contextCatchupConfig :: CatchupConfig indexer event -- ^ How far from the block should we be to bypass the catchup mechanism (in number of blocks) , _contextCatchupBufferLength :: Word64 -- ^ How many event do we have in the batch @@ -84,22 +85,24 @@ data CatchupContext event = CatchupContext Lens.makeLenses ''CatchupContext -contextCatchupBypassDistance :: Lens.Lens' (CatchupContext event) Word64 +contextCatchupBypassDistance :: Lens.Lens' (CatchupContext indexer event) Word64 contextCatchupBypassDistance = contextCatchupConfig . configCatchupBypassDistance -contextCatchupBatchSize :: Lens.Lens' (CatchupContext event) Word64 +contextCatchupBatchSize :: Lens.Lens' (CatchupContext indexer event) Word64 contextCatchupBatchSize = contextCatchupConfig . configCatchupBatchSize -contextCatchupEventHook :: Lens.Lens' (CatchupContext event) (Maybe (CatchupEvent -> IO ())) +contextCatchupEventHook :: Lens.Lens' (CatchupContext indexer event) (Maybe (CatchupEvent -> IO ())) contextCatchupEventHook = contextCatchupConfig . configCatchupEventHook {-- | WithCatchup is used to speed up the synchronisation of indexers by preparing batches of events - that will be submitted via `indexAll` to the underlying indexer. - - - Once the indexer is close enough to the tip, the transformer stops the batch to insert element + - Once the indexer is close enough to the tip, the transformer stops the batch to insert elements - one by one. -} -newtype WithCatchup indexer event = WithCatchup {_catchupWrapper :: IndexTransformer CatchupContext indexer event} +newtype WithCatchup (indexer :: Type -> Type) (event :: Type) = WithCatchup + { _catchupWrapper :: IndexTransformer (CatchupContext indexer) indexer event + } Lens.makeLenses 'WithCatchup @@ -107,7 +110,7 @@ Lens.makeLenses 'WithCatchup withCatchup :: (Point event -> event -> Word64) -- ^ The distance function - -> CatchupConfig + -> CatchupConfig indexer event -- ^ Configure how many element we put in a batch and until when we use it -> indexer event -- ^ the underlying indexer @@ -116,22 +119,22 @@ withCatchup computeDistance config = WithCatchup . IndexTransformer (CatchupContext computeDistance config 0 [] Nothing) deriving via - (IndexTransformer CatchupContext indexer) + (IndexTransformer (CatchupContext indexer) indexer) instance (IsSync m event indexer) => IsSync m event (WithCatchup indexer) deriving via - (IndexTransformer CatchupContext indexer) + (IndexTransformer (CatchupContext indexer) indexer) instance (HasDatabasePath indexer) => HasDatabasePath (WithCatchup indexer) deriving via - (IndexTransformer CatchupContext indexer) + (IndexTransformer (CatchupContext indexer) indexer) instance (Closeable m indexer) => Closeable m (WithCatchup indexer) deriving via - (IndexTransformer CatchupContext indexer) + (IndexTransformer (CatchupContext indexer) indexer) instance (Queryable m event query indexer) => Queryable m event query (WithCatchup indexer) diff --git a/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Env.hs b/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Env.hs index 49346a2fff..f11a6f5c83 100644 --- a/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Env.hs +++ b/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Env.hs @@ -89,7 +89,7 @@ querySecurityParamFromCliArgs trace CliArgs{..} = -- | Build the 'SidechainBuildIndexersConfig' from CLI arguments. mkSidechainBuildIndexersConfig - :: Trace IO Text -> CliArgs -> SecurityParam -> SidechainBuildIndexersConfig + :: Trace IO Text -> CliArgs -> SecurityParam -> SidechainBuildIndexersConfig indexer event mkSidechainBuildIndexersConfig trace CliArgs{..} securityParam = let filteredAddresses = diff --git a/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Indexers.hs b/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Indexers.hs index 2b2b74d8dc..ad917c544e 100644 --- a/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Indexers.hs +++ b/marconi-sidechain-experimental/src/Marconi/Sidechain/Experimental/Indexers.hs @@ -1,3 +1,4 @@ +{-# LANGUAGE RankNTypes #-} {-# LANGUAGE TemplateHaskell #-} module Marconi.Sidechain.Experimental.Indexers where @@ -37,10 +38,10 @@ data SidechainRunIndexersConfig = SidechainRunIndexersConfig {- | Configuration for constructing marconi indexers as used in this package. -} -data SidechainBuildIndexersConfig = SidechainBuildIndexersConfig +data SidechainBuildIndexersConfig indexer event = SidechainBuildIndexersConfig { _sidechainBuildIndexersTrace :: !(Trace IO Text) , _sidechainBuildIndexersSecurityParam :: !SecurityParam - , _sidechainBuildIndexersCatchupConfig :: !CatchupConfig + , _sidechainBuildIndexersCatchupConfig :: !(CatchupConfig indexer event) , _sidechainBuildIndexersDbPath :: !FilePath , _sidechainBuildIndexersEpochStateConfig :: !(EpochState.ExtLedgerStateWorkerConfig EpochEvent (WithDistance BlockEvent)) @@ -58,7 +59,7 @@ similarly to the marconi-cardano-chain-index application. -} sidechainBuildIndexers :: (MonadIO m) - => SidechainBuildIndexersConfig + => (forall indexer event. SidechainBuildIndexersConfig indexer event) -> m (Either IndexerError (C.ChainPoint, MarconiCardanoQueryables, SyncStatsCoordinator)) sidechainBuildIndexers config = liftIO . runExceptT $ From b2991465669b30169f19c428331195a2c380f7fb Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Tue, 30 Jan 2024 13:22:27 +0200 Subject: [PATCH 04/15] Change type signature of WithCatchup hook Signed-off-by: Ana Pantilie --- .../src/Marconi/Cardano/Indexers/BlockInfo.hs | 9 ++++--- .../Cardano/Indexers/MintTokenEvent.hs | 9 ++++--- .../src/Marconi/Cardano/Indexers/Spent.hs | 10 +++++--- .../src/Marconi/Cardano/Indexers/Utxo.hs | 8 +++--- marconi-core/src/Marconi/Core.hs | 4 +-- .../Marconi/Core/Transformer/WithCatchup.hs | 25 +++++++++++-------- 6 files changed, 36 insertions(+), 29 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs index 8a1b09fe6d..5a54969a53 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs @@ -37,7 +37,7 @@ module Marconi.Cardano.Indexers.BlockInfo ( import Cardano.Api qualified as C import Cardano.BM.Data.Trace (Trace) import Cardano.BM.Tracing qualified as BM -import Control.Lens ((&), (?~), (^.)) +import Control.Lens ((&), (.~), (^.)) import Control.Lens qualified as Lens import Control.Monad (when) import Control.Monad.Except (MonadError (throwError)) @@ -135,8 +135,8 @@ mkBlockInfoIndexer path = do blockInfoInsertQuery (Core.SQLRollbackPlan (Core.defaultRollbackPlan "blockInfo" "slotNo" C.chainPointToSlotNo)) -catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () -catchupConfigEventHook indexerName stdoutTrace dbPath Core.Synced = do +catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> indexer -> IO indexer +catchupConfigEventHook indexerName stdoutTrace dbPath indexer = do SQL.withConnection dbPath $ \c -> do let slotNoIndexName = "blockInfo__slotNo" createSlotNoIndexStatement = @@ -144,6 +144,7 @@ catchupConfigEventHook indexerName stdoutTrace dbPath Core.Synced = do <> fromString slotNoIndexName <> " ON blockInfo (slotNo)" Core.createIndexTable indexerName stdoutTrace c slotNoIndexName createSlotNoIndexStatement + pure indexer -- | Create a worker for 'BlockInfoIndexer' with catchup blockInfoWorker @@ -177,7 +178,7 @@ blockInfoBuilder securityParam catchupConfig textLogger path = catchupConfigWithTracer = catchupConfig & Core.configCatchupEventHook - ?~ catchupConfigEventHook indexerName textLogger blockInfoDbPath + .~ catchupConfigEventHook indexerName textLogger blockInfoDbPath blockInfoWorkerConfig = StandardWorkerConfig indexerName diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs index e867420d52..df5614cf93 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs @@ -106,7 +106,6 @@ import Control.Lens ( view, (%~), (.~), - (?~), (^.), (^?), ) @@ -327,7 +326,7 @@ mintTokenEventBuilder securityParam catchupConfig mintEventConfig textLogger pat catchupConfigWithTracer = catchupConfig & Core.configCatchupEventHook - ?~ catchupConfigEventHook indexerName textLogger mintDbPath + .~ catchupConfigEventHook indexerName textLogger mintDbPath extractMint :: AnyTxBody -> [MintTokenEvent] extractMint (AnyTxBody bn ix txb) = extractEventsFromTx bn ix txb mintTokenWorkerConfig = @@ -449,8 +448,8 @@ mkMintTokenIndexer dbPath = do (Core.defaultRollbackPlan "minting_policy_events" "slotNo" C.chainPointToSlotNo) ] -catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () -catchupConfigEventHook indexerName stdoutTrace dbPath Core.Synced = do +catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> indexer -> IO indexer +catchupConfigEventHook indexerName stdoutTrace dbPath indexer = do SQL.withConnection dbPath $ \c -> do let txIdPolicyIdIndexName = "minting_policy_events__txId_policyId" createMintPolicyIdIndexStatement = @@ -470,6 +469,8 @@ catchupConfigEventHook indexerName stdoutTrace dbPath Core.Synced = do <> fromString slotNoIndexName <> " ON minting_policy_events (slotNo)" Core.createIndexTable indexerName stdoutTrace c slotNoIndexName createSlotNoIndexStatement + -- return the original indexer + pure indexer fromTimedMintEvents :: Core.Timed point MintTokenBlockEvents diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index 597f360307..f5067d22e3 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -28,7 +28,7 @@ module Marconi.Cardano.Indexers.Spent ( import Cardano.Api qualified as C import Cardano.BM.Trace (Trace) import Cardano.BM.Tracing qualified as BM -import Control.Lens ((&), (?~), (^.)) +import Control.Lens ((&), (.~), (^.)) import Control.Lens qualified as Lens import Control.Monad.Except (MonadError) import Control.Monad.IO.Class (MonadIO) @@ -136,8 +136,8 @@ mkSpentIndexer path = do [spentInsert] [Core.SQLRollbackPlan (Core.defaultRollbackPlan "spent" "slotNo" C.chainPointToSlotNo)] -catchupConfigEventHook :: Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () -catchupConfigEventHook stdoutTrace dbPath Core.Synced = do +catchupConfigEventHook :: Trace IO Text -> FilePath -> indexer -> IO indexer +catchupConfigEventHook stdoutTrace dbPath indexer = do SQL.withConnection dbPath $ \c -> do let slotNoIndexName = "spent_slotNo" createSlotNoIndexStatement = @@ -159,6 +159,8 @@ catchupConfigEventHook stdoutTrace dbPath Core.Synced = do <> fromString spentAtTxIdIndexName <> " ON spent (spentAtTxId)" Core.createIndexTable "Spent" stdoutTrace c spentAtTxIdIndexName createSpentAtIndexStatement + -- return the original indexer + pure indexer -- | A minimal worker for the UTXO indexer, with catchup and filtering. spentWorker @@ -191,7 +193,7 @@ spentBuilder securityParam catchupConfig textLogger path = catchupConfigWithTracer = catchupConfig & Core.configCatchupEventHook - ?~ catchupConfigEventHook textLogger spentDbPath + .~ catchupConfigEventHook textLogger spentDbPath spentWorkerConfig = StandardWorkerConfig indexerName diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index 4e7212344a..021a02637e 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -54,7 +54,6 @@ import Cardano.BM.Tracing qualified as BM import Control.Lens ( (&), (.~), - (?~), (^.), ) import Control.Lens qualified as Lens @@ -169,8 +168,8 @@ mkUtxoIndexer path = do [insertEvent] [Core.SQLRollbackPlan (Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo)] -catchupConfigEventHook :: Trace IO Text -> FilePath -> Core.CatchupEvent -> IO () -catchupConfigEventHook stdoutTrace dbPath Core.Synced = do +catchupConfigEventHook :: Trace IO Text -> FilePath -> indexer -> IO indexer +catchupConfigEventHook stdoutTrace dbPath indexer = do SQL.withConnection dbPath $ \c -> do let addressIndexName = "utxo_address" createAddressIndexStatement = @@ -185,6 +184,7 @@ catchupConfigEventHook stdoutTrace dbPath Core.Synced = do <> fromString slotNoIndexName <> " ON utxo (slotNo)" Core.createIndexTable "Utxo" stdoutTrace c slotNoIndexName createSlotNoIndexStatement + pure indexer -- | A minimal worker for the UTXO indexer, with catchup and filtering. utxoWorker @@ -243,7 +243,7 @@ utxoBuilder securityParam catchupConfig utxoConfig textLogger path = extractUtxos (AnyTxBody _ indexInBlock txb) = getUtxosFromTxBody indexInBlock txb catchupConfigWithTracer = catchupConfig - & Core.configCatchupEventHook ?~ catchupConfigEventHook textLogger utxoDbPath + & Core.configCatchupEventHook .~ catchupConfigEventHook textLogger utxoDbPath utxoWorkerConfig = StandardWorkerConfig indexerName diff --git a/marconi-core/src/Marconi/Core.hs b/marconi-core/src/Marconi/Core.hs index 7e0157e69f..d0b6baa9e6 100644 --- a/marconi-core/src/Marconi/Core.hs +++ b/marconi-core/src/Marconi/Core.hs @@ -456,7 +456,7 @@ module Marconi.Core ( CatchupConfig (CatchupConfig), mkCatchupConfig, configCatchupEventHook, - HasCatchupConfig (catchupBypassDistance, catchupBatchSize, catchupEventHook), + HasCatchupConfig (catchupBypassDistance, catchupBatchSize), CatchupEvent (Synced), -- *** SQLite @@ -694,7 +694,7 @@ import Marconi.Core.Transformer.WithCache ( import Marconi.Core.Transformer.WithCatchup ( CatchupConfig (CatchupConfig), CatchupEvent (Synced), - HasCatchupConfig (catchupBatchSize, catchupBypassDistance, catchupEventHook), + HasCatchupConfig (catchupBatchSize, catchupBypassDistance), WithCatchup, configCatchupEventHook, mkCatchupConfig, diff --git a/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs b/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs index 98a91c9221..dde7d47710 100644 --- a/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs +++ b/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs @@ -59,14 +59,14 @@ data CatchupConfig (indexer :: Type -> Type) (event :: Type) = CatchupConfig -- ^ Maximal number of events in one batch , _configCatchupBypassDistance :: Word64 -- ^ How far from the block should we be to bypass the catchup mechanism (in number of blocks). - , _configCatchupEventHook :: Maybe (CatchupEvent -> IO ()) + , _configCatchupEventHook :: indexer event -> IO (indexer event) -- ^ Hook to execute when specific events are triggered. } data CatchupEvent = Synced mkCatchupConfig :: Word64 -> Word64 -> CatchupConfig indexer event -mkCatchupConfig batchSize bypassDistance = CatchupConfig batchSize bypassDistance Nothing +mkCatchupConfig batchSize bypassDistance = CatchupConfig batchSize bypassDistance pure Lens.makeLenses ''CatchupConfig @@ -91,7 +91,8 @@ contextCatchupBypassDistance = contextCatchupConfig . configCatchupBypassDistanc contextCatchupBatchSize :: Lens.Lens' (CatchupContext indexer event) Word64 contextCatchupBatchSize = contextCatchupConfig . configCatchupBatchSize -contextCatchupEventHook :: Lens.Lens' (CatchupContext indexer event) (Maybe (CatchupEvent -> IO ())) +contextCatchupEventHook + :: Lens.Lens' (CatchupContext indexer event) (indexer event -> IO (indexer event)) contextCatchupEventHook = contextCatchupConfig . configCatchupEventHook {-- | WithCatchup is used to speed up the synchronisation of indexers by preparing batches of events @@ -118,6 +119,12 @@ withCatchup withCatchup computeDistance config = WithCatchup . IndexTransformer (CatchupContext computeDistance config 0 [] Nothing) +runCatchupHook :: WithCatchup indexer event -> IO (WithCatchup indexer event) +runCatchupHook (WithCatchup (IndexTransformer context indexer)) = do + let catchupHook = context ^. contextCatchupEventHook + newIndexer <- catchupHook indexer + pure (WithCatchup (IndexTransformer context newIndexer)) + deriving via (IndexTransformer (CatchupContext indexer) indexer) instance @@ -150,12 +157,10 @@ instance IndexerTrans WithCatchup where class HasCatchupConfig indexer where catchupBypassDistance :: Lens.Lens' (indexer event) Word64 catchupBatchSize :: Lens.Lens' (indexer event) Word64 - catchupEventHook :: Lens.Lens' (indexer event) (Maybe (CatchupEvent -> IO ())) instance {-# OVERLAPPING #-} HasCatchupConfig (WithCatchup indexer) where catchupBypassDistance = catchupWrapper . wrapperConfig . contextCatchupBypassDistance catchupBatchSize = catchupWrapper . wrapperConfig . contextCatchupBatchSize - catchupEventHook = catchupWrapper . wrapperConfig . contextCatchupEventHook instance {-# OVERLAPPABLE #-} @@ -164,7 +169,6 @@ instance where catchupBypassDistance = unwrap . catchupBypassDistance catchupBatchSize = unwrap . catchupBatchSize - catchupEventHook = unwrap . catchupEventHook instance {-# OVERLAPPABLE #-} @@ -173,7 +177,6 @@ instance where catchupBypassDistance = unwrapMap . catchupBypassDistance catchupBatchSize = unwrapMap . catchupBatchSize - catchupEventHook = unwrapMap . catchupEventHook catchupDistance :: Lens.Lens' (WithCatchup indexer event) (Point event -> event -> Word64) catchupDistance = catchupWrapper . wrapperConfig . contextDistanceComputation @@ -211,11 +214,11 @@ instance pure $ resetBuffer ix'' in if hasCaughtUp then do - maybe (pure ()) (\f -> liftIO $ f Synced) $ indexer ^. catchupEventHook + newIndexer <- liftIO $ runCatchupHook indexer indexer' <- - if null (indexer ^. catchupBuffer) - then pure indexer - else sendBatch indexer + if null (newIndexer ^. catchupBuffer) + then pure newIndexer + else sendBatch newIndexer indexVia caughtUpIndexer timedEvent indexer' else do let indexer' = pushEvent indexer From de21b2c3b826975c39fc674a8671944b4dbd71d7 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Wed, 31 Jan 2024 15:07:33 +0200 Subject: [PATCH 05/15] Add UtxoOrSpent sum type as event Signed-off-by: Ana Pantilie --- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 31 ++++++++++++++----- 1 file changed, 23 insertions(+), 8 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 7676c21e96..aff3f014ef 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -9,6 +9,7 @@ module Marconi.Cardano.Indexers.UtxoWithSpent ( ) where import Cardano.Api qualified as C +import Control.Applicative ((<|>)) import Control.Lens ((^.)) import Control.Lens qualified as Lens import Control.Monad.Error.Class (MonadError) @@ -24,7 +25,9 @@ import Database.SQLite.Simple.ToField (ToField (toField)) import Marconi.Cardano.Core.Indexer.Worker (StandardSQLiteIndexer) import Marconi.Cardano.Core.Orphans () import Marconi.Cardano.Core.Types (TxIndexInBlock) +import Marconi.Cardano.Indexers.Spent (SpentInfo) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync +import Marconi.Cardano.Indexers.Utxo (Utxo) import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) import Marconi.Core.Indexer.SQLiteIndexer qualified as Core import Marconi.Core.Type qualified as Core @@ -111,12 +114,22 @@ instance FromRow UtxoWithSpent where txIxSpent <- SQL.field pure . pure $ C.TxIn txIdSpent txIxSpent --- | An alias for a non-empty list of @Utxo@, it's the event potentially produced on each block -type UtxoWithSpentEvent = NonEmpty UtxoWithSpent +type UtxoOrSpentEvent = NonEmpty UtxoOrSpent -type instance Core.Point UtxoWithSpentEvent = C.ChainPoint -type UtxoWithSpentIndexer = Core.SQLiteIndexer UtxoWithSpentEvent -type StandardUtxoWithSpentIndexer m = StandardSQLiteIndexer m UtxoWithSpentEvent +data UtxoOrSpent + = UtxoEvent Utxo + | SpentEvent SpentInfo + +instance SQL.ToRow (Core.Timed C.ChainPoint UtxoOrSpent) where + toRow (Core.Timed cp (UtxoEvent utxoEvent)) = toRow (Core.Timed cp utxoEvent) + toRow (Core.Timed cp (SpentEvent spentEvent)) = toRow (Core.Timed cp spentEvent) + +instance FromRow (Core.Timed C.ChainPoint UtxoOrSpent) where + fromRow = fmap UtxoEvent <$> fromRow <|> fmap SpentEvent <$> fromRow + +type instance Core.Point UtxoOrSpentEvent = C.ChainPoint +type UtxoWithSpentIndexer = Core.SQLiteIndexer UtxoOrSpentEvent +type StandardUtxoWithSpentIndexer m = StandardSQLiteIndexer m UtxoOrSpentEvent -- | Make a SQLiteIndexer for UtxoWithSpent mkUtxoWithSpentIndexer @@ -139,11 +152,13 @@ mkUtxoWithSpentIndexer path = do , blockHeaderHash BLOB NOT NULL , txIdSpent TEXT , txIxSpent INT - )|] - -- TODO: fix + ) + as SELECT * FROM utxo LEFT OUTER JOIN spent ON utxo.txId == spent.txId AND utxo.txIx == spent.txIx + |] + -- TODO: fix not good need 2 types of insert utxoInsertQuery :: SQL.Query utxoInsertQuery = - [sql|INSERT INTO utxo ( + [sql|INSERT OR REPLACE INTO utxo_with_spent ( address, txIndex, txId, From 62616b14f64da49fdb921de719ddc0ab858547ce Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Wed, 31 Jan 2024 18:27:52 +0200 Subject: [PATCH 06/15] Generalise InsertPlan Signed-off-by: Ana Pantilie --- .../doc/marconi-as-a-library/tutorials/BasicApp.hs | 2 +- .../src/Marconi/Cardano/Indexers/BlockInfo.hs | 2 +- .../src/Marconi/Cardano/Indexers/Datum.hs | 2 +- .../src/Marconi/Cardano/Indexers/EpochNonce.hs | 2 +- .../src/Marconi/Cardano/Indexers/EpochSDD.hs | 2 +- .../src/Marconi/Cardano/Indexers/MintTokenEvent.hs | 2 +- .../src/Marconi/Cardano/Indexers/Spent.hs | 2 +- .../src/Marconi/Cardano/Indexers/SyncHelper.hs | 4 ++-- .../src/Marconi/Cardano/Indexers/Utxo.hs | 2 +- .../src/Marconi/Cardano/Indexers/UtxoWithSpent.hs | 2 +- marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs | 9 +++++---- marconi-core/test/Marconi/CoreSpec.hs | 2 +- .../src/Marconi/Starter/Indexers/AddressCount.hs | 2 +- 13 files changed, 18 insertions(+), 17 deletions(-) diff --git a/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs b/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs index e79bf62642..abbb0fe0ea 100644 --- a/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs +++ b/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs @@ -336,7 +336,7 @@ mkBlockInfoSqliteIndexer dbPath = do -- single row. List.singleton -- The query that is called for each row that is the output of the previous parameter. - blockInfoInsertQuery + (pure blockInfoInsertQuery) ] ] -- Requests launched when a rollback occurs diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs index 5a54969a53..b417e2a656 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs @@ -132,7 +132,7 @@ mkBlockInfoIndexer path = do path id createBlockInfoTable - blockInfoInsertQuery + (pure blockInfoInsertQuery) (Core.SQLRollbackPlan (Core.defaultRollbackPlan "blockInfo" "slotNo" C.chainPointToSlotNo)) catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> indexer -> IO indexer diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs index 31237671d1..d5dea9a4df 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs @@ -125,7 +125,7 @@ mkDatumIndexer path = do Sync.mkSyncedSqliteIndexer path createDatumTables - [[Core.SQLInsertPlan (traverse NonEmpty.toList) datumInsertQuery]] + [[Core.SQLInsertPlan (traverse NonEmpty.toList) (pure datumInsertQuery)]] [Core.SQLRollbackPlan (Core.defaultRollbackPlan "datum" "slotNo" C.chainPointToSlotNo)] -- | A worker with catchup for a 'DatumIndexer' diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs index 322ef88d0c..2d2c2fc2ac 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs @@ -107,7 +107,7 @@ mkEpochNonceIndexer path = do , slotNo , blockHeaderHash ) VALUES (?, ?, ?, ?, ?)|] - insertEvent = [Core.SQLInsertPlan pure nonceInsertQuery] + insertEvent = [Core.SQLInsertPlan pure (pure nonceInsertQuery)] Sync.mkSyncedSqliteIndexer path [createNonce] diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs index ff78e27cc2..0fcb5e26fc 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs @@ -114,7 +114,7 @@ mkEpochSDDIndexer path = do , slotNo , blockHeaderHash ) VALUES (?, ?, ?, ?, ?, ?)|] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) sddInsertQuery] + insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure sddInsertQuery)] Sync.mkSyncedSqliteIndexer path [createSDD] diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs index df5614cf93..f5c7173a27 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs @@ -439,7 +439,7 @@ mkMintTokenIndexer dbPath = do ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] createMintPolicyEventTables = [createMintPolicyEvent] - mintInsertPlans = [Core.SQLInsertPlan fromTimedMintEvents mintEventInsertQuery] + mintInsertPlans = [Core.SQLInsertPlan fromTimedMintEvents (pure mintEventInsertQuery)] Sync.mkSyncedSqliteIndexer dbPath createMintPolicyEventTables diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index f5067d22e3..c4f2400f15 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -129,7 +129,7 @@ mkSpentIndexer path = do VALUES (?, ?, ?, ?, ?)|] createSpentTables = [createSpent] spentInsert = - [Core.SQLInsertPlan (traverse NonEmpty.toList) spentInsertQuery] + [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure spentInsertQuery)] Sync.mkSyncedSqliteIndexer path createSpentTables diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs index 9ed5d6deb5..2806e57aba 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs @@ -91,8 +91,8 @@ mkSingleInsertSyncedSqliteIndexer -> (Core.Timed (Core.Point event) event -> param) -> SQL.Query -- ^ the creation query - -> SQL.Query - -- ^ the insert query + -> IO SQL.Query + -- ^ the action producing an insert query -> Core.SQLRollbackPlan (Core.Point event) -- ^ the rollback query -> m (Core.SQLiteIndexer event) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index 021a02637e..c602f8299d 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -160,7 +160,7 @@ mkUtxoIndexer path = do ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] createUtxoTables = [createUtxo] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) utxoInsertQuery] + insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure utxoInsertQuery)] Sync.mkSyncedSqliteIndexer path diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index aff3f014ef..363dacaaf5 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -174,7 +174,7 @@ mkUtxoWithSpentIndexer path = do ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] createUtxoTables = [createUtxoWithSpent] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) utxoInsertQuery] + insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure utxoInsertQuery)] Sync.mkSyncedSqliteIndexer path diff --git a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs index 479eb77a99..f18000869a 100644 --- a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs +++ b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs @@ -94,7 +94,7 @@ data SQLInsertPlan event = forall a. SQLInsertPlan { planExtractor :: Timed (Point event) event -> [a] -- ^ How to transform the event into a type that can be handle by the database - , planInsert :: SQL.Query + , planInsert :: IO SQL.Query -- ^ The insert statement for the extracted data } @@ -225,7 +225,7 @@ mkSingleInsertSqliteIndexer -- ^ The SQL query to fetch the last stable point from the indexer. -> m (SQLiteIndexer event) mkSingleInsertSqliteIndexer path extract create insert rollback' = - mkSqliteIndexer path [create] [[SQLInsertPlan (pure . extract) insert]] [rollback'] + mkSqliteIndexer path [create] [[SQLInsertPlan (pure . extract) (pure insert)]] [rollback'] -- | Map SQLite errors to an indexer error handleSQLErrors :: IO a -> IO (Either IndexerError a) @@ -246,10 +246,11 @@ runIndexQueriesStep _ _ [] = pure () runIndexQueriesStep c events plan = let runIndexQuery (SQLInsertPlan planExtractor planInsert) = do let rows = planExtractor =<< events + query <- planInsert case rows of [] -> pure () - [x] -> SQL.execute c planInsert x - _nonEmpty -> SQL.executeMany c planInsert rows + [x] -> SQL.execute c query x + _nonEmpty -> SQL.executeMany c query rows in Async.mapConcurrently_ runIndexQuery plan -- | Run a list of insert queries in one single transaction. diff --git a/marconi-core/test/Marconi/CoreSpec.hs b/marconi-core/test/Marconi/CoreSpec.hs index 39103e2b6d..5e9f3fc9e2 100644 --- a/marconi-core/test/Marconi/CoreSpec.hs +++ b/marconi-core/test/Marconi/CoreSpec.hs @@ -676,7 +676,7 @@ sqliteModelIndexerWithFile filepath = do ) ] ) - "INSERT INTO index_model VALUES (?, ?, ?)" + (pure "INSERT INTO index_model VALUES (?, ?, ?)") ] ] [Core.SQLRollbackPlan (Core.defaultRollbackPlan "index_model" "pointSlotNo" extractor)] diff --git a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs index 81c58ab2d0..46c7245772 100644 --- a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs +++ b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs @@ -196,7 +196,7 @@ mkAddressCountSqliteIndexer dbPath = do dbPath [dbCreation] -- request launched when the indexer is created [ - [ SQLInsertPlan eventToRows addressCountInsertQuery + [ SQLInsertPlan eventToRows (pure addressCountInsertQuery) ] ] -- requests launched when an event is stored [SQLRollbackPlan (Core.defaultRollbackPlan "address_count" "slotNo" C.chainPointToSlotNo)] From 0e74988c47d2a88514db5c01ec09de17ef20a2f1 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Thu, 1 Feb 2024 13:26:05 +0200 Subject: [PATCH 07/15] Fix insert plan Signed-off-by: Ana Pantilie --- .../src/Marconi/Cardano/Indexers/SyncHelper.hs | 2 +- marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs index 2806e57aba..064193fc20 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs @@ -91,7 +91,7 @@ mkSingleInsertSyncedSqliteIndexer -> (Core.Timed (Core.Point event) event -> param) -> SQL.Query -- ^ the creation query - -> IO SQL.Query + -> ([param] -> SQL.Query) -- ^ the action producing an insert query -> Core.SQLRollbackPlan (Core.Point event) -- ^ the rollback query diff --git a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs index f18000869a..9c66bc9b95 100644 --- a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs +++ b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs @@ -94,8 +94,8 @@ data SQLInsertPlan event = forall a. SQLInsertPlan { planExtractor :: Timed (Point event) event -> [a] -- ^ How to transform the event into a type that can be handle by the database - , planInsert :: IO SQL.Query - -- ^ The insert statement for the extracted data + , planInsert :: [a] -> SQL.Query + -- ^ The insert statement builder for the extracted data } newtype InsertPointQuery = InsertPointQuery {getInsertPointQuery :: SQL.Query} @@ -246,7 +246,7 @@ runIndexQueriesStep _ _ [] = pure () runIndexQueriesStep c events plan = let runIndexQuery (SQLInsertPlan planExtractor planInsert) = do let rows = planExtractor =<< events - query <- planInsert + query = planInsert rows case rows of [] -> pure () [x] -> SQL.execute c query x From 1e2e0f3c3572074795d6b972a29f1227187d864f Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Thu, 1 Feb 2024 17:13:10 +0200 Subject: [PATCH 08/15] Fix insert strategy AGAIN Signed-off-by: Ana Pantilie --- .../tutorials/BasicApp.hs | 14 +++-- .../marconi-cardano-indexers.cabal | 1 - .../src/Marconi/Cardano/Indexers/BlockInfo.hs | 2 +- .../src/Marconi/Cardano/Indexers/Datum.hs | 3 +- .../Marconi/Cardano/Indexers/EpochNonce.hs | 4 +- .../src/Marconi/Cardano/Indexers/EpochSDD.hs | 4 +- .../Cardano/Indexers/MintTokenEvent.hs | 3 +- .../src/Marconi/Cardano/Indexers/Spent.hs | 3 +- .../Marconi/Cardano/Indexers/SyncHelper.hs | 6 +- .../src/Marconi/Cardano/Indexers/Utxo.hs | 3 +- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 57 +++++++++---------- marconi-core/src/Marconi/Core.hs | 2 +- .../src/Marconi/Core/Indexer/SQLiteIndexer.hs | 42 ++++++++------ marconi-core/test/Marconi/CoreSpec.hs | 25 +++++--- .../Marconi/Starter/Indexers/AddressCount.hs | 4 +- 15 files changed, 96 insertions(+), 77 deletions(-) diff --git a/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs b/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs index abbb0fe0ea..405ecbbf80 100644 --- a/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs +++ b/doc/read-the-docs-site/doc/marconi-as-a-library/tutorials/BasicApp.hs @@ -53,7 +53,7 @@ import Marconi.Cardano.Core.Types ( ) import Marconi.Cardano.Indexers.SyncHelper qualified as Core import Marconi.Core qualified as Core -import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation, inMemoryDB) +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation, defaultInsertPlan, inMemoryDB) import Marconi.Core.JsonRpc qualified as Core import Network.JsonRpc.Types (JsonRpc, RawJsonRpc) @@ -332,11 +332,13 @@ mkBlockInfoSqliteIndexer dbPath = do -- write the events in the SQLite database [ [ Core.SQLInsertPlan - -- Translate the event into a list of rows. In this specific indexer, each event is a - -- single row. - List.singleton - -- The query that is called for each row that is the output of the previous parameter. - (pure blockInfoInsertQuery) + ( defaultInsertPlan + -- Translate the event into a list of rows. In this specific indexer, each event is a + -- single row. + List.singleton + -- The query that is called for each row that is the output of the previous parameter. + blockInfoInsertQuery + ) ] ] -- Requests launched when a rollback occurs diff --git a/marconi-cardano-indexers/marconi-cardano-indexers.cabal b/marconi-cardano-indexers/marconi-cardano-indexers.cabal index c6df40ed7e..3e85175b9f 100644 --- a/marconi-cardano-indexers/marconi-cardano-indexers.cabal +++ b/marconi-cardano-indexers/marconi-cardano-indexers.cabal @@ -113,7 +113,6 @@ library , sqlite-simple , text , time - , transformers , vector-map test-suite marconi-cardano-indexers-test diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs index b417e2a656..5a54969a53 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs @@ -132,7 +132,7 @@ mkBlockInfoIndexer path = do path id createBlockInfoTable - (pure blockInfoInsertQuery) + blockInfoInsertQuery (Core.SQLRollbackPlan (Core.defaultRollbackPlan "blockInfo" "slotNo" C.chainPointToSlotNo)) catchupConfigEventHook :: Text -> Trace IO Text -> FilePath -> indexer -> IO indexer diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs index d5dea9a4df..44ad13c11f 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs @@ -67,6 +67,7 @@ import Marconi.Cardano.Core.Types ( import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Core (SQLiteDBLocation) import Marconi.Core qualified as Core +import Marconi.Core.Indexer.SQLiteIndexer (defaultInsertPlan) import System.FilePath (()) data DatumInfo = DatumInfo @@ -125,7 +126,7 @@ mkDatumIndexer path = do Sync.mkSyncedSqliteIndexer path createDatumTables - [[Core.SQLInsertPlan (traverse NonEmpty.toList) (pure datumInsertQuery)]] + [[Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) datumInsertQuery)]] [Core.SQLRollbackPlan (Core.defaultRollbackPlan "datum" "slotNo" C.chainPointToSlotNo)] -- | A worker with catchup for a 'DatumIndexer' diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs index 2d2c2fc2ac..39b1678795 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochNonce.hs @@ -52,7 +52,7 @@ import Marconi.Cardano.Indexers.ExtLedgerStateCoordinator ( ) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Core qualified as Core -import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation, defaultInsertPlan) import Ouroboros.Consensus.Cardano.Block qualified as O import Ouroboros.Consensus.HeaderValidation qualified as O import Ouroboros.Consensus.Ledger.Extended qualified as O @@ -107,7 +107,7 @@ mkEpochNonceIndexer path = do , slotNo , blockHeaderHash ) VALUES (?, ?, ?, ?, ?)|] - insertEvent = [Core.SQLInsertPlan pure (pure nonceInsertQuery)] + insertEvent = [Core.SQLInsertPlan (defaultInsertPlan pure nonceInsertQuery)] Sync.mkSyncedSqliteIndexer path [createNonce] diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs index 0fcb5e26fc..2a4445d7bb 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/EpochSDD.hs @@ -62,7 +62,7 @@ import Marconi.Cardano.Indexers.ExtLedgerStateCoordinator ( ) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Core qualified as Core -import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation, defaultInsertPlan) import Ouroboros.Consensus.Cardano.Block qualified as O import Ouroboros.Consensus.Ledger.Extended qualified as O import Ouroboros.Consensus.Shelley.Ledger qualified as O @@ -114,7 +114,7 @@ mkEpochSDDIndexer path = do , slotNo , blockHeaderHash ) VALUES (?, ?, ?, ?, ?, ?)|] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure sddInsertQuery)] + insertEvent = [Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) sddInsertQuery)] Sync.mkSyncedSqliteIndexer path [createSDD] diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs index f5c7173a27..eed4e062e2 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/MintTokenEvent.hs @@ -145,6 +145,7 @@ import Marconi.Cardano.Core.Types ( import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Core (SQLiteDBLocation) import Marconi.Core qualified as Core +import Marconi.Core.Indexer.SQLiteIndexer (defaultInsertPlan) import System.FilePath (()) -- | A raw SQLite indexer for 'MintTokenBlockEvents' @@ -439,7 +440,7 @@ mkMintTokenIndexer dbPath = do ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] createMintPolicyEventTables = [createMintPolicyEvent] - mintInsertPlans = [Core.SQLInsertPlan fromTimedMintEvents (pure mintEventInsertQuery)] + mintInsertPlans = [Core.SQLInsertPlan (defaultInsertPlan fromTimedMintEvents mintEventInsertQuery)] Sync.mkSyncedSqliteIndexer dbPath createMintPolicyEventTables diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index c4f2400f15..9bd76ab636 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -56,6 +56,7 @@ import Marconi.Cardano.Core.Types (AnyTxBody (AnyTxBody), SecurityParam) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Core (SQLiteDBLocation) import Marconi.Core qualified as Core +import Marconi.Core.Indexer.SQLiteIndexer (defaultInsertPlan) import System.FilePath (()) data SpentInfo = SpentInfo @@ -129,7 +130,7 @@ mkSpentIndexer path = do VALUES (?, ?, ?, ?, ?)|] createSpentTables = [createSpent] spentInsert = - [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure spentInsertQuery)] + [Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) spentInsertQuery)] Sync.mkSyncedSqliteIndexer path createSpentTables diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs index 064193fc20..96ab1aca81 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs @@ -91,8 +91,8 @@ mkSingleInsertSyncedSqliteIndexer -> (Core.Timed (Core.Point event) event -> param) -> SQL.Query -- ^ the creation query - -> ([param] -> SQL.Query) - -- ^ the action producing an insert query + -> SQL.Query + -- ^ the insert query -> Core.SQLRollbackPlan (Core.Point event) -- ^ the rollback query -> m (Core.SQLiteIndexer event) @@ -100,7 +100,7 @@ mkSingleInsertSyncedSqliteIndexer path extract tableCreation insertQuery rollbac Core.mkSqliteIndexer path [tableCreation, syncTableCreation] - [[Core.SQLInsertPlan (pure . extract) insertQuery]] + [[Core.SQLInsertPlan (SQL.defaultInsertPlan (pure . extract) insertQuery)]] [rollbackPlan] syncSetStablePoint syncLastPointQuery diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index c602f8299d..5c30c62ab2 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -91,6 +91,7 @@ import Marconi.Cardano.Core.Types ( import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Core (SQLiteDBLocation) import Marconi.Core qualified as Core +import Marconi.Core.Indexer.SQLiteIndexer (defaultInsertPlan) import System.FilePath (()) -- | Indexer representation of an UTxO @@ -160,7 +161,7 @@ mkUtxoIndexer path = do ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] createUtxoTables = [createUtxo] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure utxoInsertQuery)] + insertEvent = [Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) utxoInsertQuery)] Sync.mkSyncedSqliteIndexer path diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 363dacaaf5..7418dbff59 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -14,7 +14,6 @@ import Control.Lens ((^.)) import Control.Lens qualified as Lens import Control.Monad.Error.Class (MonadError) import Control.Monad.Except (MonadIO) -import Control.Monad.Trans.Maybe (MaybeT (MaybeT, runMaybeT)) import Data.Aeson.TH qualified as Aeson import Data.List.NonEmpty (NonEmpty) import Data.List.NonEmpty qualified as NonEmpty @@ -51,7 +50,7 @@ data UtxoWithSpent = UtxoWithSpent -- ^ the inline script hash of the tx out , _txIndex :: TxIndexInBlock -- ^ the index at which the tx is present in the block - , _spentAt :: !(Maybe C.TxIn) + , _spentAt :: !(Maybe C.TxId) -- ^ the tx id and index at which the tx out is spent } deriving (Show, Eq) @@ -65,7 +64,7 @@ instance SQL.ToRow (Core.Timed C.ChainPoint UtxoWithSpent) where let (C.TxIn txid txix) = u ^. Core.event . txIn spentAtFields = case u ^. Core.event . spentAt of - (Just (C.TxIn txIdSpent txIxSpent)) -> [toField txIdSpent, toField txIxSpent] + (Just txIdSpent) -> [toField txIdSpent] Nothing -> [] in toRow [ toField $ u ^. Core.event . address @@ -96,7 +95,7 @@ instance FromRow UtxoWithSpent where _value <- SQL.field _inlineScript <- SQL.field _inlineScriptHash <- SQL.field - _spentAt <- runMaybeT parseSpent + _spentAt <- SQL.field pure $ UtxoWithSpent { _address @@ -108,11 +107,6 @@ instance FromRow UtxoWithSpent where , _inlineScriptHash , _spentAt } - where - parseSpent = MaybeT $ do - txIdSpent <- SQL.field - txIxSpent <- SQL.field - pure . pure $ C.TxIn txIdSpent txIxSpent type UtxoOrSpentEvent = NonEmpty UtxoOrSpent @@ -151,31 +145,34 @@ mkUtxoWithSpentIndexer path = do , slotNo INT NOT NULL , blockHeaderHash BLOB NOT NULL , txIdSpent TEXT - , txIxSpent INT ) as SELECT * FROM utxo LEFT OUTER JOIN spent ON utxo.txId == spent.txId AND utxo.txIx == spent.txIx |] - -- TODO: fix not good need 2 types of insert - utxoInsertQuery :: SQL.Query - utxoInsertQuery = - [sql|INSERT OR REPLACE INTO utxo_with_spent ( - address, - txIndex, - txId, - txIx, - datumHash, - value, - inlineScript, - inlineScriptHash, - slotNo, - blockHeaderHash, - txIdSpent, - txIxSpent - ) VALUES - (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + utxoInsertQuery :: [Core.Timed C.ChainPoint UtxoOrSpent] -> SQL.Query + utxoInsertQuery [] = [sql||] + utxoInsertQuery ((Core.Timed _ utxoOrSpent) : _) = + case utxoOrSpent of + UtxoEvent _ -> + [sql|INSERT INTO utxo_with_spent ( + address, + txIndex, + txId, + txIx, + datumHash, + value, + inlineScript, + inlineScriptHash, + slotNo, + blockHeaderHash, + txIdSpent, + ) VALUES + (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + SpentEvent _ -> + [sql|UPDATE utxo_with_spent + SET txIdSpent = ? + WHERE txId = ? AND txIx = ? |] createUtxoTables = [createUtxoWithSpent] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) (pure utxoInsertQuery)] - + insertEvent = [Core.SQLInsertPlan undefined] -- (traverse NonEmpty.toList)] Sync.mkSyncedSqliteIndexer path createUtxoTables diff --git a/marconi-core/src/Marconi/Core.hs b/marconi-core/src/Marconi/Core.hs index d0b6baa9e6..512ba79282 100644 --- a/marconi-core/src/Marconi/Core.hs +++ b/marconi-core/src/Marconi/Core.hs @@ -282,7 +282,7 @@ module Marconi.Core ( inMemoryDB, parseDBLocation, connection, - SQLInsertPlan (SQLInsertPlan, planExtractor, planInsert), + SQLInsertPlan (SQLInsertPlan, planInsert), SQLRollbackPlan (SQLRollbackPlan, planRollback), defaultRollbackPlan, diff --git a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs index 9c66bc9b95..ad28635b10 100644 --- a/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs +++ b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs @@ -1,4 +1,5 @@ {-# LANGUAGE QuasiQuotes #-} +{-# LANGUAGE RankNTypes #-} {-# LANGUAGE StrictData #-} {-# LANGUAGE TemplateHaskell #-} {-# LANGUAGE UndecidableInstances #-} @@ -32,7 +33,8 @@ module Marconi.Core.Indexer.SQLiteIndexer ( querySyncedOnlySQLiteIndexerWith, handleSQLErrors, dbLastSync, - SQLInsertPlan (SQLInsertPlan, planInsert, planExtractor), + SQLInsertPlan (SQLInsertPlan, planInsert), + defaultInsertPlan, SQLRollbackPlan (SQLRollbackPlan, planRollback), defaultRollbackPlan, @@ -89,13 +91,8 @@ data ExpectedPersistentDB = ExpectedPersistentDB instance Exception ExpectedPersistentDB -- | A 'SQLInsertPlan' provides a piece information about how an event should be inserted in the database -data SQLInsertPlan event = forall a. - (SQL.ToRow a) => - SQLInsertPlan - { planExtractor :: Timed (Point event) event -> [a] - -- ^ How to transform the event into a type that can be handle by the database - , planInsert :: [a] -> SQL.Query - -- ^ The insert statement builder for the extracted data +newtype SQLInsertPlan event = SQLInsertPlan + { planInsert :: [Timed (Point event) event] -> SQL.Connection -> IO () } newtype InsertPointQuery = InsertPointQuery {getInsertPointQuery :: SQL.Query} @@ -225,7 +222,11 @@ mkSingleInsertSqliteIndexer -- ^ The SQL query to fetch the last stable point from the indexer. -> m (SQLiteIndexer event) mkSingleInsertSqliteIndexer path extract create insert rollback' = - mkSqliteIndexer path [create] [[SQLInsertPlan (pure . extract) (pure insert)]] [rollback'] + mkSqliteIndexer + path + [create] + [[SQLInsertPlan (defaultInsertPlan (pure . extract) insert)]] + [rollback'] -- | Map SQLite errors to an indexer error handleSQLErrors :: IO a -> IO (Either IndexerError a) @@ -236,6 +237,21 @@ handleSQLErrors value = , Handler (\(x :: SQL.SQLError) -> pure . Left . IndexerInternalError . Text.pack $ show x) ] +defaultInsertPlan + :: forall a q + . (SQL.ToRow q) + => (a -> [q]) + -> SQL.Query + -> [a] + -> SQL.Connection + -> IO () +defaultInsertPlan planExtractor query events c = do + let rows = planExtractor =<< events + case rows of + [] -> pure () + [x] -> SQL.execute c query x + _nonEmpty -> SQL.executeMany c query rows + -- | Run a list of insert queries in one single transaction. runIndexQueriesStep :: SQL.Connection @@ -244,13 +260,7 @@ runIndexQueriesStep -> IO () runIndexQueriesStep _ _ [] = pure () runIndexQueriesStep c events plan = - let runIndexQuery (SQLInsertPlan planExtractor planInsert) = do - let rows = planExtractor =<< events - query = planInsert rows - case rows of - [] -> pure () - [x] -> SQL.execute c query x - _nonEmpty -> SQL.executeMany c query rows + let runIndexQuery (SQLInsertPlan planInsert) = planInsert events c in Async.mapConcurrently_ runIndexQuery plan -- | Run a list of insert queries in one single transaction. diff --git a/marconi-core/test/Marconi/CoreSpec.hs b/marconi-core/test/Marconi/CoreSpec.hs index 5e9f3fc9e2..380f143104 100644 --- a/marconi-core/test/Marconi/CoreSpec.hs +++ b/marconi-core/test/Marconi/CoreSpec.hs @@ -182,7 +182,12 @@ import GHC.Generics (Generic) import Marconi.Core (Streamable (streamFrom), unwrap, withStream) import Marconi.Core qualified as Core import Marconi.Core.Coordinator qualified as Core (errorBox, threadIds) -import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation, inMemoryDB, parseDBLocation) +import Marconi.Core.Indexer.SQLiteIndexer ( + SQLiteDBLocation, + defaultInsertPlan, + inMemoryDB, + parseDBLocation, + ) import Network.Socket ( Family (AF_UNIX), SockAddr (SockAddrUnix), @@ -668,15 +673,17 @@ sqliteModelIndexerWithFile filepath = do ] [ [ Core.SQLInsertPlan - ( \t -> - [ - ( t ^. Core.point . testPointSlot - , UUID.toText $ t ^. Core.point . testPointHash - , t ^. Core.event - ) - ] + ( defaultInsertPlan + ( \t -> + [ + ( t ^. Core.point . testPointSlot + , UUID.toText $ t ^. Core.point . testPointHash + , t ^. Core.event + ) + ] + ) + "INSERT INTO index_model VALUES (?, ?, ?)" ) - (pure "INSERT INTO index_model VALUES (?, ?, ?)") ] ] [Core.SQLRollbackPlan (Core.defaultRollbackPlan "index_model" "pointSlotNo" extractor)] diff --git a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs index 46c7245772..8763d60bf4 100644 --- a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs +++ b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs @@ -57,7 +57,7 @@ import Marconi.Core ( point, ) import Marconi.Core qualified as Core -import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation, defaultInsertPlan) type SQLiteStandardIndexer event = Core.StandardIndexer IO Core.SQLiteIndexer event type AddressCountIndexer = SQLiteStandardIndexer AddressCountEvent @@ -196,7 +196,7 @@ mkAddressCountSqliteIndexer dbPath = do dbPath [dbCreation] -- request launched when the indexer is created [ - [ SQLInsertPlan eventToRows (pure addressCountInsertQuery) + [ SQLInsertPlan (defaultInsertPlan eventToRows addressCountInsertQuery) ] ] -- requests launched when an event is stored [SQLRollbackPlan (Core.defaultRollbackPlan "address_count" "slotNo" C.chainPointToSlotNo)] From 1dbd54038b10b020ff21cba7208000466676daeb Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 2 Feb 2024 13:18:06 +0200 Subject: [PATCH 09/15] Implement custom insert plan Signed-off-by: Ana Pantilie --- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 68 ++++++++++++------- 1 file changed, 42 insertions(+), 26 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 7418dbff59..51a842004b 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -1,5 +1,6 @@ {-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE NamedFieldPuns #-} +{-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE QuasiQuotes #-} {-# LANGUAGE TemplateHaskell #-} {-# OPTIONS_GHC -Wno-orphans #-} @@ -15,16 +16,17 @@ import Control.Lens qualified as Lens import Control.Monad.Error.Class (MonadError) import Control.Monad.Except (MonadIO) import Data.Aeson.TH qualified as Aeson +import Data.Foldable (traverse_) import Data.List.NonEmpty (NonEmpty) import Data.List.NonEmpty qualified as NonEmpty -import Database.SQLite.Simple (FromRow (fromRow), ToRow (toRow)) +import Database.SQLite.Simple (FromRow (fromRow), NamedParam ((:=)), ToRow (toRow)) import Database.SQLite.Simple qualified as SQL import Database.SQLite.Simple.QQ (sql) import Database.SQLite.Simple.ToField (ToField (toField)) import Marconi.Cardano.Core.Indexer.Worker (StandardSQLiteIndexer) import Marconi.Cardano.Core.Orphans () import Marconi.Cardano.Core.Types (TxIndexInBlock) -import Marconi.Cardano.Indexers.Spent (SpentInfo) +import Marconi.Cardano.Indexers.Spent (SpentInfo (SpentInfo)) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync import Marconi.Cardano.Indexers.Utxo (Utxo) import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) @@ -148,34 +150,48 @@ mkUtxoWithSpentIndexer path = do ) as SELECT * FROM utxo LEFT OUTER JOIN spent ON utxo.txId == spent.txId AND utxo.txIx == spent.txIx |] - utxoInsertQuery :: [Core.Timed C.ChainPoint UtxoOrSpent] -> SQL.Query - utxoInsertQuery [] = [sql||] - utxoInsertQuery ((Core.Timed _ utxoOrSpent) : _) = - case utxoOrSpent of - UtxoEvent _ -> - [sql|INSERT INTO utxo_with_spent ( - address, - txIndex, - txId, - txIx, - datumHash, - value, - inlineScript, - inlineScriptHash, - slotNo, - blockHeaderHash, - txIdSpent, - ) VALUES - (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] - SpentEvent _ -> - [sql|UPDATE utxo_with_spent - SET txIdSpent = ? - WHERE txId = ? AND txIx = ? |] createUtxoTables = [createUtxoWithSpent] - insertEvent = [Core.SQLInsertPlan undefined] -- (traverse NonEmpty.toList)] + insertEvent = [Core.SQLInsertPlan insertPlan] Sync.mkSyncedSqliteIndexer path createUtxoTables [insertEvent] -- TODO: fix [Core.SQLRollbackPlan (Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo)] + +insertPlan :: [Core.Timed C.ChainPoint UtxoOrSpentEvent] -> SQL.Connection -> IO () +insertPlan events conn = do + let rows = traverse NonEmpty.toList =<< events + traverse_ buildAndExecuteQuery rows + where + buildAndExecuteQuery :: Core.Timed C.ChainPoint UtxoOrSpent -> IO () + buildAndExecuteQuery row@(Core.Timed _ utxoOrSpent) = + case utxoOrSpent of + UtxoEvent _ -> do + let query = + [sql|INSERT INTO utxo_with_spent ( + address, + txIndex, + txId, + txIx, + datumHash, + value, + inlineScript, + inlineScriptHash, + slotNo, + blockHeaderHash, + txIdSpent, + ) VALUES + (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + SQL.execute conn query row + SpentEvent (SpentInfo (C.TxIn txId txIx) txIdSpent) -> do + let query = + [sql|UPDATE utxo_with_spent + SET txIdSpent = :txIdSpent + WHERE txId = :txId AND txIx = :txIx|] + params = + [ ":txId" := txId + , ":txIx" := txIx + , ":txIdSpent" := txIdSpent + ] + SQL.executeNamed conn query params From 9721fd0097eff6e4a9bc0a7ed6025ae29d9321f5 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 2 Feb 2024 13:50:48 +0200 Subject: [PATCH 10/15] Fix table creation query Signed-off-by: Ana Pantilie --- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 31 ++++++++++++++++--- 1 file changed, 27 insertions(+), 4 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 51a842004b..ace846f03c 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -54,6 +54,7 @@ data UtxoWithSpent = UtxoWithSpent -- ^ the index at which the tx is present in the block , _spentAt :: !(Maybe C.TxId) -- ^ the tx id and index at which the tx out is spent + , _spentAtSlotNo :: !(Maybe C.SlotNo) } deriving (Show, Eq) @@ -98,6 +99,7 @@ instance FromRow UtxoWithSpent where _inlineScript <- SQL.field _inlineScriptHash <- SQL.field _spentAt <- SQL.field + _spentAtSlotNo <- SQL.field pure $ UtxoWithSpent { _address @@ -108,6 +110,7 @@ instance FromRow UtxoWithSpent where , _inlineScript , _inlineScriptHash , _spentAt + , _spentAtSlotNo } type UtxoOrSpentEvent = NonEmpty UtxoOrSpent @@ -147,8 +150,25 @@ mkUtxoWithSpentIndexer path = do , slotNo INT NOT NULL , blockHeaderHash BLOB NOT NULL , txIdSpent TEXT + , slotNoSpent INT ) - as SELECT * FROM utxo LEFT OUTER JOIN spent ON utxo.txId == spent.txId AND utxo.txIx == spent.txIx + as + SELECT + utxo.address, + utxo.txIndex, + utxo.txId, + utxo.txIx, + utxo.datumHash, + utxo.value, + utxo.inlineScript, + utxo.inlineScriptHash, + utxo.slotNo, + utxo.blockHeaderHash, + spent.spentAtTxId, + spent.slotNo + FROM utxo + LEFT OUTER JOIN spent ON utxo.txId == spent.txId AND utxo.txIx == spent.txIx + WHERE ISNULL spent.spentAtTxId AND ISNULL spent.slotNo |] createUtxoTables = [createUtxoWithSpent] insertEvent = [Core.SQLInsertPlan insertPlan] @@ -165,7 +185,7 @@ insertPlan events conn = do traverse_ buildAndExecuteQuery rows where buildAndExecuteQuery :: Core.Timed C.ChainPoint UtxoOrSpent -> IO () - buildAndExecuteQuery row@(Core.Timed _ utxoOrSpent) = + buildAndExecuteQuery row@(Core.Timed chainPoint utxoOrSpent) = case utxoOrSpent of UtxoEvent _ -> do let query = @@ -185,13 +205,16 @@ insertPlan events conn = do (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] SQL.execute conn query row SpentEvent (SpentInfo (C.TxIn txId txIx) txIdSpent) -> do - let query = + let pointToSlot (C.ChainPoint slotNoSpent _) = slotNoSpent + pointToSlot _ = 0 + query = [sql|UPDATE utxo_with_spent - SET txIdSpent = :txIdSpent + SET txIdSpent = :txIdSpent, slotNoSpent = :slotNoSpent WHERE txId = :txId AND txIx = :txIx|] params = [ ":txId" := txId , ":txIx" := txIx , ":txIdSpent" := txIdSpent + , ":slotNoSpent" := pointToSlot chainPoint ] SQL.executeNamed conn query params From a0c48a57629bed998e39449ffe4dbed1321d8fa3 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 2 Feb 2024 14:21:47 +0200 Subject: [PATCH 11/15] Add custom rollback plan + fixes Signed-off-by: Ana Pantilie --- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 50 +++++++++++++------ 1 file changed, 36 insertions(+), 14 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index ace846f03c..64384f6a0f 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -6,7 +6,7 @@ {-# OPTIONS_GHC -Wno-orphans #-} module Marconi.Cardano.Indexers.UtxoWithSpent ( - + mkUtxoWithSpentIndexer, ) where import Cardano.Api qualified as C @@ -52,7 +52,7 @@ data UtxoWithSpent = UtxoWithSpent -- ^ the inline script hash of the tx out , _txIndex :: TxIndexInBlock -- ^ the index at which the tx is present in the block - , _spentAt :: !(Maybe C.TxId) + , _spentAtTx :: !(Maybe C.TxId) -- ^ the tx id and index at which the tx out is spent , _spentAtSlotNo :: !(Maybe C.SlotNo) } @@ -65,10 +65,14 @@ Lens.makeLenses ''UtxoWithSpent instance SQL.ToRow (Core.Timed C.ChainPoint UtxoWithSpent) where toRow u = let (C.TxIn txid txix) = u ^. Core.event . txIn - spentAtFields = - case u ^. Core.event . spentAt of + spentAtField = + case u ^. Core.event . spentAtTx of (Just txIdSpent) -> [toField txIdSpent] Nothing -> [] + spentAtSlotNoField = + case u ^. Core.event . spentAtSlotNo of + (Just slotNoSpent) -> [toField slotNoSpent] + Nothing -> [] in toRow [ toField $ u ^. Core.event . address , toField $ u ^. Core.event . txIndex @@ -79,14 +83,15 @@ instance SQL.ToRow (Core.Timed C.ChainPoint UtxoWithSpent) where , toField $ u ^. Core.event . inlineScript , toField $ u ^. Core.event . inlineScriptHash ] - <> toRow spentAtFields + <> toRow spentAtField + <> toRow spentAtSlotNoField <> toRow (u ^. Core.point) instance FromRow (Core.Timed C.ChainPoint UtxoWithSpent) where fromRow = do - utxo <- fromRow + utxoWithSpent <- fromRow point <- fromRow - pure $ Core.Timed point utxo + pure $ Core.Timed point utxoWithSpent instance FromRow UtxoWithSpent where fromRow = do @@ -98,7 +103,7 @@ instance FromRow UtxoWithSpent where _value <- SQL.field _inlineScript <- SQL.field _inlineScriptHash <- SQL.field - _spentAt <- SQL.field + _spentAtTx <- SQL.field _spentAtSlotNo <- SQL.field pure $ UtxoWithSpent @@ -109,7 +114,7 @@ instance FromRow UtxoWithSpent where , _value , _inlineScript , _inlineScriptHash - , _spentAt + , _spentAtTx , _spentAtSlotNo } @@ -171,13 +176,11 @@ mkUtxoWithSpentIndexer path = do WHERE ISNULL spent.spentAtTxId AND ISNULL spent.slotNo |] createUtxoTables = [createUtxoWithSpent] - insertEvent = [Core.SQLInsertPlan insertPlan] Sync.mkSyncedSqliteIndexer path createUtxoTables - [insertEvent] - -- TODO: fix - [Core.SQLRollbackPlan (Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo)] + [[Core.SQLInsertPlan insertPlan]] + [Core.SQLRollbackPlan rollbackPlan] insertPlan :: [Core.Timed C.ChainPoint UtxoOrSpentEvent] -> SQL.Connection -> IO () insertPlan events conn = do @@ -201,8 +204,9 @@ insertPlan events conn = do slotNo, blockHeaderHash, txIdSpent, + slotNoSpent, ) VALUES - (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] SQL.execute conn query row SpentEvent (SpentInfo (C.TxIn txId txIx) txIdSpent) -> do let pointToSlot (C.ChainPoint slotNoSpent _) = slotNoSpent @@ -218,3 +222,21 @@ insertPlan events conn = do , ":slotNoSpent" := pointToSlot chainPoint ] SQL.executeNamed conn query params + +rollbackPlan + :: C.ChainPoint + -> SQL.Connection + -> IO () +rollbackPlan point conn = + case C.chainPointToSlotNo point of + Nothing -> + SQL.execute_ conn [sql|DELETE FROM utxo_with_spent|] + Just slotNo -> do + let deleteNewRows = + [sql|DELETE FROM utxo_with_spent WHERE slotNo == :slotNo|] + deleteNewSpentReferences = + [sql|UPDATE utxo_with_spent + SET txIdSpent = NULL, slotNoSpent = NULL + WHERE slotNoSpent = :slotNo|] + SQL.executeNamed conn deleteNewRows [":slotNo" := slotNo] + SQL.executeNamed conn deleteNewSpentReferences [":slotNo" := slotNo] From 1b798fe06bf3645919868fbb9f74c08dec00a1b4 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 2 Feb 2024 14:25:25 +0200 Subject: [PATCH 12/15] Fix toRow Signed-off-by: Ana Pantilie --- .../src/Marconi/Cardano/Indexers/UtxoWithSpent.hs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 64384f6a0f..cbff1bf0c7 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -6,6 +6,7 @@ {-# OPTIONS_GHC -Wno-orphans #-} module Marconi.Cardano.Indexers.UtxoWithSpent ( + StandardUtxoWithSpentIndexer, mkUtxoWithSpentIndexer, ) where @@ -50,7 +51,7 @@ data UtxoWithSpent = UtxoWithSpent -- ^ the inline script of the tx out , _inlineScriptHash :: !(Maybe C.ScriptHash) -- ^ the inline script hash of the tx out - , _txIndex :: TxIndexInBlock + , _txIndex :: !TxIndexInBlock -- ^ the index at which the tx is present in the block , _spentAtTx :: !(Maybe C.TxId) -- ^ the tx id and index at which the tx out is spent @@ -83,9 +84,9 @@ instance SQL.ToRow (Core.Timed C.ChainPoint UtxoWithSpent) where , toField $ u ^. Core.event . inlineScript , toField $ u ^. Core.event . inlineScriptHash ] + <> toRow (u ^. Core.point) <> toRow spentAtField <> toRow spentAtSlotNoField - <> toRow (u ^. Core.point) instance FromRow (Core.Timed C.ChainPoint UtxoWithSpent) where fromRow = do From cef5439faa803692edd44a648f6ef3a775e9f402 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 2 Feb 2024 15:32:58 +0200 Subject: [PATCH 13/15] Add mkIndexer for two table version Signed-off-by: Ana Pantilie --- .../src/Marconi/Cardano/Indexers/Spent.hs | 64 +++++++++----- .../src/Marconi/Cardano/Indexers/Utxo.hs | 84 ++++++++++++------- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 42 ++++++++-- 3 files changed, 131 insertions(+), 59 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index 9bd76ab636..f763fd9032 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -21,6 +21,14 @@ module Marconi.Cardano.Indexers.Spent ( StandardSpentIndexer, catchupConfigEventHook, + -- * Queries + createSpent, + spentInsertQuery, + + -- * SQL plans + spentInsertPlan, + spentRollbackPlan, + -- * Extractor getInputs, ) where @@ -105,37 +113,49 @@ type SpentIndexer = Core.SQLiteIndexer SpentInfoEvent -- | A SQLite Spent indexer with Catchup type StandardSpentIndexer m = StandardSQLiteIndexer m SpentInfoEvent +createSpent, spentInsertQuery :: SQL.Query +createSpent = + [sql|CREATE TABLE IF NOT EXISTS spent + ( txId TEXT NOT NULL + , txIx INT NOT NULL + , spentAtTxId TEXT NOT NULL + , slotNo INT NOT NULL + , blockHeaderHash BLOB NOT NULL + )|] +spentInsertQuery = + [sql|INSERT OR IGNORE INTO spent + ( txId + , txIx + , spentAtTxId + , slotNo + , blockHeaderHash + ) + VALUES (?, ?, ?, ?, ?)|] + +spentInsertPlan + :: (Core.ToRow (Core.Timed (Core.Point (NonEmpty a)) a)) + => Core.SQLInsertPlan (NonEmpty a) +spentInsertPlan = + Core.SQLInsertPlan $ + defaultInsertPlan (traverse NonEmpty.toList) spentInsertQuery + +spentRollbackPlan :: Core.SQLRollbackPlan C.ChainPoint +spentRollbackPlan = + Core.SQLRollbackPlan $ + Core.defaultRollbackPlan "spent" "slotNo" C.chainPointToSlotNo + mkSpentIndexer :: (MonadIO m, MonadError Core.IndexerError m) => SQLiteDBLocation -> m (Core.SQLiteIndexer SpentInfoEvent) mkSpentIndexer path = do - let createSpent = - [sql|CREATE TABLE IF NOT EXISTS spent - ( txId TEXT NOT NULL - , txIx INT NOT NULL - , spentAtTxId TEXT NOT NULL - , slotNo INT NOT NULL - , blockHeaderHash BLOB NOT NULL - )|] - spentInsertQuery :: SQL.Query - spentInsertQuery = - [sql|INSERT OR IGNORE INTO spent - ( txId - , txIx - , spentAtTxId - , slotNo - , blockHeaderHash - ) - VALUES (?, ?, ?, ?, ?)|] - createSpentTables = [createSpent] - spentInsert = - [Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) spentInsertQuery)] + let createSpentTables = [createSpent] + spentInsert = [spentInsertPlan] Sync.mkSyncedSqliteIndexer path createSpentTables [spentInsert] - [Core.SQLRollbackPlan (Core.defaultRollbackPlan "spent" "slotNo" C.chainPointToSlotNo)] + [spentRollbackPlan] catchupConfigEventHook :: Trace IO Text -> FilePath -> indexer -> IO indexer catchupConfigEventHook stdoutTrace dbPath indexer = do diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index 5c30c62ab2..d1e307373d 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -41,6 +41,14 @@ module Marconi.Cardano.Indexers.Utxo ( trackedAddresses, includeScript, + -- * Queries + createUtxo, + utxoInsertQuery, + + -- * SQL plans + utxoInsertPlan, + utxoRollbackPlan, + -- * Extractors getUtxoEventsFromBlock, getUtxosFromTx, @@ -125,6 +133,47 @@ type instance Core.Point UtxoEvent = C.ChainPoint type UtxoIndexer = Core.SQLiteIndexer UtxoEvent type StandardUtxoIndexer m = StandardSQLiteIndexer m UtxoEvent +createUtxo, utxoInsertQuery :: SQL.Query +createUtxo = + [sql|CREATE TABLE IF NOT EXISTS utxo + ( address BLOB NOT NULL + , txIndex INT NOT NULL + , txId TEXT NOT NULL + , txIx INT NOT NULL + , datumHash BLOB + , value BLOB + , inlineScript BLOB + , inlineScriptHash BLOB + , slotNo INT NOT NULL + , blockHeaderHash BLOB NOT NULL + )|] +utxoInsertQuery = + [sql|INSERT INTO utxo ( + address, + txIndex, + txId, + txIx, + datumHash, + value, + inlineScript, + inlineScriptHash, + slotNo, + blockHeaderHash + ) VALUES + (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + +utxoInsertPlan + :: (Core.ToRow (Core.Timed (Core.Point (NonEmpty a)) a)) + => Core.SQLInsertPlan (NonEmpty a) +utxoInsertPlan = + Core.SQLInsertPlan $ + defaultInsertPlan (traverse NonEmpty.toList) utxoInsertQuery + +utxoRollbackPlan :: Core.SQLRollbackPlan C.ChainPoint +utxoRollbackPlan = + Core.SQLRollbackPlan $ + Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo + -- | Make a SQLiteIndexer for Utxos mkUtxoIndexer :: (MonadIO m, MonadError Core.IndexerError m) @@ -132,42 +181,13 @@ mkUtxoIndexer -- ^ SQL connection to database -> m UtxoIndexer mkUtxoIndexer path = do - let createUtxo = - [sql|CREATE TABLE IF NOT EXISTS utxo - ( address BLOB NOT NULL - , txIndex INT NOT NULL - , txId TEXT NOT NULL - , txIx INT NOT NULL - , datumHash BLOB - , value BLOB - , inlineScript BLOB - , inlineScriptHash BLOB - , slotNo INT NOT NULL - , blockHeaderHash BLOB NOT NULL - )|] - utxoInsertQuery :: SQL.Query - utxoInsertQuery = - [sql|INSERT INTO utxo ( - address, - txIndex, - txId, - txIx, - datumHash, - value, - inlineScript, - inlineScriptHash, - slotNo, - blockHeaderHash - ) VALUES - (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] - createUtxoTables = [createUtxo] - insertEvent = [Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) utxoInsertQuery)] - + let createUtxoTables = [createUtxo] + insertEvent = [utxoInsertPlan] Sync.mkSyncedSqliteIndexer path createUtxoTables [insertEvent] - [Core.SQLRollbackPlan (Core.defaultRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo)] + [utxoRollbackPlan] catchupConfigEventHook :: Trace IO Text -> FilePath -> indexer -> IO indexer catchupConfigEventHook stdoutTrace dbPath indexer = do diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index cbff1bf0c7..9ce34c4c5e 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -7,6 +7,7 @@ module Marconi.Cardano.Indexers.UtxoWithSpent ( StandardUtxoWithSpentIndexer, + mkUtxoOrSpentIndexer, mkUtxoWithSpentIndexer, ) where @@ -27,9 +28,14 @@ import Database.SQLite.Simple.ToField (ToField (toField)) import Marconi.Cardano.Core.Indexer.Worker (StandardSQLiteIndexer) import Marconi.Cardano.Core.Orphans () import Marconi.Cardano.Core.Types (TxIndexInBlock) -import Marconi.Cardano.Indexers.Spent (SpentInfo (SpentInfo)) +import Marconi.Cardano.Indexers.Spent ( + SpentInfo (SpentInfo), + createSpent, + spentInsertPlan, + spentRollbackPlan, + ) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync -import Marconi.Cardano.Indexers.Utxo (Utxo) +import Marconi.Cardano.Indexers.Utxo (Utxo, createUtxo, utxoInsertPlan, utxoRollbackPlan) import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) import Marconi.Core.Indexer.SQLiteIndexer qualified as Core import Marconi.Core.Type qualified as Core @@ -119,8 +125,9 @@ instance FromRow UtxoWithSpent where , _spentAtSlotNo } -type UtxoOrSpentEvent = NonEmpty UtxoOrSpent - +{- | We need to consider the following cases: either the event to be indexed represents + the creation of a new UTxO or the spending of an existing one. +-} data UtxoOrSpent = UtxoEvent Utxo | SpentEvent SpentInfo @@ -132,11 +139,34 @@ instance SQL.ToRow (Core.Timed C.ChainPoint UtxoOrSpent) where instance FromRow (Core.Timed C.ChainPoint UtxoOrSpent) where fromRow = fmap UtxoEvent <$> fromRow <|> fmap SpentEvent <$> fromRow +-- | The event type for indexing UTxOs with Spent information. +type UtxoOrSpentEvent = NonEmpty UtxoOrSpent + type instance Core.Point UtxoOrSpentEvent = C.ChainPoint type UtxoWithSpentIndexer = Core.SQLiteIndexer UtxoOrSpentEvent type StandardUtxoWithSpentIndexer m = StandardSQLiteIndexer m UtxoOrSpentEvent --- | Make a SQLiteIndexer for UtxoWithSpent +{- | Make a SQLiteIndexer which indexes data into two tables: one for UTxO + information and the other for Spent information. +-} +mkUtxoOrSpentIndexer + :: (MonadIO m, MonadError Core.IndexerError m) + => SQLiteDBLocation + -- ^ SQL connection to database + -> m UtxoWithSpentIndexer +mkUtxoOrSpentIndexer path = do + let createTables = [createUtxo, createSpent] + insertPlans = [[utxoInsertPlan, spentInsertPlan]] + rollbackPlans = [utxoRollbackPlan, spentRollbackPlan] + Sync.mkSyncedSqliteIndexer + path + createTables + insertPlans + rollbackPlans + +{- | Make a SQLiteIndexer which indexes data in a table containing both + UTxO and Spent information. +-} mkUtxoWithSpentIndexer :: (MonadIO m, MonadError Core.IndexerError m) => SQLiteDBLocation @@ -183,6 +213,7 @@ mkUtxoWithSpentIndexer path = do [[Core.SQLInsertPlan insertPlan]] [Core.SQLRollbackPlan rollbackPlan] +-- | Custom SQL insert plan for adding both UTxO and Spent data to the table. insertPlan :: [Core.Timed C.ChainPoint UtxoOrSpentEvent] -> SQL.Connection -> IO () insertPlan events conn = do let rows = traverse NonEmpty.toList =<< events @@ -224,6 +255,7 @@ insertPlan events conn = do ] SQL.executeNamed conn query params +-- | Custom SQL rollback plan for removing both UTxO and Spent data from the table. rollbackPlan :: C.ChainPoint -> SQL.Connection From 52a9507919b8b237181fd681de9a693f26ff3f73 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Fri, 2 Feb 2024 16:45:08 +0200 Subject: [PATCH 14/15] WIP: add utxoWithSpent worker Signed-off-by: Ana Pantilie --- .../src/Marconi/Cardano/Indexers/Spent.hs | 5 +- .../src/Marconi/Cardano/Indexers/Utxo.hs | 8 +- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 79 +++++++++++++++++-- 3 files changed, 82 insertions(+), 10 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs index f763fd9032..b1fa52cd24 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Spent.hs @@ -183,7 +183,7 @@ catchupConfigEventHook stdoutTrace dbPath indexer = do -- return the original indexer pure indexer --- | A minimal worker for the UTXO indexer, with catchup and filtering. +-- | A minimal worker for the spent indexer spentWorker :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) => StandardWorkerConfig m Core.SQLiteIndexer input SpentInfoEvent @@ -196,7 +196,8 @@ spentWorker config path = do mkStandardWorker config indexer {- | Convenience wrapper around 'spentWorker' with some defaults for -creating 'StandardWorkerConfig', including a preprocessor. +creating 'StandardWorkerConfig', including a preprocessor. Adds catchup +capabilities as well. -} spentBuilder :: (MonadIO n, MonadError Core.IndexerError n) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs index d1e307373d..53c6a421a4 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Utxo.hs @@ -53,6 +53,7 @@ module Marconi.Cardano.Indexers.Utxo ( getUtxoEventsFromBlock, getUtxosFromTx, getUtxosFromTxBody, + extractUtxos, ) where import Cardano.Api qualified as C @@ -260,8 +261,6 @@ utxoBuilder securityParam catchupConfig utxoConfig textLogger path = let indexerName = "Utxo" indexerEventLogger = BM.contramap (fmap (fmap $ Text.pack . show)) textLogger utxoDbPath = path "utxo.db" - extractUtxos :: AnyTxBody -> [Utxo] - extractUtxos (AnyTxBody _ indexInBlock txb) = getUtxosFromTxBody indexInBlock txb catchupConfigWithTracer = catchupConfig & Core.configCatchupEventHook .~ catchupConfigEventHook textLogger utxoDbPath @@ -364,9 +363,12 @@ instance (const utxoQuery) (\(Core.EventsMatchingQuery p) -> parseResult p) +extractUtxos :: AnyTxBody -> [Utxo] +extractUtxos (AnyTxBody _ indexInBlock txb) = getUtxosFromTxBody indexInBlock txb + {- | Extract UtxoEvents from Cardano Block - Returns @Nothing@ if the block doesn't consume or spend any utxo + Returns the empty list if the block doesn't consume or spend any utxo -} getUtxoEventsFromBlock :: (C.IsCardanoEra era) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 9ce34c4c5e..9d7bf5576c 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -12,8 +12,10 @@ module Marconi.Cardano.Indexers.UtxoWithSpent ( ) where import Cardano.Api qualified as C +import Cardano.BM.Trace (Trace) +import Cardano.BM.Tracing qualified as BM import Control.Applicative ((<|>)) -import Control.Lens ((^.)) +import Control.Lens ((&), (.~), (^.)) import Control.Lens qualified as Lens import Control.Monad.Error.Class (MonadError) import Control.Monad.Except (MonadIO) @@ -21,24 +23,41 @@ import Data.Aeson.TH qualified as Aeson import Data.Foldable (traverse_) import Data.List.NonEmpty (NonEmpty) import Data.List.NonEmpty qualified as NonEmpty +import Data.String (fromString) +import Data.Text (Text) +import Data.Text qualified as Text import Database.SQLite.Simple (FromRow (fromRow), NamedParam ((:=)), ToRow (toRow)) import Database.SQLite.Simple qualified as SQL import Database.SQLite.Simple.QQ (sql) import Database.SQLite.Simple.ToField (ToField (toField)) -import Marconi.Cardano.Core.Indexer.Worker (StandardSQLiteIndexer) +import Marconi.Cardano.Core.Indexer.Worker ( + StandardSQLiteIndexer, + StandardWorker, + StandardWorkerConfig (StandardWorkerConfig), + mkStandardWorker, + ) import Marconi.Cardano.Core.Orphans () -import Marconi.Cardano.Core.Types (TxIndexInBlock) +import Marconi.Cardano.Core.Types (AnyTxBody (AnyTxBody), SecurityParam, TxIndexInBlock) import Marconi.Cardano.Indexers.Spent ( SpentInfo (SpentInfo), createSpent, + getInputs, spentInsertPlan, spentRollbackPlan, ) import Marconi.Cardano.Indexers.SyncHelper qualified as Sync -import Marconi.Cardano.Indexers.Utxo (Utxo, createUtxo, utxoInsertPlan, utxoRollbackPlan) -import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation) +import Marconi.Cardano.Indexers.Utxo ( + Utxo, + createUtxo, + extractUtxos, + utxoInsertPlan, + utxoRollbackPlan, + ) +import Marconi.Core qualified as Core +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation (Storage)) import Marconi.Core.Indexer.SQLiteIndexer qualified as Core import Marconi.Core.Type qualified as Core +import System.FilePath (()) {- | Indexer representation of a UTxO together with the information about where it was spent. If this information is missing, then it isn't @@ -273,3 +292,53 @@ rollbackPlan point conn = WHERE slotNoSpent = :slotNo|] SQL.executeNamed conn deleteNewRows [":slotNo" := slotNo] SQL.executeNamed conn deleteNewSpentReferences [":slotNo" := slotNo] + +catchupConfigEventHook + :: Trace IO Text + -> FilePath + -> UtxoWithSpentIndexer + -> IO UtxoWithSpentIndexer +catchupConfigEventHook _ dbPath _ = mkUtxoWithSpentIndexer (Storage dbPath) + +-- | A minimal worker for the UtxoWithSpent indexer +utxoWithSpentWorker + :: (MonadIO n, MonadError Core.IndexerError n, MonadIO m) + => StandardWorkerConfig m Core.SQLiteIndexer input UtxoOrSpentEvent + -- ^ General configuration of a worker + -> SQLiteDBLocation + -- ^ SQLite database location + -> n (StandardWorker m input UtxoOrSpentEvent Core.SQLiteIndexer) +utxoWithSpentWorker config path = do + indexer <- mkUtxoOrSpentIndexer path + mkStandardWorker config indexer + +{- | Convenience wrapper around 'utxoWithSpentWorker' with some defaults for +creating 'StandardWorkerConfig', including a preprocessor. Adds catchup +capabilities as well. +-} +utxoWithSpentBuilder + :: (MonadIO n, MonadError Core.IndexerError n) + => SecurityParam + -> Core.CatchupConfig indexer event + -> Trace IO Text + -> FilePath + -> n (StandardWorker IO [AnyTxBody] UtxoOrSpentEvent Core.SQLiteIndexer) +utxoWithSpentBuilder securityParam catchupConfig textLogger path = + let indexerName = "UtxoWithSpent" + indexerEventLogger = BM.contramap (fmap (fmap $ Text.pack . show)) textLogger + spentDbPath = path "utxoWithSpent.db" + extractUtxoOrSpent :: AnyTxBody -> [UtxoOrSpent] + extractUtxoOrSpent anyTxBody@(AnyTxBody _ _ txb) = + (UtxoEvent <$> extractUtxos anyTxBody) <> (SpentEvent <$> getInputs txb) + catchupConfigWithTracer = + catchupConfig + & Core.configCatchupEventHook + .~ catchupConfigEventHook textLogger spentDbPath + utxoWithSpentWorkerConfig = + StandardWorkerConfig + indexerName + securityParam + undefined -- catchupConfigWithTracer + (pure . NonEmpty.nonEmpty . (>>= extractUtxoOrSpent)) + (BM.appendName indexerName indexerEventLogger) + in utxoWithSpentWorker utxoWithSpentWorkerConfig (Core.parseDBLocation spentDbPath) From b574b7ed678b5e2a82dfcf2d22ef32f2b6024a39 Mon Sep 17 00:00:00 2001 From: Ana Pantilie Date: Sat, 3 Feb 2024 13:20:14 +0200 Subject: [PATCH 15/15] Finish implementing indexing for UtxoWithSpent Signed-off-by: Ana Pantilie --- .../Marconi/Cardano/Indexers/UtxoWithSpent.hs | 26 +++++++++++++------ 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs index 9d7bf5576c..fecf49d0ce 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -9,27 +9,29 @@ module Marconi.Cardano.Indexers.UtxoWithSpent ( StandardUtxoWithSpentIndexer, mkUtxoOrSpentIndexer, mkUtxoWithSpentIndexer, + utxoWithSpentBuilder, ) where import Cardano.Api qualified as C import Cardano.BM.Trace (Trace) import Cardano.BM.Tracing qualified as BM import Control.Applicative ((<|>)) +import Control.Exception (throwIO) import Control.Lens ((&), (.~), (^.)) import Control.Lens qualified as Lens import Control.Monad.Error.Class (MonadError) -import Control.Monad.Except (MonadIO) +import Control.Monad.Except (MonadIO, runExceptT) import Data.Aeson.TH qualified as Aeson import Data.Foldable (traverse_) import Data.List.NonEmpty (NonEmpty) import Data.List.NonEmpty qualified as NonEmpty -import Data.String (fromString) import Data.Text (Text) import Data.Text qualified as Text import Database.SQLite.Simple (FromRow (fromRow), NamedParam ((:=)), ToRow (toRow)) import Database.SQLite.Simple qualified as SQL import Database.SQLite.Simple.QQ (sql) import Database.SQLite.Simple.ToField (ToField (toField)) +import Marconi.Cardano.Core.Extract.WithDistance (WithDistance) import Marconi.Cardano.Core.Indexer.Worker ( StandardSQLiteIndexer, StandardWorker, @@ -55,8 +57,6 @@ import Marconi.Cardano.Indexers.Utxo ( ) import Marconi.Core qualified as Core import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation (Storage)) -import Marconi.Core.Indexer.SQLiteIndexer qualified as Core -import Marconi.Core.Type qualified as Core import System.FilePath (()) {- | Indexer representation of a UTxO together with the information @@ -293,12 +293,22 @@ rollbackPlan point conn = SQL.executeNamed conn deleteNewRows [":slotNo" := slotNo] SQL.executeNamed conn deleteNewSpentReferences [":slotNo" := slotNo] +type TopUtxoWithSpentIndexer = + Core.WithTransform + Core.SQLiteIndexer + (NonEmpty UtxoOrSpent) + (WithDistance (Maybe (NonEmpty UtxoOrSpent))) + catchupConfigEventHook :: Trace IO Text -> FilePath - -> UtxoWithSpentIndexer - -> IO UtxoWithSpentIndexer -catchupConfigEventHook _ dbPath _ = mkUtxoWithSpentIndexer (Storage dbPath) + -> TopUtxoWithSpentIndexer + -> IO TopUtxoWithSpentIndexer +catchupConfigEventHook _ dbPath topIndexer = do + result <- runExceptT $ mkUtxoWithSpentIndexer (Storage dbPath) + case result of + Left err -> throwIO err + Right indexer -> pure (Lens.set Core.unwrapMap indexer topIndexer) -- | A minimal worker for the UtxoWithSpent indexer utxoWithSpentWorker @@ -338,7 +348,7 @@ utxoWithSpentBuilder securityParam catchupConfig textLogger path = StandardWorkerConfig indexerName securityParam - undefined -- catchupConfigWithTracer + catchupConfigWithTracer (pure . NonEmpty.nonEmpty . (>>= extractUtxoOrSpent)) (BM.appendName indexerName indexerEventLogger) in utxoWithSpentWorker utxoWithSpentWorkerConfig (Core.parseDBLocation spentDbPath)