forked from cerc-io/ipld-eth-server
Add cold import script
This commit is contained in:
@@ -0,0 +1,19 @@
|
||||
package ethereum
|
||||
|
||||
type DatabaseType int
|
||||
|
||||
const (
|
||||
Level DatabaseType = iota
|
||||
)
|
||||
|
||||
type DatabaseConfig struct {
|
||||
Type DatabaseType
|
||||
Path string
|
||||
}
|
||||
|
||||
func CreateDatabaseConfig(dbType DatabaseType, path string) DatabaseConfig {
|
||||
return DatabaseConfig{
|
||||
Type: dbType,
|
||||
Path: path,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
package ethereum
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/ethereum/go-ethereum/core/types"
|
||||
"github.com/ethereum/go-ethereum/ethdb"
|
||||
|
||||
"github.com/vulcanize/vulcanizedb/pkg/datastore/ethereum/level"
|
||||
)
|
||||
|
||||
type Database interface {
|
||||
GetBlock(hash []byte, blockNumber int64) *types.Block
|
||||
GetBlockHash(blockNumber int64) []byte
|
||||
GetBlockReceipts(blockHash []byte, blockNumber int64) types.Receipts
|
||||
}
|
||||
|
||||
func CreateDatabase(config DatabaseConfig) (Database, error) {
|
||||
switch config.Type {
|
||||
case Level:
|
||||
levelDBConnection, err := ethdb.NewLDBDatabase(config.Path, 128, 1024)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
levelDBReader := level.NewLevelDatabaseReader(levelDBConnection)
|
||||
levelDB := level.NewLevelDatabase(levelDBReader)
|
||||
return levelDB, nil
|
||||
default:
|
||||
return nil, fmt.Errorf("Unknown ethereum database: %s", config.Path)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package level
|
||||
|
||||
import (
|
||||
"github.com/ethereum/go-ethereum/common"
|
||||
"github.com/ethereum/go-ethereum/core/types"
|
||||
)
|
||||
|
||||
type LevelDatabase struct {
|
||||
reader Reader
|
||||
}
|
||||
|
||||
func NewLevelDatabase(ldbReader Reader) *LevelDatabase {
|
||||
return &LevelDatabase{
|
||||
reader: ldbReader,
|
||||
}
|
||||
}
|
||||
|
||||
func (l LevelDatabase) GetBlock(blockHash []byte, blockNumber int64) *types.Block {
|
||||
n := uint64(blockNumber)
|
||||
h := common.BytesToHash(blockHash)
|
||||
return l.reader.GetBlock(h, n)
|
||||
}
|
||||
|
||||
func (l LevelDatabase) GetBlockHash(blockNumber int64) []byte {
|
||||
n := uint64(blockNumber)
|
||||
h := l.reader.GetCanonicalHash(n)
|
||||
return h.Bytes()
|
||||
}
|
||||
|
||||
func (l LevelDatabase) GetBlockReceipts(blockHash []byte, blockNumber int64) types.Receipts {
|
||||
n := uint64(blockNumber)
|
||||
h := common.BytesToHash(blockHash)
|
||||
return l.reader.GetBlockReceipts(h, n)
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package level
|
||||
|
||||
import (
|
||||
"github.com/ethereum/go-ethereum/common"
|
||||
"github.com/ethereum/go-ethereum/core"
|
||||
"github.com/ethereum/go-ethereum/core/types"
|
||||
)
|
||||
|
||||
type Reader interface {
|
||||
GetBlock(hash common.Hash, number uint64) *types.Block
|
||||
GetBlockReceipts(hash common.Hash, number uint64) types.Receipts
|
||||
GetCanonicalHash(number uint64) common.Hash
|
||||
}
|
||||
|
||||
type LevelDatabaseReader struct {
|
||||
core.DatabaseReader
|
||||
}
|
||||
|
||||
func NewLevelDatabaseReader(reader core.DatabaseReader) *LevelDatabaseReader {
|
||||
return &LevelDatabaseReader{DatabaseReader: reader}
|
||||
}
|
||||
|
||||
func (ldbr *LevelDatabaseReader) GetBlock(hash common.Hash, number uint64) *types.Block {
|
||||
return core.GetBlock(ldbr.DatabaseReader, hash, number)
|
||||
}
|
||||
|
||||
func (ldbr *LevelDatabaseReader) GetBlockReceipts(hash common.Hash, number uint64) types.Receipts {
|
||||
return core.GetBlockReceipts(ldbr.DatabaseReader, hash, number)
|
||||
}
|
||||
|
||||
func (ldbr *LevelDatabaseReader) GetCanonicalHash(number uint64) common.Hash {
|
||||
return core.GetCanonicalHash(ldbr.DatabaseReader, number)
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
package level_test
|
||||
|
||||
import (
|
||||
"github.com/ethereum/go-ethereum/common"
|
||||
. "github.com/onsi/ginkgo"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/datastore/ethereum/level"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/fakes"
|
||||
)
|
||||
|
||||
var _ = Describe("Level database", func() {
|
||||
Describe("Getting a block hash", func() {
|
||||
It("converts block number to uint64 to fetch hash from reader", func() {
|
||||
mockReader := fakes.NewMockLevelDatabaseReader()
|
||||
ldb := level.NewLevelDatabase(mockReader)
|
||||
blockNumber := int64(12345)
|
||||
|
||||
ldb.GetBlockHash(blockNumber)
|
||||
|
||||
expectedBlockNumber := uint64(blockNumber)
|
||||
mockReader.AssertGetCanonicalHashCalledWith(expectedBlockNumber)
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Getting a block", func() {
|
||||
It("converts block number to uint64 and hash to common.Hash to fetch block from reader", func() {
|
||||
mockReader := fakes.NewMockLevelDatabaseReader()
|
||||
ldb := level.NewLevelDatabase(mockReader)
|
||||
blockHash := []byte{5, 4, 3, 2, 1}
|
||||
blockNumber := int64(12345)
|
||||
|
||||
ldb.GetBlock(blockHash, blockNumber)
|
||||
|
||||
expectedBlockHash := common.BytesToHash(blockHash)
|
||||
expectedBlockNumber := uint64(blockNumber)
|
||||
mockReader.AssertGetBlockCalledWith(expectedBlockHash, expectedBlockNumber)
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Getting a block's receipts", func() {
|
||||
It("converts block number to uint64 and hash to common.Hash to fetch receipts from reader", func() {
|
||||
mockReader := fakes.NewMockLevelDatabaseReader()
|
||||
ldb := level.NewLevelDatabase(mockReader)
|
||||
blockHash := []byte{5, 4, 3, 2, 1}
|
||||
blockNumber := int64(12345)
|
||||
|
||||
ldb.GetBlockReceipts(blockHash, blockNumber)
|
||||
|
||||
expectedBlockHash := common.BytesToHash(blockHash)
|
||||
expectedBlockNumber := uint64(blockNumber)
|
||||
mockReader.AssertGetBlockReceiptsCalledWith(expectedBlockHash, expectedBlockNumber)
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,13 @@
|
||||
package level_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
. "github.com/onsi/ginkgo"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
|
||||
func TestLevel(t *testing.T) {
|
||||
RegisterFailHandler(Fail)
|
||||
RunSpecs(t, "Level Suite")
|
||||
}
|
||||
@@ -9,14 +9,14 @@ type BlockRepository struct {
|
||||
*InMemory
|
||||
}
|
||||
|
||||
func (blockRepository *BlockRepository) CreateOrUpdateBlock(block core.Block) error {
|
||||
func (blockRepository *BlockRepository) CreateOrUpdateBlock(block core.Block) (int64, error) {
|
||||
blockRepository.CreateOrUpdateBlockCallCount++
|
||||
blockRepository.blocks[block.Number] = block
|
||||
for _, transaction := range block.Transactions {
|
||||
blockRepository.receipts[transaction.Hash] = transaction.Receipt
|
||||
blockRepository.logs[transaction.TxHash] = transaction.Logs
|
||||
}
|
||||
return nil
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func (blockRepository *BlockRepository) GetBlock(blockNumber int64) (core.Block, error) {
|
||||
|
||||
@@ -23,7 +23,9 @@ var _ = Describe("Postgres DB", func() {
|
||||
It("connects to the database", func() {
|
||||
var err error
|
||||
pgConfig := config.DbConnectionString(test_config.DBConfig)
|
||||
|
||||
sqlxdb, err = sqlx.Connect("postgres", pgConfig)
|
||||
|
||||
Expect(err).Should(BeNil())
|
||||
Expect(sqlxdb).ShouldNot(BeNil())
|
||||
})
|
||||
@@ -76,10 +78,10 @@ var _ = Describe("Postgres DB", func() {
|
||||
db := test_config.NewTestDB(node)
|
||||
blocksRepository := repositories.BlockRepository{DB: db}
|
||||
|
||||
err1 := blocksRepository.CreateOrUpdateBlock(badBlock)
|
||||
savedBlock, err2 := blocksRepository.GetBlock(123)
|
||||
_, err1 := blocksRepository.CreateOrUpdateBlock(badBlock)
|
||||
|
||||
Expect(err1).To(HaveOccurred())
|
||||
savedBlock, err2 := blocksRepository.GetBlock(123)
|
||||
Expect(err2).To(HaveOccurred())
|
||||
Expect(savedBlock).To(BeZero())
|
||||
})
|
||||
@@ -87,7 +89,9 @@ var _ = Describe("Postgres DB", func() {
|
||||
It("throws error when can't connect to the database", func() {
|
||||
invalidDatabase := config.Database{}
|
||||
node := core.Node{GenesisBlock: "GENESIS", NetworkID: 1, ID: "x123", ClientName: "geth"}
|
||||
|
||||
_, err := postgres.NewDB(invalidDatabase, node)
|
||||
|
||||
Expect(err).To(Equal(postgres.ErrDBConnectionFailed))
|
||||
})
|
||||
|
||||
@@ -110,10 +114,10 @@ var _ = Describe("Postgres DB", func() {
|
||||
db, _ := postgres.NewDB(test_config.DBConfig, node)
|
||||
logRepository := repositories.LogRepository{DB: db}
|
||||
|
||||
err := logRepository.CreateLogs([]core.Log{badLog})
|
||||
savedBlock := logRepository.GetLogs("x123", 1)
|
||||
err := logRepository.CreateLogs([]core.Log{badLog}, 123)
|
||||
|
||||
Expect(err).ToNot(BeNil())
|
||||
savedBlock := logRepository.GetLogs("x123", 1)
|
||||
Expect(savedBlock).To(BeNil())
|
||||
})
|
||||
|
||||
@@ -129,12 +133,11 @@ var _ = Describe("Postgres DB", func() {
|
||||
db, _ := postgres.NewDB(test_config.DBConfig, node)
|
||||
blockRepository := repositories.BlockRepository{DB: db}
|
||||
|
||||
err1 := blockRepository.CreateOrUpdateBlock(block)
|
||||
savedBlock, err2 := blockRepository.GetBlock(123)
|
||||
_, err1 := blockRepository.CreateOrUpdateBlock(block)
|
||||
|
||||
Expect(err1).To(HaveOccurred())
|
||||
savedBlock, err2 := blockRepository.GetBlock(123)
|
||||
Expect(err2).To(HaveOccurred())
|
||||
Expect(savedBlock).To(BeZero())
|
||||
})
|
||||
|
||||
})
|
||||
|
||||
@@ -3,10 +3,11 @@ package repositories
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
|
||||
"errors"
|
||||
"log"
|
||||
|
||||
"github.com/jmoiron/sqlx"
|
||||
|
||||
"github.com/vulcanize/vulcanizedb/pkg/core"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/datastore"
|
||||
"github.com/vulcanize/vulcanizedb/pkg/datastore/postgres"
|
||||
@@ -16,6 +17,8 @@ const (
|
||||
blocksFromHeadBeforeFinal = 20
|
||||
)
|
||||
|
||||
var ErrBlockExists = errors.New("Won't add block that already exists.")
|
||||
|
||||
type BlockRepository struct {
|
||||
*postgres.DB
|
||||
}
|
||||
@@ -28,22 +31,21 @@ func (blockRepository BlockRepository) SetBlocksStatus(chainHead int64) {
|
||||
cutoff)
|
||||
}
|
||||
|
||||
func (blockRepository BlockRepository) CreateOrUpdateBlock(block core.Block) error {
|
||||
func (blockRepository BlockRepository) CreateOrUpdateBlock(block core.Block) (int64, error) {
|
||||
var err error
|
||||
var blockId int64
|
||||
retrievedBlockHash, ok := blockRepository.getBlockHash(block)
|
||||
if !ok {
|
||||
err = blockRepository.insertBlock(block)
|
||||
return err
|
||||
return blockRepository.insertBlock(block)
|
||||
}
|
||||
if ok && retrievedBlockHash != block.Hash {
|
||||
err = blockRepository.removeBlock(block.Number)
|
||||
if err != nil {
|
||||
return err
|
||||
return 0, err
|
||||
}
|
||||
err = blockRepository.insertBlock(block)
|
||||
return err
|
||||
return blockRepository.insertBlock(block)
|
||||
}
|
||||
return nil
|
||||
return blockId, ErrBlockExists
|
||||
}
|
||||
|
||||
func (blockRepository BlockRepository) MissingBlockNumbers(startingBlockNumber int64, highestBlockNumber int64) []int64 {
|
||||
@@ -92,7 +94,7 @@ func (blockRepository BlockRepository) GetBlock(blockNumber int64) (core.Block,
|
||||
return savedBlock, nil
|
||||
}
|
||||
|
||||
func (blockRepository BlockRepository) insertBlock(block core.Block) error {
|
||||
func (blockRepository BlockRepository) insertBlock(block core.Block) (int64, error) {
|
||||
var blockId int64
|
||||
tx, _ := blockRepository.DB.BeginTx(context.Background(), nil)
|
||||
err := tx.QueryRow(
|
||||
@@ -104,15 +106,17 @@ func (blockRepository BlockRepository) insertBlock(block core.Block) error {
|
||||
Scan(&blockId)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return postgres.ErrDBInsertFailed
|
||||
return 0, postgres.ErrDBInsertFailed
|
||||
}
|
||||
err = blockRepository.createTransactions(tx, blockId, block.Transactions)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return postgres.ErrDBInsertFailed
|
||||
if len(block.Transactions) > 0 {
|
||||
err = blockRepository.createTransactions(tx, blockId, block.Transactions)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return 0, postgres.ErrDBInsertFailed
|
||||
}
|
||||
}
|
||||
tx.Commit()
|
||||
return nil
|
||||
return blockId, nil
|
||||
}
|
||||
|
||||
func (blockRepository BlockRepository) createTransactions(tx *sql.Tx, blockId int64, transactions []core.Transaction) error {
|
||||
|
||||
@@ -177,6 +177,22 @@ var _ = Describe("Saving blocks", func() {
|
||||
Expect(retrievedBlockTwo.Transactions[0].Hash).To(Equal("x678"))
|
||||
})
|
||||
|
||||
It("returns 'block exists' error if attempting to add duplicate block", func() {
|
||||
block := core.Block{
|
||||
Number: 12345,
|
||||
Hash: "0x12345",
|
||||
}
|
||||
|
||||
_, err := blockRepository.CreateOrUpdateBlock(block)
|
||||
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
_, err = blockRepository.CreateOrUpdateBlock(block)
|
||||
|
||||
Expect(err).To(HaveOccurred())
|
||||
Expect(err).To(MatchError(repositories.ErrBlockExists))
|
||||
})
|
||||
|
||||
It("saves the attributes associated to a transaction", func() {
|
||||
gasLimit := uint64(5000)
|
||||
gasPrice := int64(3)
|
||||
|
||||
@@ -13,14 +13,14 @@ type LogRepository struct {
|
||||
*postgres.DB
|
||||
}
|
||||
|
||||
func (logRepository LogRepository) CreateLogs(lgs []core.Log) error {
|
||||
func (logRepository LogRepository) CreateLogs(lgs []core.Log, receiptId int64) error {
|
||||
tx, _ := logRepository.DB.BeginTx(context.Background(), nil)
|
||||
for _, tlog := range lgs {
|
||||
_, err := tx.Exec(
|
||||
`INSERT INTO logs (block_number, address, tx_hash, index, topic0, topic1, topic2, topic3, data)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||
`INSERT INTO logs (block_number, address, tx_hash, index, topic0, topic1, topic2, topic3, data, receipt_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
|
||||
`,
|
||||
tlog.BlockNumber, tlog.Address, tlog.TxHash, tlog.Index, tlog.Topics[0], tlog.Topics[1], tlog.Topics[2], tlog.Topics[3], tlog.Data,
|
||||
tlog.BlockNumber, tlog.Address, tlog.TxHash, tlog.Index, tlog.Topics[0], tlog.Topics[1], tlog.Topics[2], tlog.Topics[3], tlog.Data, receiptId,
|
||||
)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
|
||||
@@ -13,36 +13,45 @@ import (
|
||||
)
|
||||
|
||||
var _ = Describe("Logs Repository", func() {
|
||||
var db *postgres.DB
|
||||
var logsRepository datastore.LogRepository
|
||||
var node core.Node
|
||||
BeforeEach(func() {
|
||||
node = core.Node{
|
||||
GenesisBlock: "GENESIS",
|
||||
NetworkID: 1,
|
||||
ID: "b6f90c0fdd8ec9607aed8ee45c69322e47b7063f0bfb7a29c8ecafab24d0a22d24dd2329b5ee6ed4125a03cb14e57fd584e67f9e53e6c631055cbbd82f080845",
|
||||
ClientName: "Geth/v1.7.2-stable-1db4ecdc/darwin-amd64/go1.9",
|
||||
}
|
||||
db = test_config.NewTestDB(node)
|
||||
logsRepository = repositories.LogRepository{DB: db}
|
||||
})
|
||||
|
||||
Describe("Saving logs", func() {
|
||||
var db *postgres.DB
|
||||
var blockRepository datastore.BlockRepository
|
||||
var logsRepository datastore.LogRepository
|
||||
var receiptRepository datastore.ReceiptRepository
|
||||
var node core.Node
|
||||
|
||||
BeforeEach(func() {
|
||||
node = core.Node{
|
||||
GenesisBlock: "GENESIS",
|
||||
NetworkID: 1,
|
||||
ID: "b6f90c0fdd8ec9607aed8ee45c69322e47b7063f0bfb7a29c8ecafab24d0a22d24dd2329b5ee6ed4125a03cb14e57fd584e67f9e53e6c631055cbbd82f080845",
|
||||
ClientName: "Geth/v1.7.2-stable-1db4ecdc/darwin-amd64/go1.9",
|
||||
}
|
||||
db = test_config.NewTestDB(node)
|
||||
blockRepository = repositories.BlockRepository{DB: db}
|
||||
logsRepository = repositories.LogRepository{DB: db}
|
||||
receiptRepository = repositories.ReceiptRepository{DB: db}
|
||||
})
|
||||
|
||||
It("returns the log when it exists", func() {
|
||||
blockNumber := int64(12345)
|
||||
blockId, err := blockRepository.CreateOrUpdateBlock(core.Block{Number: blockNumber})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
receiptId, err := receiptRepository.CreateReceipt(blockId, core.Receipt{})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
logsRepository.CreateLogs([]core.Log{{
|
||||
BlockNumber: 1,
|
||||
BlockNumber: blockNumber,
|
||||
Index: 0,
|
||||
Address: "x123",
|
||||
TxHash: "x456",
|
||||
Topics: core.Topics{0: "x777", 1: "x888", 2: "x999"},
|
||||
Data: "xabc",
|
||||
}},
|
||||
)
|
||||
}}, receiptId)
|
||||
|
||||
log := logsRepository.GetLogs("x123", 1)
|
||||
log := logsRepository.GetLogs("x123", blockNumber)
|
||||
|
||||
Expect(log).NotTo(BeNil())
|
||||
Expect(log[0].BlockNumber).To(Equal(int64(1)))
|
||||
Expect(log[0].BlockNumber).To(Equal(blockNumber))
|
||||
Expect(log[0].Address).To(Equal("x123"))
|
||||
Expect(log[0].Index).To(Equal(int64(0)))
|
||||
Expect(log[0].TxHash).To(Equal("x456"))
|
||||
@@ -58,24 +67,27 @@ var _ = Describe("Logs Repository", func() {
|
||||
})
|
||||
|
||||
It("filters to the correct block number and address", func() {
|
||||
blockNumber := int64(12345)
|
||||
blockId, err := blockRepository.CreateOrUpdateBlock(core.Block{Number: blockNumber})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
receiptId, err := receiptRepository.CreateReceipt(blockId, core.Receipt{})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
logsRepository.CreateLogs([]core.Log{{
|
||||
BlockNumber: 1,
|
||||
BlockNumber: blockNumber,
|
||||
Index: 0,
|
||||
Address: "x123",
|
||||
TxHash: "x456",
|
||||
Topics: core.Topics{0: "x777", 1: "x888", 2: "x999"},
|
||||
Data: "xabc",
|
||||
}},
|
||||
)
|
||||
}}, receiptId)
|
||||
logsRepository.CreateLogs([]core.Log{{
|
||||
BlockNumber: 1,
|
||||
BlockNumber: blockNumber,
|
||||
Index: 1,
|
||||
Address: "x123",
|
||||
TxHash: "x789",
|
||||
Topics: core.Topics{0: "x111", 1: "x222", 2: "x333"},
|
||||
Data: "xdef",
|
||||
}},
|
||||
)
|
||||
}}, receiptId)
|
||||
logsRepository.CreateLogs([]core.Log{{
|
||||
BlockNumber: 2,
|
||||
Index: 0,
|
||||
@@ -83,10 +95,9 @@ var _ = Describe("Logs Repository", func() {
|
||||
TxHash: "x456",
|
||||
Topics: core.Topics{0: "x777", 1: "x888", 2: "x999"},
|
||||
Data: "xabc",
|
||||
}},
|
||||
)
|
||||
}}, receiptId)
|
||||
|
||||
log := logsRepository.GetLogs("x123", 1)
|
||||
log := logsRepository.GetLogs("x123", blockNumber)
|
||||
|
||||
type logIndex struct {
|
||||
blockNumber int64
|
||||
@@ -111,14 +122,12 @@ var _ = Describe("Logs Repository", func() {
|
||||
Expect(len(log)).To(Equal(2))
|
||||
Expect(uniqueBlockNumbers).To(Equal(
|
||||
[]logIndex{
|
||||
{blockNumber: 1, Index: 0},
|
||||
{blockNumber: 1, Index: 1}},
|
||||
{blockNumber: blockNumber, Index: 0},
|
||||
{blockNumber: blockNumber, Index: 1}},
|
||||
))
|
||||
})
|
||||
|
||||
It("saves the logs attached to a receipt", func() {
|
||||
var blockRepository datastore.BlockRepository
|
||||
blockRepository = repositories.BlockRepository{DB: db}
|
||||
|
||||
logs := []core.Log{{
|
||||
Address: "0x8a4774fe82c63484afef97ca8d89a6ea5e21f973",
|
||||
@@ -171,7 +180,7 @@ var _ = Describe("Logs Repository", func() {
|
||||
}
|
||||
|
||||
block := core.Block{Transactions: []core.Transaction{transaction}}
|
||||
err := blockRepository.CreateOrUpdateBlock(block)
|
||||
_, err := blockRepository.CreateOrUpdateBlock(block)
|
||||
Expect(err).To(Not(HaveOccurred()))
|
||||
retrievedLogs := logsRepository.GetLogs("0x99041f808d598b782d5a3e498681c2452a31da08", 4745407)
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package repositories
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
|
||||
"github.com/vulcanize/vulcanizedb/pkg/core"
|
||||
@@ -12,6 +13,73 @@ type ReceiptRepository struct {
|
||||
*postgres.DB
|
||||
}
|
||||
|
||||
func (receiptRepository ReceiptRepository) CreateReceiptsAndLogs(blockId int64, receipts []core.Receipt) error {
|
||||
tx, err := receiptRepository.DB.BeginTx(context.Background(), nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, receipt := range receipts {
|
||||
receiptId, err := createReceipt(receipt, blockId, tx)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return err
|
||||
}
|
||||
if len(receipt.Logs) > 0 {
|
||||
err = createLogs(receipt.Logs, receiptId, tx)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
tx.Commit()
|
||||
return nil
|
||||
}
|
||||
|
||||
func createReceipt(receipt core.Receipt, blockId int64, tx *sql.Tx) (int64, error) {
|
||||
var receiptId int64
|
||||
err := tx.QueryRow(
|
||||
`INSERT INTO receipts
|
||||
(contract_address, tx_hash, cumulative_gas_used, gas_used, state_root, status, block_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
||||
RETURNING id`,
|
||||
receipt.ContractAddress, receipt.TxHash, receipt.CumulativeGasUsed, receipt.GasUsed, receipt.StateRoot, receipt.Status, blockId,
|
||||
).Scan(&receiptId)
|
||||
return receiptId, err
|
||||
}
|
||||
|
||||
func createLogs(logs []core.Log, receiptId int64, tx *sql.Tx) error {
|
||||
for _, log := range logs {
|
||||
_, err := tx.Exec(
|
||||
`INSERT INTO logs (block_number, address, tx_hash, index, topic0, topic1, topic2, topic3, data, receipt_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
|
||||
`,
|
||||
log.BlockNumber, log.Address, log.TxHash, log.Index, log.Topics[0], log.Topics[1], log.Topics[2], log.Topics[3], log.Data, receiptId,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (receiptRepository ReceiptRepository) CreateReceipt(blockId int64, receipt core.Receipt) (int64, error) {
|
||||
tx, _ := receiptRepository.DB.BeginTx(context.Background(), nil)
|
||||
var receiptId int64
|
||||
err := tx.QueryRow(
|
||||
`INSERT INTO receipts
|
||||
(contract_address, tx_hash, cumulative_gas_used, gas_used, state_root, status, block_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
||||
RETURNING id`,
|
||||
receipt.ContractAddress, receipt.TxHash, receipt.CumulativeGasUsed, receipt.GasUsed, receipt.StateRoot, receipt.Status, blockId).Scan(&receiptId)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
return receiptId, err
|
||||
}
|
||||
tx.Commit()
|
||||
return receiptId, nil
|
||||
}
|
||||
|
||||
func (receiptRepository ReceiptRepository) GetReceipt(txHash string) (core.Receipt, error) {
|
||||
row := receiptRepository.DB.QueryRow(
|
||||
`SELECT contract_address,
|
||||
|
||||
@@ -10,7 +10,9 @@ import (
|
||||
"github.com/vulcanize/vulcanizedb/test_config"
|
||||
)
|
||||
|
||||
var _ bool = Describe("Logs Repository", func() {
|
||||
var _ = Describe("Logs Repository", func() {
|
||||
var blockRepository datastore.BlockRepository
|
||||
var logRepository datastore.LogRepository
|
||||
var receiptRepository datastore.ReceiptRepository
|
||||
var db *postgres.DB
|
||||
var node core.Node
|
||||
@@ -22,14 +24,67 @@ var _ bool = Describe("Logs Repository", func() {
|
||||
ClientName: "Geth/v1.7.2-stable-1db4ecdc/darwin-amd64/go1.9",
|
||||
}
|
||||
db = test_config.NewTestDB(node)
|
||||
blockRepository = repositories.BlockRepository{DB: db}
|
||||
logRepository = repositories.LogRepository{DB: db}
|
||||
receiptRepository = repositories.ReceiptRepository{DB: db}
|
||||
})
|
||||
|
||||
Describe("Saving receipts", func() {
|
||||
Describe("Saving multiple receipts", func() {
|
||||
It("persists each receipt and its logs", func() {
|
||||
blockNumber := int64(1234567)
|
||||
blockId, err := blockRepository.CreateOrUpdateBlock(core.Block{Number: blockNumber})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
txHashOne := "0xTxHashOne"
|
||||
txHashTwo := "0xTxHashTwo"
|
||||
addressOne := "0xAddressOne"
|
||||
addressTwo := "0xAddressTwo"
|
||||
logsOne := []core.Log{{
|
||||
Address: addressOne,
|
||||
BlockNumber: blockNumber,
|
||||
TxHash: txHashOne,
|
||||
}, {
|
||||
Address: addressOne,
|
||||
BlockNumber: blockNumber,
|
||||
TxHash: txHashOne,
|
||||
}}
|
||||
logsTwo := []core.Log{{
|
||||
BlockNumber: blockNumber,
|
||||
TxHash: txHashTwo,
|
||||
Address: addressTwo,
|
||||
}}
|
||||
receiptOne := core.Receipt{
|
||||
Logs: logsOne,
|
||||
TxHash: txHashOne,
|
||||
}
|
||||
receiptTwo := core.Receipt{
|
||||
Logs: logsTwo,
|
||||
TxHash: txHashTwo,
|
||||
}
|
||||
receipts := []core.Receipt{receiptOne, receiptTwo}
|
||||
|
||||
err = receiptRepository.CreateReceiptsAndLogs(blockId, receipts)
|
||||
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
persistedReceiptOne, err := receiptRepository.GetReceipt(txHashOne)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(persistedReceiptOne).NotTo(BeNil())
|
||||
Expect(persistedReceiptOne.TxHash).To(Equal(txHashOne))
|
||||
persistedReceiptTwo, err := receiptRepository.GetReceipt(txHashTwo)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(persistedReceiptTwo).NotTo(BeNil())
|
||||
Expect(persistedReceiptTwo.TxHash).To(Equal(txHashTwo))
|
||||
persistedAddressOneLogs := logRepository.GetLogs(addressOne, blockNumber)
|
||||
Expect(persistedAddressOneLogs).NotTo(BeNil())
|
||||
Expect(len(persistedAddressOneLogs)).To(Equal(2))
|
||||
persistedAddressTwoLogs := logRepository.GetLogs(addressTwo, blockNumber)
|
||||
Expect(persistedAddressTwoLogs).NotTo(BeNil())
|
||||
Expect(len(persistedAddressTwoLogs)).To(Equal(1))
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Saving receipts on a block's transactions", func() {
|
||||
It("returns the receipt when it exists", func() {
|
||||
var blockRepository datastore.BlockRepository
|
||||
db := test_config.NewTestDB(node)
|
||||
blockRepository = repositories.BlockRepository{DB: db}
|
||||
expected := core.Receipt{
|
||||
ContractAddress: "0xde0b295669a9fd93d5f28d9ec85e40f4cb697bae",
|
||||
CumulativeGasUsed: 7996119,
|
||||
@@ -66,9 +121,6 @@ var _ bool = Describe("Logs Repository", func() {
|
||||
})
|
||||
|
||||
It("still saves receipts without logs", func() {
|
||||
var blockRepository datastore.BlockRepository
|
||||
db := test_config.NewTestDB(node)
|
||||
blockRepository = repositories.BlockRepository{DB: db}
|
||||
receipt := core.Receipt{
|
||||
TxHash: "0x002c4799161d809b23f67884eb6598c9df5894929fe1a9ead97ca175d360f547",
|
||||
}
|
||||
|
||||
@@ -13,14 +13,18 @@ import (
|
||||
|
||||
var _ = Describe("Watched Events Repository", func() {
|
||||
var db *postgres.DB
|
||||
var logRepository datastore.LogRepository
|
||||
var blocksRepository datastore.BlockRepository
|
||||
var filterRepository datastore.FilterRepository
|
||||
var logRepository datastore.LogRepository
|
||||
var receiptRepository datastore.ReceiptRepository
|
||||
var watchedEventRepository datastore.WatchedEventRepository
|
||||
|
||||
BeforeEach(func() {
|
||||
db = test_config.NewTestDB(core.Node{})
|
||||
logRepository = repositories.LogRepository{DB: db}
|
||||
blocksRepository = repositories.BlockRepository{DB: db}
|
||||
filterRepository = repositories.FilterRepository{DB: db}
|
||||
logRepository = repositories.LogRepository{DB: db}
|
||||
receiptRepository = repositories.ReceiptRepository{DB: db}
|
||||
watchedEventRepository = repositories.WatchedEventRepository{DB: db}
|
||||
})
|
||||
|
||||
@@ -56,7 +60,11 @@ var _ = Describe("Watched Events Repository", func() {
|
||||
}
|
||||
err := filterRepository.CreateFilter(filter)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
err = logRepository.CreateLogs(logs)
|
||||
blockId, err := blocksRepository.CreateOrUpdateBlock(core.Block{})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
receiptId, err := receiptRepository.CreateReceipt(blockId, core.Receipt{})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
err = logRepository.CreateLogs(logs, receiptId)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
matchingLogs, err := watchedEventRepository.GetWatchedEvents("Filter1")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
@@ -69,7 +77,6 @@ var _ = Describe("Watched Events Repository", func() {
|
||||
Expect(matchingLogs[0].Topic1).To(Equal(expectedWatchedEventLog[0].Topic1))
|
||||
Expect(matchingLogs[0].Topic2).To(Equal(expectedWatchedEventLog[0].Topic2))
|
||||
Expect(matchingLogs[0].Data).To(Equal(expectedWatchedEventLog[0].Data))
|
||||
|
||||
})
|
||||
|
||||
It("retrieves a watched event log by name", func() {
|
||||
@@ -110,7 +117,11 @@ var _ = Describe("Watched Events Repository", func() {
|
||||
}}
|
||||
err := filterRepository.CreateFilter(filter)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
err = logRepository.CreateLogs(logs)
|
||||
blockId, err := blocksRepository.CreateOrUpdateBlock(core.Block{Hash: "Ox123"})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
receiptId, err := receiptRepository.CreateReceipt(blockId, core.Receipt{TxHash: "0x123"})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
err = logRepository.CreateLogs(logs, receiptId)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
matchingLogs, err := watchedEventRepository.GetWatchedEvents("Filter1")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
@@ -12,7 +12,7 @@ var ErrBlockDoesNotExist = func(blockNumber int64) error {
|
||||
}
|
||||
|
||||
type BlockRepository interface {
|
||||
CreateOrUpdateBlock(block core.Block) error
|
||||
CreateOrUpdateBlock(block core.Block) (int64, error)
|
||||
GetBlock(blockNumber int64) (core.Block, error)
|
||||
MissingBlockNumbers(startingBlockNumber int64, endingBlockNumber int64) []int64
|
||||
SetBlocksStatus(chainHead int64)
|
||||
@@ -38,7 +38,7 @@ type FilterRepository interface {
|
||||
}
|
||||
|
||||
type LogRepository interface {
|
||||
CreateLogs(logs []core.Log) error
|
||||
CreateLogs(logs []core.Log, receiptId int64) error
|
||||
GetLogs(address string, blockNumber int64) []core.Log
|
||||
}
|
||||
|
||||
@@ -47,6 +47,8 @@ var ErrReceiptDoesNotExist = func(txHash string) error {
|
||||
}
|
||||
|
||||
type ReceiptRepository interface {
|
||||
CreateReceiptsAndLogs(blockId int64, receipts []core.Receipt) error
|
||||
CreateReceipt(blockId int64, receipt core.Receipt) (int64, error)
|
||||
GetReceipt(txHash string) (core.Receipt, error)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user