@@ -20,6 +20,8 @@ describe("RunEngine attempt failures", () => {
2020 } ,
2121 queue : {
2222 redis : redisOptions ,
23+ masterQueueConsumersDisabled : true ,
24+ processWorkerQueueDebounceMs : 50 ,
2325 } ,
2426 runLock : {
2527 redis : redisOptions ,
@@ -62,7 +64,7 @@ describe("RunEngine attempt failures", () => {
6264 traceContext : { } ,
6365 traceId : "t12345" ,
6466 spanId : "s12345" ,
65- masterQueue : "main" ,
67+ workerQueue : "main" ,
6668 queue : "task/test-task" ,
6769 isTest : false ,
6870 tags : [ ] ,
@@ -71,10 +73,10 @@ describe("RunEngine attempt failures", () => {
7173 ) ;
7274
7375 //dequeue the run
74- const dequeued = await engine . dequeueFromMasterQueue ( {
76+ await setTimeout ( 500 ) ;
77+ const dequeued = await engine . dequeueFromWorkerQueue ( {
7578 consumerId : "test_12345" ,
76- masterQueue : run . masterQueue ,
77- maxRunCount : 10 ,
79+ workerQueue : "main" ,
7880 } ) ;
7981
8082 //create an attempt
@@ -173,6 +175,8 @@ describe("RunEngine attempt failures", () => {
173175 } ,
174176 queue : {
175177 redis : redisOptions ,
178+ masterQueueConsumersDisabled : true ,
179+ processWorkerQueueDebounceMs : 50 ,
176180 } ,
177181 runLock : {
178182 redis : redisOptions ,
@@ -213,7 +217,7 @@ describe("RunEngine attempt failures", () => {
213217 traceContext : { } ,
214218 traceId : "t12345" ,
215219 spanId : "s12345" ,
216- masterQueue : "main" ,
220+ workerQueue : "main" ,
217221 queue : "task/test-task" ,
218222 isTest : false ,
219223 tags : [ ] ,
@@ -222,10 +226,10 @@ describe("RunEngine attempt failures", () => {
222226 ) ;
223227
224228 //dequeue the run
225- const dequeued = await engine . dequeueFromMasterQueue ( {
229+ await setTimeout ( 500 ) ;
230+ const dequeued = await engine . dequeueFromWorkerQueue ( {
226231 consumerId : "test_12345" ,
227- masterQueue : run . masterQueue ,
228- maxRunCount : 10 ,
232+ workerQueue : "main" ,
229233 } ) ;
230234
231235 //create an attempt
@@ -284,6 +288,8 @@ describe("RunEngine attempt failures", () => {
284288 } ,
285289 queue : {
286290 redis : redisOptions ,
291+ masterQueueConsumersDisabled : true ,
292+ processWorkerQueueDebounceMs : 50 ,
287293 } ,
288294 runLock : {
289295 redis : redisOptions ,
@@ -324,7 +330,7 @@ describe("RunEngine attempt failures", () => {
324330 traceContext : { } ,
325331 traceId : "t12345" ,
326332 spanId : "s12345" ,
327- masterQueue : "main" ,
333+ workerQueue : "main" ,
328334 queue : "task/test-task" ,
329335 isTest : false ,
330336 tags : [ ] ,
@@ -333,10 +339,10 @@ describe("RunEngine attempt failures", () => {
333339 ) ;
334340
335341 //dequeue the run
336- const dequeued = await engine . dequeueFromMasterQueue ( {
342+ await setTimeout ( 500 ) ;
343+ const dequeued = await engine . dequeueFromWorkerQueue ( {
337344 consumerId : "test_12345" ,
338- masterQueue : run . masterQueue ,
339- maxRunCount : 10 ,
345+ workerQueue : "main" ,
340346 } ) ;
341347
342348 //create an attempt
@@ -393,6 +399,8 @@ describe("RunEngine attempt failures", () => {
393399 } ,
394400 queue : {
395401 redis : redisOptions ,
402+ masterQueueConsumersDisabled : true ,
403+ processWorkerQueueDebounceMs : 50 ,
396404 } ,
397405 runLock : {
398406 redis : redisOptions ,
@@ -431,7 +439,7 @@ describe("RunEngine attempt failures", () => {
431439 traceContext : { } ,
432440 traceId : "t12345" ,
433441 spanId : "s12345" ,
434- masterQueue : "main" ,
442+ workerQueue : "main" ,
435443 queue : "task/test-task" ,
436444 isTest : false ,
437445 tags : [ ] ,
@@ -440,10 +448,10 @@ describe("RunEngine attempt failures", () => {
440448 ) ;
441449
442450 //dequeue the run
443- const dequeued = await engine . dequeueFromMasterQueue ( {
451+ await setTimeout ( 500 ) ;
452+ const dequeued = await engine . dequeueFromWorkerQueue ( {
444453 consumerId : "test_12345" ,
445- masterQueue : run . masterQueue ,
446- maxRunCount : 10 ,
454+ workerQueue : "main" ,
447455 } ) ;
448456
449457 //create an attempt
@@ -500,6 +508,8 @@ describe("RunEngine attempt failures", () => {
500508 } ,
501509 queue : {
502510 redis : redisOptions ,
511+ masterQueueConsumersDisabled : true ,
512+ processWorkerQueueDebounceMs : 50 ,
503513 } ,
504514 runLock : {
505515 redis : redisOptions ,
@@ -548,7 +558,7 @@ describe("RunEngine attempt failures", () => {
548558 traceContext : { } ,
549559 traceId : "t12345" ,
550560 spanId : "s12345" ,
551- masterQueue : "main" ,
561+ workerQueue : "main" ,
552562 queue : "task/test-task" ,
553563 isTest : false ,
554564 tags : [ ] ,
@@ -557,10 +567,10 @@ describe("RunEngine attempt failures", () => {
557567 ) ;
558568
559569 //dequeue the run
560- const dequeued = await engine . dequeueFromMasterQueue ( {
570+ await setTimeout ( 500 ) ;
571+ const dequeued = await engine . dequeueFromWorkerQueue ( {
561572 consumerId : "test_12345" ,
562- masterQueue : run . masterQueue ,
563- maxRunCount : 10 ,
573+ workerQueue : "main" ,
564574 } ) ;
565575
566576 //create an attempt
@@ -657,6 +667,8 @@ describe("RunEngine attempt failures", () => {
657667 } ,
658668 queue : {
659669 redis : redisOptions ,
670+ masterQueueConsumersDisabled : true ,
671+ processWorkerQueueDebounceMs : 50 ,
660672 } ,
661673 runLock : {
662674 redis : redisOptions ,
@@ -707,7 +719,7 @@ describe("RunEngine attempt failures", () => {
707719 traceContext : { } ,
708720 traceId : "t12345" ,
709721 spanId : "s12345" ,
710- masterQueue : "main" ,
722+ workerQueue : "main" ,
711723 queue : "task/test-task" ,
712724 isTest : false ,
713725 tags : [ ] ,
@@ -716,10 +728,10 @@ describe("RunEngine attempt failures", () => {
716728 ) ;
717729
718730 //dequeue the run
719- const dequeued = await engine . dequeueFromMasterQueue ( {
731+ await setTimeout ( 500 ) ;
732+ const dequeued = await engine . dequeueFromWorkerQueue ( {
720733 consumerId : "test_12345" ,
721- masterQueue : run . masterQueue ,
722- maxRunCount : 10 ,
734+ workerQueue : "main" ,
723735 } ) ;
724736
725737 //create first attempt
@@ -762,10 +774,10 @@ describe("RunEngine attempt failures", () => {
762774 await setTimeout ( 5_000 ) ;
763775
764776 //dequeue again
765- const dequeued2 = await engine . dequeueFromMasterQueue ( {
777+ await setTimeout ( 500 ) ;
778+ const dequeued2 = await engine . dequeueFromWorkerQueue ( {
766779 consumerId : "test_12345" ,
767- masterQueue : run . masterQueue ,
768- maxRunCount : 10 ,
780+ workerQueue : "main" ,
769781 } ) ;
770782 expect ( dequeued2 . length ) . toBe ( 1 ) ;
771783
0 commit comments