diff --git a/config/config.e2e.mainnet.yaml b/config/config.e2e.mainnet.yaml index cd9243ae8..df7b38c63 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: @@ -39,7 +39,7 @@ features: transactionPool: enabled: true transactionPoolWarmer: - enabled: true + enabled: false cronExpression: '*/1 * * * * *' ttlInSeconds: 10 updateCollectionExtraDetails: 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/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 0e9e7bb64..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'; @@ -32,7 +33,7 @@ export async function getNonce( return currentNonceResponse.data.data.nonce; } catch (e) { console.error(e); - return 0; + throw e; } } @@ -57,16 +58,39 @@ 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; } } +// 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({ @@ -95,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; } @@ -160,8 +186,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; } } @@ -287,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; } @@ -496,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}`, diff --git a/src/test/chain-simulator/utils/prepare-test-data.ts b/src/test/chain-simulator/utils/prepare-test-data.ts index e02a454b9..5a6badff3 100644 --- a/src/test/chain-simulator/utils/prepare-test-data.ts +++ b/src/test/chain-simulator/utils/prepare-test-data.ts @@ -25,7 +25,17 @@ async function prepareTestData() { await ChainSimulatorUtils.deployPingPongSc(config.aliceAddress); console.log('āœ“ Deployed PingPong smart contract'); - await new Promise((resolve) => setTimeout(resolve, 30000)); + 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 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 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 f1a70a4a8..e1ba59f66 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', ]; @@ -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 7721f4508..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[], @@ -80,8 +94,78 @@ 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. 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) => { + const deadline = Date.now() + timeoutMs; + + while (Date.now() < deadline) { + if (clients.every(client => client.connected)) { + return; + } + + for (const client of failedClients) { + failedClients.delete(client); + 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'}`); + }; + + // 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, @@ -89,75 +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, { - path: '/ws/subscription', - }); - clients.push(client); - - client.on("connect_error", (err) => { - throw new Error(`${clientId} connection failed: ${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, { - path: '/ws/subscription', + 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"); }); - clients.push(client); - - client.on("connect_error", (err) => { throw new Error(`${clientId} connection failed: ${err.message}`); }); 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 () => { @@ -185,6 +250,17 @@ describe('Websocket subscriptions e2e tests', () => { connectAndSubscribeGeneral("client4", client4SubscriptionConfig); + 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 ---"); @@ -207,7 +283,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); @@ -216,69 +300,77 @@ 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', () => { - 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); @@ -293,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);