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 } 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 863e86f..892890c 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,18 +69,23 @@ 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, 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), + 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 @@ -201,23 +207,16 @@ 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 { - tx, err := tm.db.Begin(ctx) - if err != nil { - return err - } - defer tx.Rollback(ctx) - qtx := tm.dbQuery.WithTx(tx) - +func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, tx pgx.Tx, vr *validatorRegistryBindings.ValidatorregistryUpdated) error { 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 { @@ -227,8 +226,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 := qtx.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(), @@ -247,7 +252,7 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat if validatorData.validatorValidity == data.ValidatorRegistrationValidityValid && validatorData.validatorStatus != "" { - err := qtx.CreateValidatorStatus(ctx, data.CreateValidatorStatusParams{ + err := q.CreateValidatorStatus(ctx, data.CreateValidatorStatusParams{ ValidatorIndex: dbTypes.Int64ToPgTypeInt8(validatorID), Status: validatorData.validatorStatus, }) @@ -257,12 +262,7 @@ func (tm *TxMapperDB) AddValidatorRegistryEvent(ctx context.Context, vr *validat } } } - - err = qtx.CreateValidatorRegistryEventsSyncedUntil(ctx, int64(vr.Raw.BlockNumber)) - if err != nil { - return err - } - return tx.Commit(ctx) + return nil } func (tm *TxMapperDB) UpdateValidatorStatus(ctx context.Context) error { @@ -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: } diff --git a/internal/metrics/types.go b/internal/metrics/types.go index 383b956..42b228a 100644 --- a/internal/metrics/types.go +++ b/internal/metrics/types.go @@ -3,6 +3,8 @@ 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" ) @@ -49,7 +51,7 @@ type TxExecution struct { } type TxMapper interface { - AddTransactionSubmittedEvent(ctx context.Context, tse *data.TransactionSubmittedEvent) error + AddTransactionSubmittedEvent(ctx context.Context, tx pgx.Tx, st *sequencerBindings.SequencerTransactionSubmitted) error AddDecryptionKeysAndMessages( ctx context.Context, dkam *DecKeysAndMessages, @@ -60,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/transaction_submitted_syncer.go b/internal/syncer/transaction_submitted_syncer.go new file mode 100644 index 0000000..878978d --- /dev/null +++ b/internal/syncer/transaction_submitted_syncer.go @@ -0,0 +1,152 @@ +package syncer + +import ( + "context" + "fmt" + "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/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" +) + +const ( + AssumedReorgDepth = 10 + maxRequestBlockRange = 10_000 +) + +type TransactionSubmittedSyncer struct { + contract *sequencerBindings.Sequencer + db *pgxpool.Pool + dbQuery *data.Queries + ethClient *ethclient.Client + txMapper metrics.TxMapper + syncStartBlockNumber uint64 +} + +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 { + return fmt.Errorf("failed to query transaction submitted events sync status, %v", err) + } + 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 *TransactionSubmittedSyncer) 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 fmt.Errorf("failed to get execution block header by number, %v", err) + } + 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, tx, event) + if err != nil { + log.Err(err).Msg("err adding transaction submitted event") + return err + } + 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("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). + Int("num-inserted-events", len(events)). + Msg("synced sequencer contract") + return nil +} + +func (s *TransactionSubmittedSyncer) 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, 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, 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 new file mode 100644 index 0000000..496c6c8 --- /dev/null +++ b/internal/syncer/validator_registry.go @@ -0,0 +1,149 @@ +package syncer + +import ( + "context" + "fmt" + "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/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 + db *pgxpool.Pool + dbQuery *data.Queries + ethClient *ethclient.Client + txMapper metrics.TxMapper + syncStartBlockNumber uint64 +} + +func NewValidatorRegistrySyncer( + contract *validatorRegistryBindings.Validatorregistry, + db *pgxpool.Pool, + ethClient *ethclient.Client, + txMapper metrics.TxMapper, + syncStartBlockNumber uint64, +) *ValidatorRegistrySyncer { + return &ValidatorRegistrySyncer{ + contract: contract, + db: db, + dbQuery: data.New(db), + ethClient: ethClient, + txMapper: txMapper, + syncStartBlockNumber: syncStartBlockNumber, + } +} + +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 fmt.Errorf("failed to query validator registry sync status, %v", err) + } + 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 + } + + header, err := ets.ethClient.HeaderByNumber(ctx, new(big.Int).SetUint64(end)) + if err != nil { + return fmt.Errorf("failed to get execution block header by number, %v", err) + } + 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, tx, event) + if err != nil { + log.Err(err).Msg("err adding validator registry updated event") + return err + } + log.Info(). + Uint64("block", event.Raw.BlockNumber). + Msg("new validator registry updated message") + } + + err = qtx.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 + } + 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). + 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, 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, fmt.Errorf("failed to iterate query validator registry updated events, %v", it.Error()) + } + return events, nil +} diff --git a/internal/watcher/blocks.go b/internal/watcher/blocks.go index aa0c4d3..fb9ff45 100644 --- a/internal/watcher/blocks.go +++ b/internal/watcher/blocks.go @@ -2,33 +2,47 @@ 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/observer/internal/syncer" "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 + txMapper metrics.TxMapper + transactionSubmittedSyncer *syncer.TransactionSubmittedSyncer + validatorRegistrySyncer *syncer.ValidatorRegistrySyncer -type BlockReceivedEvent struct { - Header *types.Header - Time time.Time + recentBlocksMux sync.Mutex + recentBlocks map[uint64]*types.Header + mostRecentBlock uint64 } -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, + transactionSubmittedSyncer *syncer.TransactionSubmittedSyncer, + validatorRegistrySyncer *syncer.ValidatorRegistrySyncer, +) *BlocksWatcher { return &BlocksWatcher{ - config: config, - blocksChannel: blocksChannel, - blocksChannelForProposerDuties: blocksChannelForProposerDuties, - ethClient: ethClient, + config: config, + ethClient: ethClient, + txMapper: txMapper, + transactionSubmittedSyncer: transactionSubmittedSyncer, + validatorRegistrySyncer: validatorRegistrySyncer, + recentBlocksMux: sync.Mutex{}, + recentBlocks: make(map[uint64]*types.Header), + mostRecentBlock: 0, } } @@ -45,16 +59,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 +76,92 @@ 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) + + 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 + 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/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/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/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 f274387..9a14b67 100644 --- a/internal/watcher/watcher.go +++ b/internal/watcher/watcher.go @@ -7,18 +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" "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/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/beaconapiclient" "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" ) @@ -62,13 +62,8 @@ func New( } 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,65 +102,33 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { SlotDuration, ) - blocksWatcher := NewBlocksWatcher(w.config, blocksChannel, blocksChannelForProposerDuties, ethClient) - encryptionTxWatcher := NewEncryptedTxWatcher(w.config, txSubmittedEventChannel, ethClient) + sequencerContract, err := sequencerBindings.NewSequencer(ethCommon.HexToAddress(w.config.SequencerContractAddress), ethClient) + if err != nil { + return err + } - blockNumber, err := txMapper.QueryBlockNumberFromValidatorRegistryEventsSyncedUntil(ctx) + validatorRegistryContract, err := validatorRegistryBindings.NewValidatorregistry(ethCommon.HexToAddress(w.config.ValidatorRegistryContractAddress), ethClient) if err != nil { - if err == pgx.ErrNoRows { - blockNumber = int64(ValidatorRegistryDeploymentBlockNumber) - } else { - return err - } + return err } - validatorRegisterWatcher := NewValidatorRegistryWatcher(w.config, validatorRegistryChannel, ethClient, blockNumber) + blockNumber, err := ethClient.BlockNumber(ctx) + if err != nil { + return err + } - p2pMsgsWatcher := NewP2PMsgsWatcherWatcher(w.config, blocksChannel, decryptionDataChannel, keyShareChannel, txMapper) - if err := runner.StartService(blocksWatcher, encryptionTxWatcher, p2pMsgsWatcher, validatorRegisterWatcher); err != nil { + transactionSubmittedSyncer := syncer.NewTransactionSubmittedSyncer(sequencerContract, w.db, ethClient, txMapper, blockNumber) + validatorRegistrySyncer := syncer.NewValidatorRegistrySyncer(validatorRegistryContract, w.db, ethClient, 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 } 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[:], - 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( @@ -206,12 +169,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() } 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 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{