Skip to content

Commit e77a0ae

Browse files
committed
Fix flaky persistent subscription nack tests
Terminate only after the finish event and all retry redeliveries are processed, instead of racing the asynchronous retry redeliveries against the first-pass finish event.
1 parent f20b68f commit e77a0ae

1 file changed

Lines changed: 23 additions & 2 deletions

File tree

packages/test/src/persistentSubscription/subscribeToPersistentSubscriptionToStream.test.ts

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -234,6 +234,14 @@ describe("subscribeToPersistentSubscriptionToStream", () => {
234234
const defer = new Defer();
235235

236236
const nacked: string[] = [];
237+
const ackedRetries = new Set<string>();
238+
let finishSeen = false;
239+
240+
const maybeResolve = () => {
241+
if (finishSeen && ackedRetries.size === retryCount) {
242+
defer.resolve();
243+
}
244+
};
237245

238246
const onError = jest.fn((error) => {
239247
defer.reject(error);
@@ -246,7 +254,8 @@ describe("subscribeToPersistentSubscriptionToStream", () => {
246254

247255
if (event.event.type === "finish-test") {
248256
await subscription.ack(event);
249-
defer.resolve();
257+
finishSeen = true;
258+
maybeResolve();
250259
return;
251260
}
252261

@@ -261,6 +270,10 @@ describe("subscribeToPersistentSubscriptionToStream", () => {
261270
}
262271

263272
await subscription.ack(event);
273+
if (event.event.type === "retry-event") {
274+
ackedRetries.add(event.event.id);
275+
}
276+
maybeResolve();
264277
return;
265278
});
266279

@@ -331,6 +344,8 @@ describe("subscribeToPersistentSubscriptionToStream", () => {
331344
const GROUP_NAME = "async_iter_nack_group_name";
332345
const doSomething = jest.fn();
333346
const nacked: string[] = [];
347+
const ackedRetries = new Set<string>();
348+
let finishSeen = false;
334349

335350
// Skip the first twenty events and retry the next 20 events.
336351
// we should see the number of times that the `onEvent` callback
@@ -365,7 +380,9 @@ describe("subscribeToPersistentSubscriptionToStream", () => {
365380

366381
if (resolvedEvent.event.type === "finish-test") {
367382
await subscription.ack(resolvedEvent);
368-
break;
383+
finishSeen = true;
384+
if (ackedRetries.size === retryCount) break;
385+
continue;
369386
}
370387

371388
if (!nacked.includes(resolvedEvent.event.id)) {
@@ -379,6 +396,10 @@ describe("subscribeToPersistentSubscriptionToStream", () => {
379396
}
380397

381398
await subscription.ack(resolvedEvent);
399+
if (resolvedEvent.event.type === "retry-event") {
400+
ackedRetries.add(resolvedEvent.event.id);
401+
}
402+
if (finishSeen && ackedRetries.size === retryCount) break;
382403
}
383404

384405
expect(doSomething).toBeCalledTimes(

0 commit comments

Comments
 (0)