diff --git a/.env.example b/.env.example index c801fe3..27648ba 100644 --- a/.env.example +++ b/.env.example @@ -34,6 +34,8 @@ H5_PUBLIC_BASE_URL=http://127.0.0.1:5173 # MEMIND_AGENT_RUN_AUTODISPATCH=1 # MEMIND_AGENT_RUN_WORKER_POLL_MS=1000 # MEMIND_AGENT_RUN_WORKER_BATCH_SIZE=1 +# MEMIND_AGENT_RUN_HEARTBEAT_MS=30000 +# MEMIND_AGENT_RUN_WORKER_EXPECT_RUNNING=0 # Runtime SLO 日报(默认每日 23:55 写 reports/runtime-slo)。 # MEMIND_RUNTIME_REPORT_DIR=/Users/john/Project/Memind/reports/runtime-slo diff --git a/agent-run-gateway.mjs b/agent-run-gateway.mjs index 7ba2723..ccce195 100644 --- a/agent-run-gateway.mjs +++ b/agent-run-gateway.mjs @@ -8,6 +8,7 @@ const CODE_TOOL_MODES = new Set(['code', 'code-task', 'code_task', 'code-tool', const RUN_METADATA_KEY = 'memindRun'; const DEFAULT_MAX_CONCURRENT_RUNS = 1; const DEFAULT_RUN_TIMEOUT_MS = 15 * 60 * 1000; +const DEFAULT_RUN_HEARTBEAT_MS = 30 * 1000; const TOOL_GATEWAY_SUMMARY_LIMIT = 4096; function nowMs() { @@ -218,6 +219,10 @@ export function createAgentRunGateway({ process.env.MEMIND_AGENT_RUN_TIMEOUT_MS, DEFAULT_RUN_TIMEOUT_MS, ), + heartbeatMs = positiveInteger( + process.env.MEMIND_AGENT_RUN_HEARTBEAT_MS, + DEFAULT_RUN_HEARTBEAT_MS, + ), }) { const inFlight = new Set(); const queuedDispatches = []; @@ -338,6 +343,31 @@ export function createAgentRunGateway({ await appendEvent(runId, status, fields); } + function startRunHeartbeat(runId, { attempt }) { + let stopped = false; + const writeHeartbeat = async () => { + if (stopped) return; + try { + await appendEvent(runId, 'worker_heartbeat', { + attempt, + pid: process.pid, + heartbeatMs, + }); + } catch (err) { + console.error('[AgentRun] worker heartbeat failed:', err instanceof Error ? err.message : err); + } + }; + void writeHeartbeat(); + const timer = setInterval(() => { + void writeHeartbeat(); + }, heartbeatMs); + timer.unref?.(); + return () => { + stopped = true; + clearInterval(timer); + }; + } + async function runWithTimeout(runId, task) { let timer = null; const timeout = new Promise((_, reject) => { @@ -440,6 +470,7 @@ export function createAgentRunGateway({ ); if (Number(claim?.affectedRows ?? 0) === 0) return; await appendEvent(runId, 'running', { attempt: nextAttempt }); + const stopHeartbeat = startRunHeartbeat(runId, { attempt: nextAttempt }); try { const { sessionId } = await runWithTimeout(runId, () => executeRun(row, runId)); @@ -461,6 +492,8 @@ export function createAgentRunGateway({ if (retryable && autoDispatch) { setTimeout(() => dispatchRun(runId), retryDelaysMs[nextAttempt - 1]); } + } finally { + stopHeartbeat(); } } @@ -494,25 +527,69 @@ export function createAgentRunGateway({ statusCounts[row.status] = Number(row.count ?? 0); } const [oldestRunningRows] = await pool.query( - `SELECT MIN(started_at) AS oldest_started_at - FROM h5_agent_runs - WHERE status = 'running'`, + `SELECT + r.id, + r.request_id, + r.started_at, + r.updated_at, + r.attempts, + h.latest_heartbeat_at + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' + ORDER BY COALESCE(h.latest_heartbeat_at, r.started_at) ASC + LIMIT 1`, ); - const oldestRunningStartedAt = oldestRunningRows[0]?.oldest_started_at == null + const oldestRunningStartedAt = oldestRunningRows[0]?.started_at == null ? null - : Number(oldestRunningRows[0].oldest_started_at); + : Number(oldestRunningRows[0].started_at); const oldestRunningAgeMs = oldestRunningStartedAt == null ? 0 : Math.max(0, nowMs() - oldestRunningStartedAt); + const oldestRunningHeartbeatAt = oldestRunningRows[0]?.latest_heartbeat_at == null + ? null + : Number(oldestRunningRows[0].latest_heartbeat_at); + const oldestRunningHeartbeatAgeMs = oldestRunningRows[0] + ? Math.max(0, nowMs() - Number(oldestRunningRows[0].latest_heartbeat_at ?? oldestRunningRows[0].started_at ?? nowMs())) + : 0; + const [runningWithoutHeartbeatRows] = await pool.query( + `SELECT COUNT(*) AS count + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' AND h.latest_heartbeat_at IS NULL`, + ); return { autoDispatch, maxConcurrentRuns, runTimeoutMs, + heartbeatMs, inFlight: inFlight.size, pendingDispatches: queuedDispatches.length, statusCounts, oldestRunningStartedAt, oldestRunningAgeMs, + oldestRunningHeartbeatAt, + oldestRunningHeartbeatAgeMs, + runningWithoutHeartbeatCount: Number(runningWithoutHeartbeatRows[0]?.count ?? 0), + latestRunningRun: oldestRunningRows[0] ? { + id: oldestRunningRows[0].id, + requestId: oldestRunningRows[0].request_id, + startedAt: oldestRunningStartedAt, + updatedAt: oldestRunningRows[0].updated_at == null ? null : Number(oldestRunningRows[0].updated_at), + attempts: Number(oldestRunningRows[0].attempts ?? 0), + heartbeatAt: oldestRunningHeartbeatAt, + heartbeatAgeMs: oldestRunningHeartbeatAgeMs, + } : null, terminalStatuses: [...TERMINAL_STATUSES], toolGateway: toolGateway?.getStatus ? toolGateway.getStatus() : null, }; @@ -528,32 +605,60 @@ export function createAgentRunGateway({ const normalizedLimit = positiveInteger(limit, maxConcurrentRuns); const cutoff = nowMs() - normalizedStaleMs; const [rows] = await pool.query( - `SELECT id, request_id, started_at, updated_at, attempts - FROM h5_agent_runs - WHERE status = 'running' AND started_at IS NOT NULL AND started_at <= ? - ORDER BY started_at ASC + `SELECT + r.id, + r.request_id, + r.started_at, + r.updated_at, + r.attempts, + h.latest_heartbeat_at + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' + AND r.started_at IS NOT NULL + AND COALESCE(h.latest_heartbeat_at, r.started_at) <= ? + ORDER BY COALESCE(h.latest_heartbeat_at, r.started_at) ASC LIMIT ?`, [cutoff, normalizedLimit], ); const recovered = []; for (const row of rows) { const ageMs = Math.max(0, nowMs() - Number(row.started_at ?? 0)); + const heartbeatAt = row.latest_heartbeat_at == null ? null : Number(row.latest_heartbeat_at); + const heartbeatAgeMs = Math.max(0, nowMs() - Number(row.latest_heartbeat_at ?? row.started_at ?? 0)); const item = { id: row.id, requestId: row.request_id, startedAt: Number(row.started_at ?? 0), updatedAt: Number(row.updated_at ?? 0), attempts: Number(row.attempts ?? 0), + heartbeatAt, + heartbeatAgeMs, ageMs, reason, }; if (!dryRun) { - const message = `agent run recovered from stale running state after ${ageMs}ms`; + const message = heartbeatAt == null + ? `agent run recovered from stale running state after ${ageMs}ms without heartbeat` + : `agent run recovered from stale running state after heartbeat was stale for ${heartbeatAgeMs}ms`; const completedAt = nowMs(); const [update] = await pool.query( `UPDATE h5_agent_runs SET status = 'failed', error_message = ?, updated_at = ?, completed_at = ? - WHERE id = ? AND status = 'running' AND started_at <= ?`, + WHERE id = ? + AND status = 'running' + AND started_at IS NOT NULL + AND COALESCE( + (SELECT MAX(created_at) + FROM h5_agent_run_events + WHERE run_id = h5_agent_runs.id AND event_type = 'worker_heartbeat'), + started_at + ) <= ?`, [message, completedAt, completedAt, row.id, cutoff], ); if (Number(update?.affectedRows ?? 0) === 0) continue; @@ -561,6 +666,8 @@ export function createAgentRunGateway({ reason, staleMs: normalizedStaleMs, ageMs, + heartbeatAt, + heartbeatAgeMs, status: 'failed', error: message, }); diff --git a/agent-run-gateway.test.mjs b/agent-run-gateway.test.mjs index f78ddc8..e8f71ef 100644 --- a/agent-run-gateway.test.mjs +++ b/agent-run-gateway.test.mjs @@ -9,6 +9,13 @@ function createFakePool() { const runs = new Map(); const events = []; + const latestHeartbeatAt = (runId) => { + const timestamps = events + .filter((event) => event.runId === runId && event.eventType === 'worker_heartbeat') + .map((event) => Number(event.createdAt ?? 0)); + return timestamps.length ? Math.max(...timestamps) : null; + }; + return { runs, events, @@ -33,24 +40,45 @@ function createFakePool() { } return [[...counts].map(([status, count]) => ({ status, count }))]; } - if (sql.includes('SELECT MIN(started_at) AS oldest_started_at')) { - const started = [...runs.values()] - .filter((row) => row.status === 'running' && row.started_at != null) - .map((row) => Number(row.started_at)); - return [[{ oldest_started_at: started.length ? Math.min(...started) : null }]]; - } - if (sql.includes('SELECT id, request_id, started_at, updated_at, attempts')) { - const [cutoff, limit = 1] = params; + if (sql.includes('h.latest_heartbeat_at') && sql.includes("WHERE r.status = 'running'") && sql.includes('LIMIT 1')) { return [[...runs.values()] - .filter((row) => row.status === 'running' && row.started_at != null && Number(row.started_at) <= Number(cutoff)) - .sort((a, b) => Number(a.started_at) - Number(b.started_at)) - .slice(0, Number(limit)) + .filter((row) => row.status === 'running') .map((row) => ({ id: row.id, request_id: row.request_id, started_at: row.started_at, updated_at: row.updated_at, attempts: row.attempts, + latest_heartbeat_at: latestHeartbeatAt(row.id), + })) + .sort((a, b) => Number(a.latest_heartbeat_at ?? a.started_at ?? 0) - Number(b.latest_heartbeat_at ?? b.started_at ?? 0)) + .slice(0, 1)]; + } + if (sql.includes('SELECT COUNT(*) AS count') && sql.includes('worker_heartbeat')) { + return [[{ + count: [...runs.values()] + .filter((row) => row.status === 'running' && latestHeartbeatAt(row.id) == null) + .length, + }]]; + } + if (sql.includes('SELECT') && sql.includes('latest_heartbeat_at') && sql.includes('COALESCE(h.latest_heartbeat_at, r.started_at) <= ?')) { + const [cutoff, limit = 1] = params; + return [[...runs.values()] + .map((row) => ({ row, heartbeatAt: latestHeartbeatAt(row.id) })) + .filter(({ row, heartbeatAt }) => ( + row.status === 'running' && + row.started_at != null && + Number(heartbeatAt ?? row.started_at) <= Number(cutoff) + )) + .sort((a, b) => Number(a.heartbeatAt ?? a.row.started_at) - Number(b.heartbeatAt ?? b.row.started_at)) + .slice(0, Number(limit)) + .map(({ row, heartbeatAt }) => ({ + id: row.id, + request_id: row.request_id, + started_at: row.started_at, + updated_at: row.updated_at, + attempts: row.attempts, + latest_heartbeat_at: heartbeatAt, }))]; } if (sql.includes('SELECT id') && sql.includes("status IN ('queued', 'retryable')")) { @@ -107,10 +135,15 @@ function createFakePool() { }); return [{ affectedRows: 1 }]; } - if (sql.includes("WHERE id = ? AND status = 'running' AND started_at <= ?")) { + if (sql.includes("WHERE id = ?") && sql.includes("status = 'running'") && sql.includes('worker_heartbeat')) { const [errorMessage, updatedAt, completedAt, id, cutoff] = params; const row = runs.get(id); - if (!row || row.status !== 'running' || Number(row.started_at) > Number(cutoff)) { + if ( + !row || + row.status !== 'running' || + row.started_at == null || + Number(latestHeartbeatAt(id) ?? row.started_at) > Number(cutoff) + ) { return [{ affectedRows: 0 }]; } Object.assign(row, { @@ -598,6 +631,7 @@ test('external worker dispatches queued runs through the same queue controls', a const result = await gateway.dispatchQueuedRuns({ limit: 10 }); assert.equal(result.dispatched, 1); + await waitFor(() => pool.events.some((event) => event.runId === run.id && event.eventType === 'worker_heartbeat')); await waitFor(() => pool.runs.get(run.id)?.status === 'succeeded'); assert.deepEqual(submitted, ['req-worker']); }); @@ -710,3 +744,87 @@ test('stale running recovery marks old running rows failed with an event', async true, ); }); + +test('stale running recovery ignores old running rows with fresh heartbeat', async () => { + const pool = createFakePool(); + const gateway = createAgentRunGateway({ + pool, + userAuth: {}, + tkmindProxy: {}, + autoDispatch: false, + runTimeoutMs: 1000, + }); + + const run = await gateway.createRun('user-1', { + requestId: 'req-stale-heartbeat-fresh', + userMessage: { role: 'user', content: [] }, + }); + Object.assign(pool.runs.get(run.id), { + status: 'running', + attempts: 1, + started_at: Date.now() - 5000, + updated_at: Date.now() - 5000, + }); + pool.events.push({ + id: 'heartbeat-1', + runId: run.id, + eventType: 'worker_heartbeat', + dataJson: JSON.stringify({ attempt: 1 }), + createdAt: Date.now(), + }); + + const result = await gateway.recoverStaleRunningRuns({ staleMs: 1000, dryRun: false }); + + assert.equal(result.considered, 0); + assert.equal(result.recovered, 0); + assert.equal(pool.runs.get(run.id).status, 'running'); +}); + +test('queue status reports running heartbeat age and missing heartbeat count', async () => { + const pool = createFakePool(); + const gateway = createAgentRunGateway({ + pool, + userAuth: {}, + tkmindProxy: {}, + autoDispatch: false, + runTimeoutMs: 1000, + heartbeatMs: 250, + }); + + const runWithHeartbeat = await gateway.createRun('user-1', { + requestId: 'req-heartbeat-status-1', + userMessage: { role: 'user', content: [] }, + }); + Object.assign(pool.runs.get(runWithHeartbeat.id), { + status: 'running', + attempts: 1, + started_at: Date.now() - 5000, + updated_at: Date.now() - 5000, + }); + pool.events.push({ + id: 'heartbeat-status-1', + runId: runWithHeartbeat.id, + eventType: 'worker_heartbeat', + dataJson: JSON.stringify({ attempt: 1 }), + createdAt: Date.now() - 100, + }); + + const runWithoutHeartbeat = await gateway.createRun('user-1', { + requestId: 'req-heartbeat-status-2', + userMessage: { role: 'user', content: [] }, + }); + Object.assign(pool.runs.get(runWithoutHeartbeat.id), { + status: 'running', + attempts: 1, + started_at: Date.now() - 2000, + updated_at: Date.now() - 2000, + }); + + const status = await gateway.getQueueStatus(); + + assert.equal(status.heartbeatMs, 250); + assert.equal(status.runningWithoutHeartbeatCount, 1); + assert.equal(status.latestRunningRun.id, runWithoutHeartbeat.id); + assert.equal(status.oldestRunningHeartbeatAt, null); + assert.ok(status.oldestRunningHeartbeatAgeMs >= 1900); +}); diff --git a/docs/agent-run-worker-rollout-runbook.md b/docs/agent-run-worker-rollout-runbook.md index bf804da..d01a1c3 100644 --- a/docs/agent-run-worker-rollout-runbook.md +++ b/docs/agent-run-worker-rollout-runbook.md @@ -82,6 +82,8 @@ MEMIND_AGENT_CODE_RUN_TASK_TYPES= MEMIND_AGENT_CODE_RUNS_REQUIRE_VALIDATION=1 MEMIND_AGENT_RUN_QUEUE_CONCURRENCY=2 MEMIND_AGENT_RUN_WORKER_BATCH_SIZE=2 +MEMIND_AGENT_RUN_HEARTBEAT_MS=30000 +MEMIND_AGENT_RUN_WORKER_EXPECT_RUNNING=1 MEMIND_TOOL_GATEWAY_ENABLED=0 MEMIND_TOOL_GATEWAY_DRY_RUN=0 VITE_AGENT_CODE_RUNS_ENABLED=1 @@ -96,6 +98,7 @@ MEMIND_AGENT_RUN_WORKER_START=1 MEMIND_AGENT_RUN_WORKER_POLL_MS=3000 MEMIND_AGENT_RUN_WORKER_BATCH_SIZE=2 MEMIND_AGENT_RUN_QUEUE_CONCURRENCY=2 +MEMIND_AGENT_RUN_HEARTBEAT_MS=30000 MEMIND_TOOL_GATEWAY_ENABLED=1 MEMIND_TOOL_GATEWAY_DRY_RUN=0 ``` @@ -186,6 +189,20 @@ Check stale running recovery before and after worker changes: node /Users/john/Project/Memind/scripts/agent-run-worker.mjs --recover-stale --stale-ms 900000 --limit 5 ``` +Check worker heartbeat status during all-user gray: + +```bash +MEMIND_AGENT_RUN_WORKER_EXPECT_RUNNING=1 node /Users/john/Project/Memind/scripts/check-agent-run-worker.mjs +``` + +Expected queue fields: + +```text +heartbeatMs=30000 +oldestRunningHeartbeatAgeMs=0 when no running rows +runningWithoutHeartbeatCount=0 during normal operation +``` + Apply only when the listed running rows are known stale: ```bash diff --git a/docs/architecture/memind-2-runtime-execution-assessment-20260702.md b/docs/architecture/memind-2-runtime-execution-assessment-20260702.md index fd2bd67..ddeb28e 100644 --- a/docs/architecture/memind-2-runtime-execution-assessment-20260702.md +++ b/docs/architecture/memind-2-runtime-execution-assessment-20260702.md @@ -19,12 +19,12 @@ - Stream Controller 已具备 SSE headers、abort propagation、backpressure pipeline。 - Redis Router 已启用,承担 worker runtime state。 - Goose Worker Pool 已从固定单点走向四 worker 可观测调度。 -- Aider/OpenHands 已从普通聊天默认能力中剥离,进入 code mode 和后端灰度门禁;P6.3 已把 code run 从 goosed session extension 外移到 `agent-run-v1` Tool Gateway 协议,P6.4 已完成 Aider 真实执行 canary,P6.5 已完成 OpenHands 真实执行 canary,P6.6 已完成 external worker 精确接管 code-run canary,P6.7 已加入 Tool Gateway 产物校验与输出审计,P6.8 已安装 external worker LaunchAgent,P6.9 已完成带 validation 的 external worker 灰度 canary,P6.10 已加入 external worker 只读观测脚本,P6.11 已加入放量策略门禁,P6.12 已开启全用户长期灰度并通过普通测试用户真实路径,P6.13 已安装自动暂停 guard,P6.15 已让 H5 code-run 自动补 receipt validation 并恢复后端 required validation,P6.17/P6.18 已把 external worker 并发 2 通过 canary 并固化为当前 all-user gray 策略,P6.19 已加入 task-level artifact validation,P8.1 已加入 stale running recovery,P8.2 已完成 Portal DB/Auth 瞬时错误兜底。 +- Aider/OpenHands 已从普通聊天默认能力中剥离,进入 code mode 和后端灰度门禁;P6.3 已把 code run 从 goosed session extension 外移到 `agent-run-v1` Tool Gateway 协议,P6.4 已完成 Aider 真实执行 canary,P6.5 已完成 OpenHands 真实执行 canary,P6.6 已完成 external worker 精确接管 code-run canary,P6.7 已加入 Tool Gateway 产物校验与输出审计,P6.8 已安装 external worker LaunchAgent,P6.9 已完成带 validation 的 external worker 灰度 canary,P6.10 已加入 external worker 只读观测脚本,P6.11 已加入放量策略门禁,P6.12 已开启全用户长期灰度并通过普通测试用户真实路径,P6.13 已安装自动暂停 guard,P6.15 已让 H5 code-run 自动补 receipt validation 并恢复后端 required validation,P6.17/P6.18 已把 external worker 并发 2 通过 canary 并固化为当前 all-user gray 策略,P6.19 已加入 task-level artifact validation,P8.1 已加入 stale running recovery,P8.2 已完成 Portal DB/Auth 瞬时错误兜底,P8.3 已加入 worker lease heartbeat。 - PG 和 MindSpace 仍保持生产数据边界,SLO 报告只做统计读取;P6.3-P6.15 不新增 schema migration,不删除或修改既有用户数据。 整体执行评分: 9.95 / 10。 -可以支撑当前 H5 streaming 稳定性改造的基础目标。自动采样、worker sidecar heartbeat、SLO 只读快照、SLO 日报定时器、SLO 日报保留策略、first-token EWMA、first-token p50/p95 窗口趋势、Tool Gateway Queue v0、外部 worker 接管入口、真实 worker canary、后端 code-mode canary、code-run 用户级灰度 gate、H5 页面编辑 UI canary、P6.3 Tool Gateway 协议化、P6.4 Aider 真实 canary、P6.5 OpenHands 真实 canary、P6.6 external worker code-run canary、P6.7 Tool Gateway guardrails、P6.8 worker LaunchAgent、P6.9 validated external worker canary、P6.10 worker observability、P6.11 rollout policy gates、P6.12 all-user gray、P6.13 auto-pause guard、P6.15 H5 validation metadata、P6.17/P6.18 worker concurrency 2、P6.19 task-level artifact validation、P8.1 stale running recovery 和 P8.2 Portal DB/Auth transient hardening 已经落地。主要剩余差距转为用户可见进度/失败说明,以及后续更细任务类型的校验策略扩展。 +可以支撑当前 H5 streaming 稳定性改造的基础目标。自动采样、worker sidecar heartbeat、SLO 只读快照、SLO 日报定时器、SLO 日报保留策略、first-token EWMA、first-token p50/p95 窗口趋势、Tool Gateway Queue v0、外部 worker 接管入口、真实 worker canary、后端 code-mode canary、code-run 用户级灰度 gate、H5 页面编辑 UI canary、P6.3 Tool Gateway 协议化、P6.4 Aider 真实 canary、P6.5 OpenHands 真实 canary、P6.6 external worker code-run canary、P6.7 Tool Gateway guardrails、P6.8 worker LaunchAgent、P6.9 validated external worker canary、P6.10 worker observability、P6.11 rollout policy gates、P6.12 all-user gray、P6.13 auto-pause guard、P6.15 H5 validation metadata、P6.17/P6.18 worker concurrency 2、P6.19 task-level artifact validation、P8.1 stale running recovery、P8.2 Portal DB/Auth transient hardening 和 P8.3 worker lease heartbeat 已经落地。P6.22 用户可见进度/失败说明按用户决策暂不开发,后续剩余差距转为更长时间 soak、自动报表告警和更细任务类型的校验策略扩展。 ## 实测结果 @@ -974,9 +974,58 @@ Data boundary: - guard `shouldPause=false` - `failedRecentCount=0` +### P8.3 Worker Lease Heartbeat + +结果: 通过,worker 执行 running run 时会写 `worker_heartbeat`,stale recovery、runtime/status、worker check 和 auto-pause guard 已优先使用 heartbeat age。 + +- `agent-run-gateway.mjs`: + - 新增 `MEMIND_AGENT_RUN_HEARTBEAT_MS`,默认 `30000`。 + - running 后立即写 `worker_heartbeat`,再按 interval 定时写入。 + - terminal 后停止 heartbeat timer。 + - stale recovery 使用 `COALESCE(latest worker_heartbeat, started_at)`。 + - recovery UPDATE 会复核最新 heartbeat,避免误杀刚刷新 heartbeat 的 run。 +- runtime/status queue 摘要新增: + - `heartbeatMs` + - `oldestRunningHeartbeatAt` + - `oldestRunningHeartbeatAgeMs` + - `runningWithoutHeartbeatCount` + - `latestRunningRun` +- `agent-run-guard.mjs`: + - running guard 改用 `oldestRunningHeartbeatAgeMs`。 +- `check-agent-run-worker.mjs`: + - 增加 heartbeat queue 摘要。 +- 测试: + - `node --check agent-run-gateway.mjs scripts/agent-run-worker.mjs scripts/agent-run-guard.mjs scripts/check-agent-run-worker.mjs` + - `node --test agent-run-gateway.test.mjs agent-run-routes.test.mjs`,31 tests pass。 +- 生产部署: + - backup `/Users/john/Project/memind_backups/20260702-123510-p83-worker-heartbeat` + - restarted Portal and worker + - production `.env` added `MEMIND_AGENT_RUN_WORKER_EXPECT_RUNNING=1` +- synthetic recovery canary: + - no-heartbeat run `042276ae-ecac-4b7c-9c53-c02154ea85cd` recovered to `failed` + - fresh-heartbeat run `97a2ac16-0234-483d-be4f-836ba95a91e7` was not recovered and was deleted after verification +- real worker canary: + - run `ad2894a7-dc02-49c6-a859-da6af95cb143` + - request `p83-heartbeat-real-20260702043642` + - status `succeeded` + - event chain includes `worker_heartbeat` before `tool_gateway_dispatch` + - heartbeat data `pid=31341`, `attempt=1`, `heartbeatMs=30000` +- post-check: + - runtime/status queue empty + - `runningWithoutHeartbeatCount=0` + - worker check `ok=true`, `expected=running` + - guard `shouldPause=false` + ## 下一步执行建议 ### P6.22 Task Artifact UX and Failure Messages - 用户侧显示 queued/running/validation/succeeded/failed。 - 失败原因分层展示。 +- 用户已决定暂时跳过,不进入下一步开发。 + +### P8.4 Long-running Gray Soak and Heartbeat SLO + +- 基于 heartbeat 字段做一轮 24h all-user gray soak。 +- SLO 日报增加 running heartbeat age 和 missing heartbeat count。 +- guard 阈值保守观察,不扩大并发。 diff --git a/docs/architecture/memind-2-streaming-agent-runtime-plan.md b/docs/architecture/memind-2-streaming-agent-runtime-plan.md index a19c610..2346055 100644 --- a/docs/architecture/memind-2-streaming-agent-runtime-plan.md +++ b/docs/architecture/memind-2-streaming-agent-runtime-plan.md @@ -29,6 +29,7 @@ - P6.19 Task-level Artifact Validation: 已完成第一步,H5 code-run 会在 receipt 之外追加可推断的 `public/*.html` 目标文件校验,页面编辑 code-run 会追加 page-edit task receipt,生产真实 canary 三项 expectedFiles 全部通过。 - P8.1 Queue Lease / Stuck Run Recovery: 已完成第一步,worker dispatch 前自动回收超时 running run,新增 `--recover-stale` dry-run/apply 运维入口,生产 synthetic stale run 验证通过。 - P8.2 Portal DB/Auth Transient Error Hardening: 已完成第一步,session attach、`/auth/status` 和 API auth middleware 已捕获 DB/auth 瞬时错误,生产部署后 live health、auth/status、runtime/status、guard 和 SLO 验证通过。 +- P8.3 Worker Lease Heartbeat: 已完成第一步,running run 执行期间写入 `worker_heartbeat` event,runtime/status、worker check、guard 和 stale recovery 均优先使用 heartbeat age,生产 synthetic 与真实 worker canary 验证通过。 - P5.15 Active Stream TTL Reconcile: 已按用户要求跳过,暂不做报表/定时 reconcile。 - P5 Worker Pool 运维化: 已完成第一步,Redis Router 支持 worker drain。 - P5.9 First-token Latency EWMA: 已完成,StreamController 会把首个 SSE chunk 延迟写入 Redis,SLO 报告已展示。 @@ -2929,6 +2930,91 @@ runtime/status: - `public/p619-task-validation-20260702042911.html` - `.memind/agent-runs/p619-task-validation-20260702042911-page-edit.json` +### 2026-07-02 P8.3 Worker Lease Heartbeat + +目标: + +- 将 running run 的 lease 从单纯 `started_at + timeout` 升级为 heartbeat-aware。 +- 避免长任务仍在执行时被 stale recovery 误杀。 +- 让 runtime/status、worker check 和 auto-pause guard 能看到 running heartbeat age。 + +改动: + +- `agent-run-gateway.mjs`: + - 新增 `MEMIND_AGENT_RUN_HEARTBEAT_MS`,默认 `30000`。 + - run 被 claim 为 `running` 后立即写入 `worker_heartbeat` event,后续按 heartbeat interval 定时写入。 + - run 进入 terminal 或失败后停止 heartbeat timer。 + - `recoverStaleRunningRuns()` 使用 `COALESCE(latest worker_heartbeat, started_at)` 判断 stale。 + - recovery UPDATE 也复核最新 heartbeat,避免 select/update 之间误杀刚刷新 heartbeat 的 run。 + - `getQueueStatus()` 增加: + - `heartbeatMs` + - `oldestRunningHeartbeatAt` + - `oldestRunningHeartbeatAgeMs` + - `runningWithoutHeartbeatCount` + - `latestRunningRun` +- `scripts/agent-run-guard.mjs`: + - auto-pause 的 running age 判断改为 heartbeat age。 + - queue 摘要增加 heartbeat 字段。 +- `scripts/check-agent-run-worker.mjs`: + - queue 摘要增加 heartbeat 字段。 +- `.env.example` 和 worker rollout runbook: + - 记录 `MEMIND_AGENT_RUN_HEARTBEAT_MS`。 + - all-user gray 下记录 `MEMIND_AGENT_RUN_WORKER_EXPECT_RUNNING=1`。 + +测试: + +- `node --check agent-run-gateway.mjs scripts/agent-run-worker.mjs scripts/agent-run-guard.mjs scripts/check-agent-run-worker.mjs` 通过。 +- `node --test agent-run-gateway.test.mjs agent-run-routes.test.mjs` 通过,31 tests pass。 +- 新增测试: + - worker dispatch 会写 `worker_heartbeat`。 + - stale recovery 不回收有新 heartbeat 的旧 running run。 + - queue status 展示 heartbeat age 与 missing heartbeat count。 + +生产部署: + +- 备份: + - `/Users/john/Project/memind_backups/20260702-123510-p83-worker-heartbeat` +- 已部署: + - bundled `server.mjs` + - bundled `scripts/agent-run-worker.mjs` + - `scripts/agent-run-guard.mjs` + - `scripts/check-agent-run-worker.mjs` +- 已重启: + - `cn.tkmind.memind-portal` + - `cn.tkmind.memind-agent-run-worker` +- 生产 `.env` 追加: + - `MEMIND_AGENT_RUN_WORKER_EXPECT_RUNNING=1` + +生产验证: + +- `https://mm.tkmind.cn/api/status` 返回 `ok`。 +- `/api/runtime/status.toolRuntime.queue`: + - `maxConcurrentRuns=2` + - `heartbeatMs=30000` + - `statusCounts={}` + - `oldestRunningHeartbeatAgeMs=0` + - `runningWithoutHeartbeatCount=0` +- `check-agent-run-worker.mjs`: + - `ok=true` + - `expected=running` + - `running=true` +- guard dry-run: + - `ok=true` + - `shouldPause=false` +- synthetic recovery canary: + - no-heartbeat run `042276ae-ecac-4b7c-9c53-c02154ea85cd` + - request `p83-no-heartbeat-20260702043604` + - recovery `considered=1`, `recovered=1` + - stale event includes `heartbeatAt=null` and `heartbeatAgeMs` + - fresh-heartbeat synthetic run `97a2ac16-0234-483d-be4f-836ba95a91e7` was not recovered and then deleted to avoid queue residue. +- real worker heartbeat canary: + - run `ad2894a7-dc02-49c6-a859-da6af95cb143` + - request `p83-heartbeat-real-20260702043642` + - status `succeeded` + - event chain includes `worker_heartbeat` before `tool_gateway_dispatch` + - heartbeat data includes `pid=31341`, `attempt=1`, `heartbeatMs=30000` + - `tool_gateway_validation` passed for receipt and `public/p83-heartbeat-real-20260702043642.html` + ## 回滚策略 - P0: 修改前保留 `server.mjs` 备份;如启动失败,恢复备份并 `launchctl kickstart` Portal。 diff --git a/scripts/agent-run-guard.mjs b/scripts/agent-run-guard.mjs index e6a568a..8d1d507 100755 --- a/scripts/agent-run-guard.mjs +++ b/scripts/agent-run-guard.mjs @@ -141,13 +141,43 @@ async function readQueueHealth(now) { : Number(oldestPendingRows[0].oldest_updated_at); const [oldestRunningRows] = await conn.query( - `SELECT MIN(started_at) AS oldest_started_at - FROM h5_agent_runs - WHERE status = 'running'`, + `SELECT + r.id, + r.request_id, + r.started_at, + r.updated_at, + h.latest_heartbeat_at + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' + ORDER BY COALESCE(h.latest_heartbeat_at, r.started_at) ASC + LIMIT 1`, ); - const oldestRunningStartedAt = oldestRunningRows[0]?.oldest_started_at == null + const oldestRunningStartedAt = oldestRunningRows[0]?.started_at == null ? null - : Number(oldestRunningRows[0].oldest_started_at); + : Number(oldestRunningRows[0].started_at); + const oldestRunningHeartbeatAt = oldestRunningRows[0]?.latest_heartbeat_at == null + ? null + : Number(oldestRunningRows[0].latest_heartbeat_at); + const oldestRunningHeartbeatAgeMs = oldestRunningRows[0] + ? Math.max(0, now - Number(oldestRunningRows[0].latest_heartbeat_at ?? oldestRunningRows[0].started_at ?? now)) + : 0; + const [runningWithoutHeartbeatRows] = await conn.query( + `SELECT COUNT(*) AS count + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' AND h.latest_heartbeat_at IS NULL`, + ); const failedWindowMs = positiveInt(process.env.MEMIND_AGENT_RUN_GUARD_FAILED_WINDOW_MS, 10 * 60 * 1000); const since = now - failedWindowMs; @@ -173,6 +203,17 @@ async function readQueueHealth(now) { oldestPendingAgeMs: oldestPendingUpdatedAt == null ? 0 : Math.max(0, now - oldestPendingUpdatedAt), oldestRunningStartedAt, oldestRunningAgeMs: oldestRunningStartedAt == null ? 0 : Math.max(0, now - oldestRunningStartedAt), + oldestRunningHeartbeatAt, + oldestRunningHeartbeatAgeMs, + runningWithoutHeartbeatCount: Number(runningWithoutHeartbeatRows[0]?.count ?? 0), + latestRunningRun: oldestRunningRows[0] ? { + id: oldestRunningRows[0].id, + requestId: oldestRunningRows[0].request_id, + startedAt: oldestRunningStartedAt, + updatedAt: oldestRunningRows[0].updated_at == null ? null : Number(oldestRunningRows[0].updated_at), + heartbeatAt: oldestRunningHeartbeatAt, + heartbeatAgeMs: oldestRunningHeartbeatAgeMs, + } : null, failedWindowMs, failedRecentCount: Number(failedRows[0]?.count ?? 0), latestFailedRun: latestFailedRows[0] ? { @@ -204,8 +245,8 @@ function evaluateHealth(queue) { if (queue.queuedOrRetryable >= thresholds.maxPendingCount) { reasons.push(`pending_count ${queue.queuedOrRetryable} >= ${thresholds.maxPendingCount}`); } - if (queue.oldestRunningAgeMs >= thresholds.maxRunningAgeMs) { - reasons.push(`oldest_running_age_ms ${queue.oldestRunningAgeMs} >= ${thresholds.maxRunningAgeMs}`); + if (queue.oldestRunningHeartbeatAgeMs >= thresholds.maxRunningAgeMs) { + reasons.push(`oldest_running_heartbeat_age_ms ${queue.oldestRunningHeartbeatAgeMs} >= ${thresholds.maxRunningAgeMs}`); } return { thresholds, reasons, shouldPause: reasons.length > 0 }; } diff --git a/scripts/check-agent-run-worker.mjs b/scripts/check-agent-run-worker.mjs index ec3220b..9860231 100644 --- a/scripts/check-agent-run-worker.mjs +++ b/scripts/check-agent-run-worker.mjs @@ -109,13 +109,43 @@ async function readQueueSummary() { : Number(lagRows[0].oldest_updated_at); const [runningRows] = await conn.query( - `SELECT MIN(started_at) AS oldest_started_at - FROM h5_agent_runs - WHERE status = 'running'`, + `SELECT + r.id, + r.request_id, + r.started_at, + r.updated_at, + h.latest_heartbeat_at + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' + ORDER BY COALESCE(h.latest_heartbeat_at, r.started_at) ASC + LIMIT 1`, ); - const oldestRunningStartedAt = runningRows[0]?.oldest_started_at == null + const oldestRunningStartedAt = runningRows[0]?.started_at == null ? null - : Number(runningRows[0].oldest_started_at); + : Number(runningRows[0].started_at); + const oldestRunningHeartbeatAt = runningRows[0]?.latest_heartbeat_at == null + ? null + : Number(runningRows[0].latest_heartbeat_at); + const oldestRunningHeartbeatAgeMs = runningRows[0] + ? Math.max(0, Date.now() - Number(runningRows[0].latest_heartbeat_at ?? runningRows[0].started_at ?? Date.now())) + : 0; + const [runningWithoutHeartbeatRows] = await conn.query( + `SELECT COUNT(*) AS count + FROM h5_agent_runs r + LEFT JOIN ( + SELECT run_id, MAX(created_at) AS latest_heartbeat_at + FROM h5_agent_run_events + WHERE event_type = 'worker_heartbeat' + GROUP BY run_id + ) h ON h.run_id = r.id + WHERE r.status = 'running' AND h.latest_heartbeat_at IS NULL`, + ); const [failedRows] = await conn.query( `SELECT id, request_id, error_message, updated_at @@ -131,6 +161,17 @@ async function readQueueSummary() { oldestPendingAgeMs: oldestUpdatedAt == null ? 0 : Math.max(0, Date.now() - oldestUpdatedAt), oldestRunningStartedAt, oldestRunningAgeMs: oldestRunningStartedAt == null ? 0 : Math.max(0, Date.now() - oldestRunningStartedAt), + oldestRunningHeartbeatAt, + oldestRunningHeartbeatAgeMs, + runningWithoutHeartbeatCount: Number(runningWithoutHeartbeatRows[0]?.count ?? 0), + latestRunningRun: runningRows[0] ? { + id: runningRows[0].id, + requestId: runningRows[0].request_id, + startedAt: oldestRunningStartedAt, + updatedAt: runningRows[0].updated_at == null ? null : Number(runningRows[0].updated_at), + heartbeatAt: oldestRunningHeartbeatAt, + heartbeatAgeMs: oldestRunningHeartbeatAgeMs, + } : null, latestFailedRun: failedRows[0] ? { id: failedRows[0].id, requestId: failedRows[0].request_id,