Skip to content
58 changes: 54 additions & 4 deletions packages/agent/src/audit-trail/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -182,10 +182,12 @@ refuses an over-cap bulk write.

## HTTP routes

When `auditTrail` is set, the agent exposes three routes (all behind Forest's auth, gated by
`assertCanRead` on the target collection). When the caller's role has a record-level scope on that
collection, every route below — including the correlation lookups — additionally requires the
target id to currently exist and match that scope. A scope can't be evaluated against a record that
When `auditTrail` is set, the agent exposes four routes, all behind Forest's auth. The three
record-scoped ones are gated by `assertCanRead` on the collection they name; the cross-collection
timeline has no target collection, so it checks `canRead` per collection and queries only the ones
that pass. When the caller's role has a record-level scope on that collection, every record-scoped
route below — including the correlation lookups — additionally requires the target id to currently
exist and match that scope. A scope can't be evaluated against a record that
no longer exists, so this only refuses a still-existing, out-of-scope id: once a record is genuinely
deleted, anyone who can read the collection can see its history, reconstructed state, or correlated
operations, scope aside — inspecting what was deleted is much of the point of an audit trail.
Expand Down Expand Up @@ -334,6 +336,54 @@ to `20`, capped at `100`. Out-of-bound or non-numeric values fall back to the de
Ties on equal timestamps fall back to insertion order (the auto-increment `id`), so the order is
deterministic and stable across pages whatever the direction and filters.

### `GET /forest/_audit-trail` — cross-collection timeline

The project-level feed: every collection at once, newest first, for the page that shows what
happened across the project rather than to one record. Each row carries its own `collection` and
`recordId`, so the client can name and link the record it belongs to.

```json
{ "data": [ /* rows */ ], "meta": { "cursor": { "before": "…", "excludeIds": [12, 11] } } }
```

**Detail values follow the same rule as the per-record routes:** a caller who can read a collection
gets its rows with `previousValues` and `newValues`, whatever their permission level. Someone who can
read the collection can already read each record's history one at a time, so the feed serves the same
values in aggregate. Who can open the project's Activity Logs at all is decided per team by Forest
(the `activityLogs` feature a team can be deactivated for), not by this route. The capabilities
payload advertises the route as `canUseAuditTrailTimeline`, true only when the store implements
`listTimeline`.

**Permissions.** Only the collections the caller can read are queried. A collection the caller sees
only through a **record-level scope is left out entirely**: a scope can't be evaluated across a whole
timeline without fetching every record it mentions, and the row alone would already reveal that a
record the caller cannot read exists and was touched. A caller scoped on every collection therefore
gets an empty timeline.

**Filters** are the per-record route's, minus `fields` (a column name means nothing across
collections): `userIds`, `operations`, `startDate` / `endDate` / `timezone`, `search`, parsed exactly
the same way — including the 400 on an unrecognized `operation`.

**Paging is by cursor, not offset.** A project-wide feed keeps growing at the head, so an offset
silently shifts rows across pages. `page[size]` still applies (default `20`, capped at `100`), and
`meta.cursor` carries the next page's position — `null` once the last row has been served:

| query param | format | effect |
| ------------- | ------------------------------- | --------------------------------------------------- |
| `before` | ISO instant | **inclusive** upper bound on `timestamp` |
| `excludeIds` | comma-separated integers | row ids already served at that boundary |

The bound is inclusive on purpose: rows sharing the boundary timestamp must not fall between two
pages, so they come back and `excludeIds` drops the ones already sent. Ties order by descending `id`
(insertion order), and a timestamp holding more rows than fit on one page accumulates its exclusions
across pages rather than looping. Feed the two values back verbatim; don't synthesize them.

`meta.cursor` is also `null` if the walk cannot move past the current timestamp — staying on one
must always add at least the last row's id to the exclusions, so a cursor that would come back
unchanged ends the walk instead of repeating the page forever.

A custom `AuditStore` that doesn't implement `listTimeline` simply doesn't get this route mounted.

### `GET /forest/_audit-trail/correlation/{correlationKey}`

Returns `{ data: AuditRecord[] }` — the operation(s) recorded under one `correlationKey` for a
Expand Down
42 changes: 42 additions & 0 deletions packages/agent/src/audit-trail/migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,12 @@ function qualifiedMigrationName(schema: string | undefined, tableName: string):
: `${tableName}:001-create-audit-logs`;
}

function qualifiedMigrationName002(schema: string | undefined, tableName: string): string {
return schema
? `${schema}.${tableName}:002-index-timestamp-id`
: `${tableName}:002-index-timestamp-id`;
}

