Skip to content
Merged
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
19 changes: 19 additions & 0 deletions .changeset/record-change-payload-credential-mask.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
---
'@objectstack/objectql': patch
'@objectstack/plugin-approvals': patch
'@objectstack/plugin-webhooks': patch
'@objectstack/service-knowledge': minor
---

Record-change payloads apply the same credential mask and internal-field omission as write responses.

Clause-②: yes (widening)

- **`data.record.created` / `data.record.updated` events.** The engine projects the event's `after` and `changes` bodies through `omitInternalFieldsFromWriteResponse` (`@objectstack/core`), the helper every external write response already uses: credential-class fields (`secret`, and `password` outside the exempt `managedBy` buckets) carry `SECRET_MASK` (or `null` when unset), and `internal: true` fields are omitted. The engine's own write result is unchanged, so a privileged in-process caller that reads the stored value back off `insert` / `update` still sees it.
- **Approval request snapshot.** The record snapshot an approval request stores (`payload_json`) applies the same rule when the request is opened.
- **Outbound webhook body.** The delivered body, and the delivery row that stores it, apply the same rule to `before`, `after` and `changes`.
- **Knowledge index documents.** `recordToDocument` takes the object definition as an optional fourth argument and skips credential-class and `internal` fields, under `'*'` and when a source names one explicitly. `KnowledgeService` passes the definition from the bound engine.
- **New public surface of `@objectstack/service-knowledge` (additive):** `recordToDocument` accepts the object definition as an optional fourth argument; existing three-argument calls behave as before.
- **Receivers see masked values.** Webhook receivers and realtime clients now get `SECRET_MASK` (or `null` when unset) for credential-class fields and no key for `internal` fields.
- **Existing rows are not rewritten.** Approval snapshots, webhook delivery rows and knowledge documents written before this change keep their stored bodies; reindexing a knowledge source refreshes its documents.
- The audit trail already masked these fields and is unchanged. No other accept set or public schema changes.
144 changes: 144 additions & 0 deletions packages/objectql/src/engine-realtime-credential-mask.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license.

/**
* A `data.record.*` event is an external exit of the write it describes: its
* subscribers store or forward the record bodies (`after`, `changes`) to
* readers below the write boundary. So both bodies carry what a write
* RESPONSE carries — credential-class fields masked, `internal: true` fields
* omitted, by the one helper every write response uses — while the engine's
* own write result, read by privileged in-process callers, stays whole.
*/

import { describe, it, expect, beforeEach, vi } from 'vitest';
import type { IRealtimeService, RealtimeEventPayload } from '@objectstack/spec/contracts';
import { SECRET_MASK } from '@objectstack/spec/data';
import { ObjectQL } from './engine.js';

const CREDENTIAL = 'synthetic-credential-7c1e';
const INTERNAL = 'synthetic-internal-0b42';

const widget = {
name: 'widget',
label: 'Widget',
fields: {
id: { name: 'id', type: 'text' as const, primaryKey: true },
title: { name: 'title', type: 'text' as const },
passphrase: { name: 'passphrase', type: 'password' as const },
lookup_digest: { name: 'lookup_digest', type: 'text' as const, internal: true },
},
};

function makeStubDriver() {
const stores = new Map<string, Map<string, Record<string, unknown>>>();
const storeFor = (o: string) => {
let s = stores.get(o);
if (!s) { s = new Map(); stores.set(o, s); }
return s;
};
let nextId = 0;
const driver: any = {
name: 'memory', version: '0.0.0', supports: {},
async connect() {}, async disconnect() {}, async checkHealth() { return true; }, async execute() { return null; },
async find(o: string) { return Array.from(storeFor(o).values()); },
findStream() { throw new Error('ns'); },
async findOne(o: string, ast: any) {
const where = ast?.where ?? {};
for (const r of storeFor(o).values()) {
if (Object.entries(where).every(([k, v]) => k.startsWith('$') || (r[k] ?? null) === ((v as any)?.$eq ?? v ?? null))) return r;
}
return null;
},
async create(o: string, data: Record<string, unknown>) {
nextId += 1;
const id = (data.id as string) ?? `r_${nextId}`;
const row = { ...data, id };
storeFor(o).set(id, row);
return row;
},
async update(o: string, id: string, data: Record<string, unknown>) {
const s = storeFor(o); const cur = s.get(id);
if (!cur) throw new Error(`nf ${o}/${id}`);
const up = { ...cur, ...data, id }; s.set(id, up); return up;
},
async delete(o: string, id: string) { return storeFor(o).delete(id); },
async count(o: string) { return (await this.find(o)).length; },
async bulkCreate(o: string, rows: Record<string, unknown>[]) { return Promise.all(rows.map((r) => this.create(o, r))); },
async updateMany() { return 0; }, async deleteMany() { return 0; },
async bulkUpdate() { return []; }, async bulkDelete() {},
async upsert(o: string, data: Record<string, unknown>) { return this.create(o, data); },
async beginTransaction() { return { commit: async () => {}, rollback: async () => {} }; },
async commit() {}, async rollback() {},
};
return { driver };
}

