From e26cb32c3a9f0d7a39a7898ebba7c0c0d13860d5 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Thu, 3 Sep 2026 15:06:58 +0300 Subject: [PATCH 1/8] cs 2e2 tests fixes --- .../chain-simulator/utils/testSequencer.js | 2 +- .../websocket.subscriptions.cs-e2e.ts | 53 +++++++++++++++---- 2 files changed, 45 insertions(+), 10 deletions(-) diff --git a/src/test/chain-simulator/utils/testSequencer.js b/src/test/chain-simulator/utils/testSequencer.js index f1a70a4a8..5fcdfa1b9 100644 --- a/src/test/chain-simulator/utils/testSequencer.js +++ b/src/test/chain-simulator/utils/testSequencer.js @@ -13,7 +13,7 @@ class CustomSequencer extends Sequencer { 'delegation-legacy.cs-e2e.ts', 'accounts.cs-e2e.ts', 'stake.cs-e2e.ts', - 'round.cs-e2e.ts', + 'rounds.cs-e2e.ts', 'results.cs-e2e.ts', 'miniblocks.cs-e2e.ts', ]; diff --git a/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts b/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts index 7721f4508..2d155cf18 100644 --- a/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts +++ b/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts @@ -80,6 +80,33 @@ const aliceEsdts: string[] = []; describe('Websocket subscriptions e2e tests', () => { const clients: Socket[] = []; + const connectionErrors: string[] = []; + + // auto-reconnect is disabled on purpose: the subscriptions below are emitted from the 'connect' + // handler, so every reconnect would re-subscribe and the shared response arrays would collect the + // same message twice. connections are instead retried explicitly, before any operation is sent + const socketOptions = { path: '/ws/subscription', reconnection: false }; + + const waitForConnections = async (timeoutMs: number) => { + const deadline = Date.now() + timeoutMs; + + while (Date.now() < deadline) { + if (clients.every(client => client.connected)) { + return; + } + + for (const client of clients) { + if (!client.connected) { + client.connect(); + } + } + + await new Promise(resolve => setTimeout(resolve, 500)); + } + + const pending = clients.filter(client => !client.connected).length; + throw new Error(`${pending} of ${clients.length} websocket clients did not connect within ${timeoutMs}ms. Errors: ${connectionErrors.join('; ') || 'none reported'}`); + }; // --- Connect Helper --- const connectAndSubscribe = ( @@ -97,13 +124,13 @@ describe('Websocket subscriptions e2e tests', () => { eventResponses.set(filterKey, receivedEvents); transferResponses.set(filterKey, receivedTransfers); - const client: Socket = io(WS_SERVER_URL, { - path: '/ws/subscription', - }); + const client: Socket = io(WS_SERVER_URL, socketOptions); clients.push(client); + // never throw from a socket callback: it escapes as an uncaughtException that jest attributes + // to whichever test happens to be running, in any file. waitForConnections reports it instead client.on("connect_error", (err) => { - throw new Error(`${clientId} connection failed: ${err.message}`); + connectionErrors.push(`${clientId}: ${err.message}`); }); client.on("customTransactionUpdate", (data: { transactions: any[] }) => { @@ -137,12 +164,10 @@ describe('Websocket subscriptions e2e tests', () => { }; const connectAndSubscribeGeneral = (clientId: string, subConfig: typeof client4SubscriptionConfig) => { - const client: Socket = io(WS_SERVER_URL, { - path: '/ws/subscription', - }); + const client: Socket = io(WS_SERVER_URL, socketOptions); clients.push(client); - client.on("connect_error", (err) => { throw new Error(`${clientId} connection failed: ${err.message}`); }); + client.on("connect_error", (err) => { connectionErrors.push(`${clientId}: ${err.message}`); }); client.on("poolUpdate", (data: any) => generalResponses.pool.push(data)); client.on("eventsUpdate", (data: any) => generalResponses.events.push(data)); @@ -185,6 +210,8 @@ describe('Websocket subscriptions e2e tests', () => { connectAndSubscribeGeneral("client4", client4SubscriptionConfig); + await waitForConnections(30000); + await new Promise(resolve => setTimeout(resolve, 10000)); log("\n--- Starting Operations ---"); @@ -216,7 +243,15 @@ describe('Websocket subscriptions e2e tests', () => { }); afterAll(() => { - clients.forEach(client => client.connected && client.disconnect()); + // unconditionally: a client that is disconnected or mid-handshake would otherwise be skipped + // and keep its handlers and timers alive for the rest of the run, which is shared (--runInBand) + for (const client of clients) { + client.removeAllListeners(); + client.disconnect(); + client.close(); + } + + clients.length = 0; }); it('should receive TXs sent by Alice for Client 1', () => { From 82be978dc692c41697c8457347f936af0bf52b0a Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Thu, 3 Sep 2026 15:37:54 +0300 Subject: [PATCH 2/8] fix ws subscriptions --- .../chain-simulator/websocket.subscriptions.cs-e2e.ts | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts b/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts index 2d155cf18..cf6df420b 100644 --- a/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts +++ b/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts @@ -234,7 +234,15 @@ describe('Websocket subscriptions e2e tests', () => { await axios.post(`${config.chainSimulatorUrl}/simulator/generate-blocks/10`); log("Waiting for WS messages..."); - await new Promise(resolve => setTimeout(resolve, 35000)); + + // the broadcaster commits its cursor only after it sees the round that follows the one it just + // sent. a simulator that has stopped producing blocks never provides that round, so it times + // out and replays the whole window on its next tick. one block per wait step keeps it moving, + // which is why this is a loop of short sleeps rather than a single long one + for (let i = 0; i < 30; i++) { + await axios.post(`${config.chainSimulatorUrl}/simulator/generate-blocks/1`); + await new Promise(resolve => setTimeout(resolve, 1000)); + } } catch (e: any) { console.error("Error in beforeAll:", e.message); From ca893a2221180f4998142c100fd1808687dd198e Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Thu, 3 Sep 2026 15:54:21 +0300 Subject: [PATCH 3/8] wait for api instead of hardcoded timeout --- .../utils/prepare-test-data.ts | 28 ++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/src/test/chain-simulator/utils/prepare-test-data.ts b/src/test/chain-simulator/utils/prepare-test-data.ts index e02a454b9..c78889cd6 100644 --- a/src/test/chain-simulator/utils/prepare-test-data.ts +++ b/src/test/chain-simulator/utils/prepare-test-data.ts @@ -1,7 +1,29 @@ +import axios from 'axios'; import { config } from '../config/env.config'; import { fundAddress, issueMultipleEsdts, issueMultipleMetaESDTCollections, issueMultipleNftsCollections } from './chain.simulator.operations'; import { ChainSimulatorUtils } from './test.utils'; +async function waitForApi(description: string, url: string, expected: number, timeoutMs: number = 180000) { + const deadline = Date.now() + timeoutMs; + let last = 0; + + while (Date.now() < deadline) { + try { + last = (await axios.get(url)).data; + if (last >= expected) { + console.log(`✓ ${description}: ${last}`); + return; + } + } catch (error: any) { + last = -1; + } + + await new Promise((resolve) => setTimeout(resolve, 5000)); + } + + throw new Error(`${description}: reached ${last}, expected at least ${expected}, after ${timeoutMs}ms (${url})`); +} + async function prepareTestData() { try { console.log('Starting test data preparation...'); @@ -25,7 +47,11 @@ async function prepareTestData() { await ChainSimulatorUtils.deployPingPongSc(config.aliceAddress); console.log('✓ Deployed PingPong smart contract'); - await new Promise((resolve) => setTimeout(resolve, 30000)); + await waitForApi('Tokens listed by the API', `${config.apiServiceUrl}/tokens/count`, 5); + await waitForApi('Tokens listed on the issuer account', `${config.apiServiceUrl}/accounts/${config.aliceAddress}/tokens/count`, 5); + // 2 NFT + 2 SFT collections, five items each; the meta-esdt ones come on top of that + await waitForApi('Collections listed by the API', `${config.apiServiceUrl}/collections/count`, 4); + await waitForApi('NFTs listed by the API', `${config.apiServiceUrl}/nfts/count`, 20); console.log('Test data preparation completed successfully!'); } catch (error) { From e0dc00e8d2dcee1d6a38d3b470bcc38a326924d5 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Thu, 3 Sep 2026 16:14:44 +0300 Subject: [PATCH 4/8] wati for api to index --- .../utils/prepare-test-data.ts | 29 ++++++++++++------- 1 file changed, 19 insertions(+), 10 deletions(-) diff --git a/src/test/chain-simulator/utils/prepare-test-data.ts b/src/test/chain-simulator/utils/prepare-test-data.ts index c78889cd6..9e2ce0e25 100644 --- a/src/test/chain-simulator/utils/prepare-test-data.ts +++ b/src/test/chain-simulator/utils/prepare-test-data.ts @@ -3,25 +3,28 @@ import { config } from '../config/env.config'; import { fundAddress, issueMultipleEsdts, issueMultipleMetaESDTCollections, issueMultipleNftsCollections } from './chain.simulator.operations'; import { ChainSimulatorUtils } from './test.utils'; -async function waitForApi(description: string, url: string, expected: number, timeoutMs: number = 180000) { +// the api serves most of this from caches the cache warmer fills on its own crons, so what the chain +// reports says nothing about what the tests will see. wait on the api's own view instead, per kind of +// data, rather than on a single interval that has to be long enough for the slowest of them +async function waitForApi(description: string, url: string, isReady: (data: any) => boolean, timeoutMs: number = 180000) { const deadline = Date.now() + timeoutMs; - let last = 0; + let last: any = 'no response yet'; while (Date.now() < deadline) { try { last = (await axios.get(url)).data; - if (last >= expected) { - console.log(`✓ ${description}: ${last}`); + if (isReady(last)) { + console.log(`✓ ${description}`); return; } } catch (error: any) { - last = -1; + last = error.message; } await new Promise((resolve) => setTimeout(resolve, 5000)); } - throw new Error(`${description}: reached ${last}, expected at least ${expected}, after ${timeoutMs}ms (${url})`); + throw new Error(`${description}: still not ready after ${timeoutMs}ms. ${url} last returned ${JSON.stringify(last).slice(0, 300)}`); } async function prepareTestData() { @@ -47,11 +50,17 @@ async function prepareTestData() { await ChainSimulatorUtils.deployPingPongSc(config.aliceAddress); console.log('✓ Deployed PingPong smart contract'); - await waitForApi('Tokens listed by the API', `${config.apiServiceUrl}/tokens/count`, 5); - await waitForApi('Tokens listed on the issuer account', `${config.apiServiceUrl}/accounts/${config.aliceAddress}/tokens/count`, 5); + await waitForApi('Tokens listed by the API', `${config.apiServiceUrl}/tokens/count`, count => count >= 5); + await waitForApi('Tokens listed on the issuer account', `${config.apiServiceUrl}/accounts/${config.aliceAddress}/tokens/count`, count => count >= 5); // 2 NFT + 2 SFT collections, five items each; the meta-esdt ones come on top of that - await waitForApi('Collections listed by the API', `${config.apiServiceUrl}/collections/count`, 4); - await waitForApi('NFTs listed by the API', `${config.apiServiceUrl}/nfts/count`, 20); + await waitForApi('Collections listed by the API', `${config.apiServiceUrl}/collections/count`, count => count >= 4); + await waitForApi('NFTs listed by the API', `${config.apiServiceUrl}/nfts/count`, count => count >= 20); + + // node and validator statistics are filled in by warmers on a one minute cron, and shards are + // derived from them. until those have run at least once, /shards is empty and nodes come back + // with a null rating and no status + await waitForApi('Shards reported by the API', `${config.apiServiceUrl}/shards`, shards => shards.length >= 4); + await waitForApi('Node ratings filled in by the API', `${config.apiServiceUrl}/nodes?size=1`, nodes => nodes.length > 0 && nodes[0].status !== undefined && nodes[0].tempRating !== null); console.log('Test data preparation completed successfully!'); } catch (error) { From 02f4bf59b92d6defb8d7a8e5f5214fc9a7194aa6 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 25 Sep 2026 13:03:13 +0300 Subject: [PATCH 5/8] fix flaky chain simulator e2e tests - websocket: retry a connection only after connect_error. calling connect() while the handshake is in flight sent a second CONNECT packet, the server opened a second socket and every message arrived twice - websocket: subscribe once per client, wait for the subscription acks and count distinct messages, since the broadcaster replays rounds on timeout - run the test files missing from the sequencer list in a fixed order - rethrow errors from the transaction helpers instead of returning placeholders that fail later, unrelated assertions - poll the api for bob's nft instead of a fixed sleep Co-Authored-By: Claude Opus 5.5 --- src/test/chain-simulator/accounts.cs-e2e.ts | 12 +- .../utils/chain.simulator.operations.ts | 15 +- .../utils/prepare-test-data.ts | 37 +--- src/test/chain-simulator/utils/test.utils.ts | 25 +++ .../chain-simulator/utils/testSequencer.js | 4 +- .../websocket.subscriptions.cs-e2e.ts | 195 +++++++++++------- 6 files changed, 176 insertions(+), 112 deletions(-) diff --git a/src/test/chain-simulator/accounts.cs-e2e.ts b/src/test/chain-simulator/accounts.cs-e2e.ts index f1182d86e..8e41c26d9 100644 --- a/src/test/chain-simulator/accounts.cs-e2e.ts +++ b/src/test/chain-simulator/accounts.cs-e2e.ts @@ -3,6 +3,7 @@ import { config } from "./config/env.config"; import { NftType } from "src/endpoints/nfts/entities/nft.type"; import { NftSubType } from "src/endpoints/nfts/entities/nft.sub.type"; import { transferNftFromTo } from "./utils/chain.simulator.operations"; +import { ChainSimulatorUtils } from "./utils/test.utils"; describe('Accounts e2e tests with chain simulator', () => { describe('GET /accounts with query parameters', () => { @@ -1993,10 +1994,17 @@ describe('Accounts e2e tests with chain simulator', () => { const sendNftTx = await transferNftFromTo(config.chainSimulatorUrl, config.aliceAddress, config.bobAddress, nft.collection, nft.nonce); const transaction = await axios.get(`${config.apiServiceUrl}/transactions/${sendNftTx}`); - await new Promise(resolve => setTimeout(resolve, 2000)); - expect(transaction.status).toBe(200); + // the transfer is processed on chain, but bob's nfts are read from the index, which catches up later + await ChainSimulatorUtils.waitForApi( + 'NFT received by Bob', + `${config.apiServiceUrl}/accounts/${config.bobAddress}/nfts?withReceivedAt=true`, + nfts => nfts.length >= 1 && nfts[0].receivedAt !== undefined, + 60000, + 1000, + ); + const checkBobNft = await axios.get(`${config.apiServiceUrl}/accounts/${config.bobAddress}/nfts?withReceivedAt=true`); const bobNft = checkBobNft.data; diff --git a/src/test/chain-simulator/utils/chain.simulator.operations.ts b/src/test/chain-simulator/utils/chain.simulator.operations.ts index 0e9e7bb64..274962bf9 100644 --- a/src/test/chain-simulator/utils/chain.simulator.operations.ts +++ b/src/test/chain-simulator/utils/chain.simulator.operations.ts @@ -32,7 +32,7 @@ export async function getNonce( return currentNonceResponse.data.data.nonce; } catch (e) { console.error(e); - return 0; + throw e; } } @@ -57,13 +57,17 @@ export async function deploySc(args: DeployScArgs): Promise { const scDeployLog = txResponse?.data?.data?.transaction?.logs?.events?.find( (event: { identifier: string }) => event.identifier === 'SCDeploy', ); + if (!scDeployLog) { + throw new Error(`SC deploy ${txHash} produced no SCDeploy event`); + } + console.log( - `Deployed SC. tx hash: ${txHash}. address: ${scDeployLog?.address}`, + `Deployed SC. tx hash: ${txHash}. address: ${scDeployLog.address}`, ); - return scDeployLog?.address; + return scDeployLog.address; } catch (e) { console.error(e); - return 'n/a'; + throw e; } } @@ -160,8 +164,9 @@ export async function sendTransaction( ); return txHash; } catch (e) { + // rethrown: a placeholder result only moves the failure to some later, unrelated assertion console.error(e); - return 'n/a'; + throw e; } } diff --git a/src/test/chain-simulator/utils/prepare-test-data.ts b/src/test/chain-simulator/utils/prepare-test-data.ts index 9e2ce0e25..5a6badff3 100644 --- a/src/test/chain-simulator/utils/prepare-test-data.ts +++ b/src/test/chain-simulator/utils/prepare-test-data.ts @@ -1,32 +1,7 @@ -import axios from 'axios'; import { config } from '../config/env.config'; import { fundAddress, issueMultipleEsdts, issueMultipleMetaESDTCollections, issueMultipleNftsCollections } from './chain.simulator.operations'; import { ChainSimulatorUtils } from './test.utils'; -// the api serves most of this from caches the cache warmer fills on its own crons, so what the chain -// reports says nothing about what the tests will see. wait on the api's own view instead, per kind of -// data, rather than on a single interval that has to be long enough for the slowest of them -async function waitForApi(description: string, url: string, isReady: (data: any) => boolean, timeoutMs: number = 180000) { - const deadline = Date.now() + timeoutMs; - let last: any = 'no response yet'; - - while (Date.now() < deadline) { - try { - last = (await axios.get(url)).data; - if (isReady(last)) { - console.log(`✓ ${description}`); - return; - } - } catch (error: any) { - last = error.message; - } - - await new Promise((resolve) => setTimeout(resolve, 5000)); - } - - throw new Error(`${description}: still not ready after ${timeoutMs}ms. ${url} last returned ${JSON.stringify(last).slice(0, 300)}`); -} - async function prepareTestData() { try { console.log('Starting test data preparation...'); @@ -50,17 +25,17 @@ async function prepareTestData() { await ChainSimulatorUtils.deployPingPongSc(config.aliceAddress); console.log('✓ Deployed PingPong smart contract'); - await waitForApi('Tokens listed by the API', `${config.apiServiceUrl}/tokens/count`, count => count >= 5); - await waitForApi('Tokens listed on the issuer account', `${config.apiServiceUrl}/accounts/${config.aliceAddress}/tokens/count`, count => count >= 5); + await ChainSimulatorUtils.waitForApi('Tokens listed by the API', `${config.apiServiceUrl}/tokens/count`, count => count >= 5); + await ChainSimulatorUtils.waitForApi('Tokens listed on the issuer account', `${config.apiServiceUrl}/accounts/${config.aliceAddress}/tokens/count`, count => count >= 5); // 2 NFT + 2 SFT collections, five items each; the meta-esdt ones come on top of that - await waitForApi('Collections listed by the API', `${config.apiServiceUrl}/collections/count`, count => count >= 4); - await waitForApi('NFTs listed by the API', `${config.apiServiceUrl}/nfts/count`, count => count >= 20); + await ChainSimulatorUtils.waitForApi('Collections listed by the API', `${config.apiServiceUrl}/collections/count`, count => count >= 4); + await ChainSimulatorUtils.waitForApi('NFTs listed by the API', `${config.apiServiceUrl}/nfts/count`, count => count >= 20); // node and validator statistics are filled in by warmers on a one minute cron, and shards are // derived from them. until those have run at least once, /shards is empty and nodes come back // with a null rating and no status - await waitForApi('Shards reported by the API', `${config.apiServiceUrl}/shards`, shards => shards.length >= 4); - await waitForApi('Node ratings filled in by the API', `${config.apiServiceUrl}/nodes?size=1`, nodes => nodes.length > 0 && nodes[0].status !== undefined && nodes[0].tempRating !== null); + await ChainSimulatorUtils.waitForApi('Shards reported by the API', `${config.apiServiceUrl}/shards`, shards => shards.length >= 4); + await ChainSimulatorUtils.waitForApi('Node ratings filled in by the API', `${config.apiServiceUrl}/nodes?size=1`, nodes => nodes.length > 0 && nodes[0].status !== undefined && nodes[0].tempRating !== null); console.log('Test data preparation completed successfully!'); } catch (error) { diff --git a/src/test/chain-simulator/utils/test.utils.ts b/src/test/chain-simulator/utils/test.utils.ts index 02cac1a91..08bbdac9b 100644 --- a/src/test/chain-simulator/utils/test.utils.ts +++ b/src/test/chain-simulator/utils/test.utils.ts @@ -51,6 +51,31 @@ export class ChainSimulatorUtils { } } + // the api serves most of its data from caches the cache warmer fills on its own crons, and from an + // index that is refreshed on its own interval, so what the chain reports says nothing about what the + // api will return. wait on the api's own view instead of on a fixed interval that has to be long + // enough for the slowest of them + static async waitForApi(description: string, url: string, isReady: (data: any) => boolean, timeoutMs: number = 180000, intervalMs: number = 5000) { + const deadline = Date.now() + timeoutMs; + let last: any = 'no response yet'; + + while (Date.now() < deadline) { + try { + last = (await axios.get(url)).data; + if (isReady(last)) { + console.log(`✓ ${description}`); + return; + } + } catch (error: any) { + last = error.message; + } + + await new Promise((resolve) => setTimeout(resolve, intervalMs)); + } + + throw new Error(`${description}: still not ready after ${timeoutMs}ms. ${url} last returned ${JSON.stringify(last).slice(0, 300)}`); + } + private static async checkSimulatorHealth(maxRetries: number = 50): Promise { let retries = 0; diff --git a/src/test/chain-simulator/utils/testSequencer.js b/src/test/chain-simulator/utils/testSequencer.js index 5fcdfa1b9..e1ba59f66 100644 --- a/src/test/chain-simulator/utils/testSequencer.js +++ b/src/test/chain-simulator/utils/testSequencer.js @@ -31,7 +31,9 @@ class CustomSequencer extends Sequencer { if (indexB !== -1) { return 1; } - return 0; + // the files not listed above share one chain, so the order they run in decides the state each + // of them sees. leaving it to the order jest discovered them in makes that differ between runs + return testA.path.localeCompare(testB.path); }); } } diff --git a/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts b/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts index cf6df420b..a41b408a2 100644 --- a/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts +++ b/src/test/chain-simulator/websocket.subscriptions.cs-e2e.ts @@ -23,9 +23,23 @@ const log = (...args: any[]) => { } }; -const txResponses: Map = new Map(); -const eventResponses: Map = new Map(); -const transferResponses: Map = new Map(); // New: Store transfers +// the broadcaster can push the same round twice: when it times out waiting for the following round it +// does not advance its cursor, and replays from there on the next tick. a replayed item is identical +// to the original, so the responses are kept by identity and only the distinct ones are counted +const txResponses: Map> = new Map(); +const eventResponses: Map> = new Map(); +const transferResponses: Map> = new Map(); // New: Store transfers + +const txKey = (tx: any) => tx.txHash; +const eventKey = (evt: any) => `${evt.txHash}:${evt.order}:${evt.identifier}`; + +const collect = (target: Map, items: any[], key: (item: any) => string) => { + for (const item of items) { + target.set(key(item), item); + } +}; + +const received = (responses: Map>, filterKey: string) => [...(responses.get(filterKey)?.values() ?? [])]; const generalResponses = { pool: [] as any[], @@ -81,10 +95,14 @@ const aliceEsdts: string[] = []; describe('Websocket subscriptions e2e tests', () => { const clients: Socket[] = []; const connectionErrors: string[] = []; + const failedClients: Set = new Set(); + const subscriptions: Promise[] = []; // auto-reconnect is disabled on purpose: the subscriptions below are emitted from the 'connect' - // handler, so every reconnect would re-subscribe and the shared response arrays would collect the - // same message twice. connections are instead retried explicitly, before any operation is sent + // handler, so every reconnect would re-subscribe. connections are instead retried explicitly, and + // only once connect_error says the previous attempt is over: calling connect() while a handshake is + // still in flight sends a second CONNECT packet, the server opens a second socket for it, and every + // message then arrives twice const socketOptions = { path: '/ws/subscription', reconnection: false }; const waitForConnections = async (timeoutMs: number) => { @@ -95,10 +113,9 @@ describe('Websocket subscriptions e2e tests', () => { return; } - for (const client of clients) { - if (!client.connected) { - client.connect(); - } + for (const client of failedClients) { + failedClients.delete(client); + client.connect(); } await new Promise(resolve => setTimeout(resolve, 500)); @@ -108,7 +125,47 @@ describe('Websocket subscriptions e2e tests', () => { throw new Error(`${pending} of ${clients.length} websocket clients did not connect within ${timeoutMs}ms. Errors: ${connectionErrors.join('; ') || 'none reported'}`); }; + // a subscription that the server rejects is answered with an 'error' event and never acknowledged, + // so the timeout is what reports it. the failure is resolved rather than rejected, since it may come + // before anything awaits it and would otherwise surface as an unhandled rejection + const subscribe = (client: Socket, clientId: string, event: string, ...args: any[]) => { + const subscription = client.timeout(30000).emitWithAck(event, ...args).then( + (ack: any) => { + log(` ACK ${event} ${clientId}:`, ack); + return undefined; + }, + () => `${clientId}: ${event} was not acknowledged`, + ); + + subscriptions.push(subscription); + }; + // --- Connect Helper --- + const connectClient = (clientId: string, onConnect: (client: Socket) => void) => { + const client: Socket = io(WS_SERVER_URL, socketOptions); + clients.push(client); + + // never throw from a socket callback: it escapes as an uncaughtException that jest attributes + // to whichever test happens to be running, in any file. waitForConnections reports it instead + client.on("connect_error", (err) => { + connectionErrors.push(`${clientId}: ${err.message}`); + failedClients.add(client); + }); + + let subscribed = false; + client.on("connect", () => { + log(`\n ${clientId} connected.`); + + if (!subscribed) { + subscribed = true; + onConnect(client); + } + }); + + return client; + }; + + // --- Subscribe Helpers --- const connectAndSubscribe = ( filterKey: string, txFilter: any, @@ -116,73 +173,56 @@ describe('Websocket subscriptions e2e tests', () => { transferFilter: any, clientId: string ) => { - const receivedTxs: any[] = []; - const receivedEvents: any[] = []; - const receivedTransfers: any[] = []; + const receivedTxs: Map = new Map(); + const receivedEvents: Map = new Map(); + const receivedTransfers: Map = new Map(); txResponses.set(filterKey, receivedTxs); eventResponses.set(filterKey, receivedEvents); transferResponses.set(filterKey, receivedTransfers); - const client: Socket = io(WS_SERVER_URL, socketOptions); - clients.push(client); - - // never throw from a socket callback: it escapes as an uncaughtException that jest attributes - // to whichever test happens to be running, in any file. waitForConnections reports it instead - client.on("connect_error", (err) => { - connectionErrors.push(`${clientId}: ${err.message}`); + const client = connectClient(clientId, (client) => { + if (txFilter) { + subscribe(client, clientId, "subscribeCustomTransactions", txFilter); + } + if (eventFilter) { + subscribe(client, clientId, "subscribeCustomEvents", eventFilter); + } + if (transferFilter) { + subscribe(client, clientId, "subscribeCustomTransfers", transferFilter); + } }); client.on("customTransactionUpdate", (data: { transactions: any[] }) => { log(`\n💸 ${clientId} received ${data.transactions.length} txs`); - receivedTxs.push(...data.transactions); + collect(receivedTxs, data.transactions, txKey); }); client.on("customEventUpdate", (data: { events: any[] }) => { log(`\n🔔 ${clientId} received ${data.events.length} events`); - receivedEvents.push(...data.events); + collect(receivedEvents, data.events, eventKey); }); client.on("customTransferUpdate", (data: { transfers: any[] }) => { log(`\n💎 ${clientId} received ${data.transfers.length} transfers`); - receivedTransfers.push(...data.transfers); - }); - - client.on("connect", () => { - log(`\n ${clientId} connected.`); - - if (txFilter) { - client.emit("subscribeCustomTransactions", txFilter, (ack: any) => log(` ACK TXs ${clientId}:`, ack)); - } - if (eventFilter) { - client.emit("subscribeCustomEvents", eventFilter, (ack: any) => log(` ACK Events ${clientId}:`, ack)); - } - if (transferFilter) { - client.emit("subscribeCustomTransfers", transferFilter, (ack: any) => log(` ACK Transfers ${clientId}:`, ack)); - } + collect(receivedTransfers, data.transfers, txKey); }); }; const connectAndSubscribeGeneral = (clientId: string, subConfig: typeof client4SubscriptionConfig) => { - const client: Socket = io(WS_SERVER_URL, socketOptions); - clients.push(client); - - client.on("connect_error", (err) => { connectionErrors.push(`${clientId}: ${err.message}`); }); + const client = connectClient(clientId, (client) => { + subscribe(client, clientId, "subscribePool", subConfig.pool); + subscribe(client, clientId, "subscribeEvents", subConfig.events); + subscribe(client, clientId, "subscribeTransactions", subConfig.transactions); + subscribe(client, clientId, "subscribeBlocks", subConfig.blocks); + subscribe(client, clientId, "subscribeStats"); + }); client.on("poolUpdate", (data: any) => generalResponses.pool.push(data)); client.on("eventsUpdate", (data: any) => generalResponses.events.push(data)); client.on("transactionUpdate", (data: any) => generalResponses.transactions.push(data)); client.on("blocksUpdate", (data: any) => generalResponses.blocks.push(data)); client.on("statsUpdate", (data: any) => generalResponses.stats.push(data)); - - client.on("connect", () => { - log(`\n ${clientId} connected with specific configs.`); - client.emit("subscribePool", subConfig.pool, (ack: any) => log(`ACK Pool ${clientId}:`, ack)); - client.emit("subscribeEvents", subConfig.events, (ack: any) => log(`ACK Events ${clientId}:`, ack)); - client.emit("subscribeTransactions", subConfig.transactions, (ack: any) => log(`ACK Txs ${clientId}:`, ack)); - client.emit("subscribeBlocks", subConfig.blocks, (ack: any) => log(`ACK Blocks ${clientId}:`, ack)); - client.emit("subscribeStats", (ack: any) => log(`ACK Stats ${clientId}:`, ack)); - }); }; beforeAll(async () => { @@ -212,6 +252,15 @@ describe('Websocket subscriptions e2e tests', () => { await waitForConnections(30000); + // every client is connected, so every 'connect' handler has run and queued its subscriptions + const subscriptionErrors = (await Promise.all(subscriptions)).filter(error => error !== undefined); + if (subscriptionErrors.length > 0) { + throw new Error(`Subscriptions failed: ${subscriptionErrors.join('; ')}`); + } + + // the broadcaster starts from the latest round it sees on its first tick after the subscriptions + // exist, and never goes back to the ones before it. give it time to take that position before + // any operation produces blocks await new Promise(resolve => setTimeout(resolve, 10000)); log("\n--- Starting Operations ---"); @@ -263,65 +312,65 @@ describe('Websocket subscriptions e2e tests', () => { }); it('should receive TXs sent by Alice for Client 1', () => { - const txs = txResponses.get(filterKeys.CLIENT_1); - expect(txs?.length).toBe(4); + const txs = received(txResponses, filterKeys.CLIENT_1); + expect(txs.length).toBe(4); - txs?.forEach((tx) => { + txs.forEach((tx) => { expect(tx.sender).toEqual(config.aliceAddress); }); }); it('should receive Events with identifier "pong" for Client 1', () => { - const events = eventResponses.get(filterKeys.CLIENT_1); - expect(events?.length).toBe(1); + const events = received(eventResponses, filterKeys.CLIENT_1); + expect(events.length).toBe(1); - events?.forEach((evt) => { + events.forEach((evt) => { expect(evt.identifier).toEqual('pong'); }); }); it('should receive TXs sent by Bob for Client 2', () => { - const txs = txResponses.get(filterKeys.CLIENT_2); - expect(txs?.length).toBe(1); + const txs = received(txResponses, filterKeys.CLIENT_2); + expect(txs.length).toBe(1); - txs?.forEach((tx) => { + txs.forEach((tx) => { expect(tx.sender).toEqual(config.bobAddress); }); }); it('should receive Events generated by PingPong contract (address) for Client 2', () => { - const events = eventResponses.get(filterKeys.CLIENT_2); - expect(events?.length).toBe(6); + const events = received(eventResponses, filterKeys.CLIENT_2); + expect(events.length).toBe(6); - events?.forEach((evt) => { + events.forEach((evt) => { expect(evt.address).toEqual(pingPongScAddress); }); }); it('should receive specific Alice-to-Bob TXs for Client 3', () => { - const txs = txResponses.get(filterKeys.CLIENT_3); - expect(txs?.length).toBeGreaterThanOrEqual(1); - txs?.forEach((tx) => { + const txs = received(txResponses, filterKeys.CLIENT_3); + expect(txs.length).toBeGreaterThanOrEqual(1); + txs.forEach((tx) => { expect(tx.sender).toEqual(config.aliceAddress); expect(tx.receiver).toEqual(config.bobAddress); }); }); it('should receive ANY transfer involving Alice (Client 5 - Address Filter)', () => { - const transfers = transferResponses.get(filterKeys.CLIENT_5); - expect(transfers?.length).toBeGreaterThan(0); + const transfers = received(transferResponses, filterKeys.CLIENT_5); + expect(transfers.length).toBeGreaterThan(0); - transfers?.forEach(t => { + transfers.forEach(t => { const isAliceInvolved = t.sender === config.aliceAddress || t.receiver === config.aliceAddress; expect(isAliceInvolved).toBe(true); }); }); it('should receive ONLY EGLD transfers where ALICE is involved (Client 6 - Token EGLD Filter)', () => { - const transfers = transferResponses.get(filterKeys.CLIENT_6); - expect(transfers?.length).toBeGreaterThan(0); + const transfers = received(transferResponses, filterKeys.CLIENT_6); + expect(transfers.length).toBeGreaterThan(0); - transfers?.forEach(t => { + transfers.forEach(t => { const val1 = `1${'0'.repeat(18)}`; const val2 = `2${'0'.repeat(18)}`; expect([val1, val2]).toContain(t.value); @@ -336,10 +385,10 @@ describe('Websocket subscriptions e2e tests', () => { }); it('should receive ONLY specific ESDT transfers (Client 7 - Dynamic Token Filter)', () => { - const transfers = transferResponses.get(filterKeys.CLIENT_7); - expect(transfers?.length).toBeGreaterThan(0); + const transfers = received(transferResponses, filterKeys.CLIENT_7); + expect(transfers.length).toBeGreaterThan(0); - transfers?.forEach(t => { + transfers.forEach(t => { const esdtTransfers = t.action?.arguments?.transfers; const containsAliceEsdt = esdtTransfers.filter((et: any) => et.token === aliceEsdts[0]).length > 0; expect(containsAliceEsdt).toBe(true); From 1cf0b953f40ac22bdd64bdb69316fc34b13f198d Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 25 Sep 2026 14:34:56 +0300 Subject: [PATCH 6/8] enable fast warming on e2e cs tests --- config/config.e2e.mainnet.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/config/config.e2e.mainnet.yaml b/config/config.e2e.mainnet.yaml index cd9243ae8..0918f9ece 100644 --- a/config/config.e2e.mainnet.yaml +++ b/config/config.e2e.mainnet.yaml @@ -8,7 +8,7 @@ api: websocket: true cron: cacheWarmer: true - fastWarm: false + fastWarm: true queueWorker: true elasticUpdater: false flags: From 47716f89c872168469ce9760510ee96b8e3c8b55 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 25 Sep 2026 16:40:15 +0300 Subject: [PATCH 7/8] wait for issued tokens to be indexed and turn off the pool warmer in cs tests the transaction pool warmer, enabled only in the e2e config, resolved the action of every pending transaction and smart contract result each second, including the ESDTTransfer that delivers an issued token, before the token was indexed. the unresolved token properties were then cached for an hour, so later transfers of that token came back without an action, depending on whether a warmer tick caught that result in the pool. the warmer is now off, as on mainnet and devnet, and the pool is still read on request. the issuance helpers also wait until the token or collection can be found through the same elastic query the api runs, before returning it. --- config/config.e2e.mainnet.yaml | 6 ++++- src/test/chain-simulator/config/.env.example | 1 + src/test/chain-simulator/config/env.config.ts | 1 + .../utils/chain.simulator.operations.ts | 24 +++++++++++++++++++ 4 files changed, 31 insertions(+), 1 deletion(-) diff --git a/config/config.e2e.mainnet.yaml b/config/config.e2e.mainnet.yaml index 0918f9ece..cc9d7f531 100644 --- a/config/config.e2e.mainnet.yaml +++ b/config/config.e2e.mainnet.yaml @@ -39,7 +39,11 @@ features: transactionPool: enabled: true transactionPoolWarmer: - enabled: true + # off, as on mainnet and devnet: every second it resolved the action of each pending transaction and + # smart contract result, including the ESDTTransfer through which a token is issued, before the token + # is indexed. the unresolved token properties were then cached for an hour, and every later transfer + # of that token came back without an action. the pool is still read on request + enabled: false cronExpression: '*/1 * * * * *' ttlInSeconds: 10 updateCollectionExtraDetails: diff --git a/src/test/chain-simulator/config/.env.example b/src/test/chain-simulator/config/.env.example index dc915101f..01edd2dca 100644 --- a/src/test/chain-simulator/config/.env.example +++ b/src/test/chain-simulator/config/.env.example @@ -1,5 +1,6 @@ CHAIN_SIMULATOR_URL=http://localhost:8085 API_SERVICE_URL=http://localhost:3001 SUBSCRIPTIONS_SERIVCE_URL=http://localhost:6002 +ELASTIC_URL=http://localhost:9200 ALICE_ADDRESS=erd1qyu5wthldzr8wx5c9ucg8kjagg0jfs53s8nr3zpz3hypefsdd8ssycr6th BOB_ADDRESS=erd1spyavw0956vq68xj8y4tenjpq2wd5a9p2c6j8gsz7ztyrnpxrruqzu66jx diff --git a/src/test/chain-simulator/config/env.config.ts b/src/test/chain-simulator/config/env.config.ts index 70f0a1c35..fad47991b 100644 --- a/src/test/chain-simulator/config/env.config.ts +++ b/src/test/chain-simulator/config/env.config.ts @@ -9,6 +9,7 @@ export const config = { chainSimulatorUrl: process.env.CHAIN_SIMULATOR_URL || 'http://localhost:8085', apiServiceUrl: process.env.API_SERVICE_URL || 'http://localhost:3001', subscriptionsServiceUrl: process.env.SUBSCRIPTIONS_SERVICE_URL || 'http://localhost:6002', + elasticUrl: process.env.ELASTIC_URL || 'http://localhost:9200', aliceAddress: process.env.ALICE_ADDRESS || 'erd1qyu5wthldzr8wx5c9ucg8kjagg0jfs53s8nr3zpz3hypefsdd8ssycr6th', bobAddress: process.env.BOB_ADDRESS || 'erd1spyavw0956vq68xj8y4tenjpq2wd5a9p2c6j8gsz7ztyrnpxrruqzu66jx', }; diff --git a/src/test/chain-simulator/utils/chain.simulator.operations.ts b/src/test/chain-simulator/utils/chain.simulator.operations.ts index 274962bf9..0b3ee1350 100644 --- a/src/test/chain-simulator/utils/chain.simulator.operations.ts +++ b/src/test/chain-simulator/utils/chain.simulator.operations.ts @@ -1,5 +1,6 @@ import axios from 'axios'; import { AddressUtils } from "@multiversx/sdk-nestjs-common"; +import { config } from '../config/env.config'; axios.defaults.adapter = 'fetch'; axios.defaults.headers.common['Connection'] = 'close'; @@ -71,6 +72,25 @@ export async function deploySc(args: DeployScArgs): Promise { } } +// the api resolves a token through its document in the tokens index, and caches the answer, including +// the answer that the token does not exist. anything that makes the api look the token up before it is +// indexed leaves the token unresolved for as long as that answer is cached, so an issued token is not +// handed to the tests before it can be found there, through the same query the api runs +export async function waitForTokenIndexed(identifier: string, timeoutMs: number = 60000) { + const deadline = Date.now() + timeoutMs; + + while (Date.now() < deadline) { + const response = await axios.get(`${config.elasticUrl}/tokens/_search?q=_id:${identifier}`); + if (response.data?.hits?.hits?.length > 0) { + return; + } + + await new Promise(resolve => setTimeout(resolve, 500)); + } + + throw new Error(`Token ${identifier} was not indexed within ${timeoutMs}ms`); +} + export async function issueEsdt(args: IssueEsdtArgs) { const txHash = await sendTransaction( new SendTransactionArgs({ @@ -99,6 +119,8 @@ export async function issueEsdt(args: IssueEsdtArgs) { console.log( `Issued token with ticker ${args.tokenTicker}. tx hash: ${txHash}. identifier: ${tokenIdentifier}`, ); + + await waitForTokenIndexed(tokenIdentifier); return tokenIdentifier; } @@ -292,6 +314,7 @@ export async function issueCollection(args: IssueNftArgs, type: 'NonFungible' | `Issued ${type} collection with ticker ${args.tokenTicker}. tx hash: ${txHash}. identifier: ${tokenIdentifier}` ); + await waitForTokenIndexed(tokenIdentifier); return tokenIdentifier; } @@ -501,6 +524,7 @@ export async function issueMultipleMetaESDTCollections( ).toString(); metaEsdtCollectionIdentifiers.push({ identifier: tokenIdentifier }); + await waitForTokenIndexed(tokenIdentifier); console.log( `Issued MetaESDT collection ${tokenName}. tx hash: ${txHash}. identifier: ${tokenIdentifier}`, From edaf82212db16c4fc86ad76f5c635b788cef63d0 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 25 Sep 2026 16:46:15 +0300 Subject: [PATCH 8/8] remove comm --- config/config.e2e.mainnet.yaml | 4 ---- 1 file changed, 4 deletions(-) diff --git a/config/config.e2e.mainnet.yaml b/config/config.e2e.mainnet.yaml index cc9d7f531..df7b38c63 100644 --- a/config/config.e2e.mainnet.yaml +++ b/config/config.e2e.mainnet.yaml @@ -39,10 +39,6 @@ features: transactionPool: enabled: true transactionPoolWarmer: - # off, as on mainnet and devnet: every second it resolved the action of each pending transaction and - # smart contract result, including the ESDTTransfer through which a token is issued, before the token - # is indexed. the unresolved token properties were then cached for an hour, and every later transfer - # of that token came back without an action. the pool is still read on request enabled: false cronExpression: '*/1 * * * * *' ttlInSeconds: 10