Improve I/O error propagation
This commit is contained in:
@@ -18,6 +18,7 @@ package ethereum
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"github.com/sirupsen/logrus"
|
||||
|
||||
"github.com/ethereum/go-ethereum/core/types"
|
||||
"github.com/ethereum/go-ethereum/ethdb"
|
||||
@@ -37,6 +38,7 @@ func CreateDatabase(config DatabaseConfig) (Database, error) {
|
||||
case Level:
|
||||
levelDBConnection, err := ethdb.NewLDBDatabase(config.Path, 128, 1024)
|
||||
if err != nil {
|
||||
logrus.Error("CreateDatabase: error connecting to new LDBD: ", err)
|
||||
return nil, err
|
||||
}
|
||||
levelDBReader := level.NewLevelDatabaseReader(levelDBConnection)
|
||||
|
||||
@@ -43,12 +43,16 @@ func NewBlockRepository(database *postgres.DB) *BlockRepository {
|
||||
return &BlockRepository{database: database}
|
||||
}
|
||||
|
||||
func (blockRepository BlockRepository) SetBlocksStatus(chainHead int64) {
|
||||
func (blockRepository BlockRepository) SetBlocksStatus(chainHead int64) error {
|
||||
cutoff := chainHead - blocksFromHeadBeforeFinal
|
||||
blockRepository.database.Exec(`
|
||||
_, err := blockRepository.database.Exec(`
|
||||
UPDATE blocks SET is_final = TRUE
|
||||
WHERE is_final = FALSE AND number < $1`,
|
||||
cutoff)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (blockRepository BlockRepository) CreateOrUpdateBlock(block core.Block) (int64, error) {
|
||||
@@ -70,7 +74,7 @@ func (blockRepository BlockRepository) CreateOrUpdateBlock(block core.Block) (in
|
||||
|
||||
func (blockRepository BlockRepository) MissingBlockNumbers(startingBlockNumber int64, highestBlockNumber int64, nodeId string) []int64 {
|
||||
numbers := make([]int64, 0)
|
||||
blockRepository.database.Select(&numbers,
|
||||
err := blockRepository.database.Select(&numbers,
|
||||
`SELECT all_block_numbers
|
||||
FROM (
|
||||
SELECT generate_series($1::INT, $2::INT) AS all_block_numbers) series
|
||||
@@ -79,6 +83,9 @@ func (blockRepository BlockRepository) MissingBlockNumbers(startingBlockNumber i
|
||||
) `,
|
||||
startingBlockNumber,
|
||||
highestBlockNumber, nodeId)
|
||||
if err != nil {
|
||||
log.Error("MissingBlockNumbers: error getting blocks: ", err)
|
||||
}
|
||||
return numbers
|
||||
}
|
||||
|
||||
@@ -108,6 +115,7 @@ func (blockRepository BlockRepository) GetBlock(blockNumber int64) (core.Block,
|
||||
case sql.ErrNoRows:
|
||||
return core.Block{}, datastore.ErrBlockDoesNotExist(blockNumber)
|
||||
default:
|
||||
log.Error("GetBlock: error loading blocks: ", err)
|
||||
return savedBlock, err
|
||||
}
|
||||
}
|
||||
@@ -202,6 +210,7 @@ func (blockRepository BlockRepository) createReceipt(tx *sql.Tx, blockId int64,
|
||||
RETURNING id`,
|
||||
receipt.ContractAddress, receipt.TxHash, receipt.CumulativeGasUsed, receipt.GasUsed, receipt.StateRoot, receipt.Status, blockId).Scan(&receiptId)
|
||||
if err != nil {
|
||||
log.Error("createReceipt: error inserting receipt: ", err)
|
||||
return receiptId, err
|
||||
}
|
||||
return receiptId, nil
|
||||
@@ -256,6 +265,7 @@ func (blockRepository BlockRepository) loadBlock(blockRows *sqlx.Row) (core.Bloc
|
||||
var block b
|
||||
err := blockRows.StructScan(&block)
|
||||
if err != nil {
|
||||
log.Error("loadBlock: error loading block: ", err)
|
||||
return core.Block{}, err
|
||||
}
|
||||
transactionRows, err := blockRepository.database.Queryx(`
|
||||
@@ -271,6 +281,7 @@ func (blockRepository BlockRepository) loadBlock(blockRows *sqlx.Row) (core.Bloc
|
||||
WHERE block_id = $1
|
||||
ORDER BY hash`, block.ID)
|
||||
if err != nil {
|
||||
log.Error("loadBlock: error fetting transactions: ", err)
|
||||
return core.Block{}, err
|
||||
}
|
||||
block.Transactions = blockRepository.LoadTransactions(transactionRows)
|
||||
|
||||
@@ -47,14 +47,17 @@ func (contractRepository ContractRepository) CreateContract(contract core.Contra
|
||||
return nil
|
||||
}
|
||||
|
||||
func (contractRepository ContractRepository) ContractExists(contractHash string) bool {
|
||||
func (contractRepository ContractRepository) ContractExists(contractHash string) (bool, error) {
|
||||
var exists bool
|
||||
contractRepository.DB.QueryRow(
|
||||
err := contractRepository.DB.QueryRow(
|
||||
`SELECT exists(
|
||||
SELECT 1
|
||||
FROM watched_contracts
|
||||
WHERE contract_hash = $1)`, contractHash).Scan(&exists)
|
||||
return exists
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return exists, nil
|
||||
}
|
||||
|
||||
func (contractRepository ContractRepository) GetContract(contractHash string) (core.Contract, error) {
|
||||
@@ -66,12 +69,15 @@ func (contractRepository ContractRepository) GetContract(contractHash string) (c
|
||||
if err == sql.ErrNoRows {
|
||||
return core.Contract{}, datastore.ErrContractDoesNotExist(contractHash)
|
||||
}
|
||||
savedContract := contractRepository.addTransactions(core.Contract{Hash: hash, Abi: abi})
|
||||
savedContract, err := contractRepository.addTransactions(core.Contract{Hash: hash, Abi: abi})
|
||||
if err != nil {
|
||||
return core.Contract{}, err
|
||||
}
|
||||
return savedContract, nil
|
||||
}
|
||||
|
||||
func (contractRepository ContractRepository) addTransactions(contract core.Contract) core.Contract {
|
||||
transactionRows, _ := contractRepository.DB.Queryx(`
|
||||
func (contractRepository ContractRepository) addTransactions(contract core.Contract) (core.Contract, error) {
|
||||
transactionRows, err := contractRepository.DB.Queryx(`
|
||||
SELECT hash,
|
||||
nonce,
|
||||
tx_to,
|
||||
@@ -83,8 +89,11 @@ func (contractRepository ContractRepository) addTransactions(contract core.Contr
|
||||
FROM transactions
|
||||
WHERE tx_to = $1
|
||||
ORDER BY block_id DESC`, contract.Hash)
|
||||
if err != nil {
|
||||
return core.Contract{}, err
|
||||
}
|
||||
blockRepository := &BlockRepository{contractRepository.DB}
|
||||
transactions := blockRepository.LoadTransactions(transactionRows)
|
||||
savedContract := core.Contract{Hash: contract.Hash, Transactions: transactions, Abi: contract.Abi}
|
||||
return savedContract
|
||||
return savedContract, nil
|
||||
}
|
||||
|
||||
@@ -42,6 +42,7 @@ func (repository HeaderRepository) CreateOrUpdateHeader(header core.Header) (int
|
||||
if headerDoesNotExist(err) {
|
||||
return repository.insertHeader(header)
|
||||
}
|
||||
log.Error("CreateOrUpdateHeader: error getting header hash: ", err)
|
||||
return 0, err
|
||||
}
|
||||
if headerMustBeReplaced(hash, header) {
|
||||
@@ -54,6 +55,7 @@ func (repository HeaderRepository) GetHeader(blockNumber int64) (core.Header, er
|
||||
var header core.Header
|
||||
err := repository.database.Get(&header, `SELECT id, block_number, hash, raw, block_timestamp FROM headers WHERE block_number = $1 AND eth_node_fingerprint = $2`,
|
||||
blockNumber, repository.database.Node.ID)
|
||||
log.Error("GetHeader: error getting headers: ", err)
|
||||
return header, err
|
||||
}
|
||||
|
||||
@@ -74,18 +76,6 @@ func (repository HeaderRepository) MissingBlockNumbers(startingBlockNumber, endi
|
||||
return numbers, nil
|
||||
}
|
||||
|
||||
func (repository HeaderRepository) HeaderExists(blockNumber int64) (bool, error) {
|
||||
_, err := repository.GetHeader(blockNumber)
|
||||
if err != nil {
|
||||
if headerDoesNotExist(err) {
|
||||
return false, nil
|
||||
}
|
||||
return false, err
|
||||
}
|
||||
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func headerMustBeReplaced(hash string, header core.Header) bool {
|
||||
return hash != header.Hash
|
||||
}
|
||||
@@ -98,6 +88,7 @@ func (repository HeaderRepository) getHeaderHash(header core.Header) (string, er
|
||||
var hash string
|
||||
err := repository.database.Get(&hash, `SELECT hash FROM headers WHERE block_number = $1 AND eth_node_fingerprint = $2`,
|
||||
header.BlockNumber, repository.database.Node.ID)
|
||||
log.Error("getHeaderHash: error getting headers: ", err)
|
||||
return hash, err
|
||||
}
|
||||
|
||||
@@ -106,6 +97,9 @@ func (repository HeaderRepository) insertHeader(header core.Header) (int64, erro
|
||||
err := repository.database.QueryRowx(
|
||||
`INSERT INTO public.headers (block_number, hash, block_timestamp, raw, eth_node_id, eth_node_fingerprint) VALUES ($1, $2, $3::NUMERIC, $4, $5, $6) RETURNING id`,
|
||||
header.BlockNumber, header.Hash, header.Timestamp, header.Raw, repository.database.NodeID, repository.database.Node.ID).Scan(&headerId)
|
||||
if err != nil {
|
||||
log.Error("insertHeader: error inserting header: ", err)
|
||||
}
|
||||
return headerId, err
|
||||
}
|
||||
|
||||
@@ -113,6 +107,7 @@ func (repository HeaderRepository) replaceHeader(header core.Header) (int64, err
|
||||
_, err := repository.database.Exec(`DELETE FROM headers WHERE block_number = $1 AND eth_node_fingerprint = $2`,
|
||||
header.BlockNumber, repository.database.Node.ID)
|
||||
if err != nil {
|
||||
log.Error("replaceHeader: error deleting headers: ", err)
|
||||
return 0, err
|
||||
}
|
||||
return repository.insertHeader(header)
|
||||
|
||||
@@ -18,6 +18,7 @@ package repositories
|
||||
|
||||
import (
|
||||
"context"
|
||||
"github.com/sirupsen/logrus"
|
||||
|
||||
"database/sql"
|
||||
|
||||
@@ -43,12 +44,16 @@ func (logRepository LogRepository) CreateLogs(lgs []core.Log, receiptId int64) e
|
||||
return postgres.ErrDBInsertFailed
|
||||
}
|
||||
}
|
||||
tx.Commit()
|
||||
err := tx.Commit()
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return postgres.ErrDBInsertFailed
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (logRepository LogRepository) GetLogs(address string, blockNumber int64) []core.Log {
|
||||
logRows, _ := logRepository.DB.Query(
|
||||
func (logRepository LogRepository) GetLogs(address string, blockNumber int64) ([]core.Log, error) {
|
||||
logRows, err := logRepository.DB.Query(
|
||||
`SELECT block_number,
|
||||
address,
|
||||
tx_hash,
|
||||
@@ -61,10 +66,13 @@ func (logRepository LogRepository) GetLogs(address string, blockNumber int64) []
|
||||
FROM logs
|
||||
WHERE address = $1 AND block_number = $2
|
||||
ORDER BY block_number DESC`, address, blockNumber)
|
||||
if err != nil {
|
||||
return []core.Log{}, err
|
||||
}
|
||||
return logRepository.loadLogs(logRows)
|
||||
}
|
||||
|
||||
func (logRepository LogRepository) loadLogs(logsRows *sql.Rows) []core.Log {
|
||||
func (logRepository LogRepository) loadLogs(logsRows *sql.Rows) ([]core.Log, error) {
|
||||
var lgs []core.Log
|
||||
for logsRows.Next() {
|
||||
var blockNumber int64
|
||||
@@ -73,7 +81,10 @@ func (logRepository LogRepository) loadLogs(logsRows *sql.Rows) []core.Log {
|
||||
var index int64
|
||||
var data string
|
||||
var topics core.Topics
|
||||
logsRows.Scan(&blockNumber, &address, &txHash, &index, &topics[0], &topics[1], &topics[2], &topics[3], &data)
|
||||
err := logsRows.Scan(&blockNumber, &address, &txHash, &index, &topics[0], &topics[1], &topics[2], &topics[3], &data)
|
||||
if err != nil {
|
||||
logrus.Warn("loadLogs: Error scanning a row in logRows")
|
||||
}
|
||||
lg := core.Log{
|
||||
BlockNumber: blockNumber,
|
||||
TxHash: txHash,
|
||||
@@ -86,5 +97,5 @@ func (logRepository LogRepository) loadLogs(logsRows *sql.Rows) []core.Log {
|
||||
}
|
||||
lgs = append(lgs, lg)
|
||||
}
|
||||
return lgs
|
||||
return lgs, nil
|
||||
}
|
||||
|
||||
@@ -19,6 +19,7 @@ package repositories
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"github.com/sirupsen/logrus"
|
||||
|
||||
"github.com/vulcanize/vulcanizedb/pkg/core"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/datastore"
|
||||
@@ -61,6 +62,9 @@ func createReceipt(receipt core.Receipt, blockId int64, tx *sql.Tx) (int64, erro
|
||||
RETURNING id`,
|
||||
receipt.ContractAddress, receipt.TxHash, receipt.CumulativeGasUsed, receipt.GasUsed, receipt.StateRoot, receipt.Status, blockId,
|
||||
).Scan(&receiptId)
|
||||
if err != nil {
|
||||
logrus.Error("createReceipt: Error inserting: ", err)
|
||||
}
|
||||
return receiptId, err
|
||||
}
|
||||
|
||||
@@ -90,6 +94,7 @@ func (receiptRepository ReceiptRepository) CreateReceipt(blockId int64, receipt
|
||||
receipt.ContractAddress, receipt.TxHash, receipt.CumulativeGasUsed, receipt.GasUsed, receipt.StateRoot, receipt.Status, blockId).Scan(&receiptId)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
logrus.Warning("CreateReceipt: error inserting receipt: ", err)
|
||||
return receiptId, err
|
||||
}
|
||||
tx.Commit()
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
package repositories
|
||||
|
||||
import (
|
||||
"github.com/sirupsen/logrus"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/core"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/datastore/postgres"
|
||||
)
|
||||
@@ -28,6 +29,7 @@ type WatchedEventRepository struct {
|
||||
func (watchedEventRepository WatchedEventRepository) GetWatchedEvents(name string) ([]*core.WatchedEvent, error) {
|
||||
rows, err := watchedEventRepository.DB.Queryx(`SELECT id, name, block_number, address, tx_hash, index, topic0, topic1, topic2, topic3, data FROM watched_event_logs where name=$1`, name)
|
||||
if err != nil {
|
||||
logrus.Error("GetWatchedEvents: error getting watched events: ", err)
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
@@ -37,11 +39,13 @@ func (watchedEventRepository WatchedEventRepository) GetWatchedEvents(name strin
|
||||
lg := new(core.WatchedEvent)
|
||||
err = rows.StructScan(lg)
|
||||
if err != nil {
|
||||
logrus.Warn("GetWatchedEvents: error scanning log: ", err)
|
||||
return nil, err
|
||||
}
|
||||
lgs = append(lgs, lg)
|
||||
}
|
||||
if err = rows.Err(); err != nil {
|
||||
logrus.Warn("GetWatchedEvents: error scanning logs: ", err)
|
||||
return nil, err
|
||||
}
|
||||
return lgs, nil
|
||||
|
||||
@@ -31,7 +31,7 @@ type BlockRepository interface {
|
||||
CreateOrUpdateBlock(block core.Block) (int64, error)
|
||||
GetBlock(blockNumber int64) (core.Block, error)
|
||||
MissingBlockNumbers(startingBlockNumber, endingBlockNumber int64, nodeID string) []int64
|
||||
SetBlocksStatus(chainHead int64)
|
||||
SetBlocksStatus(chainHead int64) error
|
||||
}
|
||||
|
||||
var ErrContractDoesNotExist = func(contractHash string) error {
|
||||
@@ -41,7 +41,7 @@ var ErrContractDoesNotExist = func(contractHash string) error {
|
||||
type ContractRepository interface {
|
||||
CreateContract(contract core.Contract) error
|
||||
GetContract(contractHash string) (core.Contract, error)
|
||||
ContractExists(contractHash string) bool
|
||||
ContractExists(contractHash string) (bool, error)
|
||||
}
|
||||
|
||||
var ErrFilterDoesNotExist = func(name string) error {
|
||||
@@ -57,12 +57,11 @@ type HeaderRepository interface {
|
||||
CreateOrUpdateHeader(header core.Header) (int64, error)
|
||||
GetHeader(blockNumber int64) (core.Header, error)
|
||||
MissingBlockNumbers(startingBlockNumber, endingBlockNumber int64, nodeID string) ([]int64, error)
|
||||
HeaderExists(blockNumber int64) (bool, error)
|
||||
}
|
||||
|
||||
type LogRepository interface {
|
||||
CreateLogs(logs []core.Log, receiptId int64) error
|
||||
GetLogs(address string, blockNumber int64) []core.Log
|
||||
GetLogs(address string, blockNumber int64) ([]core.Log, error)
|
||||
}
|
||||
|
||||
var ErrReceiptDoesNotExist = func(txHash string) error {
|
||||
|
||||
Reference in New Issue
Block a user