Skip to content

@nestjs-pipeline/audit

npm version License

Audit-trail behavior for @nestjs-pipeline/core — records who did what, when, and with what outcome for every command (and, when listed in captureKinds, every query or event handler), and forwards the record to a pluggable audit sink. Captured application values must satisfy that sink’s serialization requirements.

Sink-agnostic: it depends only on a tiny AuditSink interface. A zero-dependency console sink is the default; Postgres is a genuine drop-in, and your own sink (event store, Kafka, HTTP collector, …) is a one-line swap — handlers never change. Records are written on both success and failure, sensitive payload fields are redacted by default, and the actor can be resolved from the pipeline context.



An audit trail is a classic cross-cutting concern: the same “record who did what” logic is needed on dozens of handlers. Writing it inline couples every handler to your audit storage and is easy to get wrong (forgetting failures, leaking passwords, missing the actor). This behavior centralizes it:

  • Success and failure — denied/rejected attempts are audited too (the part most hand-rolled trails miss). With the default failOpen: true, a handler failure is re-thrown unchanged even if the audit sink also fails.
  • Redaction built in — password, token, secret, … are masked before anything is stored.
  • Actor resolution — pull the acting principal from context.items populated by an upstream auth behavior.
  • One seam — the AuditSink interface. Console today, Postgres or your event store tomorrow, with no handler changes.

It generalizes the Audit-Trail example from the root README into a reusable, redaction-aware, outcome-aware package.


Terminal window
pnpm add @nestjs-pipeline/audit

Peer dependencies:

Terminal window
pnpm add @nestjs-pipeline/core @nestjs/common reflect-metadata

Requires Node.js 22.12 or later, @nestjs/common ^12.1.0 and @nestjs-pipeline/core ^0.4.2.

Published as an ES module; a CommonJS application loads it with require(). Coming from 0.3.x, see Upgrading from 0.3.x.

The bundled sinks are typed structurally, so this package adds zero heavy dependencies. For the Postgres sink, add a pg Pool/Client in your app (pnpm add pg); the console sink needs nothing.


Register the module and add AuditBehavior to your global behaviors (or per-handler via @UsePipeline).

import { Module } from '@nestjs/common';
import { PipelineModule } from '@nestjs-pipeline/core';
import { AuditModule, AuditBehavior } from '@nestjs-pipeline/audit';
@Module({
imports: [
// Zero-config: audit every command to the console.
AuditModule.forRoot(),
PipelineModule.forRoot({
globalBehaviors: { scope: 'all', before: [AuditBehavior] },
}),
],
})
export class AppModule {}

Ordering: place AuditBehavior near the outside of the chain so its duration covers the whole handler, and after any auth behavior that populates context.items for the actor factory.


Every audited run produces one AuditRecord, forwarded to the sink:

{
"id": "0197…", // UUIDv7 per entry
"correlationId": "019728a3-…",
"tenantId": "acme", // context.tenantId, when present
"action": "user.create", // defaults to requestName
"severity": "medium", // 'low' | 'medium' | 'high' | 'critical'
"outcome": "success", // or 'failure'
"actor": { "id": "admin-1" }, // resolved from context (optional)
"requestKind": "command",
"requestName": "CreateUserCommand",
"handlerName": "CreateUserHandler",
"payload": { "username": "jane", "password": "[REDACTED]" },
"response": undefined, // only when captureResponse: true
"error": undefined, // present on failure
"durationMs": 12.3,
"timestamp": "2026-03-01T12:00:00.000Z",
"metadata": { "tenantId": "acme" } // metadata factory output, plus tenantId when present
}

A sink implements AuditSink: write, and optionally begin:

interface AuditSink {
write(record: AuditRecord): Promise<void> | void;
begin?(record: AuditStartRecord): Promise<void> | void;
}

begin receives a pending record before the handler runs; write receives the final record under the same id and must replace it. A durable sink should implement both (see Architecture and delivery guarantees).

Zero-dependency; writes each record as a JSON line. Successes go to log, failures to warn. Used automatically when no sink is passed.

import { Logger } from '@nestjs/common';
import { LogAuditSink } from '@nestjs-pipeline/audit';
AuditModule.forRoot({
sink: new LogAuditSink({ logger: new Logger('Audit'), pretty: true }),
});

Implements begin and write: it inserts a pending row when a command starts and completes that row (INSERT … ON CONFLICT (id) DO UPDATE) when it finishes. Create the table once with createAuditTableSql.

import { Module } from '@nestjs/common';
import { Pool } from 'pg';
import {
AuditModule,
PostgresAuditSink,
createAuditTableSql,
} from '@nestjs-pipeline/audit';
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
await pool.query(createAuditTableSql()); // → table "audit_log"
@Module({
imports: [
AuditModule.forRoot({
sink: new PostgresAuditSink(pool, { table: 'audit_log' }),
}),
],
})
export class AppModule {}

