Skip to content

Commit 94d2650

Browse files
committed
fix: close query consistency tails
1 parent ac7b3b3 commit 94d2650

9 files changed

Lines changed: 157 additions & 8 deletions

File tree

‎changelogs/unreleased.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414

1515
## Fixed / changed
1616

17+
- Fixed aggregate direct `.toArray()` to honor cache/meta execution paths, extended `find()` ObjectId auto-conversion to comparison operators, forwarded CountQueue abort signals into MongoDB count options, and added a warning when sync idempotency falls back to in-memory storage.
1718
- Added a short-lived read-cache dirty barrier around writes and transaction commits. Cached reads now bypass and avoid refilling query cache while a namespace is being invalidated, reducing stale-cache windows when a process exits between a database write and post-write invalidation.
1819
- Added optional Change Stream sync idempotency gates (`sync.idempotency`) with per-target keys and duplicate stats, so supervised restarts can skip targets already marked as applied before saving the shared resume token.
1920
- Added strict optimistic-locking support to Model `updateBatch(..., { versionMode: 'strict' })`; default `counter` behavior remains unchanged.

‎src/adapters/mongodb/common/collection-accessor.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -294,7 +294,11 @@ export class MongoCollectionAccessor<TSchema extends Document = Document> {
294294
const cacheTTL = typeof merged.cache === 'number' ? merged.cache : 0;
295295
const { cache: _cache, ...keyOptions } = merged;
296296
void _cache;
297-
const executeCount = () => countDocuments(this.collectionRef, normalizedQuery ?? {}, keyOptions as Parameters<Collection<TSchema>['countDocuments']>[1]);
297+
const executeCount = (signal?: AbortSignal) => countDocuments(
298+
this.collectionRef,
299+
normalizedQuery ?? {},
300+
(signal ? { ...keyOptions, signal } : keyOptions) as Parameters<Collection<TSchema>['countDocuments']>[1],
301+
);
298302
const countQueue = this.management.defaults?.countQueue;
299303
if (
300304
cacheTTL > 0

‎src/adapters/mongodb/queries/index.ts‎

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -413,7 +413,7 @@ class AggregateChain<TResult = unknown, TSchema extends Document = Document> imp
413413
return wrapQueryResultWithMeta(this.collection, this.defaults, op, this.options, startTs, result);
414414
}
415415

416-
toArray(): Promise<TResult[]> {
416+
private runToArray(): Promise<TResult[]> {
417417
if (this.executed) {
418418
throw createError(ErrorCodes.INVALID_OPERATION, 'Query already executed.');
419419
}
@@ -426,6 +426,13 @@ class AggregateChain<TResult = unknown, TSchema extends Document = Document> imp
426426
});
427427
}
428428

429+
toArray(): Promise<TResult[]> {
430+
if (this.executed) {
431+
throw createError(ErrorCodes.INVALID_OPERATION, 'Query already executed.');
432+
}
433+
return this.executeResult() as Promise<TResult[]>;
434+
}
435+
429436
private async executeResult(): Promise<TResult[] | ResultWithMeta<TResult[]> | NodeJS.ReadableStream> {
430437
const startTs = Date.now();
431438

@@ -453,11 +460,11 @@ class AggregateChain<TResult = unknown, TSchema extends Document = Document> imp
453460
return this.wrapResult('aggregate', startTs, cached as TResult[]) as TResult[] | ResultWithMeta<TResult[]>;
454461
}
455462
const qc = this.queryCache;
456-
const result = await this.toArray();
463+
const result = await this.runToArray();
457464
await Promise.resolve(qc.set(cacheKey, result, cacheTTL));
458465
return this.wrapResult('aggregate', startTs, result) as TResult[] | ResultWithMeta<TResult[]>;
459466
}
460-
return this.toArray().then((result) => this.wrapResult('aggregate', startTs, result) as TResult[] | ResultWithMeta<TResult[]>);
467+
return this.runToArray().then((result) => this.wrapResult('aggregate', startTs, result) as TResult[] | ResultWithMeta<TResult[]>);
461468
}
462469

