From b8e1f167b99aff44fe9b9a5ed60665532fb8e1bb Mon Sep 17 00:00:00 2001 From: root Date: Thu, 16 Jul 2026 05:52:22 +0200 Subject: [PATCH] =?UTF-8?q?fix:=20=E8=BF=87=E6=BB=A4=E8=A7=82=E5=AF=9F?= =?UTF-8?q?=E6=B5=81=E7=BC=BA=E5=A4=B1=E8=BA=AB=E4=BB=BD=E4=BA=8B=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../workbench-kafka-refresh-handoff.test.ts | 11 ++++---- .../cloud/workbench-kafka-refresh-handoff.ts | 25 ++++++++++++++----- 2 files changed, 25 insertions(+), 11 deletions(-) diff --git a/internal/cloud/workbench-kafka-refresh-handoff.test.ts b/internal/cloud/workbench-kafka-refresh-handoff.test.ts index ffdf35a3..a60443eb 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, filteredNonHwlab: 0, replayed: 3, buffered: 3, bufferedDelivered: 2, liveDelivered: 0, deduplicated: 1, resumeSkipped: 0 }); + assert.deepEqual(summary.counts, { matched: 3, filteredNonHwlab: 0, filteredMissingIdentity: 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,7 +101,7 @@ 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 () => { +test("unscoped observer filters non-deliverable history 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" } }; @@ -111,16 +111,17 @@ test("unscoped observer filters legacy non-HWLAB records without weakening scope 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" }]); + return queryResult([legacy, { ...record(1), value: { schema: "hwlab.event.v1", eventType: "legacy-without-identity" } }, record(2)], [{ partition: 0, startOffset: "0", endOffset: "3" }]); }, 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.deepEqual(delivered, ["hwlab:evt_0_2"]); + assert.equal(summary.counts.matched, 3); assert.equal(summary.counts.filteredNonHwlab, 2); + assert.equal(summary.counts.filteredMissingIdentity, 1); assert.equal(summary.counts.replayed, 1); }); diff --git a/internal/cloud/workbench-kafka-refresh-handoff.ts b/internal/cloud/workbench-kafka-refresh-handoff.ts index 5a89b199..36498e16 100644 --- a/internal/cloud/workbench-kafka-refresh-handoff.ts +++ b/internal/cloud/workbench-kafka-refresh-handoff.ts @@ -35,6 +35,7 @@ export function createWorkbenchKafkaRefreshHandoff(options = {}) { const counts = { matched: 0, filteredNonHwlab: 0, + filteredMissingIdentity: 0, replayed: 0, buffered: 0, bufferedDelivered: 0, @@ -101,9 +102,12 @@ export function createWorkbenchKafkaRefreshHandoff(options = {}) { function onLiveEnvelope(envelope, transport) { if (phase === "closed" || phase === "failed") return; - if (allowUnscoped && envelope?.schema !== "hwlab.event.v1") { - counts.filteredNonHwlab += 1; - return; + if (allowUnscoped) { + const filtered = unscopedEnvelopeFilter(envelope); + if (filtered) { + counts[filtered] += 1; + return; + } } if (!envelopeMatchesScope(envelope, sessionId, traceId)) return; if (phase !== "live") { @@ -192,9 +196,12 @@ 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 (allowUnscoped) { + const filtered = unscopedEnvelopeFilter(envelope); + if (filtered) { + counts[filtered] += 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); @@ -450,6 +457,12 @@ function envelopeMatchesScope(envelope, sessionId, traceId) { return true; } +function unscopedEnvelopeFilter(envelope) { + if (!envelope || typeof envelope !== "object" || Array.isArray(envelope) || envelope.schema !== "hwlab.event.v1") return "filteredNonHwlab"; + if (!textValue(envelope.eventId) && !textValue(envelope.sourceEventId)) return "filteredMissingIdentity"; + return null; +} + function requireTransportIdentity(value, expectedTopic, phase) { const topic = textValue(value?.topic); const partition = partitionValue(value?.partition);