@fluojs/queue
Advanced tools
| import type { Container } from '@fluojs/di'; | ||
| import { type ApplicationLogger, type CompiledModule, type OnApplicationBootstrap, type OnApplicationShutdown, type OnModuleDestroy } from '@fluojs/runtime'; | ||
| import type { ApplicationLogger, CompiledModule, OnApplicationBootstrap, OnApplicationShutdown, OnModuleDestroy } from '@fluojs/runtime'; | ||
| import type { NormalizedQueueModuleOptions, Queue } from './types.js'; | ||
@@ -62,2 +62,3 @@ /** | ||
| private shutdown; | ||
| private waitForInFlightStartup; | ||
| private closeInitializedResources; | ||
@@ -64,0 +65,0 @@ private tryCloseWorker; |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"service.d.ts","sourceRoot":"","sources":["../src/service.ts"],"names":[],"mappings":"AAEA,OAAO,KAAK,EAAE,SAAS,EAAE,MAAM,YAAY,CAAC;AAE5C,OAAO,EACL,KAAK,iBAAiB,EACtB,KAAK,cAAc,EACnB,KAAK,sBAAsB,EAC3B,KAAK,qBAAqB,EAC1B,KAAK,eAAe,EACrB,MAAM,iBAAiB,CAAC;AASzB,OAAO,KAAK,EACV,4BAA4B,EAC5B,KAAK,EAIN,MAAM,YAAY,CAAC;AAkGpB;;;;;GAKG;AACH,qBACa,qBAAsB,YAAW,KAAK,EAAE,sBAAsB,EAAE,qBAAqB,EAAE,eAAe;IAY/G,OAAO,CAAC,QAAQ,CAAC,OAAO;IACxB,OAAO,CAAC,QAAQ,CAAC,gBAAgB;IACjC,OAAO,CAAC,QAAQ,CAAC,eAAe;IAChC,OAAO,CAAC,QAAQ,CAAC,MAAM;IAdzB,OAAO,CAAC,QAAQ,CAAC,oBAAoB,CAAkD;IACvF,OAAO,CAAC,QAAQ,CAAC,eAAe,CAAoC;IACpE,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAqC;IACtE,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAA8B;IAC/D,OAAO,CAAC,QAAQ,CAAC,iBAAiB,CAAyB;IAC3D,OAAO,CAAC,cAAc,CAA+B;IACrD,OAAO,CAAC,WAAW,CAA+B;IAClD,OAAO,CAAC,YAAY,CAA4B;IAChD,OAAO,CAAC,eAAe,CAA4B;gBAGhC,OAAO,EAAE,4BAA4B,EACrC,gBAAgB,EAAE,SAAS,EAC3B,eAAe,EAAE,SAAS,cAAc,EAAE,EAC1C,MAAM,EAAE,iBAAiB;IAKtC,sBAAsB,IAAI,OAAO,CAAC,IAAI,CAAC;IAIvC,qBAAqB,IAAI,OAAO,CAAC,IAAI,CAAC;IAItC,eAAe,IAAI,OAAO,CAAC,IAAI,CAAC;IAItC;;;;;;;OAOG;IACG,OAAO,CAAC,IAAI,SAAS,MAAM,EAAE,GAAG,EAAE,IAAI,GAAG,OAAO,CAAC,MAAM,CAAC;IAuB9D;;;;OAIG;IACH,4BAA4B;YAWd,aAAa;YAwBb,cAAc;YAad,oBAAoB;YAOpB,kBAAkB;IAgBhC,OAAO,CAAC,cAAc;YAQR,iBAAiB;YAOjB,yBAAyB;IAyBvC,OAAO,CAAC,mBAAmB;IAS3B,OAAO,CAAC,oBAAoB;IAa5B,OAAO,CAAC,mBAAmB;IAkB3B,OAAO,CAAC,0BAA0B;IAkBlC,OAAO,CAAC,0BAA0B;IASlC,OAAO,CAAC,yBAAyB;YASnB,kCAAkC;YAkBlC,qBAAqB;YAiBrB,aAAa;YAOb,oBAAoB;IAyBlC,OAAO,CAAC,sBAAsB;YAQhB,QAAQ;YAsBR,yBAAyB;YAqBzB,cAAc;YAQd,aAAa;YAQb,uBAAuB;CAOtC"} | ||
| {"version":3,"file":"service.d.ts","sourceRoot":"","sources":["../src/service.ts"],"names":[],"mappings":"AAEA,OAAO,KAAK,EAAE,SAAS,EAAE,MAAM,YAAY,CAAC;AAE5C,OAAO,KAAK,EACV,iBAAiB,EACjB,cAAc,EACd,sBAAsB,EACtB,qBAAqB,EACrB,eAAe,EAChB,MAAM,iBAAiB,CAAC;AASzB,OAAO,KAAK,EACV,4BAA4B,EAC5B,KAAK,EAIN,MAAM,YAAY,CAAC;AAkGpB;;;;;GAKG;AACH,qBACa,qBAAsB,YAAW,KAAK,EAAE,sBAAsB,EAAE,qBAAqB,EAAE,eAAe;IAY/G,OAAO,CAAC,QAAQ,CAAC,OAAO;IACxB,OAAO,CAAC,QAAQ,CAAC,gBAAgB;IACjC,OAAO,CAAC,QAAQ,CAAC,eAAe;IAChC,OAAO,CAAC,QAAQ,CAAC,MAAM;IAdzB,OAAO,CAAC,QAAQ,CAAC,oBAAoB,CAAkD;IACvF,OAAO,CAAC,QAAQ,CAAC,eAAe,CAAoC;IACpE,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAAqC;IACtE,OAAO,CAAC,QAAQ,CAAC,gBAAgB,CAA8B;IAC/D,OAAO,CAAC,QAAQ,CAAC,iBAAiB,CAAyB;IAC3D,OAAO,CAAC,cAAc,CAA+B;IACrD,OAAO,CAAC,WAAW,CAA+B;IAClD,OAAO,CAAC,YAAY,CAA4B;IAChD,OAAO,CAAC,eAAe,CAA4B;gBAGhC,OAAO,EAAE,4BAA4B,EACrC,gBAAgB,EAAE,SAAS,EAC3B,eAAe,EAAE,SAAS,cAAc,EAAE,EAC1C,MAAM,EAAE,iBAAiB;IAKtC,sBAAsB,IAAI,OAAO,CAAC,IAAI,CAAC;IAIvC,qBAAqB,IAAI,OAAO,CAAC,IAAI,CAAC;IAItC,eAAe,IAAI,OAAO,CAAC,IAAI,CAAC;IAItC;;;;;;;OAOG;IACG,OAAO,CAAC,IAAI,SAAS,MAAM,EAAE,GAAG,EAAE,IAAI,GAAG,OAAO,CAAC,MAAM,CAAC;IA2B9D;;;;OAIG;IACH,4BAA4B;YAWd,aAAa;YAwBb,cAAc;YAed,oBAAoB;YASpB,kBAAkB;IAgBhC,OAAO,CAAC,cAAc;YAQR,iBAAiB;YAOjB,yBAAyB;IAyBvC,OAAO,CAAC,mBAAmB;IAS3B,OAAO,CAAC,oBAAoB;IAa5B,OAAO,CAAC,mBAAmB;IAkB3B,OAAO,CAAC,0BAA0B;IAkBlC,OAAO,CAAC,0BAA0B;IASlC,OAAO,CAAC,yBAAyB;YASnB,kCAAkC;YAkBlC,qBAAqB;YAiBrB,aAAa;YAOb,oBAAoB;IAyBlC,OAAO,CAAC,sBAAsB;YAQhB,QAAQ;YAyBR,sBAAsB;YAgBtB,yBAAyB;YAqBzB,cAAc;YAQd,aAAa;YAQb,uBAAuB;CAOtC"} |
+24
-2
@@ -107,2 +107,5 @@ let _initClass; | ||
| await this.ensureStarted(); | ||
| if (this.lifecycleState !== 'started') { | ||
| throw new Error(`Queue lifecycle state is ${this.lifecycleState}.`); | ||
| } | ||
| const descriptor = this.descriptorsByJobType.get(job.constructor); | ||
@@ -165,7 +168,11 @@ if (!descriptor) { | ||
| await this.initializeWorkers(redis); | ||
| this.lifecycleState = 'started'; | ||
| if (this.lifecycleState === 'starting') { | ||
| this.lifecycleState = 'started'; | ||
| } | ||
| } | ||
| async handleStartupFailure() { | ||
| await this.closeInitializedResources(); | ||
| this.lifecycleState = 'idle'; | ||
| if (this.lifecycleState === 'starting') { | ||
| this.lifecycleState = 'idle'; | ||
| } | ||
| this.redisClient = undefined; | ||
@@ -319,5 +326,7 @@ this.startPromise = undefined; | ||
| this.shutdownPromise = (async () => { | ||
| await this.waitForInFlightStartup(); | ||
| await this.closeInitializedResources(); | ||
| await this.deadLetterManager.drainPendingWrites(); | ||
| this.lifecycleState = 'stopped'; | ||
| this.redisClient = undefined; | ||
| this.startPromise = undefined; | ||
@@ -327,2 +336,15 @@ })(); | ||
| } | ||
| async waitForInFlightStartup() { | ||
| const startup = this.startPromise; | ||
| if (!startup) { | ||
| return; | ||
| } | ||
| try { | ||
| await startup; | ||
| } catch { | ||
| // ensureStarted() owns startup rollback and preserves the original | ||
| // bootstrap error. Shutdown still continues so partially registered | ||
| // resources cannot outlive the application lifecycle. | ||
| } | ||
| } | ||
| async closeInitializedResources() { | ||
@@ -329,0 +351,0 @@ const workers = Array.from(this.workersByJobName.values()); |
+5
-5
@@ -13,3 +13,3 @@ { | ||
| ], | ||
| "version": "1.0.0-beta.4", | ||
| "version": "1.0.0-beta.5", | ||
| "private": false, | ||
@@ -42,6 +42,6 @@ "license": "MIT", | ||
| "bullmq": "^5.58.0", | ||
| "@fluojs/core": "^1.0.0-beta.4", | ||
| "@fluojs/di": "^1.0.0-beta.6", | ||
| "@fluojs/redis": "^1.0.0-beta.3", | ||
| "@fluojs/runtime": "^1.0.0-beta.11" | ||
| "@fluojs/core": "^1.0.0-beta.5", | ||
| "@fluojs/di": "^1.0.0-beta.7", | ||
| "@fluojs/redis": "^1.0.0-beta.4", | ||
| "@fluojs/runtime": "^1.0.0-beta.12" | ||
| }, | ||
@@ -48,0 +48,0 @@ "devDependencies": { |
+10
-10
@@ -62,2 +62,11 @@ # @fluojs/queue | ||
| @Inject(QueueLifecycleService) | ||
| export class OrderService { | ||
| constructor(private readonly queue: QueueLifecycleService) {} | ||
| async placeOrder(id: string) { | ||
| await this.queue.enqueue(new ProcessOrderJob(id)); | ||
| } | ||
| } | ||
| @Module({ | ||
@@ -68,14 +77,5 @@ imports: [ | ||
| ], | ||
| providers: [OrderWorker], | ||
| providers: [OrderService, OrderWorker], | ||
| }) | ||
| export class AppModule {} | ||
| export class OrderService { | ||
| @Inject(QueueLifecycleService) | ||
| private readonly queue: QueueLifecycleService; | ||
| async placeOrder(id: string) { | ||
| await this.queue.enqueue(new ProcessOrderJob(id)); | ||
| } | ||
| } | ||
| ``` | ||
@@ -82,0 +82,0 @@ |
+10
-10
@@ -62,2 +62,11 @@ # @fluojs/queue | ||
| @Inject(QueueLifecycleService) | ||
| export class OrderService { | ||
| constructor(private readonly queue: QueueLifecycleService) {} | ||
| async placeOrder(id: string) { | ||
| await this.queue.enqueue(new ProcessOrderJob(id)); | ||
| } | ||
| } | ||
| @Module({ | ||
@@ -68,14 +77,5 @@ imports: [ | ||
| ], | ||
| providers: [OrderWorker], | ||
| providers: [OrderService, OrderWorker], | ||
| }) | ||
| export class AppModule {} | ||
| export class OrderService { | ||
| @Inject(QueueLifecycleService) | ||
| private readonly queue: QueueLifecycleService; | ||
| async placeOrder(id: string) { | ||
| await this.queue.enqueue(new ProcessOrderJob(id)); | ||
| } | ||
| } | ||
| ``` | ||
@@ -82,0 +82,0 @@ |
73136
1.04%1294
1.81%Updated
Updated
Updated