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..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,22 +332,26 @@ 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. - 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 [ 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-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/marconi-cardano-indexers.cabal b/marconi-cardano-indexers/marconi-cardano-indexers.cabal index 5a5226b3a2..3e85175b9f 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 diff --git a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/BlockInfo.hs index fbd261daed..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)) @@ -133,10 +133,10 @@ 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 +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 @@ -151,7 +152,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 +167,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) @@ -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/Datum.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/Datum.hs index 6b8a9b84fc..44ad13c11f 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, @@ -66,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 @@ -124,13 +126,13 @@ mkDatumIndexer path = do Sync.mkSyncedSqliteIndexer path createDatumTables - [[Core.SQLInsertPlan (traverse NonEmpty.toList) datumInsertQuery]] - [Core.SQLRollbackPlan "datum" "slotNo" C.chainPointToSlotNo] + [[Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) datumInsertQuery)]] + [Core.SQLRollbackPlan (Core.defaultRollbackPlan "datum" "slotNo" C.chainPointToSlotNo)] -- | 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 +148,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 153cbc84ed..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,21 +107,21 @@ mkEpochNonceIndexer path = do , slotNo , blockHeaderHash ) VALUES (?, ?, ?, ?, ?)|] - insertEvent = [Core.SQLInsertPlan pure nonceInsertQuery] + insertEvent = [Core.SQLInsertPlan (defaultInsertPlan pure nonceInsertQuery)] Sync.mkSyncedSqliteIndexer 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 } 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 2245dab475..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,21 +114,21 @@ mkEpochSDDIndexer path = do , slotNo , blockHeaderHash ) VALUES (?, ?, ?, ?, ?, ?)|] - insertEvent = [Core.SQLInsertPlan (traverse NonEmpty.toList) sddInsertQuery] + insertEvent = [Core.SQLInsertPlan (defaultInsertPlan (traverse NonEmpty.toList) sddInsertQuery)] Sync.mkSyncedSqliteIndexer 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 } 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 509c60ceb6..eed4e062e2 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, (%~), (.~), - (?~), (^.), (^?), ) @@ -146,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' @@ -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 @@ -327,7 +327,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 = @@ -440,15 +440,17 @@ mkMintTokenIndexer dbPath = do ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] createMintPolicyEventTables = [createMintPolicyEvent] - mintInsertPlans = [Core.SQLInsertPlan fromTimedMintEvents mintEventInsertQuery] + mintInsertPlans = [Core.SQLInsertPlan (defaultInsertPlan fromTimedMintEvents mintEventInsertQuery)] Sync.mkSyncedSqliteIndexer 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 +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 = @@ -468,6 +470,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/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 e8aeeac5e4..b1fa52cd24 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 @@ -28,7 +36,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) @@ -56,6 +64,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 @@ -104,40 +113,52 @@ 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 (traverse NonEmpty.toList) spentInsertQuery] + let createSpentTables = [createSpent] + spentInsert = [spentInsertPlan] Sync.mkSyncedSqliteIndexer path createSpentTables [spentInsert] - [Core.SQLRollbackPlan "spent" "slotNo" C.chainPointToSlotNo] + [spentRollbackPlan] -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,11 +180,13 @@ 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. +-- | A minimal worker for the spent indexer 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 @@ -173,12 +196,13 @@ 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) => SecurityParam - -> Core.CatchupConfig + -> Core.CatchupConfig indexer event -> BM.Trace IO Text -> FilePath -> n (StandardWorker IO [AnyTxBody] SpentInfoEvent Core.SQLiteIndexer) @@ -191,7 +215,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/SyncHelper.hs b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs index 9ed5d6deb5..96ab1aca81 100644 --- a/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/SyncHelper.hs @@ -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 bb37ccc780..53c6a421a4 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 #-} @@ -42,10 +41,19 @@ module Marconi.Cardano.Indexers.Utxo ( trackedAddresses, includeScript, + -- * Queries + createUtxo, + utxoInsertQuery, + + -- * SQL plans + utxoInsertPlan, + utxoRollbackPlan, + -- * Extractors getUtxoEventsFromBlock, getUtxosFromTx, getUtxosFromTxBody, + extractUtxos, ) where import Cardano.Api qualified as C @@ -55,7 +63,6 @@ import Cardano.BM.Tracing qualified as BM import Control.Lens ( (&), (.~), - (?~), (^.), ) import Control.Lens qualified as Lens @@ -76,7 +83,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, @@ -94,6 +100,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 @@ -106,7 +113,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 @@ -127,6 +134,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) @@ -134,45 +182,16 @@ 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 (traverse NonEmpty.toList) utxoInsertQuery] - + let createUtxoTables = [createUtxo] + insertEvent = [utxoInsertPlan] Sync.mkSyncedSqliteIndexer path createUtxoTables [insertEvent] - [Core.SQLRollbackPlan "utxo" "slotNo" C.chainPointToSlotNo] + [utxoRollbackPlan] -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 = @@ -187,11 +206,12 @@ 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 :: (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) @@ -232,7 +252,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 @@ -241,11 +261,9 @@ 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 + & Core.configCatchupEventHook .~ catchupConfigEventHook textLogger utxoDbPath utxoWorkerConfig = StandardWorkerConfig indexerName @@ -345,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 new file mode 100644 index 0000000000..fecf49d0ce --- /dev/null +++ b/marconi-cardano-indexers/src/Marconi/Cardano/Indexers/UtxoWithSpent.hs @@ -0,0 +1,354 @@ +{-# LANGUAGE FlexibleContexts #-} +{-# LANGUAGE NamedFieldPuns #-} +{-# LANGUAGE OverloadedStrings #-} +{-# LANGUAGE QuasiQuotes #-} +{-# LANGUAGE TemplateHaskell #-} +{-# OPTIONS_GHC -Wno-orphans #-} + +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, 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.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, + StandardWorkerConfig (StandardWorkerConfig), + mkStandardWorker, + ) +import Marconi.Cardano.Core.Orphans () +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, + extractUtxos, + utxoInsertPlan, + utxoRollbackPlan, + ) +import Marconi.Core qualified as Core +import Marconi.Core.Indexer.SQLiteIndexer (SQLiteDBLocation (Storage)) +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 + 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 + , _spentAtTx :: !(Maybe C.TxId) + -- ^ the tx id and index at which the tx out is spent + , _spentAtSlotNo :: !(Maybe C.SlotNo) + } + 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 + 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 + , 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 (u ^. Core.point) + <> toRow spentAtField + <> toRow spentAtSlotNoField + +instance FromRow (Core.Timed C.ChainPoint UtxoWithSpent) where + fromRow = do + utxoWithSpent <- fromRow + point <- fromRow + pure $ Core.Timed point utxoWithSpent + +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 + _spentAtTx <- SQL.field + _spentAtSlotNo <- SQL.field + pure $ + UtxoWithSpent + { _address + , _txIndex + , _txIn = C.TxIn txId txIx + , _datumHash + , _value + , _inlineScript + , _inlineScriptHash + , _spentAtTx + , _spentAtSlotNo + } + +{- | 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 + +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 + +-- | 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 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 + -- ^ 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 + , slotNoSpent INT + ) + 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] + Sync.mkSyncedSqliteIndexer + path + createUtxoTables + [[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 + traverse_ buildAndExecuteQuery rows + where + buildAndExecuteQuery :: Core.Timed C.ChainPoint UtxoOrSpent -> IO () + buildAndExecuteQuery row@(Core.Timed chainPoint 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, + slotNoSpent, + ) VALUES + (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)|] + SQL.execute conn query row + SpentEvent (SpentInfo (C.TxIn txId txIx) txIdSpent) -> do + let pointToSlot (C.ChainPoint slotNoSpent _) = slotNoSpent + pointToSlot _ = 0 + query = + [sql|UPDATE utxo_with_spent + 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 + +-- | Custom SQL rollback plan for removing both UTxO and Spent data from the table. +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] + +type TopUtxoWithSpentIndexer = + Core.WithTransform + Core.SQLiteIndexer + (NonEmpty UtxoOrSpent) + (WithDistance (Maybe (NonEmpty UtxoOrSpent))) + +catchupConfigEventHook + :: Trace IO Text + -> FilePath + -> 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 + :: (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 + catchupConfigWithTracer + (pure . NonEmpty.nonEmpty . (>>= extractUtxoOrSpent)) + (BM.appendName indexerName indexerEventLogger) + in utxoWithSpentWorker utxoWithSpentWorkerConfig (Core.parseDBLocation spentDbPath) 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.hs b/marconi-core/src/Marconi/Core.hs index 31e941bf53..512ba79282 100644 --- a/marconi-core/src/Marconi/Core.hs +++ b/marconi-core/src/Marconi/Core.hs @@ -282,8 +282,9 @@ module Marconi.Core ( inMemoryDB, parseDBLocation, connection, - SQLInsertPlan (SQLInsertPlan, planExtractor, planInsert), - SQLRollbackPlan (SQLRollbackPlan, tableName, pointName, pointExtractor), + SQLInsertPlan (SQLInsertPlan, planInsert), + SQLRollbackPlan (SQLRollbackPlan, planRollback), + defaultRollbackPlan, -- **** Reexport from SQLite ToRow (..), @@ -455,7 +456,7 @@ module Marconi.Core ( CatchupConfig (CatchupConfig), mkCatchupConfig, configCatchupEventHook, - HasCatchupConfig (catchupBypassDistance, catchupBatchSize, catchupEventHook), + HasCatchupConfig (catchupBypassDistance, catchupBatchSize), CatchupEvent (Synced), -- *** SQLite @@ -620,6 +621,7 @@ import Marconi.Core.Indexer.SQLiteIndexer ( ToRow (..), connection, dbLastSync, + defaultRollbackPlan, handleSQLErrors, inMemoryDB, mkSingleInsertSqliteIndexer, @@ -692,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/Indexer/SQLiteIndexer.hs b/marconi-core/src/Marconi/Core/Indexer/SQLiteIndexer.hs index 4d47a461d1..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,8 +33,10 @@ module Marconi.Core.Indexer.SQLiteIndexer ( querySyncedOnlySQLiteIndexerWith, handleSQLErrors, dbLastSync, - SQLInsertPlan (SQLInsertPlan, planInsert, planExtractor), - SQLRollbackPlan (SQLRollbackPlan, tableName, pointName, pointExtractor), + SQLInsertPlan (SQLInsertPlan, planInsert), + defaultInsertPlan, + SQLRollbackPlan (SQLRollbackPlan, planRollback), + defaultRollbackPlan, -- * Reexport from SQLite SQL.ToRow (..), @@ -88,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 :: SQL.Query - -- ^ The insert statement for the extracted data +newtype SQLInsertPlan event = SQLInsertPlan + { planInsert :: [Timed (Point event) event] -> SQL.Connection -> IO () } newtype InsertPointQuery = InsertPointQuery {getInsertPointQuery :: SQL.Query} @@ -102,17 +100,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. @@ -233,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) 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) @@ -244,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 @@ -252,12 +260,7 @@ runIndexQueriesStep -> IO () runIndexQueriesStep _ _ [] = pure () runIndexQueriesStep c events plan = - let runIndexQuery (SQLInsertPlan planExtractor planInsert) = do - let rows = planExtractor =<< events - case rows of - [] -> pure () - [x] -> SQL.execute c planInsert x - _nonEmpty -> SQL.executeMany c planInsert rows + let runIndexQuery (SQLInsertPlan planInsert) = planInsert events c in Async.mapConcurrently_ runIndexQuery plan -- | Run a list of insert queries in one single transaction. @@ -298,6 +301,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 +343,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/src/Marconi/Core/Transformer/WithCatchup.hs b/marconi-core/src/Marconi/Core/Transformer/WithCatchup.hs index 0107083e76..dde7d47710 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,26 +54,26 @@ 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 -- ^ 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 -mkCatchupConfig batchSize bypassDistance = CatchupConfig batchSize bypassDistance Nothing +mkCatchupConfig :: Word64 -> Word64 -> CatchupConfig indexer event +mkCatchupConfig batchSize bypassDistance = CatchupConfig batchSize bypassDistance pure 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,25 @@ 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) (indexer event -> IO (indexer event)) 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 +111,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 @@ -115,23 +119,29 @@ 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) + (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) @@ -147,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 #-} @@ -161,7 +169,6 @@ instance where catchupBypassDistance = unwrap . catchupBypassDistance catchupBatchSize = unwrap . catchupBatchSize - catchupEventHook = unwrap . catchupEventHook instance {-# OVERLAPPABLE #-} @@ -170,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 @@ -208,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 diff --git a/marconi-core/test/Marconi/CoreSpec.hs b/marconi-core/test/Marconi/CoreSpec.hs index 710e225aa0..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,18 +673,20 @@ 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 (?, ?, ?)" ) - "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-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 $ diff --git a/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs b/marconi-starter/src/Marconi/Starter/Indexers/AddressCount.hs index 5dd8bd368d..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,10 +196,10 @@ mkAddressCountSqliteIndexer dbPath = do dbPath [dbCreation] -- request launched when the indexer is created [ - [ SQLInsertPlan eventToRows addressCountInsertQuery + [ SQLInsertPlan (defaultInsertPlan 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