describe('record-change event bodies apply the write-response non-exposure rules', () => {
let engine: ObjectQL;
let published: RealtimeEventPayload[];

beforeEach(async () => {
published = [];
const realtime: IRealtimeService = {
publish: vi.fn(async (event: RealtimeEventPayload) => { published.push(event); }),
subscribe: vi.fn(async () => 'sub-1'),
unsubscribe: vi.fn(async () => undefined),
};
engine = new ObjectQL();
const { driver } = makeStubDriver();
engine.registerDriver(driver, true);
await engine.init();
engine.registry.registerObject(widget);
engine.setRealtimeService(realtime);
vi.spyOn((engine as any).logger, 'warn').mockImplementation(() => undefined);
vi.spyOn((engine as any).logger, 'debug').mockImplementation(() => undefined);
});

it('a create event masks the credential field and omits the internal field from `after`', async () => {
await engine.insert('widget', { id: 'w1', title: 'One', passphrase: CREDENTIAL, lookup_digest: INTERNAL });
expect(published).toHaveLength(1);
const payload = published[0].payload as Record<string, any>;
expect(payload.type).toBe('data.record.created');
expect(payload.after.title).toBe('One');
expect(payload.after.passphrase).toBe(SECRET_MASK);
expect('lookup_digest' in payload.after).toBe(false);
expect(JSON.stringify(payload)).not.toContain(CREDENTIAL);
expect(JSON.stringify(payload)).not.toContain(INTERNAL);
});

it('an update event projects both `after` and `changes`', async () => {
await engine.insert('widget', { id: 'w2', title: 'Two', passphrase: 'first', lookup_digest: 'first' });
published.length = 0;
await engine.update('widget', { id: 'w2', title: 'Two b', passphrase: CREDENTIAL, lookup_digest: INTERNAL });
expect(published).toHaveLength(1);
const payload = published[0].payload as Record<string, any>;
expect(payload.type).toBe('data.record.updated');
expect(payload.changes.title).toBe('Two b');
expect(payload.changes.passphrase).toBe(SECRET_MASK);
expect('lookup_digest' in payload.changes).toBe(false);
expect(payload.after.passphrase).toBe(SECRET_MASK);
expect('lookup_digest' in payload.after).toBe(false);
expect(JSON.stringify(payload)).not.toContain(CREDENTIAL);
expect(JSON.stringify(payload)).not.toContain(INTERNAL);
});

it('an unset credential field rides the event as null, not as the mask', async () => {
await engine.insert('widget', { id: 'w3', title: 'Three', passphrase: null });
const payload = published[0].payload as Record<string, any>;
expect(payload.after.passphrase).toBeNull();
});

it("the engine's own write result is unchanged for the privileged in-process caller", async () => {
const created = await engine.insert('widget', {
id: 'w4', title: 'Four', passphrase: CREDENTIAL, lookup_digest: INTERNAL,
}) as Record<string, unknown>;
expect(created.passphrase).toBe(CREDENTIAL);
expect(created.lookup_digest).toBe(INTERNAL);

const updated = await engine.update('widget', { id: 'w4', lookup_digest: `${INTERNAL}-b` }) as Record<string, unknown>;
expect(updated.passphrase).toBe(CREDENTIAL);
expect(updated.lookup_digest).toBe(`${INTERNAL}-b`);
// ...while the events those writes published carried neither value.
expect(JSON.stringify(published.map((e) => e.payload))).not.toContain(CREDENTIAL);
expect(JSON.stringify(published.map((e) => e.payload))).not.toContain(INTERNAL);
});
});
50 changes: 44 additions & 6 deletions packages/objectql/src/engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,11 @@ import {
// The data door's object-existence 404, shared for the same reason: an
// in-process verb refuses a name the registry does not resolve with it.
objectNotFoundError,
// The one write-response non-exposure helper (credential mask + `internal`
// omission); a record-change event body is an external exit of the same
// write, so it is projected through the same rule — see
// {@link projectEventRecordBody}.
omitInternalFieldsFromWriteResponse,
} from '@objectstack/core';
import { WriteEpoch, isWriteEpochOperation } from './write-epoch.js';
import { bridgeAuthzInvalidation } from './authz-invalidation-bridge.js';
Expand Down Expand Up @@ -3365,6 +3370,33 @@ function redactEventMetadataBody(
return rest;
}

