From e2437a5fdcd0c25f9fedb9e41f5b431ea26ad9ce Mon Sep 17 00:00:00 2001 From: Andrey Sobolev Date: Fri, 17 Oct 2025 12:07:47 +0700 Subject: [PATCH] Fix kafka test --- packages/kafka/src/__test__/queue.spec.ts | 32 +++++++++++++++++------ 1 file changed, 24 insertions(+), 8 deletions(-) diff --git a/packages/kafka/src/__test__/queue.spec.ts b/packages/kafka/src/__test__/queue.spec.ts index 374511db01..789d94e913 100644 --- a/packages/kafka/src/__test__/queue.spec.ts +++ b/packages/kafka/src/__test__/queue.spec.ts @@ -11,12 +11,13 @@ describe('queue', () => { const docsCount = 50 // Reduced from 100 for faster tests try { let msgCount = 0 + let consumerHandle: any const p1 = new Promise((resolve, reject) => { const to = setTimeout(() => { reject(new Error(`Timeout waiting for messages: received ${msgCount}/${docsCount}`)) }, 15000) // Reduced from 100000 - queue.createConsumer( + consumerHandle = queue.createConsumer( testCtx, 'qtest', genId, @@ -27,12 +28,19 @@ describe('queue', () => { resolve() } }, - { retryDelay: 100, maxRetryDelay: 3 } + { retryDelay: 100, maxRetryDelay: 3, fromBegining: true } ) }) - // Wait a bit for consumer to be ready - await new Promise((resolve) => setTimeout(resolve, 500)) + // Wait for consumer to be connected and subscribed + await new Promise((resolve) => { + const checkInterval = setInterval(() => { + if (consumerHandle.isConnected() === true) { + clearInterval(checkInterval) + resolve(undefined) + } + }, 50) + }) const producer = queue.getProducer(testCtx, 'qtest') // Send messages in batches for better performance @@ -60,12 +68,13 @@ describe('queue', () => { try { let counter = 2 + let consumerHandle: any const p = new Promise((resolve, reject) => { const to = setTimeout(() => { reject(new Error(`Timeout waiting for retry processing: counter=${counter}`)) }, 10000) // Added timeout - queue.createConsumer( + consumerHandle = queue.createConsumer( testCtx, 'test', genId, @@ -77,12 +86,19 @@ describe('queue', () => { clearTimeout(to) resolve() }, - { retryDelay: 50, maxRetryDelay: 3 } + { retryDelay: 50, maxRetryDelay: 3, fromBegining: true } ) }) - // Wait a bit for consumer to be ready - await new Promise((resolve) => setTimeout(resolve, 500)) + // Wait for consumer to be connected and subscribed + await new Promise((resolve) => { + const checkInterval = setInterval(() => { + if (consumerHandle.isConnected() === true) { + clearInterval(checkInterval) + resolve(undefined) + } + }, 50) + }) const producer = queue.getProducer(testCtx, 'test') await producer.send(testCtx, genId as any as WorkspaceUuid, ['msg'])