463470
then<TResult1 = TResult[], TResult2 = never>(
@@ -582,6 +589,7 @@ export async function countDocuments<TSchema extends Document = Document>(
582589
const countOptions = buildCountDriverOptions<TSchema>(rawOptions);
583590
const maxTimeMS = rawOptions.maxTimeMS as number | undefined;
584591
const comment = rawOptions.comment as string | undefined;
592+
const signal = rawOptions.signal as AbortSignal | undefined;
585593
const canUseEstimatedCount = isEmptyQuery
586594
&& !hasSessionOption(rawOptions)
587595
&& rawOptions.collation === undefined
@@ -610,6 +618,7 @@ export async function countDocuments<TSchema extends Document = Document>(
610618
const estimatedOptions: Record<string, unknown> = {};
611619
if (maxTimeMS !== undefined) estimatedOptions.maxTimeMS = maxTimeMS;
612620
if (comment) estimatedOptions.comment = comment;
621+
if (signal) estimatedOptions.signal = signal;
613622
return collection.estimatedDocumentCount(estimatedOptions as Parameters<Collection<TSchema>['estimatedDocumentCount']>[0]);
614623
}
615624

‎src/adapters/mongodb/queries/query-helpers.ts‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -144,13 +144,16 @@ function normalizeQueryArray(
144144
});
145145
}
146146

