Skip to content

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.

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.

The application layer only knows a dispatcher port. Its BullMQ adapter wraps each payload in withJobContext, inside the request, while the context is current:

api/src/users/jobs/bullmq-user-event-dispatcher.adapter.ts
/* 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 })) }),
);
}
}

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

api/src/users/jobs/send-welcome-email.processor.ts
/* 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:

api/src/auths/infrastructure/session-job-principal.ts
/* 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 };
}
}