Or build it from a DI-managed pool with forRootAsync:

AuditModule.forRootAsync({
inject: [PG_POOL],
useFactory: (pool: Pool) => new PostgresAuditSink(pool),
});

The table name is validated as a plain SQL identifier (interpolated, not parameterized); every record value is passed as a bound parameter.

PostgreSQL jsonb cannot hold a NUL character or a lone UTF-16 surrogate, and rejects the whole INSERT when a value contains one. The sink stores each such character as U+FFFD (�) instead, so the record is kept.

A row still pending long after it started is an attempt interrupted by a process stop; its outcome is unknown. duration_ms and completed_at stay null:

SELECT * FROM audit_log
WHERE outcome = 'pending' AND occurred_at < now() - interval '1 hour';

Anything that matches AuditSink works — an event store, Kafka, an HTTP collector, your domain repository:

import type { AuditRecord, AuditSink } from '@nestjs-pipeline/audit';
export class KafkaAuditSink implements AuditSink {
constructor(private readonly producer: Producer) {}
async write(record: AuditRecord): Promise<void> {
await this.producer.send({
topic: 'audit',
messages: [{ key: record.correlationId, value: JSON.stringify(record) }],
});
}
}
AuditModule.forRoot({ sink: new KafkaAuditSink(producer) });

AuditBehavior times the handler, builds an AuditRecord, and writes it to the sink. On success it returns the handler response after the sink write. On handler failure it attempts to write the failure record and then re-throws the original handler error when the sink write succeeds or failOpen: true suppresses a sink failure. If the sink throws while failOpen: false:

  • on the success path, the sink error propagates, failing the request;
  • on the handler-failure path, the caller receives the handler’s own error, unchanged, and the sink error is logged. The produced record is also stashed on context.items under AUDIT_RECORD_ITEM for any later behavior to read.

When the sink implements begin (and recordStart is not false), a pending start record is written before the handler runs and stashed under AUDIT_START_RECORD_ITEM_TOKEN; the final record reuses its id.

Opt in per handler with options:

import { CommandHandler } from '@nestjs/cqrs';
import { UsePipeline } from '@nestjs-pipeline/core';
import { audit } from '@nestjs-pipeline/audit';
@CommandHandler(DeleteUserCommand)
@UsePipeline(
audit({
action: 'user.delete',
severity: 'high',
actor: (c) => ({ id: c.items.get('currentUserId') as string }),
}),
)
export class DeleteUserHandler { /* ... */ }

The raw tuple form @UsePipeline([AuditBehavior, { ... }]) remains supported as an escape hatch.


AuditBehavior is a pipeline step around the handler. It never sees the handler’s database transaction: every sink call is a separate operation.

For each audited request (commands only, unless captureKinds lists more):

  1. Start — when the sink implements begin, build the pending AuditStartRecord (actor, action, redacted payload, outcome: 'pending') and call sink.begin(record). Stash it under AUDIT_START_RECORD_ITEM_TOKEN.
  2. Handler — next() runs the rest of the pipeline and the handler, which commits its own changes.
  3. Finish — build the final AuditRecord under the same id (success or failure, duration, optional response, error) and call sink.write(record). Stash it under AUDIT_RECORD_ITEM.

What survives a process stop at each point:

Stop happens Sink with begin (Postgres) Sink without begin (console)
before step 1 completes no row, and the handler did not run no line, and the handler did not run
during step 2 or before step 3 completes a pending row: the attempt is known, its outcome is not nothing: the attempt is lost
after step 3 the final row the final line

Two consequences follow:

  • A pending row does not say whether the handler’s changes committed. Treat it as “outcome unknown” and reconcile it against the business data.
  • With failOpen: false, a failed begin stops the request before the handler runs: no audit, no action. A failed final write fails a successful request, but the handler’s changes are already committed; the row stays pending.

These guarantees hold for any sink and any database, because the behavior does not need to join the handler’s transaction. What they do not give is an audit row that commits together with the business change. That needs the application’s cooperation, described next.


Recording the audit row atomically with the business write

Section titled “Recording the audit row atomically with the business write”

The only way to guarantee “a committed change always has its audit row, and a rolled-back change has none” is to insert the audit row in the same database transaction as the business change. A pipeline behavior cannot do that on its own: it runs outside the handler and cannot see its unit of work. The application must do it in its persistence layer.

