Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -336,26 +340,103 @@ 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')))

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 () => {
Expand Down Expand Up @@ -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
Expand All @@ -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 () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, ScrubbedRef>()
/**
* The batch currently in flight, so shutdown can interrupt it.
*
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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
Comment on lines +667 to +669

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift

🔎 Supported by static analysis

🏁 Script executed:

#!/bin/bash
set -eu

printf '%s\n' '--- changed-file diff ---'
git diff --unified=35 69a4765903776a5e70d4633590673395b1eaec2f e10294b97e4aadf25b7b541ed6e3808241c38642 -- nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts

printf '%s\n' '--- reviewed source around URL planning and hand-off ---'
sed -n '560,735p' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts
sed -n '820,925p' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts

printf '%s\n' '--- direct lifecycle references ---'
rg -n -C 8 'seenRefs|forgetUnwritten|handleBatch|write|month.?key|missing|available' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts

Repository: PostHog/posthog

Length of output: 41880


🤖 get_repo_knowledge executed:

get_repo_knowledge PostHog/posthog /tmp/coderabbit-repo-knowledge/posthog-posthog-1c883bd0/conventions

Length of output: 40997


🏁 Script executed:

#!/bin/bash
set -eu
file=nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts
test=nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.test.ts

printf '%s\n' '--- image-batcher outline ---'
ast-grep outline "$file"

printf '%s\n' '--- exact planning and batch lifecycle ---'
nl -ba "$file" | sed -n '300,560p'
nl -ba "$file" | sed -n '620,735p'
nl -ba "$file" | sed -n '820,930p'

printf '%s\n' '--- exact tests for missing keys and overlapping hand-offs ---'
nl -ba "$test" | sed -n '330,455p'
nl -ba "$test" | sed -n '1325,1380p'

printf '%s\n' '--- key reader and month-key consumers ---'
rg -n -C 12 'class .*Reader|read\\(|strongly consistent|month key|tableKeyString|sessionMonth|readKeys' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub

Repository: PostHog/posthog

Length of output: 38091


🤖 get_repo_knowledge executed:

get_repo_knowledge PostHog/posthog /tmp/coderabbit-repo-knowledge/posthog-posthog-1c883bd0

Length of output: 3910


🏁 Script executed:

#!/bin/bash
set -eu
file=nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub/image-batcher.ts

printf '%s\n' '--- handleBatch and hand-off wait path ---'
nl -ba "$file" | sed -n '230,315p'
nl -ba "$file" | sed -n '536,620p'
nl -ba "$file" | sed -n '930,1015p'

printf '%s\n' '--- key reader declarations and implementation files ---'
rg -n -F 'class MlKeyReader' nodejs/src
rg -n -F 'class KeyReader' nodejs/src
rg -n -F 'imageKeyId' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub
rg -n -F 'delete' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub

Repository: PostHog/posthog

Length of output: 13700


🤖 get_repo_knowledge executed:

get_repo_knowledge PostHog/posthog /tmp/coderabbit-repo-knowledge/posthog-posthog-1c883bd0

Length of output: 88


🏁 Script executed:

#!/bin/bash
set -u

printf '%s\n' '--- MlKeyReader source ---'
nl -ba nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys/reader.ts

printf '%s\n' '--- key schema and key-manager consumers ---'
rg -n -C 8 'imageKeyId|tableKeyString|reader\.read|MlKeyReader|sessionMonth|month key' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub || true

printf '%s\n' '--- repository month-key writes/deletes in the ML mirror area ---'
rg -n -C 8 'delete.*key|key.*delete|write.*key|put.*key|tableKeyString|imageKeyId' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror-image-scrub || true

Repository: PostHog/posthog

Length of output: 45667


🤖 get_repo_knowledge executed:

get_repo_knowledge PostHog/posthog /tmp/coderabbit-repo-knowledge/posthog-posthog-1c883bd0/conventions

Length of output: 40997


🏁 Script executed:

#!/bin/bash
set -u
printf '%s\n' '--- MlSessionKeyStore implementation ---'
rg -n -F 'class MlSessionKeyStore' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys
file=$(rg -l -F 'class MlSessionKeyStore' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys | head -n 1)
if [ -n "$file" ]; then
  nl -ba "$file" | sed -n '1,260p'
fi

printf '%s\n' '--- month-key creation and repair evidence ---'
rg -n -C 12 'imageKeyId|month.*key|key.*month|repair|conditional|deleted' nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys/key-store.ts nodejs/src/ingestion/pipelines/sessionreplay/ml-mirror/keys/key-store.test.ts 2>/dev/null || true

Repository: PostHog/posthog

Length of output: 42351


🤖 get_repo_knowledge executed:

get_repo_knowledge PostHog/posthog /tmp/coderabbit-repo-knowledge/posthog-posthog-1c883bd0

Length of output: 3891


Do not deduplicate URL refs while their write is pending.

handleBatch() can return while the first hand-off awaits MlKeyReader.read(). A later batch can then hit seenRefs.has(ref), skip the copy, and record its offset. If the first read finds no month key, write() removes the first mark and stores its offsets. Both offsets can advance without storing either URL image.

MlKeyBatch can provision a missing month key, and its overlapping-batch path allows a later read before the earlier write. Keep pending URL copies retryable until the owning hand-off succeeds, or wait for the owner outcome before skipping them. Add a regression test for overlapping hand-offs.

}
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)
Expand All @@ -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)
}
}
}
Expand Down Expand Up @@ -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'
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down
Loading