Merge pull request #2560 from pikasTech/fix/2541-unscoped-replay
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success

修复:过滤观察流历史非目标事件
This commit is contained in:
Lyon
2026-07-16 11:45:49 +08:00
committed by GitHub
2 changed files with 35 additions and 2 deletions
@@ -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({
@@ -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);
}