Files
Trapa-EurekaandClaude Opus 5 d810c81cf3 feat(core): persistence invariants, config/boundary validation, migration smoke test and correlation IDs (#10198)
Adds test harnesses for tenant isolation, TypeORM/MikroORM parity, persistence invariants and idempotency; startup validation for DB config; Nx plugin/core boundary rules; a fresh-database migration smoke test; and a correlation id carried from HTTP requests into docs queue jobs and logs.

Testing (packages/core/src/lib/core/testing):
- tenant-isolation: an in-memory tenant-aware repository, fixtures and strict assertions (a non-empty page, every row in the caller's tenant). Applied to EmployeeService and OrganizationProjectService.
- orm-conformance and persistence-invariants: suites that run the same checks under TypeORM and MikroORM (run-both-orms.sh uses yarn nx).
- idempotency: assertion helpers and specs for token cleanup (positive control), employee notifications and the Zapier timer webhook.
- database/migration-smoke.spec.ts: runs the full migration chain on a fresh better-sqlite3 database. It is excluded from the default core jest run; run it with `nx run core:test-migration-smoke`.

Idempotency:
- EmployeeNotificationService.create() takes an opt-in `absorbRedelivery`. With it set, an identical unread, unarchived notification under 60 s old (same receiver, entity, type, sender, title and message) is returned instead of inserted, and a missing key or a failed lookup inserts as usual. The event handler does not enable it (the in-process EventBus never redelivers), so every caller inserts one row per event as before.
- ZapierWebhookService skips resending a webhook that already succeeded for the same subscription, action and time log within 5 minutes. The key is reserved while in flight, failed deliveries stay retryable, pruning stops at the first unexpired entry, and the cache is capped at 10,000 entries.

Config (packages/config):
- An unknown DB_TYPE fails fast with the list of supported values; an empty DB_TYPE still means better-sqlite3; mongodb throws an Error.
- Pool and timeout variables are parsed with Number.parseInt semantics, only for postgres and mysql. They throw only where tarn already refused to start and warn otherwise. SQLite ignores them and still prints the startup values.

Observability:
- RequestContextMiddleware accepts an inbound x-correlation-id of 1-128 visible ASCII characters (otherwise it generates a UUIDv4) and echoes it on the response. CORS allows and exposes X-Correlation-Id. RequestContext.currentCorrelationId() is added.
- The docs queue carries the correlation id through job payloads (including bulk reindex) into pipeline outcome, error, dead-letter and enqueue-failure logs.

Boundaries and build:
- Every project gets a type tag. The ESLint depConstraints stop type:core depending on plugins, and plugins may depend only on core, shared libs and the known extension points (ai-chat, job-proposal, integration-ai, job-*-ui).
- core declares @gauzy/scheduler (package.json and implicitDependencies); the webapp Dockerfile copies scheduler's package.json.
- The integration-zapier jest config can import @gauzy/core (transformIgnorePatterns, allowJs, isolatedModules).
- Fixes the stale @nrwl/nx eslint-disable id in the e2e roles-permissions steps and the broken @gauzy/core mock in the docs document-scope spec.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-17 23:35:52 +02:00
..

@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)

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.

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(...).

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.

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)

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

nx build scheduler