Move event watcher to util (#262)

This commit is contained in:
prathamesh0
2022-11-25 17:19:37 +05:30
committed by GitHub
parent e47aab2ed7
commit 94c8ed9575
89 changed files with 111 additions and 590 deletions
@@ -39,5 +39,6 @@ export const handler = async (argv: any): Promise<void> => {
);
await createCheckpointCmd.initIndexer(Indexer, graphWatcher);
await createCheckpointCmd.exec();
};
@@ -35,5 +35,6 @@ export const handler = async (argv: any): Promise<void> => {
);
await verifyCheckpointCmd.initIndexer(Indexer, graphWatcher);
await verifyCheckpointCmd.exec(graphDb);
};
@@ -27,6 +27,7 @@ const main = async (): Promise<void> => {
);
await exportStateCmd.initIndexer(Indexer, graphWatcher);
await exportStateCmd.exec();
};
@@ -10,7 +10,6 @@ import { getGraphDbAndWatcher } from '@cerc-io/graph-node';
import { Database, ENTITY_QUERY_TYPE_MAP, ENTITY_TO_LATEST_ENTITY_MAP } from '../database';
import { Indexer } from '../indexer';
import { EventWatcher } from '../events';
import { State } from '../entity/State';
const log = debug('vulcanize:import-state');
@@ -28,7 +27,8 @@ export const main = async (): Promise<any> => {
ENTITY_TO_LATEST_ENTITY_MAP
);
await importStateCmd.initIndexer(Indexer, EventWatcher, graphWatcher);
await importStateCmd.initIndexer(Indexer, graphWatcher);
await importStateCmd.exec(State, graphDb);
};
@@ -27,6 +27,7 @@ const main = async (): Promise<void> => {
);
await indexBlockCmd.initIndexer(Indexer, graphWatcher);
await indexBlockCmd.exec();
};
@@ -27,6 +27,7 @@ const main = async (): Promise<void> => {
);
await inspectCIDCmd.initIndexer(Indexer, graphWatcher);
await inspectCIDCmd.exec();
};
@@ -32,5 +32,6 @@ export const handler = async (argv: any): Promise<void> => {
);
await resetWatcherCmd.initIndexer(Indexer, graphWatcher);
await resetWatcherCmd.exec();
};
@@ -27,6 +27,7 @@ const main = async (): Promise<void> => {
);
await watchContractCmd.initIndexer(Indexer, graphWatcher);
await watchContractCmd.exec();
};
-70
View File
@@ -1,70 +0,0 @@
//
// Copyright 2021 Vulcanize, Inc.
//
import assert from 'assert';
import { PubSub } from 'graphql-subscriptions';
import { EthClient } from '@cerc-io/ipld-eth-client';
import {
JobQueue,
EventWatcher as BaseEventWatcher,
EventWatcherInterface,
QUEUE_BLOCK_PROCESSING,
QUEUE_EVENT_PROCESSING,
IndexerInterface
} from '@cerc-io/util';
import { Indexer } from './indexer';
export class EventWatcher implements EventWatcherInterface {
_ethClient: EthClient
_indexer: Indexer
_subscription: ZenObservable.Subscription | undefined
_baseEventWatcher: BaseEventWatcher
_pubsub: PubSub
_jobQueue: JobQueue
constructor (ethClient: EthClient, indexer: IndexerInterface, pubsub: PubSub, jobQueue: JobQueue) {
assert(ethClient);
assert(indexer);
this._ethClient = ethClient;
this._indexer = indexer as Indexer;
this._pubsub = pubsub;
this._jobQueue = jobQueue;
this._baseEventWatcher = new BaseEventWatcher(this._ethClient, this._indexer, this._pubsub, this._jobQueue);
}
getEventIterator (): AsyncIterator<any> {
return this._baseEventWatcher.getEventIterator();
}
getBlockProgressEventIterator (): AsyncIterator<any> {
return this._baseEventWatcher.getBlockProgressEventIterator();
}
async start (): Promise<void> {
assert(!this._subscription, 'subscription already started');
await this.initBlockProcessingOnCompleteHandler();
await this.initEventProcessingOnCompleteHandler();
this._baseEventWatcher.startBlockProcessing();
}
async stop (): Promise<void> {
this._baseEventWatcher.stop();
}
async initBlockProcessingOnCompleteHandler (): Promise<void> {
this._jobQueue.onComplete(QUEUE_BLOCK_PROCESSING, async (job) => {
await this._baseEventWatcher.blockProcessingCompleteHandler(job);
});
}
async initEventProcessingOnCompleteHandler (): Promise<void> {
await this._jobQueue.onComplete(QUEUE_EVENT_PROCESSING, async (job) => {
await this._baseEventWatcher.eventProcessingCompleteHandler(job);
});
}
}
+1 -3
View File
@@ -2,7 +2,6 @@
// Copyright 2021 Vulcanize, Inc.
//
import assert from 'assert';
import 'reflect-metadata';
import debug from 'debug';
@@ -12,7 +11,6 @@ import { getGraphDbAndWatcher } from '@cerc-io/graph-node';
import { Database, ENTITY_QUERY_TYPE_MAP, ENTITY_TO_LATEST_ENTITY_MAP } from './database';
import { Indexer } from './indexer';
import { EventWatcher } from './events';
const log = debug('vulcanize:fill');
@@ -29,7 +27,7 @@ export const main = async (): Promise<any> => {
ENTITY_TO_LATEST_ENTITY_MAP
);
await fillCmd.initIndexer(Indexer, EventWatcher, graphWatcher);
await fillCmd.initIndexer(Indexer, graphWatcher);
// Get contractEntitiesMap required for fill-state
// NOTE: Assuming each entity type is only mapped to a single contract
+2 -4
View File
@@ -17,11 +17,10 @@ import {
getResultState,
setGQLCacheHints,
IndexerInterface,
EventWatcherInterface
EventWatcher
} from '@cerc-io/util';
import { Indexer } from './indexer';
import { EventWatcher } from './events';
import { Author } from './entity/Author';
import { Blog } from './entity/Blog';
@@ -29,9 +28,8 @@ import { Category } from './entity/Category';
const log = debug('vulcanize:resolver');
export const createResolvers = async (indexerArg: IndexerInterface, eventWatcherArg: EventWatcherInterface): Promise<any> => {
export const createResolvers = async (indexerArg: IndexerInterface, eventWatcher: EventWatcher): Promise<any> => {
const indexer = indexerArg as Indexer;
const eventWatcher = eventWatcherArg as EventWatcher;
const gqlCacheConfig = indexer.serverConfig.gqlCache;
+1 -2
View File
@@ -13,7 +13,6 @@ import { getGraphDbAndWatcher } from '@cerc-io/graph-node';
import { createResolvers } from './resolvers';
import { Indexer } from './indexer';
import { Database, ENTITY_QUERY_TYPE_MAP, ENTITY_TO_LATEST_ENTITY_MAP } from './database';
import { EventWatcher } from './events';
const log = debug('vulcanize:server');
@@ -30,7 +29,7 @@ export const main = async (): Promise<any> => {
ENTITY_TO_LATEST_ENTITY_MAP
);
await serverCmd.initIndexer(Indexer, EventWatcher, graphWatcher);
await serverCmd.initIndexer(Indexer, graphWatcher);
const typeDefs = fs.readFileSync(path.join(__dirname, 'schema.gql')).toString();
return serverCmd.exec(createResolvers, typeDefs);