mirror of
https://github.com/ever-co/ever-gauzy.git
synced 2026-10-02 01:54:50 +08:00
[Feat] Background scheduling system with BullMQ (#9504)
* feat: add worker application and scheduler package Initialize a new NestJS-based worker application and a scheduler library within the workspace. This includes setting up the basic application structure, Nx project configurations, and updating the root TypeScript path mappings for the scheduler package. * feat(core): add redis module exports Export the redis module from the core package to make it available for use across the application. This includes adding the barrel export in the main entry point and creating the redis module index. * build: add bullmq dependency Add bullmq to the project dependencies to support background job processing and message queuing. * feat: add bullmq dependency and update related packages in yarn.lock * chore(deps): add @nestjs/bullmq Add @nestjs/bullmq version 11.0.4 to the project dependencies. * feat(scheduler): implement queue-aware background job system Introduce a robust scheduling and queueing system powered by BullMQ and NestJS. This update adds decorators for defining scheduled jobs and queue processors, along with a discovery mechanism to automatically register them across the application. - Add `@ScheduledJob`, `@QueueWorker`, and `@QueueJobHandler` decorators. - Implement a discovery service to automatically register jobs and processors from providers and controllers. - Integrate BullMQ for distributed job processing and persistence. - Refactor the worker application from a web server to a headless application context with signal handling. - Add worker lifecycle management through startup and heartbeat jobs. - Include support for job retries, timeouts, and overlap prevention. * chore(worker): disable redis integration and update heartbeat frequency Comment out Redis-related modules and services in the worker application. Increase the heartbeat job frequency from every 10 minutes to every 30 seconds for more frequent status checks. Additionally, refactor core package exports and remove redundant tsconfig path mappings. * chore(worker): remove Redis integration and clean up app module * feat(scheduler): enhance README with detailed usage, job definitions, and options for background job scheduling * fix(scheduler): return execution result from discovery handler Update the task handler in SchedulerDiscoveryService to return the result of the executed method. This ensures that any return value from the scheduled task is correctly propagated and can be captured by the caller. * refactor(scheduler): use crypto.randomInt for random delay Replace Math.random() with crypto.randomInt() in the randomDelay helper to provide a more robust source of randomness. This change also includes minor formatting cleanup and improves the handling of error messages and stack traces during job execution. * feat(cspell): add "bullmq" to custom words list * refactor(core): remove unused RedisModule import Remove the unused RedisModule import from the core index file to clean up the public API surface. * refactor(worker): use node: prefix for built-in imports Update the imports for fs and path to use the explicit node: protocol and reorder imports for consistency. * fix(tsconfig): correct extends path and update compiler options
This commit is contained in:
+2
-1
@@ -852,7 +852,8 @@
|
||||
"Hasher",
|
||||
"VARCHAR",
|
||||
"TOCTOU",
|
||||
"Gmqh"
|
||||
"Gmqh",
|
||||
"bullmq"
|
||||
],
|
||||
"useGitignore": true,
|
||||
"ignorePaths": [
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
{
|
||||
"name": "worker",
|
||||
"$schema": "../../node_modules/nx/schemas/project-schema.json",
|
||||
"sourceRoot": "apps/worker/src",
|
||||
"projectType": "application",
|
||||
"targets": {
|
||||
"build": {
|
||||
"executor": "@nx/webpack:webpack",
|
||||
"outputs": ["{options.outputPath}"],
|
||||
"defaultConfiguration": "production",
|
||||
"options": {
|
||||
"target": "node",
|
||||
"compiler": "tsc",
|
||||
"outputPath": "dist/apps/worker",
|
||||
"main": "apps/worker/src/main.ts",
|
||||
"tsConfig": "apps/worker/tsconfig.app.json",
|
||||
"assets": ["apps/worker/src/assets"],
|
||||
"webpackConfig": "apps/worker/webpack.config.js",
|
||||
"generatePackageJson": true
|
||||
},
|
||||
"configurations": {
|
||||
"development": {
|
||||
"outputHashing": "none"
|
||||
},
|
||||
"production": {}
|
||||
}
|
||||
},
|
||||
"prune-lockfile": {
|
||||
"dependsOn": ["build"],
|
||||
"cache": true,
|
||||
"executor": "@nx/js:prune-lockfile",
|
||||
"outputs": ["{workspaceRoot}/dist/apps/worker/package.json", "{workspaceRoot}/dist/apps/worker/yarn.lock"],
|
||||
"options": {
|
||||
"buildTarget": "build"
|
||||
}
|
||||
},
|
||||
"copy-workspace-modules": {
|
||||
"dependsOn": ["build"],
|
||||
"cache": true,
|
||||
"outputs": ["{workspaceRoot}/dist/apps/worker/workspace_modules"],
|
||||
"executor": "@nx/js:copy-workspace-modules",
|
||||
"options": {
|
||||
"buildTarget": "build"
|
||||
}
|
||||
},
|
||||
"prune": {
|
||||
"dependsOn": ["prune-lockfile", "copy-workspace-modules"],
|
||||
"executor": "nx:noop"
|
||||
},
|
||||
"serve": {
|
||||
"continuous": true,
|
||||
"executor": "@nx/js:node",
|
||||
"defaultConfiguration": "development",
|
||||
"dependsOn": ["build"],
|
||||
"options": {
|
||||
"buildTarget": "worker:build",
|
||||
"runBuildTargetDependencies": false
|
||||
},
|
||||
"configurations": {
|
||||
"development": {
|
||||
"buildTarget": "worker:build:development"
|
||||
},
|
||||
"production": {
|
||||
"buildTarget": "worker:build:production"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"tags": []
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
import { SchedulerModule } from '@gauzy/scheduler';
|
||||
import { Module } from '@nestjs/common';
|
||||
import { WorkerJobsModule } from './worker-jobs.module';
|
||||
import { WORKER_DEFAULT_QUEUE, WORKER_QUEUE_ENABLED } from './worker.constants';
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
SchedulerModule.forRoot({
|
||||
enabled: process.env.WORKER_SCHEDULER_ENABLED !== 'false',
|
||||
enableQueueing: WORKER_QUEUE_ENABLED,
|
||||
defaultQueueName: WORKER_DEFAULT_QUEUE,
|
||||
defaultTimezone: process.env.WORKER_TIMEZONE,
|
||||
defaultJobOptions: {
|
||||
preventOverlap: true,
|
||||
retries: 1,
|
||||
retryDelayMs: 5000
|
||||
}
|
||||
}),
|
||||
WorkerJobsModule
|
||||
],
|
||||
})
|
||||
export class AppModule {}
|
||||
@@ -0,0 +1,34 @@
|
||||
import { ScheduledJob } from '@gauzy/scheduler';
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { CronExpression } from '@nestjs/schedule';
|
||||
import { WORKER_DEFAULT_QUEUE, WORKER_QUEUE_ENABLED } from '../worker.constants';
|
||||
|
||||
@Injectable()
|
||||
export class WorkerLifecycleJob {
|
||||
private readonly logger = new Logger(WorkerLifecycleJob.name);
|
||||
|
||||
@ScheduledJob({
|
||||
name: 'worker.startup.scheduler',
|
||||
enabled: WORKER_QUEUE_ENABLED,
|
||||
runOnStart: true,
|
||||
preventOverlap: true,
|
||||
queueName: WORKER_DEFAULT_QUEUE,
|
||||
queueJobName: 'worker.startup'
|
||||
})
|
||||
async announceStartup(): Promise<{ timestamp: string }> {
|
||||
this.logger.log('Queueing worker startup job.');
|
||||
return { timestamp: new Date().toISOString() };
|
||||
}
|
||||
|
||||
@ScheduledJob({
|
||||
name: 'worker.heartbeat.scheduler',
|
||||
enabled: WORKER_QUEUE_ENABLED && process.env.WORKER_HEARTBEAT_ENABLED !== 'false',
|
||||
cron: CronExpression.EVERY_30_SECONDS,
|
||||
queueName: WORKER_DEFAULT_QUEUE,
|
||||
queueJobName: 'worker.heartbeat'
|
||||
})
|
||||
async heartbeat(): Promise<{ timestamp: string }> {
|
||||
this.logger.log('Queueing worker heartbeat job...');
|
||||
return { timestamp: new Date().toISOString() };
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { Job } from 'bullmq';
|
||||
import { QueueJobHandler, QueueWorker, QueueWorkerHost } from '@gauzy/scheduler';
|
||||
import { WORKER_DEFAULT_QUEUE } from '../worker.constants';
|
||||
|
||||
interface WorkerLifecyclePayload {
|
||||
timestamp: string;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
@QueueWorker(WORKER_DEFAULT_QUEUE)
|
||||
export class WorkerLifecycleProcessor extends QueueWorkerHost {
|
||||
private readonly logger = new Logger(WorkerLifecycleProcessor.name);
|
||||
|
||||
@QueueJobHandler('worker.startup')
|
||||
async handleStartup(job: Job<WorkerLifecyclePayload>): Promise<void> {
|
||||
this.logger.log(`Worker startup event processed at ${job.data.timestamp}`);
|
||||
}
|
||||
|
||||
@QueueJobHandler('worker.heartbeat')
|
||||
async handleHeartbeat(job: Job<WorkerLifecyclePayload>): Promise<void> {
|
||||
this.logger.log(`Worker heartbeat processed at ${job.data.timestamp}`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { SchedulerModule } from '@gauzy/scheduler';
|
||||
import { WorkerLifecycleJob } from './jobs/worker-lifecycle.job';
|
||||
import { WorkerLifecycleProcessor } from './processors/worker-lifecycle.processor';
|
||||
import { WORKER_DEFAULT_QUEUE, WORKER_QUEUE_ENABLED } from './worker.constants';
|
||||
|
||||
const providers = WORKER_QUEUE_ENABLED
|
||||
? [WorkerLifecycleJob, WorkerLifecycleProcessor]
|
||||
: [WorkerLifecycleJob];
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
SchedulerModule.forFeature({
|
||||
jobProviders: providers,
|
||||
queues: WORKER_QUEUE_ENABLED ? [WORKER_DEFAULT_QUEUE] : []
|
||||
})
|
||||
]
|
||||
})
|
||||
export class WorkerJobsModule {}
|
||||
@@ -0,0 +1,2 @@
|
||||
export const WORKER_DEFAULT_QUEUE = process.env.WORKER_DEFAULT_QUEUE || 'worker-default';
|
||||
export const WORKER_QUEUE_ENABLED = process.env.WORKER_QUEUE_ENABLED !== 'false' && process.env.REDIS_ENABLED === 'true';
|
||||
@@ -0,0 +1,20 @@
|
||||
import * as dotenv from 'dotenv';
|
||||
import * as fs from 'node:fs';
|
||||
import * as path from 'node:path';
|
||||
|
||||
function loadEnvFile(envPath: string, options: { override?: boolean } = {}): void {
|
||||
if (!fs.existsSync(envPath)) {
|
||||
return;
|
||||
}
|
||||
|
||||
dotenv.config({ path: envPath, quiet: true, ...options });
|
||||
}
|
||||
|
||||
export function loadEnv(): void {
|
||||
const cwd = process.cwd();
|
||||
const envPath = path.resolve(cwd, '.env');
|
||||
const envLocalPath = path.resolve(cwd, '.env.local');
|
||||
|
||||
loadEnvFile(envPath, { override: true });
|
||||
loadEnvFile(envLocalPath);
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
import { Logger } from '@nestjs/common';
|
||||
import { NestFactory } from '@nestjs/core';
|
||||
import { loadEnv } from './load-env';
|
||||
|
||||
async function bootstrap() {
|
||||
loadEnv();
|
||||
const { AppModule } = await import('./app/app.module');
|
||||
|
||||
const app = await NestFactory.createApplicationContext(AppModule, {
|
||||
logger: ['log', 'error', 'warn', 'debug', 'verbose']
|
||||
});
|
||||
|
||||
Logger.log('Worker application started.', 'WorkerBootstrap');
|
||||
|
||||
const shutdown = async (signal: string): Promise<void> => {
|
||||
Logger.log(`Received ${signal}, closing worker...`, 'WorkerBootstrap');
|
||||
await app.close();
|
||||
process.exit(0);
|
||||
};
|
||||
|
||||
process.on('SIGINT', () => {
|
||||
void shutdown('SIGINT');
|
||||
});
|
||||
process.on('SIGTERM', () => {
|
||||
void shutdown('SIGTERM');
|
||||
});
|
||||
}
|
||||
|
||||
void bootstrap();
|
||||
@@ -0,0 +1,13 @@
|
||||
{
|
||||
"extends": "./tsconfig.json",
|
||||
"compilerOptions": {
|
||||
"outDir": "../../dist/out-tsc",
|
||||
"module": "commonjs",
|
||||
"types": ["node"],
|
||||
"emitDecoratorMetadata": true,
|
||||
"target": "es2021",
|
||||
"removeComments": true,
|
||||
"resolveJsonModule": true
|
||||
},
|
||||
"include": ["src/**/*.ts"]
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"extends": "../../tsconfig.json",
|
||||
"files": [],
|
||||
"include": [],
|
||||
"references": [
|
||||
{
|
||||
"path": "./tsconfig.app.json"
|
||||
}
|
||||
],
|
||||
"compilerOptions": {
|
||||
"incremental": true,
|
||||
"types": ["node", "jest"]
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
const { composePlugins, withNx } = require('@nx/webpack');
|
||||
|
||||
// Nx plugins for webpack.
|
||||
module.exports = composePlugins(
|
||||
withNx({
|
||||
target: 'node'
|
||||
}),
|
||||
(config) => {
|
||||
config.output = {
|
||||
...config.output,
|
||||
...(process.env.NODE_ENV !== 'production' && {
|
||||
clean: true,
|
||||
devtoolModuleFilenameTemplate: '[absolute-resource-path]'
|
||||
})
|
||||
};
|
||||
config.devtool = 'source-map';
|
||||
// Update the webpack config as needed here.
|
||||
// e.g. `config.plugins.push(new MyPlugin())`
|
||||
return config;
|
||||
}
|
||||
);
|
||||
@@ -580,6 +580,7 @@
|
||||
"@nebular/bootstrap": "^9.1.0-rc.6",
|
||||
"@nebular/security": "^17.0.0",
|
||||
"@nebular/theme": "^17.0.0",
|
||||
"@nestjs/bullmq": "^11.0.4",
|
||||
"@nestjs/common": "^11.1.14",
|
||||
"@nestjs/core": "^11.1.14",
|
||||
"@nestjs/mapped-types": "^2.1.0",
|
||||
@@ -593,6 +594,7 @@
|
||||
"angular2-smart-table": "^4.1.1",
|
||||
"autoprefixer": "^10.4.20",
|
||||
"bcrypt": "^5.1.1",
|
||||
"bullmq": "^5.70.1",
|
||||
"dotenv": "^17.2.4",
|
||||
"jsdom": "^23.0.1",
|
||||
"lodash-es": "^4.17.21",
|
||||
|
||||
@@ -17,6 +17,7 @@ export { LazyFileInterceptor } from './lib/core/interceptors';
|
||||
export { FileStorage, FileStorageFactory, UploadedFileStorage } from './lib/core/file-storage';
|
||||
export * from './lib/shared';
|
||||
export * from './lib/event-bus';
|
||||
export { RedisModule, EVER_REDIS_CLIENT } from './lib/redis';
|
||||
|
||||
export * from './lib/tenant';
|
||||
export { RoleModule, RoleService } from './lib/role';
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
export * from './redis.module';
|
||||
@@ -0,0 +1,197 @@
|
||||
# @gauzy/scheduler
|
||||
|
||||
Reusable background jobs for NestJS modules using:
|
||||
|
||||
- `@nestjs/schedule` for cron/interval triggers
|
||||
- `@nestjs/bullmq` + `bullmq` for queue-backed execution
|
||||
|
||||
## What you get
|
||||
|
||||
- `@ScheduledJob(...)` for declarative background jobs
|
||||
- Queue-aware jobs (`queueName`, `queueJobName`, `queueJobOptions`)
|
||||
- Worker abstraction: `@QueueWorker(...)` + `QueueWorkerHost` + `@QueueJobHandler(...)`
|
||||
- Auto-discovery of decorated jobs from providers/controllers
|
||||
- Runtime APIs via `SchedulerService` (`listJobs`, `triggerNow`, `enqueue`)
|
||||
- Safety defaults (overlap prevention, retries, timeout, jitter)
|
||||
|
||||
## Usage model
|
||||
|
||||
Use the scheduler in 2 layers:
|
||||
|
||||
1. Root setup once in your application (or worker application) with `SchedulerModule.forRoot(...)`.
|
||||
2. Feature setup per domain with `SchedulerModule.forFeature(...)` to register queues and job providers.
|
||||
|
||||
## 1) Root setup (global)
|
||||
|
||||
```ts
|
||||
import { Module } from '@nestjs/common';
|
||||
import { SchedulerModule } from '@gauzy/scheduler';
|
||||
import { TokenCleanupJobsModule } from './token-cleanup-jobs.module';
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
SchedulerModule.forRoot({
|
||||
enabled: true,
|
||||
enableQueueing: true,
|
||||
defaultTimezone: 'UTC',
|
||||
defaultQueueName: 'background',
|
||||
defaultJobOptions: {
|
||||
enabled: true,
|
||||
preventOverlap: true,
|
||||
retries: 1,
|
||||
retryDelayMs: 5000,
|
||||
timeoutMs: 30000,
|
||||
maxRandomDelayMs: 0
|
||||
}
|
||||
}),
|
||||
TokenCleanupJobsModule
|
||||
]
|
||||
})
|
||||
export class AppModule {}
|
||||
```
|
||||
|
||||
Notes:
|
||||
|
||||
- `enableQueueing` defaults to `process.env.REDIS_ENABLED === 'true'`.
|
||||
- If `queueConnection` is omitted, scheduler resolves Redis connection from:
|
||||
- `REDIS_URL`, or
|
||||
- `REDIS_HOST`, `REDIS_PORT`, `REDIS_USER`, `REDIS_PASSWORD`, `REDIS_TLS`.
|
||||
|
||||
## 2) Scheduled producer job
|
||||
|
||||
Decorate a method with `@ScheduledJob`. If `queueName` is set, the method return value is enqueued as BullMQ job data.
|
||||
|
||||
```ts
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { CronExpression } from '@nestjs/schedule';
|
||||
import { ScheduledJob } from '@gauzy/scheduler';
|
||||
|
||||
@Injectable()
|
||||
export class TokenCleanupScheduler {
|
||||
private readonly logger = new Logger(TokenCleanupScheduler.name);
|
||||
|
||||
@ScheduledJob({
|
||||
name: 'token.cleanup.expired.scheduler',
|
||||
cron: CronExpression.EVERY_HOUR,
|
||||
queueName: 'token-maintenance',
|
||||
queueJobName: 'token.cleanup.expired',
|
||||
preventOverlap: true
|
||||
})
|
||||
async enqueueExpiredCleanup(): Promise<{ requestedAt: string }> {
|
||||
const requestedAt = new Date().toISOString();
|
||||
this.logger.log(`Queue expired cleanup at ${requestedAt}`);
|
||||
return { requestedAt };
|
||||
}
|
||||
|
||||
@ScheduledJob({
|
||||
name: 'token.cleanup.inactive.scheduler',
|
||||
cron: CronExpression.EVERY_6_HOURS,
|
||||
queueName: 'token-maintenance',
|
||||
queueJobName: 'token.cleanup.inactive'
|
||||
})
|
||||
async enqueueInactiveCleanup(): Promise<{ requestedAt: string }> {
|
||||
return { requestedAt: new Date().toISOString() };
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
## 3) Queue worker (job consumer)
|
||||
|
||||
Use `QueueWorkerHost` and map queue job names to handlers with `@QueueJobHandler(...)`.
|
||||
|
||||
```ts
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { Job } from 'bullmq';
|
||||
import { QueueJobHandler, QueueWorker, QueueWorkerHost } from '@gauzy/scheduler';
|
||||
|
||||
@Injectable()
|
||||
@QueueWorker('token-maintenance')
|
||||
export class TokenCleanupWorker extends QueueWorkerHost {
|
||||
private readonly logger = new Logger(TokenCleanupWorker.name);
|
||||
|
||||
@QueueJobHandler('token.cleanup.expired')
|
||||
async handleExpired(job: Job<{ requestedAt: string }>): Promise<void> {
|
||||
this.logger.log(`Process expired cleanup requested at ${job.data.requestedAt}`);
|
||||
// Execute business logic (CQRS command, service call, etc.)
|
||||
}
|
||||
|
||||
@QueueJobHandler('token.cleanup.inactive')
|
||||
async handleInactive(job: Job<{ requestedAt: string }>): Promise<void> {
|
||||
this.logger.log(`Process inactive cleanup requested at ${job.data.requestedAt}`);
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
## 4) Feature registration
|
||||
|
||||
Register queues and providers for the feature.
|
||||
|
||||
```ts
|
||||
import { Module } from '@nestjs/common';
|
||||
import { SchedulerModule } from '@gauzy/scheduler';
|
||||
import { TokenCleanupScheduler } from './token-cleanup.scheduler';
|
||||
import { TokenCleanupWorker } from './token-cleanup.worker';
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
SchedulerModule.forFeature({
|
||||
queues: ['token-maintenance'],
|
||||
jobProviders: [TokenCleanupScheduler, TokenCleanupWorker]
|
||||
})
|
||||
]
|
||||
})
|
||||
export class TokenCleanupJobsModule {}
|
||||
```
|
||||
|
||||
## 5) Runtime control (optional)
|
||||
|
||||
```ts
|
||||
import { Controller, Get, Param, Post } from '@nestjs/common';
|
||||
import { SchedulerService } from '@gauzy/scheduler';
|
||||
|
||||
@Controller('scheduler')
|
||||
export class SchedulerController {
|
||||
constructor(private readonly scheduler: SchedulerService) {}
|
||||
|
||||
@Get('jobs')
|
||||
listJobs() {
|
||||
return this.scheduler.listJobs();
|
||||
}
|
||||
|
||||
@Post('jobs/:id/trigger')
|
||||
async trigger(@Param('id') id: string) {
|
||||
await this.scheduler.triggerNow(id);
|
||||
return { ok: true };
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
## ScheduledJob options
|
||||
|
||||
- `name`: custom job id (default: `ProviderName.methodName`)
|
||||
- `cron`: cron expression schedule
|
||||
- `intervalMs`: interval schedule (mutually exclusive with `cron`)
|
||||
- `runOnStart`: execute once at bootstrap
|
||||
- `enabled`: enable/disable a job
|
||||
- `preventOverlap`: skip new run if previous run is still active
|
||||
- `retries`, `retryDelayMs`: retry behavior
|
||||
- `timeoutMs`: max execution time per attempt
|
||||
- `maxRandomDelayMs`: jitter before each scheduled execution
|
||||
- `queueName`: queue target (if set, execution result is enqueued)
|
||||
- `queueJobName`: BullMQ job name (default: scheduler job id)
|
||||
- `queueJobOptions`: BullMQ `JobsOptions`
|
||||
|
||||
## Common pitfalls
|
||||
|
||||
- Queue job without queue registration:
|
||||
- Ensure the queue exists in `SchedulerModule.forFeature({ queues: [...] })`.
|
||||
- Queueing disabled:
|
||||
- If `enableQueueing` is `false`, jobs with `queueName` will fail by design.
|
||||
- Duplicate job names:
|
||||
- `name` must be unique across all discovered scheduled jobs.
|
||||
|
||||
## Build
|
||||
|
||||
```bash
|
||||
nx build scheduler
|
||||
```
|
||||
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"name": "@gauzy/scheduler",
|
||||
"version": "0.0.1",
|
||||
"type": "commonjs",
|
||||
"main": "./src/index.js",
|
||||
"types": "./src/index.d.ts",
|
||||
"dependencies": {
|
||||
"@nestjs/bullmq": "^11.0.4",
|
||||
"@nestjs/common": "^11.1.14",
|
||||
"@nestjs/core": "^11.1.14",
|
||||
"@nestjs/schedule": "^6.1.1",
|
||||
"bullmq": "^5.70.1",
|
||||
"cron": "^4.3.3",
|
||||
"tslib": "^2.3.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
{
|
||||
"name": "scheduler",
|
||||
"$schema": "../../node_modules/nx/schemas/project-schema.json",
|
||||
"sourceRoot": "packages/scheduler/src",
|
||||
"projectType": "library",
|
||||
"release": {
|
||||
"version": {
|
||||
"manifestRootsToUpdate": ["dist/{projectRoot}"],
|
||||
"currentVersionResolver": "git-tag",
|
||||
"fallbackCurrentVersionResolver": "disk"
|
||||
}
|
||||
},
|
||||
"tags": [],
|
||||
"targets": {
|
||||
"build": {
|
||||
"executor": "@nx/js:tsc",
|
||||
"outputs": ["{options.outputPath}"],
|
||||
"options": {
|
||||
"outputPath": "dist/packages/scheduler",
|
||||
"tsConfig": "packages/scheduler/tsconfig.lib.json",
|
||||
"packageJson": "packages/scheduler/package.json",
|
||||
"main": "packages/scheduler/src/index.ts",
|
||||
"assets": ["packages/scheduler/*.md"]
|
||||
}
|
||||
},
|
||||
"nx-release-publish": {
|
||||
"options": {
|
||||
"packageRoot": "dist/{projectRoot}"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
export * from './lib/scheduler.module';
|
||||
export * from './lib/decorators/queue-job-handler.decorator';
|
||||
export * from './lib/decorators/queue-worker.decorator';
|
||||
export * from './lib/decorators/scheduled-job.decorator';
|
||||
export * from './lib/interfaces/discovered-scheduled-job.interface';
|
||||
export * from './lib/interfaces/scheduler-feature-options.interface';
|
||||
export * from './lib/interfaces/scheduled-job-options.interface';
|
||||
export * from './lib/interfaces/scheduler-job-descriptor.interface';
|
||||
export * from './lib/interfaces/scheduler-module-options.interface';
|
||||
export * from './lib/interfaces/scheduler-queue-job.interface';
|
||||
export * from './lib/hosts/queue-worker.host';
|
||||
export * from './lib/services/scheduler-queue.service';
|
||||
export * from './lib/services/scheduler.service';
|
||||
@@ -0,0 +1,3 @@
|
||||
export const SCHEDULED_JOB_METADATA = 'gauzy:scheduler:job';
|
||||
export const SCHEDULER_MODULE_OPTIONS = Symbol('GAUZY_SCHEDULER_MODULE_OPTIONS');
|
||||
export const QUEUE_JOB_HANDLER_METADATA = 'gauzy:scheduler:queue-job-handler';
|
||||
@@ -0,0 +1,11 @@
|
||||
import { SetMetadata } from '@nestjs/common';
|
||||
import { QUEUE_JOB_HANDLER_METADATA } from '../constants/scheduler.constants';
|
||||
|
||||
export function QueueJobHandler(jobName: string): MethodDecorator {
|
||||
const normalizedJobName = jobName?.trim();
|
||||
if (!normalizedJobName) {
|
||||
throw new Error('Queue job handler name cannot be empty.');
|
||||
}
|
||||
|
||||
return SetMetadata(QUEUE_JOB_HANDLER_METADATA, normalizedJobName);
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
import { Processor } from '@nestjs/bullmq';
|
||||
|
||||
export const QueueWorker: typeof Processor = Processor;
|
||||
@@ -0,0 +1,7 @@
|
||||
import { SetMetadata } from '@nestjs/common';
|
||||
import { SCHEDULED_JOB_METADATA } from '../constants/scheduler.constants';
|
||||
import { ScheduledJobOptions } from '../interfaces/scheduled-job-options.interface';
|
||||
|
||||
export function ScheduledJob(options: ScheduledJobOptions = {}): MethodDecorator {
|
||||
return SetMetadata(SCHEDULED_JOB_METADATA, options);
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
import { WorkerHost } from '@nestjs/bullmq';
|
||||
import { Job } from 'bullmq';
|
||||
import { QUEUE_JOB_HANDLER_METADATA } from '../constants/scheduler.constants';
|
||||
|
||||
type QueueJobHandler = (job: Job, token?: string) => Promise<unknown> | unknown;
|
||||
|
||||
export abstract class QueueWorkerHost extends WorkerHost {
|
||||
private handlerMap?: Map<string, QueueJobHandler>;
|
||||
|
||||
async process(job: Job, token?: string): Promise<unknown> {
|
||||
const handler = this.getHandlers().get(job.name);
|
||||
|
||||
if (!handler) {
|
||||
const available = Array.from(this.getHandlers().keys());
|
||||
throw new Error(
|
||||
`No handler found for queue job "${job.name}" in ${this.constructor.name}. Available handlers: ${available.join(', ')}`
|
||||
);
|
||||
}
|
||||
|
||||
return handler(job, token);
|
||||
}
|
||||
|
||||
private getHandlers(): Map<string, QueueJobHandler> {
|
||||
if (this.handlerMap) {
|
||||
return this.handlerMap;
|
||||
}
|
||||
|
||||
const handlers = new Map<string, QueueJobHandler>();
|
||||
const prototype = Object.getPrototypeOf(this) as Record<string, unknown>;
|
||||
const methodNames = Object.getOwnPropertyNames(prototype);
|
||||
|
||||
for (const methodName of methodNames) {
|
||||
if (methodName === 'constructor') {
|
||||
continue;
|
||||
}
|
||||
|
||||
const methodRef = prototype[methodName];
|
||||
if (typeof methodRef !== 'function') {
|
||||
continue;
|
||||
}
|
||||
|
||||
const handlerName = Reflect.getMetadata(QUEUE_JOB_HANDLER_METADATA, methodRef) as string | undefined;
|
||||
if (!handlerName) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (handlers.has(handlerName)) {
|
||||
throw new Error(`Duplicate queue job handler "${handlerName}" in ${this.constructor.name}.`);
|
||||
}
|
||||
|
||||
const target = this as unknown as Record<string, QueueJobHandler>;
|
||||
handlers.set(handlerName, (job: Job, token?: string) => target[methodName](job, token));
|
||||
}
|
||||
|
||||
this.handlerMap = handlers;
|
||||
return handlers;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
import { ScheduledJobOptions } from './scheduled-job-options.interface';
|
||||
import { JobsOptions } from 'bullmq';
|
||||
|
||||
export interface ResolvedScheduledJobOptions {
|
||||
enabled: boolean;
|
||||
description?: string;
|
||||
cron?: string;
|
||||
intervalMs?: number;
|
||||
runOnStart: boolean;
|
||||
preventOverlap: boolean;
|
||||
retries: number;
|
||||
retryDelayMs: number;
|
||||
timeoutMs?: number;
|
||||
maxRandomDelayMs: number;
|
||||
queueName?: string;
|
||||
queueJobName?: string;
|
||||
queueJobOptions?: JobsOptions;
|
||||
}
|
||||
|
||||
export interface DiscoveredScheduledJob {
|
||||
id: string;
|
||||
providerName: string;
|
||||
methodName: string;
|
||||
options: ResolvedScheduledJobOptions;
|
||||
handler: () => Promise<unknown>;
|
||||
}
|
||||
|
||||
export interface RegisterScheduledJobInput {
|
||||
providerName: string;
|
||||
methodName: string;
|
||||
metadata: ScheduledJobOptions;
|
||||
handler: () => Promise<unknown>;
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
import { JobsOptions } from 'bullmq';
|
||||
|
||||
export interface ScheduledJobDefaults {
|
||||
enabled?: boolean;
|
||||
preventOverlap?: boolean;
|
||||
retries?: number;
|
||||
retryDelayMs?: number;
|
||||
timeoutMs?: number;
|
||||
maxRandomDelayMs?: number;
|
||||
}
|
||||
|
||||
export interface ScheduledJobOptions extends ScheduledJobDefaults {
|
||||
name?: string;
|
||||
description?: string;
|
||||
cron?: string;
|
||||
intervalMs?: number;
|
||||
runOnStart?: boolean;
|
||||
queueName?: string;
|
||||
queueJobName?: string;
|
||||
queueJobOptions?: JobsOptions;
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
import { Type } from '@nestjs/common';
|
||||
import { RegisterQueueOptions } from '@nestjs/bullmq';
|
||||
|
||||
export type SchedulerQueueRegistration = string | RegisterQueueOptions;
|
||||
|
||||
export interface SchedulerFeatureOptions {
|
||||
jobProviders?: Type<unknown>[];
|
||||
queues?: SchedulerQueueRegistration[];
|
||||
}
|
||||
@@ -0,0 +1,18 @@
|
||||
export type SchedulerJobScheduleType = 'cron' | 'interval' | 'manual';
|
||||
export type SchedulerJobExecutionTarget = 'inline' | 'queue';
|
||||
|
||||
export interface SchedulerJobDescriptor {
|
||||
id: string;
|
||||
providerName: string;
|
||||
methodName: string;
|
||||
description?: string;
|
||||
enabled: boolean;
|
||||
runOnStart: boolean;
|
||||
scheduleType: SchedulerJobScheduleType;
|
||||
executionTarget: SchedulerJobExecutionTarget;
|
||||
cron?: string;
|
||||
intervalMs?: number;
|
||||
queueName?: string;
|
||||
queueJobName?: string;
|
||||
running: boolean;
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
import { ConnectionOptions } from 'bullmq';
|
||||
import { RegisterQueueOptions } from '@nestjs/bullmq';
|
||||
import { ScheduledJobDefaults } from './scheduled-job-options.interface';
|
||||
import { SchedulerQueueRegistration } from './scheduler-feature-options.interface';
|
||||
|
||||
export interface SchedulerModuleOptions {
|
||||
enabled?: boolean;
|
||||
defaultTimezone?: string;
|
||||
logRegisteredJobs?: boolean;
|
||||
defaultJobOptions?: ScheduledJobDefaults;
|
||||
enableQueueing?: boolean;
|
||||
queueConnection?: ConnectionOptions;
|
||||
queues?: SchedulerQueueRegistration[];
|
||||
defaultQueueName?: string;
|
||||
}
|
||||
|
||||
export interface ResolvedScheduledJobDefaults {
|
||||
enabled: boolean;
|
||||
preventOverlap: boolean;
|
||||
retries: number;
|
||||
retryDelayMs: number;
|
||||
timeoutMs?: number;
|
||||
maxRandomDelayMs: number;
|
||||
}
|
||||
|
||||
export interface ResolvedSchedulerModuleOptions {
|
||||
enabled: boolean;
|
||||
defaultTimezone?: string;
|
||||
logRegisteredJobs: boolean;
|
||||
defaultJobOptions: ResolvedScheduledJobDefaults;
|
||||
enableQueueing: boolean;
|
||||
queueConnection: ConnectionOptions;
|
||||
queues: RegisterQueueOptions[];
|
||||
defaultQueueName: string;
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
import { JobsOptions } from 'bullmq';
|
||||
|
||||
export interface SchedulerQueueJobInput<TData = unknown> {
|
||||
queueName: string;
|
||||
jobName: string;
|
||||
data?: TData;
|
||||
options?: JobsOptions;
|
||||
}
|
||||
@@ -0,0 +1,69 @@
|
||||
import { BullModule, BullRootModuleOptions } from '@nestjs/bullmq';
|
||||
import { DynamicModule, Module, Type } from '@nestjs/common';
|
||||
import { DiscoveryModule } from '@nestjs/core';
|
||||
import { ScheduleModule } from '@nestjs/schedule';
|
||||
import { SCHEDULER_MODULE_OPTIONS } from './constants/scheduler.constants';
|
||||
import { SchedulerFeatureOptions } from './interfaces/scheduler-feature-options.interface';
|
||||
import { SchedulerModuleOptions } from './interfaces/scheduler-module-options.interface';
|
||||
import { ScheduledJobMetadataAccessor } from './services/scheduled-job-metadata.accessor';
|
||||
import { SchedulerDiscoveryService } from './services/scheduler-discovery.service';
|
||||
import { SchedulerJobRegistryService } from './services/scheduler-job-registry.service';
|
||||
import { SchedulerJobRunnerService } from './services/scheduler-job-runner.service';
|
||||
import { SchedulerQueueService } from './services/scheduler-queue.service';
|
||||
import { SchedulerService } from './services/scheduler.service';
|
||||
import { normalizeQueueRegistrations } from './utils/normalize-queue-registrations';
|
||||
import { normalizeSchedulerModuleOptions } from './utils/normalize-scheduler-options';
|
||||
|
||||
@Module({})
|
||||
export class SchedulerModule {
|
||||
static forRoot(options: SchedulerModuleOptions = {}): DynamicModule {
|
||||
const normalizedOptions = normalizeSchedulerModuleOptions(options);
|
||||
const imports: DynamicModule['imports'] = [DiscoveryModule, ScheduleModule.forRoot()];
|
||||
|
||||
if (normalizedOptions.enableQueueing) {
|
||||
const rootOptions: BullRootModuleOptions = {
|
||||
connection: normalizedOptions.queueConnection
|
||||
};
|
||||
imports.push(BullModule.forRoot(rootOptions));
|
||||
|
||||
if (normalizedOptions.queues.length > 0) {
|
||||
imports.push(BullModule.registerQueue(...normalizedOptions.queues));
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
module: SchedulerModule,
|
||||
global: true,
|
||||
imports,
|
||||
providers: [
|
||||
{
|
||||
provide: SCHEDULER_MODULE_OPTIONS,
|
||||
useValue: normalizedOptions
|
||||
},
|
||||
ScheduledJobMetadataAccessor,
|
||||
SchedulerJobRegistryService,
|
||||
SchedulerJobRunnerService,
|
||||
SchedulerQueueService,
|
||||
SchedulerDiscoveryService,
|
||||
SchedulerService
|
||||
],
|
||||
exports: [SchedulerService, SchedulerQueueService]
|
||||
};
|
||||
}
|
||||
|
||||
static forFeature(optionsOrProviders: SchedulerFeatureOptions | Type<unknown>[] = []): DynamicModule {
|
||||
const options = Array.isArray(optionsOrProviders)
|
||||
? { jobProviders: optionsOrProviders }
|
||||
: optionsOrProviders;
|
||||
const queueOptions = normalizeQueueRegistrations(options.queues ?? []);
|
||||
const imports = queueOptions.length > 0 ? [BullModule.registerQueue(...queueOptions)] : [];
|
||||
const jobProviders = options.jobProviders ?? [];
|
||||
|
||||
return {
|
||||
module: SchedulerModule,
|
||||
imports,
|
||||
providers: [...jobProviders],
|
||||
exports: [...jobProviders]
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { Reflector } from '@nestjs/core';
|
||||
import { SCHEDULED_JOB_METADATA } from '../constants/scheduler.constants';
|
||||
import { ScheduledJobOptions } from '../interfaces/scheduled-job-options.interface';
|
||||
|
||||
@Injectable()
|
||||
export class ScheduledJobMetadataAccessor {
|
||||
constructor(private readonly reflector: Reflector) {}
|
||||
|
||||
get(target: (...args: unknown[]) => unknown): ScheduledJobOptions | undefined {
|
||||
return this.reflector.get<ScheduledJobOptions>(SCHEDULED_JOB_METADATA, target);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,216 @@
|
||||
import {
|
||||
Inject,
|
||||
Injectable,
|
||||
Logger,
|
||||
OnApplicationBootstrap,
|
||||
OnApplicationShutdown,
|
||||
OnModuleInit
|
||||
} from '@nestjs/common';
|
||||
import { DiscoveryService, MetadataScanner } from '@nestjs/core';
|
||||
import { InstanceWrapper } from '@nestjs/core/injector/instance-wrapper';
|
||||
import { SchedulerRegistry } from '@nestjs/schedule';
|
||||
import { CronJob } from 'cron';
|
||||
import * as crypto from 'node:crypto';
|
||||
import { SCHEDULER_MODULE_OPTIONS } from '../constants/scheduler.constants';
|
||||
import { DiscoveredScheduledJob } from '../interfaces/discovered-scheduled-job.interface';
|
||||
import { ResolvedSchedulerModuleOptions } from '../interfaces/scheduler-module-options.interface';
|
||||
import { ScheduledJobMetadataAccessor } from './scheduled-job-metadata.accessor';
|
||||
import { SchedulerJobRegistryService } from './scheduler-job-registry.service';
|
||||
import { SchedulerJobRunnerService } from './scheduler-job-runner.service';
|
||||
|
||||
type RegisteredScheduleKind = 'cron' | 'interval';
|
||||
|
||||
@Injectable()
|
||||
export class SchedulerDiscoveryService implements OnModuleInit, OnApplicationBootstrap, OnApplicationShutdown {
|
||||
private readonly logger = new Logger(SchedulerDiscoveryService.name);
|
||||
private readonly registeredScheduleKinds = new Map<string, RegisteredScheduleKind>();
|
||||
|
||||
constructor(
|
||||
private readonly discoveryService: DiscoveryService,
|
||||
private readonly metadataScanner: MetadataScanner,
|
||||
private readonly metadataAccessor: ScheduledJobMetadataAccessor,
|
||||
private readonly jobRegistry: SchedulerJobRegistryService,
|
||||
private readonly jobRunner: SchedulerJobRunnerService,
|
||||
private readonly schedulerRegistry: SchedulerRegistry,
|
||||
@Inject(SCHEDULER_MODULE_OPTIONS)
|
||||
private readonly moduleOptions: ResolvedSchedulerModuleOptions
|
||||
) {}
|
||||
|
||||
onModuleInit(): void {
|
||||
this.discoverJobs();
|
||||
this.registerSchedules();
|
||||
}
|
||||
|
||||
async onApplicationBootstrap(): Promise<void> {
|
||||
const startupJobs = this.jobRegistry.getAll().filter((job) => job.options.enabled && job.options.runOnStart);
|
||||
|
||||
for (const job of startupJobs) {
|
||||
await this.executeJob(job);
|
||||
}
|
||||
}
|
||||
|
||||
onApplicationShutdown(): void {
|
||||
this.unregisterSchedules();
|
||||
}
|
||||
|
||||
private discoverJobs(): void {
|
||||
const wrappers: Array<InstanceWrapper> = [
|
||||
...this.discoveryService.getProviders(),
|
||||
...this.discoveryService.getControllers()
|
||||
];
|
||||
|
||||
for (const wrapper of wrappers) {
|
||||
if (typeof wrapper.isDependencyTreeStatic === 'function' && !wrapper.isDependencyTreeStatic()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const instance = wrapper.instance as Record<string, unknown> | undefined;
|
||||
if (!instance) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const prototype = Object.getPrototypeOf(instance) as object | undefined;
|
||||
if (!prototype) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const providerName = this.resolveProviderName(wrapper, instance);
|
||||
const methodNames = this.metadataScanner.getAllMethodNames(prototype);
|
||||
|
||||
for (const methodName of methodNames) {
|
||||
const methodCandidate = instance[methodName];
|
||||
if (typeof methodCandidate !== 'function') {
|
||||
continue;
|
||||
}
|
||||
|
||||
const metadata = this.metadataAccessor.get(methodCandidate as (...args: unknown[]) => unknown);
|
||||
if (!metadata) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const job = this.jobRegistry.register({
|
||||
providerName,
|
||||
methodName,
|
||||
metadata,
|
||||
handler: async () => {
|
||||
return await Promise.resolve(methodCandidate.call(instance));
|
||||
}
|
||||
});
|
||||
|
||||
if (this.moduleOptions.logRegisteredJobs) {
|
||||
this.logger.log(`Registered scheduled job "${job.id}" (${providerName}.${methodName}).`);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private registerSchedules(): void {
|
||||
const jobs = this.jobRegistry.getAll();
|
||||
for (const job of jobs) {
|
||||
if (!job.options.enabled || !this.moduleOptions.enabled) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (job.options.cron) {
|
||||
this.registerCronJob(job);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (job.options.intervalMs !== undefined) {
|
||||
this.registerIntervalJob(job);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (this.moduleOptions.logRegisteredJobs) {
|
||||
this.logger.log(`Job "${job.id}" is registered for manual/startup execution.`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private registerCronJob(job: DiscoveredScheduledJob): void {
|
||||
const cronJob = new CronJob(
|
||||
job.options.cron as string,
|
||||
() => {
|
||||
void this.executeJobWithJitter(job);
|
||||
},
|
||||
null,
|
||||
false,
|
||||
this.moduleOptions.defaultTimezone
|
||||
);
|
||||
|
||||
this.schedulerRegistry.addCronJob(job.id, cronJob);
|
||||
this.registeredScheduleKinds.set(job.id, 'cron');
|
||||
cronJob.start();
|
||||
}
|
||||
|
||||
private registerIntervalJob(job: DiscoveredScheduledJob): void {
|
||||
const interval = setInterval(() => {
|
||||
void this.executeJobWithJitter(job);
|
||||
}, job.options.intervalMs as number);
|
||||
|
||||
this.schedulerRegistry.addInterval(job.id, interval);
|
||||
this.registeredScheduleKinds.set(job.id, 'interval');
|
||||
}
|
||||
|
||||
private async executeJobWithJitter(job: DiscoveredScheduledJob): Promise<void> {
|
||||
if (job.options.maxRandomDelayMs > 0) {
|
||||
const delayMs = randomDelay(job.options.maxRandomDelayMs);
|
||||
await delay(delayMs);
|
||||
}
|
||||
|
||||
await this.executeJob(job);
|
||||
}
|
||||
|
||||
private async executeJob(job: DiscoveredScheduledJob): Promise<void> {
|
||||
try {
|
||||
await this.jobRunner.execute(job);
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? (error.stack ?? error.message) : String(error);
|
||||
this.logger.error(`Job "${job.id}" execution failed.`, message);
|
||||
}
|
||||
}
|
||||
|
||||
private unregisterSchedules(): void {
|
||||
for (const [jobId, scheduleKind] of this.registeredScheduleKinds.entries()) {
|
||||
try {
|
||||
if (scheduleKind === 'cron') {
|
||||
this.schedulerRegistry.deleteCronJob(jobId);
|
||||
continue;
|
||||
}
|
||||
|
||||
this.schedulerRegistry.deleteInterval(jobId);
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
this.logger.warn(`Failed to remove "${jobId}" from scheduler registry: ${message}`);
|
||||
}
|
||||
}
|
||||
|
||||
this.registeredScheduleKinds.clear();
|
||||
}
|
||||
|
||||
private resolveProviderName(wrapper: InstanceWrapper, instance: Record<string, unknown>): string {
|
||||
if (typeof wrapper.name === 'string' && wrapper.name.trim().length > 0) {
|
||||
return wrapper.name;
|
||||
}
|
||||
|
||||
const constructorName = (instance.constructor as { name?: string }).name;
|
||||
if (constructorName && constructorName.trim().length > 0) {
|
||||
return constructorName;
|
||||
}
|
||||
|
||||
return 'AnonymousProvider';
|
||||
}
|
||||
}
|
||||
|
||||
function delay(ms: number): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
setTimeout(resolve, ms);
|
||||
});
|
||||
}
|
||||
|
||||
function randomDelay(maxMs: number): number {
|
||||
if (maxMs <= 0) {
|
||||
return 0;
|
||||
}
|
||||
return Math.floor(crypto.randomInt(0, maxMs + 1));
|
||||
}
|
||||
@@ -0,0 +1,125 @@
|
||||
import { Inject, Injectable } from '@nestjs/common';
|
||||
import { SCHEDULER_MODULE_OPTIONS } from '../constants/scheduler.constants';
|
||||
import {
|
||||
DiscoveredScheduledJob,
|
||||
RegisterScheduledJobInput,
|
||||
ResolvedScheduledJobOptions
|
||||
} from '../interfaces/discovered-scheduled-job.interface';
|
||||
import { ResolvedSchedulerModuleOptions } from '../interfaces/scheduler-module-options.interface';
|
||||
import { ScheduledJobOptions } from '../interfaces/scheduled-job-options.interface';
|
||||
|
||||
@Injectable()
|
||||
export class SchedulerJobRegistryService {
|
||||
private readonly jobs = new Map<string, DiscoveredScheduledJob>();
|
||||
|
||||
constructor(
|
||||
@Inject(SCHEDULER_MODULE_OPTIONS)
|
||||
private readonly moduleOptions: ResolvedSchedulerModuleOptions
|
||||
) {}
|
||||
|
||||
register(input: RegisterScheduledJobInput): DiscoveredScheduledJob {
|
||||
const id = this.resolveJobId(input.providerName, input.methodName, input.metadata.name);
|
||||
|
||||
if (this.jobs.has(id)) {
|
||||
throw new Error(`Duplicate scheduled job id "${id}".`);
|
||||
}
|
||||
|
||||
const options = this.resolveJobOptions(input.metadata, id);
|
||||
const job: DiscoveredScheduledJob = {
|
||||
id,
|
||||
providerName: input.providerName,
|
||||
methodName: input.methodName,
|
||||
options,
|
||||
handler: input.handler
|
||||
};
|
||||
|
||||
this.jobs.set(id, job);
|
||||
return job;
|
||||
}
|
||||
|
||||
getAll(): DiscoveredScheduledJob[] {
|
||||
return Array.from(this.jobs.values());
|
||||
}
|
||||
|
||||
getById(id: string): DiscoveredScheduledJob | undefined {
|
||||
return this.jobs.get(id);
|
||||
}
|
||||
|
||||
private resolveJobId(providerName: string, methodName: string, customName?: string): string {
|
||||
const defaultId = `${providerName}.${methodName}`;
|
||||
const name = customName?.trim();
|
||||
return name && name.length > 0 ? name : defaultId;
|
||||
}
|
||||
|
||||
private resolveJobOptions(metadata: ScheduledJobOptions, jobId: string): ResolvedScheduledJobOptions {
|
||||
const cron = metadata.cron?.trim();
|
||||
const intervalMs = metadata.intervalMs;
|
||||
const queueNameInput = metadata.queueName?.trim();
|
||||
const queueJobName = metadata.queueJobName?.trim();
|
||||
const queueName =
|
||||
queueNameInput && queueNameInput.length > 0
|
||||
? queueNameInput
|
||||
: queueJobName || metadata.queueJobOptions
|
||||
? this.moduleOptions.defaultQueueName
|
||||
: undefined;
|
||||
const enabled = metadata.enabled ?? this.moduleOptions.defaultJobOptions.enabled;
|
||||
|
||||
if (cron && intervalMs !== undefined) {
|
||||
throw new Error(`Job "${jobId}" cannot define both "cron" and "intervalMs".`);
|
||||
}
|
||||
|
||||
if (intervalMs !== undefined && (!Number.isFinite(intervalMs) || intervalMs <= 0)) {
|
||||
throw new Error(`Job "${jobId}" has invalid "intervalMs" value "${intervalMs}".`);
|
||||
}
|
||||
|
||||
if (enabled && queueName && !this.moduleOptions.enableQueueing) {
|
||||
throw new Error(`Job "${jobId}" targets queue "${queueName}" but queueing is disabled.`);
|
||||
}
|
||||
|
||||
const retries = this.toNonNegativeInteger(metadata.retries ?? this.moduleOptions.defaultJobOptions.retries, 'retries', jobId);
|
||||
const retryDelayMs = this.toNonNegativeInteger(
|
||||
metadata.retryDelayMs ?? this.moduleOptions.defaultJobOptions.retryDelayMs,
|
||||
'retryDelayMs',
|
||||
jobId
|
||||
);
|
||||
const maxRandomDelayMs = this.toNonNegativeInteger(
|
||||
metadata.maxRandomDelayMs ?? this.moduleOptions.defaultJobOptions.maxRandomDelayMs,
|
||||
'maxRandomDelayMs',
|
||||
jobId
|
||||
);
|
||||
const timeoutMs =
|
||||
metadata.timeoutMs !== undefined
|
||||
? this.toPositiveInteger(metadata.timeoutMs, 'timeoutMs', jobId)
|
||||
: this.moduleOptions.defaultJobOptions.timeoutMs;
|
||||
|
||||
return {
|
||||
enabled,
|
||||
description: metadata.description,
|
||||
cron: cron && cron.length > 0 ? cron : undefined,
|
||||
intervalMs,
|
||||
runOnStart: metadata.runOnStart ?? false,
|
||||
preventOverlap: metadata.preventOverlap ?? this.moduleOptions.defaultJobOptions.preventOverlap,
|
||||
retries,
|
||||
retryDelayMs,
|
||||
timeoutMs,
|
||||
maxRandomDelayMs,
|
||||
queueName,
|
||||
queueJobName: queueJobName && queueJobName.length > 0 ? queueJobName : undefined,
|
||||
queueJobOptions: metadata.queueJobOptions
|
||||
};
|
||||
}
|
||||
|
||||
private toNonNegativeInteger(value: number, fieldName: string, jobId: string): number {
|
||||
if (!Number.isFinite(value) || value < 0) {
|
||||
throw new Error(`Job "${jobId}" has invalid "${fieldName}" value "${value}".`);
|
||||
}
|
||||
return Math.floor(value);
|
||||
}
|
||||
|
||||
private toPositiveInteger(value: number, fieldName: string, jobId: string): number {
|
||||
if (!Number.isFinite(value) || value <= 0) {
|
||||
throw new Error(`Job "${jobId}" has invalid "${fieldName}" value "${value}".`);
|
||||
}
|
||||
return Math.floor(value);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
import { Inject, Injectable, Logger } from '@nestjs/common';
|
||||
import { SCHEDULER_MODULE_OPTIONS } from '../constants/scheduler.constants';
|
||||
import { DiscoveredScheduledJob } from '../interfaces/discovered-scheduled-job.interface';
|
||||
import { ResolvedSchedulerModuleOptions } from '../interfaces/scheduler-module-options.interface';
|
||||
import { SchedulerQueueService } from './scheduler-queue.service';
|
||||
|
||||
@Injectable()
|
||||
export class SchedulerJobRunnerService {
|
||||
private readonly logger = new Logger(SchedulerJobRunnerService.name);
|
||||
private readonly runningJobs = new Set<string>();
|
||||
|
||||
constructor(
|
||||
@Inject(SCHEDULER_MODULE_OPTIONS)
|
||||
private readonly moduleOptions: ResolvedSchedulerModuleOptions,
|
||||
private readonly queueService: SchedulerQueueService
|
||||
) {}
|
||||
|
||||
isRunning(jobId: string): boolean {
|
||||
return this.runningJobs.has(jobId);
|
||||
}
|
||||
|
||||
async execute(job: DiscoveredScheduledJob): Promise<void> {
|
||||
if (!this.moduleOptions.enabled || !job.options.enabled) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (job.options.preventOverlap && this.runningJobs.has(job.id)) {
|
||||
this.logger.warn(`Skipping "${job.id}" because the previous run is still in progress.`);
|
||||
return;
|
||||
}
|
||||
|
||||
const startedAt = Date.now();
|
||||
this.runningJobs.add(job.id);
|
||||
|
||||
try {
|
||||
await this.executeWithRetry(job);
|
||||
this.logger.debug(`Finished "${job.id}" in ${Date.now() - startedAt}ms`);
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.stack ?? error.message : String(error);
|
||||
this.logger.error(`Scheduled job "${job.id}" failed.`, message);
|
||||
throw error;
|
||||
} finally {
|
||||
this.runningJobs.delete(job.id);
|
||||
}
|
||||
}
|
||||
|
||||
private async executeWithRetry(job: DiscoveredScheduledJob): Promise<void> {
|
||||
const attempts = job.options.retries + 1;
|
||||
let attempt = 1;
|
||||
let lastError: unknown;
|
||||
|
||||
while (attempt <= attempts) {
|
||||
try {
|
||||
await this.executeSingleAttempt(job);
|
||||
return;
|
||||
} catch (error) {
|
||||
lastError = error;
|
||||
const hasNextAttempt = attempt < attempts;
|
||||
if (!hasNextAttempt) {
|
||||
break;
|
||||
}
|
||||
|
||||
this.logger.warn(
|
||||
`"${job.id}" failed on attempt ${attempt}/${attempts}. Retrying in ${job.options.retryDelayMs}ms.`
|
||||
);
|
||||
await sleep(job.options.retryDelayMs);
|
||||
attempt += 1;
|
||||
}
|
||||
}
|
||||
|
||||
throw lastError;
|
||||
}
|
||||
|
||||
private async executeSingleAttempt(job: DiscoveredScheduledJob): Promise<void> {
|
||||
const execution = this.executeJobHandler(job);
|
||||
if (job.options.timeoutMs === undefined) {
|
||||
await execution;
|
||||
return;
|
||||
}
|
||||
|
||||
await Promise.race([execution, timeout(job.options.timeoutMs, job.id)]);
|
||||
}
|
||||
|
||||
private async executeJobHandler(job: DiscoveredScheduledJob): Promise<void> {
|
||||
const data = await Promise.resolve(job.handler());
|
||||
|
||||
if (!job.options.queueName) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.queueService.enqueue({
|
||||
queueName: job.options.queueName,
|
||||
jobName: job.options.queueJobName ?? job.id,
|
||||
data,
|
||||
options: job.options.queueJobOptions
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function sleep(ms: number): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
setTimeout(resolve, ms);
|
||||
});
|
||||
}
|
||||
|
||||
function timeout(ms: number, jobId: string): Promise<never> {
|
||||
return new Promise((_, reject) => {
|
||||
setTimeout(() => {
|
||||
reject(new Error(`Scheduled job "${jobId}" timed out after ${ms}ms.`));
|
||||
}, ms);
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
import { getQueueToken } from '@nestjs/bullmq';
|
||||
import { Inject, Injectable } from '@nestjs/common';
|
||||
import { ModuleRef } from '@nestjs/core';
|
||||
import { Queue } from 'bullmq';
|
||||
import { SCHEDULER_MODULE_OPTIONS } from '../constants/scheduler.constants';
|
||||
import { SchedulerQueueJobInput } from '../interfaces/scheduler-queue-job.interface';
|
||||
import { ResolvedSchedulerModuleOptions } from '../interfaces/scheduler-module-options.interface';
|
||||
|
||||
@Injectable()
|
||||
export class SchedulerQueueService {
|
||||
constructor(
|
||||
private readonly moduleRef: ModuleRef,
|
||||
@Inject(SCHEDULER_MODULE_OPTIONS)
|
||||
private readonly moduleOptions: ResolvedSchedulerModuleOptions
|
||||
) {}
|
||||
|
||||
async enqueue<TData = unknown>(input: SchedulerQueueJobInput<TData>): Promise<void> {
|
||||
if (!this.moduleOptions.enableQueueing) {
|
||||
throw new Error(
|
||||
`Queueing is disabled. Cannot enqueue job "${input.jobName}" for queue "${input.queueName}".`
|
||||
);
|
||||
}
|
||||
|
||||
const queue = this.getQueue(input.queueName);
|
||||
await queue.add(input.jobName, input.data, input.options);
|
||||
}
|
||||
|
||||
private getQueue(queueName: string): Queue {
|
||||
const token = getQueueToken(queueName);
|
||||
let queue: Queue | undefined;
|
||||
try {
|
||||
queue = this.moduleRef.get<Queue>(token, { strict: false });
|
||||
} catch {
|
||||
queue = undefined;
|
||||
}
|
||||
|
||||
if (!queue) {
|
||||
throw new Error(
|
||||
`Queue "${queueName}" is not registered. Add it via SchedulerModule.forFeature({ queues: [...] }).`
|
||||
);
|
||||
}
|
||||
|
||||
return queue;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
import { Injectable, NotFoundException } from '@nestjs/common';
|
||||
import { SchedulerJobDescriptor, SchedulerJobScheduleType } from '../interfaces/scheduler-job-descriptor.interface';
|
||||
import { SchedulerQueueJobInput } from '../interfaces/scheduler-queue-job.interface';
|
||||
import { SchedulerJobRegistryService } from './scheduler-job-registry.service';
|
||||
import { SchedulerJobRunnerService } from './scheduler-job-runner.service';
|
||||
import { SchedulerQueueService } from './scheduler-queue.service';
|
||||
|
||||
@Injectable()
|
||||
export class SchedulerService {
|
||||
constructor(
|
||||
private readonly jobRegistry: SchedulerJobRegistryService,
|
||||
private readonly jobRunner: SchedulerJobRunnerService,
|
||||
private readonly queueService: SchedulerQueueService
|
||||
) {}
|
||||
|
||||
listJobs(): SchedulerJobDescriptor[] {
|
||||
return this.jobRegistry.getAll().map((job) => ({
|
||||
id: job.id,
|
||||
providerName: job.providerName,
|
||||
methodName: job.methodName,
|
||||
description: job.options.description,
|
||||
enabled: job.options.enabled,
|
||||
runOnStart: job.options.runOnStart,
|
||||
scheduleType: this.resolveScheduleType(job.options.cron, job.options.intervalMs),
|
||||
executionTarget: job.options.queueName ? 'queue' : 'inline',
|
||||
cron: job.options.cron,
|
||||
intervalMs: job.options.intervalMs,
|
||||
queueName: job.options.queueName,
|
||||
queueJobName: job.options.queueJobName,
|
||||
running: this.jobRunner.isRunning(job.id)
|
||||
}));
|
||||
}
|
||||
|
||||
async triggerNow(jobId: string): Promise<void> {
|
||||
const job = this.jobRegistry.getById(jobId);
|
||||
if (!job) {
|
||||
throw new NotFoundException(`Scheduled job "${jobId}" not found.`);
|
||||
}
|
||||
|
||||
await this.jobRunner.execute(job);
|
||||
}
|
||||
|
||||
async enqueue<TData = unknown>(input: SchedulerQueueJobInput<TData>): Promise<void> {
|
||||
await this.queueService.enqueue(input);
|
||||
}
|
||||
|
||||
private resolveScheduleType(cron?: string, intervalMs?: number): SchedulerJobScheduleType {
|
||||
if (cron) {
|
||||
return 'cron';
|
||||
}
|
||||
if (intervalMs !== undefined) {
|
||||
return 'interval';
|
||||
}
|
||||
return 'manual';
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
import { RegisterQueueOptions } from '@nestjs/bullmq';
|
||||
import { SchedulerQueueRegistration } from '../interfaces/scheduler-feature-options.interface';
|
||||
|
||||
const DEFAULT_QUEUE_NAME = 'default';
|
||||
|
||||
export function normalizeQueueRegistrations(queueRegistrations: SchedulerQueueRegistration[]): RegisterQueueOptions[] {
|
||||
const queueMap = new Map<string, RegisterQueueOptions>();
|
||||
|
||||
for (const registration of queueRegistrations) {
|
||||
const option: RegisterQueueOptions =
|
||||
typeof registration === 'string' ? { name: registration } : { ...registration };
|
||||
const queueName = normalizeQueueName(option.name);
|
||||
queueMap.set(queueName, {
|
||||
...option,
|
||||
name: queueName
|
||||
});
|
||||
}
|
||||
|
||||
return Array.from(queueMap.values());
|
||||
}
|
||||
|
||||
export function normalizeQueueName(name?: string): string {
|
||||
const queueName = name?.trim() ?? DEFAULT_QUEUE_NAME;
|
||||
if (!queueName) {
|
||||
throw new Error('Queue name cannot be empty.');
|
||||
}
|
||||
return queueName;
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
import {
|
||||
ResolvedScheduledJobDefaults,
|
||||
ResolvedSchedulerModuleOptions,
|
||||
SchedulerModuleOptions
|
||||
} from '../interfaces/scheduler-module-options.interface';
|
||||
import { normalizeQueueName, normalizeQueueRegistrations } from './normalize-queue-registrations';
|
||||
import { resolveBullConnection } from './resolve-bull-connection';
|
||||
|
||||
const DEFAULT_JOB_OPTIONS: ResolvedScheduledJobDefaults = {
|
||||
enabled: true,
|
||||
preventOverlap: true,
|
||||
retries: 0,
|
||||
retryDelayMs: 0,
|
||||
timeoutMs: undefined,
|
||||
maxRandomDelayMs: 0
|
||||
};
|
||||
|
||||
const DEFAULT_MODULE_OPTIONS: Omit<ResolvedSchedulerModuleOptions, 'defaultJobOptions'> = {
|
||||
enabled: true,
|
||||
logRegisteredJobs: true,
|
||||
defaultTimezone: undefined,
|
||||
enableQueueing: process.env['REDIS_ENABLED'] === 'true',
|
||||
queueConnection: resolveBullConnection(),
|
||||
queues: [],
|
||||
defaultQueueName: 'default'
|
||||
};
|
||||
|
||||
export function normalizeSchedulerModuleOptions(
|
||||
options: SchedulerModuleOptions = {}
|
||||
): ResolvedSchedulerModuleOptions {
|
||||
const enableQueueing = options.enableQueueing ?? DEFAULT_MODULE_OPTIONS.enableQueueing;
|
||||
const defaultQueueName = normalizeQueueName(options.defaultQueueName ?? DEFAULT_MODULE_OPTIONS.defaultQueueName);
|
||||
const queues = enableQueueing ? normalizeQueueRegistrations(options.queues ?? []) : [];
|
||||
|
||||
return {
|
||||
enabled: options.enabled ?? DEFAULT_MODULE_OPTIONS.enabled,
|
||||
logRegisteredJobs: options.logRegisteredJobs ?? DEFAULT_MODULE_OPTIONS.logRegisteredJobs,
|
||||
defaultTimezone: options.defaultTimezone ?? DEFAULT_MODULE_OPTIONS.defaultTimezone,
|
||||
enableQueueing,
|
||||
queueConnection: resolveBullConnection(options.queueConnection),
|
||||
queues,
|
||||
defaultQueueName,
|
||||
defaultJobOptions: {
|
||||
enabled: options.defaultJobOptions?.enabled ?? DEFAULT_JOB_OPTIONS.enabled,
|
||||
preventOverlap: options.defaultJobOptions?.preventOverlap ?? DEFAULT_JOB_OPTIONS.preventOverlap,
|
||||
retries: clampNonNegativeInteger(options.defaultJobOptions?.retries ?? DEFAULT_JOB_OPTIONS.retries),
|
||||
retryDelayMs: clampNonNegativeInteger(
|
||||
options.defaultJobOptions?.retryDelayMs ?? DEFAULT_JOB_OPTIONS.retryDelayMs
|
||||
),
|
||||
timeoutMs:
|
||||
options.defaultJobOptions?.timeoutMs !== undefined
|
||||
? clampPositiveInteger(options.defaultJobOptions.timeoutMs)
|
||||
: DEFAULT_JOB_OPTIONS.timeoutMs,
|
||||
maxRandomDelayMs: clampNonNegativeInteger(
|
||||
options.defaultJobOptions?.maxRandomDelayMs ?? DEFAULT_JOB_OPTIONS.maxRandomDelayMs
|
||||
)
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
function clampNonNegativeInteger(value: number): number {
|
||||
if (!Number.isFinite(value) || value < 0) {
|
||||
throw new Error(`Invalid non-negative integer value "${value}".`);
|
||||
}
|
||||
return Math.floor(value);
|
||||
}
|
||||
|
||||
function clampPositiveInteger(value: number): number {
|
||||
if (!Number.isFinite(value) || value <= 0) {
|
||||
throw new Error(`Invalid positive integer value "${value}".`);
|
||||
}
|
||||
return Math.floor(value);
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
import { ConnectionOptions } from 'bullmq';
|
||||
|
||||
export function resolveBullConnection(connection?: ConnectionOptions): ConnectionOptions {
|
||||
if (connection) {
|
||||
return connection;
|
||||
}
|
||||
|
||||
if (process.env['REDIS_URL']) {
|
||||
return parseRedisUrl(process.env['REDIS_URL']);
|
||||
}
|
||||
|
||||
return {
|
||||
host: process.env['REDIS_HOST'] || '127.0.0.1',
|
||||
port: parsePort(process.env['REDIS_PORT'], 6379),
|
||||
username: process.env['REDIS_USER'] || undefined,
|
||||
password: process.env['REDIS_PASSWORD'] || undefined,
|
||||
...(process.env['REDIS_TLS'] === 'true' ? { tls: {} } : {})
|
||||
};
|
||||
}
|
||||
|
||||
function parseRedisUrl(redisUrl: string): ConnectionOptions {
|
||||
const parsed = new URL(redisUrl);
|
||||
const isTls = parsed.protocol === 'rediss:' || process.env['REDIS_TLS'] === 'true';
|
||||
|
||||
return {
|
||||
host: parsed.hostname,
|
||||
port: parsePort(parsed.port, 6379),
|
||||
username: parsed.username || undefined,
|
||||
password: parsed.password || undefined,
|
||||
...(isTls ? { tls: {} } : {})
|
||||
};
|
||||
}
|
||||
|
||||
function parsePort(value: string | undefined, fallback: number): number {
|
||||
const parsed = Number.parseInt(value ?? '', 10);
|
||||
if (!Number.isFinite(parsed) || parsed <= 0) {
|
||||
return fallback;
|
||||
}
|
||||
return parsed;
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
{
|
||||
"extends": "../../tsconfig.base.json",
|
||||
"compilerOptions": {
|
||||
"module": "commonjs",
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"strict": true,
|
||||
"importHelpers": true,
|
||||
"noImplicitOverride": true,
|
||||
"noImplicitReturns": true,
|
||||
"noFallthroughCasesInSwitch": true,
|
||||
"noPropertyAccessFromIndexSignature": true
|
||||
},
|
||||
"files": [],
|
||||
"include": [],
|
||||
"references": [
|
||||
{
|
||||
"path": "./tsconfig.lib.json"
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
{
|
||||
"extends": "./tsconfig.json",
|
||||
"compilerOptions": {
|
||||
"outDir": "../../dist/out-tsc",
|
||||
"declaration": true,
|
||||
"types": ["node"],
|
||||
"target": "es2021",
|
||||
"experimentalDecorators": true,
|
||||
"emitDecoratorMetadata": true,
|
||||
"strictNullChecks": true,
|
||||
"noImplicitAny": true,
|
||||
"strictBindCallApply": true,
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"noFallthroughCasesInSwitch": true
|
||||
},
|
||||
"include": ["src/**/*.ts"]
|
||||
}
|
||||
@@ -62,6 +62,7 @@
|
||||
"@gauzy/plugin-camshot": ["./packages/plugins/camshot/src/index.ts"],
|
||||
"@gauzy/plugin-soundshot": ["./packages/plugins/soundshot/src/index.ts"],
|
||||
"@gauzy/mcp-server": ["./packages/mcp-server/src/index.ts"],
|
||||
"@gauzy/scheduler": ["./packages/scheduler/src/index.ts"],
|
||||
"@gauzy/ui-auth": ["./packages/ui-auth/src/index.ts"],
|
||||
"@gauzy/ui-config": ["./packages/ui-config/src/index.ts"],
|
||||
"@gauzy/ui-core": ["./packages/ui-core/src/index.ts"],
|
||||
|
||||
@@ -8528,6 +8528,21 @@
|
||||
version "4.0.1"
|
||||
resolved "https://codeload.github.com/ever-co/nestjs-axios/tar.gz/6c7a33d5147f2218e76808f628fabbf4d4f616ad"
|
||||
|
||||
"@nestjs/bull-shared@^11.0.4":
|
||||
version "11.0.4"
|
||||
resolved "https://registry.npmjs.org/@nestjs/bull-shared/-/bull-shared-11.0.4.tgz#6179bb0a0ae705193a113ea60021d28732c6038d"
|
||||
integrity sha512-VBJcDHSAzxQnpcDfA0kt9MTGUD1XZzfByV70su0W0eDCQ9aqIEBlzWRW21tv9FG9dIut22ysgDidshdjlnczLw==
|
||||
dependencies:
|
||||
tslib "2.8.1"
|
||||
|
||||
"@nestjs/bullmq@^11.0.4":
|
||||
version "11.0.4"
|
||||
resolved "https://registry.npmjs.org/@nestjs/bullmq/-/bullmq-11.0.4.tgz#aa0d62f949a9bfa7456b43780aa7b404dc5536b3"
|
||||
integrity sha512-wBzK9raAVG0/6NTMdvLGM4/FQ1lsB35/pYS8L6a0SDgkTiLpd7mAjQ8R692oMx5s7IjvgntaZOuTUrKYLNfIkA==
|
||||
dependencies:
|
||||
"@nestjs/bull-shared" "^11.0.4"
|
||||
tslib "2.8.1"
|
||||
|
||||
"@nestjs/cache-manager@^3.1.0":
|
||||
version "3.1.0"
|
||||
resolved "https://registry.yarnpkg.com/@nestjs/cache-manager/-/cache-manager-3.1.0.tgz#65f3be68154e48557790c0a58a55aadb92439dd0"
|
||||
@@ -18441,6 +18456,19 @@ builtin-status-codes@^3.0.0:
|
||||
resolved "https://registry.yarnpkg.com/builtin-status-codes/-/builtin-status-codes-3.0.0.tgz#85982878e21b98e1c66425e03d0174788f569ee8"
|
||||
integrity sha512-HpGFw18DgFWlncDfjTa2rcQ4W88O1mC8e8yZ2AvQY5KDaktSTwo+KRf6nHK6FRI5FyRyb/5T6+TSxfP7QyGsmQ==
|
||||
|
||||
bullmq@^5.70.1:
|
||||
version "5.70.1"
|
||||
resolved "https://registry.npmjs.org/bullmq/-/bullmq-5.70.1.tgz#34c64dbd17bc05b1299cc07e85079ad1c94eb1ed"
|
||||
integrity sha512-HjfGHfICkAClrFL0Y07qNbWcmiOCv1l+nusupXUjrvTPuDEyPEJ23MP0lUwUs/QEy1a3pWt/P/sCsSZ1RjRK+w==
|
||||
dependencies:
|
||||
cron-parser "4.9.0"
|
||||
ioredis "5.9.3"
|
||||
msgpackr "1.11.5"
|
||||
node-abort-controller "3.1.1"
|
||||
semver "7.7.4"
|
||||
tslib "2.8.1"
|
||||
uuid "11.1.0"
|
||||
|
||||
bundle-name@^4.1.0:
|
||||
version "4.1.0"
|
||||
resolved "https://registry.yarnpkg.com/bundle-name/-/bundle-name-4.1.0.tgz#f3b96b34160d6431a19d7688135af7cfb8797889"
|
||||
@@ -20553,14 +20581,14 @@ create-require@^1.1.0:
|
||||
resolved "https://registry.yarnpkg.com/create-require/-/create-require-1.1.1.tgz#c1d7e8f1e5f6cfc9ff65f9cd352d37348756c333"
|
||||
integrity sha512-dcKFX3jn0MpIaXjisoRvexIJVEKzaq7z2rZKxf+MSr9TkdmHmsU4m2lcLojrj/FHl8mk5VxMmYA+ftRkP/3oKQ==
|
||||
|
||||
cron-parser@^4.2.0:
|
||||
cron-parser@4.9.0, cron-parser@^4.2.0:
|
||||
version "4.9.0"
|
||||
resolved "https://registry.yarnpkg.com/cron-parser/-/cron-parser-4.9.0.tgz#0340694af3e46a0894978c6f52a6dbb5c0f11ad5"
|
||||
integrity sha512-p0SaNjrHOnQeR8/VnfGbmg9te2kfyYSQ7Sc/j/6DtPL3JQvKxmjO9TSjNFpujqV3vEYYBvNNvXSxzyksBWAx1Q==
|
||||
dependencies:
|
||||
luxon "^3.2.1"
|
||||
|
||||
cron@4.4.0:
|
||||
cron@4.4.0, cron@^4.3.3:
|
||||
version "4.4.0"
|
||||
resolved "https://registry.npmjs.org/cron/-/cron-4.4.0.tgz#1488444a23ea7134e2b7686c17711abdffcebba8"
|
||||
integrity sha512-fkdfq+b+AHI4cKdhZlppHveI/mgz2qpiYxcm+t5E5TsxX7QrLS1VE0+7GENEk9z0EeGPcpSciGv6ez24duWhwQ==
|
||||
@@ -27086,7 +27114,7 @@ ionicons@^4.6.3:
|
||||
resolved "https://registry.yarnpkg.com/ionicons/-/ionicons-4.6.3.tgz#e4f3a3e9b66761eb8c0db27798cc1a8ce97777d2"
|
||||
integrity sha512-cgP+VIr2cTJpMfFyVHTerq6n2jeoiGboVoe3GlaAo5zoSBDAEXORwUZhv6m+lCyxlsHCS3nqPUE+MKyZU71t8Q==
|
||||
|
||||
ioredis@^5.7.0:
|
||||
ioredis@5.9.3, ioredis@^5.7.0:
|
||||
version "5.9.3"
|
||||
resolved "https://registry.yarnpkg.com/ioredis/-/ioredis-5.9.3.tgz#e897af9f87ee4b7bc61d8bd6373f466aca43d4e0"
|
||||
integrity sha512-VI5tMCdeoxZWU5vjHWsiE/Su76JGhBvWF1MJnV9ZtGltHk9BmD48oDq8Tj8haZ85aceXZMxLNDQZRVo5QKNgXA==
|
||||
@@ -32374,6 +32402,13 @@ msgpackr-extract@^3.0.2:
|
||||
"@msgpackr-extract/msgpackr-extract-linux-x64" "3.0.3"
|
||||
"@msgpackr-extract/msgpackr-extract-win32-x64" "3.0.3"
|
||||
|
||||
msgpackr@1.11.5:
|
||||
version "1.11.5"
|
||||
resolved "https://registry.npmjs.org/msgpackr/-/msgpackr-1.11.5.tgz#edf0b9d9cb7d8ed6897dd0e42cfb865a2f4b602e"
|
||||
integrity sha512-UjkUHN0yqp9RWKy0Lplhh+wlpdt9oQBYgULZOiFhV3VclSF1JnSQWZ5r9gORQlNYaUKQoR8itv7g7z1xDDuACA==
|
||||
optionalDependencies:
|
||||
msgpackr-extract "^3.0.2"
|
||||
|
||||
msgpackr@^1.11.2:
|
||||
version "1.11.8"
|
||||
resolved "https://registry.yarnpkg.com/msgpackr/-/msgpackr-1.11.8.tgz#8283c79eb6e5d488f6fb3fac4996006baa390614"
|
||||
@@ -38818,16 +38853,16 @@ semver@7.7.3:
|
||||
resolved "https://registry.yarnpkg.com/semver/-/semver-7.7.3.tgz#4b5f4143d007633a8dc671cd0a6ef9147b8bb946"
|
||||
integrity sha512-SdsKMrI9TdgjdweUSR9MweHA4EJ8YxHn8DFaDisvhVlUOe4BF1tLD7GAj0lIqWVl+dPb/rExr0Btby5loQm20Q==
|
||||
|
||||
semver@7.7.4, semver@^7.0.0, semver@^7.1.1, semver@^7.1.2, semver@^7.1.3, semver@^7.3.2, semver@^7.3.4, semver@^7.3.5, semver@^7.3.7, semver@^7.3.8, semver@^7.5.2, semver@^7.5.3, semver@^7.5.4, semver@^7.6.0, semver@^7.6.2, semver@^7.6.3, semver@^7.7.2, semver@^7.7.3, semver@^7.7.4, semver@~7.7.3:
|
||||
version "7.7.4"
|
||||
resolved "https://registry.yarnpkg.com/semver/-/semver-7.7.4.tgz#28464e36060e991fa7a11d0279d2d3f3b57a7e8a"
|
||||
integrity sha512-vFKC2IEtQnVhpT78h1Yp8wzwrf8CM+MzKMHGJZfBtzhZNycRFnXsHk6E5TxIkkMsgNS7mdX3AGB7x2QM2di4lA==
|
||||
|
||||
semver@^6.0.0, semver@^6.2.0, semver@^6.3.0, semver@^6.3.1:
|
||||
version "6.3.1"
|
||||
resolved "https://registry.yarnpkg.com/semver/-/semver-6.3.1.tgz#556d2ef8689146e46dcea4bfdd095f3434dffcb4"
|
||||
integrity sha512-BR7VvDCVHO+q2xBEWskxS6DJE1qRnb7DxzUrogb71CWoSficBxYsiAGd+Kl0mmq/MprG9yArRkyrQxTO6XjMzA==
|
||||
|
||||
semver@^7.0.0, semver@^7.1.1, semver@^7.1.2, semver@^7.1.3, semver@^7.3.2, semver@^7.3.4, semver@^7.3.5, semver@^7.3.7, semver@^7.3.8, semver@^7.5.2, semver@^7.5.3, semver@^7.5.4, semver@^7.6.0, semver@^7.6.2, semver@^7.6.3, semver@^7.7.2, semver@^7.7.3, semver@^7.7.4, semver@~7.7.3:
|
||||
version "7.7.4"
|
||||
resolved "https://registry.yarnpkg.com/semver/-/semver-7.7.4.tgz#28464e36060e991fa7a11d0279d2d3f3b57a7e8a"
|
||||
integrity sha512-vFKC2IEtQnVhpT78h1Yp8wzwrf8CM+MzKMHGJZfBtzhZNycRFnXsHk6E5TxIkkMsgNS7mdX3AGB7x2QM2di4lA==
|
||||
|
||||
send@^1.1.0, send@^1.2.0, send@latest:
|
||||
version "1.2.1"
|
||||
resolved "https://registry.yarnpkg.com/send/-/send-1.2.1.tgz#9eab743b874f3550f40a26867bf286ad60d3f3ed"
|
||||
|
||||
Reference in New Issue
Block a user