mirror of
https://github.com/cerc-io/watcher-ts
synced 2026-09-12 10:37:32 +00:00
Update codegen with changes implemented in mobymask watcher (#148)
* Update codegen with index-block CLI and remove graph-node * Add filter logs by contract flag * Skip generating GQL API for immutable variables * Add config for maxEventsBlockRange * Add new flags in existing watchers
This commit is contained in:
@@ -17,3 +17,4 @@ export * from './src/graph-decimal';
|
||||
export * from './src/ipld-indexer';
|
||||
export * from './src/ipld-database';
|
||||
export * from './src/ipfs';
|
||||
export * from './src/index-block';
|
||||
|
||||
@@ -1,9 +1,13 @@
|
||||
import debug from 'debug';
|
||||
import assert from 'assert';
|
||||
|
||||
import { JOB_KIND_PRUNE, QUEUE_BLOCK_PROCESSING, JOB_KIND_INDEX } from './constants';
|
||||
import { JOB_KIND_PRUNE, QUEUE_BLOCK_PROCESSING, JOB_KIND_INDEX, UNKNOWN_EVENT_NAME } from './constants';
|
||||
import { JobQueue } from './job-queue';
|
||||
import { IndexerInterface } from './types';
|
||||
import { BlockProgressInterface, IndexerInterface } from './types';
|
||||
import { wait } from './misc';
|
||||
import { OrderDirection } from './database';
|
||||
|
||||
const DEFAULT_EVENTS_IN_BATCH = 50;
|
||||
|
||||
const log = debug('vulcanize:common');
|
||||
|
||||
@@ -98,3 +102,93 @@ export const processBlockByNumber = async (
|
||||
await wait(blockDelayInMilliSecs);
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
* Process events in batches for a block.
|
||||
* @param indexer
|
||||
* @param block
|
||||
* @param eventsInBatch
|
||||
*/
|
||||
export const processBatchEvents = async (indexer: IndexerInterface, block: BlockProgressInterface, eventsInBatch: number): Promise<void> => {
|
||||
// Check if block processing is complete.
|
||||
while (!block.isComplete) {
|
||||
console.time('time:common#processBacthEvents-fetching_events_batch');
|
||||
|
||||
// Fetch events in batches
|
||||
const events = await indexer.getBlockEvents(
|
||||
block.blockHash,
|
||||
{
|
||||
index: [
|
||||
{ value: block.lastProcessedEventIndex + 1, operator: 'gte', not: false }
|
||||
]
|
||||
},
|
||||
{
|
||||
limit: eventsInBatch || DEFAULT_EVENTS_IN_BATCH,
|
||||
orderBy: 'index',
|
||||
orderDirection: OrderDirection.asc
|
||||
}
|
||||
);
|
||||
|
||||
console.timeEnd('time:common#processBacthEvents-fetching_events_batch');
|
||||
|
||||
if (events.length) {
|
||||
log(`Processing events batch from index ${events[0].index} to ${events[0].index + events.length - 1}`);
|
||||
}
|
||||
|
||||
console.time('time:common#processBacthEvents-processing_events_batch');
|
||||
|
||||
for (let event of events) {
|
||||
// Process events in loop
|
||||
|
||||
const eventIndex = event.index;
|
||||
// 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.
|
||||
// Skip check if logs fetched are filtered by contract address.
|
||||
if (!indexer.serverConfig.filterLogs) {
|
||||
const prevIndex = eventIndex - 1;
|
||||
|
||||
if (prevIndex !== block.lastProcessedEventIndex) {
|
||||
throw new Error(`Events received out of order for block number ${block.blockNumber} hash ${block.blockHash},` +
|
||||
` prev event index ${prevIndex}, got event index ${event.index} and lastProcessedEventIndex ${block.lastProcessedEventIndex}, aborting`);
|
||||
}
|
||||
}
|
||||
|
||||
let watchedContract;
|
||||
|
||||
if (!indexer.isWatchedContract) {
|
||||
// uni-info-watcher indexer doesn't have watched contracts implementation.
|
||||
watchedContract = true;
|
||||
} else {
|
||||
watchedContract = await indexer.isWatchedContract(event.contract);
|
||||
}
|
||||
|
||||
if (watchedContract) {
|
||||
// We might not have parsed this event yet. This can happen if the contract was added
|
||||
// as a result of a previous event in the same block.
|
||||
if (event.eventName === UNKNOWN_EVENT_NAME) {
|
||||
const logObj = JSON.parse(event.extraInfo);
|
||||
|
||||
assert(indexer.parseEventNameAndArgs);
|
||||
assert(typeof watchedContract !== 'boolean');
|
||||
const { eventName, eventInfo } = indexer.parseEventNameAndArgs(watchedContract.kind, logObj);
|
||||
|
||||
event.eventName = eventName;
|
||||
event.eventInfo = JSON.stringify(eventInfo);
|
||||
event = await indexer.saveEventEntity(event);
|
||||
}
|
||||
|
||||
await indexer.processEvent(event);
|
||||
}
|
||||
|
||||
block = await indexer.updateBlockProgress(block, event.index);
|
||||
}
|
||||
|
||||
console.timeEnd('time:common#processBacthEvents-processing_events_batch');
|
||||
}
|
||||
};
|
||||
|
||||
@@ -11,7 +11,7 @@ import { ConnectionOptions } from 'typeorm';
|
||||
|
||||
import { Config as CacheConfig, getCache } from '@vulcanize/cache';
|
||||
import { EthClient } from '@vulcanize/ipld-eth-client';
|
||||
import { BaseProvider, JsonRpcProvider } from '@ethersproject/providers';
|
||||
import { JsonRpcProvider } from '@ethersproject/providers';
|
||||
|
||||
import { getCustomProvider } from './misc';
|
||||
|
||||
@@ -35,6 +35,7 @@ export interface ServerConfig {
|
||||
subgraphPath: string;
|
||||
wasmRestartBlocksInterval: number;
|
||||
filterLogs: boolean;
|
||||
maxEventsBlockRange: number;
|
||||
}
|
||||
|
||||
export interface UpstreamConfig {
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
//
|
||||
// Copyright 2022 Vulcanize, Inc.
|
||||
//
|
||||
|
||||
import assert from 'assert';
|
||||
|
||||
import { BlockProgressInterface, IndexerInterface } from './types';
|
||||
import { processBatchEvents } from './common';
|
||||
|
||||
export const indexBlock = async (
|
||||
indexer: IndexerInterface,
|
||||
eventsInBatch: number,
|
||||
argv: {
|
||||
block: number,
|
||||
}
|
||||
): Promise<any> => {
|
||||
let blockProgressEntities: Partial<BlockProgressInterface>[] = await indexer.getBlocksAtHeight(argv.block, false);
|
||||
|
||||
if (!blockProgressEntities.length) {
|
||||
console.time('time:index-block#getBlocks-ipld-eth-server');
|
||||
const blocks = await indexer.getBlocks({ blockNumber: argv.block });
|
||||
|
||||
blockProgressEntities = blocks.map((block: any): Partial<BlockProgressInterface> => {
|
||||
block.blockTimestamp = block.timestamp;
|
||||
|
||||
return block;
|
||||
});
|
||||
|
||||
console.timeEnd('time:index-block#getBlocks-ipld-eth-server');
|
||||
}
|
||||
|
||||
assert(blockProgressEntities.length, `No blocks fetched for block number ${argv.block}.`);
|
||||
|
||||
for (const partialblockProgress of blockProgressEntities) {
|
||||
let blockProgress: BlockProgressInterface;
|
||||
|
||||
// Check if blockProgress fetched from database.
|
||||
if (!partialblockProgress.id) {
|
||||
blockProgress = await indexer.fetchBlockEvents(partialblockProgress);
|
||||
} else {
|
||||
blockProgress = partialblockProgress as BlockProgressInterface;
|
||||
}
|
||||
|
||||
assert(indexer.processBlock);
|
||||
await indexer.processBlock(blockProgress.blockHash, blockProgress.blockNumber);
|
||||
|
||||
await processBatchEvents(indexer, blockProgress, eventsInBatch);
|
||||
}
|
||||
};
|
||||
@@ -13,17 +13,13 @@ import {
|
||||
JOB_KIND_EVENTS,
|
||||
JOB_KIND_CONTRACT,
|
||||
MAX_REORG_DEPTH,
|
||||
UNKNOWN_EVENT_NAME,
|
||||
QUEUE_BLOCK_PROCESSING,
|
||||
QUEUE_EVENT_PROCESSING
|
||||
} from './constants';
|
||||
import { JobQueue } from './job-queue';
|
||||
import { EventInterface, IndexerInterface, IPLDIndexerInterface, SyncStatusInterface } from './types';
|
||||
import { wait } from './misc';
|
||||
import { createPruningJob } from './common';
|
||||
import { OrderDirection } from './database';
|
||||
|
||||
const DEFAULT_EVENTS_IN_BATCH = 50;
|
||||
import { createPruningJob, processBatchEvents } from './common';
|
||||
|
||||
const log = debug('vulcanize:job-runner');
|
||||
|
||||
@@ -241,92 +237,13 @@ export class JobRunner {
|
||||
const { blockHash } = job.data;
|
||||
|
||||
console.time('time:job-runner#_processEvents-get-block-progress');
|
||||
let block = await this._indexer.getBlockProgress(blockHash);
|
||||
const block = await this._indexer.getBlockProgress(blockHash);
|
||||
console.timeEnd('time:job-runner#_processEvents-get-block-progress');
|
||||
assert(block);
|
||||
|
||||
console.time('time:job-runner#_processEvents-events');
|
||||
|
||||
while (!block.isComplete) {
|
||||
console.time('time:job-runner#_processEvents-fetching_events_batch');
|
||||
|
||||
// Fetch events in batches
|
||||
const events: EventInterface[] = await this._indexer.getBlockEvents(
|
||||
blockHash,
|
||||
{
|
||||
index: [
|
||||
{ value: block.lastProcessedEventIndex + 1, operator: 'gte', not: false }
|
||||
]
|
||||
},
|
||||
{
|
||||
limit: this._jobQueueConfig.eventsInBatch || DEFAULT_EVENTS_IN_BATCH,
|
||||
orderBy: 'index',
|
||||
orderDirection: OrderDirection.asc
|
||||
}
|
||||
);
|
||||
|
||||
console.timeEnd('time:job-runner#_processEvents-fetching_events_batch');
|
||||
|
||||
if (events.length) {
|
||||
log(`Processing events batch from index ${events[0].index} to ${events[0].index + events.length - 1}`);
|
||||
}
|
||||
|
||||
console.time('time:job-runner#_processEvents-processing_events_batch');
|
||||
|
||||
for (let event of events) {
|
||||
// Process events in loop
|
||||
|
||||
const eventIndex = event.index;
|
||||
// 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.
|
||||
// Skip check if logs fetched are filtered by contract address.
|
||||
if (!this._indexer.serverConfig.filterLogs) {
|
||||
const prevIndex = eventIndex - 1;
|
||||
|
||||
if (prevIndex !== block.lastProcessedEventIndex) {
|
||||
throw new Error(`Events received out of order for block number ${block.blockNumber} hash ${block.blockHash},` +
|
||||
` prev event index ${prevIndex}, got event index ${event.index} and lastProcessedEventIndex ${block.lastProcessedEventIndex}, aborting`);
|
||||
}
|
||||
}
|
||||
|
||||
let watchedContract;
|
||||
|
||||
if (!this._indexer.isWatchedContract) {
|
||||
// uni-info-watcher indexer doesn't have watched contracts implementation.
|
||||
watchedContract = true;
|
||||
} else {
|
||||
watchedContract = await this._indexer.isWatchedContract(event.contract);
|
||||
}
|
||||
|
||||
if (watchedContract) {
|
||||
// We might not have parsed this event yet. This can happen if the contract was added
|
||||
// as a result of a previous event in the same block.
|
||||
if (event.eventName === UNKNOWN_EVENT_NAME) {
|
||||
const logObj = JSON.parse(event.extraInfo);
|
||||
|
||||
assert(this._indexer.parseEventNameAndArgs);
|
||||
assert(typeof watchedContract !== 'boolean');
|
||||
const { eventName, eventInfo } = this._indexer.parseEventNameAndArgs(watchedContract.kind, logObj);
|
||||
|
||||
event.eventName = eventName;
|
||||
event.eventInfo = JSON.stringify(eventInfo);
|
||||
event = await this._indexer.saveEventEntity(event);
|
||||
}
|
||||
|
||||
await this._indexer.processEvent(event);
|
||||
}
|
||||
|
||||
block = await this._indexer.updateBlockProgress(block, event.index);
|
||||
}
|
||||
|
||||
console.timeEnd('time:job-runner#_processEvents-processing_events_batch');
|
||||
}
|
||||
await processBatchEvents(this._indexer, block, this._jobQueueConfig.eventsInBatch);
|
||||
|
||||
console.timeEnd('time:job-runner#_processEvents-events');
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user