Add filterLogs flag to fetch logs by contract (#144)

This commit is contained in:
nikugogoi 2022-07-18 14:04:59 +05:30 committed by GitHub
parent 9a35955166
commit e2668690b5
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
33 changed files with 157 additions and 47 deletions

View File

@ -154,6 +154,10 @@ export class Indexer implements IPLDIndexerInterface {
this._populateRelationsMap(); this._populateRelationsMap();
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
await this._baseIndexer.fetchIPLDStatus(); await this._baseIndexer.fetchIPLDStatus();

View File

@ -191,6 +191,10 @@ export class Indexer implements IPLDIndexerInterface {
this._populateRelationsMap(); this._populateRelationsMap();
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
await this._baseIndexer.fetchIPLDStatus(); await this._baseIndexer.fetchIPLDStatus();

View File

@ -43,7 +43,7 @@ export const handler = async (argv: any): Promise<void> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue, config.server.mode); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
const syncStatus = await indexer.getSyncStatus(); const syncStatus = await indexer.getSyncStatus();
assert(syncStatus, 'Missing syncStatus'); assert(syncStatus, 'Missing syncStatus');

View File

@ -43,7 +43,7 @@ import { CONTRACT_KIND } from '../utils/index';
}).argv; }).argv;
const config: Config = await getConfig(argv.configFile); const config: Config = await getConfig(argv.configFile);
const { database: dbConfig, server: { mode }, jobQueue: jobQueueConfig } = config; const { database: dbConfig, server, jobQueue: jobQueueConfig } = config;
const { ethClient, ethProvider } = await initClients(config); const { ethClient, ethProvider } = await initClients(config);
assert(dbConfig); assert(dbConfig);
@ -59,7 +59,7 @@ import { CONTRACT_KIND } from '../utils/index';
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue, mode); const indexer = new Indexer(server, db, ethClient, ethProvider, jobQueue);
await indexer.watchContract(argv.address, CONTRACT_KIND, argv.checkpoint, argv.startingBlock); await indexer.watchContract(argv.address, CONTRACT_KIND, argv.checkpoint, argv.startingBlock);

View File

@ -72,7 +72,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue, config.server.mode); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue); const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue);

View File

@ -12,7 +12,7 @@ import { BaseProvider } from '@ethersproject/providers';
import { EthClient } from '@vulcanize/ipld-eth-client'; import { EthClient } from '@vulcanize/ipld-eth-client';
import { StorageLayout } from '@vulcanize/solidity-mapper'; import { StorageLayout } from '@vulcanize/solidity-mapper';
import { IndexerInterface, Indexer as BaseIndexer, ValueResult, UNKNOWN_EVENT_NAME, JobQueue, Where, QueryOptions } from '@vulcanize/util'; import { IndexerInterface, Indexer as BaseIndexer, ValueResult, UNKNOWN_EVENT_NAME, JobQueue, Where, QueryOptions, ServerConfig } from '@vulcanize/util';
import { Database } from './database'; import { Database } from './database';
import { Event } from './entity/Event'; import { Event } from './entity/Event';
@ -46,20 +46,22 @@ export class Indexer implements IndexerInterface {
_ethClient: EthClient _ethClient: EthClient
_ethProvider: BaseProvider _ethProvider: BaseProvider
_baseIndexer: BaseIndexer _baseIndexer: BaseIndexer
_serverConfig: ServerConfig
_abi: JsonFragment[] _abi: JsonFragment[]
_storageLayout: StorageLayout _storageLayout: StorageLayout
_contract: ethers.utils.Interface _contract: ethers.utils.Interface
_serverMode: string _serverMode: string
constructor (db: Database, ethClient: EthClient, ethProvider: BaseProvider, jobQueue: JobQueue, serverMode: string) { constructor (serverConfig: ServerConfig, db: Database, ethClient: EthClient, ethProvider: BaseProvider, jobQueue: JobQueue) {
assert(db); assert(db);
assert(ethClient); assert(ethClient);
this._db = db; this._db = db;
this._ethClient = ethClient; this._ethClient = ethClient;
this._ethProvider = ethProvider; this._ethProvider = ethProvider;
this._serverMode = serverMode; this._serverConfig = serverConfig;
this._serverMode = serverConfig.mode;
this._baseIndexer = new BaseIndexer(this._db, this._ethClient, this._ethProvider, jobQueue); this._baseIndexer = new BaseIndexer(this._db, this._ethClient, this._ethProvider, jobQueue);
const { abi, storageLayout } = artifacts; const { abi, storageLayout } = artifacts;
@ -73,6 +75,10 @@ export class Indexer implements IndexerInterface {
this._contract = new ethers.utils.Interface(this._abi); this._contract = new ethers.utils.Interface(this._abi);
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
} }

View File

@ -82,7 +82,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue, config.server.mode); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
const jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue); const jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue);