147+
const QUERY_OBJECTID_ARRAY_OPERATORS = new Set(['$in', '$nin', '$all']);
148+
const QUERY_OBJECTID_SCALAR_OPERATORS = new Set(['$eq', '$ne', '$gt', '$gte', '$lt', '$lte']);
149+
147150
/**
148151
* Recursively normalizes a query filter, converting eligible 24-char hex strings to ObjectId.
149152
*
150153
* Notes:
151154
* - This is the query-path-specific filter normalizer.
152155
* - Unlike `utils/objectid-converter.ts` (general object traversal), this one preserves
153-
* query operator semantics: $and/$or/$nor/$elemMatch/$in/$nin/$all/$eq/$ne.
156+
* query operator semantics: $and/$or/$nor/$elemMatch/$in/$nin/$all/$eq/$ne/$gt/$gte/$lt/$lte.
154157
*/
155158
export function normalizeQueryFilter(
156159
filter: Record<string, unknown>,
@@ -188,13 +191,13 @@ export function normalizeQueryFilter(
188191
if (hasOperators) {
189192
const nestedResult: Record<string, unknown> = {};
190193
for (const [op, opVal] of Object.entries(nested)) {
191-
if (shouldConvert && (op === '$in' || op === '$nin' || op === '$all') && Array.isArray(opVal)) {
194+
if (shouldConvert && QUERY_OBJECTID_ARRAY_OPERATORS.has(op) && Array.isArray(opVal)) {
192195
nestedResult[op] = opVal.map((item) =>
193196
typeof item === 'string' && item.length === 24 && ObjectId.isValid(item)
194197
? new ObjectId(item)
195198
: item,
196199
);
197-
} else if (shouldConvert && (op === '$eq' || op === '$ne') && typeof opVal === 'string' && opVal.length === 24 && ObjectId.isValid(opVal)) {
200+
} else if (shouldConvert && QUERY_OBJECTID_SCALAR_OPERATORS.has(op) && typeof opVal === 'string' && opVal.length === 24 && ObjectId.isValid(opVal)) {
198201
nestedResult[op] = new ObjectId(opVal);
199202
} else if (op === '$elemMatch' && opVal && typeof opVal === 'object' && !Array.isArray(opVal)) {
200203
nestedResult[op] = normalizeQueryFilter(opVal as Record<string, unknown>, autoConvert, currentPath, depth + 1);
@@ -488,6 +491,7 @@ export function buildCountDriverOptions<TSchema extends Document = Document>(
488491
'session',
489492
'readConcern',
490493
'readPreference',
494+
'signal',
491495
]);
492496
for (const key of Object.keys(driverOptions)) {
493497
if (!allowed.has(key)) {

‎src/capabilities/sync/index.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -409,6 +409,9 @@ export class ChangeStreamSyncManager {
409409
});
410410
this.clientFactory = options.clientFactory ?? defaultClientFactory;
411411
this.idempotency = resolveSyncIdempotencyRuntime(options.config.idempotency);
412+
if (options.config.idempotency?.enabled === true && options.config.idempotency.store === undefined) {
413+
this.logger?.warn?.('[Sync] sync.idempotency is enabled without a durable store; the memory fallback only protects replay within the current process.');
414+
}
412415
}
413416

414417
/**

‎test/integration/mongodb/aggregate.test.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,6 +220,36 @@ describe('aggregate() / distinct() / count() / explain()', () => {
220220
assert.equal(third[0].total, 21);
221221
});
222222

223+
it('direct .toArray() honors aggregate cache option', async () => {
224+
const db = runtime._adapter.db;
225+
const pipeline = [
226+
{ $match: { category: 'clothing' } },
227+
{ $group: { _id: null, total: { $sum: 1 } } },
228+
];
229+
230+
const first = await col.aggregate(pipeline, { cache: 60_000 }).toArray();
231+
assert.equal(first[0].total, 20);
232+
233+
await db.collection('products').insertOne({
234+
sku: 'SKU-CACHE-AGG-TOARRAY',
235+
category: 'clothing',
236+
price: 1,
237+
stock: 1,
238+
rating: 5,
239+
active: true,
240+
});
241+
242+
const second = await col.aggregate(pipeline, { cache: 60_000 }).toArray();
243+
assert.equal(second[0].total, 20);
244+
});
245+
246+
it('direct .toArray() honors aggregate meta option', async () => {
247+
const result = await col.aggregate([{ $match: { category: 'food' } }], { meta: true }).toArray() as any;
248+
assert.ok(Array.isArray(result.data));
249+
assert.equal(result.meta.op, 'aggregate');
250+
assert.equal(result.meta.ns.coll, 'products');
251+
});
252+
223253
it('invalidates target collection cache after $merge', async () => {
224254
const db = runtime._adapter.db;
225255
const summary = runtime.collection('product_summary');

‎test/integration/mongodb/queries-chain-extended.test.ts‎

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -289,6 +289,7 @@ describe('ObjectId auto-conversion — integration', () => {
289289
let runtime2: any;
290290
let col2: any;
291291
const VALID_HEX = '507f1f77bcf86cd799439011';
292+
const NEXT_HEX = '507f1f77bcf86cd799439012';
292293

293294
before(async () => {
294295
const ctx = await bootstrap2.setup();
@@ -300,7 +301,7 @@ describe('ObjectId auto-conversion — integration', () => {
300301
const { ObjectId } = require('mongodb');
301302
await col2.insertMany([
302303
{ name: 'X', userId: new ObjectId(VALID_HEX), score: 5 },
303-
{ name: 'Y', userId: new ObjectId('507f1f77bcf86cd799439012'), score: 3 },
304+
{ name: 'Y', userId: new ObjectId(NEXT_HEX), score: 3 },
304305
{ name: 'Z', score: 1, tags: [new ObjectId(VALID_HEX)] },
305306
]);
306307
});
@@ -323,6 +324,20 @@ describe('ObjectId auto-conversion — integration', () => {
323324
assert.ok(docs.length >= 0); // might or might not match depending on conversion
324325
});
325326

327+
it('find converts ObjectId comparison operator values', async () => {
328+
const gtDocs = await col2.find({ userId: { $gt: VALID_HEX } }).toArray();
329+
assert.deepEqual(gtDocs.map((doc: { name: string }) => doc.name), ['Y']);
330+
331+
const gteDocs = await col2.find({ userId: { $gte: VALID_HEX } }).sort({ name: 1 }).toArray();
332+
assert.deepEqual(gteDocs.map((doc: { name: string }) => doc.name), ['X', 'Y']);
333+
334+
const ltDocs = await col2.find({ userId: { $lt: NEXT_HEX } }).toArray();
335+
assert.deepEqual(ltDocs.map((doc: { name: string }) => doc.name), ['X']);
336+
337+
const lteDocs = await col2.find({ userId: { $lte: VALID_HEX } }).toArray();
338+
assert.deepEqual(lteDocs.map((doc: { name: string }) => doc.name), ['X']);
339+
});
340+
326341
it('find with $expr operator — SPECIAL_OPERATORS skips conversion', async () => {
327342
const docs = await col2.find({ $expr: { $gt: ['$score', 2] } }).toArray();
328343
assert.ok(Array.isArray(docs));

‎test/unit/runtime/runtime-compat.test.ts‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -525,6 +525,33 @@ describe('P6 runtime compat mock path', () => {
525525
assert.equal(calls.distinct, 1);
526526
});
527527

528+
it('passes CountQueue AbortSignal into count driver options', async () => {
529+
const controller = new AbortController();
530+
let capturedOptions: Record<string, unknown> | undefined;
531+
const nativeCollection = {
532+
namespace: 'compat_db.users',
533+
countDocuments(_query: unknown, options?: Record<string, unknown>) {
534+
capturedOptions = options;
535+
return resolved(7);
536+
},
537+
};
538+
const accessor = new MongoCollectionAccessor(
539+
'compat_db',
540+
'users',
541+
nativeCollection as never,
542+
{
543+
defaults: {
544+
countQueue: {
545+
execute: <T>(fn: (signal?: AbortSignal) => Promise<T>) => fn(controller.signal),
546+
},
547+
},
548+
},
549+
);
550+
551+
assert.equal(await accessor.count({ active: true } as never), 7);
552+
assert.equal(capturedOptions?.signal, controller.signal);
553+
});
554+
528555
it('records write cache invalidation on the active transaction session', async () => {
529556
const recorded: string[] = [];
530557
let deletedPatterns = 0;

‎test/unit/sync/sync.test.ts‎

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -490,6 +490,62 @@ describe('P4-C sync', () => {
490490
assert.equal(stats.duplicateEventCount, 1);
491491
});
492492

493+
it('warns when sync idempotency uses the in-memory fallback store', () => {
494+
const warnings: string[] = [];
495+
const db = {
496+
databaseName: 'source_db',
497+
watch: () => ({ close: () => Promise.resolve(true) }),
498+
};
499+
500+
new MonSQLize.ChangeStreamSyncManager({
501+
db,
502+
config: {
503+
enabled: true,
504+
idempotency: { enabled: true },
505+
targets: [{ name: 'target', apply: async () => undefined }],
506+
},
507+
logger: {
508+
warn: (message: string) => warnings.push(message),
509+
error: () => undefined,
510+
info: () => undefined,
511+
debug: () => undefined,
512+
},
513+
});
514+
515+
assert.ok(warnings.some((message) => message.includes('memory fallback only protects replay within the current process')));
516+
});
517+
518+
it('does not warn when sync idempotency receives an explicit store', () => {
519+
const warnings: string[] = [];
520+
const db = {
521+
databaseName: 'source_db',
522+
watch: () => ({ close: () => Promise.resolve(true) }),
523+
};
524+
525+
new MonSQLize.ChangeStreamSyncManager({
526+
db,
527+
config: {
528+
enabled: true,
529+
idempotency: {
530+
enabled: true,
531+
store: {
532+
get: () => undefined,
533+
set: () => undefined,
534+
},
535+
},
536+
targets: [{ name: 'target', apply: async () => undefined }],
537+
},
538+
logger: {
539+
warn: (message: string) => warnings.push(message),
540+
error: () => undefined,
541+
info: () => undefined,
542+
debug: () => undefined,
543+
},
544+
});
545+
546+
assert.equal(warnings.length, 0);
547+
});
548+
493549
it('stops processing changes when resume token persistence fails', async () => {
494550
const liveStream = new EventEmitter() as EventEmitter & { close(): Promise<boolean> };
495551
let closed = false;

0 commit comments

Comments
 (0)