From 557591e5162553aac0632ab158b16607ba51ba27 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 16:06:50 +0000 Subject: [PATCH 1/5] wip: stage-1 principal-less producers take the explicit system opt-in Claude-Session: https://claude.ai/code/session_01WMQprn46CND82KmY8sZWBu Co-authored-by: Claude --- .../src/sys-metadata-repository.ts | 12 +++++- .../service-messaging/src/email-channel.ts | 11 ++++- .../src/fan-out-system-context.ts | 18 +++++++-- .../src/inbox-system-context.ts | 38 ++++++++++++++++++ .../src/messaging-service.ts | 35 +++++++++++----- .../src/outbox-dispatcher-scope.ts | 27 +++++++++++++ .../src/recipient-resolver.ts | 34 ++++++++++++++-- .../service-messaging/src/sms-channel.ts | 11 ++++- .../service-messaging/src/sql-http-outbox.ts | 30 +++++++++----- .../service-messaging/src/sql-outbox.ts | 40 +++++++++++++++---- .../src/template-renderer.ts | 5 ++- .../src/settings-service-plugin.ts | 19 ++++++++- .../service-settings/src/settings-service.ts | 14 +++++-- .../service-storage/src/metadata-store.ts | 28 ++++++++++++- 14 files changed, 275 insertions(+), 47 deletions(-) create mode 100644 packages/services/service-messaging/src/inbox-system-context.ts diff --git a/packages/metadata-protocol/src/sys-metadata-repository.ts b/packages/metadata-protocol/src/sys-metadata-repository.ts index dc58518587c..e0e1a888676 100644 --- a/packages/metadata-protocol/src/sys-metadata-repository.ts +++ b/packages/metadata-protocol/src/sys-metadata-repository.ts @@ -496,7 +496,8 @@ export class SysMetadataRepository implements MetadataRepository { // // [#21911, ADR-0096] This read, and every store call of `put`, `delete`, // `promoteDraft`, `restoreVersion`, `listDrafts` and the two lineage - // counters, carries the explicit system opt-in (`{ ...ctx, isSystem: true }` + // counters — and [#21908] the reads of `getByHash`, `list`, `history` and + // `watch`'s replay — carries the explicit system opt-in (`{ ...ctx, isSystem: true }` // inside a transaction, so the handle rides along): the repository is // platform plumbing under a door that already authorized the caller, and // it scopes its own rows by organization. None of them reaches the data @@ -520,6 +521,7 @@ export class SysMetadataRepository implements MetadataRepository { async getByHash(ref: MetaRef, hash: string): Promise { this.assertOpen(); const full = this.fullRef(ref); + // [#21908] The explicit system opt-in, as `get` carries — see its note. const row = await this.engine.findOne(this.historyTable, { where: { organization_id: this.organizationId, @@ -527,6 +529,7 @@ export class SysMetadataRepository implements MetadataRepository { name: full.name, checksum: hash, }, + context: { isSystem: true }, }); if (!row) return null; const rawBody = (row as any).metadata; @@ -1189,9 +1192,11 @@ export class SysMetadataRepository implements MetadataRepository { state: 'active', }; if (filter.type) where.type = filter.type; + // [#21908] The explicit system opt-in, as `get` carries — see its note. const rows = await this.engine.find('sys_metadata', { where, limit: filter.limit, + context: { isSystem: true }, }); for (const row of rows) { if (filter.nameContains && !String(row.name).includes(filter.nameContains)) continue; @@ -1305,7 +1310,8 @@ export class SysMetadataRepository implements MetadataRepository { type: full.type, name: full.name, }; - const rows = await this.engine.find(this.historyTable, { where }); + // [#21908] The explicit system opt-in, as `get` carries — see its note. + const rows = await this.engine.find(this.historyTable, { where, context: { isSystem: true } }); rows.sort((a: any, b: any) => { const va = typeof a.event_seq === 'number' ? a.event_seq : 0; const vb = typeof b.event_seq === 'number' ? b.event_seq : 0; @@ -1378,8 +1384,10 @@ export class SysMetadataRepository implements MetadataRepository { filter: WatchFilter, since: number, ): Promise { + // [#21908] The explicit system opt-in, as `get` carries — see its note. const rows = await this.engine.find(this.historyTable, { where: { organization_id: this.organizationId }, + context: { isSystem: true }, }); const out: MetadataEvent[] = []; for (const row of rows as Array>) { diff --git a/packages/services/service-messaging/src/email-channel.ts b/packages/services/service-messaging/src/email-channel.ts index a35f24537f6..258368b7e17 100644 --- a/packages/services/service-messaging/src/email-channel.ts +++ b/packages/services/service-messaging/src/email-channel.ts @@ -17,6 +17,7 @@ import { DEFAULT_LOCALE, } from './template-renderer.js'; import { RECIPIENT_LOCALE_FIELD, USER_OBJECT, resolveRecipientLocale } from './recipient-locale.js'; +import { FAN_OUT_SYSTEM_CONTEXT } from './fan-out-system-context.js'; /** The user identity object a recipient id is resolved to an address against (re-exported from the locale seam, #13881). */ export { USER_OBJECT }; @@ -200,10 +201,12 @@ export function createEmailChannel(opts: EmailChannelOptions): MessagingChannel // under that key is not a read — rung 2 applies. let localeRead = false; try { + // The explicit system opt-in — see FAN_OUT_SYSTEM_CONTEXT: the + // recipient's address and locale, read for the delivery only. user = await data.findOne(userObject, { where: { id: recipient }, fields: ['email', RECIPIENT_LOCALE_FIELD], - }); + }, { context: FAN_OUT_SYSTEM_CONTEXT }); localeRead = true; } catch (err) { // Ruling item 3: NO path may dead-letter because of the locale @@ -215,7 +218,11 @@ export function createEmailChannel(opts: EmailChannelOptions): MessagingChannel `[email] recipient lookup for '${recipient}' with '${RECIPIENT_LOCALE_FIELD}' failed (${(err as Error).message}); retrying address-only`, ); try { - user = await data.findOne(userObject, { where: { id: recipient }, fields: ['email'] }); + user = await data.findOne( + userObject, + { where: { id: recipient }, fields: ['email'] }, + { context: FAN_OUT_SYSTEM_CONTEXT }, + ); } catch (retryErr) { ctx.logger.warn(`[email] address lookup for '${recipient}' failed (${(retryErr as Error).message})`); return undefined; diff --git a/packages/services/service-messaging/src/fan-out-system-context.ts b/packages/services/service-messaging/src/fan-out-system-context.ts index be688ec7f35..a3b0ed02e05 100644 --- a/packages/services/service-messaging/src/fan-out-system-context.ts +++ b/packages/services/service-messaging/src/fan-out-system-context.ts @@ -10,6 +10,14 @@ * `PreferenceResolver.loadRows` (`sys_notification_preference`) and * `RecipientResolver.resolveEmail` (`sys_user`). * + * [#21908] Also by the rest of the fan-out: `RecipientResolver.resolveRole` + * (`sys_member`), `resolveTeam` (`sys_team_member`) and `resolveOwnerOf` (the + * named record, projected to `id` plus its owner fields — maintainer ruling + * Q2), the email and SMS channels' recipient reads (`sys_user`: the address or + * number, and the locale), `NotificationTemplateStore.load` + * (`sys_notification_template`), and `emit()`'s dedup lookup + * (`sys_notification` by `dedup_key`). + * * `emit()` takes no caller context. The door in front of an emitting caller * has already decided whether it may notify, and every row these calls read or * write belongs to a RECIPIENT, not to the emitter — so no caller's grants @@ -20,9 +28,11 @@ * ⚠️ The reads are made on another principal's behalf. What they return — a * user id for an address, a locale, a preference row — is consumed inside the * fan-out. `emit()` answers the notification id, counts and per-delivery - * outcomes; the one read-derived value in those is a delivery's recipient id, - * resolved from an address the emitter itself named, and no in-repo caller of - * `emit()` relays the outcomes to a door. Keep it that way. A read under this - * context whose result reaches the caller would widen what that caller can see. + * outcomes; the read-derived values in those are a delivery's recipient id, + * resolved from an audience the emitter itself named, and — on a dedup hit — + * the id of the event already written under the dedup key the emitter itself + * named. No in-repo caller of `emit()` relays the outcomes to a door. Keep it + * that way. A read under this context whose result reaches the caller would + * widen what that caller can see. */ export const FAN_OUT_SYSTEM_CONTEXT = { isSystem: true } as const; diff --git a/packages/services/service-messaging/src/inbox-system-context.ts b/packages/services/service-messaging/src/inbox-system-context.ts new file mode 100644 index 00000000000..fe27b436d38 --- /dev/null +++ b/packages/services/service-messaging/src/inbox-system-context.ts @@ -0,0 +1,38 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +/** + * [#21908] The execution context the inbox's own read-state calls run under: + * the explicit system opt-in, taken INSIDE this service (maintainer ruling Q1). + * + * Carried by `MessagingService.listInbox` (the `sys_inbox_message` window and + * `countUnreadTotal`'s total), `readReceiptStates` + * (`sys_notification_receipt`), and mark-read / mark-all-read: + * `unreadNotificationIds`, `upsertReadReceipt` (its read, its update and its + * insert) and `notificationOrganization`. + * + * Before this, each of those calls reached the data engine with no context at + * all — no principal and no opt-in — and passed the security middleware only + * through its principal-less hand-off (ADR-0096 E1), which D5 closes. Carrying + * the caller's principal instead is not open: `member_default` grants a member + * READ ONLY on `sys_inbox_message` / `sys_notification_receipt` (the inbox + * channel is their writer), so mark-read would break. + * + * ## The boundary is this service's own user scope + * + * RLS does not read these calls, and it did not read them before either. What + * keeps them to the caller's own rows is the user id the REST door derives + * from the authenticated session (`runtime/src/domains/notifications.ts`), or + * that `*AsCaller` derives from the caller's execution context: + * + * - every read of the user's rows is keyed `where: { user_id }`; + * - the receipt it inserts is stamped with that `user_id`; + * - the receipt it updates is addressed by the id that a `user_id`-keyed read + * of the SAME call returned, and by nothing else; + * - `notificationOrganization` reads one `sys_notification` EVENT row (which + * names no recipient) by id, projecting `id` and `organization_id`, and its + * only use is the organization stamp on the caller's own receipt. + * + * ⛔ Never let a call under this context drop the `user_id` term or address a + * row a `user_id`-keyed read did not return: here that term IS the isolation. + */ +export const INBOX_SYSTEM_CONTEXT = { isSystem: true } as const; diff --git a/packages/services/service-messaging/src/messaging-service.ts b/packages/services/service-messaging/src/messaging-service.ts index d978e84d929..98e45c6f63a 100644 --- a/packages/services/service-messaging/src/messaging-service.ts +++ b/packages/services/service-messaging/src/messaging-service.ts @@ -23,6 +23,7 @@ import type { import { INBOX_OBJECT, RECEIPT_OBJECT } from './inbox-channel.js'; import { type InboxCaller, resolveInboxRecipient } from './inbox-caller.js'; import { FAN_OUT_SYSTEM_CONTEXT } from './fan-out-system-context.js'; +import { INBOX_SYSTEM_CONTEXT } from './inbox-system-context.js'; import { assertActorReferenceResolves } from './actor-reference.js'; /** The L2 event object every `emit()` writes one row to (ADR-0030). */ @@ -545,7 +546,12 @@ export class MessagingService { if (opts.type) where.topic = opts.type; const [rows, stateByNotif] = await Promise.all([ - data.find(INBOX_OBJECT, { where, orderBy: [{ field: 'created_at', order: 'desc' }], limit }) as Promise>>, + // The explicit system opt-in, scoped by `where.user_id` — see INBOX_SYSTEM_CONTEXT. + data.find( + INBOX_OBJECT, + { where, orderBy: [{ field: 'created_at', order: 'desc' }], limit }, + { context: INBOX_SYSTEM_CONTEXT }, + ) as Promise>>, this.readReceiptStates(data, userId), ]); @@ -641,10 +647,11 @@ export class MessagingService { where: Record, stateByNotif: ReadonlyMap, ): Promise { + // `where` is listInbox's own, so it carries `user_id` — see INBOX_SYSTEM_CONTEXT. const ids = (await data.find(INBOX_OBJECT, { where, fields: ['notification_id'], - })) as Array>; + }, { context: INBOX_SYSTEM_CONTEXT })) as Array>; let unread = 0; for (const row of ids) { @@ -671,7 +678,7 @@ export class MessagingService { private async readReceiptStates(data: IDataEngine, userId: string): Promise> { const receipts = await (data.find(RECEIPT_OBJECT, { where: { user_id: userId, channel: 'inbox' }, - }) as Promise>>).catch(() => [] as Array>); + }, { context: INBOX_SYSTEM_CONTEXT }) as Promise>>).catch(() => [] as Array>); const stateByNotif = new Map(); for (const r of receipts) { @@ -860,7 +867,7 @@ export class MessagingService { data.find(INBOX_OBJECT, { where: { user_id: userId }, fields: ['notification_id'], - }) as Promise>>, + }, { context: INBOX_SYSTEM_CONTEXT }) as Promise>>, this.readReceiptStates(data, userId), ]); @@ -886,9 +893,15 @@ export class MessagingService { ): Promise { const where = { notification_id: notificationId, user_id: userId, channel: 'inbox' }; const flipToRead = async (): Promise => { - const existing = await data.findOne(RECEIPT_OBJECT, { where, fields: ['id'] }); + // Under INBOX_SYSTEM_CONTEXT: the read is keyed on `user_id`, and the + // update addresses only the id that read returned. + const existing = await data.findOne(RECEIPT_OBJECT, { where, fields: ['id'] }, { context: INBOX_SYSTEM_CONTEXT }); if (!existing?.id) return false; - await data.update(RECEIPT_OBJECT, { state: 'read', at }, { where: { id: existing.id } } as never); + await data.update( + RECEIPT_OBJECT, + { state: 'read', at }, + { where: { id: existing.id }, context: INBOX_SYSTEM_CONTEXT } as never, + ); return true; }; @@ -931,7 +944,7 @@ export class MessagingService { at, organization_id: organizationId, created_at: at, - }); + }, { context: INBOX_SYSTEM_CONTEXT }); return 1; } catch (err) { if (isUniqueViolationError(err) && (await flipToRead())) return 1; @@ -955,10 +968,12 @@ export class MessagingService { notificationId: string, ): Promise { try { + // One event row, two columns, used only as the stamp on the + // caller's own receipt — see INBOX_SYSTEM_CONTEXT. const row = await data.findOne(NOTIFICATION_EVENT_OBJECT, { where: { id: notificationId }, fields: ['id', 'organization_id'], - }); + }, { context: INBOX_SYSTEM_CONTEXT }); const org = (row as Record | undefined)?.organization_id; return org != null && String(org) !== '' ? String(org) : null; } catch (err) { @@ -1295,10 +1310,12 @@ export class MessagingService { /** Find an existing event id by its dedup key, tolerating lookup failure. */ private async findEventByDedupKey(data: IDataEngine, dedupKey: string): Promise { try { + // The explicit system opt-in — see FAN_OUT_SYSTEM_CONTEXT. The id it + // finds is the one `emit()` answers for a key the emitter named. const row = await data.findOne(NOTIFICATION_EVENT_OBJECT, { where: { dedup_key: dedupKey }, fields: ['id'], - }); + }, { context: FAN_OUT_SYSTEM_CONTEXT }); const id = row?.id; return id != null && String(id).length > 0 ? String(id) : undefined; } catch (err) { diff --git a/packages/services/service-messaging/src/outbox-dispatcher-scope.ts b/packages/services/service-messaging/src/outbox-dispatcher-scope.ts index 4f9acdc9c4d..19acc8a5943 100644 --- a/packages/services/service-messaging/src/outbox-dispatcher-scope.ts +++ b/packages/services/service-messaging/src/outbox-dispatcher-scope.ts @@ -80,6 +80,12 @@ export function dispatcherSweepOptions( * `SqlNotificationOutbox.claim` / `claimDigest` / `reapExpired` and * `SqlHttpOutbox.claim` / `reapExpired` issue. * + * [#21908] And by the ACK that closes each attempt — `SqlNotificationOutbox.ack` + * and `SqlHttpOutbox.ack` (both arities, `ackById` included): the state read, + * the compare-and-set write and the read-back. Same warrant, re-derived: `ack` + * has exactly the two callers {@link dispatcherAckOptions} names, both inside a + * dispatcher's `runPartition()` tick. + * * The warrant is {@link dispatcherSweepOptions}'s, read for authorization * instead of tenancy: no request, session or principal exists on the * `setInterval` tick that reaches these sites, so no caller's grants could @@ -94,6 +100,27 @@ export function dispatcherSweepOptions( export const DISPATCHER_SYSTEM_CONTEXT = { isSystem: true } as const; +/** + * [#21908] The execution context the outboxes' PRODUCER-side calls run under: + * the explicit system opt-in, carried by `SqlNotificationOutbox.enqueue` and + * `SqlHttpOutbox.enqueue` / `recordUndeliverable` (the dedup read, the insert, + * and the read that resolves a lost dedup race) and by both outboxes' `list`. + * + * The warrant is the outbox contract's: `INotificationOutbox` / + * `IHttpOutbox` carry no caller at all. `enqueue` is reached from the emit + * fan-out and from the HTTP producers (webhook auto-enqueue, flow callouts), + * each of which has already decided that a delivery should exist; the row it + * writes is the outbox's own bookkeeping, stamped with the PRODUCER's + * organization on the row (#13546), not a caller's. Without the opt-in these + * calls reach the data engine with no principal and no system opt-in, which is + * the principal-less hand-off ADR-0096 D5 closes. + * + * ⛔ Never on `redeliver`: it is the one request-reachable call on these + * objects, and it threads the caller's tenant (see {@link DISPATCHER_SYSTEM_CONTEXT}). + */ +export const OUTBOX_SYSTEM_CONTEXT = { isSystem: true } as const; + + /** * The write options for a dispatcher **`ack`** — the single-record * (`multi: false`) write that records one delivery attempt's outcome on diff --git a/packages/services/service-messaging/src/recipient-resolver.ts b/packages/services/service-messaging/src/recipient-resolver.ts index d69aa303b3b..dc4a8d618d7 100644 --- a/packages/services/service-messaging/src/recipient-resolver.ts +++ b/packages/services/service-messaging/src/recipient-resolver.ts @@ -152,7 +152,13 @@ export class RecipientResolver { const where: Record = { role }; if (ctx.organizationId) where.organization_id = ctx.organizationId; try { - const rows = await data.find(this.memberObject, { where, fields: ['user_id'], limit: 10000 }); + // The explicit system opt-in — see FAN_OUT_SYSTEM_CONTEXT: a + // membership read whose only use is the recipient ids. + const rows = await data.find( + this.memberObject, + { where, fields: ['user_id'], limit: 10000 }, + { context: FAN_OUT_SYSTEM_CONTEXT }, + ); return userIds(rows); } catch (err) { this.opts.logger.warn(`[recipients] role '${role}' lookup failed (${msg(err)}); 0 recipients`); @@ -164,11 +170,12 @@ export class RecipientResolver { private async resolveTeam(teamId: string, data: IDataEngine | undefined): Promise { if (!teamId || !data) return []; try { + // The explicit system opt-in — see FAN_OUT_SYSTEM_CONTEXT. const rows = await data.find(this.teamMemberObject, { where: { team_id: teamId }, fields: ['user_id'], limit: 10000, - }); + }, { context: FAN_OUT_SYSTEM_CONTEXT }); return userIds(rows); } catch (err) { this.opts.logger.warn(`[recipients] team '${teamId}' lookup failed (${msg(err)}); 0 recipients`); @@ -176,11 +183,30 @@ export class RecipientResolver { } } - /** `owner_of:` → the owner/assignee field of the referenced record. */ + /** + * `owner_of:` → the owner/assignee field of the referenced record. + * + * [#21908, maintainer ruling Q2] The explicit system opt-in, the posture its + * sibling {@link resolveEmail} carries (see FAN_OUT_SYSTEM_CONTEXT): the + * read projects `id` and the owner fields and nothing else, and its result + * leaves this method only as recipient ids. Before this the read reached the + * engine with no principal, and the sharing middleware answers such a read + * of a `private` object with deny-all, so an `owner_of:` audience on one + * silently resolved to nobody. Under the opt-in it resolves the owner, + * whatever the object's sharing model; nothing the read returns reaches the + * emitter. + * + * ⛔ Never widen the projection, and never return anything but the id: the + * record is read on the platform's authority, not the emitter's. + */ private async resolveOwnerOf(object: string, id: string, data: IDataEngine | undefined): Promise { if (!object || !id || !data) return []; try { - const rec = await data.findOne(object, { where: { id }, fields: ['id', ...this.ownerFields] }); + const rec = await data.findOne( + object, + { where: { id }, fields: ['id', ...this.ownerFields] }, + { context: FAN_OUT_SYSTEM_CONTEXT }, + ); if (!rec) return []; for (const f of this.ownerFields) { const v = rec[f]; diff --git a/packages/services/service-messaging/src/sms-channel.ts b/packages/services/service-messaging/src/sms-channel.ts index c55cc69f2f9..4fbbcaf78e5 100644 --- a/packages/services/service-messaging/src/sms-channel.ts +++ b/packages/services/service-messaging/src/sms-channel.ts @@ -17,6 +17,7 @@ import { DEFAULT_LOCALE, } from './template-renderer.js'; import { RECIPIENT_LOCALE_FIELD, USER_OBJECT, resolveRecipientLocale } from './recipient-locale.js'; +import { FAN_OUT_SYSTEM_CONTEXT } from './fan-out-system-context.js'; /** * Structural view of the SMS service (`@objectstack/service-sms`'s @@ -148,10 +149,12 @@ export function createSmsChannel(opts: SmsChannelOptions): MessagingChannel { // retry the column was never asked for, so rung 2 applies. let localeRead = false; try { + // The explicit system opt-in — see FAN_OUT_SYSTEM_CONTEXT: the + // recipient's number and locale, read for the delivery only. user = await data.findOne(userObject, { where: { id: recipient }, fields: ['phone_number', RECIPIENT_LOCALE_FIELD], - }); + }, { context: FAN_OUT_SYSTEM_CONTEXT }); localeRead = true; } catch (err) { // Ruling item 3 (#13881): the locale read must never cost the @@ -160,7 +163,11 @@ export function createSmsChannel(opts: SmsChannelOptions): MessagingChannel { `[sms] recipient lookup for '${recipient}' with '${RECIPIENT_LOCALE_FIELD}' failed (${(err as Error).message}); retrying phone-only`, ); try { - user = await data.findOne(userObject, { where: { id: recipient }, fields: ['phone_number'] }); + user = await data.findOne( + userObject, + { where: { id: recipient }, fields: ['phone_number'] }, + { context: FAN_OUT_SYSTEM_CONTEXT }, + ); } catch (retryErr) { ctx.logger.warn(`[sms] phone lookup for '${recipient}' failed (${(retryErr as Error).message})`); return undefined; diff --git a/packages/services/service-messaging/src/sql-http-outbox.ts b/packages/services/service-messaging/src/sql-http-outbox.ts index 27b16c1a119..8fb4f6c6755 100644 --- a/packages/services/service-messaging/src/sql-http-outbox.ts +++ b/packages/services/service-messaging/src/sql-http-outbox.ts @@ -6,6 +6,7 @@ import { hashPartition } from './backoff.js'; import { toEpochMs } from './audit-timestamp.js'; import { DISPATCHER_SYSTEM_CONTEXT, + OUTBOX_SYSTEM_CONTEXT, dispatcherAckCasOptions, dispatcherAckOptions, dispatcherSweepOptions, @@ -163,10 +164,12 @@ export class SqlHttpOutbox implements IHttpOutbox { input: Omit, terminal: { signature: string | undefined; status: HttpDeliveryStatus; error?: string }, ): Promise { + // The explicit system opt-in on every call of the producer side — see + // OUTBOX_SYSTEM_CONTEXT. const existing = await this.engine.findOne(this.objectName, { where: { source: input.source, dedup_key: input.dedupKey }, fields: ['id'], - }); + }, { context: OUTBOX_SYSTEM_CONTEXT }); if (existing?.id) return existing.id as string; const id = randomUUID(); @@ -202,13 +205,13 @@ export class SqlHttpOutbox implements IHttpOutbox { updated_at: now, }; try { - await this.engine.insert(this.objectName, row); + await this.engine.insert(this.objectName, row, { context: OUTBOX_SYSTEM_CONTEXT }); return id; } catch (err) { const winner = await this.engine.findOne(this.objectName, { where: { source: input.source, dedup_key: input.dedupKey }, fields: ['id'], - }); + }, { context: OUTBOX_SYSTEM_CONTEXT }); if (winner?.id) return winner.id as string; throw err; } @@ -354,10 +357,12 @@ export class SqlHttpOutbox implements IHttpOutbox { // The runtime half of the credential contract, for JS callers and // casts: refused before any IO. assertHttpClaimCredential(id, claimed); + // Every call of the ack carries the dispatcher's opt-in — see + // DISPATCHER_SYSTEM_CONTEXT. const current = (await this.engine.findOne(this.objectName, { where: { id }, fields: ['status', 'attempts', 'claimed_by', 'claimed_at'], - })) as Pick | null; + }, { context: DISPATCHER_SYSTEM_CONTEXT })) as Pick | null; // An id matching no row: no state to corrupt, no claim to lose. if (!current) return; // Not claimed at all — reaped back to the queue, already terminal, or @@ -388,7 +393,10 @@ export class SqlHttpOutbox implements IHttpOutbox { // but the id is silently discarded (#11009). Declared a global-sweep // site — no request context exists on the tick that reaches here. // Warrant in `outbox-dispatcher-scope.ts`. - dispatcherAckCasOptions(id, 'in_flight', claimed.claimedBy, claimed.claimedAt), + { + ...dispatcherAckCasOptions(id, 'in_flight', claimed.claimedBy, claimed.claimedAt), + context: DISPATCHER_SYSTEM_CONTEXT, + }, ); // Did the conditional write land? `IDataEngine.update` declares its @@ -399,7 +407,7 @@ export class SqlHttpOutbox implements IHttpOutbox { const after = (await this.engine.findOne(this.objectName, { where: { id }, fields: ['status', 'attempts'], - })) as Pick | null; + }, { context: DISPATCHER_SYSTEM_CONTEXT })) as Pick | null; if (!after || after.status !== patch.status || (after.attempts ?? 0) !== attempts) { throw new HttpAckError(httpAckLostClaimMessage(id, after?.status ?? 'unknown'), 'DELIVERY_NOT_ELIGIBLE'); } @@ -414,7 +422,7 @@ export class SqlHttpOutbox implements IHttpOutbox { const current = (await this.engine.findOne(this.objectName, { where: { id }, fields: ['attempts'], - })) as { attempts?: number } | null; + }, { context: DISPATCHER_SYSTEM_CONTEXT })) as { attempts?: number } | null; if (!current) return; await this.engine.update( @@ -423,7 +431,7 @@ export class SqlHttpOutbox implements IHttpOutbox { // Single-record write, audited under the `update` op. Declared a // global-sweep site — no request context reaches it. Warrant in // `outbox-dispatcher-scope.ts`. - dispatcherAckOptions(id), + { ...dispatcherAckOptions(id), context: DISPATCHER_SYSTEM_CONTEXT }, ); } @@ -431,7 +439,11 @@ export class SqlHttpOutbox implements IHttpOutbox { const where: Record = {}; if (filter?.status) where.status = filter.status; if (filter?.source) where.source = filter.source; - const rows = (await this.engine.find(this.objectName, { where })) as DeliveryRow[]; + const rows = (await this.engine.find( + this.objectName, + { where }, + { context: OUTBOX_SYSTEM_CONTEXT }, + )) as DeliveryRow[]; return rows.map((r) => this.toDelivery(r)); } diff --git a/packages/services/service-messaging/src/sql-outbox.ts b/packages/services/service-messaging/src/sql-outbox.ts index 464d2e00346..60f597aa658 100644 --- a/packages/services/service-messaging/src/sql-outbox.ts +++ b/packages/services/service-messaging/src/sql-outbox.ts @@ -14,7 +14,12 @@ import type { } from './outbox.js'; import { hashPartition } from './backoff.js'; import { toEpochMs } from './audit-timestamp.js'; -import { DISPATCHER_SYSTEM_CONTEXT, dispatcherAckCasOptions, dispatcherSweepOptions } from './outbox-dispatcher-scope.js'; +import { + DISPATCHER_SYSTEM_CONTEXT, + OUTBOX_SYSTEM_CONTEXT, + dispatcherAckCasOptions, + dispatcherSweepOptions, +} from './outbox-dispatcher-scope.js'; import { NotificationAckError, notificationAckLostClaimMessage, @@ -92,7 +97,13 @@ export class SqlNotificationOutbox implements INotificationOutbox { recipient_id: input.recipientId, channel: input.channel, }; - const existing = await this.engine.findOne(this.objectName, { where: dedup, fields: ['id'] }); + // The explicit system opt-in on every call of the producer side — see + // OUTBOX_SYSTEM_CONTEXT. + const existing = await this.engine.findOne( + this.objectName, + { where: dedup, fields: ['id'] }, + { context: OUTBOX_SYSTEM_CONTEXT }, + ); if (existing?.id) return String(existing.id); const id = randomUUID(); @@ -120,11 +131,15 @@ export class SqlNotificationOutbox implements INotificationOutbox { updated_at: now, }; try { - await this.engine.insert(this.objectName, row); + await this.engine.insert(this.objectName, row, { context: OUTBOX_SYSTEM_CONTEXT }); return id; } catch (err) { // Unique-index collision (dedup race) → return the winner. - const winner = await this.engine.findOne(this.objectName, { where: dedup, fields: ['id'] }); + const winner = await this.engine.findOne( + this.objectName, + { where: dedup, fields: ['id'] }, + { context: OUTBOX_SYSTEM_CONTEXT }, + ); if (winner?.id) return String(winner.id); throw err; } @@ -222,10 +237,12 @@ export class SqlNotificationOutbox implements INotificationOutbox { if (typeof claimed.claimedBy !== 'string' || typeof claimed.claimedAt !== 'number') { throw new NotificationAckError(notificationAckNoCredentialMessage(id), 'DELIVERY_NOT_ELIGIBLE'); } + // Every call of the ack carries the dispatcher's opt-in — see + // DISPATCHER_SYSTEM_CONTEXT. const current = (await this.engine.findOne(this.objectName, { where: { id }, fields: ['status', 'attempts', 'claimed_by', 'claimed_at'], - })) as { + }, { context: DISPATCHER_SYSTEM_CONTEXT })) as { status?: DeliveryStatus; attempts?: number; claimed_by?: string | null; @@ -306,7 +323,10 @@ export class SqlNotificationOutbox implements INotificationOutbox { // Predicate write (`updateMany`), audited under that op. Declared a // global-sweep site — no request context exists on the tick that // reaches here. Warrant in `outbox-dispatcher-scope.ts`. - dispatcherAckCasOptions(id, 'in_flight', claimed.claimedBy, claimed.claimedAt) as any, + { + ...dispatcherAckCasOptions(id, 'in_flight', claimed.claimedBy, claimed.claimedAt), + context: DISPATCHER_SYSTEM_CONTEXT, + } as any, ); // Did the conditional write land? `IDataEngine.update` declares its @@ -322,7 +342,7 @@ export class SqlNotificationOutbox implements INotificationOutbox { const after = (await this.engine.findOne(this.objectName, { where: { id }, fields: ['status', 'attempts'], - })) as { status?: DeliveryStatus; attempts?: number } | null; + }, { context: DISPATCHER_SYSTEM_CONTEXT })) as { status?: DeliveryStatus; attempts?: number } | null; if (!after || after.status !== status || (after.attempts ?? 0) !== attempts) { throw new NotificationAckError( notificationAckLostClaimMessage(id, after?.status ?? 'unknown'), @@ -353,7 +373,11 @@ export class SqlNotificationOutbox implements INotificationOutbox { const where: Record = {}; if (filter?.status) where.status = filter.status; if (filter?.notificationId) where.notification_id = filter.notificationId; - const rows = (await this.engine.find(this.objectName, { where })) as DeliveryRow[]; + const rows = (await this.engine.find( + this.objectName, + { where }, + { context: OUTBOX_SYSTEM_CONTEXT }, + )) as DeliveryRow[]; return rows.map((r) => this.toRecord(r)); } diff --git a/packages/services/service-messaging/src/template-renderer.ts b/packages/services/service-messaging/src/template-renderer.ts index 3d5084c71d1..df3c24c99ef 100644 --- a/packages/services/service-messaging/src/template-renderer.ts +++ b/packages/services/service-messaging/src/template-renderer.ts @@ -1,6 +1,7 @@ // Copyright (c) 2025 ObjectStack. Licensed under the Apache-2.0 license. import type { IDataEngine } from '@objectstack/spec/contracts'; +import { FAN_OUT_SYSTEM_CONTEXT } from './fan-out-system-context.js'; /** The object notification templates live in. */ export const TEMPLATE_OBJECT = 'sys_notification_template'; @@ -116,10 +117,12 @@ export class NotificationTemplateStore { const candidates = localeCandidates(locale); for (const loc of candidates) { try { + // The explicit system opt-in — see FAN_OUT_SYSTEM_CONTEXT: the + // platform's own template row, read to render a delivery. const row = await data.findOne(this.objectName, { where: { topic, channel, locale: loc, is_active: true }, fields: ['subject', 'body', 'format'], - }); + }, { context: FAN_OUT_SYSTEM_CONTEXT }); if (row) return row as NotificationTemplateRow; } catch { return null; // best-effort — fall back to generic rendering diff --git a/packages/services/service-settings/src/settings-service-plugin.ts b/packages/services/service-settings/src/settings-service-plugin.ts index d468299688f..c15ba9a4583 100644 --- a/packages/services/service-settings/src/settings-service-plugin.ts +++ b/packages/services/service-settings/src/settings-service-plugin.ts @@ -392,12 +392,20 @@ export class SettingsServicePlugin implements Plugin { * `SettingsSecretStore`. The store bypasses the tenant audit * warning because secrets are scoped through their owning * `sys_setting` row (which already carries the tenant context). + * + * [#21908] Every call carries the explicit system opt-in, as `delete` below + * already did: `sys_secret` is a platform-owned cipher store, and each + * call is the settings service acting on its own row after its own + * capability and lock gates passed. Without it these calls reach the data + * engine with no principal and no opt-in — the principal-less hand-off + * ADR-0096 D5 closes — and every encrypted setting would stop reading back + * once that hand-off denies. */ private buildSecretStore(engine: IDataEngine): SettingsSecretStore { const eng: any = engine; return { async insert(row) { - await eng.insert('sys_secret', row, { bypassTenantAudit: true }); + await eng.insert('sys_secret', row, { bypassTenantAudit: true, context: { isSystem: true } }); return { id: row.id }; }, async get(id) { @@ -405,6 +413,7 @@ export class SettingsServicePlugin implements Plugin { where: { id }, limit: 1, bypassTenantAudit: true, + context: { isSystem: true }, }); const row = Array.isArray(rows) ? rows[0] : rows?.data?.[0]; return row ?? null; @@ -415,10 +424,16 @@ export class SettingsServicePlugin implements Plugin { // `options.where.id`). Passing `{ where, data, ... }` as the // data argument left id=undefined and tripped // "Update requires an ID or options.multi=true". + // + // [#21908] Under the opt-in the engine's `readonly` strip no longer + // runs on this write, so a re-wrap's `ciphertext` (declared + // `readonly`) is written, which is what this member promises; without + // a context the strip took it. No caller in this repository reaches + // `update` today. await eng.update( 'sys_secret', { id, ...patch }, - { bypassTenantAudit: true }, + { bypassTenantAudit: true, context: { isSystem: true } }, ); }, async delete(id) { diff --git a/packages/services/service-settings/src/settings-service.ts b/packages/services/service-settings/src/settings-service.ts index 745d1930ec9..0a13d9abbf9 100644 --- a/packages/services/service-settings/src/settings-service.ts +++ b/packages/services/service-settings/src/settings-service.ts @@ -64,8 +64,9 @@ const DEFAULT_OBJECT = 'sys_setting'; * gates. See {@link SettingsService.upsertRow} for the full argument and for * why the field stays `readonly` for everybody else. * - * The reads ({@link SettingsService.loadRows}, `upsertRow`'s existence probe) - * and `upsertRow`'s insert carry it too. They are plumbing: the door in front + * The reads ({@link SettingsService.loadRows}, `upsertRow`'s existence probe, + * and [#21908] `readStoredHandle`'s re-read of the stored handle) and + * `upsertRow`'s insert carry it too. They are plumbing: the door in front * of them, when there is one, has already authorized the caller, and * `loadRows` runs on every request's execution-context build. Without it they * reach the data engine with no principal and no system opt-in — the @@ -2637,7 +2638,14 @@ export class SettingsService { ): Promise<{ found: boolean; handle: string | null }> { if (this.engine) { const { where, bypass } = this.rowIdentity(row); - const rows = await this.engine.find(this.objectName, { where, limit: 1, ...bypass }); + // The explicit system opt-in, as the write it verifies carries: see + // SETTINGS_SYSTEM_CONTEXT. + const rows = await this.engine.find(this.objectName, { + where, + limit: 1, + ...bypass, + context: SETTINGS_SYSTEM_CONTEXT, + }); const current = Array.isArray(rows) ? rows[0] : undefined; if (!current) return { found: false, handle: null }; return { diff --git a/packages/services/service-storage/src/metadata-store.ts b/packages/services/service-storage/src/metadata-store.ts index aaedd7b2240..05bdd6e7dca 100644 --- a/packages/services/service-storage/src/metadata-store.ts +++ b/packages/services/service-storage/src/metadata-store.ts @@ -130,6 +130,27 @@ function writeOptionsFor( return { context: { tenantId: organizationId } }; } +/** + * [#21908] The execution context {@link StorageMetadataStore.createFile}'s + * `sys_file` insert runs under: the explicit system opt-in, taken inside this + * store (maintainer ruling Q1), merged with the acting organization's + * `tenantId` when there is one. + * + * Before this the insert reached the data engine with a tenant-only context + * (or none), no principal and no opt-in, and passed the security middleware + * only through its principal-less hand-off (ADR-0096 E1), which D5 closes. + * Carrying the caller's principal instead is not open: no member grant exists + * on `sys_file`, so uploads would break. + * + * What it keeps is the door-derived scope, unchanged: the row is stamped + * `owner_id` with the uploading user the door resolved from the session, and + * `tenantId` still reaches `DriverOptions.tenantId`, where the driver's + * `injectTenantOnInsert` stamps the organization exactly as before + * ({@link StorageWriteContext}). ⛔ Never drop the `tenantId` from this + * context: under the opt-in no other layer stamps the organization. + */ +const CREATE_FILE_SYSTEM_CONTEXT = { isSystem: true } as const; + /** * Persisted upload-session record (matches `sys_upload_session` object schema). */ @@ -320,6 +341,11 @@ export class StorageMetadataStore { const now = new Date().toISOString(); const full: FileRecord = { created_at: now, updated_at: now, ...rec }; const options = writeOptionsFor(context); + // [#21908] The engine insert carries the explicit system opt-in, keeping the + // acting organization beside it — see CREATE_FILE_SYSTEM_CONTEXT. + const insertOptions = { + context: { ...(options?.context ?? {}), ...CREATE_FILE_SYSTEM_CONTEXT }, + }; if (!this.engine) { // The engine-absent stand-in has no schema to ask, so it records what it // was told rather than deriving a column: a no-engine deployment has no @@ -334,7 +360,7 @@ export class StorageMetadataStore { return stamped; } await this.engineOp('sys_file', 'insert', FILE_INSERT_CONSEQUENCE, (engine) => - engine.insert('sys_file', full, options), + engine.insert('sys_file', full, insertOptions), ); return full; } From 2bea8b684b6a2766bb7cb480f12eb750e66cbaa2 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 16:20:36 +0000 Subject: [PATCH 2/5] wip: pins for the stage-1 producers; createSession moves with createFile Claude-Session: https://claude.ai/code/session_01WMQprn46CND82KmY8sZWBu Co-authored-by: Claude --- ...tocol.platform-store-system-opt-in.test.ts | 35 +++ .../src/system-context.pin.test.ts | 272 +++++++++++++++++- .../src/settings-system-context.pin.test.ts | 85 +++++- .../service-storage/src/metadata-store.ts | 47 +-- .../sys-file-organization-stamping.test.ts | 61 +++- ...load-session-organization-stamping.test.ts | 25 +- ...t-audit-update-delete-half-repairs.test.ts | 14 + 7 files changed, 483 insertions(+), 56 deletions(-) diff --git a/packages/metadata-protocol/src/protocol.platform-store-system-opt-in.test.ts b/packages/metadata-protocol/src/protocol.platform-store-system-opt-in.test.ts index ad0c930b461..43154d9ea11 100644 --- a/packages/metadata-protocol/src/protocol.platform-store-system-opt-in.test.ts +++ b/packages/metadata-protocol/src/protocol.platform-store-system-opt-in.test.ts @@ -28,6 +28,7 @@ import { describe, expect, it } from 'vitest'; import { assertEngineDeleteDispatch, assertEngineUpdateDispatch, assertEngineFindOnePredicate } from '@objectstack/metadata-core'; import { ObjectStackProtocolImplementation } from './protocol.js'; +import { SysMetadataRepository } from './sys-metadata-repository.js'; interface Row { id: string; @@ -294,3 +295,37 @@ describe('platform-store calls carry the explicit system opt-in (#21911)', () => }); }); }); + +describe('[#21908] row 23 — the repository reads the engine-lane slice left carry the opt-in', () => { + it('SysMetadataRepository.getByHash, list, history and watch’s replay', async () => { + const { engine, calls, historyRows } = makeStubEngine(); + const protocol = new ObjectStackProtocolImplementation(engine); + await protocol.saveMetaItem({ + type: 'view', name: 'proj_task_grid', item: viewBody('proj_task_grid'), mode: 'publish', + }); + // The population the reads below must find, written by the save above. + expect(historyRows.length).toBeGreaterThan(0); + const hash = String(historyRows[0].checksum); + const ref = { type: 'view', name: 'proj_task_grid' } as any; + const repo = new SysMetadataRepository({ engine, organizationId: null }); + + await expectSystemOptIn(calls, 'getByHash', async () => { + expect(await repo.getByHash(ref, hash)).not.toBeNull(); + }); + await expectSystemOptIn(calls, 'list', async () => { + const headers: unknown[] = []; + for await (const h of repo.list({ type: 'view' } as any)) headers.push(h); + expect(headers).toHaveLength(1); + }); + await expectSystemOptIn(calls, 'history', async () => { + const events: unknown[] = []; + for await (const e of repo.history(ref)) events.push(e); + expect(events.length).toBeGreaterThan(0); + }); + await expectSystemOptIn(calls, 'replayFromHistory', async () => { + const it = repo.watch({} as any, 0)[Symbol.asyncIterator](); + expect((await it.next()).done).toBe(false); + await it.return?.(); + }); + }); +}); diff --git a/packages/services/service-messaging/src/system-context.pin.test.ts b/packages/services/service-messaging/src/system-context.pin.test.ts index 3442a014053..aad0fae1f4c 100644 --- a/packages/services/service-messaging/src/system-context.pin.test.ts +++ b/packages/services/service-messaging/src/system-context.pin.test.ts @@ -26,9 +26,13 @@ import { SqlNotificationOutbox } from './sql-outbox.js'; import { SqlHttpOutbox } from './sql-http-outbox.js'; import { MessagingService } from './messaging-service.js'; import { createInboxChannel } from './inbox-channel.js'; +import { RecipientResolver } from './recipient-resolver.js'; +import { createEmailChannel } from './email-channel.js'; +import { createSmsChannel } from './sms-channel.js'; +import { NotificationTemplateStore } from './template-renderer.js'; import { assertEngineFindOnePredicate, assertEngineUpdateDispatch } from '@objectstack/metadata-core'; -type Call = { verb: string; object: string; context: unknown }; +type Call = { verb: string; object: string; context: unknown; query?: any; data?: any; options?: any; answered?: unknown }; function silentLogger() { return { info: () => {}, warn: () => {}, error: () => {}, debug: () => {} }; @@ -44,21 +48,28 @@ function recordingEngine(answer: (verb: string, object: string, query: any) => u const ctxOf = (query: any, options: any) => options?.context ?? query?.context; const engine = { async find(object: string, query: any, options?: any) { - calls.push({ verb: 'find', object, context: ctxOf(query, options) }); - return (answer('find', object, query) as unknown[]) ?? []; + // Recorded before it is answered, so a read the double fails is counted too. + const call: Call = { verb: 'find', object, context: ctxOf(query, options), query }; + calls.push(call); + const rows = (answer('find', object, query) as unknown[]) ?? []; + // [#21908] The caller's bound holds here as it does on the engine. + call.answered = typeof query?.limit === 'number' ? rows.slice(0, query.limit) : rows; + return call.answered; }, async findOne(object: string, query: any, options?: any) { assertEngineFindOnePredicate(object, query); - calls.push({ verb: 'findOne', object, context: ctxOf(query, options) }); - return answer('findOne', object, query) ?? null; + const call: Call = { verb: 'findOne', object, context: ctxOf(query, options), query }; + calls.push(call); + call.answered = answer('findOne', object, query) ?? null; + return call.answered; }, async insert(object: string, data: Record, options?: any) { - calls.push({ verb: 'insert', object, context: options?.context }); + calls.push({ verb: 'insert', object, context: options?.context, data, options }); return { id: `${object}_1`, ...data }; }, async update(object: string, data: Record, options?: any) { assertEngineUpdateDispatch(data, options); - calls.push({ verb: 'update', object, context: options?.context }); + calls.push({ verb: 'update', object, context: options?.context, data, options }); return 1; }, async count() { return 0; }, @@ -162,3 +173,250 @@ describe('[#21913] the emit fan-out carries the explicit system opt-in', () => { expectAllSystem(calls); }); }); + +// --------------------------------------------------------------------------- +// [#21908] Stage 1 of the closure: the producers the services-lane slice left. +// --------------------------------------------------------------------------- + +/** Plain-equality `where` over seeded rows — anything else is refused, loudly. */ +function matchRows(rows: Array>, where: Record = {}) { + for (const k of Object.keys(where)) { + if (k.startsWith('$')) throw new Error(`recording engine: unimplemented combinator ${k}`); + } + return rows.filter((r) => Object.entries(where).every(([k, v]) => (r[k] ?? null) === (v ?? null))); +} + +const USER_SCOPED = new Set(['sys_inbox_message', 'sys_notification_receipt']); + +/** + * Two users' inbox rows and receipts, and the events they are about. Every read + * answers only what its `where` selects, so a call that dropped its `user_id` + * term would be answered the other user's rows — and the negative pin sees it. + */ +function twoUserInbox() { + const rows: Record>> = { + sys_inbox_message: [ + { id: 'm1', user_id: 'u_a', notification_id: 'n1', topic: 't', title: 'a1', created_at: '2026-01-03' }, + { id: 'm2', user_id: 'u_b', notification_id: 'n2', topic: 't', title: 'b1', created_at: '2026-01-02' }, + { id: 'm3', user_id: 'u_a', notification_id: 'n3', topic: 't', title: 'a2', created_at: '2026-01-01' }, + ], + sys_notification_receipt: [ + { id: 'r1', notification_id: 'n1', user_id: 'u_a', channel: 'inbox', state: 'delivered' }, + { id: 'r2', notification_id: 'n2', user_id: 'u_b', channel: 'inbox', state: 'delivered' }, + ], + sys_notification: [ + { id: 'n1', organization_id: 'org_1' }, + { id: 'n2', organization_id: 'org_2' }, + { id: 'n3', organization_id: 'org_1' }, + ], + }; + return recordingEngine((verb, object, query) => { + const hits = matchRows(rows[object] ?? [], query?.where); + return verb === 'find' ? hits : (hits[0] ?? null); + }); +} + +async function driveInboxAs(engine: any, userId: string) { + const service = new MessagingService({ logger: silentLogger(), getData: () => engine } as any); + // A saturated window, so `countUnreadTotal` runs too. + const listed = await service.listInbox(userId, { limit: 2 }); + // n1 has a receipt (the update branch); n3 has none (the insert branch, + // which reads the event's organization). + await service.markRead(userId, ['n1']); + await service.markRead(userId, ['n3']); + const swept = await service.markAllRead(userId); + return { listed, swept }; +} + +describe('[#21908] Q1 — the inbox read-state producers: the opt-in inside the service, the door-derived user scope kept', () => { + it('every call carries isSystem, and every call on the user’s rows carries the user id in its where or its stamp', async () => { + const { engine, calls } = twoUserInbox(); + const { listed } = await driveInboxAs(engine, 'u_a'); + expect(listed.notifications.map((n) => n.id)).toEqual(['n1', 'n3']); + + // The population first: every producer the ruling names actually ran — + // listInbox's window and countUnreadTotal, readReceiptStates, + // unreadNotificationIds, upsertReadReceipt's read, update and insert, + // and notificationOrganization. + const seen = (verb: string, object: string) => calls.filter((c) => c.verb === verb && c.object === object); + expect(seen('find', 'sys_inbox_message').length).toBeGreaterThanOrEqual(3); + expect(seen('find', 'sys_notification_receipt').length).toBeGreaterThanOrEqual(2); + expect(seen('findOne', 'sys_notification_receipt').length).toBeGreaterThanOrEqual(2); + // (The double keeps no write state, so mark-all-read re-flips n1 and + // re-inserts n3's receipt: each branch runs at least once.) + expect(seen('update', 'sys_notification_receipt').length).toBeGreaterThanOrEqual(1); + expect(seen('insert', 'sys_notification_receipt').length).toBeGreaterThanOrEqual(1); + expect(seen('findOne', 'sys_notification').length).toBeGreaterThanOrEqual(1); + expectAllSystem(calls); + + // Reads of the user's rows: keyed on the door-derived user id. + for (const c of calls.filter((x) => (x.verb === 'find' || x.verb === 'findOne') && USER_SCOPED.has(x.object))) { + expect(c.query?.where?.user_id, `${c.verb} on ${c.object}`).toBe('u_a'); + } + // The receipt it inserts: stamped with it. + for (const insert of seen('insert', 'sys_notification_receipt')) { + expect(insert.data).toMatchObject({ user_id: 'u_a', notification_id: 'n3', organization_id: 'org_1' }); + } + // The receipt it updates: addressed by an id a user-keyed read returned, and by nothing else. + const returnedIds = seen('findOne', 'sys_notification_receipt').map((c) => (c.answered as any)?.id).filter(Boolean); + for (const update of seen('update', 'sys_notification_receipt')) { + expect(update.options.where).toEqual({ id: 'r1' }); + } + expect(returnedIds).toContain('r1'); + // The event read: one row by id, two columns, used only as the stamp. + for (const event of seen('findOne', 'sys_notification')) { + expect(event.query).toEqual({ where: { id: 'n3' }, fields: ['id', 'organization_id'] }); + } + }); + + it('⛔ negative: no call made for one user reads, writes or addresses another user’s rows', async () => { + const { engine, calls } = twoUserInbox(); + await driveInboxAs(engine, 'u_a'); + + // Every row any read answered belongs to the caller… + for (const c of calls.filter((x) => USER_SCOPED.has(x.object) && (x.verb === 'find' || x.verb === 'findOne'))) { + const answered = c.verb === 'find' ? (c.answered as any[]) : [c.answered].filter(Boolean); + for (const row of answered) expect(row.user_id, `${c.verb} on ${c.object}`).toBe('u_a'); + } + // …and nothing is written as, or onto, u_b: not their receipt r2, not their inbox row. + for (const c of calls.filter((x) => x.verb === 'insert' || x.verb === 'update')) { + expect(c.data?.user_id ?? 'u_a').toBe('u_a'); + expect(c.options?.where?.id).not.toBe('r2'); + } + }); +}); + +describe('[#21908] Q2 — `owner_of:` carries the opt-in, reads only the owner fields, and yields only recipient ids', () => { + it('resolveOwnerOf: isSystem, a projection of id plus the owner fields, and the owner id alone back', async () => { + const { engine, calls } = recordingEngine((verb, object) => + verb === 'findOne' && object === 'deal' + ? { id: 'd1', owner_id: 'u_owner', name: 'Confidential deal', amount: 900_000 } + : null); + const resolver = new RecipientResolver({ getData: () => engine, logger: silentLogger() }); + + const ids = await resolver.resolve(['owner_of:deal:d1', { ownerOf: { object: 'deal', id: 'd1' } } as any]); + + // Only the recipient id leaves the resolver — never a value the record carried. + expect(ids).toEqual(['u_owner']); + const reads = calls.filter((c) => c.object === 'deal'); + expect(reads).toHaveLength(2); + for (const read of reads) { + expect(read.context).toEqual({ isSystem: true }); + expect(read.query).toEqual({ + where: { id: 'd1' }, + fields: ['id', 'owner_id', 'assigned_to', 'assignee_id', 'owner', 'assignee'], + }); + } + }); + + it('nothing the read returns reaches the emitter: emit() answers recipient ids and counts only', async () => { + const { engine } = recordingEngine((verb, object) => + verb === 'findOne' && object === 'deal' + ? { id: 'd1', owner_id: 'u_owner', name: 'Confidential deal' } + : (verb === 'find' ? [] : null)); + const service = new MessagingService({ logger: silentLogger(), getData: () => engine } as any); + service.registerChannel(createInboxChannel({ getData: () => engine })); + + const result = await service.emit({ topic: 'deal.won', audience: ['owner_of:deal:d1'], channels: ['inbox'], title: 't', body: 'b' } as any); + + expect(result.delivered).toBe(1); + expect(JSON.stringify(result)).not.toContain('Confidential deal'); + }); +}); + +describe('[#21908] rows 24 onward — the remaining fan-out and outbox calls carry the explicit system opt-in', () => { + it('resolveRole and resolveTeam', async () => { + const { engine, calls } = recordingEngine((verb) => (verb === 'find' ? [{ user_id: 'u1' }] : null)); + const resolver = new RecipientResolver({ getData: () => engine, logger: silentLogger() }); + expect(await resolver.resolve(['role:admin', 'team:t1'], { organizationId: 'org_1' })).toEqual(['u1']); + expect(calls.map((c) => c.object)).toEqual(['sys_member', 'sys_team_member']); + expectAllSystem(calls); + }); + + it('emit()’s dedup lookup', async () => { + const { engine, calls } = recordingEngine((verb, object, query) => + verb === 'findOne' && object === 'sys_notification' && query?.where?.dedup_key ? { id: 'n_prior' } : null); + const service = new MessagingService({ logger: silentLogger(), getData: () => engine } as any); + const result = await service.emit({ topic: 't', audience: ['u1'], channels: ['inbox'], dedupKey: 'k1', title: 't', body: 'b' } as any); + expect(result).toMatchObject({ notificationId: 'n_prior', deduped: true }); + expect(calls).toHaveLength(1); + expectAllSystem(calls); + }); + + it('SqlNotificationOutbox: enqueue (dedup read, insert, race read-back), ack (state read, CAS write, read-back) and list', async () => { + let inserted = false; + const { engine, calls } = recordingEngine((verb, _object, query) => { + if (verb === 'findOne' && query?.where?.notification_id) return inserted ? { id: 'd_winner' } : null; + if (verb === 'findOne' && Array.isArray(query?.fields) && query.fields.includes('claimed_by')) { + return { status: 'in_flight', attempts: 0, claimed_by: 'node-a', claimed_at: NOW }; + } + if (verb === 'findOne') return { status: 'success', attempts: 1 }; + return []; + }); + // The insert loses a dedup race once, so the winner read-back runs too. + const realInsert = engine.insert; + engine.insert = async (...args: any[]) => { await realInsert(...args); inserted = true; throw new Error('unique violation'); }; + const outbox = new SqlNotificationOutbox(engine, { partitionCount: 1 }); + + expect(await outbox.enqueue({ notificationId: 'n1', recipientId: 'u1', channel: 'inbox' } as any)).toBe('d_winner'); + await outbox.ack({ id: 'd1', claimedBy: 'node-a', claimedAt: NOW } as any, { success: true } as any); + await outbox.list({ status: 'success' }); + + expect(calls.map((c) => c.verb)).toEqual(['findOne', 'insert', 'findOne', 'findOne', 'update', 'findOne', 'find']); + expectAllSystem(calls); + }); + + it('SqlHttpOutbox: enqueue and recordUndeliverable, ack (both arities) and list', async () => { + const { engine, calls } = recordingEngine((verb, _object, query) => { + if (verb === 'findOne' && query?.where?.dedup_key) return null; + if (verb === 'findOne' && Array.isArray(query?.fields) && query.fields.includes('claimed_by')) { + return { status: 'in_flight', attempts: 0, claimed_by: 'node-a', claimed_at: NOW }; + } + if (verb === 'findOne' && Array.isArray(query?.fields) && query.fields.length === 1) return { attempts: 0 }; + if (verb === 'findOne') return { status: 'success', attempts: 1 }; + return []; + }); + const outbox = new SqlHttpOutbox(engine, { partitionCount: 1 }); + const base = { source: 'webhook', refId: 'wh1', url: 'https://example.test/hook', payload: {} }; + + await outbox.enqueue({ ...base, dedupKey: 'k1' } as any); + await outbox.recordUndeliverable({ ...base, dedupKey: 'k2', reason: 'no secret' } as any); + await outbox.ack('h1', { success: true } as any, { claimedBy: 'node-a', claimedAt: NOW } as any); + await outbox.ack('h2', { success: false, error: 'x' } as any); + await outbox.list({ source: 'webhook' }); + + expect(calls.map((c) => c.verb)).toEqual([ + 'findOne', 'insert', 'findOne', 'insert', // enqueue, recordUndeliverable + 'findOne', 'update', 'findOne', // the credentialed ack + 'findOne', 'update', // the by-id ack + 'find', // list + ]); + expect(calls.every((c) => c.object === 'sys_http_delivery')).toBe(true); + expectAllSystem(calls); + }); + + it('the email and SMS recipient reads — the first read and the address-only retry — and the template read', async () => { + let failFirst = true; + const { engine, calls } = recordingEngine((verb, object, query) => { + if (verb !== 'findOne') return []; + if (object === 'sys_user' && failFirst && query?.fields?.length === 2) { + failFirst = false; + throw new Error('no locale column'); + } + if (object === 'sys_user') return { email: 'ada@example.com', phone_number: '+15555550100' }; + return null; + }); + const store = new NotificationTemplateStore({ getData: () => engine }); + const email = createEmailChannel({ getEmail: () => ({ async send() { return { id: 'e1' }; } }) as any, getData: () => engine, store }); + const sms = createSmsChannel({ getSms: () => ({ async send() { return { id: 's1' }; } }) as any, getData: () => engine, store }); + const notification = { title: 't', body: 'b', recipients: ['u1'], topic: 'deal.won' }; + + expect((await email.send({ logger: silentLogger() }, { notification, channel: 'email', recipient: 'u1' } as any)).ok).toBe(true); + expect((await sms.send({ logger: silentLogger() }, { notification, channel: 'sms', recipient: 'u1' } as any)).ok).toBe(true); + + // email: failed read + retry; sms: one read; and the template reads. + expect(calls.filter((c) => c.object === 'sys_user')).toHaveLength(3); + expect(calls.filter((c) => c.object === 'sys_notification_template').length).toBeGreaterThanOrEqual(2); + expectAllSystem(calls); + }); +}); diff --git a/packages/services/service-settings/src/settings-system-context.pin.test.ts b/packages/services/service-settings/src/settings-system-context.pin.test.ts index f6ca92a5d8e..707adb67901 100644 --- a/packages/services/service-settings/src/settings-system-context.pin.test.ts +++ b/packages/services/service-settings/src/settings-system-context.pin.test.ts @@ -18,8 +18,10 @@ */ import { describe, it, expect } from 'vitest'; +import { assertEngineDeleteDispatch, assertEngineUpdateDispatch } from '@objectstack/objectql'; import { SettingsService } from './settings-service.js'; -import { buildSettingAuditWriter, wrapEngineAsSettingsEngine } from './settings-service-plugin.js'; +import { SettingsServicePlugin, buildSettingAuditWriter, wrapEngineAsSettingsEngine } from './settings-service-plugin.js'; +import { LocalCryptoProvider } from './local-crypto-provider.js'; type Call = { verb: 'find' | 'insert' | 'update' | 'delete'; object: string; context: unknown }; @@ -35,8 +37,10 @@ function matches(row: Record, where: Record): /** * An `IDataEngine`-shaped double that records the context each call carried. - * It implements only the verbs the moved calls use — `find` and `insert` — so - * a call this pin does not expect fails loudly instead of being answered. + * It implements only the verbs the moved calls use — `find` and `insert`, and + * [#21908] the `update` / `delete` a rotation and the `sys_secret` store issue, + * each routed through the engine's own dispatch predicate — so a call this pin + * does not expect fails loudly instead of being answered. */ function recordingEngine() { const rows: Array> = []; @@ -49,14 +53,26 @@ function recordingEngine() { calls.push({ verb: 'find', object, context: readContext(query, options) }); // The user-scope insert first proves its `user_id` names a user (#21913). if (object === 'sys_user') return [{ id: query?.where?.id }]; - const hits = rows.filter((r) => matches(r, query?.where ?? {})); + const hits = rows.filter((r) => (r.__object ?? 'sys_setting') === object && matches(r, query?.where ?? {})); return typeof query?.limit === 'number' ? hits.slice(0, query.limit) : hits; }, async insert(object: string, data: Record, options?: any) { calls.push({ verb: 'insert', object, context: options?.context }); - if (object === 'sys_setting') rows.push({ ...data }); + if (object === 'sys_setting' || object === 'sys_secret') rows.push({ __object: object, ...data }); return { ...data }; }, + async update(object: string, data: Record, options?: any) { + assertEngineUpdateDispatch(data, options); + calls.push({ verb: 'update', object, context: options?.context }); + const where = options?.where ?? { id: data.id }; + for (const row of rows.filter((r) => r.__object === object && matches(r, where))) Object.assign(row, data); + return 1; + }, + async delete(object: string, options?: any) { + assertEngineDeleteDispatch(options); + calls.push({ verb: 'delete', object, context: options?.context }); + return 1; + }, }; return { engine, calls, rows }; } @@ -109,3 +125,62 @@ describe('[#21913] SettingsService engine calls carry the explicit system opt-in expect(calls).toEqual([{ verb: 'insert', object: 'sys_setting_audit', context: { isSystem: true } }]); }); }); + +// --------------------------------------------------------------------------- +// [#21908] Stage 1 of the closure: the `sys_secret` store and the rotation's +// verification read. +// --------------------------------------------------------------------------- + +const SECRET_MANIFEST = { + namespace: 'sms', + version: 1, + label: 'SMS', + scope: 'tenant', + specifiers: [{ type: 'password', key: 'twilio_auth_token', label: 'Auth token', required: false, encrypted: true }], +} as any; + +describe('[#21908] the sys_secret store and readStoredHandle carry the explicit system opt-in', () => { + it('the store the plugin builds: insert, get, update and delete each pass isSystem', async () => { + const { engine, calls } = recordingEngine(); + const store = (new SettingsServicePlugin() as any).buildSecretStore(engine); + const row = { id: 'sec_1', namespace: 'sms', key: 'k', kms_key_id: 'local', alg: 'aes-256-gcm', version: 1, ciphertext: 'c1' }; + + await store.insert(row); + expect(await store.get('sec_1')).toMatchObject({ id: 'sec_1', ciphertext: 'c1' }); + await store.update('sec_1', { ciphertext: 'c2', version: 2 }); + await store.delete('sec_1'); + + expect(calls.map((c) => `${c.verb}:${c.object}`)).toEqual([ + 'insert:sys_secret', 'find:sys_secret', 'update:sys_secret', 'delete:sys_secret', + ]); + for (const call of calls) { + expect(call.context, `${call.verb} on ${call.object}`).toEqual({ isSystem: true }); + } + }); + + it('a rotation: the secret writes, the row update and the verification read after it all pass isSystem', async () => { + const { engine, calls } = recordingEngine(); + const svc = new SettingsService({ + env: {}, + engine: wrapEngineAsSettingsEngine(engine as any), + cryptoProvider: new LocalCryptoProvider(), + secretStore: (new SettingsServicePlugin() as any).buildSecretStore(engine), + } as any); + svc.registerManifest(SECRET_MANIFEST); + + await svc.set('sms', 'twilio_auth_token', 'alpha'); + await svc.set('sms', 'twilio_auth_token', 'beta'); + expect((await svc.get('sms', 'twilio_auth_token')).value).toBe('beta'); + + // The population: the second write updated the row, then re-read it + // (`readStoredHandle`) before it reaped the rotated-away secret. + const verbs = calls.map((c) => `${c.verb}:${c.object}`); + const updateAt = verbs.indexOf('update:sys_setting'); + expect(updateAt).toBeGreaterThan(-1); + expect(verbs.slice(updateAt + 1)).toEqual(expect.arrayContaining(['find:sys_setting', 'delete:sys_secret'])); + expect(verbs.filter((v) => v === 'insert:sys_secret')).toHaveLength(2); + for (const call of calls) { + expect(call.context, `${call.verb} on ${call.object}`).toEqual({ isSystem: true }); + } + }); +}); diff --git a/packages/services/service-storage/src/metadata-store.ts b/packages/services/service-storage/src/metadata-store.ts index 05bdd6e7dca..382b09b2612 100644 --- a/packages/services/service-storage/src/metadata-store.ts +++ b/packages/services/service-storage/src/metadata-store.ts @@ -131,25 +131,33 @@ function writeOptionsFor( } /** - * [#21908] The execution context {@link StorageMetadataStore.createFile}'s - * `sys_file` insert runs under: the explicit system opt-in, taken inside this - * store (maintainer ruling Q1), merged with the acting organization's - * `tenantId` when there is one. + * [#21908] The engine options the two INSERT doors of this store — + * {@link StorageMetadataStore.createFile} (`sys_file`) and + * {@link StorageMetadataStore.createSession} (`sys_upload_session`) — run + * under: the explicit system opt-in, taken inside this store (maintainer + * ruling Q1), with the acting organization's `tenantId` beside it when there + * is one. * - * Before this the insert reached the data engine with a tenant-only context + * Before this each insert reached the data engine with a tenant-only context * (or none), no principal and no opt-in, and passed the security middleware * only through its principal-less hand-off (ADR-0096 E1), which D5 closes. * Carrying the caller's principal instead is not open: no member grant exists - * on `sys_file`, so uploads would break. + * on `sys_file` / `sys_upload_session`, so uploads would break. * - * What it keeps is the door-derived scope, unchanged: the row is stamped - * `owner_id` with the uploading user the door resolved from the session, and - * `tenantId` still reaches `DriverOptions.tenantId`, where the driver's - * `injectTenantOnInsert` stamps the organization exactly as before - * ({@link StorageWriteContext}). ⛔ Never drop the `tenantId` from this - * context: under the opt-in no other layer stamps the organization. + * What it keeps is the door-derived scope, unchanged. The `sys_file` row is + * stamped `owner_id` with the uploading user the door resolved from the + * session; the `sys_upload_session` row carries no user column and is bound to + * that file by `file_id`. `tenantId` still reaches `DriverOptions.tenantId`, + * where the driver's `injectTenantOnInsert` stamps the organization exactly as + * before ({@link StorageWriteContext}). ⛔ Never drop the `tenantId` here: + * under the opt-in no other layer stamps the organization. */ -const CREATE_FILE_SYSTEM_CONTEXT = { isSystem: true } as const; +function systemInsertOptionsFor( + context?: StorageWriteContext, +): { context: { isSystem: true; tenantId?: string } } { + const tenant = writeOptionsFor(context)?.context; + return { context: { ...(tenant ?? {}), isSystem: true } }; +} /** * Persisted upload-session record (matches `sys_upload_session` object schema). @@ -341,11 +349,6 @@ export class StorageMetadataStore { const now = new Date().toISOString(); const full: FileRecord = { created_at: now, updated_at: now, ...rec }; const options = writeOptionsFor(context); - // [#21908] The engine insert carries the explicit system opt-in, keeping the - // acting organization beside it — see CREATE_FILE_SYSTEM_CONTEXT. - const insertOptions = { - context: { ...(options?.context ?? {}), ...CREATE_FILE_SYSTEM_CONTEXT }, - }; if (!this.engine) { // The engine-absent stand-in has no schema to ask, so it records what it // was told rather than deriving a column: a no-engine deployment has no @@ -359,8 +362,10 @@ export class StorageMetadataStore { this.files.set(stamped.id, stamped); return stamped; } + // [#21908] The explicit system opt-in, the organization kept beside it — + // see systemInsertOptionsFor. await this.engineOp('sys_file', 'insert', FILE_INSERT_CONSEQUENCE, (engine) => - engine.insert('sys_file', full, insertOptions), + engine.insert('sys_file', full, systemInsertOptionsFor(context)), ); return full; } @@ -486,8 +491,10 @@ export class StorageMetadataStore { this.sessions.set(stamped.id, stamped); return stamped; } + // [#21908] The explicit system opt-in, the organization kept beside it — + // see systemInsertOptionsFor. await this.engineOp('sys_upload_session', 'insert', SESSION_INSERT_CONSEQUENCE, (engine) => - engine.insert('sys_upload_session', full, options), + engine.insert('sys_upload_session', full, systemInsertOptionsFor(context)), ); return full; } diff --git a/packages/services/service-storage/src/sys-file-organization-stamping.test.ts b/packages/services/service-storage/src/sys-file-organization-stamping.test.ts index 3738b53c0ac..2656289ec9f 100644 --- a/packages/services/service-storage/src/sys-file-organization-stamping.test.ts +++ b/packages/services/service-storage/src/sys-file-organization-stamping.test.ts @@ -94,7 +94,8 @@ describe('StorageMetadataStore.createFile: the acting organization reaches the e // The platform's insert-side chokepoint reads exactly this: // `context.tenantId` → `buildDriverOptions` → `DriverOptions.tenantId` → // `SqlDriver.injectTenantOnInsert` → the object's tenant column. - expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_A' } }); + // [#21908] …beside the explicit system opt-in (ruling Q1). + expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); }); it('⛔ does NOT stamp the column onto the payload — that decision is the driver’s', async () => { @@ -110,7 +111,7 @@ describe('StorageMetadataStore.createFile: the acting organization reaches the e expect(engine._inserts[0].data).not.toHaveProperty('organization_id'); }); - it('passes NO options at all when there is no organization — the pre-#12745 shape', async () => { + it('carries the opt-in and NO tenant when there is no organization — ⛔ none is invented', async () => { const engine = createRecordingEngine(); const store = new StorageMetadataStore(engine); @@ -119,13 +120,14 @@ describe('StorageMetadataStore.createFile: the acting organization reaches the e await store.createFile(fileRec('f3'), { organizationId: '' }); await store.createFile(fileRec('f4'), { organizationId: null }); - // An empty context is still a context; handing one to the engine would - // change what every other option resolver on that call sees. + // [#21908] Every insert now carries the explicit system opt-in (ruling Q1), + // and still no `tenantId` where the caller has no organization: the opt-in + // is the only key, never a guessed organization. expect(engine._inserts.map((i: { options: unknown }) => i.options)).toEqual([ - undefined, - undefined, - undefined, - undefined, + { context: { isSystem: true } }, + { context: { isSystem: true } }, + { context: { isSystem: true } }, + { context: { isSystem: true } }, ]); }); @@ -225,13 +227,45 @@ describe('storage upload routes: a file created with a session lands with that o // The owner was already threaded before this card; the organization is // what was being dropped two lines away from it. expect(insert.data.owner_id).toBe('u1'); - expect(insert.options).toEqual({ context: { tenantId: 'org_A' } }); + expect(insert.options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); } finally { await fs.rm(rootDir, { recursive: true, force: true }); } }); } + it('[#21908] ⛔ the owner stamp is the door-derived user alone — a body naming another user does not reach the insert', async () => { + // Under the explicit system opt-in (ruling Q1) RLS no longer reads this + // insert, so the user the door resolved from the session is the boundary. + const rootDir = join(tmpdir(), `os-21908-owner-${Math.random().toString(36).slice(2)}`); + await fs.mkdir(rootDir, { recursive: true }); + try { + const adapter = new LocalStorageAdapter({ rootDir, signingSecret: 'test-secret' }); + const engine = createRecordingEngine(); + const store = new StorageMetadataStore(engine); + const httpServer = createMockHttpServer(); + registerStorageRoutes(httpServer as any, adapter, store, { + basePath: '/api/v1/storage', + resolveSession: async () => ({ userId: 'u1', organizationId: 'org_A' }), + }); + + const handler = httpServer._getHandler('POST', '/api/v1/storage/upload/presigned')!; + const res = createMockRes(); + await handler( + createMockReq({ body: { filename: 'a.txt', mimeType: 'text/plain', size: 3, owner_id: 'u_other', organization_id: 'org_B' } }), + res, + ); + + expect(res._status).toBe(200); + const insert = engine._inserts.find((i: any) => i.object === 'sys_file'); + expect(insert.data.owner_id).toBe('u1'); + expect(insert.data).not.toHaveProperty('organization_id'); + expect(insert.options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); + } finally { + await fs.rm(rootDir, { recursive: true, force: true }); + } + }); + it('a session with no active organization stamps nothing — ⛔ no guess is substituted', async () => { const rootDir = join(tmpdir(), `os-12745-noorg-${Math.random().toString(36).slice(2)}`); await fs.mkdir(rootDir, { recursive: true }); @@ -254,7 +288,8 @@ describe('storage upload routes: a file created with a session lands with that o expect(res._status).toBe(200); // Such a row is precisely what the backfill reports rather than repairs. - expect(engine._inserts[0].options).toBeUndefined(); + // [#21908] The opt-in rides alone: no organization is substituted. + expect(engine._inserts[0].options).toEqual({ context: { isSystem: true } }); } finally { await fs.rm(rootDir, { recursive: true, force: true }); } @@ -325,7 +360,7 @@ describe('StorageServicePlugin: the session bridge reports the ACTIVE organizati }); expect(res._status).toBe(200); - expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_A' } }); + expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); }); it('accepts the flattened shape a host may hand back directly', async () => { @@ -335,7 +370,7 @@ describe('StorageServicePlugin: the session bridge reports the ACTIVE organizati }); expect(res._status).toBe(200); - expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_B' } }); + expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_B', isSystem: true } }); }); it('⛔ invents nothing when the session carries no active organization', async () => { @@ -351,6 +386,6 @@ describe('StorageServicePlugin: the session bridge reports the ACTIVE organizati expect(res._status).toBe(200); expect(engine._inserts[0].data.owner_id).toBe('u1'); - expect(engine._inserts[0].options).toBeUndefined(); + expect(engine._inserts[0].options).toEqual({ context: { isSystem: true } }); }); }); diff --git a/packages/services/service-storage/src/sys-upload-session-organization-stamping.test.ts b/packages/services/service-storage/src/sys-upload-session-organization-stamping.test.ts index 8a858e5e4d4..e70d64e1344 100644 --- a/packages/services/service-storage/src/sys-upload-session-organization-stamping.test.ts +++ b/packages/services/service-storage/src/sys-upload-session-organization-stamping.test.ts @@ -111,7 +111,8 @@ describe('StorageMetadataStore.createSession: the acting organization reaches th // The platform's insert-side chokepoint reads exactly this: // `context.tenantId` → `buildDriverOptions` → `DriverOptions.tenantId` → // `SqlDriver.injectTenantOnInsert` → the object's tenant column. - expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_A' } }); + // [#21908] …beside the explicit system opt-in (ruling Q1). + expect(engine._inserts[0].options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); }); it('⛔ does NOT stamp the column onto the payload — that decision is the driver’s', async () => { @@ -127,7 +128,7 @@ describe('StorageMetadataStore.createSession: the acting organization reaches th expect(engine._inserts[0].data).not.toHaveProperty('organization_id'); }); - it('passes NO options at all when there is no organization — the pre-#12928 shape', async () => { + it('carries the opt-in and NO tenant when there is no organization — ⛔ none is invented', async () => { const engine = createRecordingEngine(); const store = new StorageMetadataStore(engine); @@ -137,13 +138,14 @@ describe('StorageMetadataStore.createSession: the acting organization reaches th await store.createSession(sessionRec('s4'), { organizationId: null }); // "NULL only where the caller genuinely has none" — the ruled other half. - // An empty context is still a context; handing one to the engine would - // change what every other option resolver on that call sees. + // [#21908] Every insert now carries the explicit system opt-in (ruling Q1), + // and still no `tenantId` where the caller has no organization: the opt-in + // is the only key, never a guessed organization. expect(engine._inserts.map((i: { options: unknown }) => i.options)).toEqual([ - undefined, - undefined, - undefined, - undefined, + { context: { isSystem: true } }, + { context: { isSystem: true } }, + { context: { isSystem: true } }, + { context: { isSystem: true } }, ]); }); @@ -268,7 +270,7 @@ describe('the chunked-upload door stamps its sys_upload_session row from the ses expect(res._status).toBe(200); const insert = engine._inserts.find((i: any) => i.object === 'sys_upload_session'); expect(insert).toBeTruthy(); - expect(insert.options).toEqual({ context: { tenantId: 'org_A' } }); + expect(insert.options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); }); it('stamps the file and the session row from the SAME session value', async () => { @@ -279,7 +281,7 @@ describe('the chunked-upload door stamps its sys_upload_session row from the ses // pair together is what keeps them from drifting apart again. const file = engine._inserts.find((i: any) => i.object === 'sys_file'); const session = engine._inserts.find((i: any) => i.object === 'sys_upload_session'); - expect(file.options).toEqual({ context: { tenantId: 'org_A' } }); + expect(file.options).toEqual({ context: { tenantId: 'org_A', isSystem: true } }); expect(session.options).toEqual(file.options); }); @@ -291,7 +293,8 @@ describe('the chunked-upload door stamps its sys_upload_session row from the ses expect(insert).toBeTruthy(); // Ruled: NULL only where the caller genuinely has none. Such a row is what // the TTL sweep reaps rather than what a backfill repairs. - expect(insert.options).toBeUndefined(); + // [#21908] The opt-in rides alone: no organization is substituted. + expect(insert.options).toEqual({ context: { isSystem: true } }); expect(insert.data).not.toHaveProperty('organization_id'); }); }); diff --git a/packages/services/service-storage/src/tenant-audit-update-delete-half-repairs.test.ts b/packages/services/service-storage/src/tenant-audit-update-delete-half-repairs.test.ts index 1ccf0ed344a..dffd590803e 100644 --- a/packages/services/service-storage/src/tenant-audit-update-delete-half-repairs.test.ts +++ b/packages/services/service-storage/src/tenant-audit-update-delete-half-repairs.test.ts @@ -583,6 +583,20 @@ describe('[#13178] the engine leg: the store’s context becomes DriverOptions.t expect(del!.options?.tenantId).toBe('org_A'); }); + it('[#21908] create carries it beside the explicit system opt-in — the organization stamp survives', async () => { + // Under the opt-in the organizations plugin's fill-only stamp stands down; + // what stamps the column is this channel, so it has to still arrive. + const { engine, calls } = await makeEngine(); + const store = new StorageMetadataStore(engine as any); + + await store.createFile({ ...fileRec('f9'), owner_id: 'u1' }, { organizationId: 'org_A' }); + await store.createSession(sessionRec('s9'), { organizationId: 'org_A' }); + + const creates = calls.filter((c) => c.method === 'create'); + expect(creates.map((c) => c.object)).toEqual(['sys_file', 'sys_upload_session']); + for (const c of creates) expect(c.options?.tenantId).toBe('org_A'); + }); + it('and WITHOUT the repair’s context there is none — the shape the audit exists to catch', async () => { const { engine, calls } = await makeEngine(); const store = new StorageMetadataStore(engine as any); From 9c844ac4d3b87311e87e86695906f0c82aaf04e5 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 16:31:02 +0000 Subject: [PATCH 3/5] wip: changeset; isSystem and tenant-audit census pages follow the census Claude-Session: https://claude.ai/code/session_01WMQprn46CND82KmY8sZWBu Co-authored-by: Claude --- .../21908-principal-less-producers-final.md | 18 +++++++++++ content/docs/permissions/system-context.mdx | 2 +- .../docs/permissions/tenant-audit-census.mdx | 30 +++++++++---------- ...08-tenant-audit-write-call-sites.counts.md | 30 +++++++++---------- 4 files changed, 49 insertions(+), 31 deletions(-) create mode 100644 .changeset/21908-principal-less-producers-final.md diff --git a/.changeset/21908-principal-less-producers-final.md b/.changeset/21908-principal-less-producers-final.md new file mode 100644 index 00000000000..ee1db8afa59 --- /dev/null +++ b/.changeset/21908-principal-less-producers-final.md @@ -0,0 +1,18 @@ +--- +"@objectstack/service-messaging": patch +"@objectstack/service-storage": patch +"@objectstack/service-settings": patch +"@objectstack/metadata-protocol": patch +--- + +The remaining platform producers in these four packages now pass the explicit system opt-in (`{ isSystem: true }`) on their data-engine calls. Until now they reached the engine with no principal and no opt-in, and the security middleware let that through only because of its principal-less hand-off. + +Clause-②: no + +- **service-messaging, the inbox read state.** `listInbox` (and its unread total), the receipt read behind it, and mark-read / mark-all-read take the opt-in inside the service. Their scope is unchanged: every read of a user's rows is keyed on the user id the door derived from the session, the receipt a mark-read inserts is stamped with it, and the receipt it updates is one a user-keyed read returned. +- **service-messaging, `owner_of:` audiences.** The record read takes the opt-in, the same posture as the email lookup beside it. It reads only `id` and the owner fields, and only the owner id leaves the resolver. An `owner_of:` audience on an object whose sharing model is `private` now resolves its owner; before, it resolved to nobody. +- **service-messaging, the rest of the fan-out and the outboxes.** The `role:` and `team:` membership reads, the email and SMS recipient reads, the notification template read, the dedup lookup in `emit()`, and both outboxes' enqueue, ack and list. +- **service-storage.** `StorageMetadataStore.createFile` and `createSession` insert under the opt-in. The organization still reaches the driver beside it, so the stored organization is unchanged, and the file's `owner_id` is still the uploading user. +- **service-settings.** The `sys_secret` store the plugin builds (insert, get, update), and the read that verifies a rotation before the old secret is reaped. A store `update` now writes the `ciphertext` it is given; without a context the engine's read-only strip dropped it. No caller in this repository uses `update`. +- **metadata-protocol.** `SysMetadataRepository.getByHash`, `list`, `history` and the history replay of `watch()`. +- None of the gates the middleware runs before its hand-off applies to these calls. ⛔ No new export on any package entry, and no new elevation API. diff --git a/content/docs/permissions/system-context.mdx b/content/docs/permissions/system-context.mdx index db91cb97888..92b6179ba09 100644 --- a/content/docs/permissions/system-context.mdx +++ b/content/docs/permissions/system-context.mdx @@ -353,7 +353,7 @@ still holds equal to the census on every pull request: | — in tests | 1013 | — | | — in non-test sources | 798 | — | | Appearances of the bare identifier `isSystem` in non-test sources | 813 | — | -| — parsed as a declaration | 26 | ✅ | +| — parsed as a declaration | 27 | ✅ | | — parsed as an object-literal / type key (producers and option objects) | 310 | — | | — parsed as a property **read** | 120 | ✅ | | — parsed in some other syntactic position (a local, a cast, a conditional) | 9 | ✅ | diff --git a/content/docs/permissions/tenant-audit-census.mdx b/content/docs/permissions/tenant-audit-census.mdx index 2fe206b417a..7afde0e2bf5 100644 --- a/content/docs/permissions/tenant-audit-census.mdx +++ b/content/docs/permissions/tenant-audit-census.mdx @@ -122,7 +122,7 @@ are reported as `undecidable` rather than assumed either way. The same holds twice over for the context. An options argument spelled as a literal can be read; one spelled `options`, `{ ...opts }`, or handed through a -forwarding shim cannot, and **60 of the 233 sites are spelled that way**. A +forwarding shim cannot, and **54 of the 233 sites are spelled that way**. A context resolved from an inline literal or a local `const` can be tested for `isSystem`; one arriving from a helper call cannot. @@ -150,8 +150,8 @@ now **0**: nothing on this surface threads a context that provably lacks the fla **"No tenant context" counted sites it had not read.** An options argument the walker could not parse was folded into the same bucket as one it had read and -found empty. That published **69 sites "carrying no tenant context at all"** -when 9 said so and 60 were simply unread — an over-claim in the *alarming* +found empty. That published **62 sites "carrying no tenant context at all"** +when 8 said so and 54 were simply unread — an over-claim in the *alarming* direction, on the very figure this page tells other cards to cite. `carries` is now three-valued, and an unreadable argument can never contribute to the provable count. @@ -188,9 +188,9 @@ reproduce them. Where it disagrees, it disagrees on the page: | carried figure | where it survives | this census | | :--- | :--- | ---: | | 175 write call sites | quoted in the merged changeset | **233** | -| 24 carrying no tenant context | quoted in the merged changeset | **2** provable and tenancy-enabled; **33** more whose options argument is unreadable | +| 24 carrying no tenant context | quoted in the merged changeset | **2** provable and tenancy-enabled; **31** more whose options argument is unreadable | | 127 of 175 statically decidable, 48 runtime-parameter-name sites | restated on the `isSystem`-scoping card | **155 of 233** decidable, **78** undecidable | -| 135 (77%) silenced by the `isSystem` guard before the posture gate | the lost issue body — **no surviving corroboration** | **not reproduced**: 121 decidably elevated, 0 decidably not, 103 undecidable | +| 135 (77%) silenced by the `isSystem` guard before the posture gate | the lost issue body — **no surviving corroboration** | **not reproduced**: 123 decidably elevated, 0 decidably not, 102 undecidable | | 141 and 132, two independent re-derivations | the card that filed this work | — | **The differences are not reconciled, and deliberately so.** The old census's @@ -207,14 +207,14 @@ would report a smaller number and would not say so. The fourth row is the one worth flagging to anyone citing it. **The 135 / 77% figure has no surviving corroboration anywhere in the tree.** This census reads -121 of 233 (52%) as decidably elevated, with 103 more whose elevation is a +123 of 233 (53%) as decidably elevated, with 102 more whose elevation is a run-time fact — so the claim is neither confirmed nor refuted, and the honest answer is that a static reading cannot settle it. ⇒ **Cite `2 / 233`, and say what it is**: the sites whose options argument was READ and holds no tenant context, against a decidably tenancy-enabled object. That is the control's provable yield surface. ⛔ Do not cite it as "the sites -without tenant context" — **33 further sites** have an options argument this +without tenant context" — **31 further sites** have an options argument this cannot read, and they are neither in nor out. {/* BEGIN GENERATED: tenant-audit-census (scripts/tenant-audit-census.mjs) — DO NOT EDIT */} @@ -228,14 +228,14 @@ cannot read, and they are neither in nor out. | …whose object name is chosen at run time | 78 | | …against an object with tenancy ENABLED | 154 | | …against an object that declares tenancy off | 1 | -| threading a tenant context | 164 | -| PROVABLY carrying none (options read, no context key) | **9** | +| threading a tenant context | 171 | +| PROVABLY carrying none (options read, no context key) | **8** | | …of those, against a decidably tenancy-enabled object | **2** | -| options argument UNREADABLE — may or may not carry one | 60 | -| …of those, against a decidably tenancy-enabled object | 33 | -| threading a decidably ELEVATED (`isSystem`) context | 121 | +| options argument UNREADABLE — may or may not carry one | 54 | +| …of those, against a decidably tenancy-enabled object | 31 | +| threading a decidably ELEVATED (`isSystem`) context | 123 | | threading a context that is decidably NOT elevated | 0 | -| threading a context whose elevation is a run-time fact | 103 | +| threading a context whose elevation is a run-time fact | 102 | | how the instrument reached the site | count | | :--- | ---: | @@ -297,11 +297,11 @@ holds still. They are required to be HERE and to say WHEN they were true; their values are not compared. The reasoning, and the measurement behind it, are in `scripts/check-tenant-audit-census.mjs`. -Measured on 2026-10-06 at `3832674ac`. +Measured on 2026-10-06 at `2bea8b684`. | corpus scale (not enforced) | count | | :--- | ---: | -| tracked non-test sources scanned | 609 | +| tracked non-test sources scanned | 612 | | engine-shaped types recognised | 70 | | declared objects in the registry | 116 | | same-named calls subtracted as non-engine | 159 | diff --git a/docs/audits/2026-08-tenant-audit-write-call-sites.counts.md b/docs/audits/2026-08-tenant-audit-write-call-sites.counts.md index f5eaac2edcf..598d6061259 100644 --- a/docs/audits/2026-08-tenant-audit-write-call-sites.counts.md +++ b/docs/audits/2026-08-tenant-audit-write-call-sites.counts.md @@ -38,14 +38,14 @@ silent, and `node scripts/tenant-audit-census.mjs --write` is the resolution. | Object name chosen at run time | 78 | | Against a tenancy-enabled object | 154 | | Against an object declaring tenancy off | 1 | -| Threading a tenant context | 164 | -| Provably carrying none | 9 | +| Threading a tenant context | 171 | +| Provably carrying none | 8 | | …and decidably tenancy-enabled | 2 | -| Options argument unreadable | 60 | -| …and decidably tenancy-enabled | 33 | -| Threading a decidably elevated context | 121 | +| Options argument unreadable | 54 | +| …and decidably tenancy-enabled | 31 | +| Threading a decidably elevated context | 123 | | Threading a decidably non-elevated context | 0 | -| Threading a context of undecidable elevation | 103 | +| Threading a context of undecidable elevation | 102 | ## Subtractions the census could NOT defend — enforced @@ -90,11 +90,11 @@ holds still. They are required to be HERE and to say WHEN they were true; their values are not compared. The reasoning, and the measurement behind it, are in `scripts/check-tenant-audit-census.mjs`. -Measured on 2026-10-06 at `3832674ac`. +Measured on 2026-10-06 at `2bea8b684`. | corpus scale (not enforced) | count | | :--- | ---: | -| tracked non-test sources scanned | 609 | +| tracked non-test sources scanned | 612 | | engine-shaped types recognised | 70 | | declared objects in the registry | 116 | | same-named calls subtracted as non-engine | 159 | @@ -219,13 +219,13 @@ Measured on 2026-10-06 at `3832674ac`. | `packages/services/service-job/src/db-job-adapter.ts` | `update` | `sys_job_run` | enabled | elevated | 1 | | `packages/services/service-messaging/src/inbox-channel.ts` | `insert` | `objectName` | undecidable | context, elevation undecidable | 1 | | `packages/services/service-messaging/src/inbox-channel.ts` | `insert` | `receiptObject` | undecidable | context, elevation undecidable | 1 | -| `packages/services/service-messaging/src/messaging-service.ts` | `insert` | `RECEIPT_OBJECT` | undecidable | PROVABLY NONE | 1 | +| `packages/services/service-messaging/src/messaging-service.ts` | `insert` | `RECEIPT_OBJECT` | undecidable | context, elevation undecidable | 1 | | `packages/services/service-messaging/src/messaging-service.ts` | `update` | `RECEIPT_OBJECT` | undecidable | options unreadable | 1 | | `packages/services/service-messaging/src/messaging-service.ts` | `insert` | `sys_notification` | enabled | context, elevation undecidable | 1 | -| `packages/services/service-messaging/src/sql-http-outbox.ts` | `insert` | `this.objectName` | undecidable | options unreadable | 1 | -| `packages/services/service-messaging/src/sql-http-outbox.ts` | `update` | `this.objectName` | undecidable | context, elevation undecidable | 2 | -| `packages/services/service-messaging/src/sql-http-outbox.ts` | `update` | `this.objectName` | undecidable | options unreadable | 3 | -| `packages/services/service-messaging/src/sql-outbox.ts` | `insert` | `this.objectName` | undecidable | options unreadable | 1 | +| `packages/services/service-messaging/src/sql-http-outbox.ts` | `insert` | `this.objectName` | undecidable | context, elevation undecidable | 1 | +| `packages/services/service-messaging/src/sql-http-outbox.ts` | `update` | `this.objectName` | undecidable | context, elevation undecidable | 4 | +| `packages/services/service-messaging/src/sql-http-outbox.ts` | `update` | `this.objectName` | undecidable | options unreadable | 1 | +| `packages/services/service-messaging/src/sql-outbox.ts` | `insert` | `this.objectName` | undecidable | context, elevation undecidable | 1 | | `packages/services/service-messaging/src/sql-outbox.ts` | `update` | `this.objectName` | undecidable | context, elevation undecidable | 3 | | `packages/services/service-messaging/src/sql-outbox.ts` | `update` | `this.objectName` | undecidable | options unreadable | 1 | | `packages/services/service-queue/src/db-queue-adapter.ts` | `delete` | `sys_job_queue` | enabled | context, elevation undecidable | 2 | @@ -235,8 +235,8 @@ Measured on 2026-10-06 at `3832674ac`. | `packages/services/service-settings/src/settings-service-plugin.ts` | `insert` | `objectName` | undecidable | options unreadable | 1 | | `packages/services/service-settings/src/settings-service-plugin.ts` | `update` | `objectName` | undecidable | options unreadable | 2 | | `packages/services/service-settings/src/settings-service-plugin.ts` | `delete` | `sys_secret` | enabled | elevated | 1 | -| `packages/services/service-settings/src/settings-service-plugin.ts` | `insert` | `sys_secret` | enabled | options unreadable | 1 | -| `packages/services/service-settings/src/settings-service-plugin.ts` | `update` | `sys_secret` | enabled | options unreadable | 1 | +| `packages/services/service-settings/src/settings-service-plugin.ts` | `insert` | `sys_secret` | enabled | elevated | 1 | +| `packages/services/service-settings/src/settings-service-plugin.ts` | `update` | `sys_secret` | enabled | elevated | 1 | | `packages/services/service-settings/src/settings-service-plugin.ts` | `insert` | `sys_setting_audit` | enabled | elevated | 1 | | `packages/services/service-settings/src/settings-service.ts` | `insert` | `this.objectName` | undecidable | options unreadable | 1 | | `packages/services/service-settings/src/settings-service.ts` | `update` | `this.objectName` | undecidable | options unreadable | 1 | From 8064c39ee5bdcde659472d1d3f2c3e8198edad4b Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 16:33:57 +0000 Subject: [PATCH 4/5] wip: the engine-double ledger records the settings pin's update and delete doubles Claude-Session: https://claude.ai/code/session_01WMQprn46CND82KmY8sZWBu Co-authored-by: Claude --- scripts/engine-double-contract.pinned.json | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/scripts/engine-double-contract.pinned.json b/scripts/engine-double-contract.pinned.json index fcda1534fa0..15ed413e19f 100644 --- a/scripts/engine-double-contract.pinned.json +++ b/scripts/engine-double-contract.pinned.json @@ -4281,6 +4281,16 @@ "verb": "update", "pinned": 2 }, + { + "file": "packages/services/service-settings/src/settings-system-context.pin.test.ts", + "verb": "delete", + "pinned": 1 + }, + { + "file": "packages/services/service-settings/src/settings-system-context.pin.test.ts", + "verb": "update", + "pinned": 1 + }, { "file": "packages/services/service-settings/src/sys-secret-orphan-report.test.ts", "verb": "update", From e2aa5ecaa0548e791ba50de3b9aec6838eaf241c Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 17:01:49 +0000 Subject: [PATCH 5/5] test(service-messaging): the stage-1 pin's matcher refuses a combinator inside its callback Claude-Session: https://claude.ai/code/session_01WMQprn46CND82KmY8sZWBu Co-authored-by: Claude --- .../service-messaging/src/system-context.pin.test.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/services/service-messaging/src/system-context.pin.test.ts b/packages/services/service-messaging/src/system-context.pin.test.ts index aad0fae1f4c..8ce5d4654fd 100644 --- a/packages/services/service-messaging/src/system-context.pin.test.ts +++ b/packages/services/service-messaging/src/system-context.pin.test.ts @@ -180,10 +180,10 @@ describe('[#21913] the emit fan-out carries the explicit system opt-in', () => { /** Plain-equality `where` over seeded rows — anything else is refused, loudly. */ function matchRows(rows: Array>, where: Record = {}) { - for (const k of Object.keys(where)) { + return rows.filter((r) => Object.entries(where).every(([k, v]) => { if (k.startsWith('$')) throw new Error(`recording engine: unimplemented combinator ${k}`); - } - return rows.filter((r) => Object.entries(where).every(([k, v]) => (r[k] ?? null) === (v ?? null))); + return (r[k] ?? null) === (v ?? null); + })); } const USER_SCOPED = new Set(['sys_inbox_message', 'sys_notification_receipt']);