mirror of
https://github.com/cerc-io/watcher-ts
synced 2026-09-05 15:34:07 +00:00
Implement switching endpoints after slow eth_getLogs RPC requests (#525)
* Switch upstream endpoint if getLogs requests are too slow * Refactor methods for switching client to indexer * Update codegen indexer template * Add dummy methods in graph-node test Indexer * Upgrade package versions to 0.2.101 --------- Co-authored-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@cerc-io/cli",
|
||||
"version": "0.2.100",
|
||||
"version": "0.2.101",
|
||||
"main": "dist/index.js",
|
||||
"license": "AGPL-3.0",
|
||||
"scripts": {
|
||||
@@ -15,13 +15,13 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@apollo/client": "^3.7.1",
|
||||
"@cerc-io/cache": "^0.2.100",
|
||||
"@cerc-io/ipld-eth-client": "^0.2.100",
|
||||
"@cerc-io/cache": "^0.2.101",
|
||||
"@cerc-io/ipld-eth-client": "^0.2.101",
|
||||
"@cerc-io/libp2p": "^0.42.2-laconic-0.1.4",
|
||||
"@cerc-io/nitro-node": "^0.1.15",
|
||||
"@cerc-io/peer": "^0.2.100",
|
||||
"@cerc-io/rpc-eth-client": "^0.2.100",
|
||||
"@cerc-io/util": "^0.2.100",
|
||||
"@cerc-io/peer": "^0.2.101",
|
||||
"@cerc-io/rpc-eth-client": "^0.2.101",
|
||||
"@cerc-io/util": "^0.2.101",
|
||||
"@ethersproject/providers": "^5.4.4",
|
||||
"@graphql-tools/utils": "^9.1.1",
|
||||
"@ipld/dag-cbor": "^8.0.0",
|
||||
|
||||
@@ -96,7 +96,7 @@ export class BaseCmd {
|
||||
this._jobQueue = new JobQueue({ dbConnectionString, maxCompletionLag: maxCompletionLagInSecs });
|
||||
await this._jobQueue.start();
|
||||
|
||||
const { ethClient, ethProvider } = await initClients(this._config);
|
||||
const { ethClient, ethProvider } = await initClients(this._config.upstream);
|
||||
this._ethProvider = ethProvider;
|
||||
this._clients = { ethClient, ...clients };
|
||||
}
|
||||
|
||||
@@ -7,8 +7,6 @@ import { hideBin } from 'yargs/helpers';
|
||||
import 'reflect-metadata';
|
||||
import assert from 'assert';
|
||||
import { ConnectionOptions } from 'typeorm';
|
||||
import { errors } from 'ethers';
|
||||
import debug from 'debug';
|
||||
|
||||
import { JsonRpcProvider } from '@ethersproject/providers';
|
||||
import {
|
||||
@@ -22,15 +20,10 @@ import {
|
||||
GraphWatcherInterface,
|
||||
startMetricsServer,
|
||||
Config,
|
||||
UpstreamConfig,
|
||||
NEW_BLOCK_MAX_RETRIES_ERROR,
|
||||
setActiveUpstreamEndpointMetric
|
||||
UpstreamConfig
|
||||
} from '@cerc-io/util';
|
||||
|
||||
import { BaseCmd } from './base';
|
||||
import { initClients } from './utils/index';
|
||||
|
||||
const log = debug('vulcanize:job-runner');
|
||||
|
||||
interface Arguments {
|
||||
configFile: string;
|
||||
@@ -40,10 +33,6 @@ export class JobRunnerCmd {
|
||||
_argv?: Arguments;
|
||||
_baseCmd: BaseCmd;
|
||||
|
||||
_currentEndpointIndex = {
|
||||
rpcProviderEndpoint: 0
|
||||
};
|
||||
|
||||
constructor () {
|
||||
this._baseCmd = new BaseCmd();
|
||||
}
|
||||
@@ -124,26 +113,8 @@ export class JobRunnerCmd {
|
||||
const jobRunner = new JobRunner(
|
||||
config.jobQueue,
|
||||
indexer,
|
||||
jobQueue,
|
||||
async (error: any) => {
|
||||
// Check if it is a server error or timeout from ethers.js
|
||||
// https://docs.ethers.org/v5/api/utils/logger/#errors--server-error
|
||||
// https://docs.ethers.org/v5/api/utils/logger/#errors--timeout
|
||||
if (error.code === errors.SERVER_ERROR || error.code === errors.TIMEOUT || error.message === NEW_BLOCK_MAX_RETRIES_ERROR) {
|
||||
const oldRpcEndpoint = config.upstream.ethServer.rpcProviderEndpoints[this._currentEndpointIndex.rpcProviderEndpoint];
|
||||
++this._currentEndpointIndex.rpcProviderEndpoint;
|
||||
|
||||
if (this._currentEndpointIndex.rpcProviderEndpoint === config.upstream.ethServer.rpcProviderEndpoints.length) {
|
||||
this._currentEndpointIndex.rpcProviderEndpoint = 0;
|
||||
}
|
||||
|
||||
const { ethClient, ethProvider } = await initClients(config, this._currentEndpointIndex);
|
||||
indexer.switchClients({ ethClient, ethProvider });
|
||||
setActiveUpstreamEndpointMetric(config, this._currentEndpointIndex.rpcProviderEndpoint);
|
||||
|
||||
log(`RPC endpoint ${oldRpcEndpoint} is not working; failing over to new RPC endpoint ${ethProvider.connection.url}`);
|
||||
}
|
||||
});
|
||||
jobQueue
|
||||
);
|
||||
|
||||
// Delete all active and pending (before completed) jobs to start job-runner without old queued jobs
|
||||
await jobRunner.jobQueue.deleteAllJobs('completed');
|
||||
@@ -154,7 +125,7 @@ export class JobRunnerCmd {
|
||||
await startJobRunner(jobRunner);
|
||||
jobRunner.handleShutdown();
|
||||
|
||||
await startMetricsServer(config, jobQueue, indexer, this._currentEndpointIndex);
|
||||
await startMetricsServer(config, jobQueue, indexer);
|
||||
}
|
||||
|
||||
_getArgv (): any {
|
||||
|
||||
@@ -9,7 +9,7 @@ import { providers } from 'ethers';
|
||||
|
||||
// @ts-expect-error https://github.com/microsoft/TypeScript/issues/49721#issuecomment-1319854183
|
||||
import { PeerIdObj } from '@cerc-io/peer';
|
||||
import { Config, EthClient, getCustomProvider } from '@cerc-io/util';
|
||||
import { EthClient, UpstreamConfig, getCustomProvider } from '@cerc-io/util';
|
||||
import { getCache } from '@cerc-io/cache';
|
||||
import { EthClient as GqlEthClient } from '@cerc-io/ipld-eth-client';
|
||||
import { EthClient as RpcEthClient } from '@cerc-io/rpc-eth-client';
|
||||
@@ -22,16 +22,10 @@ export function readPeerId (filePath: string): PeerIdObj {
|
||||
return JSON.parse(peerIdJson);
|
||||
}
|
||||
|
||||
export const initClients = async (config: Config, endpointIndexes = { rpcProviderEndpoint: 0 }): Promise<{
|
||||
export const initClients = async (upstreamConfig: UpstreamConfig, endpointIndexes = { rpcProviderEndpoint: 0 }): Promise<{
|
||||
ethClient: EthClient,
|
||||
ethProvider: providers.JsonRpcProvider
|
||||
}> => {
|
||||
const { database: dbConfig, upstream: upstreamConfig, server: serverConfig } = config;
|
||||
|
||||
assert(serverConfig, 'Missing server config');
|
||||
assert(dbConfig, 'Missing database config');
|
||||
assert(upstreamConfig, 'Missing upstream config');
|
||||
|
||||
const { ethServer: { gqlApiEndpoint, rpcProviderEndpoints, rpcClient = false }, cache: cacheConfig } = upstreamConfig;
|
||||
|
||||
assert(rpcProviderEndpoints, 'Missing upstream ethServer.rpcProviderEndpoints');
|
||||
|
||||
Reference in New Issue
Block a user