477 lines
22 KiB
TypeScript
477 lines
22 KiB
TypeScript
/*
|
|
* SPEC: PJ2026-0104010803 Workbench事件流可见性 draft-2026-07-09-p0-kafka-authority.
|
|
* Responsibility: debug-only SSE passthrough for HWLAB Kafka event envelopes.
|
|
*/
|
|
import { sendJson } from "./server-http-utils.ts";
|
|
import { openKafkaEventStream } from "./kafka-event-bridge.ts";
|
|
import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts";
|
|
import { authenticateWorkbenchRead } from "./server-workbench-read-http.ts";
|
|
import { workbenchKafkaDebugCapability } from "./workbench-kafka-debug-capability.ts";
|
|
import {
|
|
completeWorkbenchKafkaDebugReplay,
|
|
createWorkbenchKafkaDebugReplayTracker,
|
|
markWorkbenchKafkaDebugReplayBarrierComplete,
|
|
markWorkbenchKafkaDebugReplayConsumerReady,
|
|
normalizeWorkbenchKafkaDebugOffsetRange,
|
|
observeWorkbenchKafkaDebugReplayRecord,
|
|
recordWorkbenchKafkaDebugReplayDelivery,
|
|
workbenchKafkaDebugOffsetRangeWidth
|
|
} from "./workbench-kafka-debug-replay-contract.ts";
|
|
|
|
const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1";
|
|
const DEFAULT_STREAM = "hwlab";
|
|
const ALLOWED_STREAMS = new Set(["stdio", "agentrun", "hwlab", "hwlab-debug"]);
|
|
|
|
export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, options = {}) {
|
|
const route = routeSuffix(url.pathname);
|
|
try {
|
|
const auth = await authenticateWorkbenchRead(request, response, options);
|
|
if (!auth) return;
|
|
if (auth.actor?.role !== "admin") {
|
|
return sendJson(response, 403, debugError("admin_required", "Only admin users can access raw Workbench Kafka debug streams."));
|
|
}
|
|
const stream = streamFromUrl(url);
|
|
const isolatedDebug = stream === "hwlab-debug" ? isolatedDebugConfig(response, url, options.env ?? process.env, { requireReplay: route === "/events" }) : null;
|
|
if (isolatedDebug === false) return;
|
|
if (route === "") {
|
|
if (request.method !== "GET") return methodNotAllowed(response, "GET");
|
|
return sendJson(response, 200, describeKafkaSseDebug(url, options.env ?? process.env, isolatedDebug));
|
|
}
|
|
if (route === "/events") {
|
|
if (request.method !== "GET") return methodNotAllowed(response, "GET");
|
|
return openKafkaDebugSse(request, response, url, options, isolatedDebug);
|
|
}
|
|
return sendJson(response, 404, debugError("workbench_debug_kafka_sse_route_not_found", "Workbench debug Kafka SSE route is not implemented.", { route }));
|
|
} catch (error) {
|
|
options.logger?.warn?.({
|
|
event: "workbench_debug_kafka_sse_failed",
|
|
route,
|
|
errorName: error?.name ?? "Error",
|
|
message: error instanceof Error ? error.message : String(error ?? "unknown"),
|
|
valuesRedacted: true
|
|
});
|
|
if (response.headersSent) return writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error) });
|
|
return sendJson(response, 500, debugError("workbench_debug_kafka_sse_failed", "Workbench debug Kafka SSE request failed."));
|
|
}
|
|
}
|
|
|
|
function describeKafkaSseDebug(url, env, isolatedDebug = null) {
|
|
const stream = streamFromUrl(url);
|
|
const requestedFilters = filtersFromUrl(url);
|
|
const resolvedFilters = stream === "hwlab-debug" ? compactObject({ traceId: requestedFilters.traceId, replayId: requestedFilters.replayId }) : resolveKafkaDebugFilters(requestedFilters, { traceStore: defaultCodeAgentTraceStore });
|
|
return {
|
|
ok: true,
|
|
contractVersion: CONTRACT_VERSION,
|
|
stream,
|
|
topic: isolatedDebug?.topic ?? topicForStream(stream, env),
|
|
debugIsolation: stream === "hwlab-debug",
|
|
liveOnly: false,
|
|
replay: stream === "hwlab-debug",
|
|
replayLimit: isolatedDebug?.replayLimit ?? null,
|
|
replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null,
|
|
groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null,
|
|
replayRequest: isolatedDebug?.replayRequest ?? null,
|
|
filters: requestedFilters,
|
|
resolvedFilters,
|
|
eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`,
|
|
valuesPrinted: false
|
|
};
|
|
}
|
|
|
|
async function openKafkaDebugSse(request, response, url, options, isolatedDebug = null) {
|
|
const env = options.env ?? process.env;
|
|
const stream = streamFromUrl(url);
|
|
const filters = filtersFromUrl(url);
|
|
const isolatedReplay = stream === "hwlab-debug";
|
|
const resolvedFilters = isolatedReplay ? compactObject({ traceId: filters.traceId, replayId: filters.replayId }) : resolveKafkaDebugFilters(filters, options);
|
|
const replayRequest = isolatedDebug?.replayRequest ?? replayRequestFromUrl(url);
|
|
const correlatedReplay = isolatedReplay && replayRequest.mode === "correlated-v2";
|
|
const replayTracker = isolatedReplay ? createWorkbenchKafkaDebugReplayTracker({
|
|
mode: replayRequest.mode,
|
|
replayId: resolvedFilters.replayId,
|
|
traceId: resolvedFilters.traceId,
|
|
topic: isolatedDebug.topic,
|
|
offsetRange: replayRequest.offsetRange,
|
|
producer: replayRequest.producer
|
|
}) : null;
|
|
response.writeHead(200, {
|
|
"content-type": "text/event-stream; charset=utf-8",
|
|
"cache-control": "no-store, no-transform",
|
|
connection: "keep-alive",
|
|
"x-accel-buffering": "no",
|
|
"x-content-type-options": "nosniff"
|
|
});
|
|
if (typeof response.flushHeaders === "function") response.flushHeaders();
|
|
let kafkaStream = null;
|
|
let closed = false;
|
|
let consumerReady = false;
|
|
let replayFinished = false;
|
|
let replayTimer = null;
|
|
let replayCount = 0;
|
|
let replayTerminalObserved = false;
|
|
let barrierCompletionPending = false;
|
|
const bufferedRecords = [];
|
|
const close = () => {
|
|
if (closed) return;
|
|
closed = true;
|
|
if (replayTimer) clearTimeout(replayTimer);
|
|
void kafkaStream?.stop?.();
|
|
};
|
|
response.once("close", close);
|
|
request.once?.("aborted", close);
|
|
request.socket?.once?.("close", close);
|
|
const finishReplay = (reason) => {
|
|
if (!isolatedReplay || replayFinished || closed) return;
|
|
replayFinished = true;
|
|
if (replayTimer) clearTimeout(replayTimer);
|
|
const replayResult = completeWorkbenchKafkaDebugReplay(replayTracker, {
|
|
reason,
|
|
terminalObserved: replayTerminalObserved
|
|
});
|
|
writeSse(response, "hwlab.kafka.replay-complete", {
|
|
ok: true,
|
|
contractVersion: CONTRACT_VERSION,
|
|
stream,
|
|
topic: isolatedDebug.topic,
|
|
traceId: resolvedFilters.traceId,
|
|
replayComplete: true,
|
|
replayMode: replayResult.mode,
|
|
replayId: replayResult.replayId,
|
|
phase: replayResult.phase,
|
|
code: replayResult.code,
|
|
result: replayResult,
|
|
reason,
|
|
count: replayCount,
|
|
counts: replayResult.counts,
|
|
rejectedByReason: replayResult.rejectedByReason,
|
|
offsetRange: replayResult.offsetRange,
|
|
sourceLineage: replayResult.sourceLineage,
|
|
limit: isolatedDebug.replayLimit,
|
|
timeoutMs: isolatedDebug.replayTimeoutMs,
|
|
terminalObserved: replayTerminalObserved,
|
|
barrierComplete: replayResult.barrier.completed,
|
|
stopBoundary: "sse-end",
|
|
clientCountsAvailable: false,
|
|
valuesPrinted: false
|
|
});
|
|
if (!response.writableEnded) response.end();
|
|
};
|
|
const writeRecord = (record) => {
|
|
if (closed || replayFinished) return;
|
|
const delivered = writeSse(response, "hwlab.kafka.event", {
|
|
ok: true,
|
|
contractVersion: CONTRACT_VERSION,
|
|
stream,
|
|
topic: record.topic,
|
|
partition: record.partition,
|
|
offset: record.offset,
|
|
key: record.key,
|
|
timestamp: record.timestamp,
|
|
value: record.value,
|
|
serverSentAt: new Date().toISOString(),
|
|
valuesPrinted: false
|
|
});
|
|
if (!isolatedReplay) return;
|
|
recordWorkbenchKafkaDebugReplayDelivery(replayTracker, record, delivered);
|
|
replayCount = replayTracker.counts.delivered;
|
|
replayTerminalObserved ||= debugRecordIsTerminal(record);
|
|
if (!correlatedReplay && replayTerminalObserved) finishReplay("terminal");
|
|
else if (!correlatedReplay && replayCount >= isolatedDebug.replayLimit) finishReplay("limit");
|
|
};
|
|
if (isolatedReplay) {
|
|
replayTimer = setTimeout(() => finishReplay("timeout"), isolatedDebug.replayTimeoutMs);
|
|
replayTimer.unref?.();
|
|
}
|
|
try {
|
|
kafkaStream = await openKafkaEventStream({
|
|
env,
|
|
stream,
|
|
topic: isolatedDebug?.topic ?? null,
|
|
groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null,
|
|
fromBeginning: isolatedReplay ? !correlatedReplay : url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true",
|
|
...resolvedFilters,
|
|
offsetRange: replayRequest.offsetRange,
|
|
kafkaFactory: options.kafkaFactory,
|
|
onRecord: isolatedReplay ? async (observation) => {
|
|
observeWorkbenchKafkaDebugReplayRecord(replayTracker, observation);
|
|
} : null,
|
|
onBarrierComplete: correlatedReplay ? async () => {
|
|
markWorkbenchKafkaDebugReplayBarrierComplete(replayTracker);
|
|
if (consumerReady) finishReplay("barrier");
|
|
else barrierCompletionPending = true;
|
|
} : null,
|
|
onEvent: async (record) => {
|
|
if (isolatedReplay && !consumerReady) {
|
|
if (bufferedRecords.length < isolatedDebug.replayLimit) bufferedRecords.push(record);
|
|
return;
|
|
}
|
|
writeRecord(record);
|
|
}
|
|
});
|
|
if (closed) {
|
|
await kafkaStream.stop?.();
|
|
return;
|
|
}
|
|
if (isolatedReplay) markWorkbenchKafkaDebugReplayConsumerReady(replayTracker, { groupId: kafkaStream.groupId, seekApplied: kafkaStream.seekApplied });
|
|
writeSse(response, "hwlab.kafka.connected", {
|
|
ok: true,
|
|
contractVersion: CONTRACT_VERSION,
|
|
consumerReady: true,
|
|
groupId: kafkaStream.groupId,
|
|
stream,
|
|
topic: kafkaStream.topic ?? isolatedDebug?.topic ?? topicForStream(stream, env),
|
|
debugIsolation: isolatedReplay,
|
|
deliverySemantics: isolatedReplay ? "debug-replay" : "diagnostic",
|
|
liveOnly: false,
|
|
replay: isolatedReplay,
|
|
replayLimit: isolatedDebug?.replayLimit ?? null,
|
|
replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null,
|
|
replayId: replayTracker?.replayId ?? null,
|
|
replayMode: replayTracker?.mode ?? null,
|
|
offsetRange: replayTracker?.requestedOffsetRange ?? null,
|
|
replayContractVersion: replayTracker?.mode === "correlated-v2" ? "workbench-kafka-debug-replay-v2" : CONTRACT_VERSION,
|
|
seekApplied: kafkaStream.seekApplied === true,
|
|
filters,
|
|
resolvedFilters,
|
|
serverSentAt: new Date().toISOString(),
|
|
valuesPrinted: false
|
|
});
|
|
consumerReady = true;
|
|
for (const record of bufferedRecords.splice(0)) {
|
|
if (replayFinished || closed) break;
|
|
writeRecord(record);
|
|
}
|
|
if (barrierCompletionPending && !replayFinished && !closed) finishReplay("barrier");
|
|
} catch (error) {
|
|
if (replayTimer) clearTimeout(replayTimer);
|
|
writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error), valuesPrinted: false });
|
|
if (isolatedReplay) finishReplay("error");
|
|
}
|
|
}
|
|
|
|
function resolveKafkaDebugFilters(filters = {}, options = {}) {
|
|
const traceId = textValue(filters.traceId);
|
|
if (!traceId) return filters;
|
|
const traceStore = options.traceStore ?? defaultCodeAgentTraceStore;
|
|
const snapshot = typeof traceStore?.snapshot === "function" ? traceStore.snapshot(traceId) : null;
|
|
const resolved = collectTraceLinkedIds(snapshot);
|
|
const hasResolvedKafkaKey = Boolean(resolved.runId || resolved.commandId || resolved.sessionId);
|
|
if (!hasResolvedKafkaKey) return filters;
|
|
const resolvedSessionId = resolved.runId || resolved.commandId ? null : resolved.sessionId;
|
|
return compactObject({
|
|
sessionId: filters.sessionId || resolvedSessionId,
|
|
runId: filters.runId || resolved.runId,
|
|
commandId: filters.commandId || resolved.commandId
|
|
});
|
|
}
|
|
|
|
function collectTraceLinkedIds(snapshot) {
|
|
const out = { sessionId: null, runId: null, commandId: null };
|
|
const events = Array.isArray(snapshot?.events) ? snapshot.events : [];
|
|
for (const event of events) collectIdsFromRecord(event, out);
|
|
collectIdsFromRecord(snapshot?.lastEvent, out);
|
|
return out;
|
|
}
|
|
|
|
function collectIdsFromRecord(record, out) {
|
|
const value = record && typeof record === "object" && !Array.isArray(record) ? record : {};
|
|
const agentRun = value.agentRun && typeof value.agentRun === "object" && !Array.isArray(value.agentRun) ? value.agentRun : {};
|
|
const payload = value.payload && typeof value.payload === "object" && !Array.isArray(value.payload) ? value.payload : {};
|
|
out.sessionId ||= textValue(value.sessionId ?? value.sourceSessionId ?? agentRun.sessionId ?? payload.sessionId);
|
|
out.runId ||= textValue(value.runId ?? value.sourceRunId ?? agentRun.runId ?? payload.runId);
|
|
out.commandId ||= textValue(value.commandId ?? value.sourceCommandId ?? agentRun.commandId ?? payload.commandId);
|
|
}
|
|
|
|
function writeSse(response, eventName, payload) {
|
|
if (response.destroyed || response.writableEnded) return false;
|
|
try {
|
|
response.write(`event: ${eventName}\n`);
|
|
const id = sseEventId(payload);
|
|
if (id) response.write(`id: ${id}\n`);
|
|
response.write(`data: ${JSON.stringify(payload)}\n\n`);
|
|
return true;
|
|
} catch {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
function routeSuffix(pathname) {
|
|
return String(pathname || "").replace(/^\/v1\/workbench\/debug\/kafka-sse/u, "");
|
|
}
|
|
|
|
function streamFromUrl(url) {
|
|
const stream = String(url.searchParams.get("stream") || DEFAULT_STREAM).trim();
|
|
return ALLOWED_STREAMS.has(stream) ? stream : DEFAULT_STREAM;
|
|
}
|
|
|
|
function filtersFromUrl(url) {
|
|
return compactObject({
|
|
traceId: safeId(url.searchParams.get("traceId") || url.searchParams.get("trace-id")),
|
|
sessionId: safeId(url.searchParams.get("sessionId") || url.searchParams.get("session-id")),
|
|
runId: safeId(url.searchParams.get("runId") || url.searchParams.get("run-id")),
|
|
commandId: safeId(url.searchParams.get("commandId") || url.searchParams.get("command-id")),
|
|
replayId: safeReplayId(url.searchParams.get("replayId") || url.searchParams.get("replay-id"))
|
|
});
|
|
}
|
|
|
|
function topicForStream(stream, env) {
|
|
if (stream === "stdio") return textValue(env.HWLAB_KAFKA_STDIO_TOPIC ?? env.AGENTRUN_KAFKA_STDIO_TOPIC) || "codex-stdio.raw.v1";
|
|
if (stream === "agentrun") return textValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC) || "agentrun.event.v1";
|
|
if (stream === "hwlab-debug") return workbenchKafkaDebugCapability(env).topic;
|
|
return textValue(env.HWLAB_KAFKA_EVENT_TOPIC) || "hwlab.event.v1";
|
|
}
|
|
|
|
function isolatedDebugConfig(response, url, env, options = {}) {
|
|
let config;
|
|
try {
|
|
config = workbenchKafkaDebugCapability(env);
|
|
} catch (error) {
|
|
sendJson(response, Number(error?.statusCode) || 503, debugError(error?.code || "workbench_kafka_debug_capability_invalid", error instanceof Error ? error.message : "Workbench isolated Kafka debug capability is invalid."));
|
|
return false;
|
|
}
|
|
if (!config.enabled) {
|
|
sendJson(response, 503, debugError("workbench_kafka_debug_disabled", "Workbench isolated Kafka debug capability is disabled."));
|
|
return false;
|
|
}
|
|
if (options.requireReplay && !safeId(url.searchParams.get("traceId") || url.searchParams.get("trace-id"))) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_trace_required", "Isolated Workbench Kafka replay requires traceId."));
|
|
return false;
|
|
}
|
|
if (options.requireReplay && !["1", "true"].includes(String(url.searchParams.get("fromBeginning") || "").toLowerCase())) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_replay_required", "Isolated Workbench Kafka debug requires explicit fromBeginning=true."));
|
|
return false;
|
|
}
|
|
const rawReplayId = url.searchParams.get("replayId") || url.searchParams.get("replay-id");
|
|
if (textValue(rawReplayId) && !safeReplayId(rawReplayId)) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_replay_id_invalid", "Replay correlation requires replayId to start with rpl_ and contain only safe identifier characters."));
|
|
return false;
|
|
}
|
|
const replayRequest = replayRequestFromUrl(url);
|
|
if (options.requireReplay && replayRequest.offsetRangeInvalid) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_offset_range_invalid", "Replay offset barrier requires partition, firstOffset, and lastOffset with firstOffset <= lastOffset."));
|
|
return false;
|
|
}
|
|
if (options.requireReplay && replayRequest.mode === "correlated-v2") {
|
|
if (!replayRequest.replayId || !replayRequest.offsetRange) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_v2_correlation_incomplete", "Correlated V2 replay requires replayId and partition/firstOffset/lastOffset."));
|
|
return false;
|
|
}
|
|
if (replayRequest.producerMetadataInvalid) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_v2_producer_metadata_invalid", "Optional producerInvoked, sourceMatched, and publishedCount must use valid boolean/integer values."));
|
|
return false;
|
|
}
|
|
if (replayRequest.producer.invoked === false) {
|
|
sendJson(response, 409, debugError("workbench_kafka_debug_producer_not_invoked", "Correlated V2 replay cannot use producerInvoked=false."));
|
|
return false;
|
|
}
|
|
if (replayRequest.producer.sourceMatched === false) {
|
|
sendJson(response, 409, debugError("workbench_kafka_debug_source_trace_missing", "Correlated V2 replay cannot start because the producer did not match the requested source trace."));
|
|
return false;
|
|
}
|
|
const rangeWidth = workbenchKafkaDebugOffsetRangeWidth(replayRequest.offsetRange);
|
|
if (!rangeWidth || (replayRequest.producer.publishedCount !== null && replayRequest.producer.publishedCount !== rangeWidth)) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_producer_cardinality_mismatch", "publishedCount must equal the requested offset range width."));
|
|
return false;
|
|
}
|
|
if (rangeWidth > config.replayLimit) {
|
|
sendJson(response, 400, debugError("workbench_kafka_debug_replay_limit_exceeded", "The correlated replay batch exceeds the YAML-owned replay limit."));
|
|
return false;
|
|
}
|
|
}
|
|
return { ...config, replayRequest };
|
|
}
|
|
|
|
function replayRequestFromUrl(url) {
|
|
const rawReplayId = url.searchParams.get("replayId") || url.searchParams.get("replay-id");
|
|
const rawProducerInvoked = url.searchParams.get("producerInvoked") || url.searchParams.get("producer-invoked");
|
|
const rawSourceMatched = url.searchParams.get("sourceMatched") || url.searchParams.get("source-matched");
|
|
const rawPublishedCount = url.searchParams.get("publishedCount") || url.searchParams.get("published-count");
|
|
const rawRange = {
|
|
partition: url.searchParams.get("partition"),
|
|
firstOffset: url.searchParams.get("firstOffset") || url.searchParams.get("first-offset") || url.searchParams.get("startOffset") || url.searchParams.get("start-offset"),
|
|
lastOffset: url.searchParams.get("lastOffset") || url.searchParams.get("last-offset") || url.searchParams.get("endOffset") || url.searchParams.get("end-offset")
|
|
};
|
|
const offsetRangeProvided = Object.values(rawRange).some((entry) => textValue(entry));
|
|
const offsetRange = normalizeWorkbenchKafkaDebugOffsetRange(rawRange);
|
|
const correlationProvided = Boolean(textValue(rawReplayId) || offsetRangeProvided || textValue(rawProducerInvoked) || textValue(rawSourceMatched) || textValue(rawPublishedCount));
|
|
const producerInvoked = optionalBoolean(rawProducerInvoked);
|
|
const sourceMatched = optionalBoolean(rawSourceMatched);
|
|
const publishedCount = optionalNonNegativeInteger(rawPublishedCount);
|
|
return {
|
|
mode: correlationProvided ? "correlated-v2" : "trace-only-v1",
|
|
replayId: safeReplayId(rawReplayId),
|
|
offsetRange,
|
|
offsetRangeInvalid: offsetRangeProvided && !offsetRange,
|
|
producerMetadataInvalid: Boolean(
|
|
(textValue(rawProducerInvoked) && producerInvoked === null)
|
|
|| (textValue(rawSourceMatched) && sourceMatched === null)
|
|
|| (textValue(rawPublishedCount) && publishedCount === null)
|
|
),
|
|
producer: {
|
|
invoked: producerInvoked,
|
|
sourceMatched,
|
|
publishedCount
|
|
},
|
|
valuesPrinted: false
|
|
};
|
|
}
|
|
|
|
function methodNotAllowed(response, allow) {
|
|
response.setHeader("allow", allow);
|
|
return sendJson(response, 405, debugError("method_not_allowed", `Method not allowed; expected ${allow}.`));
|
|
}
|
|
|
|
function debugError(code, message, extra = {}) {
|
|
return { ok: false, error: { code, message, ...extra }, valuesPrinted: false };
|
|
}
|
|
|
|
function errorMessagePayload(error) {
|
|
return { code: "kafka_sse_error", message: error instanceof Error ? error.message : String(error ?? "unknown") };
|
|
}
|
|
|
|
function debugRecordIsTerminal(record) {
|
|
const value = record?.value && typeof record.value === "object" && !Array.isArray(record.value) ? record.value : {};
|
|
const event = value.event && typeof value.event === "object" && !Array.isArray(value.event) ? value.event : {};
|
|
const eventType = textValue(value.eventType ?? event.eventType ?? event.type);
|
|
return event.terminal === true || eventType === "terminal" || eventType === "result";
|
|
}
|
|
|
|
function safeId(value) {
|
|
const text = textValue(value);
|
|
return text && /^[A-Za-z0-9_.:-]{3,220}$/u.test(text) ? text : null;
|
|
}
|
|
|
|
function safeReplayId(value) {
|
|
const text = safeId(value);
|
|
return text && /^rpl_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null;
|
|
}
|
|
|
|
function optionalBoolean(value) {
|
|
const text = textValue(value)?.toLowerCase();
|
|
if (text === "true" || text === "1") return true;
|
|
if (text === "false" || text === "0") return false;
|
|
return null;
|
|
}
|
|
|
|
function optionalNonNegativeInteger(value) {
|
|
if (!textValue(value)) return null;
|
|
const parsed = Number(value);
|
|
return Number.isSafeInteger(parsed) && parsed >= 0 ? parsed : null;
|
|
}
|
|
|
|
function textValue(value) {
|
|
const text = typeof value === "string" ? value.trim() : value === null || value === undefined ? "" : String(value).trim();
|
|
return text.length > 0 ? text : null;
|
|
}
|
|
|
|
function compactObject(value) {
|
|
return Object.fromEntries(Object.entries(value).filter(([, entry]) => textValue(entry)));
|
|
}
|
|
|
|
function sseEventId(payload) {
|
|
const topic = textValue(payload?.topic);
|
|
const partition = textValue(payload?.partition);
|
|
const offset = textValue(payload?.offset);
|
|
return topic && partition && offset ? `${topic}:${partition}:${offset}` : null;
|
|
}
|