Implement single query for transactions and blockByMhKey
This commit is contained in:
@@ -26,6 +26,7 @@ import (
|
||||
"github.com/jmoiron/sqlx"
|
||||
"github.com/lib/pq"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"github.com/thoas/go-funk"
|
||||
|
||||
"github.com/vulcanize/ipld-eth-server/v3/pkg/shared"
|
||||
)
|
||||
@@ -592,6 +593,17 @@ func (ecr *CIDRetriever) RetrieveTxCIDsByHeaderID(tx *sqlx.Tx, headerID string)
|
||||
return txCIDs, tx.Select(&txCIDs, pgStr, headerID)
|
||||
}
|
||||
|
||||
func (ecr *CIDRetriever) RetrieveTxCIDsByBlockNumber(tx *sqlx.Tx, blockNumber int64) ([]models.TxModel, error) {
|
||||
log.Debug("retrieving tx cids for block number ", blockNumber)
|
||||
pgStr := `SELECT CAST(block_number as Text), header_id, index, tx_hash, cid, mh_key,
|
||||
dst, src, tx_data, tx_type, value
|
||||
FROM eth.transaction_cids
|
||||
WHERE block_number = $1
|
||||
ORDER BY index`
|
||||
var txCIDs []models.TxModel
|
||||
return txCIDs, tx.Select(&txCIDs, pgStr, blockNumber)
|
||||
}
|
||||
|
||||
// RetrieveReceiptCIDsByTxIDs retrieves receipt CIDs by their associated tx IDs
|
||||
func (ecr *CIDRetriever) RetrieveReceiptCIDsByTxIDs(tx *sqlx.Tx, txHashes []string) ([]models.ReceiptModel, error) {
|
||||
log.Debugf("retrieving receipt cids for tx hashes %v", txHashes)
|
||||
@@ -635,13 +647,30 @@ func (ecr *CIDRetriever) RetrieveHeaderAndTxCIDsByBlockNumber(blockNumber int64)
|
||||
}
|
||||
|
||||
var allTxCIDs [][]models.TxModel
|
||||
txCIDs, err := ecr.RetrieveTxCIDsByBlockNumber(tx, blockNumber)
|
||||
if err != nil {
|
||||
log.Error("tx cid retrieval error")
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
txCIDsByHeaderID := funk.Reduce(
|
||||
txCIDs,
|
||||
func(acc map[string][]models.TxModel, txCID models.TxModel) map[string][]models.TxModel {
|
||||
if _, ok := acc[txCID.HeaderID]; !ok {
|
||||
acc[txCID.HeaderID] = []models.TxModel{}
|
||||
}
|
||||
|
||||
txCIDs = append(acc[txCID.HeaderID], txCID)
|
||||
acc[txCID.HeaderID] = txCIDs
|
||||
return acc
|
||||
},
|
||||
make(map[string][]models.TxModel),
|
||||
)
|
||||
|
||||
txCIDsByHeaderIDMap := txCIDsByHeaderID.(map[string][]models.TxModel)
|
||||
|
||||
for _, headerCID := range headerCIDs {
|
||||
var txCIDs []models.TxModel
|
||||
txCIDs, err = ecr.RetrieveTxCIDsByHeaderID(tx, headerCID.BlockHash)
|
||||
if err != nil {
|
||||
log.Error("tx cid retrieval error")
|
||||
return nil, nil, err
|
||||
}
|
||||
txCIDs := txCIDsByHeaderIDMap[headerCID.BlockHash]
|
||||
allTxCIDs = append(allTxCIDs, txCIDs)
|
||||
}
|
||||
|
||||
@@ -676,7 +705,6 @@ func (ecr *CIDRetriever) RetrieveHeaderAndTxCIDsByBlockHash(blockHash common.Has
|
||||
if err != nil {
|
||||
return models.HeaderModel{}, nil, err
|
||||
}
|
||||
fmt.Println("RetrieveHeaderAndTxCIDsByBlockHash", headerCID.ParentHash, headerCID.Timestamp)
|
||||
|
||||
var txCIDs []models.TxModel
|
||||
txCIDs, err = ecr.RetrieveTxCIDsByHeaderID(tx, headerCID.BlockHash)
|
||||
|
||||
+35
-1
@@ -20,11 +20,13 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"math/big"
|
||||
"strconv"
|
||||
|
||||
"github.com/ethereum/go-ethereum/common"
|
||||
"github.com/ethereum/go-ethereum/statediff/indexer/models"
|
||||
"github.com/jmoiron/sqlx"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"github.com/thoas/go-funk"
|
||||
"github.com/vulcanize/ipld-eth-server/v3/pkg/shared"
|
||||
)
|
||||
|
||||
@@ -99,7 +101,7 @@ func (f *IPLDFetcher) Fetch(cids CIDWrapper) (*IPLDs, error) {
|
||||
return iplds, err
|
||||
}
|
||||
|
||||
// FetchHeaders fetches headers
|
||||
// FetchHeader fetches header
|
||||
func (f *IPLDFetcher) FetchHeader(tx *sqlx.Tx, c models.HeaderModel) (models.IPLDModel, error) {
|
||||
log.Debug("fetching header ipld")
|
||||
headerBytes, err := shared.FetchIPLDByMhKey(tx, c.MhKey)
|
||||
@@ -112,6 +114,38 @@ func (f *IPLDFetcher) FetchHeader(tx *sqlx.Tx, c models.HeaderModel) (models.IPL
|
||||
}, nil
|
||||
}
|
||||
|
||||
// FetchHeaders fetches headers
|
||||
func (f *IPLDFetcher) FetchHeaders(tx *sqlx.Tx, cids []models.HeaderModel) ([]models.IPLDModel, error) {
|
||||
log.Debug("fetching header iplds")
|
||||
headerIPLDs := make([]models.IPLDModel, len(cids))
|
||||
|
||||
blockNumbers := make([]uint64, len(cids))
|
||||
mhKeys := make([]string, len(cids))
|
||||
for i, c := range cids {
|
||||
var err error
|
||||
mhKeys[i] = c.MhKey
|
||||
blockNumbers[i], err = strconv.ParseUint(c.BlockNumber, 10, 64)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
fetchedIPLDs, err := shared.FetchIPLDsByMhKeysAndBlockNumbers(tx, mhKeys, blockNumbers)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
for i, c := range cids {
|
||||
headerIPLD := funk.Find(fetchedIPLDs, func(ipld models.IPLDModel) bool {
|
||||
return ipld.Key == c.MhKey
|
||||
}).(models.IPLDModel)
|
||||
|
||||
headerIPLDs[i] = headerIPLD
|
||||
}
|
||||
|
||||
return headerIPLDs, nil
|
||||
}
|
||||
|
||||
// FetchUncles fetches uncles
|
||||
func (f *IPLDFetcher) FetchUncles(tx *sqlx.Tx, cids []models.UncleModel) ([]models.IPLDModel, error) {
|
||||
log.Debug("fetching uncle iplds")
|
||||
|
||||
Reference in New Issue
Block a user