Merge pull request #2754 from pikasTech/fix/2753-agent-observer-retention-window
Pipelines as Code CI / hwlab-nc01-v03-ci-poll-2c7e867c74d3b2ac0497c9d11f79a8fc806cb246 Success
Pipelines as Code CI / hwlab-nc01-v03-ci-poll-2c7e867c74d3b2ac0497c9d11f79a8fc806cb246 Success
修复 Agent observer Kafka retention 到 live 交接
This commit is contained in:
@@ -6,7 +6,7 @@ import { test } from "bun:test";
|
|||||||
|
|
||||||
import { createCloudApiBunServer } from "./bun-server.ts";
|
import { createCloudApiBunServer } from "./bun-server.ts";
|
||||||
import { buildCloudApiReadiness } from "./health-contract.ts";
|
import { buildCloudApiReadiness } from "./health-contract.ts";
|
||||||
import { decodeCanonicalAgentRunKafkaMessage, kafkaDnsLookup, kafkaEventBridgeConfig, kafkaMessageKeyCanMatchPartitionKey, kafkaPartitionForKey, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts";
|
import { decodeCanonicalAgentRunKafkaMessage, kafkaDnsLookup, kafkaEventBridgeConfig, kafkaMessageKeyCanMatchPartitionKey, kafkaPartitionForKey, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, queryKafkaEventStream, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts";
|
||||||
|
|
||||||
test("Workbench refresh replay resolves the same Kafka partition as the session-keyed producer", () => {
|
test("Workbench refresh replay resolves the same Kafka partition as the session-keyed producer", () => {
|
||||||
const offsets = [
|
const offsets = [
|
||||||
@@ -30,6 +30,78 @@ test("Workbench refresh replay rejects foreign keyed values before JSON parsing"
|
|||||||
assert.equal(kafkaMessageKeyCanMatchPartitionKey("ses_target", null), true);
|
assert.equal(kafkaMessageKeyCanMatchPartitionKey("ses_target", null), true);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("Agent observer retains a bounded latest window while scanning through the Kafka barrier", async () => {
|
||||||
|
const messages = Array.from({ length: 5 }, (_, offset) => ({
|
||||||
|
offset: String(offset),
|
||||||
|
key: Buffer.from(`run-${offset}`),
|
||||||
|
value: Buffer.from(JSON.stringify({ schema: "hwlab.event.v1", eventId: `event-${offset}` }))
|
||||||
|
}));
|
||||||
|
const queried = await queryKafkaEventStream({
|
||||||
|
env: { HWLAB_KAFKA_BOOTSTRAP_SERVERS: "kafka.test:9092", HWLAB_KAFKA_CLIENT_ID: "hwlab-test" },
|
||||||
|
topic: "hwlab.event.v1",
|
||||||
|
limit: 2,
|
||||||
|
scanLimit: 10,
|
||||||
|
retainLatestMatches: true,
|
||||||
|
groupIdPrefix: "hwlab-agent-observer-test",
|
||||||
|
kafkaFactory: kafkaQueryTestFactory(messages)
|
||||||
|
});
|
||||||
|
|
||||||
|
assert.equal(queried.completionReason, "end-offset");
|
||||||
|
assert.equal(queried.completion.complete, true);
|
||||||
|
assert.equal(queried.scannedCount, 5);
|
||||||
|
assert.equal(queried.matchedCount, 2);
|
||||||
|
assert.equal(queried.totalMatchedCount, 5);
|
||||||
|
assert.equal(queried.droppedMatchedCount, 3);
|
||||||
|
assert.equal(queried.retainLatestMatches, true);
|
||||||
|
assert.deepEqual(queried.events.map((event) => event.offset), ["3", "4"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("Workbench Kafka query keeps the existing stop-at-match-limit behavior", async () => {
|
||||||
|
const messages = Array.from({ length: 5 }, (_, offset) => ({
|
||||||
|
offset: String(offset),
|
||||||
|
value: Buffer.from(JSON.stringify({ schema: "hwlab.event.v1", eventId: `event-${offset}` }))
|
||||||
|
}));
|
||||||
|
const queried = await queryKafkaEventStream({
|
||||||
|
env: { HWLAB_KAFKA_BOOTSTRAP_SERVERS: "kafka.test:9092", HWLAB_KAFKA_CLIENT_ID: "hwlab-test" },
|
||||||
|
topic: "hwlab.event.v1",
|
||||||
|
limit: 2,
|
||||||
|
scanLimit: 10,
|
||||||
|
groupIdPrefix: "hwlab-workbench-query-test",
|
||||||
|
kafkaFactory: kafkaQueryTestFactory(messages)
|
||||||
|
});
|
||||||
|
|
||||||
|
assert.equal(queried.completionReason, "limit");
|
||||||
|
assert.equal(queried.completion.complete, false);
|
||||||
|
assert.equal(queried.scannedCount, 2);
|
||||||
|
assert.equal(queried.totalMatchedCount, 2);
|
||||||
|
assert.equal(queried.droppedMatchedCount, 0);
|
||||||
|
assert.equal(queried.retainLatestMatches, false);
|
||||||
|
assert.deepEqual(queried.events.map((event) => event.offset), ["0", "1"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
function kafkaQueryTestFactory(messages) {
|
||||||
|
return () => ({
|
||||||
|
admin() {
|
||||||
|
return {
|
||||||
|
async connect() {},
|
||||||
|
async fetchTopicOffsets() { return [{ partition: 0, low: "0", high: String(messages.length) }]; },
|
||||||
|
async disconnect() {}
|
||||||
|
};
|
||||||
|
},
|
||||||
|
consumer() {
|
||||||
|
return {
|
||||||
|
async connect() {},
|
||||||
|
async subscribe() {},
|
||||||
|
async run({ eachBatch }) {
|
||||||
|
await eachBatch({ batch: { topic: "hwlab.event.v1", partition: 0, messages } });
|
||||||
|
},
|
||||||
|
async stop() {},
|
||||||
|
async disconnect() {}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
test("Workbench session index uses explicit budgets and defaults missing non-core values with warnings", () => {
|
test("Workbench session index uses explicit budgets and defaults missing non-core values with warnings", () => {
|
||||||
const explicit = kafkaEventBridgeConfig({
|
const explicit = kafkaEventBridgeConfig({
|
||||||
...LIVE_REFRESH_ENV,
|
...LIVE_REFRESH_ENV,
|
||||||
|
|||||||
@@ -1021,7 +1021,7 @@ function requireKafkaProjectorStore(runtimeStore, capabilities = {}) {
|
|||||||
if (missing.length > 0) throw contractError("hwlab_kafka_projector_store_invalid", `Kafka durable capabilities require a runtime store: ${missing.join(", ")}`);
|
if (missing.length > 0) throw contractError("hwlab_kafka_projector_store_invalid", `Kafka durable capabilities require a runtime store: ${missing.join(", ")}`);
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function queryKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, traceId = null, sessionId = null, runId = null, commandId = null, partitionKey = null, limit = DEFAULT_QUERY_LIMIT, scanLimit = null, timeoutMs = DEFAULT_QUERY_TIMEOUT_MS, fromBeginning = true, groupIdPrefix = null, signal = null, kafkaFactory = defaultKafkaFactory } = {}) {
|
export async function queryKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, traceId = null, sessionId = null, runId = null, commandId = null, partitionKey = null, limit = DEFAULT_QUERY_LIMIT, scanLimit = null, timeoutMs = DEFAULT_QUERY_TIMEOUT_MS, fromBeginning = true, groupIdPrefix = null, retainLatestMatches = false, signal = null, kafkaFactory = defaultKafkaFactory } = {}) {
|
||||||
const queryStartedAtMs = Date.now();
|
const queryStartedAtMs = Date.now();
|
||||||
const timing = {
|
const timing = {
|
||||||
endOffsetSnapshotMs: 0,
|
endOffsetSnapshotMs: 0,
|
||||||
@@ -1036,6 +1036,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
const clientId = stringValue(env.HWLAB_KAFKA_CLIENT_ID) || DEFAULT_CLIENT_ID;
|
const clientId = stringValue(env.HWLAB_KAFKA_CLIENT_ID) || DEFAULT_CLIENT_ID;
|
||||||
const resolvedTopic = stringValue(topic) || kafkaTopicForStream(stream, env);
|
const resolvedTopic = stringValue(topic) || kafkaTopicForStream(stream, env);
|
||||||
const maxEvents = Math.max(1, integerValue(limit) || DEFAULT_QUERY_LIMIT);
|
const maxEvents = Math.max(1, integerValue(limit) || DEFAULT_QUERY_LIMIT);
|
||||||
|
const keepLatestMatches = retainLatestMatches === true;
|
||||||
const resolvedPartitionKey = stringValue(partitionKey);
|
const resolvedPartitionKey = stringValue(partitionKey);
|
||||||
const maxScannedRecords = scanLimit === null || scanLimit === undefined
|
const maxScannedRecords = scanLimit === null || scanLimit === undefined
|
||||||
? null
|
? null
|
||||||
@@ -1067,6 +1068,9 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
maxWaitTimeInMs: 100
|
maxWaitTimeInMs: 100
|
||||||
});
|
});
|
||||||
const events = [];
|
const events = [];
|
||||||
|
let retainedEventCursor = 0;
|
||||||
|
let totalMatchedCount = 0;
|
||||||
|
let droppedMatchedCount = 0;
|
||||||
const firstScannedOffsetByPartition = new Map();
|
const firstScannedOffsetByPartition = new Map();
|
||||||
const lastScannedOffsetByPartition = new Map();
|
const lastScannedOffsetByPartition = new Map();
|
||||||
let scannedCount = 0;
|
let scannedCount = 0;
|
||||||
@@ -1167,7 +1171,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
parsedCount += 1;
|
parsedCount += 1;
|
||||||
if (!eventMatchesFilters(value, { traceId, sessionId, runId, commandId })) filterRejectedCount += 1;
|
if (!eventMatchesFilters(value, { traceId, sessionId, runId, commandId })) filterRejectedCount += 1;
|
||||||
else {
|
else {
|
||||||
events.push({
|
retainMatchedEvent({
|
||||||
topic: messageTopic,
|
topic: messageTopic,
|
||||||
partition,
|
partition,
|
||||||
offset: message.offset,
|
offset: message.offset,
|
||||||
@@ -1184,7 +1188,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
finish(verifiedStartPartitions.size === endOffsetByPartition.size ? "end-offset" : "retention-start-unverified");
|
finish(verifiedStartPartitions.size === endOffsetByPartition.size ? "end-offset" : "retention-start-unverified");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (events.length >= maxEvents) {
|
if (!keepLatestMatches && events.length >= maxEvents) {
|
||||||
finish("limit");
|
finish("limit");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
@@ -1210,18 +1214,37 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
}
|
}
|
||||||
return kafkaQueryResult();
|
return kafkaQueryResult();
|
||||||
|
|
||||||
|
function retainMatchedEvent(event) {
|
||||||
|
totalMatchedCount += 1;
|
||||||
|
if (!keepLatestMatches || events.length < maxEvents) {
|
||||||
|
events.push(event);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
events[retainedEventCursor] = event;
|
||||||
|
retainedEventCursor = (retainedEventCursor + 1) % maxEvents;
|
||||||
|
droppedMatchedCount += 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
function retainedEvents() {
|
||||||
|
if (!keepLatestMatches || droppedMatchedCount === 0 || retainedEventCursor === 0) return events;
|
||||||
|
return [...events.slice(retainedEventCursor), ...events.slice(0, retainedEventCursor)];
|
||||||
|
}
|
||||||
|
|
||||||
function kafkaQueryResult() {
|
function kafkaQueryResult() {
|
||||||
const retentionStartVerified = endOffsetSnapshot.available && verifiedStartPartitions.size === endOffsetByPartition.size;
|
const retentionStartVerified = endOffsetSnapshot.available && verifiedStartPartitions.size === endOffsetByPartition.size;
|
||||||
const complete = completionReason === "end-offset" && retentionStartVerified;
|
const complete = completionReason === "end-offset" && retentionStartVerified;
|
||||||
|
const resultEvents = retainedEvents();
|
||||||
return {
|
return {
|
||||||
ok: true,
|
ok: true,
|
||||||
stream,
|
stream,
|
||||||
topic: resolvedTopic,
|
topic: resolvedTopic,
|
||||||
groupId,
|
groupId,
|
||||||
count: events.length,
|
count: resultEvents.length,
|
||||||
scannedCount,
|
scannedCount,
|
||||||
parsedCount,
|
parsedCount,
|
||||||
matchedCount: events.length,
|
matchedCount: resultEvents.length,
|
||||||
|
totalMatchedCount,
|
||||||
|
droppedMatchedCount,
|
||||||
invalidJsonCount,
|
invalidJsonCount,
|
||||||
filterRejectedCount,
|
filterRejectedCount,
|
||||||
keyRejectedCount,
|
keyRejectedCount,
|
||||||
@@ -1248,6 +1271,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
limit: maxEvents,
|
limit: maxEvents,
|
||||||
scanLimit: maxScannedRecords,
|
scanLimit: maxScannedRecords,
|
||||||
timeoutMs: budgetMs,
|
timeoutMs: budgetMs,
|
||||||
|
retainLatestMatches: keepLatestMatches,
|
||||||
filters: compactObject({ traceId, sessionId, runId, commandId }),
|
filters: compactObject({ traceId, sessionId, runId, commandId }),
|
||||||
partitionKeyScoped: targetPartition !== null,
|
partitionKeyScoped: targetPartition !== null,
|
||||||
targetPartition,
|
targetPartition,
|
||||||
@@ -1256,7 +1280,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab
|
|||||||
totalMs: Date.now() - queryStartedAtMs,
|
totalMs: Date.now() - queryStartedAtMs,
|
||||||
valuesPrinted: false
|
valuesPrinted: false
|
||||||
},
|
},
|
||||||
events,
|
events: resultEvents,
|
||||||
valuesPrinted: false
|
valuesPrinted: false
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -87,6 +87,7 @@ export async function handleAgentObserverHttp(request, response, options = {}) {
|
|||||||
timeoutMs: policy.retentionTimeoutMs,
|
timeoutMs: policy.retentionTimeoutMs,
|
||||||
groupIdPrefix: policy.retentionGroupPrefix,
|
groupIdPrefix: policy.retentionGroupPrefix,
|
||||||
fromBeginning: true,
|
fromBeginning: true,
|
||||||
|
retainLatestMatches: true,
|
||||||
signal
|
signal
|
||||||
}),
|
}),
|
||||||
deliverEvent: (envelope) => enqueue("hwlab.event.v1", envelope),
|
deliverEvent: (envelope) => enqueue("hwlab.event.v1", envelope),
|
||||||
|
|||||||
@@ -406,6 +406,8 @@ function boundedQueryDiagnostics(result) {
|
|||||||
scannedCount: numericValue(result?.scannedCount),
|
scannedCount: numericValue(result?.scannedCount),
|
||||||
parsedCount: numericValue(result?.parsedCount),
|
parsedCount: numericValue(result?.parsedCount),
|
||||||
matchedCount: numericValue(result?.matchedCount),
|
matchedCount: numericValue(result?.matchedCount),
|
||||||
|
totalMatchedCount: numericValue(result?.totalMatchedCount),
|
||||||
|
droppedMatchedCount: numericValue(result?.droppedMatchedCount),
|
||||||
filterRejectedCount: numericValue(result?.filterRejectedCount),
|
filterRejectedCount: numericValue(result?.filterRejectedCount),
|
||||||
keyRejectedCount: numericValue(result?.keyRejectedCount),
|
keyRejectedCount: numericValue(result?.keyRejectedCount),
|
||||||
partitionKeyScoped: result?.partitionKeyScoped === true,
|
partitionKeyScoped: result?.partitionKeyScoped === true,
|
||||||
|
|||||||
Reference in New Issue
Block a user