Jobs that keep the request's context
A job runs after the request that enqueued it is gone. With
@nestjs-pipeline/job-context, the
job carries the request’s tenant, correlation id and principal, and restores them before it
runs, so a pipeline it dispatches makes the decisions the request would have made. This
recipe follows the welcome email of A command through every
layer.
1. Register the module
Section titled “1. Register the module”AppModule registers it once, with the application’s principal port, the tenants a job may
run in, and the stores the tenant and correlation id are read from:
JobContextModule.forRoot({ principal: SessionJobPrincipal, tenants: () => persistenceConfig().tenants, sources: contextSources, imports: [AuthsModule],}),tenants is a function, read when the application starts, so the list comes from the same
configuration as the database.
2. Stamp the payload when enqueuing
Section titled “2. Stamp the payload when enqueuing”The application layer only knows a dispatcher port. Its BullMQ adapter wraps each payload in
withJobContext, inside the request, while the context is current:
/* Copyright (C) 2026-present Aristotelis — see repository license. */
import { InjectQueue } from '@nestjs/bullmq';import { Injectable } from '@nestjs/common';import { withJobContext } from '@nestjs-pipeline/job-context';import type { Queue } from 'bullmq';import type { IUserBatchDispatcher, IWelcomeEmailDispatcher, UserBatchDispatchItem, WelcomeEmailDispatch,} from '../application/ports/user-event-dispatcher.port.js';import { BATCH_UPDATE_USERS_QUEUE, type BatchUpdateUsersJobData,} from './batch-update-users.processor.js';import { WELCOME_EMAIL_QUEUE, type WelcomeEmailJobData,} from './send-welcome-email.processor.js';
/** * BullMQ infrastructure adapter for user-event application dispatch ports. * Both queues stamp the caller's tenant, correlation id and principal into the * job payload with `withJobContext`; the processors restore them with * `@InJobContext`. */@Injectable()export class BullMqUserEventDispatcher implements IWelcomeEmailDispatcher, IUserBatchDispatcher{ constructor( @InjectQueue(WELCOME_EMAIL_QUEUE) private readonly welcomeEmailQueue: Queue<WelcomeEmailJobData>, @InjectQueue(BATCH_UPDATE_USERS_QUEUE) private readonly batchUpdateQueue: Queue<BatchUpdateUsersJobData>, ) {}
async enqueueWelcomeEmail(message: WelcomeEmailDispatch): Promise<void> { await this.welcomeEmailQueue.add('send', withJobContext(message)); }
async enqueueUserBatch( items: readonly UserBatchDispatchItem[], ): Promise<void> { await this.batchUpdateQueue.add( 'batch-update', withJobContext({ items: items.map((item) => ({ ...item })) }), ); }}3. Restore it in the processor
Section titled “3. Restore it in the processor”@InJobContext() validates job.data.jobContext and runs process inside the tenant, the
correlation id and the restored principal. TenantSchemaContext and getCorrelationId()
then answer as they would in the request:
/* Copyright (C) 2026-present Aristotelis — see repository license. */
import { Processor, WorkerHost } from '@nestjs/bullmq';import { Logger, OnModuleDestroy } from '@nestjs/common';import { getCorrelationId } from '@nestjs-pipeline/correlation';import { InJobContext, type WithJobContext,} from '@nestjs-pipeline/job-context';import { TenantSchemaContext } from '@persistence/tenant-schema.context.js';import type { Job } from 'bullmq';
export const WELCOME_EMAIL_QUEUE = 'welcome-email';
export type WelcomeEmailJobData = WithJobContext<{ userId: string; username: string; email: string;}>;
export interface SimulatedWelcomeEmailResult { readonly simulated: true; readonly emailSent: false; readonly recipient: string; readonly userId: string;}
@Processor(WELCOME_EMAIL_QUEUE)export class SendWelcomeEmailProcessor extends WorkerHost implements OnModuleDestroy{ private readonly logger = new Logger(SendWelcomeEmailProcessor.name);
constructor(private readonly tenantContext: TenantSchemaContext) { super(); }
async onModuleDestroy(): Promise<void> { try { await this.worker.close(); } catch { // Worker was not initialized or already closed. } }
@InJobContext() async process( job: Job<WelcomeEmailJobData>, ): Promise<SimulatedWelcomeEmailResult> { this.logger.log( `[Simulated] Demonstrating welcome email dispatch for ${job.data.email} ` + `(user: ${job.data.username}, tenant: ${this.tenantContext.schema}, correlationId: ${getCorrelationId()}). No external email sent.`, );
return { simulated: true, emailSent: false, recipient: job.data.email, userId: job.data.userId, }; }}4. Re-check the principal when the job runs
Section titled “4. Re-check the principal when the job runs”A payload is data: anyone who can write to the queue can write it. So it carries an identity
reference only, never permissions, and the principal port decides what that identity may do
now. SessionJobPrincipal reloads the session, so a revoked session or a deleted user
refuses the job:
/* Copyright (C) 2026-present Aristotelis — see repository license. */
import { type IQueryRepository, type IWriteSideAggregateRepository, requireTenant,} from '@cqrs-ddd/core/application';import { Inject, Injectable } from '@nestjs/common';import type { Capability } from '@nestjs-pipeline/casl';import { type IJobPrincipal, InvalidJobContextError, type PrincipalReference,} from '@nestjs-pipeline/job-context';import { getSessionPrincipal, sessionPrincipalStore,} from '../../common/context/session-principal.store.js';import { API_CLIENTS } from '../../common/environment/api-clients.config.js';import { isPrincipalType, isSessionPrincipalValid, type SessionPrincipal,} from '../../common/types/session-principal.js';import { GetUserQuery } from '../../users/application/cqrs/queries/get-user.query.js';import type { User } from '../../users/domain/models/user.entity.js';import { EXT_USER_QUERY_REPOSITORY } from '../../users/persistence/repository.tokens.js';import type { Auth } from '../domain/models/auth.entity.js';import { COMMAND_REPOSITORY } from '../persistence/repository.tokens.js';
/** * The api's `IJobPrincipal`: a job acts for the principal of the request that * enqueued it, bound through `sessionPrincipalStore` the way * `SessionPrincipalContextInterceptor` binds a request's. * * Nothing the payload says is trusted beyond identity. A user is bound only * while its `Auth` session exists, belongs to it, is not revoked and has not * expired, and its user row exists; it is bound without grants, so * `CaslPermissionSource` reads its current rules. An API client is bound only * while `API_CLIENTS` still lists it for the job's tenant, with the grants * listed there now. `@AsSystem` grants are bound as declared. */@Injectable()export class SessionJobPrincipal implements IJobPrincipal<Capability> { constructor( @Inject(COMMAND_REPOSITORY.updateAuth) private readonly sessions: IWriteSideAggregateRepository<Auth>, @Inject(EXT_USER_QUERY_REPOSITORY.getUser) private readonly users: IQueryRepository<GetUserQuery, User | null>, ) {}
capture(): PrincipalReference | undefined { const principal = getSessionPrincipal(); if (!principal || !isSessionPrincipalValid(principal)) return undefined; return { id: principal.id, type: principal.type, sessionId: principal.sid }; }
async restore<T>( reference: PrincipalReference, work: () => Promise<T>, grants?: readonly Capability[], ): Promise<T> { const principal = await this.resolve(reference, grants); return sessionPrincipalStore.run(principal, work); }
private async resolve( { id, type, sessionId }: PrincipalReference, grants: readonly Capability[] | undefined, ): Promise<SessionPrincipal> { if (!isPrincipalType(type)) { throw new InvalidJobContextError(`principal type "${type}" is unknown`); } const tenant = requireTenant('restoring a job principal'); if (grants) return { id, type, tenant, grants: [...grants] };
if (type === 'service') { const client = API_CLIENTS.get(id); if (!client?.tenants.has(tenant)) { throw new InvalidJobContextError( 'the API client is no longer allowed in this tenant', ); } return { id, type, tenant, grants: client.grants }; }
const session = sessionId ? await this.sessions.findById(sessionId) : null; if ( !session || session.userId !== id || session.revokedAt !== null || Date.now() >= session.expiresAt ) { throw new InvalidJobContextError( 'the user session is unknown, revoked or expired', ); } const user = await this.users.find( new GetUserQuery({ userId: id }, { refresh: true }), ); if (!user) { throw new InvalidJobContextError('the user no longer exists'); } return { id, type, tenant, sid: sessionId }; }}