diff --git a/internal/cloud/code-agent-agentrun-adapter.ts b/internal/cloud/code-agent-agentrun-adapter.ts index 86cbf35d..f0933420 100644 --- a/internal/cloud/code-agent-agentrun-adapter.ts +++ b/internal/cloud/code-agent-agentrun-adapter.ts @@ -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 }; } diff --git a/internal/cloud/server-code-agent-http.ts b/internal/cloud/server-code-agent-http.ts index 011d6499..22a2571a 100644 --- a/internal/cloud/server-code-agent-http.ts +++ b/internal/cloud/server-code-agent-http.ts @@ -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; diff --git a/internal/cloud/workbench-projection-finalizer.ts b/internal/cloud/workbench-projection-finalizer.ts index 3c7c6ebd..42841a66 100644 --- a/internal/cloud/workbench-projection-finalizer.ts +++ b/internal/cloud/workbench-projection-finalizer.ts @@ -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();