diff --git a/internal/cloud/workbench-kafka-refresh-handoff.test.ts b/internal/cloud/workbench-kafka-refresh-handoff.test.ts index 4c97719f..ffdf35a3 100644 --- a/internal/cloud/workbench-kafka-refresh-handoff.test.ts +++ b/internal/cloud/workbench-kafka-refresh-handoff.test.ts @@ -41,7 +41,7 @@ test("retention handoff replays once, deduplicates overlap and late pre-barrier "buffered-live:4:hwlab:evt_0_4", "connected:1" ]); - assert.deepEqual(summary.counts, { matched: 3, replayed: 3, buffered: 3, bufferedDelivered: 2, liveDelivered: 0, deduplicated: 1, resumeSkipped: 0 }); + assert.deepEqual(summary.counts, { matched: 3, filteredNonHwlab: 0, replayed: 3, buffered: 3, bufferedDelivered: 2, liveDelivered: 0, deduplicated: 1, resumeSkipped: 0 }); publishLive(record(2).value, transport(2)); publishLive(record(5).value, transport(5)); @@ -101,6 +101,29 @@ test("retention emits handoff before buffered live flush and connected live", as assert.deepEqual(phases, ["replay:0", "replay:1", "handoff", "buffered-live:2", "live"]); }); +test("unscoped observer filters legacy non-HWLAB records without weakening scoped handoff", async () => { + let publishLive: any = null; + const delivered: string[] = []; + const legacy = { ...record(0), value: { schema: "platform-infra.kafka.shadow-produce.v1", eventId: "legacy-smoke" } }; + const handoff = createWorkbenchKafkaRefreshHandoff({ + allowUnscoped: true, + liveBufferLimit: 4, + subscribeLive(listener: any) { publishLive = listener; return () => {}; }, + async queryRetention() { + publishLive({ schema: "platform-infra.kafka.shadow-produce.v1", eventId: "legacy-live" }, transport(2)); + return queryResult([legacy, record(1)], [{ partition: 0, startOffset: "0", endOffset: "2" }]); + }, + async deliverEvent(envelope: any) { delivered.push(envelope.eventId); return true; }, + async deliverConnected() { return true; } + }); + + const summary = await handoff.start(); + assert.deepEqual(delivered, ["hwlab:evt_0_1"]); + assert.equal(summary.counts.matched, 2); + assert.equal(summary.counts.filteredNonHwlab, 2); + assert.equal(summary.counts.replayed, 1); +}); + test("retention handoff fails closed on a pre-barrier live gap", async () => { let publishLive: any = null; const handoff = createWorkbenchKafkaRefreshHandoff({ diff --git a/internal/cloud/workbench-kafka-refresh-handoff.ts b/internal/cloud/workbench-kafka-refresh-handoff.ts index 4d6dd614..5a89b199 100644 --- a/internal/cloud/workbench-kafka-refresh-handoff.ts +++ b/internal/cloud/workbench-kafka-refresh-handoff.ts @@ -34,6 +34,7 @@ export function createWorkbenchKafkaRefreshHandoff(options = {}) { let pendingLiveDeliveries = 0; const counts = { matched: 0, + filteredNonHwlab: 0, replayed: 0, buffered: 0, bufferedDelivered: 0, @@ -99,7 +100,12 @@ export function createWorkbenchKafkaRefreshHandoff(options = {}) { } function onLiveEnvelope(envelope, transport) { - if (phase === "closed" || phase === "failed" || !envelopeMatchesScope(envelope, sessionId, traceId)) return; + if (phase === "closed" || phase === "failed") return; + if (allowUnscoped && envelope?.schema !== "hwlab.event.v1") { + counts.filteredNonHwlab += 1; + return; + } + if (!envelopeMatchesScope(envelope, sessionId, traceId)) return; if (phase !== "live") { if (bufferedLive.length >= liveBufferLimit) { fail(handoffError("workbench_kafka_refresh_live_buffer_overflow", "Kafka refresh live buffer exceeded its YAML-owned capacity before handoff.", phase, { liveBufferLimit })); @@ -186,6 +192,10 @@ export function createWorkbenchKafkaRefreshHandoff(options = {}) { const normalized = []; for (const record of records) { const envelope = record?.value; + if (allowUnscoped && envelope?.schema !== "hwlab.event.v1") { + counts.filteredNonHwlab += 1; + continue; + } if (!envelopeMatchesScope(envelope, sessionId, traceId) || envelope?.schema !== "hwlab.event.v1") { throw handoffError("workbench_kafka_refresh_scope_mismatch", "Kafka retention query returned an envelope outside the authorized scope.", phase); }