diff --git a/.changeset/record-change-payload-credential-mask.md b/.changeset/record-change-payload-credential-mask.md new file mode 100644 index 00000000000..1ca54a8783c --- /dev/null +++ b/.changeset/record-change-payload-credential-mask.md @@ -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. diff --git a/packages/objectql/src/engine-realtime-credential-mask.test.ts b/packages/objectql/src/engine-realtime-credential-mask.test.ts new file mode 100644 index 00000000000..d997cbda23a --- /dev/null +++ b/packages/objectql/src/engine-realtime-credential-mask.test.ts @@ -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>>(); + 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) { + 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) { + 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[]) { 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) { 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; + 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; + 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; + 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; + 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; + 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); + }); +}); diff --git a/packages/objectql/src/engine.ts b/packages/objectql/src/engine.ts index 397b5ffb636..48b8210911c 100644 --- a/packages/objectql/src/engine.ts +++ b/packages/objectql/src/engine.ts @@ -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'; @@ -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 | undefined, +): Record | 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; @@ -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}`, diff --git a/packages/plugins/plugin-approvals/src/approval-service.ts b/packages/plugins/plugin-approvals/src/approval-service.ts index 3f5ec962d00..6ea9ca90046 100644 --- a/packages/plugins/plugin-approvals/src/approval-service.ts +++ b/packages/plugins/plugin-approvals/src/approval-service.ts @@ -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, @@ -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), @@ -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) }; + omitInternalFieldsFromWriteResponse(schema, copy); + return copy; + } + // ── Display enrichment ─────────────────────────────────────── /** diff --git a/packages/plugins/plugin-approvals/src/approval-snapshot-credential-mask.test.ts b/packages/plugins/plugin-approvals/src/approval-snapshot-credential-mask.test.ts new file mode 100644 index 00000000000..d7cc73cc094 --- /dev/null +++ b/packages/plugins/plugin-approvals/src/approval-snapshot-credential-mask.test.ts @@ -0,0 +1,126 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +/** + * The approval request's record snapshot (`payload_json`) is stored, and served + * to approvers and submitters below the write boundary. It is built from the + * flow's `$record`, which carries the engine's own write result — kept whole + * for privileged in-process callers. So the snapshot applies the same rule a + * write response does, by the same helper: credential-class fields masked, + * `internal: true` fields omitted. The record the caller handed in is not + * touched. Fixtures are synthetic. + */ + +import { describe, it, expect, beforeEach } from 'vitest'; +import { SECRET_MASK } from '@objectstack/spec/data'; +import { ApprovalService } from './approval-service.js'; + +const OBJECT = 'probe_vault'; +const CREDENTIAL = 'synthetic-credential-5d0a'; +const INTERNAL = 'synthetic-internal-91fe'; +const CTX = { userId: 'u1', tenantId: 't1', positions: [], permissions: [] } as any; +const SYS = { isSystem: true, positions: [], permissions: [] } as any; + +interface FakeRow { [k: string]: any } + +function makeEngine(withSchema = true) { + const tables: Record = {}; + const ensure = (n: string) => (tables[n] ??= []); + const matches = (row: FakeRow, filter: any): boolean => { + if (!filter || typeof filter !== 'object') return true; + return Object.entries(filter).every(([k, v]) => { + if (k === '$or') return (v as any[]).some((s) => matches(row, s)); + if (k === '$and') return (v as any[]).every((s) => matches(row, s)); + if (v != null && typeof v === 'object' && '$in' in (v as any)) return (v as any).$in.includes(row[k]); + return row[k] === v; + }); + }; + const engine: any = { + _tables: tables, + async find(object: string, options?: any) { + const rows = ensure(object).filter((r) => matches(r, options?.filter ?? options?.where)); + return rows.slice(options?.offset ?? 0, (options?.offset ?? 0) + (options?.limit ?? 1000)); + }, + async insert(object: string, data: any) { ensure(object).push({ ...data }); return { ...data }; }, + // Opening a request and reading it back reach no update/delete verb, so + // the double declares none (the same reasoning as + // `approval-payload-masked-field.test.ts`). + }; + if (withSchema) { + engine.getSchema = (object: string) => object !== OBJECT ? undefined : { + name: OBJECT, + fields: { + id: { name: 'id', type: 'text' }, + title: { name: 'title', type: 'text' }, + access_token: { name: 'access_token', type: 'secret' }, + passphrase: { name: 'passphrase', type: 'password' }, + lookup_digest: { name: 'lookup_digest', type: 'text', internal: true }, + }, + }; + } + return engine; +} + +function openInput(record: Record) { + return { + object: OBJECT, + recordId: 'pv_1', + runId: 'run_1', + nodeId: 'approve_step', + flowName: 'vault_approval', + config: { approvers: [{ type: 'user' as const, value: 'u9' }], behavior: 'first_response' as const, lockRecord: true }, + record, + }; +} + +const submitted = () => ({ + id: 'pv_1', + title: 'Synthetic vault', + access_token: CREDENTIAL, + passphrase: `${CREDENTIAL}-p`, + lookup_digest: INTERNAL, +}); + +describe('approval record snapshot applies the write-response non-exposure rules', () => { + let engine: any; + let svc: ApprovalService; + + beforeEach(() => { + engine = makeEngine(); + svc = new ApprovalService({ engine, clock: { now: () => new Date('2026-01-01T00:00:00Z') } }); + }); + + it('stores credential-class fields masked and internal fields omitted', async () => { + await svc.openNodeRequest(openInput(submitted()), CTX); + const raw = engine._tables['sys_approval_request'][0]; + expect(raw.payload_json).not.toContain(CREDENTIAL); + expect(raw.payload_json).not.toContain(INTERNAL); + const snapshot = JSON.parse(raw.payload_json); + expect(snapshot).toEqual({ + id: 'pv_1', + title: 'Synthetic vault', + access_token: SECRET_MASK, + passphrase: SECRET_MASK, + }); + }); + + it('serves the masked snapshot on the service read door', async () => { + const req = await svc.openNodeRequest(openInput(submitted()), CTX); + const read = await svc.getRequest((req as any).id, SYS); + expect(JSON.stringify(read)).not.toContain(CREDENTIAL); + expect(JSON.stringify(read)).not.toContain(INTERNAL); + expect((read as any).payload.access_token).toBe(SECRET_MASK); + }); + + it('leaves the record the caller handed in untouched', async () => { + const record = submitted(); + await svc.openNodeRequest(openInput(record), CTX); + expect(record).toEqual(submitted()); + }); + + it('an engine with no schema surface stores the snapshot as handed in', async () => { + const bare = makeEngine(false); + const bareSvc = new ApprovalService({ engine: bare, clock: { now: () => new Date('2026-01-01T00:00:00Z') } }); + await bareSvc.openNodeRequest(openInput({ id: 'pv_1', title: 'Plain' }), CTX); + expect(JSON.parse(bare._tables['sys_approval_request'][0].payload_json)).toEqual({ id: 'pv_1', title: 'Plain' }); + }); +}); diff --git a/packages/plugins/plugin-webhooks/src/auto-enqueuer.test.ts b/packages/plugins/plugin-webhooks/src/auto-enqueuer.test.ts index f84bbb22ae6..839afcca0d2 100644 --- a/packages/plugins/plugin-webhooks/src/auto-enqueuer.test.ts +++ b/packages/plugins/plugin-webhooks/src/auto-enqueuer.test.ts @@ -19,6 +19,7 @@ import { randomUUID } from 'node:crypto'; import { describe, expect, it, vi } from 'vitest'; import { BulkDataEventSchema, DataEventSchema } from '@objectstack/spec/api'; +import { SECRET_MASK } from '@objectstack/spec/data'; import type { IDataEngine, IRealtimeService, @@ -1099,3 +1100,61 @@ describe('AutoEnqueuer — organization dimension (#13566)', () => { await ae.stop(); }); }); + +// The outbound body (sent off-box and stored as the delivery row's payload) +// applies the write-response non-exposure rules to the record bodies, by the +// same helper: credential-class fields masked, `internal: true` fields +// omitted. Defence in depth beside the engine's own event-body projection, so +// these events are built UNPROJECTED on purpose. Synthetic values. +describe('AutoEnqueuer outbound body — credential mask and internal-field omission', () => { + const CREDENTIAL = 'synthetic-credential-e83b'; + const INTERNAL = 'synthetic-internal-4c6d'; + const schemaFor = (name: string) => name !== 'contact' ? undefined : { + name: 'contact', + fields: { + id: { name: 'id', type: 'text' }, + name: { name: 'name', type: 'text' }, + api_token: { name: 'api_token', type: 'secret' }, + lookup_digest: { name: 'lookup_digest', type: 'text', internal: true }, + }, + }; + + async function boot() { + const engine = new FakeEngine({ sys_webhook: [webhook()] }); + (engine as unknown as { getSchema: typeof schemaFor }).getSchema = schemaFor; + const realtime = new FakeRealtime(); + const { enqueue, calls } = makeRecorder(); + const ae = new AutoEnqueuer(engine, realtime, enqueue, { refreshIntervalMs: 0 }); + await ae.start(); + return { realtime, calls, ae }; + } + + it('masks credential fields and omits internal fields in `after`', async () => { + const { realtime, calls, ae } = await boot(); + await realtime.publish(event('created', 'contact', { + id: 'c-1', name: 'Alice', api_token: CREDENTIAL, lookup_digest: INTERNAL, + })); + await flush(); + expect(calls).toHaveLength(1); + const body = calls[0].payload as any; + expect(body.after).toEqual({ id: 'c-1', name: 'Alice', api_token: SECRET_MASK }); + expect(JSON.stringify(body)).not.toContain(CREDENTIAL); + expect(JSON.stringify(body)).not.toContain(INTERNAL); + await ae.stop(); + }); + + it('projects `changes` too, and leaves the shared event untouched', async () => { + const { realtime, calls, ae } = await boot(); + const ev = event('updated', 'contact', { id: 'c-2', name: 'Bob', api_token: CREDENTIAL }); + (ev.payload as any).changes = { api_token: CREDENTIAL, lookup_digest: INTERNAL }; + await realtime.publish(ev); + await flush(); + const body = calls[0].payload as any; + expect(body.changes).toEqual({ api_token: SECRET_MASK }); + expect(body.after.api_token).toBe(SECRET_MASK); + // The realtime bus hands the same event to every subscriber. + expect((ev.payload as any).after.api_token).toBe(CREDENTIAL); + expect((ev.payload as any).changes.lookup_digest).toBe(INTERNAL); + await ae.stop(); + }); +}); diff --git a/packages/plugins/plugin-webhooks/src/auto-enqueuer.ts b/packages/plugins/plugin-webhooks/src/auto-enqueuer.ts index 34388595f65..c36bc9b1713 100644 --- a/packages/plugins/plugin-webhooks/src/auto-enqueuer.ts +++ b/packages/plugins/plugin-webhooks/src/auto-enqueuer.ts @@ -3,6 +3,7 @@ import type { IDataEngine, IRealtimeService, RealtimeEventPayload } from '@objectstack/spec/contracts'; import type { WebhookTriggerType } from '@objectstack/spec/automation'; import type { EnqueueHttpInput } from '@objectstack/service-messaging'; +import { omitInternalFieldsFromWriteResponse } from '@objectstack/core'; import { WEBHOOK_SECRET_FIELD, WEBHOOK_SECRET_REFUSAL_CODE, @@ -946,6 +947,10 @@ export class AutoEnqueuer { // don't accidentally dedup. const eventId = `${event.object}:${recordId}:${action}:${event.timestamp}`; + // The outbound body — stored on the delivery row and sent off-box — with + // the record bodies projected once, for every subscription below. + const outboundBody = this.projectRecordBodies(event.object, payload as Record); + for (const sub of subs) { if (!sub.triggers.has(trigger)) continue; // [#13566] The organization dimension of the match. Decided BEFORE @@ -993,7 +998,7 @@ export class AutoEnqueuer { // fields into the payload would have silently rewritten the // `object` / `action` / `timestamp` a subscriber receives. payload: { - ...payload, + ...outboundBody, object: event.object, recordId, action, @@ -1003,6 +1008,37 @@ export class AutoEnqueuer { } } + /** + * The event payload with its record bodies (`after`, `changes`) projected + * through the generic-data-path non-exposure rules — credential-class + * fields MASKED, `internal: true` fields OMITTED — by the one helper every + * external write response goes through + * (`omitInternalFieldsFromWriteResponse`, `@objectstack/core`). + * + * The engine already projects both bodies where it publishes the event; + * this is defence in depth at the one place the body leaves the box and is + * stored (`sys_http_delivery.payload_json`, readable by administrators who + * may sit below the write boundary). Idempotent over an already-projected + * body. Pure: fresh copies, never the event the realtime bus shares with + * every other subscriber. No schema (an engine without `getSchema`, an + * unregistered object) projects nothing. + */ + private projectRecordBodies(object: string, payload: Record): Record { + let schema: unknown; + try { + schema = (this.engine as { getSchema?: (name: string) => unknown }).getSchema?.(object); + } catch { /* schema unavailable — nothing to project against */ } + const out: Record = { ...payload }; + for (const key of ['before', 'after', 'changes'] as const) { + const body = out[key]; + if (!body || typeof body !== 'object' || Array.isArray(body)) continue; + const copy = { ...(body as Record) }; + omitInternalFieldsFromWriteResponse(schema, copy); + out[key] = copy; + } + return out; + } + /** * Handler for aggregate `data.records.*` events — a predicate write * (`multi: true`) that the driver reports only as an affected-row count diff --git a/packages/qa/dogfood/test/approval-snapshot-credential-field.dogfood.test.ts b/packages/qa/dogfood/test/approval-snapshot-credential-field.dogfood.test.ts new file mode 100644 index 00000000000..a461fce1f58 --- /dev/null +++ b/packages/qa/dogfood/test/approval-snapshot-credential-field.dogfood.test.ts @@ -0,0 +1,197 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. +// +// The record snapshot an approval request stores applies the write-response +// non-exposure rules, on a real boot: a credential-class field is stored +// masked and an `internal: true` field is not stored at all, while the +// engine's own write result keeps both values for the privileged in-process +// caller that performed the write. +// +// ## The composition +// +// `bootStack` with the real `SecurityPlugin`, `ObjectQL`, SQL driver, REST and +// auth layers, plus automation, the record-change trigger and the approvals +// plugin. One synthetic object carries a `password` field and an `internal` +// field; a flow routes each new record to a position one approver holds. The +// flow's `$record` is built off the engine's write result, so without the rule +// the snapshot would carry both stored values. +// +// ## What is asserted, by class +// +// - The scene is real: the flow opened a request for each record. +// - At rest, the stored snapshot carries the mask for the credential field and +// no key for the internal field — for a record created through the HTTP data +// door and for one written in-process. +// - The approver's inbox item serves the same masked snapshot. +// - Control: the in-process engine write result still carries both stored +// values. +// +// Fixtures are synthetic. `@objectstack/plugin-approvals` and +// `@objectstack/trigger-record-change` resolve to SOURCE here (this project's +// alias), so the verdict is about the checkout, not the last build. + +import { describe, it, expect } from 'vitest'; +import { bootStack } from '@objectstack/verify'; +import { ApprovalsServicePlugin } from '@objectstack/plugin-approvals'; +import { RecordChangeTriggerPlugin } from '@objectstack/trigger-record-change'; +import { defineStack, defineFlow } from '@objectstack/spec'; +import { ObjectSchema, Field, SECRET_MASK } from '@objectstack/spec/data'; +import { PermissionSetSchema } from '@objectstack/spec/security'; +import { SecurityPlugin, securityDefaultPermissionSets } from '@objectstack/plugin-security'; + +const OBJECT = 'credsnap_item'; +const POSITION = 'credsnap_reviewer'; +const CREDENTIAL_KEY = 'credsnap_pass'; +const INTERNAL_KEY = 'credsnap_digest'; +/** Synthetic stored values. */ +const CREDENTIAL = 'SYNTH-CRED-71B2'; +const INTERNAL = 'SYNTH-DIGEST-0C9E'; +const SYS = { context: { isSystem: true } } as const; + +const CredsnapItem = ObjectSchema.create({ + name: OBJECT, + label: 'Credential Snapshot Item', + pluralLabel: 'Credential Snapshot Items', + sharingModel: 'public_read_write', + fields: { + name: Field.text({ label: 'Name', required: true }), + [CREDENTIAL_KEY]: Field.password({ label: 'Pass' }), + [INTERNAL_KEY]: Field.text({ label: 'Digest', internal: true }), + }, +}); + +const CredsnapFlow = defineFlow({ + name: 'credsnap_flow', + label: 'Credential Snapshot Flow', + description: 'Routes each new item to the reviewer position.', + type: 'autolaunched', + status: 'active', + nodes: [ + { id: 'start', type: 'start', label: 'On Create', config: { objectName: OBJECT, triggerType: 'record-after-create' } }, + { + id: 'gate', + type: 'approval', + label: 'Gate', + config: { approvers: [{ type: 'position', value: POSITION }], behavior: 'first_response' }, + }, + { id: 'approved', type: 'end', label: 'Approved' }, + { id: 'rejected', type: 'end', label: 'Rejected' }, + ], + edges: [ + { id: 'e1', source: 'start', target: 'gate' }, + { id: 'e2', source: 'gate', target: 'approved', label: 'approve' }, + { id: 'e3', source: 'gate', target: 'rejected', label: 'reject' }, + ], +}); + +const credsnapStack = defineStack({ + manifest: { + id: 'com.dogfood.credential-snapshot', + namespace: 'credsnap', + version: '0.0.0', + type: 'app', + name: 'Credential Snapshot Fixture', + description: 'One object with a credential field and an internal field, one approval flow.', + }, + // ADR-0097: the flow's record-change start node needs both capabilities. + requires: ['automation', 'triggers'], + objects: [CredsnapItem], + flows: [CredsnapFlow], +}); + +/** Every fresh member: read and create on the object, read on the request object. */ +const baselineSet = PermissionSetSchema.parse({ + name: 'credsnap_baseline', + label: 'Credential snapshot baseline', + objects: { + [OBJECT]: { allowRead: true, allowCreate: true, allowEdit: false, allowDelete: false }, + sys_approval_request: { allowRead: true, allowCreate: false, allowEdit: false, allowDelete: false }, + }, +}); + +async function snapshotAtRest(ql: any, recordId: string): Promise> { + const row = await ql.findOne('sys_approval_request', { + where: { object_name: OBJECT, record_id: recordId }, + context: { isSystem: true }, + }); + expect(row, `the flow opened an approval request for ${recordId}`).toBeTruthy(); + const raw = String(row.payload_json); + expect(raw).not.toContain(CREDENTIAL); + expect(raw).not.toContain(INTERNAL); + return JSON.parse(raw) as Record; +} + +describe('approval snapshot of a record with credential-class and internal fields', () => { + it( + 'stores the credential masked and the internal field omitted, while the engine write result keeps both', + async () => { + const stack = await bootStack(credsnapStack as unknown as Parameters[0], { + automation: true, + security: new SecurityPlugin({ + defaultPermissionSets: [...securityDefaultPermissionSets, baselineSet], + fallbackPermissionSet: baselineSet.name, + }), + extraPlugins: [new RecordChangeTriggerPlugin(), new ApprovalsServicePlugin()], + }); + try { + const ql = (await stack.kernel.getServiceAsync('objectql')) as any; + const idOf = async (email: string) => + String((await ql.findOne('sys_user', { where: { email }, context: { isSystem: true } }))?.id ?? ''); + + const submitterToken = await stack.signUp('credsnap-submitter@verify.test'); + const approverToken = await stack.signUp('credsnap-approver@verify.test'); + const approverId = await idOf('credsnap-approver@verify.test'); + expect(approverId).toBeTruthy(); + + // Staff the position BEFORE the records exist: the slate resolves at + // request creation. + await ql.insert('sys_position', { id: 'pos_credsnap', name: POSITION, label: 'Reviewer', active: true }, SYS); + await ql.insert('sys_user_position', { id: 'hold_credsnap', user_id: approverId, position: POSITION }, SYS); + + // 1) Through the HTTP data door. + const createRes = await stack.apiAs(submitterToken, 'POST', `/data/${OBJECT}`, { + name: 'Item one', [CREDENTIAL_KEY]: CREDENTIAL, [INTERNAL_KEY]: INTERNAL, + }); + expect(createRes.status).toBe(201); + const created = (await createRes.json()) as { id?: string; record?: { id?: string } }; + const viaHttp = String(created.id ?? created.record?.id ?? ''); + expect(viaHttp).toBeTruthy(); + + const httpSnapshot = await snapshotAtRest(ql, viaHttp); + expect(httpSnapshot).toHaveProperty('name', 'Item one'); + expect(httpSnapshot[CREDENTIAL_KEY]).toBe(SECRET_MASK); + expect(httpSnapshot).not.toHaveProperty(INTERNAL_KEY); + + // 2) In-process, as a privileged server-side writer. Control: the + // engine's write result still carries both stored values. + const written = await ql.insert(OBJECT, { + name: 'Item two', [CREDENTIAL_KEY]: CREDENTIAL, [INTERNAL_KEY]: INTERNAL, + }, SYS); + expect(written[CREDENTIAL_KEY]).toBe(CREDENTIAL); + expect(written[INTERNAL_KEY]).toBe(INTERNAL); + + const engineSnapshot = await snapshotAtRest(ql, String(written.id)); + expect(engineSnapshot).toHaveProperty('name', 'Item two'); + expect(engineSnapshot[CREDENTIAL_KEY]).toBe(SECRET_MASK); + expect(engineSnapshot).not.toHaveProperty(INTERNAL_KEY); + + // 3) The approver's inbox serves the same masked snapshot. + const listRes = await stack.apiAs(approverToken, 'GET', '/approvals/requests?status=pending'); + expect(listRes.status).toBe(200); + const listBody = (await listRes.json()) as { data: Array> }; + const rows = listBody.data.filter((r) => r.object_name === OBJECT); + expect(rows).toHaveLength(2); + for (const row of rows) { + const itemRes = await stack.apiAs(approverToken, 'GET', `/approvals/requests/${String(row.id)}`); + expect(itemRes.status).toBe(200); + const item = (await itemRes.json()) as Record; + expect(JSON.stringify(item)).not.toContain(CREDENTIAL); + expect(JSON.stringify(item)).not.toContain(INTERNAL); + expect((item.payload as Record)[CREDENTIAL_KEY]).toBe(SECRET_MASK); + } + } finally { + await stack.stop(); + } + }, + 120_000, + ); +}); diff --git a/packages/services/service-knowledge/src/__tests__/knowledge-service.test.ts b/packages/services/service-knowledge/src/__tests__/knowledge-service.test.ts index 4cb4f0ec043..523cf052f2f 100644 --- a/packages/services/service-knowledge/src/__tests__/knowledge-service.test.ts +++ b/packages/services/service-knowledge/src/__tests__/knowledge-service.test.ts @@ -339,3 +339,83 @@ describe('Pure helpers', () => { expect(doc.content).not.toContain('r1'); // id excluded }); }); + +// An index document is searchable text served to every reader its source +// admits, so it never carries a credential-class field (the set the engine +// masks on read) or an `internal: true` field (the set it omits) — under `*`, +// or when a source names one explicitly. The field sets come from the same +// collectors the write-response helper uses. Synthetic values. +describe('recordToDocument withholds credential-class and internal fields', () => { + const CREDENTIAL = 'synthetic-credential-a27f'; + const INTERNAL = 'synthetic-internal-6e30'; + const schema = { + name: 'probe_doc', + fields: { + id: { name: 'id', type: 'text' }, + title: { name: 'title', type: 'text' }, + notes: { name: 'notes', type: 'textarea' }, + api_token: { name: 'api_token', type: 'secret' }, + passphrase: { name: 'passphrase', type: 'password' }, + lookup_digest: { name: 'lookup_digest', type: 'text', internal: true }, + }, + }; + const record = () => ({ + id: 'r1', title: 'Visible', notes: 'Body', + api_token: CREDENTIAL, passphrase: `${CREDENTIAL}-p`, lookup_digest: INTERNAL, + }); + const source = (contentFields: string[], metadataFields: string[] = []): KnowledgeSource => ({ + id: 's1', label: 's1', adapter: 'memory', + source: { kind: 'object', object: 'probe_doc', contentFields, metadataFields } as ObjectKnowledgeSource, + }); + + it("skips them under '*'", () => { + const src = source(['*']); + const doc = recordToDocument(src, src.source as ObjectKnowledgeSource, record(), schema); + expect(doc.content).toBe('Visible\n\nBody'); + expect(JSON.stringify(doc)).not.toContain(CREDENTIAL); + expect(JSON.stringify(doc)).not.toContain(INTERNAL); + }); + + it('skips them when named explicitly as content or metadata', () => { + const src = source(['title', 'api_token', 'lookup_digest'], ['passphrase', 'notes']); + const doc = recordToDocument(src, src.source as ObjectKnowledgeSource, record(), schema); + expect(doc.content).toBe('Visible'); + expect(doc.metadata).toEqual({ notes: 'Body' }); + expect(JSON.stringify(doc)).not.toContain(CREDENTIAL); + expect(JSON.stringify(doc)).not.toContain(INTERNAL); + }); + + it('the event-sync path reads the schema off the bound engine', async () => { + const getSchema = vi.fn((name: string) => (name === 'probe_doc' ? schema : undefined)); + const svc = new KnowledgeService({ dataEngine: { find: vi.fn(), getSchema } as unknown as IDataEngine }); + const adapter = makeAdapter('memory'); + svc.registerAdapter('memory', adapter); + svc.registerSource(source(['*'])); + await svc.handleRecordUpsert('probe_doc', record()); + expect(getSchema).toHaveBeenCalledWith('probe_doc'); + const [docs] = adapter.upsertSpy.mock.calls[0] as unknown as [KnowledgeDocument[]]; + expect(docs[0].content).toBe('Visible\n\nBody'); + expect(JSON.stringify(docs)).not.toContain(CREDENTIAL); + expect(JSON.stringify(docs)).not.toContain(INTERNAL); + }); + + it('the reindex walk reads the schema off the bound engine', async () => { + const getSchema = vi.fn((name: string) => (name === 'probe_doc' ? schema : undefined)); + const find = vi.fn(async () => [record()]); + const svc = new KnowledgeService({ dataEngine: { find, getSchema } as unknown as IDataEngine }); + const adapter = makeAdapter('memory'); + svc.registerAdapter('memory', adapter); + svc.registerSource(source(['*'])); + const res = await svc.reindexSource('s1'); + expect(res.indexed).toBe(1); + const [docs] = adapter.upsertSpy.mock.calls[0] as unknown as [KnowledgeDocument[]]; + expect(JSON.stringify(docs)).not.toContain(CREDENTIAL); + expect(JSON.stringify(docs)).not.toContain(INTERNAL); + }); + + it('without a schema nothing is withheld here', () => { + const src = source(['title', 'notes']); + const doc = recordToDocument(src, src.source as ObjectKnowledgeSource, { id: 'r1', title: 'A', notes: 'B' }); + expect(doc.content).toBe('A\n\nB'); + }); +}); diff --git a/packages/services/service-knowledge/src/knowledge-service.ts b/packages/services/service-knowledge/src/knowledge-service.ts index 93069e83cae..9f0e9e62383 100644 --- a/packages/services/service-knowledge/src/knowledge-service.ts +++ b/packages/services/service-knowledge/src/knowledge-service.ts @@ -15,6 +15,10 @@ import type { KnowledgeSource, ObjectKnowledgeSource, } from '@objectstack/spec/ai'; +import { + collectCredentialWriteResponseFields, + collectInternalWriteResponseFields, +} from '@objectstack/core'; /** * Minimal logger shape; falls back to no-op when none is provided. @@ -179,7 +183,8 @@ export class KnowledgeService implements IKnowledgeService { }; } - const docs = records.map((rec) => recordToDocument(source, objSource, rec)); + const schema = this.schemaFor(objSource.object); + const docs = records.map((rec) => recordToDocument(source, objSource, rec, schema)); if (docs.length > 0) { await adapter.upsert(docs, { source, reason: 'reindex' }); } @@ -239,7 +244,7 @@ export class KnowledgeService implements IKnowledgeService { for (const source of targets) { try { const objSource = source.source as ObjectKnowledgeSource; - const doc = recordToDocument(source, objSource, record); + const doc = recordToDocument(source, objSource, record, this.schemaFor(object)); const adapter = this.getAdapter(source.adapter); await adapter.upsert([doc], { source, reason: 'event-sync' }); } catch (err) { @@ -250,6 +255,20 @@ export class KnowledgeService implements IKnowledgeService { } } + /** + * The registered object definition, when the bound engine exposes one + * (`getSchema` is not on the `IDataEngine` contract; the ObjectQL engine + * carries it). `undefined` — nothing to withhold by — otherwise. + */ + private schemaFor(object: string): unknown { + try { + return (this.options.dataEngine as { getSchema?: (name: string) => unknown } | undefined) + ?.getSchema?.(object); + } catch { + return undefined; + } + } + /** Apply an ObjectQL `data.record.deleted` event (the id is the payload's `recordId`). */ async handleRecordDelete(object: string, recordId: string): Promise { const targets = this.sourcesForObject(object); @@ -407,26 +426,44 @@ export function documentIdFor(sourceId: string, recordId: string): string { /** * Project an ObjectQL record into a `KnowledgeDocument` per the * source's `contentFields` / `metadataFields` config. Pure function. + * + * `schema` (the record's registered object definition) names the fields an + * index never carries: credential-class fields (the set the engine masks on + * read) and `internal: true` fields (the set it omits), asked through the same + * collectors the write-response helper uses (`@objectstack/core`). They are + * SKIPPED — under `'*'`, and when a source names one explicitly as content or + * metadata — rather than indexed masked: an index document is searchable text + * served to every reader the source admits, and a mask placeholder is no + * content. Without a schema nothing is withheld here (the engine's event body + * and read path still apply their own projection upstream). */ export function recordToDocument( source: KnowledgeSource, objSource: ObjectKnowledgeSource, record: Record, + schema?: unknown, ): KnowledgeDocument { + const withheld = new Set([ + ...collectCredentialWriteResponseFields(schema), + ...collectInternalWriteResponseFields(schema), + ]); const recordId = String(record.id ?? (record as any)._id ?? ''); const contentParts: string[] = []; for (const field of objSource.contentFields) { if (field === '*') { for (const [k, v] of Object.entries(record)) { + if (withheld.has(k)) continue; if (typeof v === 'string' && v.length > 0 && k !== 'id') contentParts.push(v); } } else { + if (withheld.has(field)) continue; const v = record[field]; if (v != null) contentParts.push(String(v)); } } const metadata: Record = {}; for (const field of objSource.metadataFields ?? []) { + if (withheld.has(field)) continue; if (record[field] !== undefined) metadata[field] = record[field]; } return { @@ -434,7 +471,7 @@ export function recordToDocument( sourceId: source.id, sourceRecordId: recordId || undefined, content: contentParts.join('\n\n'), - title: typeof record.title === 'string' ? record.title : undefined, + title: typeof record.title === 'string' && !withheld.has('title') ? record.title : undefined, metadata, }; }