The pattern (a transactional audit record):

  1. Use a sink whose begin stores nothing and whose write completes a row by id (an upsert, like PostgresAuditSink.write). begin must exist so that the pending record is built and stashed.
  2. In the handler’s repository, inside the transaction that writes the business change, read the pending record with getPipelineItem(context, AUDIT_START_RECORD_ITEM_TOKEN) (or pass it in from the handler) and insert it into the audit table.
  3. After the handler returns, the behavior’s write completes that row with the outcome. If the transaction rolled back, write inserts a failure row; if the process stops before write, the committed row stays pending but its business change is known to have committed.
import {
type AuditRecord,
type AuditSink,
PostgresAuditSink,
} from '@nestjs-pipeline/audit';
class TransactionalAuditSink implements AuditSink {
constructor(private readonly completing: PostgresAuditSink) {}
begin(): void {
// The repository inserts the pending row in its own transaction.
}
write(record: AuditRecord): Promise<void> {
return this.completing.write(record); // upsert by id
}
}

Inside the repository, the pending record is read from the pipeline context:

import { getPipelineItem } from '@nestjs-pipeline/core';
import { AUDIT_START_RECORD_ITEM_TOKEN } from '@nestjs-pipeline/audit';
const pending = getPipelineItem(context, AUDIT_START_RECORD_ITEM_TOKEN);
if (pending) {
await tx.query(
'INSERT INTO audit_log (id, correlation_id, action, severity, outcome, request_kind, request_name, handler_name, occurred_at) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)',
[pending.id, pending.correlationId, pending.action, pending.severity, pending.outcome,
pending.requestKind, pending.requestName, pending.handlerName, pending.timestamp],
);
}

Checklist for doing it properly:

  • The audit table lives in the same database as the business data, and the insert runs on the same connection and transaction as the business write.
  • Completion is idempotent by id (an upsert), so a retried write, or a write without a stored start row, is safe.
  • The row is redacted before storage: use the record from the token, never the raw request.
  • A reconciliation job handles rows left pending (compare with the business data, then mark them).
  • If a message broker must also receive the audit event, publish it from the stored row (a transactional outbox), not from the request path.

Limit with autocommit-only repositories: some write sides refuse to run inside an outer transaction, because they acknowledge the write (advance a version, update a cache) as soon as their statement succeeds, which would be wrong after a rollback they cannot observe. Such a repository cannot share a transaction with the audit insert until it gains commit hooks (acknowledgment and cache work run after the commit). Until then, rely on the two-phase guarantees above.


Per-handler options (AuditBehaviorOptions) shallow-merge over the module defaults passed to AuditModule.forRoot({ defaults }):

Option Type Default Description
action string context.requestName Logical action name on the record
severity 'low' | 'medium' | 'high' | 'critical' 'medium' ('low' for queries) Importance, for filtering/alerting
actor (ctx) => AuditActor | undefined — Resolve the acting principal
captureRequest boolean true Record the (redacted) request payload
captureResponse boolean false Record the (redacted) handler response
recordStart boolean true Write a pending start record first, when the sink implements begin
captureKinds AuditRequestKind[] ['command'] Request kinds to audit; list 'query' or 'event' to audit them too
redactKeys string[] — Extra field names to mask (merged with defaults)
redact (value) => unknown — Full custom redactor (replaces key-masking)
metadata (ctx) => object — Extra metadata merged into the record
includeStack boolean true Include the error stack on failure records
failOpen boolean true Log/ignore sink failures (true) or propagate the sink error (false)

Before a payload or response is stored, the values of sensitive keys are replaced with '[REDACTED]'. The built-in DEFAULT_REDACT_KEYS cover common secrets (password, pwd, token, accessToken, refreshToken, secret, apiKey, authorization, cookie, ssn, creditCard, cardNumber, cvv). Matching is case-insensitive and recurses into nested objects and arrays.

// Add app-specific keys (merged with the defaults):
@UsePipeline([AuditBehavior, { redactKeys: ['pin', 'iban'] }])
// Or take full control:
@UsePipeline([AuditBehavior, {
redact: (payload) => ({ summary: summarize(payload) }),
}])

Non-plain values are cloned rather than returned by reference. Map entries, Set values, and enumerable Error properties are traversed recursively, so a sensitive string key such as token is masked there as well. Dates, regular expressions, buffers, array buffers, and typed views retain their value in the clone. Cyclic references are rendered as '[Circular]'.

The bundled JSON sinks encode values that native JSON.stringify() would collapse using tagged objects such as { "$type": "Map", "entries": [...] }, { "$type": "Set", "values": [...] }, and explicit RegExp, Error, binary, non-finite-number, and array-hole representations. The console and Postgres sinks therefore preserve the same information after redaction.


The behavior itself doesn’t know who the caller is — resolve it from trusted session context or pipeline context.

