diff --git a/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts b/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts index 56af639f0dcc..e1adac7fc8cc 100644 --- a/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts +++ b/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts @@ -3,7 +3,11 @@ import { gzipSync } from 'node:zlib' import { MlKeyReader } from '~/ingestion/pipelines/sessionreplay/ml-mirror/keys/reader' import { MlKeyManager } from '~/ingestion/pipelines/sessionreplay/ml-mirror/keys/runtime' -import { INGESTION_VERSION_HEADER } from '~/ingestion/pipelines/sessionreplay/ml-mirror/keys/schema' +import { + INGESTION_VERSION_HEADER, + imageKeyId, + tableKeyString, +} from '~/ingestion/pipelines/sessionreplay/ml-mirror/keys/schema' import { MlKafkaTransport } from '~/ingestion/pipelines/sessionreplay/ml-mirror/keys/transport' import { hashImageBytes, imageRef, urlRef } from './content-ref' @@ -336,15 +340,21 @@ describe('ImageBatcher', () => { ]) }) - it('scrubs a newer URL version after an earlier version was persisted by the same pod', async () => { + it('skips a URL ref in a later batch after this pod persisted it', async () => { const png = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]) const first = Buffer.concat([png, Buffer.from('first')]) const second = Buffer.concat([png, Buffer.from('second')]) const store = new FakeStore() + let scrubs = 0 const batcher = new ImageBatcher( store as unknown as ImageShardStore, new FakeOffsets(), - { scrub: (bytes: Buffer) => Promise.resolve(bytes) } as unknown as ScrubClient, + { + scrub: (bytes: Buffer) => { + scrubs += 1 + return Promise.resolve(bytes) + }, + } as unknown as ScrubClient, options ) const ref = urlRef(hashImageBytes(CONTENT_KEY, Buffer.from('https://example.com/image.png'))) @@ -352,10 +362,81 @@ describe('ImageBatcher', () => { await handleAndWrite(batcher, [msg(0, 10, pt(1), first, ref, [{ 'content-type': Buffer.from('image/png') }])]) await handleAndWrite(batcher, [msg(0, 11, pt(1), second, ref, [{ 'content-type': Buffer.from('image/png') }])]) - expect(store.urlWrites.map((image) => [image.sourceOffset, image.bytes])).toEqual([ - [10, first], - [11, second], - ]) + expect(scrubs).toBe(1) + expect(store.urlWrites.map((image) => [image.sourceOffset, image.bytes])).toEqual([[10, first]]) + }) + + it('stores a URL image in a later batch when its month key was missing at the first write', async () => { + const png = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]) + const ref = `imageurl:v3:7:2026-09:${'a'.repeat(22)}` + const monthKey = tableKeyString(imageKeyId(7, '2026-09')) + const readKeys = jest + .fn() + .mockResolvedValueOnce(new Map()) + .mockResolvedValue(new Map([[monthKey, { identity: { teamId: 7, sessionMonth: '2026-09' } }]])) + const keyManager = { + kafka: { + read: (messages: Message[]) => + Promise.resolve(messages.map((message) => ({ message, original: message, version: 2 }))), + }, + reader: { read: readKeys }, + } as unknown as MlKeyManager + const store = new FakeStore() + const batcher = new ImageBatcher( + store as unknown as ImageShardStore, + new FakeOffsets(), + { scrub: (bytes: Buffer) => Promise.resolve(bytes) } as unknown as ScrubClient, + options, + null, + keyManager + ) + const copy = (offset: number): Message => + msg(0, offset, pt(1), png, ref, [{ 'content-type': Buffer.from('image/png') }]) + + await handleAndWrite(batcher, [copy(10)]) + await handleAndWrite(batcher, [copy(11)]) + + expect(store.urlWrites.map((image) => image.sourceOffset)).toEqual([11]) + }) + + it('keeps a URL ref seen when a later copy of it is dropped after an earlier copy was stored', async () => { + const png = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]) + const ref = `imageurl:v3:7:2026-09:${'a'.repeat(22)}` + const monthKey = tableKeyString(imageKeyId(7, '2026-09')) + const readKeys = jest + .fn() + .mockResolvedValueOnce(new Map([[monthKey, { identity: { teamId: 7, sessionMonth: '2026-09' } }]])) + .mockResolvedValue(new Map()) + const keyManager = { + kafka: { + read: (messages: Message[]) => + Promise.resolve(messages.map((message) => ({ message, original: message, version: 2 }))), + }, + reader: { read: readKeys }, + } as unknown as MlKeyManager + const store = new FakeStore() + let scrubs = 0 + const batcher = new ImageBatcher( + store as unknown as ImageShardStore, + new FakeOffsets(), + { + scrub: (bytes: Buffer) => { + scrubs += 1 + return Promise.resolve(bytes) + }, + } as unknown as ScrubClient, + { ...options, maxImages: 1, scrubConcurrency: 1 }, + null, + keyManager + ) + const copy = (offset: number): Message => + msg(0, offset, pt(1), png, ref, [{ 'content-type': Buffer.from('image/png') }]) + + await handleAndWrite(batcher, [copy(10), copy(11)]) + await handleAndWrite(batcher, [copy(12)]) + + expect(store.urlWrites.map((image) => image.sourceOffset)).toEqual([10]) + expect(scrubs).toBe(2) }) it('bounds concurrent URL-image writes by the write concurrency, not the scrub slots', async () => { @@ -673,7 +754,21 @@ describe('ImageBatcher', () => { expect(offsets.received.at(-1)).toEqual([{ topic: 'session_replay_image_scrub', partition: 0, offset: 5 }]) }) - it('rescrubs an image whose batch failed rather than pod-deduping the replay away', async () => { + it.each([ + ['an inline image', () => msg(0, 0, pt(1), Buffer.from('sprite'))], + [ + 'a URL image', + () => + msg( + 0, + 0, + pt(1), + Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]), + urlRef(hashImageBytes(CONTENT_KEY, Buffer.from('https://example.com/image.png'))), + [{ 'content-type': Buffer.from('image/png') }] + ), + ], + ])('rescrubs %s whose batch failed rather than pod-deduping the replay away', async (_name, message) => { // A ref is marked seen only once its image is buffered. Marking it at plan time, or inside // scrubOne, reads as a harmless simplification and silently drops the image instead: the // replay is pod-deduped, its offset advances, and nothing ever writes it. No error, no @@ -687,12 +782,11 @@ describe('ImageBatcher', () => { } as unknown as ScrubClient const store = new FakeStore() const batcher = new ImageBatcher(store as unknown as ImageShardStore, new FakeOffsets(), flakyClient, options) - const sprite = Buffer.from('sprite') - await expect(batcher.handleBatch([msg(0, 0, pt(1), sprite)])).rejects.toThrow('sidecar down') - await handleAndWrite(batcher, [msg(0, 0, pt(1), sprite)]) + await expect(batcher.handleBatch([message()])).rejects.toThrow('sidecar down') + await handleAndWrite(batcher, [message()]) - expect(store.writes.flat()).toHaveLength(1) + expect(store.writes.flat().length + store.urlWrites.length).toBe(1) }) it('advances every partition past the duplicates it skipped, not just the one it scrubbed on', async () => { diff --git a/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts b/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts index 0dada6106c5d..4f211e0e0546 100644 --- a/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts +++ b/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts @@ -206,6 +206,8 @@ export class ImageBatcher { * a batch that throws discards what it staged. Sizing is a throughput question, not a correctness one. */ private readonly seenRefs: RefDedupCache + /** The copy that holds each ref's seen mark until its hand-off is written, so only the current owner can clear the mark, even after the cache evicts and re-marks the ref. */ + private readonly unwrittenMarkOwners = new Map() /** * The batch currently in flight, so shutdown can interrupt it. * @@ -481,8 +483,9 @@ export class ImageBatcher { // Marked here rather than on completion: a staged image is a local that a thrown // batch discards, so a ref marked before retirement could be skipped on replay // without ever having been persisted. - if (ready.source === 'bytes') { + if (!this.seenRefs.has(ready.ref)) { this.seenRefs.add(ready.ref) + this.unwrittenMarkOwners.set(ready.ref, ready) } staged[retired] = null stagedCount -= 1 @@ -597,6 +600,11 @@ export class ImageBatcher { ImageScrubConsumerMetrics.incBatchFailed('write') throw error } + for (const image of handoff.images) { + if (this.unwrittenMarkOwners.get(image.ref) === image) { + this.unwrittenMarkOwners.delete(image.ref) + } + } // Observed on success only, so an S3 incident's retry budgets do not read as slow writes. ImageScrubConsumerMetrics.observeWrite((performance.now() - startedAt) / 1000) } @@ -655,16 +663,21 @@ export class ImageBatcher { } else { urlLocationByRef.set(ref, candidate) } + // Copies in one batch all stay planned, so a later valid copy still stores when an earlier one fails validation. + if (this.seenRefs.has(ref)) { + ImageScrubConsumerMetrics.incDeduped('pod', 'url') + continue + } planned.push(candidate) continue } if (inlineRefs.has(ref)) { - ImageScrubConsumerMetrics.incDeduped('batch') + ImageScrubConsumerMetrics.incDeduped('batch', 'inline') continue } inlineRefs.add(ref) if (this.seenRefs.has(ref)) { - ImageScrubConsumerMetrics.incDeduped('pod') + ImageScrubConsumerMetrics.incDeduped('pod', 'inline') continue } planned.push(candidate) @@ -674,9 +687,10 @@ export class ImageBatcher { /** A ref that was marked seen but never persisted would be deduped away unwritten if its partition came back here. */ private forgetUnwritten(images: ScrubbedRef[]): void { - for (const { ref, source } of images) { - if (source === 'bytes') { - this.seenRefs.delete(ref) + for (const image of images) { + if (this.unwrittenMarkOwners.get(image.ref) === image) { + this.unwrittenMarkOwners.delete(image.ref) + this.seenRefs.delete(image.ref) } } } @@ -870,11 +884,14 @@ export class ImageBatcher { ) ) : new Map() + // A strongly consistent read finds no month key only after a team or month deletion, which a conditional put never reverses, so every later copy of these refs is dropped the same way. const keyed = handoff.images.filter( ({ image }) => image.sessionMonth === undefined || imageKeys.has(tableKeyString(imageKeyId(Number(image.teamId), image.sessionMonth!))) ) + const keyedSet = new Set(keyed) + this.forgetUnwritten(handoff.images.filter((item) => !keyedSet.has(item))) const inlineItems = keyed.filter( (item): item is ScrubbedRef & { image: ScrubbedImage } => item.source === 'bytes' ) diff --git a/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/metrics.ts b/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/metrics.ts index d1db1035acac..ac2bbf1f8219 100644 --- a/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/metrics.ts +++ b/nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/metrics.ts @@ -26,8 +26,8 @@ export class ImageScrubConsumerMetrics { }) private static readonly deduped = new Counter({ name: 'ml_mirror_image_scrub_consumer_deduped_total', - help: 'Messages skipped as duplicate produces of a ref, by scope: "batch" (another copy in the same poll batch) or "pod" (this pod scrubbed it earlier). Dedup hit rate = deduped / (deduped + scrubbed + skipped); the batch/pod split says how much the retained seen-ref cache is earning over free intra-batch dedup', - labelNames: ['scope'], + help: 'Messages skipped as duplicate produces of a ref, by scope: "batch" (another copy in the same poll batch) or "pod" (this pod scrubbed it earlier), and by source: "inline" or "url". URL refs dedup only by pod, because every copy in one batch stays planned. Dedup hit rate = deduped / (deduped + scrubbed + skipped); the batch/pod split says how much the retained seen-ref cache is earning over free intra-batch dedup', + labelNames: ['scope', 'source'], }) /** * Intra-batch dedup can only collapse copies that arrive in the same poll batch, so its ceiling is @@ -232,8 +232,8 @@ export class ImageScrubConsumerMetrics { public static incOffsetsDiscarded(count: number): void { this.offsetsDiscarded.inc(count) } - public static incDeduped(scope: 'batch' | 'pod'): void { - this.deduped.labels(scope).inc() + public static incDeduped(scope: 'batch' | 'pod', source: ImageScrubSource): void { + this.deduped.labels(scope, source).inc() } public static incUrlImageWrite(outcome: UrlImageWriteOutcome): void { this.urlImageWrites.labels(outcome).inc()