// Every table name claimed by an audit-trail store configured in this process, keyed by
// `schema\0name` — both its own data table and the migration table Umzug derives from it.
// A second store's data table can otherwise land on the exact name a first store's migration
Expand Down Expand Up @@ -198,6 +204,42 @@ function buildMigrations(schema: string | undefined, tableName: string) {
);
},
},
{
// Its own migration rather than an index added to 001: the table has shipped, so a database
// out there has already recorded 001 as applied. The timeline route reads across every
// record, ordered and paged by (timestamp, id): without this, each page sorts the whole table.
name: qualifiedMigrationName002(schema, tableName),
up: async ({ context }: { context: MigrationContext }) => {
const table = { tableName: context.tableName, schema: context.schema };
const name = `${context.tableName}_timestamp_id`;
const hasIndex = () =>
indexNames(context.queryInterface, table, context.transaction).then(names =>
names.has(name),
);

// Idempotent for the same reason 001 is: a process losing a concurrent-boot race retries.
if (await hasIndex()) return;

try {
await context.queryInterface.addIndex(table, {
fields: ['timestamp', 'id'],
name,
transaction: context.transaction,
});
} catch (error) {
// Without Postgres's advisory lock, another agent booting at the same moment can add the
// index between the check above and this statement. Anything else is a real failure.
if (!(await hasIndex().catch(() => false))) throw error;
}
},
down: async ({ context }: { context: MigrationContext }) => {
await context.queryInterface.removeIndex(
{ tableName: context.tableName, schema: context.schema },
`${context.tableName}_timestamp_id`,
{ transaction: context.transaction },
);
},
},
];
}

Expand Down
130 changes: 130 additions & 0 deletions packages/agent/src/audit-trail/query-params.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
import type { AuditOperation } from './types';

import { ValidationError } from '@forestadmin/datasource-toolkit';
import { DateTime } from 'luxon';

const DATE_ONLY = /^\d{4}-\d{2}-\d{2}$/;
const DATE_TIME = /^(\d{4}-\d{2}-\d{2})[T ](\d{2}):(\d{2})(?::(\d{2}))?$/;
// ISO 8601 instant: carries its own timezone designator (`Z` or `±HH:mm` / `±HHMM`).
const ISO_INSTANT = /[Zz]$|[+-]\d{2}:?\d{2}$/;

export const DEFAULT_PAGE_SIZE = 20;
export const MAX_PAGE_SIZE = 100;

export const AUDIT_OPERATIONS: readonly AuditOperation[] = [
'create',
'update',
'delete',
'action',
'action_failed',
];

// Comma-separated integer ids; non-numeric tokens are dropped. Empty after parsing → no filter.
export function parseUserIds(raw?: string): number[] | undefined {
if (!raw) return undefined;

const ids = raw
.split(',')
.map(token => token.trim())
.filter(token => /^\d+$/.test(token))
.map(token => Number.parseInt(token, 10));

return ids.length > 0 ? ids : undefined;
}

// Comma-separated, from a closed set. An unrecognized value is rejected rather than dropped: a
// silently ignored filter returns unfiltered rows into a list the caller believes is filtered,
// which is worse than an error.
export function parseOperations(raw?: string): AuditOperation[] | undefined {
if (!raw) return undefined;

const tokens = raw
.split(',')
.map(token => token.trim())
.filter(token => token.length > 0);

const unknown = tokens.find(token => !AUDIT_OPERATIONS.includes(token as AuditOperation));

if (unknown) {
throw new ValidationError(
`Invalid operation: "${unknown}" (expected one of ${AUDIT_OPERATIONS.join(', ')})`,
);
}

return tokens.length > 0 ? (tokens as AuditOperation[]) : undefined;
}

export function parseFields(raw?: string): string[] | undefined {
if (!raw) return undefined;

const fields = raw
.split(',')
.map(token => token.trim())
.filter(token => token.length > 0);

return fields.length > 0 ? fields : undefined;
}

// Trimmed; empty after trimming is treated the same as absent.
export function parseSearch(raw?: string): string | undefined {
const trimmed = raw?.trim();

return trimmed || undefined;
}

// 1-based `page[size]` (default 20, capped at 100). Invalid values fall back to the default rather
// than erroring.
export function parsePageSize(raw?: string): number {
const size = Number.parseInt(raw ?? '', 10);

if (Number.isNaN(size) || size < 1) return DEFAULT_PAGE_SIZE;

return Math.min(size, MAX_PAGE_SIZE);
}

function toLocalInstant(raw: string, timezone: string, boundary: 'start' | 'end'): DateTime {
// An embedded offset already pins the instant — the request timezone and start/end boundary
// don't apply.
if (ISO_INSTANT.test(raw)) return DateTime.fromISO(raw, { setZone: true });

if (DATE_ONLY.test(raw)) {
const day = DateTime.fromISO(raw, { zone: timezone });

return boundary === 'end' ? day.endOf('day') : day.startOf('day');
}

const match = DATE_TIME.exec(raw);

if (!match) return DateTime.invalid('unparsable');

const [, date, hours, minutes, seconds] = match;
const base = DateTime.fromISO(`${date}T${hours}:${minutes}`, { zone: timezone });

if (seconds !== undefined) return base.set({ second: Number(seconds), millisecond: 0 });

// Minutes-only: end snaps to :59.999 to stay inclusive; start stays at :00.000.
return boundary === 'end' ? base.set({ second: 59, millisecond: 999 }) : base;
}

