Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions packages/agent/src/audit-trail/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ The `forest.audit_logs` table has one row per audited change:
| `operation` | `create` / `update` / `delete` / `action` / `action_failed` |
| `collection` | audited collection name |
| `record_id` | packed record id (primary keys joined with `\|`); `null` for a `create` row still `pending` (the record's id isn't assigned yet) |
| `previous_record_id` | the id the row was filed under before a confirmed `update`, set whether or not the key moved; `null` on every other operation and on any row written before this column existed. Internal: never served to a client |
| `user_id` | id of the Forest user who made the change |
| `user_first_name` | that user's first name, denormalized at write time |
| `user_last_name` | that user's last name, denormalized at write time |
Expand Down Expand Up @@ -382,7 +383,14 @@ across pages rather than looping. Feed the two values back verbatim; don't synth
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.

**Authors.** On the first page only (no `before`), `meta` also carries `availableUsers`: the
distinct authors matching the active filters across the collections queried, independent of the
cursor, in the per-record route's shape. Later pages omit the key rather than send `[]`, so a client
keeps the list it already saw.

A custom `AuditStore` that doesn't implement `listTimeline` simply doesn't get this route mounted.
One that implements `listTimeline` without `listTimelineUsers` serves the rows with no
`availableUsers`.

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

Expand Down
4 changes: 4 additions & 0 deletions packages/agent/src/audit-trail/instrument.ts
Original file line number Diff line number Diff line change
Expand Up @@ -480,6 +480,10 @@ function instrumentCollection(
return recorder.confirm(pendingId, {
operation: 'update',
recordId: toPackedRecordId(updated, primaryKeys),
// Written whether or not the key moved. Only a value here lets the route trust that the
// previous side's id is known: a null has to keep meaning "this row predates the column",
// or an old row whose key did move would be judged by the id it moved to.
previousRecordId: toPackedRecordId(record, primaryKeys),
previousValues: redactValues(previousValues, redactedFields),
newValues: redactValues(newValues, redactedFields),
});
Expand Down
44 changes: 44 additions & 0 deletions packages/agent/src/audit-trail/migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,12 @@ function qualifiedMigrationName002(schema: string | undefined, tableName: string
: `${tableName}:002-index-timestamp-id`;
}

function qualifiedMigrationName003(schema: string | undefined, tableName: string): string {
return schema
? `${schema}.${tableName}:003-add-previous-record-id`
: `${tableName}:003-add-previous-record-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 @@ -240,6 +246,44 @@ function buildMigrations(schema: string | undefined, tableName: string) {
);
},
},
{
// Its own migration rather than a column added to 001: the table has shipped, so a database
// out there has already recorded 001 as applied and would never see the edit.
name: qualifiedMigrationName003(schema, tableName),
up: async ({ context }: { context: MigrationContext }) => {
const table = { tableName: context.tableName, schema: context.schema };
const existing = await columnNames(context.queryInterface, table, context.transaction);

// Idempotent for the same reason 001 is: a process losing a concurrent-boot race retries.
if (existing.has('previous_record_id')) return;

try {
await context.queryInterface.addColumn(
table,
'previous_record_id',
{ type: DataTypes.TEXT, allowNull: true },
{ transaction: context.transaction },
);
} catch (error) {
// Without Postgres's advisory lock, another agent booting at the same moment can add the
// column between the check above and this statement. Its duplicate-column error then
// means the work is done. Anything else, or a column still missing, is a real failure.
const added = await columnNames(context.queryInterface, table, context.transaction).then(
columns => columns.has('previous_record_id'),
() => false,
);

if (!added) throw error;
}
},
down: async ({ context }: { context: MigrationContext }) => {
await context.queryInterface.removeColumn(
{ tableName: context.tableName, schema: context.schema },
'previous_record_id',
{ transaction: context.transaction },
);
},
},
];
}

Expand Down
76 changes: 49 additions & 27 deletions packages/agent/src/audit-trail/sql-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import type {
AuditStorageOptions,
AuditStore,
AuditTimelineQuery,
AuditUserSummary,
PendingAuditRecord,
} from './types';
import type { Model, ModelStatic } from 'sequelize';
Expand Down Expand Up @@ -38,6 +39,9 @@ export function defineAuditLogModel(
// migration creates the column as such — this must match. Nullable: a pending create's row
// has no id yet, since the record doesn't exist until the write resolves.
recordId: { type: DataTypes.TEXT, allowNull: true },
// Set on every confirmed update, so a null distinguishes a row older than the column from
// one whose key held still. TEXT for the same reason as `recordId`.
previousRecordId: { type: DataTypes.TEXT, allowNull: true },
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
userId: { type: DataTypes.INTEGER, allowNull: true },
// Denormalised from the caller at write time — who acted then, not who holds that id today.
userFirstName: { type: DataTypes.TEXT, allowNull: true },
Expand Down Expand Up @@ -79,6 +83,9 @@ export function toRow(
operation: record.operation,
collection: record.collection,
recordId: record.recordId,
// Optional at insert, but kept when given, as the in-memory store keeps it: a later `confirm`
// that omits it must not read as a row older than the column.
previousRecordId: record.previousRecordId ?? null,
userId: record.userId,
userFirstName: record.userFirstName,
userLastName: record.userLastName,
Expand Down Expand Up @@ -243,7 +250,7 @@ function buildTimelineWhereClause(
startTimestamp,
endTimestamp,
search,
}: AuditTimelineQuery,
}: Omit<AuditTimelineQuery, 'limit'>,
sequelize: Sequelize,
): Record<string | symbol, unknown> {
const where: Record<string | symbol, unknown> = { collection: { [Op.in]: collections } };
Expand All @@ -269,6 +276,38 @@ function buildTimelineWhereClause(
return where;
}

async function listAuthors(
model: ModelStatic<Model>,
where: Record<string | symbol, unknown>,
): Promise<AuditUserSummary[]> {
// MAX() rather than a bare column: grouping by `user_id` alone is invalid in strict SQL
// unless every selected column is either grouped or aggregated.
const rows = (await model.findAll({
where,
attributes: [
'userId',
[Sequelize.fn('MAX', Sequelize.col('user_first_name')), 'userFirstName'],
[Sequelize.fn('MAX', Sequelize.col('user_last_name')), 'userLastName'],
[Sequelize.fn('MAX', Sequelize.col('user_email')), 'userEmail'],
],
group: ['userId'],
raw: true,
transaction: null,
})) as unknown as Array<{
userId: number;
userFirstName: string | null;
userLastName: string | null;
userEmail: string | null;
}>;

return rows.map(row => ({
id: row.userId,
firstName: row.userFirstName ?? null,
lastName: row.userLastName ?? null,
email: row.userEmail ?? null,
}));
}

export function fromRow(row: Model): AuditRecord {
const plain = row.get({ plain: true }) as Record<string, unknown>;
const { timestamp } = plain;
Expand All @@ -279,6 +318,7 @@ export function fromRow(row: Model): AuditRecord {
operation: plain.operation as AuditRecord['operation'],
collection: plain.collection as string,
recordId: (plain.recordId as string) ?? null,
previousRecordId: (plain.previousRecordId as string) ?? null,
userId: plain.userId as number,
userFirstName: (plain.userFirstName as string) ?? null,
userLastName: (plain.userLastName as string) ?? null,
Expand Down Expand Up @@ -420,32 +460,14 @@ export function createSqlAuditStore(options: AuditStorageOptions): {
async listDistinctUsers(query) {
const { model, connection } = await init();

// MAX() rather than a bare column: grouping by `user_id` alone is invalid in strict SQL
// unless every selected column is either grouped or aggregated.
const rows = (await model.findAll({
where: buildHistoryWhereClause(query as AuditHistoryQuery, connection),
attributes: [
'userId',
[Sequelize.fn('MAX', Sequelize.col('user_first_name')), 'userFirstName'],
[Sequelize.fn('MAX', Sequelize.col('user_last_name')), 'userLastName'],
[Sequelize.fn('MAX', Sequelize.col('user_email')), 'userEmail'],
],
group: ['userId'],
raw: true,
transaction: null,
})) as unknown as Array<{
userId: number;
userFirstName: string | null;
userLastName: string | null;
userEmail: string | null;
}>;

return rows.map(row => ({
id: row.userId,
firstName: row.userFirstName ?? null,
lastName: row.userLastName ?? null,
email: row.userEmail ?? null,
}));
return listAuthors(model, buildHistoryWhereClause(query as AuditHistoryQuery, connection));
},
async listTimelineUsers(query) {
if (!query.collections.length) return [];

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

return listAuthors(model, buildTimelineWhereClause(query, connection));
},
async listByCorrelation({ collection, recordId, correlationKey }) {
const { model } = await init();
Expand Down
26 changes: 23 additions & 3 deletions packages/agent/src/audit-trail/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,13 @@ export type AuditRecord = {
collection: string;
/** Null for a pending create — the record's primary key isn't assigned yet. */
recordId: string | null;
/**
* The id this record was filed under before a confirmed update, written whether or not the key
* moved. Null therefore means "written before this column existed", not "the key held still": a
* row from an earlier agent cannot claim its id answers for the previous side of an update.
* Internal — it is how the agent follows a record across a rename, never served to a client.
*/
previousRecordId: string | null;
userId: number;
/** Denormalised from the caller at write time: who acted then, not who holds that id today. */
userFirstName: string | null;
Expand All @@ -30,13 +37,18 @@ export type AuditRecord = {
};

/** The subset known before the write runs, when the pending row is first inserted. */
export type PendingAuditRecord = Omit<AuditRecord, 'id' | 'status'>;
export type PendingAuditRecord = Omit<AuditRecord, 'id' | 'status' | 'previousRecordId'> &
Partial<Pick<AuditRecord, 'previousRecordId'>>;

/** What `confirm` updates once the write (or action) has resolved. */
/**
* What `confirm` updates once the write (or action) has resolved. `previousRecordId` is optional
* because only an update has a previous side to file under an id of its own.
*/
export type AuditRecordConfirmation = Pick<
AuditRecord,
'operation' | 'recordId' | 'previousValues' | 'newValues'
>;
> &
Partial<Pick<AuditRecord, 'previousRecordId'>>;

export type AuditHistoryQuery = {
collection: string;
Expand Down Expand Up @@ -140,6 +152,14 @@ export interface AuditStore {
* written before this existed simply doesn't serve the project-level route.
*/
listTimeline?(query: AuditTimelineQuery): AuditRecord[] | Promise<AuditRecord[]>;
/**
* Distinct authors matching the timeline's filters, independent of its cursor. Optional like
* `listTimeline`: without it the route serves the rows but no author list. An empty
* `collections` must match nothing.
*/
listTimelineUsers?(
query: Omit<AuditTimelineQuery, 'limit' | 'before' | 'excludeIds'>,
): AuditUserSummary[] | Promise<AuditUserSummary[]>;
/** Distinct authors matching the query filters, independent of pagination. */
listDistinctUsers(
query: Omit<AuditHistoryQuery, 'skip' | 'limit' | 'order'>,
Expand Down
38 changes: 28 additions & 10 deletions packages/agent/src/audit-trail/withhold.ts
Original file line number Diff line number Diff line change
Expand Up @@ -195,22 +195,40 @@ export default function withholdOutsidePermissionScope(
return entries.map(entry => {
if (entry.operation === 'action' || entry.operation === 'action_failed') return entry;

// Decoded once per row rather than once per side: the two sides read the same id, and a row
// whose id no longer decodes should say so once.
const decoded = decodePrimaryKeys(entry.recordId, withholding.collection, withholding.logger);

const side = (values: Record<string, unknown>, idAnswersForSide: IdAnswersForSide) =>
const side = (
values: Record<string, unknown>,
decoded: ReturnType<typeof decodePrimaryKeys>,
idAnswersForSide: IdAnswersForSide,
) =>
permissionScopeAccepts(answerableSnapshot(values, decoded, idAnswersForSide), withholding)
? values ?? {}
: {};

// An update that recorded where it came from is judged side by side: the previous state
// against the id it was filed under then, the new state against the id it ended up with. A row
// without that column is older than it, so its id still answers for the new side alone.
// `?? null` first: the column is optional on the write types, so an absent one has to read as
// unknown exactly like a null, or a row that never recorded it would answer from the id it
// ended up with.
const previousRecordId = entry.previousRecordId ?? null;
const previousIsKnown = entry.operation !== 'update' || previousRecordId !== null;

// Decoded once per distinct id, not once per side: a row whose id no longer decodes should say
// so once, and the two sides read the same id unless the key actually moved.
const decoded = decodePrimaryKeys(entry.recordId, withholding.collection, withholding.logger);
const decodedPrevious =
previousRecordId === null
? decoded
: decodePrimaryKeys(previousRecordId, withholding.collection, withholding.logger);

// A pending update is filed under the id the record had before the write, which says nothing
// about the state it was moving to, so the new side answers only with what it captured.
const newIsKnown = entry.status !== 'pending';

return {
...entry,
// The row is filed under the identity the record ended up with, so its id answers for the
// new side of an update and for a create or a delete — never for what an update moved away
// from.
previousValues: side(entry.previousValues, entry.operation !== 'update'),
newValues: side(entry.newValues, true),
previousValues: side(entry.previousValues, decodedPrevious, previousIsKnown),
newValues: side(entry.newValues, decoded, newIsKnown),
};
});
}
18 changes: 14 additions & 4 deletions packages/agent/src/routes/access/audit-trail-timeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,16 +54,26 @@ export default class AuditTrailTimelineRoute extends BaseRoute {
const collections = await this.readableCollections(context);
const { store } = this.options.auditTrail;

// Authors on the first page only, like the per-record route and Forest's activity-logs route:
// they do not change along the walk, and later pages omit the key rather than send `[]`.
// One row over the page: the only way to know whether a further page exists without a count.
const fetched = collections.length
? await store.listTimeline({ ...filters, ...cursor, collections, limit: limit + 1 })
: [];
const [fetched, availableUsers] = await Promise.all([
collections.length
? store.listTimeline({ ...filters, ...cursor, collections, limit: limit + 1 })
: [],
!cursor.before && store.listTimelineUsers
? store.listTimelineUsers({ ...filters, collections })
: undefined,
]);
const page = fetched.slice(0, limit);

context.response.body = {
data: page,
// `previousRecordId` stays out: it is how the agent follows a record across a rename, not
// something a client reads.
data: page.map(({ previousRecordId, ...served }) => served),
meta: {
cursor: fetched.length > limit ? AuditTrailTimelineRoute.nextCursor(page, cursor) : null,
...(availableUsers && { availableUsers }),
},
};
}
Expand Down
6 changes: 4 additions & 2 deletions packages/agent/src/routes/access/audit-trail.ts
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,9 @@ export default class AuditTrailRoute extends CollectionRoute {
});

context.response.body = {
data: matched.page,
// `previousRecordId` stays out: it is how the agent follows a record across a rename, not
// something a client reads.
data: matched.page.map(({ previousRecordId, ...served }) => served),
meta: {
count: matched.count,
...(isFirstFetch && { availableUsers: [...matched.authors.values()] }),
Expand Down Expand Up @@ -165,7 +167,7 @@ export default class AuditTrailRoute extends CollectionRoute {
permissionScope && gone ? this.withhold(rawData, permissionScope, context) : rawData;

context.response.body = {
data,
data: data.map(({ previousRecordId, ...served }) => served),
meta: { count, ...(availableUsers && { availableUsers }) },
};
}
Expand Down
1 change: 1 addition & 0 deletions packages/agent/test/audit-trail/in-memory-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ const record = (
operation: 'update',
collection: 'accounts',
recordId: '1',
previousRecordId: null,
userId: 1,
userFirstName: null,
userLastName: null,
Expand Down
Loading
Loading