View File

@ -38,7 +38,7 @@ export const main = async (): Promise<any> => {
const config: Config = await getConfig(argv.f); const config: Config = await getConfig(argv.f);
const { ethClient, ethProvider } = await initClients(config); const { ethClient, ethProvider } = await initClients(config);
const { host, port, mode, kind: watcherKind } = config.server; const { host, port, kind: watcherKind } = config.server;
const db = new Database(config.database); const db = new Database(config.database);
await db.init(); await db.init();
@ -55,7 +55,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue, mode); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue); const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue);

View File

@ -137,6 +137,10 @@ export class Indexer implements IPLDIndexerInterface {
this._relationsMap = new Map(); this._relationsMap = new Map();
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
await this._baseIndexer.fetchIPLDStatus(); await this._baseIndexer.fetchIPLDStatus();

View File

@ -5,10 +5,15 @@ import {
IndexerInterface, IndexerInterface,
BlockProgressInterface, BlockProgressInterface,
EventInterface, EventInterface,
SyncStatusInterface SyncStatusInterface,
ServerConfig as ServerConfigInterface
} from '@vulcanize/util'; } from '@vulcanize/util';
export class Indexer implements IndexerInterface { export class Indexer implements IndexerInterface {
get serverConfig () {
return new ServerConfig();
}
async getBlockProgress (blockHash: string): Promise<BlockProgressInterface | undefined> { async getBlockProgress (blockHash: string): Promise<BlockProgressInterface | undefined> {
assert(blockHash); assert(blockHash);
@ -140,3 +145,29 @@ class SyncStatus implements SyncStatusInterface {
this.latestCanonicalBlockNumber = 0; this.latestCanonicalBlockNumber = 0;
} }
} }
class ServerConfig implements ServerConfigInterface {
host: string;
port: number;
mode: string;
kind: string;
checkpointing: boolean;
checkpointInterval: number;
ipfsApiAddr: string;
subgraphPath: string;
wasmRestartBlocksInterval: number;
filterLogs: boolean;
constructor () {
this.host = '';
this.port = 0;
this.mode = '';
this.kind = '';
this.checkpointing = false;
this.checkpointInterval = 0;
this.ipfsApiAddr = '';
this.subgraphPath = '';
this.wasmRestartBlocksInterval = 0;
this.filterLogs = false;
}
}

View File

@ -160,6 +160,10 @@ export class Indexer implements IPLDIndexerInterface {
this._populateRelationsMap(); this._populateRelationsMap();
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
await this._baseIndexer.fetchIPLDStatus(); await this._baseIndexer.fetchIPLDStatus();

View File

@ -12,6 +12,9 @@
# IPFS API address (can be taken from the output on running the IPFS daemon). # IPFS API address (can be taken from the output on running the IPFS daemon).
# ipfsApiAddr = "/ip4/127.0.0.1/tcp/5001" # ipfsApiAddr = "/ip4/127.0.0.1/tcp/5001"
# Boolean to filter logs by contract.
filterLogs = true
[database] [database]
type = "postgres" type = "postgres"
host = "localhost" host = "localhost"

View File

@ -138,8 +138,14 @@ main().catch(err => {
const processEvent = async (indexer: Indexer, block: BlockProgress, event: Event) => { const processEvent = async (indexer: Indexer, block: BlockProgress, event: Event) => {
const eventIndex = event.index; const eventIndex = event.index;
// Check that events are processed in order.
if (eventIndex <= block.lastProcessedEventIndex) {
throw new Error(`Events received out of order for block number ${block.blockNumber} hash ${block.blockHash}, got event index ${eventIndex} and lastProcessedEventIndex ${block.lastProcessedEventIndex}, aborting`);
}
// Check if previous event in block has been processed exactly before this and abort if not. // Check if previous event in block has been processed exactly before this and abort if not.
if (eventIndex > 0) { // Skip the first event in the block. // Skip check if logs fetched are filtered by contract address.
if (!indexer.serverConfig.filterLogs) {
const prevIndex = eventIndex - 1; const prevIndex = eventIndex - 1;
if (prevIndex !== block.lastProcessedEventIndex) { if (prevIndex !== block.lastProcessedEventIndex) {

View File

@ -137,6 +137,10 @@ export class Indexer implements IPLDIndexerInterface {
this._relationsMap = new Map(); this._relationsMap = new Map();
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
await this._baseIndexer.fetchIPLDStatus(); await this._baseIndexer.fetchIPLDStatus();
@ -799,13 +803,33 @@ export class Indexer implements IPLDIndexerInterface {
async _fetchAndSaveEvents ({ cid: blockCid, blockHash }: DeepPartial<BlockProgress>): Promise<BlockProgress> { async _fetchAndSaveEvents ({ cid: blockCid, blockHash }: DeepPartial<BlockProgress>): Promise<BlockProgress> {
assert(blockHash); assert(blockHash);
let block: any, logs: any[];
const logsPromise = this._ethClient.getLogs({ blockHash }); if (this._serverConfig.filterLogs) {
const transactionsPromise = this._ethClient.getBlockWithTransactions({ blockHash }); const watchedContracts = this._baseIndexer.getWatchedContracts();
let [ // TODO: Query logs by multiple contracts.
{ block, logs }, const contractlogsWithBlockPromises = watchedContracts.map((watchedContract): Promise<any> => this._ethClient.getLogs({
{ blockHash,
contract: watchedContract.address
}));
const contractlogsWithBlock = await Promise.all(contractlogsWithBlockPromises);
// Flatten logs by contract and sort by index.
logs = contractlogsWithBlock.map(data => {
return data.logs;
}).flat()
.sort((a, b) => {
return a.index - b.index;
});
({ block } = await this._ethClient.getBlockByHash(blockHash));
} else {
({ block, logs } = await this._ethClient.getLogs({ blockHash }));
}
const {
allEthHeaderCids: { allEthHeaderCids: {
nodes: [ nodes: [
{ {
@ -815,8 +839,7 @@ export class Indexer implements IPLDIndexerInterface {
} }
] ]
} }
} } = await this._ethClient.getBlockWithTransactions({ blockHash });
] = await Promise.all([logsPromise, transactionsPromise]);
const transactionMap = transactions.reduce((acc: {[key: string]: any}, transaction: {[key: string]: any}) => { const transactionMap = transactions.reduce((acc: {[key: string]: any}, transaction: {[key: string]: any}) => {
acc[transaction.txHash] = transaction; acc[transaction.txHash] = transaction;

View File

@ -6,7 +6,7 @@
// "incremental": true, /* Enable incremental compilation */ // "incremental": true, /* Enable incremental compilation */
"target": "es5", /* Specify ECMAScript target version: 'ES3' (default), 'ES5', 'ES2015', 'ES2016', 'ES2017', 'ES2018', 'ES2019', 'ES2020', 'ES2021', or 'ESNEXT'. */ "target": "es5", /* Specify ECMAScript target version: 'ES3' (default), 'ES5', 'ES2015', 'ES2016', 'ES2017', 'ES2018', 'ES2019', 'ES2020', 'ES2021', or 'ESNEXT'. */
"module": "commonjs", /* Specify module code generation: 'none', 'commonjs', 'amd', 'system', 'umd', 'es2015', 'es2020', or 'ESNext'. */ "module": "commonjs", /* Specify module code generation: 'none', 'commonjs', 'amd', 'system', 'umd', 'es2015', 'es2020', or 'ESNext'. */
// "lib": [], /* Specify library files to be included in the compilation. */ "lib": ["es2019"], /* Specify library files to be included in the compilation. */
// "allowJs": true, /* Allow javascript files to be compiled. */ // "allowJs": true, /* Allow javascript files to be compiled. */
// "checkJs": true, /* Report errors in .js files. */ // "checkJs": true, /* Report errors in .js files. */
// "jsx": "preserve", /* Specify JSX code generation: 'preserve', 'react-native', 'react', 'react-jsx' or 'react-jsxdev'. */ // "jsx": "preserve", /* Specify JSX code generation: 'preserve', 'react-native', 'react', 'react-jsx' or 'react-jsxdev'. */

View File

@ -61,7 +61,7 @@ export const handler = async (argv: any): Promise<void> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, uniClient, erc20Client, ethClient, ethProvider, jobQueue, config.server.mode); const indexer = new Indexer(config.server, db, uniClient, erc20Client, ethClient, ethProvider, jobQueue);
const syncStatus = await indexer.getSyncStatus(); const syncStatus = await indexer.getSyncStatus();
assert(syncStatus, 'Missing syncStatus'); assert(syncStatus, 'Missing syncStatus');

View File

@ -79,7 +79,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, uniClient, erc20Client, ethClient, ethProvider, jobQueue, config.server.mode); const indexer = new Indexer(config.server, db, uniClient, erc20Client, ethClient, ethProvider, jobQueue);
const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue); const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue);

View File

@ -11,7 +11,7 @@ import { providers, utils, BigNumber } from 'ethers';
import { Client as UniClient } from '@vulcanize/uni-watcher'; import { Client as UniClient } from '@vulcanize/uni-watcher';
import { Client as ERC20Client } from '@vulcanize/erc20-watcher'; import { Client as ERC20Client } from '@vulcanize/erc20-watcher';
import { EthClient } from '@vulcanize/ipld-eth-client'; import { EthClient } from '@vulcanize/ipld-eth-client';
import { IndexerInterface, Indexer as BaseIndexer, QueryOptions, OrderDirection, BlockHeight, Relation, GraphDecimal, JobQueue, Where } from '@vulcanize/util'; import { IndexerInterface, Indexer as BaseIndexer, QueryOptions, OrderDirection, BlockHeight, Relation, GraphDecimal, JobQueue, Where, ServerConfig } from '@vulcanize/util';
import { findEthPerToken, getEthPriceInUSD, getTrackedAmountUSD, sqrtPriceX96ToTokenPrices, WHITELIST_TOKENS } from './utils/pricing'; import { findEthPerToken, getEthPriceInUSD, getTrackedAmountUSD, sqrtPriceX96ToTokenPrices, WHITELIST_TOKENS } from './utils/pricing';
import { updatePoolDayData, updatePoolHourData, updateTickDayData, updateTokenDayData, updateTokenHourData, updateUniswapDayData } from './utils/interval-updates'; import { updatePoolDayData, updatePoolHourData, updateTickDayData, updateTokenDayData, updateTokenHourData, updateUniswapDayData } from './utils/interval-updates';
@ -46,8 +46,9 @@ export class Indexer implements IndexerInterface {
_ethClient: EthClient _ethClient: EthClient
_baseIndexer: BaseIndexer _baseIndexer: BaseIndexer
_isDemo: boolean _isDemo: boolean
_serverConfig: ServerConfig
constructor (db: Database, uniClient: UniClient, erc20Client: ERC20Client, ethClient: EthClient, ethProvider: providers.BaseProvider, jobQueue: JobQueue, mode: string) { constructor (serverConfig: ServerConfig, db: Database, uniClient: UniClient, erc20Client: ERC20Client, ethClient: EthClient, ethProvider: providers.BaseProvider, jobQueue: JobQueue) {
assert(db); assert(db);
assert(uniClient); assert(uniClient);
assert(erc20Client); assert(erc20Client);
@ -57,8 +58,13 @@ export class Indexer implements IndexerInterface {
this._uniClient = uniClient; this._uniClient = uniClient;
this._erc20Client = erc20Client; this._erc20Client = erc20Client;
this._ethClient = ethClient; this._ethClient = ethClient;
this._serverConfig = serverConfig;
this._baseIndexer = new BaseIndexer(this._db, this._ethClient, ethProvider, jobQueue); this._baseIndexer = new BaseIndexer(this._db, this._ethClient, ethProvider, jobQueue);
this._isDemo = mode === 'demo'; this._isDemo = serverConfig.mode === 'demo';
}
get serverConfig () {
return this._serverConfig;
} }
getResultEvent (event: Event): ResultEvent { getResultEvent (event: Event): ResultEvent {

View File

@ -95,7 +95,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, uniClient, erc20Client, ethClient, ethProvider, jobQueue, config.server.mode); const indexer = new Indexer(config.server, db, uniClient, erc20Client, ethClient, ethProvider, jobQueue);
const jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue); const jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue);
await jobRunner.start(); await jobRunner.start();

View File

@ -60,7 +60,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, uniClient, erc20Client, ethClient, ethProvider, jobQueue, mode); const indexer = new Indexer(config.server, db, uniClient, erc20Client, ethClient, ethProvider, jobQueue);
const pubSub = new PubSub(); const pubSub = new PubSub();
const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubSub, jobQueue); const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubSub, jobQueue);

View File

@ -61,7 +61,7 @@ describe('chain pruning', () => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
indexer = new Indexer(db, ethClient, ethProvider, jobQueue); indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
assert(indexer, 'Could not create indexer object.'); assert(indexer, 'Could not create indexer object.');
jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue); jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue);

View File

@ -40,7 +40,7 @@ export const handler = async (argv: any): Promise<void> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
const syncStatus = await indexer.getSyncStatus(); const syncStatus = await indexer.getSyncStatus();
assert(syncStatus, 'Missing syncStatus'); assert(syncStatus, 'Missing syncStatus');

View File

@ -64,7 +64,7 @@ import { Indexer } from '../indexer';
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
await indexer.watchContract(argv.address, argv.kind, argv.checkpoint, argv.startingBlock); await indexer.watchContract(argv.address, argv.kind, argv.checkpoint, argv.startingBlock);

View File

@ -72,7 +72,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue); const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue);

View File

@ -3,13 +3,13 @@
// //
import debug from 'debug'; import debug from 'debug';
import { DeepPartial, FindConditions, FindManyOptions, QueryRunner } from 'typeorm'; import { DeepPartial, FindConditions, FindManyOptions, QueryRunner, Server } from 'typeorm';
import JSONbig from 'json-bigint'; import JSONbig from 'json-bigint';
import { ethers } from 'ethers'; import { ethers } from 'ethers';
import assert from 'assert'; import assert from 'assert';
import { EthClient } from '@vulcanize/ipld-eth-client'; import { EthClient } from '@vulcanize/ipld-eth-client';
import { IndexerInterface, Indexer as BaseIndexer, ValueResult, JobQueue, Where, QueryOptions } from '@vulcanize/util'; import { IndexerInterface, Indexer as BaseIndexer, ValueResult, JobQueue, Where, QueryOptions, ServerConfig } from '@vulcanize/util';
import { Database } from './database'; import { Database } from './database';
import { Event, UNKNOWN_EVENT_NAME } from './entity/Event'; import { Event, UNKNOWN_EVENT_NAME } from './entity/Event';
@ -40,15 +40,17 @@ export class Indexer implements IndexerInterface {
_ethClient: EthClient _ethClient: EthClient
_baseIndexer: BaseIndexer _baseIndexer: BaseIndexer
_ethProvider: ethers.providers.BaseProvider _ethProvider: ethers.providers.BaseProvider
_serverConfig: ServerConfig
_factoryContract: ethers.utils.Interface _factoryContract: ethers.utils.Interface
_poolContract: ethers.utils.Interface _poolContract: ethers.utils.Interface
_nfpmContract: ethers.utils.Interface _nfpmContract: ethers.utils.Interface
constructor (db: Database, ethClient: EthClient, ethProvider: ethers.providers.BaseProvider, jobQueue: JobQueue) { constructor (serverConfig: ServerConfig, db: Database, ethClient: EthClient, ethProvider: ethers.providers.BaseProvider, jobQueue: JobQueue) {
this._db = db; this._db = db;
this._ethClient = ethClient; this._ethClient = ethClient;
this._ethProvider = ethProvider; this._ethProvider = ethProvider;
this._serverConfig = serverConfig;
this._baseIndexer = new BaseIndexer(this._db, this._ethClient, this._ethProvider, jobQueue); this._baseIndexer = new BaseIndexer(this._db, this._ethClient, this._ethProvider, jobQueue);
this._factoryContract = new ethers.utils.Interface(factoryABI); this._factoryContract = new ethers.utils.Interface(factoryABI);
@ -56,6 +58,10 @@ export class Indexer implements IndexerInterface {
this._nfpmContract = new ethers.utils.Interface(nfpmABI); this._nfpmContract = new ethers.utils.Interface(nfpmABI);
} }
get serverConfig () {
return this._serverConfig;
}
async init (): Promise<void> { async init (): Promise<void> {
await this._baseIndexer.fetchContracts(); await this._baseIndexer.fetchContracts();
} }

View File

@ -84,7 +84,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
const jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue); const jobRunner = new JobRunner(jobQueueConfig, indexer, jobQueue);

View File

@ -56,7 +56,7 @@ export const main = async (): Promise<any> => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue); const eventWatcher = new EventWatcher(config.upstream, ethClient, indexer, pubsub, jobQueue);

View File

@ -127,7 +127,7 @@ describe('uni-watcher', () => {
factory = new Contract(factoryContract.address, FACTORY_ABI, signer); factory = new Contract(factoryContract.address, FACTORY_ABI, signer);
// Verifying with the db. // Verifying with the db.
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
assert(await indexer.isWatchedContract(factory.address), 'Factory contract not added to the database.'); assert(await indexer.isWatchedContract(factory.address), 'Factory contract not added to the database.');
}); });
@ -263,7 +263,7 @@ describe('uni-watcher', () => {
nfpm = new Contract(nfpmContract.address, NFPM_ABI, signer); nfpm = new Contract(nfpmContract.address, NFPM_ABI, signer);
// Verifying with the db. // Verifying with the db.
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
await indexer.init(); await indexer.init();
assert(await indexer.isWatchedContract(nfpm.address), 'NFPM contract not added to the database.'); assert(await indexer.isWatchedContract(nfpm.address), 'NFPM contract not added to the database.');
}); });

View File

@ -81,7 +81,7 @@ const main = async () => {
const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs }); const jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
await jobQueue.start(); await jobQueue.start();
const indexer = new Indexer(db, ethClient, ethProvider, jobQueue); const indexer = new Indexer(config.server, db, ethClient, ethProvider, jobQueue);
let factory: Contract; let factory: Contract;
// Checking whether factory is deployed. // Checking whether factory is deployed.

View File

@ -34,6 +34,7 @@ export interface ServerConfig {
ipfsApiAddr: string; ipfsApiAddr: string;
subgraphPath: string; subgraphPath: string;
wasmRestartBlocksInterval: number; wasmRestartBlocksInterval: number;
filterLogs: boolean;
} }
export interface UpstreamConfig { export interface UpstreamConfig {

View File

@ -315,6 +315,10 @@ export class Indexer {
return watchedContracts; return watchedContracts;
} }
getWatchedContracts (): ContractInterface[] {
return Object.values(this._watchedContracts);
}
async watchContract (address: string, kind: string, checkpoint: boolean, startingBlock: number): Promise<void> { async watchContract (address: string, kind: string, checkpoint: boolean, startingBlock: number): Promise<void> {
assert(this._db.saveContract); assert(this._db.saveContract);
const dbTx = await this._db.createTransactionRunner(); const dbTx = await this._db.createTransactionRunner();

View File

@ -279,8 +279,14 @@ export class JobRunner {
const eventIndex = event.index; const eventIndex = event.index;
// log(`Processing event ${event.id} index ${eventIndex}`); // log(`Processing event ${event.id} index ${eventIndex}`);
// Check that events are processed in order.
if (eventIndex <= block.lastProcessedEventIndex) {
throw new Error(`Events received out of order for block number ${block.blockNumber} hash ${block.blockHash}, got event index ${eventIndex} and lastProcessedEventIndex ${block.lastProcessedEventIndex}, aborting`);
}
// Check if previous event in block has been processed exactly before this and abort if not. // Check if previous event in block has been processed exactly before this and abort if not.
if (eventIndex > 0) { // Skip the first event in the block. // Skip check if logs fetched are filtered by contract address.
if (!this._indexer.serverConfig.filterLogs) {
const prevIndex = eventIndex - 1; const prevIndex = eventIndex - 1;
if (prevIndex !== block.lastProcessedEventIndex) { if (prevIndex !== block.lastProcessedEventIndex) {

View File

@ -4,6 +4,7 @@
import { Connection, DeepPartial, FindConditions, FindManyOptions, QueryRunner } from 'typeorm'; import { Connection, DeepPartial, FindConditions, FindManyOptions, QueryRunner } from 'typeorm';
import { ServerConfig } from './config';
import { Where, QueryOptions } from './database'; import { Where, QueryOptions } from './database';
import { IpldStatus } from './ipld-indexer'; import { IpldStatus } from './ipld-indexer';
@ -76,6 +77,7 @@ export interface IPLDBlockInterface {
} }
export interface IndexerInterface { export interface IndexerInterface {
readonly serverConfig: ServerConfig
getBlockProgress (blockHash: string): Promise<BlockProgressInterface | undefined> getBlockProgress (blockHash: string): Promise<BlockProgressInterface | undefined>
getBlockProgressEntities (where: FindConditions<BlockProgressInterface>, options: FindManyOptions<BlockProgressInterface>): Promise<BlockProgressInterface[]> getBlockProgressEntities (where: FindConditions<BlockProgressInterface>, options: FindManyOptions<BlockProgressInterface>): Promise<BlockProgressInterface[]>
getEvent (id: string): Promise<EventInterface | undefined> getEvent (id: string): Promise<EventInterface | undefined>