From 04df7bf47fb38e6a36821f6db72c9da6a92e5ae3 Mon Sep 17 00:00:00 2001 From: liujing Date: Mon, 20 Jul 2026 20:51:37 +0800 Subject: [PATCH] fix: harden live workflow log streaming --- src/workflow-web-assets.ts | 66 +- src/workflow-web.test.ts | 1498 +++++++++++++++++++++++++++++++++++- src/workflow-web.ts | 103 ++- 3 files changed, 1640 insertions(+), 27 deletions(-) diff --git a/src/workflow-web-assets.ts b/src/workflow-web-assets.ts index 3838e5a..669b7e2 100644 --- a/src/workflow-web-assets.ts +++ b/src/workflow-web-assets.ts @@ -1788,6 +1788,31 @@ export function getAssetsHtml(): string { if (container && parsed.length) { const frag = document.createDocumentFragment(); parsed.forEach(entry => { + if (entry.type === 'progress') { + // Container already has a progress row from a newer page — + // skip older progress so we never overwrite newest values. + const containerRow = container.querySelector('[data-progress-key="progress"]'); + if (containerRow) { + // The container row holds the most recent progress value; + // historical entries must not overwrite it. + return; + } + // No container row — check whether this batch already queued + // one (fragment rows are invisible to container.querySelector). + const fragRow = frag.querySelector('[data-progress-key="progress"]'); + if (fragRow) { + // Update fragment-local row — entries arrive in file order + // so the last one is most recent. + const content = fragRow.querySelector('.item-content'); + if (content) content.textContent = entry.content || entry.raw || ''; + return; + } + // First progress entry we have seen — create a keyed row. + entry._progressKey = 'progress'; + const row = buildTimelineRow(entry); + if (row) frag.appendChild(row); + return; + } const row = buildTimelineRow(entry); if (row) frag.appendChild(row); }); @@ -1899,18 +1924,30 @@ export function getAssetsHtml(): string { const container = document.getElementById('timeline-items-container'); if (!container || !entries || !entries.length) return; const frag = document.createDocumentFragment(); + // Track the most recent progress row in THIS batch so multiple progress + // entries arriving in one call still merge into one row. + let batchProgressRow = null; entries.forEach(entry => { // Merge high-frequency progress counters into one updatable row. if (entry.type === 'progress') { const key = 'progress'; entry._progressKey = key; - const existing = container.querySelector('[data-progress-key="progress"]'); + // Check batch-local row first (same call), then container (prior calls). + const existing = batchProgressRow || container.querySelector('[data-progress-key="progress"]'); if (existing) { const content = existing.querySelector('.item-content'); if (content) content.textContent = entry.content || entry.raw || ''; timelineProgressRow = existing; + batchProgressRow = existing; return; } + // First progress in this batch – create the row. + const row = buildTimelineRow(entry); + if (row) { + batchProgressRow = row; + frag.appendChild(row); + } + return; } const row = buildTimelineRow(entry); if (row) frag.appendChild(row); @@ -1945,13 +1982,30 @@ export function getAssetsHtml(): string { container.innerHTML = ''; timelineProgressRow = null; const frag = document.createDocumentFragment(); + // Merge high-frequency progress entries into a single updatable row + // (same logic as appendTimelineEntries). + let progressMerged = false; timelineEntries.forEach(entry => { - if (entry.type === 'progress') entry._progressKey = 'progress'; - const row = buildTimelineRow(entry); - if (row) { - if (entry.type === 'progress') timelineProgressRow = row; - frag.appendChild(row); + if (entry.type === 'progress') { + if (!progressMerged) { + progressMerged = true; + entry._progressKey = 'progress'; + const row = buildTimelineRow(entry); + if (row) { + timelineProgressRow = row; + frag.appendChild(row); + } + } else { + // Update the existing row's text instead of creating a new one. + if (timelineProgressRow) { + const contentEl = timelineProgressRow.querySelector('.item-content'); + if (contentEl) contentEl.textContent = entry.content || entry.raw || ''; + } + } + return; } + const row = buildTimelineRow(entry); + if (row) frag.appendChild(row); }); container.appendChild(frag); timelineRenderedCount = timelineEntries.length; diff --git a/src/workflow-web.test.ts b/src/workflow-web.test.ts index 03a1377..42aef7b 100644 --- a/src/workflow-web.test.ts +++ b/src/workflow-web.test.ts @@ -4,7 +4,7 @@ import os from "node:os"; import path from "node:path"; import http from "node:http"; import { Buffer } from "node:buffer"; -import { startWorkflowWebServer, splitCompleteUtf8, hasSymlinkInPath } from "./workflow-web.js"; +import { startWorkflowWebServer, splitCompleteUtf8, leadingIncompleteUtf8Bytes, readLogSlice, hasSymlinkInPath } from "./workflow-web.js"; import { writeState, appendEvent, attemptLogPath, readState } from "./workflow-state.js"; import type { WorkflowState } from "./workflow-types.js"; @@ -37,6 +37,7 @@ describe("workflow-web", () => { let tmpDir: string; let server: http.Server; let originalHome: string | undefined; + let savedAgentDir: string | undefined; let testHome: string; let port = 14319; let baseUrl = `http://127.0.0.1:${port}`; @@ -47,6 +48,7 @@ describe("workflow-web", () => { beforeAll(async () => { // Setup temp workflows root directory tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "agent-workflow-web-test-")); + savedAgentDir = process.env.AGENT_WORKFLOW_DIR; process.env.AGENT_WORKFLOW_DIR = tmpDir; originalHome = process.env.HOME; testHome = path.join(tmpDir, "home"); @@ -115,6 +117,8 @@ printf '%s\\n' '{"providers":[{"name":"test-provider","models":["test-model"]}]} afterAll(() => { return new Promise((resolve) => { server.close(() => { + if (savedAgentDir === undefined) delete process.env.AGENT_WORKFLOW_DIR; + else process.env.AGENT_WORKFLOW_DIR = savedAgentDir; if (originalHome === undefined) delete process.env.HOME; else process.env.HOME = originalHome; fs.rmSync(tmpDir, { recursive: true, force: true }); @@ -435,7 +439,7 @@ printf '%s\\n' '{"providers":[{"name":"test-provider","models":["test-model"]}]} fetch: async () => ({ ok: true, json: async () => ({}), text: async () => "" }), EventSource: class { addEventListener() {} }, setInterval: () => {}, - console: { log: () => {}, error: () => {} }, + console: { log: (...args: any[]) => {}, error: (...args: any[]) => console.log("ERR:", ...args.map((a: any) => a?.message || a)) }, }; vm.createContext(sandbox); @@ -449,6 +453,1016 @@ printf '%s\\n' '{"providers":[{"name":"test-provider","models":["test-model"]}]} expect(typeof sandbox.resetRawLogs).toBe("function"); }); + it("progress entries merge into a single row via appendTimelineEntries", async () => { + const { getAssetsHtml } = await import("./workflow-web-assets.js"); + const html = getAssetsHtml(); + const scriptMatch = html.match(/