// Bare day (`YYYY-MM-DD`) or wall-clock datetime (`YYYY-MM-DD[T| ]HH:mm[:ss]`), interpreted as
// local time in the request timezone and returned as a UTC instant so the store can compare it
// to stored timestamps.
export function parseDateBoundary(
raw: string | undefined,
timezone: string,
boundary: 'start' | 'end',
): string | undefined {
if (!raw) return undefined;

const instant = toLocalInstant(raw, timezone, boundary);

if (!instant.isValid) {
throw new ValidationError(
instant.invalidReason === 'unsupported zone'
? `Invalid timezone: "${timezone}"`
: `Invalid date: "${raw}" (expected YYYY-MM-DD, YYYY-MM-DDTHH:mm, or an ISO 8601 instant)`,
);
}

return instant.toUTC().toISO() ?? undefined;
}
54 changes: 54 additions & 0 deletions packages/agent/src/audit-trail/sql-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type {
AuditStatus,
AuditStorageOptions,
AuditStore,
AuditTimelineQuery,
PendingAuditRecord,
} from './types';
import type { Model, ModelStatic } from 'sequelize';
Expand Down Expand Up @@ -232,6 +233,42 @@ function buildHistoryWhereClause(
return where;
}

function buildTimelineWhereClause(
{
collections,
before,
excludeIds,
userIds,
operations,
startTimestamp,
endTimestamp,
search,
}: AuditTimelineQuery,
sequelize: Sequelize,
): Record<string | symbol, unknown> {
const where: Record<string | symbol, unknown> = { collection: { [Op.in]: collections } };

if (userIds) where.userId = { [Op.in]: userIds };
if (operations?.length) where.operation = { [Op.in]: operations };
if (excludeIds?.length) where.id = { [Op.notIn]: excludeIds };

const timestampRange: Record<symbol, Date> = {};
if (startTimestamp) timestampRange[Op.gte] = new Date(startTimestamp);

// The cursor bound and the `endDate` filter are both inclusive upper bounds; the tighter wins.
const upperBounds = [endTimestamp, before].filter(Boolean).map(bound => new Date(bound));

if (upperBounds.length) {
timestampRange[Op.lte] = new Date(Math.min(...upperBounds.map(date => date.getTime())));
}

if (Object.getOwnPropertySymbols(timestampRange).length) where.timestamp = timestampRange;

if (search) where[Op.and] = [searchCondition(sequelize, search)];

return where;
}

export function fromRow(row: Model): AuditRecord {
const plain = row.get({ plain: true }) as Record<string, unknown>;
const { timestamp } = plain;
Expand Down Expand Up @@ -363,6 +400,23 @@ export function createSqlAuditStore(options: AuditStorageOptions): {
transaction: null,
});
},
async listTimeline(query) {
if (!query.collections.length) return [];

const { model, connection } = await init();

const rows = await model.findAll({
where: buildTimelineWhereClause(query, connection),
order: [
['timestamp', 'DESC'],
['id', 'DESC'],
],
limit: query.limit,
transaction: null,
});

return rows.map(fromRow);
},
async listDistinctUsers(query) {
const { model, connection } = await init();

Expand Down
29 changes: 29 additions & 0 deletions packages/agent/src/audit-trail/types.ts
Comment thread
bexchauveto marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,30 @@ export type AuditHistoryQuery = {
after?: Pick<AuditRecord, 'timestamp' | 'id'>;
};

/**
* Cross-collection timeline, newest first. Paged by cursor rather than by offset: a project-wide
* feed keeps growing at the head, so an offset silently shifts rows across pages.
*/
export type AuditTimelineQuery = {
/** Collections the caller is allowed to see; an empty list must match nothing. */
collections: string[];
limit: number;
/**
* Inclusive upper bound on `timestamp`, from the previous page's cursor. Inclusive so rows
* sharing the boundary timestamp aren't skipped; `excludeIds` drops the ones already returned.
*/
before?: string;
excludeIds?: number[];
userIds?: number[];
operations?: AuditOperation[];
/** Inclusive lower bound on `timestamp` as a UTC ISO instant. */
startTimestamp?: string;
/** Inclusive upper bound on `timestamp` as a UTC ISO instant. */
endTimestamp?: string;
/** Same free-text match as `AuditHistoryQuery.search`. */
search?: string;
};

export type AuditCorrelationQuery = {
collection: string;
recordId: string;
Expand Down Expand Up @@ -111,6 +135,11 @@ export interface AuditStore {
listByCorrelation(query: AuditCorrelationQuery): AuditRecord[] | Promise<AuditRecord[]>;
/** Flat list of entries recorded under any of `correlationKeys` for a record, oldest first. */
listByCorrelations(query: AuditCorrelationsQuery): AuditRecord[] | Promise<AuditRecord[]>;
/**
* Cross-collection timeline, newest first (ties broken by descending `id`). Optional: a store
* written before this existed simply doesn't serve the project-level route.
*/
listTimeline?(query: AuditTimelineQuery): AuditRecord[] | Promise<AuditRecord[]>;
/** Distinct authors matching the query filters, independent of pagination. */
listDistinctUsers(
query: Omit<AuditHistoryQuery, 'skip' | 'limit' | 'order'>,
Expand Down
Loading
Loading