✅ test: Prove Tenant Scoping at the Driver, Below the Engine Binding (#15578)

Tenant isolation is enforced by the storage engine's own binding — Mongoose
middleware today. No amount of reading the binding answers whether it covers
every path, including the ones nobody enumerated, so this checks from below:
`attachTenantProbe` observes the commands that actually reach the wire via the
driver's `commandStarted` monitoring. The check does not depend on the binding
being correct, which is the point of putting it there.

It asks one narrow, decidable question — did the binding's own contribution
arrive, for the tenant active when the command was issued? — rather than trying
to decide tenant-safety for arbitrary filter algebra. Anything it cannot judge
is reported unscoped, so an unfamiliar command shape surfaces as a failure
instead of a silent pass.

`probe.spec.ts` runs the whole enforcement surface through it: 13 query
operations, 3 write paths, and the document-level paths (`doc.save()`,
`doc.updateOne()`, `doc.deleteOne()`, `populate()`) that make schema
middleware irreplaceable by a wrapper around the model's statics.

Test-only; nothing here is bundled.

One gap found and pinned for the next change: `doc.save()` on a persisted
document is filtered on `_id` alone. `estimatedDocumentCount()` is pinned as a
known unscoped side channel (one call site, boot path, metadata only).
This commit is contained in:
Danny Avila 2026-09-06 11:44:39 -04:00 • committed by GitHub
parent fdd5f9ce65
commit 85c90bf4c7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 906 additions and 0 deletions

View file

@ -0,0 +1,555 @@
import mongoose from 'mongoose';
import { MongoMemoryServer } from 'mongodb-memory-server';
import type { TenantProbe, TenantCommandRecord } from './probe';
import { tenantStorage, SYSTEM_TENANT_ID } from '~/config/tenantContext';
import { applyTenantIsolation } from '~/models/plugins/tenantIsolation';
import { tenantSafeBulkWrite } from '~/utils/tenantBulkWrite';
import { attachTenantProbe, unscoped } from './probe';
/**
* Completeness proof for the Mongoose tenant-isolation binding.
*
* Every other tenant test asserts what a query *returns*. These assert what
* reaches the database, which is the only way to catch a path the binding never
* hooked. In particular this pins the reason the middleware cannot be replaced
* by a wrapper around the model's static methods: document-level writes and
* populate never pass through such a wrapper, but they do reach the wire.
*/
interface WidgetDocument {
name: string;
tenantId?: string;
parts: mongoose.Types.ObjectId[];
}
interface PartDocument {
label: string;
tenantId?: string;
}
const TENANT = 'tenant-a';
let server: MongoMemoryServer;
let probe: TenantProbe;
let Widget: mongoose.Model<WidgetDocument>;
let Part: mongoose.Model<PartDocument>;
const asTenant = <T>(fn: () => Promise<T>): Promise<T> =>
tenantStorage.run({ tenantId: TENANT }, fn);
const describeLeaks = (records: readonly TenantCommandRecord[]): string =>
unscoped(records)
.map((record) => `${record.commandName} on ${record.collection}: ${record.predicate}`)
.join('\n');
beforeAll(async () => {
server = await MongoMemoryServer.create();
await mongoose.connect(server.getUri(), { monitorCommands: true });
const partSchema = new mongoose.Schema<PartDocument>({
label: String,
tenantId: { type: String, index: true },
});
applyTenantIsolation(partSchema);
Part = mongoose.model<PartDocument>('ProbePart', partSchema);
const widgetSchema = new mongoose.Schema<WidgetDocument>({
name: String,
tenantId: { type: String, index: true },
parts: [{ type: mongoose.Schema.Types.ObjectId, ref: 'ProbePart' }],
});
applyTenantIsolation(widgetSchema);
Widget = mongoose.model<WidgetDocument>('ProbeWidget', widgetSchema);
probe = attachTenantProbe(mongoose.connection, [
Widget.collection.collectionName,
Part.collection.collectionName,
]);
});
afterAll(async () => {
probe.close();
await mongoose.disconnect();
await server.stop();
});
beforeEach(async () => {
await tenantStorage.run({ tenantId: SYSTEM_TENANT_ID }, async () => {
await Widget.deleteMany({});
await Part.deleteMany({});
});
});
describe('the probe itself', () => {
it('detects a command that reached the wire unscoped', async () => {
const records = await probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.find({})
.toArray();
}),
);
expect(unscoped(records)).toHaveLength(1);
expect(records[0].commandName).toBe('find');
});
const rawFind = (filter: Record<string, unknown>) =>
probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.find(filter)
.toArray();
}),
);
/**
* The probe accepts only the binding's own contribution — `tenantId` equal to
* the active tenant, at the top level of the filter. Everything else is
* reported unscoped, because deciding tenant-safety for arbitrary filter
* algebra is unbounded and every operator left open is a way for a regression
* to look scoped.
*/
it.each([
['a different tenant entirely', { tenantId: 'tenant-b' }],
['a permissive $or branch', { $or: [{ tenantId: TENANT }, {}] }],
['tenantId nested under an unrelated field', { meta: { tenantId: TENANT } }],
['$ne, which selects every other tenant', { tenantId: { $ne: TENANT } }],
['a multi-valued $in', { tenantId: { $in: [TENANT, 'tenant-b'] } }],
['a regex inside a singleton $in', { tenantId: { $in: [/^tenant-/] } }],
['a regex over tenants', { tenantId: { $regex: '.*' } }],
['$exists: true, which matches any tenant', { tenantId: { $exists: true } }],
['operator forms the binding never emits', { tenantId: { $eq: TENANT } }],
])('does not accept %s as scoping', async (_label, filter) => {
expect(unscoped(await rawFind(filter))).toHaveLength(1);
});
it('accepts the shape the binding actually emits', async () => {
const records = await rawFind({ tenantId: TENANT, name: 'x' });
expect(records).toHaveLength(1);
expect(unscoped(records)).toHaveLength(0);
expect(records[0].tenantId).toBe(TENANT);
});
const rawAggregate = (pipeline: Record<string, unknown>[]) =>
probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.aggregate(pipeline)
.toArray();
}),
);
/**
* Position is the whole point. A `$match` that lands after data has been
* combined, limited or skipped mentions the tenant without having restricted
* what the earlier stages saw.
*/
it.each([
['after a $group', [{ $group: { _id: null } }, { $match: { tenantId: TENANT } }]],
['after a $limit', [{ $limit: 1 }, { $match: { tenantId: TENANT } }]],
['after a $skip', [{ $skip: 1 }, { $match: { tenantId: TENANT } }]],
['after a $sort', [{ $sort: { name: 1 } }, { $match: { tenantId: TENANT } }]],
])('does not accept a tenant $match %s', async (_label, pipeline) => {
expect(unscoped(await rawAggregate(pipeline))).toHaveLength(1);
});
it('accepts a pipeline whose first stage is the tenant $match', async () => {
const records = await rawAggregate([
{ $match: { tenantId: TENANT } },
{ $group: { _id: null, total: { $sum: 1 } } },
]);
expect(unscoped(records)).toHaveLength(0);
});
it('does not accept an outer $match as scoping a joined collection', async () => {
const records = await rawAggregate([
{ $match: { tenantId: TENANT } },
{
$lookup: {
from: Part.collection.collectionName,
localField: 'parts',
foreignField: '_id',
as: 'joined',
},
},
]);
expect(unscoped(records)).toHaveLength(1);
});
it('accepts a $lookup whose sub-pipeline leads with the tenant $match', async () => {
const records = await rawAggregate([
{ $match: { tenantId: TENANT } },
{
$lookup: {
from: Part.collection.collectionName,
pipeline: [{ $match: { tenantId: TENANT } }],
as: 'joined',
},
},
]);
expect(unscoped(records)).toHaveLength(0);
});
/**
* `$out` and `$merge` rewrite another collection from inside the aggregate,
* so the destination write never surfaces as its own command to inspect.
*/
it('rejects a pipeline that writes to another collection', async () => {
const records = await rawAggregate([
{ $match: { tenantId: TENANT } },
{ $out: 'probe_out_target' },
]);
expect(unscoped(records)).toHaveLength(1);
});
/**
* A cursor continuation carries no predicate, and the scope its cursor was
* opened under is not visible here — so it must be reported rather than
* silently dropped, or a cursor primed under one scope could be drained under
* another with the probe seeing nothing at all.
*/
it('reports a cursor continuation rather than discarding it', async () => {
await tenantStorage.run({ tenantId: SYSTEM_TENANT_ID }, async () => {
await Widget.insertMany(
Array.from({ length: 6 }, (_, index) => ({ name: `batched-${index}`, tenantId: TENANT })),
);
});
const records = await probe.record(() =>
asTenant(async () => {
const cursor = mongoose.connection
.db!.collection(Widget.collection.collectionName)
.find({ tenantId: TENANT })
.batchSize(2);
await cursor.toArray();
}),
);
expect(records.some((record) => record.commandName === 'getMore')).toBe(true);
});
/**
* A collation can broaden equality — under a case-insensitive one,
* `tenantId: 'tenant-a'` also matches `TENANT-A` — so the literal predicate
* no longer proves isolation.
*/
it('does not accept a predicate broadened by a collation', async () => {
const records = await probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.find({ tenantId: TENANT }, { collation: { locale: 'en', strength: 2 } })
.toArray();
}),
);
expect(unscoped(records)).toHaveLength(1);
});
/**
* `.explain()` nests the real command inside an `explain` envelope. Without
* unwrapping it the operation is invisible, though `executionStats` still
* reveals cross-tenant cardinality.
*/
it('unwraps an explain envelope instead of dropping it', async () => {
const records = await probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.find({ name: 'x' })
.explain('queryPlanner');
}),
);
expect(records).toHaveLength(1);
expect(records[0].commandName).toBe('find');
expect(records[0].scoped).toBe(false);
});
/**
* Diagnostics run inside the driver's synchronous event handler, so a
* serialization throw would escape into the database operation itself.
*/
it('does not throw out of the listener on a value JSON cannot serialize', async () => {
const records = await probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.find({ tenantId: TENANT, big: { $gt: BigInt(1) } } as unknown as Record<string, unknown>)
.toArray()
.catch(() => undefined);
}),
);
expect(records).toHaveLength(1);
expect(records[0].predicate).toContain('tenantId');
});
/**
* Dropping a collection destroys every tenant's rows. It carries no predicate
* to inspect, so the only safe report is an unscoped one — silence would be a
* green result for the most destructive escape hatch there is.
*/
it('reports a collection drop as unscoped', async () => {
const records = await probe.record(() =>
asTenant(async () => {
await mongoose.connection
.db!.collection(Widget.collection.collectionName)
.drop()
.catch(() => undefined);
}),
);
expect(unscoped(records)).toHaveLength(1);
expect(records[0].commandName).toBe('drop');
});
/** `update` and `delete` carry collation per batched operation, not on the envelope. */
it.each([
[
'updateOne',
async () =>
mongoose.connection.db!.collection(Widget.collection.collectionName).updateOne(
{ tenantId: TENANT },
{ $set: { name: 'x' } },
{
collation: { locale: 'en', strength: 2 },
},
),
],
[
'deleteOne',
async () =>
mongoose.connection
.db!.collection(Widget.collection.collectionName)
.deleteOne({ tenantId: TENANT }, { collation: { locale: 'en', strength: 2 } }),
],
])('does not accept a per-operation collation on %s', async (_label, operation) => {
const records = await probe.record(() => asTenant(async () => void (await operation())));
expect(unscoped(records)).toHaveLength(1);
});
it('reports nothing for collections it does not watch', async () => {
const records = await probe.record(() =>
asTenant(async () => {
await mongoose.connection.db!.collection('unwatched').find({}).toArray();
}),
);
expect(records).toHaveLength(0);
});
});
describe('query paths reach the wire scoped', () => {
it.each([
['find', () => Widget.find({ name: 'x' })],
['findOne', () => Widget.findOne({ name: 'x' })],
['countDocuments', () => Widget.countDocuments({ name: 'x' })],
['distinct', () => Widget.distinct('name')],
['updateOne', () => Widget.updateOne({ name: 'x' }, { $set: { name: 'y' } })],
['updateMany', () => Widget.updateMany({ name: 'x' }, { $set: { name: 'y' } })],
['deleteOne', () => Widget.deleteOne({ name: 'x' })],
['deleteMany', () => Widget.deleteMany({ name: 'x' })],
['findOneAndUpdate', () => Widget.findOneAndUpdate({ name: 'x' }, { $set: { name: 'y' } })],
['findOneAndDelete', () => Widget.findOneAndDelete({ name: 'x' })],
['replaceOne', () => Widget.replaceOne({ name: 'x' }, { name: 'y' })],
['findOneAndReplace', () => Widget.findOneAndReplace({ name: 'x' }, { name: 'y' })],
['aggregate', () => Widget.aggregate([{ $group: { _id: '$name' } }])],
])('%s', async (_label, operation) => {
const records = await probe.record(() => asTenant(async () => void (await operation())));
expect(records.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
});
/**
* `estimatedDocumentCount()` reads collection metadata and takes no filter, so
* no binding can scope it — it is a genuine global side channel, currently used
* once at `models/plugins/mongoMeili.ts` for a sync progress log on a boot path.
* Pinned so it stays visible; making scoped callers fail closed is a runtime
* change and belongs in its own PR.
*/
describe('estimatedDocumentCount is a known unscoped side channel', () => {
it('reaches the wire with no tenant predicate', async () => {
const records = await probe.record(() =>
asTenant(async () => void (await Widget.estimatedDocumentCount())),
);
expect(records).toHaveLength(1);
expect(records[0].scoped).toBe(false);
});
});
describe('write paths reach the wire scoped', () => {
it('Model.create', async () => {
const records = await probe.record(() =>
asTenant(async () => void (await Widget.create({ name: 'created', parts: [] }))),
);
expect(records.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
it('insertMany', async () => {
const records = await probe.record(() =>
asTenant(async () => void (await Widget.insertMany([{ name: 'a' }, { name: 'b' }]))),
);
expect(records.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
it('tenantSafeBulkWrite', async () => {
const records = await probe.record(() =>
asTenant(async () =>
tenantSafeBulkWrite(Widget, [
{ insertOne: { document: { name: 'bulk' } } },
{ updateOne: { filter: { name: 'bulk' }, update: { $set: { name: 'bulk2' } } } },
{ deleteOne: { filter: { name: 'gone' } } },
]),
),
);
expect(records.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
});
/**
* These are the paths that make schema middleware irreplaceable: none of them
* pass through a wrapper around the model's static methods, but all of them
* reach the database.
*/
describe('document-level paths reach the wire scoped', () => {
it('doc.save() on a new document', async () => {
const records = await probe.record(() =>
asTenant(async () => {
const widget = new Widget({ name: 'saved' });
await widget.save();
}),
);
expect(records.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
/**
* KNOWN GAP, found by this probe. Saving an already-persisted document issues
* `update` filtered on `_id` alone — the tenant is stamped onto the payload but
* never asserted in the predicate. It is scoped by provenance (you can only
* fetch a document your tenant can see), so the normal flow is safe, but an
* `_id` obtained from an unscoped source — `runAsSystem`, a cached id, a
* client-supplied id — writes across tenants unguarded.
*
* Pinned here rather than fixed: closing it changes runtime behaviour and
* belongs in its own change. The assertion is written to FAIL once the
* predicate is added, so the fix cannot land without updating this test.
*/
it('doc.save() on a fetched document is scoped by provenance, not by predicate', async () => {
await asTenant(async () => void (await Widget.create({ name: 'fetched', parts: [] })));
const records = await probe.record(() =>
asTenant(async () => {
const widget = await Widget.findOne({ name: 'fetched' });
widget!.name = 'renamed';
await widget!.save();
}),
);
const updates = records.filter((record) => record.commandName === 'update');
expect(updates).toHaveLength(1);
expect(updates[0].scoped).toBe(false);
expect(updates[0].predicate).toContain('_id');
});
it('doc.updateOne()', async () => {
await asTenant(async () => void (await Widget.create({ name: 'target', parts: [] })));
const records = await probe.record(() =>
asTenant(async () => {
const widget = await Widget.findOne({ name: 'target' });
await widget!.updateOne({ $set: { name: 'updated' } });
}),
);
expect(records.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
/**
* `Document.prototype.deleteOne()` is called out in `schema/auditLog.ts` as an
* escape hatch around *document* middleware. It is not one for tenant
* filtering — it builds a Query, so the query-level hook still fires — but
* nothing pinned that, so a future change could silently make it one.
*/
it('doc.deleteOne() reaches the wire scoped', async () => {
await asTenant(async () => void (await Widget.create({ name: 'doomed', parts: [] })));
const records = await probe.record(() =>
asTenant(async () => {
const widget = await Widget.findOne({ name: 'doomed' });
await widget!.deleteOne();
}),
);
const deletes = records.filter((record) => record.commandName === 'delete');
expect(deletes).toHaveLength(1);
expect(describeLeaks(records)).toBe('');
});
it('doc.deleteOne() cannot remove another tenant document', async () => {
const foreign = await tenantStorage.run({ tenantId: 'tenant-b' }, async () =>
Widget.create({ name: 'foreign-doomed', parts: [] }),
);
const smuggled = await tenantStorage.run({ tenantId: SYSTEM_TENANT_ID }, async () =>
Widget.findById(foreign._id),
);
await asTenant(async () => void (await smuggled!.deleteOne()));
const survivor = await tenantStorage.run({ tenantId: SYSTEM_TENANT_ID }, async () =>
Widget.findById(foreign._id).lean(),
);
expect(survivor).not.toBeNull();
});
it('populate() issues its own query', async () => {
await asTenant(async () => {
const part = await Part.create({ label: 'part-1' });
await Widget.create({ name: 'parent', parts: [part._id] });
});
const records = await probe.record(() =>
asTenant(async () => void (await Widget.findOne({ name: 'parent' }).populate('parts'))),
);
const partReads = records.filter(
(record) => record.collection === Part.collection.collectionName,
);
expect(partReads.length).toBeGreaterThan(0);
expect(describeLeaks(records)).toBe('');
});
});
describe('system scope', () => {
it('deliberately reaches the wire unscoped', async () => {
const records = await probe.record(() =>
tenantStorage.run({ tenantId: SYSTEM_TENANT_ID }, async () => {
await Widget.find({ name: 'x' });
}),
);
expect(unscoped(records)).toHaveLength(1);
});
});

View file

@ -0,0 +1,351 @@
import type { Connection } from 'mongoose';
import { getTenantId, SYSTEM_TENANT_ID } from '~/config/tenantContext';
/**
* Driver-level completeness probe for tenant isolation.
*
* Tenant scoping is enforced by the storage engine's own binding — Mongoose
* middleware today. That leaves one question no amount of reading the binding
* can answer: does it cover *every* path, including the ones nobody enumerated?
*
* This probe answers it from below, by observing the commands that actually
* reach the wire. The check does not depend on the binding being correct, which
* is the whole point of putting it here.
*
* It asks one narrow, decidable question: **did the binding's own contribution
* arrive, for the tenant that was active when the command was issued?** It does
* not try to decide tenant-safety for arbitrary filter algebra — that problem is
* unbounded, and every operator left open is a way for a regression to look
* scoped. The binding emits `{ tenantId: <active> }` at the top level of a
* filter, or unshifts `{ $match: { tenantId: <active> } }` onto a pipeline, so
* that is exactly what is checked. Anything else is reported unscoped.
*
* Requires the connection to have been opened with `monitorCommands: true`.
*/
/** One command observed reaching the database. */
export interface TenantCommandRecord {
readonly commandName: string;
readonly collection: string;
/** The tenant active when the command was issued, if any. */
readonly tenantId?: string;
/** True when the command provably restricts itself to that tenant. */
readonly scoped: boolean;
/** The predicates as sent, for test failure output. */
readonly predicate: string;
}
export interface TenantProbe {
/**
* Runs `fn` and returns every observed command against a watched collection.
* Not re-entrant: one recording at a time per probe.
*/
record(fn: () => Promise<unknown>): Promise<readonly TenantCommandRecord[]>;
/** Detaches the listener. */
close(): void;
}
type CommandDocument = Record<string, unknown>;
/**
* Predicate locations per command, as the wire protocol names them.
* `list` marks commands whose payload is an array of sub-operations.
*/
const PREDICATE_PATHS: ReadonlyMap<string, { readonly key: string; readonly list?: string }> =
new Map([
['find', { key: 'filter' }],
['count', { key: 'query' }],
['distinct', { key: 'query' }],
['findAndModify', { key: 'query' }],
['update', { key: 'q', list: 'updates' }],
['delete', { key: 'q', list: 'deletes' }],
['insert', { key: '', list: 'documents' }],
]);
/**
* Commands that touch a collection but carry no tenant-bearing data. Anything
* outside this list and `PREDICATE_PATHS` is reported unscoped rather than
* dropped, so a command shape the probe does not understand fails closed
* instead of passing unseen.
*
* `drop` is deliberately absent: unlike index inspection, dropping a collection
* destroys every tenant's rows, and a leak detector must never be silent about
* that.
*/
const UNJUDGED_COMMANDS: ReadonlySet<string> = new Set([
'createIndexes',
'listIndexes',
'dropIndexes',
'listCollections',
'collMod',
'create',
]);
/**
* `JSON.stringify` throws on BSON values it does not know, and this runs inside
* the driver's synchronous event handler — a throw would escape into the
* database operation itself. Diagnostics must never do that.
*/
function safeStringify(value: unknown): string {
try {
return JSON.stringify(value, (_key, item) => (typeof item === 'bigint' ? `${item}n` : item));
} catch {
return '(unserializable predicate)';
}
}
/** Stages that rewrite another collection from inside the aggregate command. */
const COLLECTION_WRITING_STAGES = ['$out', '$merge'] as const;
/**
* Whether a filter carries the binding's own contribution: `tenantId` equal to
* the active tenant, at the top level, as a plain equality.
*/
function pinsTenant(filter: unknown, tenantId: string): boolean {
if (filter == null || typeof filter !== 'object' || Array.isArray(filter)) {
return false;
}
return (filter as CommandDocument).tenantId === tenantId;
}
/** A `$lookup` reads its foreign collection inside the outer command. */
function lookupTargetsTenant(lookup: unknown, tenantId: string): boolean {
if (lookup == null || typeof lookup !== 'object') {
return false;
}
return pipelineTargetsTenant((lookup as CommandDocument).pipeline, tenantId);
}
/**
* A pipeline is scoped when its *first* stage is the tenant `$match` the
* binding unshifts, it writes to no other collection, and every joined read
* scopes itself.
*
* Position is the whole point: `[{ $limit: 1 }, { $match: { tenantId } }]`
* mentions the tenant but lets another tenant's row take the limited slot, and
* a `$group` before the match has already combined every tenant's data.
*/
function pipelineTargetsTenant(pipeline: unknown, tenantId: string): boolean {
if (!Array.isArray(pipeline) || pipeline.length === 0) {
return false;
}
for (const stage of pipeline) {
if (stage == null || typeof stage !== 'object') {
continue;
}
const record = stage as CommandDocument;
for (const name of COLLECTION_WRITING_STAGES) {
if (name in record) {
return false;
}
}
if ('$lookup' in record && !lookupTargetsTenant(record.$lookup, tenantId)) {
return false;
}
}
const first = pipeline[0];
if (first == null || typeof first !== 'object') {
return false;
}
return pinsTenant((first as CommandDocument).$match, tenantId);
}
/** Every predicate a command carries, so a partially-scoped batch still fails. */
function predicatesOf(command: CommandDocument, commandName: string): unknown[] {
const path = PREDICATE_PATHS.get(commandName);
if (!path) {
return [];
}
if (!path.list) {
return [command[path.key]];
}
const operations = command[path.list];
if (!Array.isArray(operations)) {
return [];
}
return path.key ? operations.map((operation) => operation?.[path.key]) : operations;
}
/** A command reduced to the operation the probe should actually judge. */
interface ResolvedCommand {
readonly command: CommandDocument;
readonly commandName: string;
readonly collection: string;
}
/**
* Resolves the operation a wire command represents, unwrapping `explain` —
* whose payload is the nested command rather than a collection name, and which
* would otherwise be dropped even though `executionStats` reveals cross-tenant
* cardinality.
*/
function resolveCommand(
command: CommandDocument,
commandName: string,
): ResolvedCommand | undefined {
if (commandName === 'explain') {
const inner = command.explain;
if (inner == null || typeof inner !== 'object') {
return undefined;
}
const innerCommand = inner as CommandDocument;
const [innerName] = Object.keys(innerCommand);
return innerName ? resolveCommand(innerCommand, innerName) : undefined;
}
// `getMore` carries the cursor id under its own name and the collection separately.
const value = commandName === 'getMore' ? command.collection : command[commandName];
if (typeof value !== 'string') {
return undefined;
}
return { command, commandName, collection: value };
}
/** The predicates a command carries, for test failure output. */
function describePredicates(command: CommandDocument, commandName: string): string {
if (commandName === 'aggregate') {
return safeStringify(command.pipeline);
}
if (commandName === 'getMore') {
return '(cursor continuation — no predicate)';
}
return safeStringify(predicatesOf(command, commandName));
}
/**
* Whether a collation applies, checking the batched operations as well as the
* envelope: `update` and `delete` carry collation inside each `updates[]` /
* `deletes[]` entry, where a command-level check never sees it.
*/
function carriesCollation(command: CommandDocument, commandName: string): boolean {
if ('collation' in command) {
return true;
}
const path = PREDICATE_PATHS.get(commandName);
if (path?.list == null) {
return false;
}
const operations = command[path.list];
if (!Array.isArray(operations)) {
return false;
}
return operations.some(
(operation) => operation != null && typeof operation === 'object' && 'collation' in operation,
);
}
/** Whether this command is one the probe can judge at all. */
function carriesPredicate(command: CommandDocument, commandName: string): boolean {
if (UNJUDGED_COMMANDS.has(commandName)) {
return false;
}
if (commandName === 'aggregate') {
return Array.isArray(command.pipeline);
}
if (!PREDICATE_PATHS.has(commandName)) {
return true;
}
return predicatesOf(command, commandName).length > 0;
}
/**
* Whether every predicate this command carries targets `tenantId`.
* Only called for commands `carriesPredicate` has already accepted.
*/
function isScoped(command: CommandDocument, commandName: string, tenantId: string): boolean {
// A collation can broaden equality — under a case-insensitive one,
// `tenantId: 'tenant-a'` also matches `TENANT-A` — so literal equality no
// longer proves isolation and the command cannot be called scoped.
if (carriesCollation(command, commandName)) {
return false;
}
if (commandName === 'aggregate') {
return pipelineTargetsTenant(command.pipeline, tenantId);
}
// A cursor continuation carries no predicate, and the scope it was opened
// under is not visible here, so it fails closed rather than passing unseen.
if (commandName === 'getMore') {
return false;
}
if (!PREDICATE_PATHS.has(commandName)) {
return false;
}
return predicatesOf(command, commandName).every((predicate) => pinsTenant(predicate, tenantId));
}
export function attachTenantProbe(
connection: Connection,
collections: Iterable<string>,
): TenantProbe {
const watched = new Set(collections);
const databaseName = connection.db?.databaseName;
let recording: TenantCommandRecord[] | null = null;
const onCommandStarted = (event: {
commandName: string;
command: CommandDocument;
databaseName?: string;
}): void => {
if (recording == null) {
return;
}
// A shared MongoClient emits for every database it serves, and collection
// names are only unique within one.
if (event.databaseName != null && event.databaseName !== databaseName) {
return;
}
const resolved = resolveCommand(event.command, event.commandName);
if (resolved == null || !watched.has(resolved.collection)) {
return;
}
const { command, commandName, collection } = resolved;
if (!carriesPredicate(command, commandName)) {
return;
}
// Async context reaches this handler, so each command is judged against the
// tenant that was actually active when it was issued.
const tenantId = getTenantId();
const predicate = describePredicates(command, commandName);
if (tenantId == null || tenantId === SYSTEM_TENANT_ID) {
recording.push({ commandName, collection, scoped: false, predicate });
return;
}
recording.push({
commandName,
collection,
tenantId,
scoped: isScoped(command, commandName, tenantId),
predicate,
});
};
connection.getClient().on('commandStarted', onCommandStarted);
return {
async record(fn) {
const captured: TenantCommandRecord[] = [];
recording = captured;
try {
await fn();
} finally {
recording = null;
}
return captured;
},
close() {
connection.getClient().off('commandStarted', onCommandStarted);
},
};
}
/** Convenience for assertions: the commands that reached the wire unscoped. */
export function unscoped(records: readonly TenantCommandRecord[]): readonly TenantCommandRecord[] {
return records.filter((record) => !record.scoped);
}