/**
* Project a `data.record.*` event body (`after` / `changes`) through the
* generic-data-path non-exposure rules — credential-class fields MASKED,
* `internal: true` fields OMITTED — by the ONE helper every external write
* response already goes through (`omitInternalFieldsFromWriteResponse`,
* `@objectstack/core`). ⛔ Never a second copy of that rule here.
*
* The event is an external exit of the write: its subscribers (outbound
* deliveries, search indexes, flows that snapshot the record) store or forward
* what they receive to readers below the write boundary, so the body carries
* what a read of the same row would answer — never the stored value the
* engine's own write result keeps whole for the privileged in-process caller.
*
* Pure with respect to the write result: the helper deletes and assigns in
* place, so it runs on a fresh shallow copy and the caller's record (the
* engine's return value) is never touched.
*/
function projectEventRecordBody(
schema: unknown,
body: Record<string, unknown> | undefined,
): Record<string, unknown> | undefined {
if (body === undefined) return body;
const copy = { ...body };
omitInternalFieldsFromWriteResponse(schema, copy);
return copy;
}

/** `DataEvent.userId` — the acting user, when the execution context names one. */
function eventUserId(execCtx?: ExecutionContext): string | undefined {
const userId = execCtx?.userId;
Expand Down Expand Up @@ -7572,17 +7604,23 @@ export class ObjectQL implements IObjectQLEngine {

try {
const timestamp = new Date().toISOString();
const changes = redactEventMetadataBody(object, eventRecordBody(input.changes), input.after);
const after = redactEventMetadataBody(object, eventRecordBody(input.after), input.after);
const schema = this._registry.getObject(object);
// Credential mask + `internal` omission on both bodies, the same rule
// the write response gets ({@link projectEventRecordBody}).
const changes = projectEventRecordBody(
schema,
redactEventMetadataBody(object, eventRecordBody(input.changes), input.after),
);
const after = projectEventRecordBody(
schema,
redactEventMetadataBody(object, eventRecordBody(input.after), input.after),
);
const userId = eventUserId(input.context);
// [#14970] The RECORD's organization, off the row itself — ⛔ never
// `input.context.tenantId`, which is the CALLER's. See
// {@link eventOrganizationId}; omitted, never `''`/`undefined`, because
// absence has exactly one spelling in the schema.
const organizationId = eventOrganizationId(
this._registry.getObject(object),
input.organizationRow,
);
const organizationId = eventOrganizationId(schema, input.organizationRow);
const event: DataEvent = DataEventSchema.parse({
id: generateEventUuid(),
type: `data.record.${action}`,
Expand Down
35 changes: 33 additions & 2 deletions packages/plugins/plugin-approvals/src/approval-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ import { isFileIdToken, referenceTargetOf } from '@objectstack/spec/data';
// mechanism for a second producer, and commit aa5994e17 landed this service's key
// (`approval_recall_not_submitter`) into it ahead of this consumer half.
import { renderOperationMessage, type ValidationMessageTranslator } from '@objectstack/spec/system';
import { isGrantActive } from '@objectstack/core';
import { isGrantActive, omitInternalFieldsFromWriteResponse } from '@objectstack/core';
import {
filterApproversWhoCanRead,
resolveApproverDirectoryOrg,
Expand Down Expand Up @@ -3038,7 +3038,7 @@ export class ApprovalService implements IApprovalService {
current_step: input.nodeId,
current_step_index: 0,
pending_approvers: approvers.join(','),
payload_json: input.record != null ? JSON.stringify(input.record) : null,
payload_json: input.record != null ? JSON.stringify(this.snapshotRecord(input.object, input.record)) : null,
flow_run_id: input.runId,
flow_node_id: input.nodeId,
node_config_json: JSON.stringify(configSnapshot),
Expand Down Expand Up @@ -5915,6 +5915,37 @@ export class ApprovalService implements IApprovalService {
}
}

// ── Record snapshot ──────────────────────────────────────────

/**
* The subject record as `payload_json` stores it: credential-class fields
* MASKED and `internal: true` fields OMITTED, by the one helper every
* external write response goes through (`omitInternalFieldsFromWriteResponse`,
* `@objectstack/core`) — the same answer a read of the row gives.
*
* The record arrives from the flow's `$record`, which the record-change
* trigger builds off the engine's own write result, and that result keeps the
* stored row whole for privileged in-process callers. The snapshot is the
* opposite case: it is stored, and served to approvers and submitters below
* the write boundary (the serve-time redaction in `payload-redaction.ts`
* narrows by field-level security, and cannot know a value is a stored
* credential). Defence in depth beside the engine's own event-body projection.
*
* Pure: projects a shallow copy, never the caller's record. No schema (an
* engine double without `getSchema`, an unregistered object) projects
* nothing — the same posture as the helper itself.
*/
private snapshotRecord(object: string, record: unknown): unknown {
if (!record || typeof record !== 'object' || Array.isArray(record)) return record;
let schema: unknown;
try {
schema = this.engine.getSchema?.(object);
} catch { /* schema unavailable — nothing to project against */ }
const copy = { ...(record as Record<string, unknown>) };
omitInternalFieldsFromWriteResponse(schema, copy);
return copy;
}

// ── Display enrichment ───────────────────────────────────────

/**
Expand Down
Loading
Loading