fix: seal AgentRun terminal projection (#1820)
This commit is contained in:
@@ -2217,17 +2217,34 @@ async function fetchAgentRunEventsForTrace({ fetchImpl, managerUrl, timeoutMs, e
|
||||
const limit = parsePositiveInteger(mapping.eventsPageLimit ?? env?.HWLAB_CODE_AGENT_AGENTRUN_EVENTS_PAGE_LIMIT, 500);
|
||||
const path = `/api/v1/runs/${encodeURIComponent(runId)}/events?afterSeq=${encodeURIComponent(String(afterSeq))}&limit=${encodeURIComponent(String(limit))}`;
|
||||
const response = await agentRunJson(fetchImpl, managerUrl, path, { method: "GET", timeoutMs, env });
|
||||
void emitCodeAgentOtelSpan("projection_sync", mapping.traceId ?? mapping.hwlabTraceId ?? mapping.businessTraceId, env, { attributes: { runId, commandId: currentCommandId || null, afterSeq, rawEventCount: Array.isArray(response?.items) ? response.items.length : 0 } });
|
||||
const rawEvents = Array.isArray(response?.items) ? response.items : [];
|
||||
const events = rawEvents.filter((event) => agentRunEventBelongsToTrace(event, { currentCommandId, afterSeq, endSeq }));
|
||||
const maxSeq = Math.max(afterSeq, ...rawEvents.map((event) => Number(event?.seq ?? 0)).filter(Number.isFinite));
|
||||
const traceLastSeq = Math.max(afterSeq, ...events.map((event) => Number(event?.seq ?? 0)).filter(Number.isFinite));
|
||||
void emitCodeAgentOtelSpan("projection_sync", mapping.traceId ?? mapping.hwlabTraceId ?? mapping.businessTraceId, env, {
|
||||
attributes: {
|
||||
runId,
|
||||
commandId: currentCommandId || null,
|
||||
afterSeq,
|
||||
endSeq,
|
||||
limit,
|
||||
rawEventCount: rawEvents.length,
|
||||
eventCount: events.length,
|
||||
commandFiltered: Boolean(currentCommandId),
|
||||
maxSeq,
|
||||
traceLastSeq,
|
||||
terminalFromRawEvents: agentRunTerminalStatusFromEvents(rawEvents),
|
||||
terminalFromEvents: agentRunTerminalStatusFromEvents(events)
|
||||
}
|
||||
});
|
||||
return {
|
||||
events,
|
||||
afterSeq,
|
||||
endSeq,
|
||||
rawEventCount: rawEvents.length,
|
||||
commandFiltered: Boolean(currentCommandId),
|
||||
maxSeq: Math.max(afterSeq, ...rawEvents.map((event) => Number(event?.seq ?? 0))),
|
||||
traceLastSeq: Math.max(afterSeq, ...events.map((event) => Number(event?.seq ?? 0)).filter(Number.isFinite))
|
||||
maxSeq,
|
||||
traceLastSeq
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -1155,7 +1155,7 @@ function scheduleAgentRunProjectionSync({ traceId, params = {}, options = {}, tr
|
||||
pollIntervalMs,
|
||||
getCachedResult: () => options.codeAgentChatResults?.get?.(traceId),
|
||||
isCanceled: isCodeAgentResultCanceled,
|
||||
isTerminal: isTraceCommandTerminalStatus,
|
||||
isTerminal: agentRunProjectionObservedTerminal,
|
||||
activitySignature: agentRunProjectionActivitySignature,
|
||||
sleep: sleepAgentRunProjectionSync,
|
||||
performanceStore: options.backendPerformanceStore,
|
||||
@@ -1166,7 +1166,7 @@ function scheduleAgentRunProjectionSync({ traceId, params = {}, options = {}, tr
|
||||
traceStore,
|
||||
forceResultSync: agentRunProjectionObservedTerminal(result, traceStore?.snapshot?.(traceId))
|
||||
}),
|
||||
onTerminal: (payload, finalizerOptions = {}) => scheduleCodeAgentTerminalTurnStatusEffects({ payload, params, options, preserveLastTraceId: finalizerOptions.preserveLastTraceId === true })
|
||||
onTerminal: (payload, finalizerOptions = {}) => scheduleCodeAgentTerminalTurnStatusEffects({ payload: codeAgentPayloadWithObservedTerminalStatus(payload, traceStore?.snapshot?.(traceId)) ?? payload, params, options, preserveLastTraceId: finalizerOptions.preserveLastTraceId === true })
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1293,7 +1293,7 @@ async function runAgentRunProjectionResumePass({ options = {}, traceStore = defa
|
||||
try {
|
||||
const result = await resumeAgentRunProjectionCandidate({ candidate, options, traceStore });
|
||||
resumed += 1;
|
||||
if (isTraceCommandTerminalStatus(result?.status)) {
|
||||
if (agentRunProjectionObservedTerminal(result, result?.runnerTrace ?? traceStore?.snapshot?.(candidate.traceId))) {
|
||||
terminal += 1;
|
||||
candidateRetryState.delete(retryKey);
|
||||
} else {
|
||||
@@ -1513,7 +1513,7 @@ async function resumeAgentRunProjectionCandidate({ candidate, options = {}, trac
|
||||
forceResultSync: true,
|
||||
refreshEvents: true
|
||||
});
|
||||
const payload = synced?.result ?? currentResult;
|
||||
const payload = codeAgentPayloadWithObservedTerminalStatus(synced?.result ?? currentResult, synced?.runnerTrace ?? traceStore?.snapshot?.(candidate.traceId)) ?? (synced?.result ?? currentResult);
|
||||
if (isTraceCommandTerminalStatus(payload?.status)) {
|
||||
await recordCodeAgentTerminalTurnStatusEffects({ payload, params: candidate.params, options, preserveLastTraceId: true });
|
||||
}
|
||||
@@ -1661,30 +1661,48 @@ function agentRunProjectionActivitySignature(result, runnerTrace) {
|
||||
}
|
||||
|
||||
function agentRunProjectionObservedTerminal(result, runnerTrace) {
|
||||
return Boolean(codeAgentObservedTerminalStatus(result, runnerTrace));
|
||||
}
|
||||
|
||||
function codeAgentPayloadWithObservedTerminalStatus(payload, runnerTrace) {
|
||||
const terminalStatus = codeAgentObservedTerminalStatus(payload, runnerTrace);
|
||||
if (!terminalStatus || !payload || typeof payload !== "object") return null;
|
||||
const agentRun = payload.agentRun && typeof payload.agentRun === "object" ? payload.agentRun : null;
|
||||
const agentRunTerminalStatus = terminalStatus === "canceled" ? "cancelled" : terminalStatus;
|
||||
return {
|
||||
...payload,
|
||||
status: terminalStatus,
|
||||
runnerTrace: runnerTrace && typeof runnerTrace === "object" ? runnerTrace : payload.runnerTrace,
|
||||
...(agentRun ? { agentRun: { ...agentRun, terminalStatus: agentRun.terminalStatus ?? agentRunTerminalStatus, valuesPrinted: false } } : {})
|
||||
};
|
||||
}
|
||||
|
||||
function codeAgentObservedTerminalStatus(result, runnerTrace) {
|
||||
const agentRun = result?.agentRun && typeof result.agentRun === "object" ? result.agentRun : {};
|
||||
if (codeAgentPayloadHasSealedFinalResponse(result)) return true;
|
||||
if (
|
||||
isTraceCommandTerminalStatus(result?.status) ||
|
||||
isTraceCommandTerminalStatus(agentRun.status) ||
|
||||
isTraceCommandTerminalStatus(agentRun.runStatus) ||
|
||||
isTraceCommandTerminalStatus(agentRun.commandState) ||
|
||||
isTraceCommandTerminalStatus(agentRun.terminalStatus)
|
||||
) {
|
||||
return true;
|
||||
if (codeAgentPayloadHasSealedFinalResponse(result)) return "completed";
|
||||
for (const value of [result?.status, agentRun.status, agentRun.runStatus, agentRun.commandState, agentRun.terminalStatus]) {
|
||||
const terminalStatus = normalizeCodeAgentTerminalStatus(value);
|
||||
if (terminalStatus) return terminalStatus;
|
||||
}
|
||||
const commandId = textValue(agentRun.commandId);
|
||||
const trace = runnerTrace && typeof runnerTrace === "object" ? runnerTrace : result?.runnerTrace;
|
||||
const events = Array.isArray(trace?.events) ? trace.events : [];
|
||||
return events.some((event) => {
|
||||
for (const event of [...events].reverse()) {
|
||||
const payload = event?.payload && typeof event.payload === "object" ? event.payload : {};
|
||||
const eventCommandId = textValue(event?.commandId ?? payload.commandId);
|
||||
if (commandId && eventCommandId && eventCommandId !== commandId) return false;
|
||||
return event?.terminal === true ||
|
||||
isTraceCommandTerminalStatus(event?.status) ||
|
||||
isTraceCommandTerminalStatus(payload.status) ||
|
||||
isTraceCommandTerminalStatus(payload.terminalStatus) ||
|
||||
payload.phase === "command-terminal";
|
||||
});
|
||||
if (commandId && eventCommandId && eventCommandId !== commandId) continue;
|
||||
for (const value of [event?.status, payload.status, payload.terminalStatus, payload.phase === "command-terminal" ? payload.terminalStatus : null, event?.terminal === true ? "completed" : null]) {
|
||||
const terminalStatus = normalizeCodeAgentTerminalStatus(value);
|
||||
if (terminalStatus) return terminalStatus;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function normalizeCodeAgentTerminalStatus(value) {
|
||||
const status = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-");
|
||||
if (!CODE_AGENT_TERMINAL_STATUSES.has(status)) return null;
|
||||
return status === "cancelled" ? "canceled" : status;
|
||||
}
|
||||
|
||||
function sleepAgentRunProjectionSync(ms) {
|
||||
@@ -2594,6 +2612,7 @@ function recordCodeAgentConversationFact(payload = {}, options = {}) {
|
||||
}
|
||||
|
||||
async function recordCodeAgentTerminalTurnStatusEffects({ payload = {}, params = {}, options = {}, preserveLastTraceId = false } = {}) {
|
||||
payload = codeAgentPayloadWithObservedTerminalStatus(payload, payload?.runnerTrace) ?? payload;
|
||||
if (!payload || typeof payload !== "object" || !isTraceCommandTerminalStatus(payload.status)) return null;
|
||||
if (payload.turnStatusTerminalEffects?.recorded === true) return payload.turnStatusTerminalEffects;
|
||||
const billing = await finalizeCodeAgentBillingUsage({ payload, params, options });
|
||||
@@ -2613,11 +2632,25 @@ async function recordCodeAgentTerminalTurnStatusEffects({ payload = {}, params =
|
||||
};
|
||||
payload.turnStatusTerminalEffects = effects;
|
||||
const traceId = safeTraceId(payload.traceId ?? params.traceId);
|
||||
if (traceId) void emitCodeAgentOtelSpan("projection_write", traceId, options.env ?? process.env, {
|
||||
attributes: {
|
||||
projection: "terminal-effects",
|
||||
status: payload.status,
|
||||
terminal: true,
|
||||
billingSettled,
|
||||
ownerSettled,
|
||||
sessionId: safeSessionId(payload.sessionId ?? params.sessionId) || null,
|
||||
turnId: codeAgentTurnLifecycleFields(traceId, payload).turnId,
|
||||
runId: payload.agentRun?.runId ?? null,
|
||||
commandId: payload.agentRun?.commandId ?? null
|
||||
}
|
||||
});
|
||||
if (traceId) options.codeAgentChatResults?.set?.(traceId, payload);
|
||||
return effects;
|
||||
}
|
||||
|
||||
function scheduleCodeAgentTerminalTurnStatusEffects({ payload = {}, params = {}, options = {}, preserveLastTraceId = false } = {}) {
|
||||
payload = codeAgentPayloadWithObservedTerminalStatus(payload, payload?.runnerTrace) ?? payload;
|
||||
if (!payload || typeof payload !== "object" || !isTraceCommandTerminalStatus(payload.status)) return null;
|
||||
const existing = payload.turnStatusTerminalEffects;
|
||||
if (existing?.recorded === true || existing?.pending === true) return existing;
|
||||
|
||||
@@ -18,7 +18,7 @@ export function scheduleWorkbenchProjectionFinalizer({ traceId, currentResult =
|
||||
while (true) {
|
||||
const cached = getCachedResult?.(traceId);
|
||||
if (isCanceled?.(cached)) return;
|
||||
if (isTerminal?.(cached?.status)) {
|
||||
if (cached && isTerminal?.(cached, traceStore?.snapshot?.(traceId))) {
|
||||
performanceStore?.recordWorkbenchProjectorCandidate?.({ status: "completed", reason: "cached_terminal" });
|
||||
onTerminal?.(cached, { preserveLastTraceId: true });
|
||||
return;
|
||||
@@ -28,12 +28,13 @@ export function scheduleWorkbenchProjectionFinalizer({ traceId, currentResult =
|
||||
const synced = await syncResult({ result });
|
||||
performanceStore?.recordWorkbenchProjectorBatch?.({ phase: "candidate_scan", status: "ok", durationMs: Date.now() - syncStartedAt });
|
||||
result = synced?.result ?? result;
|
||||
if (result && isTerminal?.(result.status)) {
|
||||
const runnerTrace = synced?.runnerTrace ?? traceStore?.snapshot?.(traceId);
|
||||
if (result && isTerminal?.(result, runnerTrace)) {
|
||||
performanceStore?.recordWorkbenchProjectorCandidate?.({ status: "completed", reason: "synced_terminal" });
|
||||
onTerminal?.(result, { preserveLastTraceId: true });
|
||||
return;
|
||||
}
|
||||
const signature = activitySignature?.(result, synced?.runnerTrace ?? traceStore?.snapshot?.(traceId));
|
||||
const signature = activitySignature?.(result, runnerTrace);
|
||||
if (signature && signature !== lastActivitySignature) {
|
||||
lastActivitySignature = signature;
|
||||
lastActivityAt = Date.now();
|
||||
|
||||
Reference in New Issue
Block a user