Configure a trusted actor factory once in AuditModule.forRoot({ defaults: { actor: ... } }) so that handlers don’t need to duplicate actor resolution:

AuditModule.forRoot({
defaults: {
actor: (ctx) => {
const principal = getSessionPrincipal();
if (!principal) return { authenticated: false };
return {
id: principal.id,
authenticated: true,
principalType: principal.type,
email: principal.email,
};
},
},
});

Security requirements for actor resolution

Section titled “Security requirements for actor resolution”
  • Never trust caller-supplied request body fields: An actor factory must resolve the principal from trusted session context, tokens, or context.items set by an upstream authentication behavior/interceptor. Using a request body field (like req.email) allows unverified identities to pollute the audit trail.
  • Explicit unauthenticated status: When no principal is in scope, return { authenticated: false } (or omit id) rather than fabricating an 'anonymous' identity that could be misread as an authenticated principal.
  • Per-handler override: Individual handlers can still override the actor factory if a specific operation has different principal semantics.

Every record carries correlationId and, when the pipeline has one, tenantId, both read from the pipeline context. The pipeline takes them from the sources configured on PipelineModule.forRoot():

import { Module } from '@nestjs/common';
import { PipelineModule } from '@nestjs-pipeline/core';
import { correlationSource } from '@nestjs-pipeline/correlation';
import { tenantSource } from '@nestjs-pipeline/tenant';
import { AuditBehavior, AuditModule } from '@nestjs-pipeline/audit';
@Module({
imports: [
AuditModule.forRoot(),
PipelineModule.forRoot({
sources: { tenantId: tenantSource, correlationId: correlationSource },
globalBehaviors: { scope: 'all', before: [AuditBehavior] },
}),
],
})
export class AppModule {}

The audit package imports neither @nestjs-pipeline/tenant nor @nestjs-pipeline/correlation. The tenant is also merged into metadata as tenantId; PostgresAuditSink has no tenant column and stores it there (metadata->>'tenantId').


When the sink itself throws (e.g. the audit DB is down):

  • failOpen: true (default) — the failure is logged as a warning and ignored. A successful handler still returns its response, and if the handler had failed, its original error remains the error seen by the caller. Favors availability.
  • failOpen: false — strictly enforces audit persistence:
    • If the handler succeeded, the sink error is propagated, rejecting the request because the required audit trail could not be recorded.
    • If the handler had already failed, AuditBehavior re-throws the handler’s own error, unchanged, and logs the sink error. The error object is never modified.

Record construction and sink failures both follow failOpen. When handling an already failed request, the original request error is preserved.

The start record follows the same rule, one step earlier: when building it or sink.begin fails, failOpen: true logs the failure and runs the handler, and failOpen: false rejects the request before the handler runs.

Diagnostic logging of an audit failure is itself fail-open: a logger that throws never replaces the handler’s result or error.

actor, metadata and redact must be functions. A non-function value fails application bootstrap with a PipelineConfigurationError. Module defaults of a request-scoped AuditBehavior have no instance at bootstrap; an invalid default there is rejected before the handler runs (logged and ignored with failOpen: true).


Export Kind Description
AuditBehavior class The pipeline behavior
audit fn Type-safe intent builder returning [AuditBehavior, options]
AuditIntentOptions type Options for audit(...)
AuditModule class forRoot / forRootAsync registration
AUDIT_RECORD_ITEM symbol context.items exported unique Symbol key holding the produced record
AUDIT_RECORD_ITEM_TOKEN PipelineItemToken<AuditRecord> Typed token over the same key, for getPipelineItem
AUDIT_START_RECORD_ITEM_TOKEN PipelineItemToken<AuditStartRecord> The pending start record, set before the handler runs
AUDIT_SINK / AUDIT_DEFAULT_OPTIONS token DI tokens
AUDIT_SEVERITY / AUDIT_OUTCOMES / AUDIT_REQUEST_KINDS const Named values for severities, outcomes (including pending) and request kinds
LogAuditSink class Default zero-dep sink
PostgresAuditSink class Postgres drop-in sink
createAuditTableSql fn CREATE TABLE DDL for the Postgres sink
buildAuditRecord / buildAuditStartRecord fn Pure builders of the final and the pending record (used by the behavior)
redactValue / DEFAULT_REDACT_KEYS / REDACTED fn/const Redaction helpers
AuditSink, AuditRecord, AuditStartRecord, AuditBehaviorOptions, AuditModuleOptions, AuditModuleAsyncOptions, LogAuditSinkOptions, PostgresAuditSinkOptions, PostgresQueryableLike, BuildAuditRecordInput, … type Public types

Dual-licensed under AGPL-3.0-or-later or a Commercial License. See LICENSE and COMMERCIAL_LICENSE.txt at the repository root.