Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/manual_release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -51,5 +51,5 @@ jobs:
file: Dockerfile
platforms: linux/amd64,linux/arm64
push: true
tags: |
tags: |
ghcr.io/${{ github.repository }}/${{ env.IMAGE_NAME }}:${{ github.event.inputs.tag }}
1 change: 0 additions & 1 deletion .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -94,4 +94,3 @@ jobs:
tags: |
ghcr.io/${{ github.repository }}/ewx-worker-node-server:latest
ghcr.io/${{ github.repository }}/ewx-worker-node-server:${{ needs.release-module.outputs.new_version }}

75 changes: 38 additions & 37 deletions docs/env-vars.md

Large diffs are not rendered by default.

5 changes: 5 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,11 @@ export const ENV_SCHEMA = z.object({
.number()
.default(20000)
.describe('How often should refresh EWX Solutions information from chain (in miliseconds).'),
INDEXER_SYNC_BUFFER_BLOCKS: z.coerce
.number()
.nonnegative()
.default(50)
.describe('The maximum block delay allowed for the indexer.'),
LOCAL_SOLUTIONS_PATH: z
.string()
.optional()
Expand Down
27 changes: 22 additions & 5 deletions src/solution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import {
type SolutionGroupId,
} from './polkadot/polka';
import { type SolutionGroup } from './polkadot/polka-types';
import { isIndexerSynced, queryStakeFromIndexer } from './util/indexer-queries';
import { sleep } from './util/sleep';
import { createLogger } from './util/logger';
import { createReadPalletApi } from './util/pallet-api';
Expand Down Expand Up @@ -270,11 +271,27 @@ const hasValidGroupConfiguration = async (
return false;
}

const hasStake: QueryStakeResult = await queryStake(
api,
operatorAddress,
solutionGroup.namespace,
);
let hasStake: QueryStakeResult | null = null;

try {
const isSynced = await isIndexerSynced(currentBlockNumber);
if (isSynced) {
hasStake = await queryStakeFromIndexer(operatorAddress, solutionGroup.namespace);
}
} catch (error) {
logger.warn(
{ solutionGroupId: solutionGroup.namespace, error },
'failed to retrieve stake from indexer',
);
}

if (hasStake == null) {
logger.info(
{ solutionGroupId: solutionGroup.namespace },
'falling back to direct RPC query for stake',
);
hasStake = await queryStake(api, operatorAddress, solutionGroup.namespace);
}

if (hasStake.currentStake < BigInt(solutionGroup.operatorsConfig.stakingAmounts.min)) {
logger.info(
Expand Down
156 changes: 156 additions & 0 deletions src/util/indexer-queries.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,156 @@
import axios from 'axios';
import { z } from 'zod';
import { type QueryStakeResult } from '../polkadot/polka';
import { createLogger } from './logger';
import { getBaseUrls } from './base-urls';
import { MAIN_CONFIG } from '../config';

const logger = createLogger('Indexer');

export const OperatorStakeSchema = z.object({
data: z.object({
operatorSubscribedStakes: z.array(
z.object({
currentStake: z.string().transform((val) => BigInt(val)),
nextStake: z.string().transform((val) => BigInt(val)),
rewardPeriodIndex: z.number().transform((val) => BigInt(val)),
}),
),
}),
});

export type OperatorStakeResponse = z.infer<typeof OperatorStakeSchema>;

export const SquidStatusSchema = z.object({
data: z.object({
squidStatus: z.object({
finalizedHash: z.string(),
finalizedHeight: z.number().transform((val) => BigInt(val)),
hash: z.string(),
height: z.number().transform((val) => BigInt(val)),
}),
}),
});

export type SquidStatusResponse = z.infer<typeof SquidStatusSchema>;

const GET_SQUID_STATUS_QUERY = `
query GetSquidStatus {
squidStatus {
finalizedHash
finalizedHeight
hash
height
}
}
`;

const GET_OPERATOR_STAKE_QUERY = `
query GetOperatorStake($operatorAddress: String!, $solutionGroupId: String!) {
operatorSubscribedStakes(
where: {
operator: { id_eq: $operatorAddress },
solutionGroup: { namespace_eq: $solutionGroupId }
},
orderBy: rewardPeriodIndex_DESC,
limit: 1
) {
currentStake
nextStake
rewardPeriodIndex
}
}
`;

export const queryStakeFromIndexer = async (
operatorAddress: string,
solutionGroupId: string,
): Promise<QueryStakeResult | null> => {
try {
const baseUrls = await getBaseUrls();

logger.info({ baseUrls, operatorAddress, solutionGroupId }, 'querying stake from indexer');

const { data } = await axios.post<OperatorStakeResponse>(
`${baseUrls.baseIndexerUrl}/core/graphql`,
{
query: GET_OPERATOR_STAKE_QUERY,
variables: {
operatorAddress,
solutionGroupId,
},
},
{
headers: {
'Content-Type': 'application/json',
},
timeout: 10000,
},
);

const parsed = OperatorStakeSchema.safeParse(data);

if (!parsed.success) {
logger.warn(
{ error: parsed.error.flatten(), operatorAddress, solutionGroupId },
'indexer response failed validation',
);
return null;
}

const stakes = parsed.data.data.operatorSubscribedStakes;

if (stakes.length > 0) {
const stakeNode = stakes[0];
return {
currentStake: stakeNode.currentStake,
nextStake: stakeNode.nextStake,
period: stakeNode.rewardPeriodIndex,
};
}
return null;
} catch (error) {
logger.warn({ error, operatorAddress, solutionGroupId }, 'failed to query stake from indexer');
return null;
}
};

export const isIndexerSynced = async (currentBlock: number): Promise<boolean> => {
try {
const baseUrls = await getBaseUrls();

const { data } = await axios.post<SquidStatusResponse>(
`${baseUrls.baseIndexerUrl}/core/graphql`,
{
query: GET_SQUID_STATUS_QUERY,
},
{
headers: {
'Content-Type': 'application/json',
},
timeout: 10000,
},
);

const parsed = SquidStatusSchema.safeParse(data);

if (!parsed.success) {
logger.warn({ error: parsed.error.flatten() }, 'squid status response failed validation');
return false;
}

const finalizedHeight = parsed.data.data.squidStatus.finalizedHeight;
const buffer = BigInt(MAIN_CONFIG.INDEXER_SYNC_BUFFER_BLOCKS);
const finalizedHeightWithBuffer = finalizedHeight + buffer;

logger.info(
{ finalizedHeight, currentBlock, buffer, finalizedHeightWithBuffer },
'indexer finalized height and current block',
);

return finalizedHeightWithBuffer >= BigInt(currentBlock);
} catch (error) {
logger.warn({ error }, 'failed to query squid status from indexer');
return false;
}
};
Loading