Add graphql server (#27)

* Add graphql server

* Update Makefile

* Update log_filters constraint

* Add GetLogFilter to repo

* Update travis (use Makefile, go fmt, go vet)

* Add logFilter schema and resolvers

* Add GetWatchedEvent to watched_events_repo

* Add watchedEventLog schema and resolvers
This commit is contained in:
Matt K
2018-02-08 10:12:08 -06:00
committed by GitHub
parent d5852654bb
commit 605b0a96ae
113 changed files with 16180 additions and 153 deletions
+1 -1
View File
@@ -3,9 +3,9 @@ package config_test
import (
"path/filepath"
cfg "github.com/vulcanize/vulcanizedb/pkg/config"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
cfg "github.com/vulcanize/vulcanizedb/pkg/config"
)
var _ = Describe("Loading the config", func() {
+1 -1
View File
@@ -3,8 +3,8 @@ package contract_summary
import (
"fmt"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum/common"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
func GenerateConsoleOutput(summary ContractSummary) string {
+1 -1
View File
@@ -18,7 +18,7 @@ type ContractSummary struct {
}
func NewSummary(blockchain core.Blockchain, repository repositories.Repository, contractHash string, blockNumber *big.Int) (ContractSummary, error) {
contract, err := repository.FindContract(contractHash)
contract, err := repository.GetContract(contractHash)
if err != nil {
return ContractSummary{}, err
} else {
+14
View File
@@ -0,0 +1,14 @@
package core
type WatchedEvent struct {
Name string `json:"name"` // name
BlockNumber int64 `json:"block_number" db:"block_number"` // block_number
Address string `json:"address"` // address
TxHash string `json:"tx_hash" db:"tx_hash"` // tx_hash
Index int64 `json:"index"` // index
Topic0 string `json:"topic0"` // topic0
Topic1 string `json:"topic1"` // topic1
Topic2 string `json:"topic2"` // topic2
Topic3 string `json:"topic3"` // topic3
Data string `json:"data"` // data
}
+3 -3
View File
@@ -5,17 +5,17 @@ import (
"errors"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
type LogFilters []LogFilter
type LogFilter struct {
Name string `json:"name"`
FromBlock int64 `json:"fromBlock"`
ToBlock int64 `json:"toBlock"`
FromBlock int64 `json:"fromBlock" db:"from_block"`
ToBlock int64 `json:"toBlock" db:"to_block"`
Address string `json:"address"`
core.Topics `json:"topics"`
}
+2 -2
View File
@@ -3,10 +3,10 @@ package filters_test
import (
"encoding/json"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/filters"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/filters"
)
var _ = Describe("Log filters", func() {
+2 -2
View File
@@ -9,12 +9,12 @@ import (
"log"
cfg "github.com/vulcanize/vulcanizedb/pkg/config"
"github.com/vulcanize/vulcanizedb/pkg/geth"
"github.com/ethereum/go-ethereum/accounts/abi"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/onsi/gomega/ghttp"
cfg "github.com/vulcanize/vulcanizedb/pkg/config"
"github.com/vulcanize/vulcanizedb/pkg/geth"
)
var _ = Describe("ABI files", func() {
+1 -1
View File
@@ -1,9 +1,9 @@
package geth
import (
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/params"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
func CalcUnclesReward(block core.Block, uncles []*types.Header) float64 {
+1 -1
View File
@@ -5,12 +5,12 @@ import (
"context"
"github.com/vulcanize/vulcanizedb/pkg/geth"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/core/types"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/vulcanize/vulcanizedb/pkg/geth"
)
type FakeGethClient struct {
+2 -2
View File
@@ -7,13 +7,13 @@ import (
"log"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/geth/node"
"github.com/ethereum/go-ethereum"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/ethclient"
"github.com/ethereum/go-ethereum/rpc"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/geth/node"
"golang.org/x/net/context"
)
+2 -2
View File
@@ -8,9 +8,9 @@ import (
"context"
"math/big"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum"
"github.com/ethereum/go-ethereum/common"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
var (
@@ -50,7 +50,7 @@ func (blockchain *Blockchain) GetAttributes(contract core.Contract) (core.Contra
for _, abiElement := range parsed.Methods {
if (len(abiElement.Outputs) > 0) && (len(abiElement.Inputs) == 0) && abiElement.Const {
attributeType := abiElement.Outputs[0].Type.String()
contractAttributes = append(contractAttributes, core.ContractAttribute{abiElement.Name, attributeType})
contractAttributes = append(contractAttributes, core.ContractAttribute{Name: abiElement.Name, Type: attributeType})
}
}
sort.Sort(contractAttributes)
+1 -1
View File
@@ -3,10 +3,10 @@ package geth
import (
"strings"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/core/types"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
func ToCoreLogs(gethLogs []types.Log) []core.Log {
+2 -2
View File
@@ -3,13 +3,13 @@ package geth_test
import (
"strings"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/geth"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/core/types"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/geth"
)
var _ = Describe("Conversion of GethLog to core.Log", func() {
+1 -1
View File
@@ -5,10 +5,10 @@ import (
"strconv"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/rpc"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
func Info(client *rpc.Client) core.Node {
+1 -1
View File
@@ -5,10 +5,10 @@ import (
"bytes"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/core/types"
"github.com/vulcanize/vulcanizedb/pkg/core"
)
func BigTo64(n *big.Int) int64 {
+2 -2
View File
@@ -3,13 +3,13 @@ package geth_test
import (
"math/big"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/geth"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/core/types"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/geth"
)
var _ = Describe("Conversion of GethReceipt to core.Receipt", func() {
@@ -0,0 +1,13 @@
package graphql_server_test
import (
"testing"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
)
func TestGraphqlServer(t *testing.T) {
RegisterFailHandler(Fail)
RunSpecs(t, "GraphqlServer Suite")
}
+162
View File
@@ -0,0 +1,162 @@
package graphql_server
import (
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/filters"
"github.com/vulcanize/vulcanizedb/pkg/repositories"
)
var Schema = `
schema {
query: Query
}
type Query {
logFilter(name: String!): LogFilter
watchedEvents(name: String!): WatchedEventList
}
type LogFilter {
name: String!
fromBlock: Int
toBlock: Int
address: String!
topics: [String]!
}
type WatchedEventList{
total: Int!
watchedEvents: [WatchedEvent]!
}
type WatchedEvent {
name: String!
blockNumber: Int!
address: String!
tx_hash: String!
topic0: String!
topic1: String!
topic2: String!
topic3: String!
data: String!
}
`
type Resolver struct {
repository repositories.Repository
}
func NewResolver(repository repositories.Repository) *Resolver {
return &Resolver{repository: repository}
}
func (r *Resolver) LogFilter(args struct {
Name string
}) (*logFilterResolver, error) {
logFilter, err := r.repository.GetFilter(args.Name)
if err != nil {
return &logFilterResolver{}, err
}
return &logFilterResolver{&logFilter}, nil
}
type logFilterResolver struct {
lf *filters.LogFilter
}
func (lfr *logFilterResolver) Name() string {
return lfr.lf.Name
}
func (lfr *logFilterResolver) FromBlock() *int32 {
fromBlock := int32(lfr.lf.FromBlock)
return &fromBlock
}
func (lfr *logFilterResolver) ToBlock() *int32 {
toBlock := int32(lfr.lf.ToBlock)
return &toBlock
}
func (lfr *logFilterResolver) Address() string {
return lfr.lf.Address
}
func (lfr *logFilterResolver) Topics() []*string {
var topics = make([]*string, 4)
for i := range topics {
if lfr.lf.Topics[i] != "" {
topics[i] = &lfr.lf.Topics[i]
}
}
return topics
}
func (r *Resolver) WatchedEvents(args struct {
Name string
}) (*watchedEventsResolver, error) {
watchedEvents, err := r.repository.GetWatchedEvents(args.Name)
if err != nil {
return &watchedEventsResolver{}, err
}
return &watchedEventsResolver{watchedEvents: watchedEvents}, err
}
type watchedEventsResolver struct {
watchedEvents []*core.WatchedEvent
}
func (wesr watchedEventsResolver) WatchedEvents() []*watchedEventResolver {
return resolveWatchedEvents(wesr.watchedEvents)
}
func (wesr watchedEventsResolver) Total() int32 {
return int32(len(wesr.watchedEvents))
}
func resolveWatchedEvents(watchedEvents []*core.WatchedEvent) []*watchedEventResolver {
watchedEventResolvers := make([]*watchedEventResolver, 0)
for _, watchedEvent := range watchedEvents {
watchedEventResolvers = append(watchedEventResolvers, &watchedEventResolver{watchedEvent})
}
return watchedEventResolvers
}
type watchedEventResolver struct {
we *core.WatchedEvent
}
func (wer watchedEventResolver) Name() string {
return wer.we.Name
}
func (wer watchedEventResolver) BlockNumber() int32 {
return int32(wer.we.BlockNumber)
}
func (wer watchedEventResolver) Address() string {
return wer.we.Address
}
func (wer watchedEventResolver) TxHash() string {
return wer.we.TxHash
}
func (wer watchedEventResolver) Topic0() string {
return wer.we.Topic0
}
func (wer watchedEventResolver) Topic1() string {
return wer.we.Topic1
}
func (wer watchedEventResolver) Topic2() string {
return wer.we.Topic2
}
func (wer watchedEventResolver) Topic3() string {
return wer.we.Topic3
}
func (wer watchedEventResolver) Data() string {
return wer.we.Data
}
+168
View File
@@ -0,0 +1,168 @@
package graphql_server_test
import (
"log"
"encoding/json"
"context"
"github.com/neelance/graphql-go"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/vulcanize/vulcanizedb/pkg/config"
"github.com/vulcanize/vulcanizedb/pkg/core"
"github.com/vulcanize/vulcanizedb/pkg/filters"
"github.com/vulcanize/vulcanizedb/pkg/graphql_server"
"github.com/vulcanize/vulcanizedb/pkg/repositories"
"github.com/vulcanize/vulcanizedb/pkg/repositories/postgres"
)
func formatJSON(data []byte) []byte {
var v interface{}
if err := json.Unmarshal(data, &v); err != nil {
log.Fatalf("invalid JSON: %s", err)
}
formatted, err := json.Marshal(v)
if err != nil {
log.Fatal(err)
}
return formatted
}
var _ = Describe("GraphQL", func() {
var cfg config.Config
var repository repositories.Repository
BeforeEach(func() {
cfg, _ = config.NewConfig("private")
node := core.Node{GenesisBlock: "GENESIS", NetworkId: 1, Id: "x123", ClientName: "geth"}
repository = postgres.BuildRepository(node)
e := repository.CreateFilter(filters.LogFilter{
Name: "TestFilter1",
FromBlock: 1,
ToBlock: 10,
Address: "0x123456789",
Topics: core.Topics{0: "topic=1", 2: "topic=2"},
})
if e != nil {
log.Fatal(e)
}
f, e := repository.GetFilter("TestFilter1")
if e != nil {
log.Println(f)
log.Fatal(e)
}
matchingEvent := core.Log{
BlockNumber: 5,
TxHash: "0xTX1",
Address: "0x123456789",
Topics: core.Topics{0: "topic=1", 2: "topic=2"},
Index: 0,
Data: "0xDATADATADATA",
}
nonMatchingEvent := core.Log{
BlockNumber: 5,
TxHash: "0xTX2",
Address: "0xOTHERADDRESS",
Topics: core.Topics{0: "topic=1", 2: "topic=2"},
Index: 0,
Data: "0xDATADATADATA",
}
e = repository.CreateLogs([]core.Log{matchingEvent, nonMatchingEvent})
if e != nil {
log.Fatal(e)
}
})
It("Queries example schema for specific log filter", func() {
var variables map[string]interface{}
r := graphql_server.NewResolver(repository)
var schema = graphql.MustParseSchema(graphql_server.Schema, r)
response := schema.Exec(context.Background(),
`{
logFilter(name: "TestFilter1") {
name
fromBlock
toBlock
address
topics
}
}`,
"",
variables)
expected := `{
"logFilter": {
"name": "TestFilter1",
"fromBlock": 1,
"toBlock": 10,
"address": "0x123456789",
"topics": ["topic=1", null, "topic=2", null]
}
}`
var v interface{}
if len(response.Errors) != 0 {
log.Fatal(response.Errors)
}
err := json.Unmarshal(response.Data, &v)
Expect(err).ToNot(HaveOccurred())
a := formatJSON(response.Data)
e := formatJSON([]byte(expected))
Expect(a).To(Equal(e))
})
It("Queries example schema for specific watched event log", func() {
var variables map[string]interface{}
r := graphql_server.NewResolver(repository)
var schema = graphql.MustParseSchema(graphql_server.Schema, r)
response := schema.Exec(context.Background(),
`{
watchedEvents(name: "TestFilter1") {
total
watchedEvents{
name
blockNumber
address
tx_hash
topic0
topic1
topic2
topic3
data
}
}
}`,
"",
variables)
expected := `{
"watchedEvents":
{
"total": 1,
"watchedEvents": [
{"name":"TestFilter1",
"blockNumber": 5,
"address": "0x123456789",
"tx_hash": "0xTX1",
"topic0": "topic=1",
"topic1": "",
"topic2": "topic=2",
"topic3": "",
"data": "0xDATADATADATA"
}
]
}
}`
var v interface{}
if len(response.Errors) != 0 {
log.Fatal(response.Errors)
}
err := json.Unmarshal(response.Data, &v)
Expect(err).ToNot(HaveOccurred())
a := formatJSON(response.Data)
e := formatJSON([]byte(expected))
Expect(a).To(Equal(e))
})
})
+6 -6
View File
@@ -21,7 +21,7 @@ var _ = Describe("Populating blocks", func() {
repository.CreateOrUpdateBlock(core.Block{Number: 2})
blocksAdded := history.PopulateMissingBlocks(blockchain, repository, 1)
_, err := repository.FindBlockByNumber(1)
_, err := repository.GetBlock(1)
Expect(blocksAdded).To(Equal(1))
Expect(err).ToNot(HaveOccurred())
@@ -54,15 +54,15 @@ var _ = Describe("Populating blocks", func() {
Expect(blocksAdded).To(Equal(3))
Expect(repository.BlockCount()).To(Equal(11))
_, err := repository.FindBlockByNumber(4)
_, err := repository.GetBlock(4)
Expect(err).To(HaveOccurred())
_, err = repository.FindBlockByNumber(5)
_, err = repository.GetBlock(5)
Expect(err).ToNot(HaveOccurred())
_, err = repository.FindBlockByNumber(8)
_, err = repository.GetBlock(8)
Expect(err).ToNot(HaveOccurred())
_, err = repository.FindBlockByNumber(10)
_, err = repository.GetBlock(10)
Expect(err).ToNot(HaveOccurred())
_, err = repository.FindBlockByNumber(13)
_, err = repository.GetBlock(13)
Expect(err).To(HaveOccurred())
})
+4 -4
View File
@@ -36,7 +36,7 @@ var _ = Describe("Blocks validator", func() {
})
It("returns the window size", func() {
window := history.ValidationWindow{1, 3}
window := history.ValidationWindow{LowerBound: 1, UpperBound: 3}
Expect(window.Size()).To(Equal(2))
})
@@ -51,15 +51,15 @@ var _ = Describe("Blocks validator", func() {
validator := history.NewBlockValidator(blockchain, repository, 2)
window := validator.ValidateBlocks()
Expect(window).To(Equal(history.ValidationWindow{5, 7}))
Expect(window).To(Equal(history.ValidationWindow{LowerBound: 5, UpperBound: 7}))
Expect(repository.BlockCount()).To(Equal(2))
Expect(repository.CreateOrUpdateBlockCallCount).To(Equal(2))
})
It("logs window message", func() {
expectedMessage := &bytes.Buffer{}
window := history.ValidationWindow{5, 7}
history.ParsedWindowTemplate.Execute(expectedMessage, history.ValidationWindow{5, 7})
window := history.ValidationWindow{LowerBound: 5, UpperBound: 7}
history.ParsedWindowTemplate.Execute(expectedMessage, history.ValidationWindow{LowerBound: 5, UpperBound: 7})
blockchain := fakes.NewBlockchainWithBlocks([]core.Block{})
repository := inmemory.NewInMemory()
+14 -6
View File
@@ -23,7 +23,15 @@ type InMemory struct {
CreateOrUpdateBlockCallCount int
}
func (repository *InMemory) AddFilter(filter filters.LogFilter) error {
func (repository *InMemory) GetWatchedEvents(name string) ([]*core.WatchedEvent, error) {
panic("implement me")
}
func (repository *InMemory) GetFilter(name string) (filters.LogFilter, error) {
panic("implement me")
}
func (repository *InMemory) CreateFilter(filter filters.LogFilter) error {
key := filter.Name
if _, ok := repository.logFilters[key]; ok || key == "" {
return errors.New("filter name not unique")
@@ -43,7 +51,7 @@ func NewInMemory() *InMemory {
}
}
func (repository *InMemory) FindReceipt(txHash string) (core.Receipt, error) {
func (repository *InMemory) GetReceipt(txHash string) (core.Receipt, error) {
if receipt, ok := repository.receipts[txHash]; ok {
return receipt, nil
}
@@ -62,14 +70,14 @@ func (repository *InMemory) SetBlocksStatus(chainHead int64) {
func (repository *InMemory) CreateLogs(logs []core.Log) error {
for _, log := range logs {
key := fmt.Sprintf("%s%s", log.BlockNumber, log.Index)
key := fmt.Sprintf("%d%d", log.BlockNumber, log.Index)
var logs []core.Log
repository.logs[key] = append(logs, log)
}
return nil
}
func (repository *InMemory) FindLogs(address string, blockNumber int64) []core.Log {
func (repository *InMemory) GetLogs(address string, blockNumber int64) []core.Log {
var matchingLogs []core.Log
for _, logs := range repository.logs {
for _, log := range logs {
@@ -91,7 +99,7 @@ func (repository *InMemory) ContractExists(contractHash string) bool {
return present
}
func (repository *InMemory) FindContract(contractHash string) (core.Contract, error) {
func (repository *InMemory) GetContract(contractHash string) (core.Contract, error) {
contract, ok := repository.contracts[contractHash]
if !ok {
return core.Contract{}, repositories.ErrContractDoesNotExist(contractHash)
@@ -130,7 +138,7 @@ func (repository *InMemory) BlockCount() int {
return len(repository.blocks)
}
func (repository *InMemory) FindBlockByNumber(blockNumber int64) (core.Block, error) {
func (repository *InMemory) GetBlock(blockNumber int64) (core.Block, error) {
if block, ok := repository.blocks[blockNumber]; ok {
return block, nil
}
@@ -56,7 +56,7 @@ func (db DB) MissingBlockNumbers(startingBlockNumber int64, highestBlockNumber i
return numbers
}
func (db DB) FindBlockByNumber(blockNumber int64) (core.Block, error) {
func (db DB) GetBlock(blockNumber int64) (core.Block, error) {
blockRows := db.DB.QueryRowx(
`SELECT id,
number,
@@ -36,7 +36,7 @@ var _ = Describe("Saving blocks", func() {
}
repositoryTwo := postgres.BuildRepository(nodeTwo)
_, err := repositoryTwo.FindBlockByNumber(123)
_, err := repositoryTwo.GetBlock(123)
Expect(err).To(HaveOccurred())
})
@@ -74,7 +74,7 @@ var _ = Describe("Saving blocks", func() {
repository.CreateOrUpdateBlock(block)
savedBlock, err := repository.FindBlockByNumber(blockNumber)
savedBlock, err := repository.GetBlock(blockNumber)
Expect(err).NotTo(HaveOccurred())
Expect(savedBlock.Reward).To(Equal(blockReward))
Expect(savedBlock.Difficulty).To(Equal(difficulty))
@@ -93,7 +93,7 @@ var _ = Describe("Saving blocks", func() {
})
It("does not find a block when searching for a number that does not exist", func() {
_, err := repository.FindBlockByNumber(111)
_, err := repository.GetBlock(111)
Expect(err).To(HaveOccurred())
})
@@ -106,7 +106,7 @@ var _ = Describe("Saving blocks", func() {
repository.CreateOrUpdateBlock(block)
savedBlock, _ := repository.FindBlockByNumber(123)
savedBlock, _ := repository.GetBlock(123)
Expect(len(savedBlock.Transactions)).To(Equal(1))
})
@@ -118,7 +118,7 @@ var _ = Describe("Saving blocks", func() {
repository.CreateOrUpdateBlock(block)
savedBlock, _ := repository.FindBlockByNumber(123)
savedBlock, _ := repository.GetBlock(123)
Expect(len(savedBlock.Transactions)).To(Equal(2))
})
@@ -138,7 +138,7 @@ var _ = Describe("Saving blocks", func() {
repository.CreateOrUpdateBlock(blockOne)
repository.CreateOrUpdateBlock(blockTwo)
savedBlock, _ := repository.FindBlockByNumber(123)
savedBlock, _ := repository.GetBlock(123)
Expect(len(savedBlock.Transactions)).To(Equal(2))
Expect(savedBlock.Transactions[0].Hash).To(Equal("x678"))
Expect(savedBlock.Transactions[1].Hash).To(Equal("x9ab"))
@@ -163,8 +163,8 @@ var _ = Describe("Saving blocks", func() {
repository.CreateOrUpdateBlock(blockOne)
repositoryTwo.CreateOrUpdateBlock(blockTwo)
retrievedBlockOne, _ := repository.FindBlockByNumber(123)
retrievedBlockTwo, _ := repositoryTwo.FindBlockByNumber(123)
retrievedBlockOne, _ := repository.GetBlock(123)
retrievedBlockTwo, _ := repositoryTwo.GetBlock(123)
Expect(retrievedBlockOne.Transactions[0].Hash).To(Equal("x123"))
Expect(retrievedBlockTwo.Transactions[0].Hash).To(Equal("x678"))
@@ -196,7 +196,7 @@ var _ = Describe("Saving blocks", func() {
repository.CreateOrUpdateBlock(block)
savedBlock, _ := repository.FindBlockByNumber(123)
savedBlock, _ := repository.GetBlock(123)
Expect(len(savedBlock.Transactions)).To(Equal(1))
savedTransaction := savedBlock.Transactions[0]
Expect(savedTransaction.Data).To(Equal(transaction.Data))
@@ -271,10 +271,10 @@ var _ = Describe("Saving blocks", func() {
repository.SetBlocksStatus(int64(blockNumberOfChainHead))
blockOne, err := repository.FindBlockByNumber(1)
blockOne, err := repository.GetBlock(1)
Expect(err).ToNot(HaveOccurred())
Expect(blockOne.IsFinal).To(Equal(true))
blockTwo, err := repository.FindBlockByNumber(24)
blockTwo, err := repository.GetBlock(24)
Expect(err).ToNot(HaveOccurred())
Expect(blockTwo.IsFinal).To(BeFalse())
})
@@ -36,7 +36,7 @@ func (db DB) ContractExists(contractHash string) bool {
return exists
}
func (db DB) FindContract(contractHash string) (core.Contract, error) {
func (db DB) GetContract(contractHash string) (core.Contract, error) {
var hash string
var abi string
contract := db.DB.QueryRow(
@@ -25,7 +25,7 @@ var _ = Describe("Creating contracts", func() {
It("returns the contract when it exists", func() {
repository.CreateContract(core.Contract{Hash: "x123"})
contract, err := repository.FindContract("x123")
contract, err := repository.GetContract("x123")
Expect(err).NotTo(HaveOccurred())
Expect(contract.Hash).To(Equal("x123"))
@@ -34,13 +34,13 @@ var _ = Describe("Creating contracts", func() {
})
It("returns err if contract does not exist", func() {
_, err := repository.FindContract("x123")
_, err := repository.GetContract("x123")
Expect(err).To(HaveOccurred())
})
It("returns empty array when no transactions 'To' a contract", func() {
repository.CreateContract(core.Contract{Hash: "x123"})
contract, err := repository.FindContract("x123")
contract, err := repository.GetContract("x123")
Expect(err).ToNot(HaveOccurred())
Expect(contract.Transactions).To(BeEmpty())
})
@@ -59,7 +59,7 @@ var _ = Describe("Creating contracts", func() {
blockRepository.CreateOrUpdateBlock(block)
repository.CreateContract(core.Contract{Hash: "x123"})
contract, err := repository.FindContract("x123")
contract, err := repository.GetContract("x123")
Expect(err).ToNot(HaveOccurred())
sort.Slice(contract.Transactions, func(i, j int) bool {
return contract.Transactions[i].Hash < contract.Transactions[j].Hash
@@ -76,7 +76,7 @@ var _ = Describe("Creating contracts", func() {
Abi: "{\"some\": \"json\"}",
Hash: "x123",
})
contract, err := repository.FindContract("x123")
contract, err := repository.GetContract("x123")
Expect(err).ToNot(HaveOccurred())
Expect(contract.Abi).To(Equal("{\"some\": \"json\"}"))
})
@@ -90,7 +90,7 @@ var _ = Describe("Creating contracts", func() {
Abi: "{\"some\": \"different json\"}",
Hash: "x123",
})
contract, err := repository.FindContract("x123")
contract, err := repository.GetContract("x123")
Expect(err).ToNot(HaveOccurred())
Expect(contract.Abi).To(Equal("{\"some\": \"different json\"}"))
})
@@ -1,8 +1,16 @@
package postgres
import "github.com/vulcanize/vulcanizedb/pkg/filters"
import (
"database/sql"
func (db DB) AddFilter(query filters.LogFilter) error {
"encoding/json"
"errors"
"github.com/vulcanize/vulcanizedb/pkg/filters"
"github.com/vulcanize/vulcanizedb/pkg/repositories"
)
func (db DB) CreateFilter(query filters.LogFilter) error {
_, err := db.DB.Exec(
`INSERT INTO log_filters
(name, from_block, to_block, address, topic0, topic1, topic2, topic3)
@@ -13,3 +21,53 @@ func (db DB) AddFilter(query filters.LogFilter) error {
}
return nil
}
func (db DB) GetFilter(name string) (filters.LogFilter, error) {
lf := DBLogFilter{}
err := db.DB.Get(&lf,
`SELECT
id,
name,
from_block,
to_block,
address,
json_build_array(topic0, topic1, topic2, topic3) AS topics
FROM log_filters
WHERE name = $1`, name)
if err != nil {
switch err {
case sql.ErrNoRows:
return filters.LogFilter{}, repositories.ErrFilterDoesNotExist(name)
default:
return filters.LogFilter{}, err
}
}
dbLogFilterToCoreLogFilter(lf)
return *lf.LogFilter, nil
}
type DBTopics []*string
func (t *DBTopics) Scan(src interface{}) error {
asBytes, ok := src.([]byte)
if !ok {
return error(errors.New("scan source was not []byte"))
}
json.Unmarshal(asBytes, &t)
return nil
}
type DBLogFilter struct {
ID int
*filters.LogFilter
Topics DBTopics
}
func dbLogFilterToCoreLogFilter(lf DBLogFilter) {
for i, v := range lf.Topics {
if v != nil {
lf.LogFilter.Topics[i] = *v
}
}
}
@@ -37,7 +37,7 @@ var _ = Describe("Logs Repository", func() {
"",
},
}
err := repository.AddFilter(logFilter)
err := repository.CreateFilter(logFilter)
Expect(err).ToNot(HaveOccurred())
})
@@ -54,8 +54,52 @@ var _ = Describe("Logs Repository", func() {
"",
},
}
err := repository.AddFilter(logFilter)
err := repository.CreateFilter(logFilter)
Expect(err).To(HaveOccurred())
})
It("gets a log filter", func() {
logFilter1 := filters.LogFilter{
Name: "TestFilter1",
FromBlock: 1,
ToBlock: 2,
Address: "0x8888f1f195afa192cfee860698584c030f4c9db1",
Topics: core.Topics{
"0x000000000000000000000000a94f5374fce5edbc8e2a8697c15331677e6ebf0b",
"",
"0x000000000000000000000000a94f5374fce5edbc8e2a8697c15331677e6ebf0b",
"",
},
}
err := repository.CreateFilter(logFilter1)
Expect(err).ToNot(HaveOccurred())
logFilter2 := filters.LogFilter{
Name: "TestFilter2",
FromBlock: 10,
ToBlock: 20,
Address: "0x8888f1f195afa192cfee860698584c030f4c9db1",
Topics: core.Topics{
"0x000000000000000000000000a94f5374fce5edbc8e2a8697c15331677e6ebf0b",
"",
"0x000000000000000000000000a94f5374fce5edbc8e2a8697c15331677e6ebf0b",
"",
},
}
err = repository.CreateFilter(logFilter2)
Expect(err).ToNot(HaveOccurred())
logFilter1, err = repository.GetFilter("TestFilter1")
Expect(err).ToNot(HaveOccurred())
Expect(logFilter1).To(Equal(logFilter1))
logFilter1, err = repository.GetFilter("TestFilter1")
Expect(err).ToNot(HaveOccurred())
Expect(logFilter2).To(Equal(logFilter2))
})
It("returns ErrFilterDoesNotExist error when log does not exist", func() {
_, err := repository.GetFilter("TestFilter1")
Expect(err).To(Equal(repositories.ErrFilterDoesNotExist("TestFilter1")))
})
})
})
+1 -1
View File
@@ -26,7 +26,7 @@ func (db DB) CreateLogs(logs []core.Log) error {
return nil
}
func (db DB) FindLogs(address string, blockNumber int64) []core.Log {
func (db DB) GetLogs(address string, blockNumber int64) []core.Log {
logRows, _ := db.DB.Query(
`SELECT block_number,
address,
+4 -4
View File
@@ -35,7 +35,7 @@ var _ = Describe("Logs Repository", func() {
}},
)
log := repository.FindLogs("x123", 1)
log := repository.GetLogs("x123", 1)
Expect(log).NotTo(BeNil())
Expect(log[0].BlockNumber).To(Equal(int64(1)))
@@ -49,7 +49,7 @@ var _ = Describe("Logs Repository", func() {
})
It("returns nil if log does not exist", func() {
log := repository.FindLogs("x123", 1)
log := repository.GetLogs("x123", 1)
Expect(log).To(BeNil())
})
@@ -82,7 +82,7 @@ var _ = Describe("Logs Repository", func() {
}},
)
log := repository.FindLogs("x123", 1)
log := repository.GetLogs("x123", 1)
type logIndex struct {
blockNumber int64
@@ -168,7 +168,7 @@ var _ = Describe("Logs Repository", func() {
block := core.Block{Transactions: []core.Transaction{transaction}}
err := blockRepository.CreateOrUpdateBlock(block)
Expect(err).To(Not(HaveOccurred()))
retrievedLogs := repository.FindLogs("0x99041f808d598b782d5a3e498681c2452a31da08", 4745407)
retrievedLogs := repository.GetLogs("0x99041f808d598b782d5a3e498681c2452a31da08", 4745407)
expected := logs[1:]
Expect(retrievedLogs).To(Equal(expected))
+3 -3
View File
@@ -94,7 +94,7 @@ var _ = Describe("Postgres repository", func() {
repository, _ := postgres.NewDB(cfg.Database, node)
err1 := repository.CreateOrUpdateBlock(badBlock)
savedBlock, err2 := repository.FindBlockByNumber(123)
savedBlock, err2 := repository.GetBlock(123)
Expect(err1).To(HaveOccurred())
Expect(err2).To(HaveOccurred())
@@ -129,7 +129,7 @@ var _ = Describe("Postgres repository", func() {
repository, _ := postgres.NewDB(cfg.Database, node)
err := repository.CreateLogs([]core.Log{badLog})
savedBlock := repository.FindLogs("x123", 1)
savedBlock := repository.GetLogs("x123", 1)
Expect(err).ToNot(BeNil())
Expect(savedBlock).To(BeNil())
@@ -148,7 +148,7 @@ var _ = Describe("Postgres repository", func() {
repository, _ := postgres.NewDB(cfg.Database, node)
err1 := repository.CreateOrUpdateBlock(block)
savedBlock, err2 := repository.FindBlockByNumber(123)
savedBlock, err2 := repository.GetBlock(123)
Expect(err1).To(HaveOccurred())
Expect(err2).To(HaveOccurred())
@@ -7,7 +7,7 @@ import (
"github.com/vulcanize/vulcanizedb/pkg/repositories"
)
func (db DB) FindReceipt(txHash string) (core.Receipt, error) {
func (db DB) GetReceipt(txHash string) (core.Receipt, error) {
row := db.DB.QueryRow(
`SELECT contract_address,
tx_hash,
@@ -42,7 +42,7 @@ var _ = Describe("Logs Repository", func() {
block := core.Block{Transactions: []core.Transaction{transaction}}
blockRepository.CreateOrUpdateBlock(block)
receipt, err := repository.FindReceipt("0xe340558980f89d5f86045ac11e5cc34e4bcec20f9f1e2a427aa39d87114e8223")
receipt, err := repository.GetReceipt("0xe340558980f89d5f86045ac11e5cc34e4bcec20f9f1e2a427aa39d87114e8223")
Expect(err).ToNot(HaveOccurred())
//Not currently serializing bloom logs
@@ -55,7 +55,7 @@ var _ = Describe("Logs Repository", func() {
})
It("returns ErrReceiptDoesNotExist when receipt does not exist", func() {
receipt, err := repository.FindReceipt("DOES NOT EXIST")
receipt, err := repository.GetReceipt("DOES NOT EXIST")
Expect(err).To(HaveOccurred())
Expect(receipt).To(BeZero())
})
@@ -76,7 +76,7 @@ var _ = Describe("Logs Repository", func() {
}
blockRepository.CreateOrUpdateBlock(block)
_, err := repository.FindReceipt(receipt.TxHash)
_, err := repository.GetReceipt(receipt.TxHash)
Expect(err).To(Not(HaveOccurred()))
})
+7 -20
View File
@@ -1,32 +1,19 @@
package postgres
type WatchedEventLog struct {
Name string `json:"name"` // name
BlockNumber int64 `json:"block_number" db:"block_number"` // block_number
Address string `json:"address"` // address
TxHash string `json:"tx_hash" db:"tx_hash"` // tx_hash
Index int64 `json:"index"` // index
Topic0 string `json:"topic0"` // topic0
Topic1 string `json:"topic1"` // topic1
Topic2 string `json:"topic2"` // topic2
Topic3 string `json:"topic3"` // topic3
Data string `json:"data"` // data
}
import (
"github.com/vulcanize/vulcanizedb/pkg/core"
)
type WatchedEventLogs interface {
AllWatchedEventLogs() ([]*WatchedEventLog, error)
}
func (db *DB) AllWatchedEventLogs() ([]*WatchedEventLog, error) {
rows, err := db.DB.Queryx(`SELECT name, block_number, address, tx_hash, index, topic0, topic1, topic2, topic3, data FROM watched_event_logs`)
func (db DB) GetWatchedEvents(name string) ([]*core.WatchedEvent, error) {
rows, err := db.DB.Queryx(`SELECT name, block_number, address, tx_hash, index, topic0, topic1, topic2, topic3, data FROM watched_event_logs where name=$1`, name)
if err != nil {
return nil, err
}
defer rows.Close()
lgs := make([]*WatchedEventLog, 0)
lgs := make([]*core.WatchedEvent, 0)
for rows.Next() {
lg := new(WatchedEventLog)
lg := new(core.WatchedEvent)
err := rows.StructScan(lg)
if err != nil {
return nil, err
@@ -27,7 +27,7 @@ var _ = Describe("Watched Events Repository", func() {
postgres.ClearData(repository)
})
It("retrieves watched logs that match the event filter", func() {
It("retrieves watched event logs that match the event filter", func() {
filter := filters.LogFilter{
Name: "Filter1",
FromBlock: 0,
@@ -45,7 +45,7 @@ var _ = Describe("Watched Events Repository", func() {
Data: "",
},
}
expectedWatchedEventLog := []*postgres.WatchedEventLog{
expectedWatchedEventLog := []*core.WatchedEvent{
{
Name: "Filter1",
BlockNumber: 0,
@@ -57,11 +57,57 @@ var _ = Describe("Watched Events Repository", func() {
Data: "",
},
}
err := repository.AddFilter(filter)
err := repository.CreateFilter(filter)
Expect(err).ToNot(HaveOccurred())
err = repository.CreateLogs(logs)
Expect(err).ToNot(HaveOccurred())
matchingLogs, err := repository.AllWatchedEventLogs()
matchingLogs, err := repository.GetWatchedEvents("Filter1")
Expect(err).ToNot(HaveOccurred())
Expect(matchingLogs).To(Equal(expectedWatchedEventLog))
})
It("retrieves a watched event log by name", func() {
filter := filters.LogFilter{
Name: "Filter1",
FromBlock: 0,
ToBlock: 10,
Address: "0x123",
Topics: core.Topics{0: "event1=10", 2: "event3=hello"},
}
logs := []core.Log{
{
BlockNumber: 0,
TxHash: "0x1",
Address: "0x123",
Topics: core.Topics{0: "event1=10", 2: "event3=hello"},
Index: 0,
Data: "",
},
{
BlockNumber: 100,
TxHash: "",
Address: "",
Topics: core.Topics{},
Index: 0,
Data: "",
},
}
expectedWatchedEventLog := []*core.WatchedEvent{{
Name: "Filter1",
BlockNumber: 0,
TxHash: "0x1",
Address: "0x123",
Topic0: "event1=10",
Topic2: "event3=hello",
Index: 0,
Data: "",
}}
err := repository.CreateFilter(filter)
Expect(err).ToNot(HaveOccurred())
err = repository.CreateLogs(logs)
Expect(err).ToNot(HaveOccurred())
matchingLogs, err := repository.GetWatchedEvents("Filter1")
Expect(err).ToNot(HaveOccurred())
Expect(matchingLogs).To(Equal(expectedWatchedEventLog))
+15 -5
View File
@@ -14,6 +14,7 @@ type Repository interface {
LogsRepository
ReceiptRepository
FilterRepository
WatchedEventsRepository
}
var ErrBlockDoesNotExist = func(blockNumber int64) error {
@@ -22,7 +23,7 @@ var ErrBlockDoesNotExist = func(blockNumber int64) error {
type BlockRepository interface {
CreateOrUpdateBlock(block core.Block) error
FindBlockByNumber(blockNumber int64) (core.Block, error)
GetBlock(blockNumber int64) (core.Block, error)
MissingBlockNumbers(startingBlockNumber int64, endingBlockNumber int64) []int64
SetBlocksStatus(chainHead int64)
}
@@ -33,17 +34,22 @@ var ErrContractDoesNotExist = func(contractHash string) error {
type ContractRepository interface {
CreateContract(contract core.Contract) error
GetContract(contractHash string) (core.Contract, error)
ContractExists(contractHash string) bool
FindContract(contractHash string) (core.Contract, error)
}
var ErrFilterDoesNotExist = func(name string) error {
return errors.New(fmt.Sprintf("filter %s does not exist", name))
}
type FilterRepository interface {
AddFilter(filter filters.LogFilter) error
CreateFilter(filter filters.LogFilter) error
GetFilter(name string) (filters.LogFilter, error)
}
type LogsRepository interface {
FindLogs(address string, blockNumber int64) []core.Log
CreateLogs(logs []core.Log) error
GetLogs(address string, blockNumber int64) []core.Log
}
var ErrReceiptDoesNotExist = func(txHash string) error {
@@ -51,5 +57,9 @@ var ErrReceiptDoesNotExist = func(txHash string) error {
}
type ReceiptRepository interface {
FindReceipt(txHash string) (core.Receipt, error)
GetReceipt(txHash string) (core.Receipt, error)
}
type WatchedEventsRepository interface {
GetWatchedEvents(name string) ([]*core.WatchedEvent, error)
}