From 72ded9dbfb98f8b18ad77c4bcfac239bca3892b7 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Wed, 19 Mar 2025 00:38:34 +0100 Subject: [PATCH 01/13] close db connection --- cmd/cli/cli.go | 21 ++++++++++++++++++--- 1 file changed, 18 insertions(+), 3 deletions(-) diff --git a/cmd/cli/cli.go b/cmd/cli/cli.go index 0001bc6..1d83289 100644 --- a/cmd/cli/cli.go +++ b/cmd/cli/cli.go @@ -136,17 +136,32 @@ func Start() error { Port: 4000, }) services = append(services, metricsServer) + db, err := database.NewDB(ctx) if err != nil { - return err + return fmt.Errorf("failed to create database connection: %w", err) } + // Ensure db is always closed + defer func() { + if db != nil { + log.Info().Msg("Closing database connection") + db.Close() + } + }() err = database.PerformMigration(ctx) if err != nil { - return err + return fmt.Errorf("failed to perform database migration: %w", err) } + watcher := watcher.New(&config, db) scheduler := scheduler.New(&config, db) services = append(services, watcher, scheduler) - return service.RunWithSighandler(ctx, services...) + + // RunWithSighandler will handle SIGINT/SIGTERM + if err := service.RunWithSighandler(ctx, services...); err != nil { + return fmt.Errorf("service error: %w", err) + } + + return nil } From f0e5708c904f00f3b9da088f66201eeb8d59c860 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Wed, 19 Mar 2025 00:39:17 +0100 Subject: [PATCH 02/13] improve go routines --- internal/metrics/tx_mapper_db.go | 193 ++++++++++++++++--------------- 1 file changed, 101 insertions(+), 92 deletions(-) diff --git a/internal/metrics/tx_mapper_db.go b/internal/metrics/tx_mapper_db.go index 863e86f..8d08c7c 100644 --- a/internal/metrics/tx_mapper_db.go +++ b/internal/metrics/tx_mapper_db.go @@ -395,6 +395,7 @@ func (tm *TxMapperDB) processTransactionExecution( slot := te.DecKeysAndMessages[0].Slot + var wg sync.WaitGroup for index, txSubEvent := range txSubEvents { decryptionKeyID, err := getDecryptionKeyID(txSubEvent, identityPreimageToDecKeyAndMsg) if err != nil { @@ -426,71 +427,77 @@ func (tm *TxMapperDB) processTransactionExecution( Uint8("tx-type", decryptedTx.Type()). Msg("tx-data") - // channel to signal the other routine to stop waiting for receipt - txErrorSignalCh := make(chan bool) - - go func(ctx context.Context, inclusionDelay int64, decryptedTx *types.Transaction, txSubEvent data.TransactionSubmittedEvent, slot int64, decryptionKeyID int64, txErrorSignalCh chan bool) { - // send tx to public mempool since keys are already public with a delay - time.Sleep(time.Duration(inclusionDelay) * time.Second) - defer close(txErrorSignalCh) - - err = tm.ethClient.SendTransaction(ctx, decryptedTx) - if err != nil { - log.Err(err).Msg("failed to send transaction") - if err.Error() == "AlreadyKnown" { - log.Debug().Hex("tx-hash", decryptedTx.Hash().Bytes()).Msg("already known") - err := tm.dbQuery.CreateDecryptedTX(ctx, data.CreateDecryptedTXParams{ - Slot: slot, - TxIndex: txSubEvent.TxIndex, - TxHash: decryptedTx.Hash().Bytes(), - TxStatus: data.TxStatusValPending, - DecryptionKeyID: decryptionKeyID, - TransactionSubmittedEventID: txSubEvent.ID, - }) - if err != nil { - log.Err(err).Msg("failed to create decrypted tx") - txErrorSignalCh <- true + // channel to propagate errors between goroutines + txErrorSignalCh := make(chan error, 1) + + wg.Add(2) + + // First goroutine: Send transaction + go func(ctx context.Context, inclusionDelay int64, decryptedTx *types.Transaction, txSubEvent data.TransactionSubmittedEvent, slot int64, decryptionKeyID int64, txErrorSignalCh chan error) { + defer wg.Done() + + select { + case <-ctx.Done(): + txErrorSignalCh <- fmt.Errorf("transaction send cancelled due to context: %w", ctx.Err()) + return + case <-time.After(time.Duration(tm.config.InclusionDelay) * time.Second): + if err := tm.ethClient.SendTransaction(ctx, decryptedTx); err != nil { + log.Err(err).Msg("failed to send transaction") + if err.Error() == "AlreadyKnown" { + log.Debug().Hex("tx-hash", decryptedTx.Hash().Bytes()).Msg("already known") + err := tm.dbQuery.CreateDecryptedTX(ctx, data.CreateDecryptedTXParams{ + Slot: slot, + TxIndex: txSubEvent.TxIndex, + TxHash: decryptedTx.Hash().Bytes(), + TxStatus: data.TxStatusValPending, + DecryptionKeyID: decryptionKeyID, + TransactionSubmittedEventID: txSubEvent.ID, + }) + if err != nil { + txErrorSignalCh <- fmt.Errorf("failed to create decrypted tx: %w", err) + return + } + } else { + err := tm.dbQuery.CreateDecryptedTX(ctx, data.CreateDecryptedTXParams{ + Slot: slot, + TxIndex: txSubEvent.TxIndex, + TxHash: decryptedTx.Hash().Bytes(), + TxStatus: data.TxStatusValInvalid, + DecryptionKeyID: decryptionKeyID, + TransactionSubmittedEventID: txSubEvent.ID, + }) + if err != nil { + log.Err(err).Msg("failed to create decrypted tx") + } + txErrorSignalCh <- fmt.Errorf("failed to send transaction: %w", err) return } } else { + log.Info().Hex("tx-hash", decryptedTx.Hash().Bytes()).Msg("transaction sent") err := tm.dbQuery.CreateDecryptedTX(ctx, data.CreateDecryptedTXParams{ Slot: slot, TxIndex: txSubEvent.TxIndex, TxHash: decryptedTx.Hash().Bytes(), - TxStatus: data.TxStatusValInvalid, + TxStatus: data.TxStatusValPending, DecryptionKeyID: decryptionKeyID, TransactionSubmittedEventID: txSubEvent.ID, }) if err != nil { - log.Err(err).Msg("failed to create decrypted tx") + txErrorSignalCh <- fmt.Errorf("failed to create decrypted tx: %w", err) + return } - txErrorSignalCh <- true - return - } - } else { - log.Info().Hex("tx-hash", decryptedTx.Hash().Bytes()).Msg("transaction sent") - err := tm.dbQuery.CreateDecryptedTX(ctx, data.CreateDecryptedTXParams{ - Slot: slot, - TxIndex: txSubEvent.TxIndex, - TxHash: decryptedTx.Hash().Bytes(), - TxStatus: data.TxStatusValPending, - DecryptionKeyID: decryptionKeyID, - TransactionSubmittedEventID: txSubEvent.ID, - }) - if err != nil { - log.Err(err).Msg("failed to create decrypted tx") - txErrorSignalCh <- true - return } } }(ctx, tm.config.InclusionDelay, decryptedTx, txSubEvent, slot, decryptionKeyID, txErrorSignalCh) - // Fire off a goroutine to wait for the transaction receipt - go func(ctx context.Context, index int, txHash common.Hash, txIndex int64, slot int64, decryptionKeyID int64, txSubEventID int64, txErrorSignalCh chan bool) { + // Second goroutine: Wait for receipt + go func(ctx context.Context, index int, txHash common.Hash, txIndex int64, slot int64, decryptionKeyID int64, txSubEventID int64, txErrorSignalCh chan error) { + defer wg.Done() + // Wait for the receipt with a timeout receipt, err := tm.waitForReceiptWithTimeout(ctx, txHash, ReceiptWaitTimeout, txErrorSignalCh) if err != nil { - log.Err(err).Msgf("failed to get receipt for transaction %s", txHash.Hex()) + log.Err(err).Msg("") // update/create status to not included err := tm.dbQuery.UpsertTX(ctx, data.UpsertTXParams{ Slot: slot, @@ -501,57 +508,59 @@ func (tm *TxMapperDB) processTransactionExecution( TransactionSubmittedEventID: txSubEventID, }) if err != nil { - log.Err(err).Msg("failed to update decrypted tx") - return + log.Err(err).Msg("failed to upsert decrypted tx") } - } else { - // receipt found - log.Info().Hex("tx-hash", receipt.TxHash.Bytes()). - Uint64("receipt-status", receipt.Status). - Msg("transaction receipt found") + return + } - block, err := tm.ethClient.BlockByNumber(ctx, receipt.BlockNumber) - if err != nil { - log.Err(err).Uint64("block-number", receipt.BlockNumber.Uint64()).Msg("failed to retrieve block") - return - } + // receipt found + log.Info().Hex("tx-hash", receipt.TxHash.Bytes()). + Uint64("receipt-status", receipt.Status). + Msg("transaction receipt found") - inclusionSlot := utils.GetSlotForBlock(block.Header().Time, tm.genesisTimestamp, tm.slotDuration) - txStatus := data.TxStatusValShieldedinclusion + block, err := tm.ethClient.BlockByNumber(ctx, receipt.BlockNumber) + if err != nil { + log.Err(err).Uint64("block-number", receipt.BlockNumber.Uint64()).Msg("failed to retrieve block") + return + } - log.Info().Uint("tx-index", receipt.TransactionIndex). - Uint64("inclusion-slot", inclusionSlot). - Msg("receipt data") + inclusionSlot := utils.GetSlotForBlock(block.Header().Time, tm.genesisTimestamp, tm.slotDuration) + txStatus := data.TxStatusValShieldedinclusion - log.Info().Int("index", index). - Int64("inclusion-slot", slot). - Msg("local data") + log.Info().Uint("tx-index", receipt.TransactionIndex). + Uint64("inclusion-slot", inclusionSlot). + Msg("receipt data") - if receipt.TransactionIndex != uint(index) { - log.Info().Uint("tx-index", receipt.TransactionIndex).Msg("transaction index mismatch") - txStatus = data.TxStatusValUnshieldedinclusion - } - if inclusionSlot != uint64(slot) { - log.Info().Int64("slot", slot).Msg("transaction slot mismatch") - txStatus = data.TxStatusValUnshieldedinclusion - } + log.Info().Int("index", index). + Int64("inclusion-slot", slot). + Msg("local data") - err = tm.dbQuery.UpsertTX(ctx, data.UpsertTXParams{ - Slot: slot, - TxIndex: txIndex, - TxHash: receipt.TxHash.Bytes(), - TxStatus: txStatus, - DecryptionKeyID: decryptionKeyID, - TransactionSubmittedEventID: txSubEventID, - BlockNumber: pgtype.Int8{Int64: receipt.BlockNumber.Int64(), Valid: true}, - }) - if err != nil { - log.Err(err).Msg("failed to update decrypted tx") - return - } + if receipt.TransactionIndex != uint(index) { + log.Info().Uint("tx-index", receipt.TransactionIndex).Msg("transaction index mismatch") + txStatus = data.TxStatusValUnshieldedinclusion + } + if inclusionSlot != uint64(slot) { + log.Info().Int64("slot", slot).Msg("transaction slot mismatch") + txStatus = data.TxStatusValUnshieldedinclusion + } + + err = tm.dbQuery.UpsertTX(ctx, data.UpsertTXParams{ + Slot: slot, + TxIndex: txIndex, + TxHash: receipt.TxHash.Bytes(), + TxStatus: txStatus, + DecryptionKeyID: decryptionKeyID, + TransactionSubmittedEventID: txSubEventID, + BlockNumber: pgtype.Int8{Int64: receipt.BlockNumber.Int64(), Valid: true}, + }) + if err != nil { + log.Err(err).Msg("failed to update decrypted tx") } }(ctx, index, decryptedTx.Hash(), txSubEvent.TxIndex, slot, decryptionKeyID, txSubEvent.ID, txErrorSignalCh) } + + // Wait for all routines to end + wg.Wait() return nil } @@ -740,7 +749,7 @@ func decryptTransaction(key []byte, encrypted []byte) (*types.Transaction, error } // waitForReceiptWithTimeout waits for a transaction receipt with a provided timeout. -func (tm *TxMapperDB) waitForReceiptWithTimeout(ctx context.Context, txHash common.Hash, receiptWaitTimeout time.Duration, txErrorSignalCh chan bool) (*types.Receipt, error) { +func (tm *TxMapperDB) waitForReceiptWithTimeout(ctx context.Context, txHash common.Hash, receiptWaitTimeout time.Duration, txErrorSignalCh chan error) (*types.Receipt, error) { ctx, cancel := context.WithTimeout(ctx, receiptWaitTimeout) defer cancel() @@ -753,15 +762,15 @@ func (tm *TxMapperDB) waitForReceiptWithTimeout(ctx context.Context, txHash comm } // waitForReceipt polls the Ethereum network for the transaction receipt until it's available or the context is canceled. -func (tm *TxMapperDB) waitForReceipt(ctx context.Context, txHash common.Hash, txErrorSignalCh chan bool) (*types.Receipt, error) { +func (tm *TxMapperDB) waitForReceipt(ctx context.Context, txHash common.Hash, txErrorSignalCh chan error) (*types.Receipt, error) { for { // check if the context has been canceled or timed out select { case <-ctx.Done(): return nil, ctx.Err() - case errSignal, ok := <-txErrorSignalCh: // Listen for a signal from the txErrorSignalCh - if ok && errSignal { - return nil, fmt.Errorf("error encountered during transaction execution %s", txHash.Hex()) + case err := <-txErrorSignalCh: // Listen for errors from the sending goroutine + if err != nil { + return nil, err } default: } From 1682ca62dde2786b33c1d746ee313bd40024c024 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Thu, 27 Mar 2025 16:58:00 +0100 Subject: [PATCH 03/13] add tx submitted event syncer --- internal/data/metrics.sql.go | 60 +++++-- internal/data/models.sqlc.gen.go | 7 + internal/data/sql/queries/metrics.sql | 16 +- internal/metrics/tx_mapper_db.go | 9 +- internal/syncer/encrypted_tx.go | 146 ++++++++++++++++++ ...0327143729_add_tx_sub_event_sync_table.sql | 16 ++ 6 files changed, 238 insertions(+), 16 deletions(-) create mode 100644 internal/syncer/encrypted_tx.go create mode 100644 migrations/20250327143729_add_tx_sub_event_sync_table.sql diff --git a/internal/data/metrics.sql.go b/internal/data/metrics.sql.go index 9cdbd33..646591c 100644 --- a/internal/data/metrics.sql.go +++ b/internal/data/metrics.sql.go @@ -268,14 +268,35 @@ func (q *Queries) CreateTransactionSubmittedEvent(ctx context.Context, arg Creat return err } +const createTransactionSubmittedEventsSyncedUntil = `-- name: CreateTransactionSubmittedEventsSyncedUntil :exec +INSERT INTO transaction_submitted_events_synced_until (block_hash, block_number) VALUES ($1, $2) +ON CONFLICT (enforce_one_row) DO UPDATE +SET block_hash = $1, block_number = $2 +` + +type CreateTransactionSubmittedEventsSyncedUntilParams struct { + BlockHash []byte + BlockNumber int64 +} + +func (q *Queries) CreateTransactionSubmittedEventsSyncedUntil(ctx context.Context, arg CreateTransactionSubmittedEventsSyncedUntilParams) error { + _, err := q.db.Exec(ctx, createTransactionSubmittedEventsSyncedUntil, arg.BlockHash, arg.BlockNumber) + return err +} + const createValidatorRegistryEventsSyncedUntil = `-- name: CreateValidatorRegistryEventsSyncedUntil :exec -INSERT INTO validator_registry_events_synced_until (block_number) VALUES ($1) +INSERT INTO validator_registry_events_synced_until (block_hash, block_number) VALUES ($1, $2) ON CONFLICT (enforce_one_row) DO UPDATE -SET block_number = $1 +SET block_hash = $1, block_number = $2 ` -func (q *Queries) CreateValidatorRegistryEventsSyncedUntil(ctx context.Context, blockNumber int64) error { - _, err := q.db.Exec(ctx, createValidatorRegistryEventsSyncedUntil, blockNumber) +type CreateValidatorRegistryEventsSyncedUntilParams struct { + BlockHash []byte + BlockNumber int64 +} + +func (q *Queries) CreateValidatorRegistryEventsSyncedUntil(ctx context.Context, arg CreateValidatorRegistryEventsSyncedUntilParams) error { + _, err := q.db.Exec(ctx, createValidatorRegistryEventsSyncedUntil, arg.BlockHash, arg.BlockNumber) return err } @@ -500,6 +521,22 @@ func (q *Queries) QueryTransactionSubmittedEvent(ctx context.Context, arg QueryT return items, nil } +const queryTransactionSubmittedEventsSyncedUntil = `-- name: QueryTransactionSubmittedEventsSyncedUntil :one +SELECT block_hash, block_number FROM transaction_submitted_events_synced_until LIMIT 1 +` + +type QueryTransactionSubmittedEventsSyncedUntilRow struct { + BlockHash []byte + BlockNumber int64 +} + +func (q *Queries) QueryTransactionSubmittedEventsSyncedUntil(ctx context.Context) (QueryTransactionSubmittedEventsSyncedUntilRow, error) { + row := q.db.QueryRow(ctx, queryTransactionSubmittedEventsSyncedUntil) + var i QueryTransactionSubmittedEventsSyncedUntilRow + err := row.Scan(&i.BlockHash, &i.BlockNumber) + return i, err +} + const queryValidatorRegistrationMessageNonceBefore = `-- name: QueryValidatorRegistrationMessageNonceBefore :one SELECT nonce FROM validator_registration_message WHERE validator_index = $1 AND event_block_number <= $2 AND event_tx_index <= $3 AND event_log_index <= $4 ORDER BY event_block_number DESC, event_tx_index DESC, event_log_index DESC FOR UPDATE ` @@ -524,14 +561,19 @@ func (q *Queries) QueryValidatorRegistrationMessageNonceBefore(ctx context.Conte } const queryValidatorRegistryEventsSyncedUntil = `-- name: QueryValidatorRegistryEventsSyncedUntil :one -SELECT block_number FROM validator_registry_events_synced_until LIMIT 1 +SELECT block_hash, block_number FROM validator_registry_events_synced_until LIMIT 1 ` -func (q *Queries) QueryValidatorRegistryEventsSyncedUntil(ctx context.Context) (int64, error) { +type QueryValidatorRegistryEventsSyncedUntilRow struct { + BlockHash []byte + BlockNumber int64 +} + +func (q *Queries) QueryValidatorRegistryEventsSyncedUntil(ctx context.Context) (QueryValidatorRegistryEventsSyncedUntilRow, error) { row := q.db.QueryRow(ctx, queryValidatorRegistryEventsSyncedUntil) - var block_number int64 - err := row.Scan(&block_number) - return block_number, err + var i QueryValidatorRegistryEventsSyncedUntilRow + err := row.Scan(&i.BlockHash, &i.BlockNumber) + return i, err } const queryValidatorStatuses = `-- name: QueryValidatorStatuses :many diff --git a/internal/data/models.sqlc.gen.go b/internal/data/models.sqlc.gen.go index d2da014..1361757 100644 --- a/internal/data/models.sqlc.gen.go +++ b/internal/data/models.sqlc.gen.go @@ -184,6 +184,12 @@ type TransactionSubmittedEvent struct { EventTxHash []byte } +type TransactionSubmittedEventsSyncedUntil struct { + EnforceOneRow bool + BlockHash []byte + BlockNumber int64 +} + type ValidatorRegistrationMessage struct { ID int64 Version pgtype.Int8 @@ -204,6 +210,7 @@ type ValidatorRegistrationMessage struct { type ValidatorRegistryEventsSyncedUntil struct { EnforceOneRow bool BlockNumber int64 + BlockHash []byte } type ValidatorStatus struct { diff --git a/internal/data/sql/queries/metrics.sql b/internal/data/sql/queries/metrics.sql index ae14293..d4b09f9 100644 --- a/internal/data/sql/queries/metrics.sql +++ b/internal/data/sql/queries/metrics.sql @@ -125,15 +125,15 @@ SELECT * FROM transaction_submitted_event WHERE eon = $1 AND tx_index >= $2 AND tx_index < $2 + $3 ORDER BY tx_index ASC; -- name: CreateValidatorRegistryEventsSyncedUntil :exec -INSERT INTO validator_registry_events_synced_until (block_number) VALUES ($1) +INSERT INTO validator_registry_events_synced_until (block_hash, block_number) VALUES ($1, $2) ON CONFLICT (enforce_one_row) DO UPDATE -SET block_number = $1; +SET block_hash = $1, block_number = $2; -- name: QueryValidatorRegistrationMessageNonceBefore :one SELECT nonce FROM validator_registration_message WHERE validator_index = $1 AND event_block_number <= $2 AND event_tx_index <= $3 AND event_log_index <= $4 ORDER BY event_block_number DESC, event_tx_index DESC, event_log_index DESC FOR UPDATE; -- name: QueryValidatorRegistryEventsSyncedUntil :one -SELECT block_number FROM validator_registry_events_synced_until LIMIT 1; +SELECT block_hash, block_number FROM validator_registry_events_synced_until LIMIT 1; -- name: CreateValidatorStatus :exec INSERT into validator_status( @@ -174,4 +174,12 @@ ON CONFLICT (slot, tx_index) DO UPDATE SET tx_status = $4, block_number = $7, - updated_at = NOW(); \ No newline at end of file + updated_at = NOW(); + +-- name: CreateTransactionSubmittedEventsSyncedUntil :exec +INSERT INTO transaction_submitted_events_synced_until (block_hash, block_number) VALUES ($1, $2) +ON CONFLICT (enforce_one_row) DO UPDATE +SET block_hash = $1, block_number = $2; + +-- name: QueryTransactionSubmittedEventsSyncedUntil :one +SELECT block_hash, block_number FROM transaction_submitted_events_synced_until LIMIT 1; diff --git a/internal/metrics/tx_mapper_db.go b/internal/metrics/tx_mapper_db.go index 8d08c7c..fc74d03 100644 --- a/internal/metrics/tx_mapper_db.go +++ b/internal/metrics/tx_mapper_db.go @@ -201,11 +201,11 @@ func (tm *TxMapperDB) AddBlock( } func (tm *TxMapperDB) QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx context.Context) (int64, error) { - blockNumber, err := tm.dbQuery.QueryValidatorRegistryEventsSyncedUntil(ctx) + data, err := tm.dbQuery.QueryValidatorRegistryEventsSyncedUntil(ctx) if err != nil { return 0, err } - return blockNumber, nil + return data.BlockNumber, nil } func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validatorRegistryBindings.ValidatorregistryUpdated) error { @@ -258,7 +258,10 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat } } - err = qtx.CreateValidatorRegistryEventsSyncedUntil(ctx, int64(vr.Raw.BlockNumber)) + err = qtx.CreateValidatorRegistryEventsSyncedUntil(ctx, data.CreateValidatorRegistryEventsSyncedUntilParams{ + BlockHash: vr.Raw.BlockHash[:], + BlockNumber: int64(vr.Raw.BlockNumber), + }) if err != nil { return err } diff --git a/internal/syncer/encrypted_tx.go b/internal/syncer/encrypted_tx.go new file mode 100644 index 0000000..2d25a3a --- /dev/null +++ b/internal/syncer/encrypted_tx.go @@ -0,0 +1,146 @@ +package syncer + +import ( + "context" + "math/big" + + "github.com/ethereum/go-ethereum/accounts/abi/bind" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/ethclient" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/pkg/errors" + "github.com/rs/zerolog/log" + sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" + "github.com/shutter-network/observer/internal/data" + "github.com/shutter-network/rolling-shutter/rolling-shutter/medley" +) + +const ( + AssumedReorgDepth = 10 + maxRequestBlockRange = 10_000 +) + +type EncrptedTxSyncer struct { + contract *sequencerBindings.Sequencer + db *pgxpool.Pool + dbQuery *data.Queries + ethClient *ethclient.Client + syncStartBlockNumber uint64 +} + +func (ets *EncrptedTxSyncer) Sync(ctx context.Context, header *types.Header) error { + // TODO: handle reorgs + syncedUntil, err := ets.dbQuery.QueryTransactionSubmittedEventsSyncedUntil(ctx) + if err != nil && err != pgx.ErrNoRows { + return errors.Wrap(err, "failed to query identity registered events sync status") + } + var start uint64 + if err == pgx.ErrNoRows { + start = ets.syncStartBlockNumber + } else { + start = uint64(syncedUntil.BlockNumber + 1) + } + endBlock := header.Number.Uint64() + log.Debug(). + Uint64("start-block", start). + Uint64("end-block", endBlock). + Msg("syncing transaction submitted events") + syncRanges := medley.GetSyncRanges(start, endBlock, maxRequestBlockRange) + for _, r := range syncRanges { + err = ets.syncRange(ctx, r[0], r[1]) + if err != nil { + return err + } + } + return nil +} + +func (ets *EncrptedTxSyncer) syncRange( + ctx context.Context, + start, + end uint64, +) error { + events, err := ets.fetchEvents(ctx, start, end) + if err != nil { + return err + } + + header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) + if err != nil { + return errors.Wrap(err, "failed to get execution block header by number") + } + + tx, err := ets.db.Begin(ctx) + if err != nil { + return err + } + defer tx.Rollback(ctx) + qtx := ets.dbQuery.WithTx(tx) + + for _, event := range events { + err := qtx.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ + EventBlockHash: event.Raw.BlockHash[:], + EventBlockNumber: int64(event.Raw.BlockNumber), + EventTxIndex: int64(event.Raw.TxIndex), + EventLogIndex: int64(event.Raw.Index), + Eon: int64(event.Eon), + TxIndex: int64(event.TxIndex), + IdentityPrefix: event.IdentityPrefix[:], + Sender: event.Sender[:], + EncryptedTransaction: event.EncryptedTransaction, + EventTxHash: event.Raw.TxHash[:], + }) + if err != nil { + log.Err(err).Msg("err adding transaction submitted event") + return nil + } + log.Info(). + Uint64("block", event.Raw.BlockNumber). + Hex("encrypted transaction (hex)", event.EncryptedTransaction). + Msg("new encrypted transaction") + } + err = qtx.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ + BlockNumber: int64(end), + BlockHash: header.Hash().Bytes(), + }) + if err != nil { + log.Err(err).Msg("error adding transaction submitted event until") + return nil + } + if err := tx.Commit(ctx); err != nil { + log.Err(err).Msg("unable to commit db transaction") + return err + } + + log.Info(). + Uint64("start-block", start). + Uint64("end-block", end). + Int("num-inserted-events", len(events)). + Msg("synced registry contract") + return nil +} + +func (s *EncrptedTxSyncer) fetchEvents( + ctx context.Context, + start, + end uint64, +) ([]*sequencerBindings.SequencerTransactionSubmitted, error) { + opts := bind.FilterOpts{ + Start: start, + End: &end, + Context: ctx, + } + it, err := s.contract.SequencerFilterer.FilterTransactionSubmitted(&opts) + if err != nil { + return nil, errors.Wrap(err, "failed to query transaction submitted events") + } + events := []*sequencerBindings.SequencerTransactionSubmitted{} + for it.Next() { + events = append(events, it.Event) + } + if it.Error() != nil { + return nil, errors.Wrap(it.Error(), "failed to iterate query transaction submitted events") + } + return events, nil +} diff --git a/migrations/20250327143729_add_tx_sub_event_sync_table.sql b/migrations/20250327143729_add_tx_sub_event_sync_table.sql new file mode 100644 index 0000000..dcca881 --- /dev/null +++ b/migrations/20250327143729_add_tx_sub_event_sync_table.sql @@ -0,0 +1,16 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE transaction_submitted_events_synced_until( + enforce_one_row bool PRIMARY KEY DEFAULT true, + block_hash bytea NOT NULL, + block_number bigint NOT NULL CHECK (block_number >= 0) +); + +ALTER TABLE validator_registry_events_synced_until +ADD COLUMN block_hash BYTEA; +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE transaction_submitted_events_synced_until; +-- +goose StatementEnd From f1b55e0fc8da5a1cdb415d0b1b4537526390e6ae Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Thu, 27 Mar 2025 21:47:49 +0100 Subject: [PATCH 04/13] add validator registry syncer --- internal/syncer/encrypted_tx.go | 2 +- internal/syncer/validator_registry.go | 102 ++++++++++++++++++++++++++ 2 files changed, 103 insertions(+), 1 deletion(-) create mode 100644 internal/syncer/validator_registry.go diff --git a/internal/syncer/encrypted_tx.go b/internal/syncer/encrypted_tx.go index 2d25a3a..94f7f6b 100644 --- a/internal/syncer/encrypted_tx.go +++ b/internal/syncer/encrypted_tx.go @@ -33,7 +33,7 @@ func (ets *EncrptedTxSyncer) Sync(ctx context.Context, header *types.Header) err // TODO: handle reorgs syncedUntil, err := ets.dbQuery.QueryTransactionSubmittedEventsSyncedUntil(ctx) if err != nil && err != pgx.ErrNoRows { - return errors.Wrap(err, "failed to query identity registered events sync status") + return errors.Wrap(err, "failed to query transaction submitted events sync status") } var start uint64 if err == pgx.ErrNoRows { diff --git a/internal/syncer/validator_registry.go b/internal/syncer/validator_registry.go new file mode 100644 index 0000000..568484e --- /dev/null +++ b/internal/syncer/validator_registry.go @@ -0,0 +1,102 @@ +package syncer + +import ( + "context" + + "github.com/ethereum/go-ethereum/accounts/abi/bind" + "github.com/ethereum/go-ethereum/core/types" + "github.com/jackc/pgx/v5" + "github.com/pkg/errors" + "github.com/rs/zerolog/log" + validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" + "github.com/shutter-network/observer/internal/data" + "github.com/shutter-network/observer/internal/metrics" + "github.com/shutter-network/rolling-shutter/rolling-shutter/medley" +) + +type ValidatorRegistrySyncer struct { + contract *validatorRegistryBindings.Validatorregistry + dbQuery *data.Queries + txMapper metrics.TxMapper + syncStartBlockNumber uint64 +} + +func (vts *ValidatorRegistrySyncer) Sync(ctx context.Context, header *types.Header) error { + // TODO: handle reorgs + syncedUntil, err := vts.dbQuery.QueryValidatorRegistryEventsSyncedUntil(ctx) + if err != nil && err != pgx.ErrNoRows { + return errors.Wrap(err, "failed to query validator registry sync status") + } + var start uint64 + if err == pgx.ErrNoRows { + start = vts.syncStartBlockNumber + } else { + start = uint64(syncedUntil.BlockNumber + 1) + } + endBlock := header.Number.Uint64() + log.Debug(). + Uint64("start-block", start). + Uint64("end-block", endBlock). + Msg("syncing validator registry updated events") + syncRanges := medley.GetSyncRanges(start, endBlock, maxRequestBlockRange) + for _, r := range syncRanges { + err = vts.syncRange(ctx, r[0], r[1]) + if err != nil { + return err + } + } + return nil +} + +func (ets *ValidatorRegistrySyncer) syncRange( + ctx context.Context, + start, + end uint64, +) error { + events, err := ets.fetchEvents(ctx, start, end) + if err != nil { + return err + } + + for _, event := range events { + err := ets.txMapper.AddValidatorRegistryEvent(ctx, event) + if err != nil { + log.Err(err).Msg("err adding validator registry updated event") + return nil + } + log.Info(). + Uint64("block", event.Raw.BlockNumber). + Msg("new validator registry updated message") + } + + log.Info(). + Uint64("start-block", start). + Uint64("end-block", end). + Int("num-inserted-events", len(events)). + Msg("synced validator registry contract") + return nil +} + +func (s *ValidatorRegistrySyncer) fetchEvents( + ctx context.Context, + start, + end uint64, +) ([]*validatorRegistryBindings.ValidatorregistryUpdated, error) { + opts := bind.FilterOpts{ + Start: start, + End: &end, + Context: ctx, + } + it, err := s.contract.ValidatorregistryFilterer.FilterUpdated(&opts) + if err != nil { + return nil, errors.Wrap(err, "failed to query validator registry updated events") + } + events := []*validatorRegistryBindings.ValidatorregistryUpdated{} + for it.Next() { + events = append(events, it.Event) + } + if it.Error() != nil { + return nil, errors.Wrap(it.Error(), "failed to iterate query validator registry updated events") + } + return events, nil +} From 476af3205eaa8acc35a39500550b500937fc7247 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Fri, 28 Mar 2025 00:03:18 +0100 Subject: [PATCH 05/13] refactor block processing --- internal/watcher/blocks.go | 129 +++++++++++++++++++++++----- internal/watcher/decryption_keys.go | 96 ++------------------- internal/watcher/p2p_msgs.go | 25 ++---- internal/watcher/watcher.go | 24 +----- 4 files changed, 122 insertions(+), 152 deletions(-) diff --git a/internal/watcher/blocks.go b/internal/watcher/blocks.go index aa0c4d3..e87d2c3 100644 --- a/internal/watcher/blocks.go +++ b/internal/watcher/blocks.go @@ -2,33 +2,41 @@ package watcher import ( "context" - "time" + "sync" "github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/ethclient" "github.com/rs/zerolog/log" "github.com/shutter-network/observer/common" + "github.com/shutter-network/observer/common/utils" + "github.com/shutter-network/observer/internal/data" + "github.com/shutter-network/observer/internal/metrics" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" ) type BlocksWatcher struct { - config *common.Config - blocksChannel chan *BlockReceivedEvent - blocksChannelForProposerDuties chan *BlockReceivedEvent - ethClient *ethclient.Client -} + config *common.Config + ethClient *ethclient.Client + + recentBlocksMux sync.Mutex + recentBlocks map[uint64]*types.Header + mostRecentBlock uint64 -type BlockReceivedEvent struct { - Header *types.Header - Time time.Time + txMapper metrics.TxMapper } -func NewBlocksWatcher(config *common.Config, blocksChannel chan *BlockReceivedEvent, blocksChannelForProposerDuties chan *BlockReceivedEvent, ethClient *ethclient.Client) *BlocksWatcher { +func NewBlocksWatcher( + config *common.Config, + ethClient *ethclient.Client, + txMapper metrics.TxMapper, +) *BlocksWatcher { return &BlocksWatcher{ - config: config, - blocksChannel: blocksChannel, - blocksChannelForProposerDuties: blocksChannelForProposerDuties, - ethClient: ethClient, + config: config, + ethClient: ethClient, + recentBlocksMux: sync.Mutex{}, + recentBlocks: make(map[uint64]*types.Header), + mostRecentBlock: 0, + txMapper: txMapper, } } @@ -45,16 +53,15 @@ func (bw *BlocksWatcher) Start(ctx context.Context, runner service.Runner) error case <-ctx.Done(): return ctx.Err() case head := <-newHeads: + err := bw.processBlock(ctx, head) + if err != nil { + log.Err(err).Msg("err processing new block") + return err + } log.Info(). Int64("number", head.Number.Int64()). Hex("hash", head.Hash().Bytes()). Msg("new head") - ev := &BlockReceivedEvent{ - Header: head, - Time: time.Now(), - } - bw.blocksChannel <- ev - bw.blocksChannelForProposerDuties <- ev case err := <-sub.Err(): return err } @@ -63,3 +70,85 @@ func (bw *BlocksWatcher) Start(ctx context.Context, runner service.Runner) error return nil } + +func (bw *BlocksWatcher) processBlock(ctx context.Context, header *types.Header) error { + err := bw.insertBlock(ctx, header) + if err != nil { + return err + } + bw.clearOldBlocks(header) + + epoch := utils.GetEpochForBlock(header.Time, GenesisTimestamp, SlotDuration, SlotsPerEpoch) + if epoch > CurrentEpoch { + CurrentEpoch = epoch + nextEpoch := epoch + 1 + err := bw.txMapper.AddProposerDuties(ctx, nextEpoch) + if err != nil { + return err + } + log.Info(). + Uint64("current epoch", epoch). + Uint64("next epoch", nextEpoch). + Msg("new proposer duties added") + } + return nil +} + +func (bw *BlocksWatcher) insertBlock(ctx context.Context, header *types.Header) error { + bw.recentBlocksMux.Lock() + defer bw.recentBlocksMux.Unlock() + bw.recentBlocks[header.Number.Uint64()] = header + if header.Number.Uint64() > bw.mostRecentBlock { + bw.mostRecentBlock = header.Number.Uint64() + } + + err := bw.txMapper.AddBlock(ctx, &data.Block{ + BlockHash: header.Hash().Bytes(), + BlockNumber: header.Number.Int64(), + BlockTimestamp: int64(header.Time), + Slot: int64(utils.GetSlotForBlock(header.Time, GenesisTimestamp, SlotDuration)), + }) + if err != nil { + log.Err(err).Msg("err adding block") + } + return err +} + +func (bw *BlocksWatcher) clearOldBlocks(latestHeader *types.Header) { + bw.recentBlocksMux.Lock() + defer bw.recentBlocksMux.Unlock() + + tooOld := []uint64{} + for block := range bw.recentBlocks { + if block < latestHeader.Number.Uint64()-100 { + tooOld = append(tooOld, block) + } + } + for _, block := range tooOld { + delete(bw.recentBlocks, block) + } +} + +func (bw *BlocksWatcher) getBlockHeaderFromSlot(slot uint64) (*types.Header, bool) { + bw.recentBlocksMux.Lock() + defer bw.recentBlocksMux.Unlock() + + slotTimestamp := utils.GetSlotTimestamp(slot, GenesisTimestamp, SlotDuration) + if header, ok := bw.recentBlocks[bw.mostRecentBlock]; ok { + if header.Time == slotTimestamp { + return header, ok + } else if header.Time < slotTimestamp { + return nil, false + } + } + + for blockNumber := range bw.recentBlocks { + if header, ok := bw.recentBlocks[blockNumber]; ok { + if header.Time == slotTimestamp { + return header, ok + } + } + } + + return nil, false +} diff --git a/internal/watcher/decryption_keys.go b/internal/watcher/decryption_keys.go index 3ba3834..3286b73 100644 --- a/internal/watcher/decryption_keys.go +++ b/internal/watcher/decryption_keys.go @@ -1,18 +1,12 @@ package watcher import ( - "context" - "fmt" - "time" - "github.com/rs/zerolog/log" "github.com/shutter-network/observer/common/utils" - "github.com/shutter-network/observer/internal/data" "github.com/shutter-network/rolling-shutter/rolling-shutter/p2pmsg" ) func (pmw *P2PMsgsWatcher) handleDecryptionKeyMsg(msg *p2pmsg.DecryptionKeys) ([]p2pmsg.Message, error) { - t := time.Now() extra := msg.Extra.(*p2pmsg.DecryptionKeys_Gnosis).Gnosis pmw.decryptionDataChannel <- &DecryptionKeysEvent{ Eon: int64(msg.Eon), @@ -22,113 +16,35 @@ func (pmw *P2PMsgsWatcher) handleDecryptionKeyMsg(msg *p2pmsg.DecryptionKeys) ([ TxPointer: int64(extra.TxPointer), } - ev, ok := pmw.getBlockReceivedEventFromSlot(extra.Slot) + _, ok := pmw.blocksWatcher.getBlockHeaderFromSlot(extra.Slot) + blocksWatcher := pmw.blocksWatcher if !ok { - if mostRecentBlock, ok := pmw.recentBlocks[pmw.mostRecentBlock]; ok { - mostRecentSlot := uint64(utils.GetSlotForBlock(mostRecentBlock.Header.Time, GenesisTimestamp, SlotDuration)) + if mostRecentBlockHeader, ok := blocksWatcher.recentBlocks[blocksWatcher.mostRecentBlock]; ok { + mostRecentSlot := uint64(utils.GetSlotForBlock(mostRecentBlockHeader.Time, GenesisTimestamp, SlotDuration)) if extra.Slot > mostRecentSlot+1 { log.Warn(). Uint64("slot", extra.Slot). Uint64("expected-slot", mostRecentSlot+1). - Uint64("most-recent-block", pmw.mostRecentBlock). + Uint64("most-recent-block", blocksWatcher.mostRecentBlock). Msg("received keys for a slot greater than expected slot") } log.Info(). Uint64("slot", extra.Slot). Int("num-keys", len(msg.Keys)). - Uint64("most-recent-block", pmw.mostRecentBlock). + Uint64("most-recent-block", blocksWatcher.mostRecentBlock). Uint64("most-recent-slot", mostRecentSlot). Msg("received keys for future slot") } return []p2pmsg.Message{}, nil } - dt := t.Sub(ev.Time) log.Warn(). Uint64("slot", extra.Slot). Int("num-keys", len(msg.Keys)). - Str("latency", fmt.Sprintf("%.2fs", dt.Seconds())). Msg("received keys for a known slot") return []p2pmsg.Message{}, nil } -func (pmw *P2PMsgsWatcher) insertBlocks(ctx context.Context) error { - for { - select { - case <-ctx.Done(): - return ctx.Err() - case ev, ok := <-pmw.blocksChannel: - if !ok { - return nil - } - err := pmw.insertBlock(ctx, ev) - if err != nil { - return err - } - pmw.clearOldBlocks(ev) - } - } -} - -func (pmw *P2PMsgsWatcher) insertBlock(ctx context.Context, ev *BlockReceivedEvent) error { - pmw.recentBlocksMux.Lock() - defer pmw.recentBlocksMux.Unlock() - pmw.recentBlocks[ev.Header.Number.Uint64()] = ev - if ev.Header.Number.Uint64() > pmw.mostRecentBlock { - pmw.mostRecentBlock = ev.Header.Number.Uint64() - } - - err := pmw.txMapper.AddBlock(ctx, &data.Block{ - BlockHash: ev.Header.Hash().Bytes(), - BlockNumber: ev.Header.Number.Int64(), - BlockTimestamp: int64(ev.Header.Time), - Slot: int64(utils.GetSlotForBlock(ev.Header.Time, GenesisTimestamp, SlotDuration)), - }) - if err != nil { - log.Err(err).Msg("err adding block") - } - return err -} - -func (pmw *P2PMsgsWatcher) clearOldBlocks(latestEv *BlockReceivedEvent) { - pmw.recentBlocksMux.Lock() - defer pmw.recentBlocksMux.Unlock() - - tooOld := []uint64{} - for block := range pmw.recentBlocks { - if block < latestEv.Header.Number.Uint64()-100 { - tooOld = append(tooOld, block) - } - } - for _, block := range tooOld { - delete(pmw.recentBlocks, block) - } -} - -func (pmw *P2PMsgsWatcher) getBlockReceivedEventFromSlot(slot uint64) (*BlockReceivedEvent, bool) { - pmw.recentBlocksMux.Lock() - defer pmw.recentBlocksMux.Unlock() - - slotTimestamp := utils.GetSlotTimestamp(slot, GenesisTimestamp, SlotDuration) - if ev, ok := pmw.recentBlocks[pmw.mostRecentBlock]; ok { - if ev.Header.Time == slotTimestamp { - return ev, ok - } else if ev.Header.Time < slotTimestamp { - return nil, false - } - } - - for blockNumber := range pmw.recentBlocks { - if ev, ok := pmw.recentBlocks[blockNumber]; ok { - if ev.Header.Time == slotTimestamp { - return ev, ok - } - } - } - - return nil, false -} - func getDecryptionKeysAndIdentities(p2pMsgs []*p2pmsg.Key) ([][]byte, [][]byte) { var keys [][]byte var identities [][]byte diff --git a/internal/watcher/p2p_msgs.go b/internal/watcher/p2p_msgs.go index 3574124..87c5745 100644 --- a/internal/watcher/p2p_msgs.go +++ b/internal/watcher/p2p_msgs.go @@ -3,13 +3,11 @@ package watcher import ( "context" "math" - "sync" pubsub "github.com/libp2p/go-libp2p-pubsub" "github.com/pkg/errors" "github.com/rs/zerolog/log" "github.com/shutter-network/observer/common" - "github.com/shutter-network/observer/internal/metrics" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" "github.com/shutter-network/rolling-shutter/rolling-shutter/p2p" "github.com/shutter-network/rolling-shutter/rolling-shutter/p2pmsg" @@ -17,15 +15,9 @@ import ( type P2PMsgsWatcher struct { config *common.Config - blocksChannel chan *BlockReceivedEvent decryptionDataChannel chan *DecryptionKeysEvent keyShareChannel chan *KeyShareEvent - - recentBlocksMux sync.Mutex - recentBlocks map[uint64]*BlockReceivedEvent - mostRecentBlock uint64 - - txMapper metrics.TxMapper + blocksWatcher *BlocksWatcher } type DecryptionKeysEvent struct { @@ -45,20 +37,15 @@ type KeyShareEvent struct { func NewP2PMsgsWatcherWatcher( config *common.Config, - blocksChannel chan *BlockReceivedEvent, decryptionDataChannel chan *DecryptionKeysEvent, keyShareChannel chan *KeyShareEvent, - txMapper metrics.TxMapper, + blocksWatcher *BlocksWatcher, ) *P2PMsgsWatcher { return &P2PMsgsWatcher{ config: config, - blocksChannel: blocksChannel, decryptionDataChannel: decryptionDataChannel, keyShareChannel: keyShareChannel, - recentBlocksMux: sync.Mutex{}, - recentBlocks: make(map[uint64]*BlockReceivedEvent), - mostRecentBlock: 0, - txMapper: txMapper, + blocksWatcher: blocksWatcher, } } @@ -69,8 +56,6 @@ func (pmw *P2PMsgsWatcher) Start(ctx context.Context, runner service.Runner) err } p2pService.AddMessageHandler(pmw) - runner.Go(func() error { return pmw.insertBlocks(ctx) }) - return runner.StartService(p2pService) } @@ -88,7 +73,7 @@ func (pmw *P2PMsgsWatcher) ValidateMessage(_ context.Context, msgUntyped p2pmsg. if extra == nil { log.Warn(). Int("num-keys", len(msg.Keys)). - Uint64("most-recent-block", pmw.mostRecentBlock). + Uint64("most-recent-block", pmw.blocksWatcher.mostRecentBlock). Msg("received DecryptionKeys without any slot") return pubsub.ValidationReject, nil } @@ -103,7 +88,7 @@ func (pmw *P2PMsgsWatcher) ValidateMessage(_ context.Context, msgUntyped p2pmsg. if extra == nil { log.Warn(). Int("num-keyshares", len(msg.Shares)). - Uint64("most-recent-block", pmw.mostRecentBlock). + Uint64("most-recent-block", pmw.blocksWatcher.mostRecentBlock). Msg("received DecryptionKeyShares without any slot") return pubsub.ValidationReject, nil } diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index f274387..77e996f 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -16,7 +16,6 @@ import ( sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" "github.com/shutter-network/observer/common" - "github.com/shutter-network/observer/common/utils" "github.com/shutter-network/observer/internal/data" "github.com/shutter-network/observer/internal/metrics" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/beaconapiclient" @@ -65,10 +64,8 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { txSubmittedEventChannel := make(chan *sequencerBindings.SequencerTransactionSubmitted) validatorRegistryChannel := make(chan *validatorRegistryBindings.ValidatorregistryUpdated) - blocksChannel := make(chan *BlockReceivedEvent) decryptionDataChannel := make(chan *DecryptionKeysEvent) keyShareChannel := make(chan *KeyShareEvent) - blocksChannelForProposerDuties := make(chan *BlockReceivedEvent) dialer := rpc.WithWebsocketDialer(websocket.Dialer{ HandshakeTimeout: 45 * time.Second, @@ -107,7 +104,7 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { SlotDuration, ) - blocksWatcher := NewBlocksWatcher(w.config, blocksChannel, blocksChannelForProposerDuties, ethClient) + blocksWatcher := NewBlocksWatcher(w.config, ethClient, txMapper) encryptionTxWatcher := NewEncryptedTxWatcher(w.config, txSubmittedEventChannel, ethClient) blockNumber, err := txMapper.QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx) @@ -121,7 +118,7 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { validatorRegisterWatcher := NewValidatorRegistryWatcher(w.config, validatorRegistryChannel, ethClient, blockNumber) - p2pMsgsWatcher := NewP2PMsgsWatcherWatcher(w.config, blocksChannel, decryptionDataChannel, keyShareChannel, txMapper) + p2pMsgsWatcher := NewP2PMsgsWatcherWatcher(w.config, decryptionDataChannel, keyShareChannel, blocksWatcher) if err := runner.StartService(blocksWatcher, encryptionTxWatcher, p2pMsgsWatcher, validatorRegisterWatcher); err != nil { return err } @@ -129,23 +126,6 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { runner.Go(func() error { for { select { - case ev, ok := <-blocksChannelForProposerDuties: - if !ok { - return nil - } - epoch := utils.GetEpochForBlock(ev.Header.Time, GenesisTimestamp, SlotDuration, SlotsPerEpoch) - if epoch > CurrentEpoch { - CurrentEpoch = epoch - nextEpoch := epoch + 1 - err := txMapper.AddProposerDuties(ctx, nextEpoch) - if err != nil { - return err - } - log.Info(). - Uint64("current epoch", epoch). - Uint64("next epoch", nextEpoch). - Msg("new proposer duties added") - } case txEvent := <-txSubmittedEventChannel: err := txMapper.AddTransactionSubmittedEvent(ctx, &data.TransactionSubmittedEvent{ EventBlockHash: txEvent.Raw.BlockHash[:], From d53230f5f3ab5b78f9f5d1608d05b710679aa33e Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Fri, 28 Mar 2025 00:10:02 +0100 Subject: [PATCH 06/13] get rid of ununsed watchers --- internal/watcher/encrypted_tx.go | 53 ------------- internal/watcher/validator_registry.go | 102 ------------------------- internal/watcher/watcher.go | 47 +----------- 3 files changed, 1 insertion(+), 201 deletions(-) delete mode 100644 internal/watcher/encrypted_tx.go delete mode 100644 internal/watcher/validator_registry.go diff --git a/internal/watcher/encrypted_tx.go b/internal/watcher/encrypted_tx.go deleted file mode 100644 index 1be253c..0000000 --- a/internal/watcher/encrypted_tx.go +++ /dev/null @@ -1,53 +0,0 @@ -package watcher - -import ( - "context" - - "github.com/ethereum/go-ethereum/accounts/abi/bind" - "github.com/ethereum/go-ethereum/common" - "github.com/ethereum/go-ethereum/ethclient" - "github.com/rs/zerolog/log" - sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" - metricsCommon "github.com/shutter-network/observer/common" - "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" -) - -type EncryptedTxWatcher struct { - config *metricsCommon.Config - txSubmittedEventChannel chan *sequencerBindings.SequencerTransactionSubmitted - ethClient *ethclient.Client -} - -func NewEncryptedTxWatcher(config *metricsCommon.Config, txSubmittedEventChannel chan *sequencerBindings.SequencerTransactionSubmitted, ethClient *ethclient.Client) *EncryptedTxWatcher { - return &EncryptedTxWatcher{ - config: config, - txSubmittedEventChannel: txSubmittedEventChannel, - ethClient: ethClient, - } -} - -func (etw *EncryptedTxWatcher) Start(ctx context.Context, runner service.Runner) error { - sequencerContract, err := sequencerBindings.NewSequencer(common.HexToAddress(etw.config.SequencerContractAddress), etw.ethClient) - if err != nil { - return err - } - watchOpts := &bind.WatchOpts{Context: ctx, Start: nil} - sub, err := sequencerContract.WatchTransactionSubmitted(watchOpts, etw.txSubmittedEventChannel) - if err != nil { - return err - } - runner.Defer(sub.Unsubscribe) - - log.Debug().Msg("Successfully subscribed to TransactionSubmitted events") - runner.Go(func() error { - for { - select { - case <-ctx.Done(): - return ctx.Err() - case err := <-sub.Err(): - return err - } - } - }) - return nil -} diff --git a/internal/watcher/validator_registry.go b/internal/watcher/validator_registry.go deleted file mode 100644 index 4d60755..0000000 --- a/internal/watcher/validator_registry.go +++ /dev/null @@ -1,102 +0,0 @@ -package watcher - -import ( - "context" - - "github.com/ethereum/go-ethereum/accounts/abi/bind" - "github.com/ethereum/go-ethereum/common" - "github.com/ethereum/go-ethereum/ethclient" - "github.com/pkg/errors" - "github.com/rs/zerolog/log" - validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" - metricsCommon "github.com/shutter-network/observer/common" - "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" -) - -type ValidatorRegistryWatcher struct { - config *metricsCommon.Config - validatorRegistryChannel chan *validatorRegistryBindings.ValidatorregistryUpdated - ethClient *ethclient.Client - startBlock int64 -} - -func NewValidatorRegistryWatcher( - config *metricsCommon.Config, - validatorRegistryChannel chan *validatorRegistryBindings.ValidatorregistryUpdated, - ethClient *ethclient.Client, - startBlock int64, -) *ValidatorRegistryWatcher { - return &ValidatorRegistryWatcher{ - config: config, - validatorRegistryChannel: validatorRegistryChannel, - ethClient: ethClient, - startBlock: startBlock, - } -} - -func (vrw *ValidatorRegistryWatcher) Start(ctx context.Context, runner service.Runner) error { - newValidatorRegistryUpdatesMsgs := make(chan *validatorRegistryBindings.ValidatorregistryUpdated) - validatorRegistryContract, err := validatorRegistryBindings.NewValidatorregistry(common.HexToAddress(vrw.config.ValidatorRegistryContractAddress), vrw.ethClient) - if err != nil { - return err - } - - //sync previous blocks which have not been processed yet - runner.Go(func() error { - err = vrw.syncPreviousBlocks(ctx, validatorRegistryContract) - if err != nil { - log.Err(err).Msg("err syncing previous blocks for validator registry") - return err - } - return nil - }) - - watchOpts := &bind.WatchOpts{Context: ctx, Start: nil} - sub, err := validatorRegistryContract.WatchUpdated(watchOpts, newValidatorRegistryUpdatesMsgs) - if err != nil { - return err - } - runner.Defer(sub.Unsubscribe) - - log.Debug().Msg("Successfully subscribed to validator register event") - runner.Go(func() error { - for { - select { - case <-ctx.Done(): - return ctx.Err() - case vruMsg := <-newValidatorRegistryUpdatesMsgs: - if vruMsg.Raw.BlockNumber > uint64(vrw.startBlock) { - // process only if block greater then start block - // since event have been or will be indexed by syncPreviousBlocks - // till startBlock - vrw.validatorRegistryChannel <- vruMsg - } - case err := <-sub.Err(): - return err - } - } - }) - return nil -} - -func (vrw *ValidatorRegistryWatcher) syncPreviousBlocks( - ctx context.Context, - validatorRegistryContract *validatorRegistryBindings.Validatorregistry, -) error { - filterOpts := &bind.FilterOpts{Context: ctx, Start: uint64(vrw.startBlock)} - events, err := validatorRegistryContract.FilterUpdated(filterOpts) - if err != nil { - return err - } - - for events.Next() { - event := events.Event - vrw.validatorRegistryChannel <- event - } - - if events.Error() != nil { - return errors.Wrap(events.Error(), "failed to iterate validator registry updated events") - } - - return nil -} diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index 77e996f..61c4e5a 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -10,11 +10,8 @@ import ( "github.com/ethereum/go-ethereum/ethclient" "github.com/ethereum/go-ethereum/rpc" "github.com/gorilla/websocket" - "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" "github.com/rs/zerolog/log" - sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" - validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" "github.com/shutter-network/observer/common" "github.com/shutter-network/observer/internal/data" "github.com/shutter-network/observer/internal/metrics" @@ -61,9 +58,6 @@ func New( } func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { - txSubmittedEventChannel := make(chan *sequencerBindings.SequencerTransactionSubmitted) - validatorRegistryChannel := make(chan *validatorRegistryBindings.ValidatorregistryUpdated) - decryptionDataChannel := make(chan *DecryptionKeysEvent) keyShareChannel := make(chan *KeyShareEvent) @@ -105,47 +99,14 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { ) blocksWatcher := NewBlocksWatcher(w.config, ethClient, txMapper) - encryptionTxWatcher := NewEncryptedTxWatcher(w.config, txSubmittedEventChannel, ethClient) - - blockNumber, err := txMapper.QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx) - if err != nil { - if err == pgx.ErrNoRows { - blockNumber = int64(ValidatorRegistryDeploymentBlockNumber) - } else { - return err - } - } - - validatorRegisterWatcher := NewValidatorRegistryWatcher(w.config, validatorRegistryChannel, ethClient, blockNumber) - p2pMsgsWatcher := NewP2PMsgsWatcherWatcher(w.config, decryptionDataChannel, keyShareChannel, blocksWatcher) - if err := runner.StartService(blocksWatcher, encryptionTxWatcher, p2pMsgsWatcher, validatorRegisterWatcher); err != nil { + if err := runner.StartService(blocksWatcher, p2pMsgsWatcher); err != nil { return err } runner.Go(func() error { for { select { - case txEvent := <-txSubmittedEventChannel: - err := txMapper.AddTransactionSubmittedEvent(ctx, &data.TransactionSubmittedEvent{ - EventBlockHash: txEvent.Raw.BlockHash[:], - EventBlockNumber: int64(txEvent.Raw.BlockNumber), - EventTxIndex: int64(txEvent.Raw.TxIndex), - EventLogIndex: int64(txEvent.Raw.Index), - Eon: int64(txEvent.Eon), - TxIndex: int64(txEvent.TxIndex), - IdentityPrefix: txEvent.IdentityPrefix[:], - Sender: txEvent.Sender[:], - EncryptedTransaction: txEvent.EncryptedTransaction, - EventTxHash: txEvent.Raw.TxHash[:], - }) - if err != nil { - log.Err(err).Msg("err adding encrypting transaction") - return err - } - log.Info(). - Hex("encrypted transaction (hex)", txEvent.EncryptedTransaction). - Msg("new encrypted transaction") case dd := <-decryptionDataChannel: keys, identites := getDecryptionKeysAndIdentities(dd.Keys) err := txMapper.AddDecryptionKeysAndMessages( @@ -186,12 +147,6 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { Int64("slot", ks.Slot). Msg("new key shares") } - case vr := <-validatorRegistryChannel: - err = txMapper.AddValidatorRegistryEvent(ctx, vr) - if err != nil { - log.Err(err).Msg("err adding validator registry") - return err - } case <-ctx.Done(): return ctx.Err() } From 7cd5a071ff0b384ab228c572e5e712ca387f7452 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Fri, 28 Mar 2025 00:39:10 +0100 Subject: [PATCH 07/13] clean syncers --- internal/metrics/tx_mapper_db.go | 45 ++++++++++++++++++++------- internal/metrics/types.go | 3 +- internal/syncer/encrypted_tx.go | 42 ++----------------------- internal/syncer/validator_registry.go | 2 ++ 4 files changed, 40 insertions(+), 52 deletions(-) diff --git a/internal/metrics/tx_mapper_db.go b/internal/metrics/tx_mapper_db.go index fc74d03..f0bd9b3 100644 --- a/internal/metrics/tx_mapper_db.go +++ b/internal/metrics/tx_mapper_db.go @@ -17,6 +17,7 @@ import ( "github.com/jackc/pgx/v5/pgxpool" "github.com/pkg/errors" "github.com/rs/zerolog/log" + sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" metricsCommon "github.com/shutter-network/observer/common" dbTypes "github.com/shutter-network/observer/common/database" @@ -68,22 +69,42 @@ func NewTxMapperDB( } } -func (tm *TxMapperDB) AddTransactionSubmittedEvent(ctx context.Context, tse *data.TransactionSubmittedEvent) error { - err := tm.dbQuery.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ - EventBlockHash: tse.EventBlockHash, - EventBlockNumber: tse.EventBlockNumber, - EventTxIndex: tse.EventTxIndex, - EventLogIndex: tse.EventLogIndex, - Eon: tse.Eon, - TxIndex: tse.TxIndex, - IdentityPrefix: tse.IdentityPrefix, - Sender: tse.Sender, - EncryptedTransaction: tse.EncryptedTransaction, - EventTxHash: tse.EventTxHash, +func (tm *TxMapperDB) AddTransactionSubmittedEvent(ctx context.Context, st *sequencerBindings.SequencerTransactionSubmitted) error { + tx, err := tm.db.Begin(ctx) + if err != nil { + return err + } + defer tx.Rollback(ctx) + qtx := tm.dbQuery.WithTx(tx) + + err = tm.dbQuery.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ + EventBlockHash: st.Raw.BlockHash.Bytes(), + EventBlockNumber: int64(st.Raw.BlockNumber), + EventTxIndex: int64(st.Raw.TxIndex), + EventLogIndex: int64(st.Raw.Index), + Eon: int64(st.Eon), + TxIndex: int64(st.TxIndex), + IdentityPrefix: st.IdentityPrefix[:], + Sender: st.Sender.Bytes(), + EncryptedTransaction: st.EncryptedTransaction, + EventTxHash: st.Raw.TxHash.Bytes(), }) if err != nil { return err } + err = qtx.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ + BlockHash: st.Raw.BlockHash[:], + BlockNumber: int64(st.Raw.BlockNumber), + }) + if err != nil { + log.Err(err).Msg("error adding transaction submitted event until") + return err + } + err = tx.Commit(ctx) + if err != nil { + log.Err(err).Msg("error commiting data in the db") + return err + } metricsEncTxReceived.Inc() return nil } diff --git a/internal/metrics/types.go b/internal/metrics/types.go index 383b956..a13ae0b 100644 --- a/internal/metrics/types.go +++ b/internal/metrics/types.go @@ -3,6 +3,7 @@ package metrics import ( "context" + sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" "github.com/shutter-network/observer/internal/data" ) @@ -49,7 +50,7 @@ type TxExecution struct { } type TxMapper interface { - AddTransactionSubmittedEvent(ctx context.Context, tse *data.TransactionSubmittedEvent) error + AddTransactionSubmittedEvent(ctx context.Context, st *sequencerBindings.SequencerTransactionSubmitted) error AddDecryptionKeysAndMessages( ctx context.Context, dkam *DecKeysAndMessages, diff --git a/internal/syncer/encrypted_tx.go b/internal/syncer/encrypted_tx.go index 94f7f6b..93a738b 100644 --- a/internal/syncer/encrypted_tx.go +++ b/internal/syncer/encrypted_tx.go @@ -2,7 +2,6 @@ package syncer import ( "context" - "math/big" "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/core/types" @@ -13,6 +12,7 @@ import ( "github.com/rs/zerolog/log" sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" "github.com/shutter-network/observer/internal/data" + "github.com/shutter-network/observer/internal/metrics" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley" ) @@ -26,6 +26,7 @@ type EncrptedTxSyncer struct { db *pgxpool.Pool dbQuery *data.Queries ethClient *ethclient.Client + txMapper metrics.TxMapper syncStartBlockNumber uint64 } @@ -65,32 +66,8 @@ func (ets *EncrptedTxSyncer) syncRange( if err != nil { return err } - - header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) - if err != nil { - return errors.Wrap(err, "failed to get execution block header by number") - } - - tx, err := ets.db.Begin(ctx) - if err != nil { - return err - } - defer tx.Rollback(ctx) - qtx := ets.dbQuery.WithTx(tx) - for _, event := range events { - err := qtx.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ - EventBlockHash: event.Raw.BlockHash[:], - EventBlockNumber: int64(event.Raw.BlockNumber), - EventTxIndex: int64(event.Raw.TxIndex), - EventLogIndex: int64(event.Raw.Index), - Eon: int64(event.Eon), - TxIndex: int64(event.TxIndex), - IdentityPrefix: event.IdentityPrefix[:], - Sender: event.Sender[:], - EncryptedTransaction: event.EncryptedTransaction, - EventTxHash: event.Raw.TxHash[:], - }) + err := ets.txMapper.AddTransactionSubmittedEvent(ctx, event) if err != nil { log.Err(err).Msg("err adding transaction submitted event") return nil @@ -100,19 +77,6 @@ func (ets *EncrptedTxSyncer) syncRange( Hex("encrypted transaction (hex)", event.EncryptedTransaction). Msg("new encrypted transaction") } - err = qtx.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ - BlockNumber: int64(end), - BlockHash: header.Hash().Bytes(), - }) - if err != nil { - log.Err(err).Msg("error adding transaction submitted event until") - return nil - } - if err := tx.Commit(ctx); err != nil { - log.Err(err).Msg("unable to commit db transaction") - return err - } - log.Info(). Uint64("start-block", start). Uint64("end-block", end). diff --git a/internal/syncer/validator_registry.go b/internal/syncer/validator_registry.go index 568484e..356a116 100644 --- a/internal/syncer/validator_registry.go +++ b/internal/syncer/validator_registry.go @@ -6,6 +6,7 @@ import ( "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/core/types" "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" "github.com/pkg/errors" "github.com/rs/zerolog/log" validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" @@ -16,6 +17,7 @@ import ( type ValidatorRegistrySyncer struct { contract *validatorRegistryBindings.Validatorregistry + db *pgxpool.Pool dbQuery *data.Queries txMapper metrics.TxMapper syncStartBlockNumber uint64 From 2783c2a964d10497cda8ff18bde4e15c49c51384 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Fri, 28 Mar 2025 11:14:35 +0100 Subject: [PATCH 08/13] hook syncers --- ..._tx.go => transaction_submitted_syncer.go} | 25 +++++++++++--- internal/syncer/validator_registry.go | 15 +++++++++ internal/watcher/blocks.go | 33 +++++++++++++------ internal/watcher/watcher.go | 24 +++++++++++++- 4 files changed, 82 insertions(+), 15 deletions(-) rename internal/syncer/{encrypted_tx.go => transaction_submitted_syncer.go} (80%) diff --git a/internal/syncer/encrypted_tx.go b/internal/syncer/transaction_submitted_syncer.go similarity index 80% rename from internal/syncer/encrypted_tx.go rename to internal/syncer/transaction_submitted_syncer.go index 93a738b..1b4781d 100644 --- a/internal/syncer/encrypted_tx.go +++ b/internal/syncer/transaction_submitted_syncer.go @@ -21,7 +21,7 @@ const ( maxRequestBlockRange = 10_000 ) -type EncrptedTxSyncer struct { +type TransactionSubmittedSyncer struct { contract *sequencerBindings.Sequencer db *pgxpool.Pool dbQuery *data.Queries @@ -30,7 +30,24 @@ type EncrptedTxSyncer struct { syncStartBlockNumber uint64 } -func (ets *EncrptedTxSyncer) Sync(ctx context.Context, header *types.Header) error { +func NewTransactionSubmittedSyncer( + contract *sequencerBindings.Sequencer, + db *pgxpool.Pool, + ethClient *ethclient.Client, + txMapper metrics.TxMapper, + syncStartBlockNumber uint64, +) *TransactionSubmittedSyncer { + return &TransactionSubmittedSyncer{ + contract: contract, + db: db, + dbQuery: data.New(db), + ethClient: ethClient, + txMapper: txMapper, + syncStartBlockNumber: syncStartBlockNumber, + } +} + +func (ets *TransactionSubmittedSyncer) Sync(ctx context.Context, header *types.Header) error { // TODO: handle reorgs syncedUntil, err := ets.dbQuery.QueryTransactionSubmittedEventsSyncedUntil(ctx) if err != nil && err != pgx.ErrNoRows { @@ -57,7 +74,7 @@ func (ets *EncrptedTxSyncer) Sync(ctx context.Context, header *types.Header) err return nil } -func (ets *EncrptedTxSyncer) syncRange( +func (ets *TransactionSubmittedSyncer) syncRange( ctx context.Context, start, end uint64, @@ -85,7 +102,7 @@ func (ets *EncrptedTxSyncer) syncRange( return nil } -func (s *EncrptedTxSyncer) fetchEvents( +func (s *TransactionSubmittedSyncer) fetchEvents( ctx context.Context, start, end uint64, diff --git a/internal/syncer/validator_registry.go b/internal/syncer/validator_registry.go index 356a116..4af2e7e 100644 --- a/internal/syncer/validator_registry.go +++ b/internal/syncer/validator_registry.go @@ -23,6 +23,21 @@ type ValidatorRegistrySyncer struct { syncStartBlockNumber uint64 } +func NewValidatorRegistrySyncer( + contract *validatorRegistryBindings.Validatorregistry, + db *pgxpool.Pool, + txMapper metrics.TxMapper, + syncStartBlockNumber uint64, +) *ValidatorRegistrySyncer { + return &ValidatorRegistrySyncer{ + contract: contract, + db: db, + dbQuery: data.New(db), + txMapper: txMapper, + syncStartBlockNumber: syncStartBlockNumber, + } +} + func (vts *ValidatorRegistrySyncer) Sync(ctx context.Context, header *types.Header) error { // TODO: handle reorgs syncedUntil, err := vts.dbQuery.QueryValidatorRegistryEventsSyncedUntil(ctx) diff --git a/internal/watcher/blocks.go b/internal/watcher/blocks.go index e87d2c3..fb9ff45 100644 --- a/internal/watcher/blocks.go +++ b/internal/watcher/blocks.go @@ -11,32 +11,38 @@ import ( "github.com/shutter-network/observer/common/utils" "github.com/shutter-network/observer/internal/data" "github.com/shutter-network/observer/internal/metrics" + "github.com/shutter-network/observer/internal/syncer" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" ) type BlocksWatcher struct { - config *common.Config - ethClient *ethclient.Client + config *common.Config + ethClient *ethclient.Client + txMapper metrics.TxMapper + transactionSubmittedSyncer *syncer.TransactionSubmittedSyncer + validatorRegistrySyncer *syncer.ValidatorRegistrySyncer recentBlocksMux sync.Mutex recentBlocks map[uint64]*types.Header mostRecentBlock uint64 - - txMapper metrics.TxMapper } func NewBlocksWatcher( config *common.Config, ethClient *ethclient.Client, txMapper metrics.TxMapper, + transactionSubmittedSyncer *syncer.TransactionSubmittedSyncer, + validatorRegistrySyncer *syncer.ValidatorRegistrySyncer, ) *BlocksWatcher { return &BlocksWatcher{ - config: config, - ethClient: ethClient, - recentBlocksMux: sync.Mutex{}, - recentBlocks: make(map[uint64]*types.Header), - mostRecentBlock: 0, - txMapper: txMapper, + config: config, + ethClient: ethClient, + txMapper: txMapper, + transactionSubmittedSyncer: transactionSubmittedSyncer, + validatorRegistrySyncer: validatorRegistrySyncer, + recentBlocksMux: sync.Mutex{}, + recentBlocks: make(map[uint64]*types.Header), + mostRecentBlock: 0, } } @@ -78,6 +84,13 @@ func (bw *BlocksWatcher) processBlock(ctx context.Context, header *types.Header) } bw.clearOldBlocks(header) + if err := bw.transactionSubmittedSyncer.Sync(ctx, header); err != nil { + return err + } + + if err := bw.validatorRegistrySyncer.Sync(ctx, header); err != nil { + return err + } epoch := utils.GetEpochForBlock(header.Time, GenesisTimestamp, SlotDuration, SlotsPerEpoch) if epoch > CurrentEpoch { CurrentEpoch = epoch diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index 61c4e5a..9150caf 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -7,14 +7,18 @@ import ( "net" "time" + ethCommon "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/ethclient" "github.com/ethereum/go-ethereum/rpc" "github.com/gorilla/websocket" "github.com/jackc/pgx/v5/pgxpool" "github.com/rs/zerolog/log" + sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" + validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" "github.com/shutter-network/observer/common" "github.com/shutter-network/observer/internal/data" "github.com/shutter-network/observer/internal/metrics" + "github.com/shutter-network/observer/internal/syncer" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/beaconapiclient" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" ) @@ -98,7 +102,25 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { SlotDuration, ) - blocksWatcher := NewBlocksWatcher(w.config, ethClient, txMapper) + sequencerContract, err := sequencerBindings.NewSequencer(ethCommon.HexToAddress(w.config.SequencerContractAddress), ethClient) + if err != nil { + return err + } + + validatorRegistryContract, err := validatorRegistryBindings.NewValidatorregistry(ethCommon.HexToAddress(w.config.ValidatorRegistryContractAddress), ethClient) + if err != nil { + return err + } + + blockNumber, err := ethClient.BlockNumber(ctx) + if err != nil { + return err + } + + transactionSubmittedSyncer := syncer.NewTransactionSubmittedSyncer(sequencerContract, w.db, ethClient, txMapper, blockNumber) + validatorRegistrySyncer := syncer.NewValidatorRegistrySyncer(validatorRegistryContract, w.db, txMapper, ValidatorRegistryDeploymentBlockNumber) + + blocksWatcher := NewBlocksWatcher(w.config, ethClient, txMapper, transactionSubmittedSyncer, validatorRegistrySyncer) p2pMsgsWatcher := NewP2PMsgsWatcherWatcher(w.config, decryptionDataChannel, keyShareChannel, blocksWatcher) if err := runner.StartService(blocksWatcher, p2pMsgsWatcher); err != nil { return err From 77df323fbc9dd8ed2d09e34fb4fe8c2943cf57c9 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Sun, 30 Mar 2025 20:43:26 +0200 Subject: [PATCH 09/13] fix query --- internal/metrics/tx_mapper_db.go | 45 +++---------------- .../syncer/transaction_submitted_syncer.go | 17 ++++++- internal/syncer/validator_registry.go | 20 ++++++++- internal/watcher/watcher.go | 2 +- 4 files changed, 40 insertions(+), 44 deletions(-) diff --git a/internal/metrics/tx_mapper_db.go b/internal/metrics/tx_mapper_db.go index f0bd9b3..d23ac2f 100644 --- a/internal/metrics/tx_mapper_db.go +++ b/internal/metrics/tx_mapper_db.go @@ -70,14 +70,7 @@ func NewTxMapperDB( } func (tm *TxMapperDB) AddTransactionSubmittedEvent(ctx context.Context, st *sequencerBindings.SequencerTransactionSubmitted) error { - tx, err := tm.db.Begin(ctx) - if err != nil { - return err - } - defer tx.Rollback(ctx) - qtx := tm.dbQuery.WithTx(tx) - - err = tm.dbQuery.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ + err := tm.dbQuery.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ EventBlockHash: st.Raw.BlockHash.Bytes(), EventBlockNumber: int64(st.Raw.BlockNumber), EventTxIndex: int64(st.Raw.TxIndex), @@ -92,19 +85,6 @@ func (tm *TxMapperDB) AddTransactionSubmittedEvent(ctx context.Context, st *sequ if err != nil { return err } - err = qtx.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ - BlockHash: st.Raw.BlockHash[:], - BlockNumber: int64(st.Raw.BlockNumber), - }) - if err != nil { - log.Err(err).Msg("error adding transaction submitted event until") - return err - } - err = tx.Commit(ctx) - if err != nil { - log.Err(err).Msg("error commiting data in the db") - return err - } metricsEncTxReceived.Inc() return nil } @@ -230,15 +210,8 @@ func (tm *TxMapperDB) QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx } func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validatorRegistryBindings.ValidatorregistryUpdated) error { - tx, err := tm.db.Begin(ctx) - if err != nil { - return err - } - defer tx.Rollback(ctx) - qtx := tm.dbQuery.WithTx(tx) - regMessage := &validatorregistry.AggregateRegistrationMessage{} - err = regMessage.Unmarshal(vr.Message) + err := regMessage.Unmarshal(vr.Message) if err != nil { log.Err(err).Hex("tx-hash", vr.Raw.TxHash.Bytes()).Msg("error unmarshalling registration message") } else { @@ -249,7 +222,7 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat } for validatorID, validatorData := range validatorIDtoValidity { - err := qtx.CreateValidatorRegistryMessage(ctx, data.CreateValidatorRegistryMessageParams{ + err := tm.dbQuery.CreateValidatorRegistryMessage(ctx, data.CreateValidatorRegistryMessageParams{ Version: dbTypes.Uint64ToPgTypeInt8(uint64(regMessage.Version)), ChainID: dbTypes.Uint64ToPgTypeInt8(regMessage.ChainID), ValidatorRegistryAddress: regMessage.ValidatorRegistryAddress.Bytes(), @@ -268,7 +241,7 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat if validatorData.validatorValidity == data.ValidatorRegistrationValidityValid && validatorData.validatorStatus != "" { - err := qtx.CreateValidatorStatus(ctx, data.CreateValidatorStatusParams{ + err := tm.dbQuery.CreateValidatorStatus(ctx, data.CreateValidatorStatusParams{ ValidatorIndex: dbTypes.Int64ToPgTypeInt8(validatorID), Status: validatorData.validatorStatus, }) @@ -278,15 +251,7 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat } } } - - err = qtx.CreateValidatorRegistryEventsSyncedUntil(ctx, data.CreateValidatorRegistryEventsSyncedUntilParams{ - BlockHash: vr.Raw.BlockHash[:], - BlockNumber: int64(vr.Raw.BlockNumber), - }) - if err != nil { - return err - } - return tx.Commit(ctx) + return nil } func (tm *TxMapperDB) UpdateValidatorStatus(ctx context.Context) error { diff --git a/internal/syncer/transaction_submitted_syncer.go b/internal/syncer/transaction_submitted_syncer.go index 1b4781d..af17cf5 100644 --- a/internal/syncer/transaction_submitted_syncer.go +++ b/internal/syncer/transaction_submitted_syncer.go @@ -2,6 +2,7 @@ package syncer import ( "context" + "math/big" "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/core/types" @@ -83,22 +84,34 @@ func (ets *TransactionSubmittedSyncer) syncRange( if err != nil { return err } + header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) + if err != nil { + return errors.Wrap(err, "failed to get execution block header by number") + } for _, event := range events { err := ets.txMapper.AddTransactionSubmittedEvent(ctx, event) if err != nil { log.Err(err).Msg("err adding transaction submitted event") - return nil + return err } log.Info(). Uint64("block", event.Raw.BlockNumber). Hex("encrypted transaction (hex)", event.EncryptedTransaction). Msg("new encrypted transaction") } + err = ets.dbQuery.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ + BlockNumber: int64(end), + BlockHash: header.Hash().Bytes(), + }) + if err != nil { + log.Err(err).Msg("err adding transaction submit event sync until") + return err + } log.Info(). Uint64("start-block", start). Uint64("end-block", end). Int("num-inserted-events", len(events)). - Msg("synced registry contract") + Msg("synced sequencer contract") return nil } diff --git a/internal/syncer/validator_registry.go b/internal/syncer/validator_registry.go index 4af2e7e..b845e65 100644 --- a/internal/syncer/validator_registry.go +++ b/internal/syncer/validator_registry.go @@ -2,9 +2,11 @@ package syncer import ( "context" + "math/big" "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/ethclient" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" "github.com/pkg/errors" @@ -19,6 +21,7 @@ type ValidatorRegistrySyncer struct { contract *validatorRegistryBindings.Validatorregistry db *pgxpool.Pool dbQuery *data.Queries + ethClient *ethclient.Client txMapper metrics.TxMapper syncStartBlockNumber uint64 } @@ -26,6 +29,7 @@ type ValidatorRegistrySyncer struct { func NewValidatorRegistrySyncer( contract *validatorRegistryBindings.Validatorregistry, db *pgxpool.Pool, + ethClient *ethclient.Client, txMapper metrics.TxMapper, syncStartBlockNumber uint64, ) *ValidatorRegistrySyncer { @@ -33,6 +37,7 @@ func NewValidatorRegistrySyncer( contract: contract, db: db, dbQuery: data.New(db), + ethClient: ethClient, txMapper: txMapper, syncStartBlockNumber: syncStartBlockNumber, } @@ -75,17 +80,30 @@ func (ets *ValidatorRegistrySyncer) syncRange( return err } + header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) + if err != nil { + return errors.Wrap(err, "failed to get execution block header by number") + } for _, event := range events { err := ets.txMapper.AddValidatorRegistryEvent(ctx, event) if err != nil { log.Err(err).Msg("err adding validator registry updated event") - return nil + return err } log.Info(). Uint64("block", event.Raw.BlockNumber). Msg("new validator registry updated message") } + err = ets.dbQuery.CreateValidatorRegistryEventsSyncedUntil(ctx, data.CreateValidatorRegistryEventsSyncedUntilParams{ + BlockNumber: int64(end), + BlockHash: header.Hash().Bytes(), + }) + if err != nil { + log.Err(err).Msg("err adding validator registry event sync until") + return err + } + log.Info(). Uint64("start-block", start). Uint64("end-block", end). diff --git a/internal/watcher/watcher.go b/internal/watcher/watcher.go index 9150caf..9a14b67 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -118,7 +118,7 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { } transactionSubmittedSyncer := syncer.NewTransactionSubmittedSyncer(sequencerContract, w.db, ethClient, txMapper, blockNumber) - validatorRegistrySyncer := syncer.NewValidatorRegistrySyncer(validatorRegistryContract, w.db, txMapper, ValidatorRegistryDeploymentBlockNumber) + validatorRegistrySyncer := syncer.NewValidatorRegistrySyncer(validatorRegistryContract, w.db, ethClient, txMapper, ValidatorRegistryDeploymentBlockNumber) blocksWatcher := NewBlocksWatcher(w.config, ethClient, txMapper, transactionSubmittedSyncer, validatorRegistrySyncer) p2pMsgsWatcher := NewP2PMsgsWatcherWatcher(w.config, decryptionDataChannel, keyShareChannel, blocksWatcher) From 9e1539dd976efa5cf8e7898e183f2f9aaeb06853 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Tue, 1 Apr 2025 14:54:07 +0200 Subject: [PATCH 10/13] add optional tx param --- internal/metrics/tx_mapper_db.go | 12 +++++++++--- internal/metrics/types.go | 3 ++- internal/syncer/validator_registry.go | 16 ++++++++++++++-- 3 files changed, 25 insertions(+), 6 deletions(-) diff --git a/internal/metrics/tx_mapper_db.go b/internal/metrics/tx_mapper_db.go index d23ac2f..fbe8a5d 100644 --- a/internal/metrics/tx_mapper_db.go +++ b/internal/metrics/tx_mapper_db.go @@ -209,7 +209,7 @@ func (tm *TxMapperDB) QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx return data.BlockNumber, nil } -func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validatorRegistryBindings.ValidatorregistryUpdated) error { +func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, tx pgx.Tx, vr *validatorRegistryBindings.ValidatorregistryUpdated) error { regMessage := &validatorregistry.AggregateRegistrationMessage{} err := regMessage.Unmarshal(vr.Message) if err != nil { @@ -221,8 +221,14 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat return err } + q := tm.dbQuery + if tx != nil { + // Use transaction if available + q = tm.dbQuery.WithTx(tx) + } + for validatorID, validatorData := range validatorIDtoValidity { - err := tm.dbQuery.CreateValidatorRegistryMessage(ctx, data.CreateValidatorRegistryMessageParams{ + err := q.CreateValidatorRegistryMessage(ctx, data.CreateValidatorRegistryMessageParams{ Version: dbTypes.Uint64ToPgTypeInt8(uint64(regMessage.Version)), ChainID: dbTypes.Uint64ToPgTypeInt8(regMessage.ChainID), ValidatorRegistryAddress: regMessage.ValidatorRegistryAddress.Bytes(), @@ -241,7 +247,7 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat if validatorData.validatorValidity == data.ValidatorRegistrationValidityValid && validatorData.validatorStatus != "" { - err := tm.dbQuery.CreateValidatorStatus(ctx, data.CreateValidatorStatusParams{ + err := q.CreateValidatorStatus(ctx, data.CreateValidatorStatusParams{ ValidatorIndex: dbTypes.Int64ToPgTypeInt8(validatorID), Status: validatorData.validatorStatus, }) diff --git a/internal/metrics/types.go b/internal/metrics/types.go index a13ae0b..6bb3fc3 100644 --- a/internal/metrics/types.go +++ b/internal/metrics/types.go @@ -3,6 +3,7 @@ package metrics import ( "context" + "github.com/jackc/pgx/v5" sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" "github.com/shutter-network/observer/internal/data" @@ -61,7 +62,7 @@ type TxMapper interface { b *data.Block, ) error QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx context.Context) (int64, error) - AddValidatorRegistryEvent(ctx context.Context, vr *validatorRegistryBindings.ValidatorregistryUpdated) error + AddValidatorRegistryEvent(ctx context.Context, tx pgx.Tx, vr *validatorRegistryBindings.ValidatorregistryUpdated) error UpdateValidatorStatus(ctx context.Context) error AddProposerDuties(ctx context.Context, epoch uint64) error } diff --git a/internal/syncer/validator_registry.go b/internal/syncer/validator_registry.go index b845e65..2dc58f8 100644 --- a/internal/syncer/validator_registry.go +++ b/internal/syncer/validator_registry.go @@ -84,8 +84,15 @@ func (ets *ValidatorRegistrySyncer) syncRange( if err != nil { return errors.Wrap(err, "failed to get execution block header by number") } + tx, err := ets.db.Begin(ctx) + if err != nil { + return err + } + defer tx.Rollback(ctx) + qtx := ets.dbQuery.WithTx(tx) + for _, event := range events { - err := ets.txMapper.AddValidatorRegistryEvent(ctx, event) + err := ets.txMapper.AddValidatorRegistryEvent(ctx, tx, event) if err != nil { log.Err(err).Msg("err adding validator registry updated event") return err @@ -95,7 +102,7 @@ func (ets *ValidatorRegistrySyncer) syncRange( Msg("new validator registry updated message") } - err = ets.dbQuery.CreateValidatorRegistryEventsSyncedUntil(ctx, data.CreateValidatorRegistryEventsSyncedUntilParams{ + err = qtx.CreateValidatorRegistryEventsSyncedUntil(ctx, data.CreateValidatorRegistryEventsSyncedUntilParams{ BlockNumber: int64(end), BlockHash: header.Hash().Bytes(), }) @@ -103,6 +110,11 @@ func (ets *ValidatorRegistrySyncer) syncRange( log.Err(err).Msg("err adding validator registry event sync until") return err } + err = tx.Commit(ctx) + if err != nil { + log.Err(err).Msg("unable to commit db transaction") + return err + } log.Info(). Uint64("start-block", start). From b35fb347cb3161c0a8ebe46206ba08dfee9cd93e Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Tue, 1 Apr 2025 14:59:23 +0200 Subject: [PATCH 11/13] txify tx submitted event --- internal/metrics/tx_mapper_db.go | 9 +++++++-- internal/metrics/types.go | 2 +- internal/syncer/transaction_submitted_syncer.go | 16 ++++++++++++++-- 3 files changed, 22 insertions(+), 5 deletions(-) diff --git a/internal/metrics/tx_mapper_db.go b/internal/metrics/tx_mapper_db.go index fbe8a5d..892890c 100644 --- a/internal/metrics/tx_mapper_db.go +++ b/internal/metrics/tx_mapper_db.go @@ -69,8 +69,13 @@ func NewTxMapperDB( } } -func (tm *TxMapperDB) AddTransactionSubmittedEvent(ctx context.Context, st *sequencerBindings.SequencerTransactionSubmitted) error { - err := tm.dbQuery.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ +func (tm *TxMapperDB) AddTransactionSubmittedEvent(ctx context.Context, tx pgx.Tx, st *sequencerBindings.SequencerTransactionSubmitted) error { + q := tm.dbQuery + if tx != nil { + // Use transaction if available + q = tm.dbQuery.WithTx(tx) + } + err := q.CreateTransactionSubmittedEvent(ctx, data.CreateTransactionSubmittedEventParams{ EventBlockHash: st.Raw.BlockHash.Bytes(), EventBlockNumber: int64(st.Raw.BlockNumber), EventTxIndex: int64(st.Raw.TxIndex), diff --git a/internal/metrics/types.go b/internal/metrics/types.go index 6bb3fc3..42b228a 100644 --- a/internal/metrics/types.go +++ b/internal/metrics/types.go @@ -51,7 +51,7 @@ type TxExecution struct { } type TxMapper interface { - AddTransactionSubmittedEvent(ctx context.Context, st *sequencerBindings.SequencerTransactionSubmitted) error + AddTransactionSubmittedEvent(ctx context.Context, tx pgx.Tx, st *sequencerBindings.SequencerTransactionSubmitted) error AddDecryptionKeysAndMessages( ctx context.Context, dkam *DecKeysAndMessages, diff --git a/internal/syncer/transaction_submitted_syncer.go b/internal/syncer/transaction_submitted_syncer.go index af17cf5..c111c22 100644 --- a/internal/syncer/transaction_submitted_syncer.go +++ b/internal/syncer/transaction_submitted_syncer.go @@ -88,8 +88,14 @@ func (ets *TransactionSubmittedSyncer) syncRange( if err != nil { return errors.Wrap(err, "failed to get execution block header by number") } + tx, err := ets.db.Begin(ctx) + if err != nil { + return err + } + defer tx.Rollback(ctx) + qtx := ets.dbQuery.WithTx(tx) for _, event := range events { - err := ets.txMapper.AddTransactionSubmittedEvent(ctx, event) + err := ets.txMapper.AddTransactionSubmittedEvent(ctx, tx, event) if err != nil { log.Err(err).Msg("err adding transaction submitted event") return err @@ -99,7 +105,7 @@ func (ets *TransactionSubmittedSyncer) syncRange( Hex("encrypted transaction (hex)", event.EncryptedTransaction). Msg("new encrypted transaction") } - err = ets.dbQuery.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ + err = qtx.CreateTransactionSubmittedEventsSyncedUntil(ctx, data.CreateTransactionSubmittedEventsSyncedUntilParams{ BlockNumber: int64(end), BlockHash: header.Hash().Bytes(), }) @@ -107,6 +113,12 @@ func (ets *TransactionSubmittedSyncer) syncRange( log.Err(err).Msg("err adding transaction submit event sync until") return err } + + err = tx.Commit(ctx) + if err != nil { + log.Err(err).Msg("unable to commit db transaction") + return err + } log.Info(). Uint64("start-block", start). Uint64("end-block", end). From 6d7ad3af82168c4d4d275890c617797ed9d69ac4 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Fri, 4 Apr 2025 11:40:04 +0200 Subject: [PATCH 12/13] fix tests --- tests/transaction_test.go | 26 ++++++++++++++++---------- tests/tx_mapper_test.go | 26 +++++++++++++++----------- tests/validator_test.go | 4 ++-- 3 files changed, 33 insertions(+), 23 deletions(-) diff --git a/tests/transaction_test.go b/tests/transaction_test.go index 06c21ff..c7bfb94 100644 --- a/tests/transaction_test.go +++ b/tests/transaction_test.go @@ -5,6 +5,9 @@ import ( cryptoRand "crypto/rand" "math/rand" + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/core/types" + "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" "github.com/shutter-network/observer/internal/data" "github.com/shutter-network/observer/internal/metrics" ) @@ -134,18 +137,21 @@ func (s *TestMetricsSuite) TestAddFullTransaction() { eventTxHash, err := generateRandomBytes(32) s.Require().NoError(err) - err = s.txMapperDB.AddTransactionSubmittedEvent(ctx, &data.TransactionSubmittedEvent{ - EventBlockHash: eventBlockHash, - EventBlockNumber: eventBlockNumber, - EventTxIndex: eventTxIndex, - EventLogIndex: eventLogIndex, - Eon: eon, - TxIndex: txIndex, - IdentityPrefix: identityPrefix, - Sender: sender, + s.txMapperDB.AddTransactionSubmittedEvent(ctx, nil, &sequencer.SequencerTransactionSubmitted{ + Eon: uint64(eon), + TxIndex: uint64(txIndex), + IdentityPrefix: [32]byte(identityPrefix), + Sender: common.Address(sender), EncryptedTransaction: ectx, - EventTxHash: eventTxHash, + Raw: types.Log{ + BlockHash: common.Hash(eventBlockHash), + BlockNumber: uint64(eventBlockNumber), + TxIndex: uint(eventTxIndex), + Index: uint(eventLogIndex), + TxHash: common.Hash(eventTxHash), + }, }) + s.Require().NoError(err) err = s.txMapperDB.AddKeyShare(ctx, &data.DecryptionKeyShare{ diff --git a/tests/tx_mapper_test.go b/tests/tx_mapper_test.go index 7188742..4cf7741 100644 --- a/tests/tx_mapper_test.go +++ b/tests/tx_mapper_test.go @@ -7,7 +7,9 @@ import ( cryptorand "crypto/rand" - "github.com/shutter-network/observer/internal/data" + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/core/types" + "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" "github.com/shutter-network/observer/internal/metrics" "github.com/shutter-network/shutter/shlib/shcrypto" blst "github.com/supranational/blst/bindings/go" @@ -100,17 +102,19 @@ func (s *TestMetricsSuite) TestAddTransactionSubmittedEventAndDecryptionData() { eventTxHash, err := generateRandomBytes(32) s.Require().NoError(err) - err = s.txMapperDB.AddTransactionSubmittedEvent(ctx, &data.TransactionSubmittedEvent{ - EventBlockHash: eventBlockHash, - EventBlockNumber: eventBlockNumber, - EventTxIndex: eventTxIndex, - EventLogIndex: eventLogIndex, - Eon: eon, - TxIndex: txIndex, - IdentityPrefix: identityPrefix, - Sender: sender, + err = s.txMapperDB.AddTransactionSubmittedEvent(ctx, nil, &sequencer.SequencerTransactionSubmitted{ + Eon: uint64(eon), + TxIndex: uint64(txIndex), + IdentityPrefix: [32]byte(identityPrefix), + Sender: common.Address(sender), EncryptedTransaction: encrypedTxBytes, - EventTxHash: eventTxHash, + Raw: types.Log{ + BlockHash: common.Hash(eventBlockHash), + BlockNumber: uint64(eventBlockNumber), + TxIndex: uint(eventTxIndex), + Index: uint(eventLogIndex), + TxHash: common.Hash(eventTxHash), + }, }) s.Require().NoError(err) diff --git a/tests/validator_test.go b/tests/validator_test.go index 25b3558..66409df 100644 --- a/tests/validator_test.go +++ b/tests/validator_test.go @@ -68,7 +68,7 @@ func (s *TestMetricsSuite) TestAggregateValidatorRegistrationMessage() { }, } s.txMapperDB = metrics.NewTxMapperDB(ctx, s.testDB.DbInstance, &observerCommon.Config{ValidatorRegistryContractAddress: ValidatorRegistryContract}, ðclient.Client{}, cl, 2, rand.Uint64(), rand.Uint64()) - err = s.txMapperDB.AddValidatorRegistryEvent(ctx, &event) + err = s.txMapperDB.AddValidatorRegistryEvent(ctx, nil, &event) s.Require().NoError(err) currentNonce, err := s.dbQuery.QueryValidatorRegistrationMessageNonceBefore(ctx, data.QueryValidatorRegistrationMessageNonceBeforeParams{ @@ -119,7 +119,7 @@ func (s *TestMetricsSuite) TestLegacyValidatorRegistrationMessage() { }, } s.txMapperDB = metrics.NewTxMapperDB(ctx, s.testDB.DbInstance, &observerCommon.Config{ValidatorRegistryContractAddress: ValidatorRegistryContract}, ðclient.Client{}, cl, 2, rand.Uint64(), rand.Uint64()) - err = s.txMapperDB.AddValidatorRegistryEvent(ctx, &event) + err = s.txMapperDB.AddValidatorRegistryEvent(ctx, nil, &event) s.Require().NoError(err) currentNonce, err := s.dbQuery.QueryValidatorRegistrationMessageNonceBefore(ctx, data.QueryValidatorRegistrationMessageNonceBeforeParams{ From 9216b8747c24ca7dd84068b1ea036f9f27742989 Mon Sep 17 00:00:00 2001 From: faheelsattar Date: Mon, 7 Apr 2025 10:16:01 +0200 Subject: [PATCH 13/13] replace err fmt --- internal/syncer/transaction_submitted_syncer.go | 10 +++++----- internal/syncer/validator_registry.go | 10 +++++----- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/internal/syncer/transaction_submitted_syncer.go b/internal/syncer/transaction_submitted_syncer.go index c111c22..878978d 100644 --- a/internal/syncer/transaction_submitted_syncer.go +++ b/internal/syncer/transaction_submitted_syncer.go @@ -2,6 +2,7 @@ package syncer import ( "context" + "fmt" "math/big" "github.com/ethereum/go-ethereum/accounts/abi/bind" @@ -9,7 +10,6 @@ import ( "github.com/ethereum/go-ethereum/ethclient" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" - "github.com/pkg/errors" "github.com/rs/zerolog/log" sequencerBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/sequencer" "github.com/shutter-network/observer/internal/data" @@ -52,7 +52,7 @@ func (ets *TransactionSubmittedSyncer) Sync(ctx context.Context, header *types.H // TODO: handle reorgs syncedUntil, err := ets.dbQuery.QueryTransactionSubmittedEventsSyncedUntil(ctx) if err != nil && err != pgx.ErrNoRows { - return errors.Wrap(err, "failed to query transaction submitted events sync status") + return fmt.Errorf("failed to query transaction submitted events sync status, %v", err) } var start uint64 if err == pgx.ErrNoRows { @@ -86,7 +86,7 @@ func (ets *TransactionSubmittedSyncer) syncRange( } header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) if err != nil { - return errors.Wrap(err, "failed to get execution block header by number") + return fmt.Errorf("failed to get execution block header by number, %v", err) } tx, err := ets.db.Begin(ctx) if err != nil { @@ -139,14 +139,14 @@ func (s *TransactionSubmittedSyncer) fetchEvents( } it, err := s.contract.SequencerFilterer.FilterTransactionSubmitted(&opts) if err != nil { - return nil, errors.Wrap(err, "failed to query transaction submitted events") + return nil, fmt.Errorf("failed to query transaction submitted events, %v", err) } events := []*sequencerBindings.SequencerTransactionSubmitted{} for it.Next() { events = append(events, it.Event) } if it.Error() != nil { - return nil, errors.Wrap(it.Error(), "failed to iterate query transaction submitted events") + return nil, fmt.Errorf("failed to iterate query transaction submitted events, %v", it.Error()) } return events, nil } diff --git a/internal/syncer/validator_registry.go b/internal/syncer/validator_registry.go index 2dc58f8..496c6c8 100644 --- a/internal/syncer/validator_registry.go +++ b/internal/syncer/validator_registry.go @@ -2,6 +2,7 @@ package syncer import ( "context" + "fmt" "math/big" "github.com/ethereum/go-ethereum/accounts/abi/bind" @@ -9,7 +10,6 @@ import ( "github.com/ethereum/go-ethereum/ethclient" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" - "github.com/pkg/errors" "github.com/rs/zerolog/log" validatorRegistryBindings "github.com/shutter-network/gnosh-contracts/gnoshcontracts/validatorregistry" "github.com/shutter-network/observer/internal/data" @@ -47,7 +47,7 @@ func (vts *ValidatorRegistrySyncer) Sync(ctx context.Context, header *types.Head // TODO: handle reorgs syncedUntil, err := vts.dbQuery.QueryValidatorRegistryEventsSyncedUntil(ctx) if err != nil && err != pgx.ErrNoRows { - return errors.Wrap(err, "failed to query validator registry sync status") + return fmt.Errorf("failed to query validator registry sync status, %v", err) } var start uint64 if err == pgx.ErrNoRows { @@ -82,7 +82,7 @@ func (ets *ValidatorRegistrySyncer) syncRange( header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) if err != nil { - return errors.Wrap(err, "failed to get execution block header by number") + return fmt.Errorf("failed to get execution block header by number, %v", err) } tx, err := ets.db.Begin(ctx) if err != nil { @@ -136,14 +136,14 @@ func (s *ValidatorRegistrySyncer) fetchEvents( } it, err := s.contract.ValidatorregistryFilterer.FilterUpdated(&opts) if err != nil { - return nil, errors.Wrap(err, "failed to query validator registry updated events") + return nil, fmt.Errorf("failed to query validator registry updated events, %v", err) } events := []*validatorRegistryBindings.ValidatorregistryUpdated{} for it.Next() { events = append(events, it.Event) } if it.Error() != nil { - return nil, errors.Wrap(it.Error(), "failed to iterate query validator registry updated events") + return nil, fmt.Errorf("failed to iterate query validator registry updated events, %v", it.Error()) } return events, nil }