From 4420cca340602a9425d1418c155e5dd5ac538e1a Mon Sep 17 00:00:00 2001 From: John Date: Thu, 2 Jul 2026 07:04:50 +0800 Subject: [PATCH] chore: add runtime observability ops --- .runtime/portal/RUNBOOK.txt | 15 +-- .runtime/portal/package.json | 18 +++- .../portal/scripts/check-stream-runtime.mjs | 98 +++++++++++++++++ .../portal/scripts/runtime-worker-drain.mjs | 101 ++++++++++++++++++ .runtime/portal/server.mjs | 68 ++++++++---- .../memind-2-streaming-agent-runtime-plan.md | 27 +++++ scripts/build-portal-runtime.mjs | 25 ++++- scripts/check-stream-runtime.mjs | 98 +++++++++++++++++ scripts/runtime-worker-drain.mjs | 101 ++++++++++++++++++ server.mjs | 2 +- tkmind-proxy.mjs | 47 ++++++-- 11 files changed, 556 insertions(+), 44 deletions(-) create mode 100755 .runtime/portal/scripts/check-stream-runtime.mjs create mode 100755 .runtime/portal/scripts/runtime-worker-drain.mjs create mode 100755 scripts/check-stream-runtime.mjs create mode 100755 scripts/runtime-worker-drain.mjs diff --git a/.runtime/portal/RUNBOOK.txt b/.runtime/portal/RUNBOOK.txt index 3edce79..403e96b 100644 --- a/.runtime/portal/RUNBOOK.txt +++ b/.runtime/portal/RUNBOOK.txt @@ -28,10 +28,13 @@ Key runtime differences must stay in .env, not in the artifact: H5_USERS_ROOT / MINDSPACE_STORAGE_ROOT / MEMIND_SHARED_PUBLISH_ROOT Deployment and operations transport: - 103 / Studio fixed IP: 58.38.22.103 H5 public domain: mm.tkmind.cn - 2026-07-02 routing decision: - - Temporarily move H5 public base from m.tkmind.cn to mm.tkmind.cn. - - Future H5 public traffic must not depend on the 105 nginx -> 127.0.0.1:19081 -> reverse SSH tunnel -> Portal :8081 path. - - Keep legacy 105/tunnel scripts only for rollback or explicitly requested migration work. - Do not switch back to 10.10.* LAN paths unless explicitly required. + Current public path: mm.tkmind.cn -> local nginx -> Portal :8081 + Legacy rollback-only path: m.tkmind.cn -> 105 nginx -> reverse SSH tunnel -> Portal :8081 + Future H5 traffic must not depend on 105 forwarding unless explicitly rolling back. + +Streaming runtime operations: + node scripts/check-stream-runtime.mjs + node scripts/runtime-worker-drain.mjs status + node scripts/runtime-worker-drain.mjs drain goosed-3 + node scripts/runtime-worker-drain.mjs undrain goosed-3 diff --git a/.runtime/portal/package.json b/.runtime/portal/package.json index 5497b5a..efb75cb 100644 --- a/.runtime/portal/package.json +++ b/.runtime/portal/package.json @@ -1,5 +1,21 @@ { "name": "tkmind-h5-portal-runtime", "private": true, - "type": "module" + "type": "module", + "engines": { + "node": ">=22" + }, + "dependencies": { + "@node-rs/argon2": "^2.0.2", + "@resvg/resvg-js": "^2.6.2", + "debug": "^4.4.3", + "express": "^4.21.2", + "http-proxy-middleware": "^3.0.3", + "jsonrepair": "^3.14.0", + "mysql2": "^3.22.5", + "qrcode": "^1.5.4", + "redis": "^4.7.1", + "sharp": "^0.35.2", + "undici": "^6.26.0" + } } diff --git a/.runtime/portal/scripts/check-stream-runtime.mjs b/.runtime/portal/scripts/check-stream-runtime.mjs new file mode 100755 index 0000000..4cb9074 --- /dev/null +++ b/.runtime/portal/scripts/check-stream-runtime.mjs @@ -0,0 +1,98 @@ +#!/usr/bin/env node +import fs from 'node:fs'; +import path from 'node:path'; +import { Agent, fetch as undiciFetch } from 'undici'; + +function loadEnvFile(filePath) { + if (!fs.existsSync(filePath)) return; + for (const line of fs.readFileSync(filePath, 'utf8').split('\n')) { + const trimmed = line.trim(); + if (!trimmed || trimmed.startsWith('#')) continue; + const idx = trimmed.indexOf('='); + if (idx < 0) continue; + const key = trimmed.slice(0, idx).trim(); + const value = trimmed.slice(idx + 1).trim(); + if (!process.env[key]) process.env[key] = value; + } +} + +loadEnvFile(path.join(process.cwd(), '.env')); + +const insecureDispatcher = new Agent({ connect: { rejectUnauthorized: false } }); +const publicBase = (process.env.H5_PUBLIC_BASE_URL || 'https://mm.tkmind.cn').replace(/\/$/, ''); +const targets = (process.env.TKMIND_API_TARGETS || process.env.TKMIND_API_TARGET || '') + .split(',') + .map((value) => value.trim()) + .filter(Boolean); + +async function fetchText(url, init = {}) { + const res = await undiciFetch(url, { + ...init, + dispatcher: url.startsWith('https://127.0.0.1') ? insecureDispatcher : undefined, + }); + const text = await res.text(); + return { res, text }; +} + +async function checkJson(pathname) { + const { res, text } = await fetchText(`${publicBase}${pathname}`); + let json = null; + try { + json = JSON.parse(text); + } catch { + json = null; + } + return { ok: res.ok, status: res.status, json, text: json ? undefined : text.slice(0, 160) }; +} + +async function checkSseHeaders(pathname) { + const { res } = await fetchText(`${publicBase}${pathname}`, { + headers: { Accept: 'text/event-stream' }, + }); + return { + ok: res.status === 401 || res.ok, + status: res.status, + xAccelBuffering: res.headers.get('x-accel-buffering'), + contentType: res.headers.get('content-type'), + cacheControl: res.headers.get('cache-control'), + }; +} + +async function checkTarget(target) { + const { res, text } = await fetchText(new URL('/status', target).toString()); + return { target, ok: res.ok, status: res.status, text: text.slice(0, 80) }; +} + +const result = { + ok: true, + publicBase, + checkedAt: new Date().toISOString(), + status: await checkJson('/api/status').catch((err) => ({ ok: false, error: err.message })), + runtime: await checkJson('/api/runtime/status').catch((err) => ({ ok: false, error: err.message })), + sse: { + sessions: await checkSseHeaders('/api/sessions/check-stream-runtime/events').catch((err) => ({ + ok: false, + error: err.message, + })), + agentRuns: await checkSseHeaders('/api/agent/runs/check-stream-runtime/events').catch((err) => ({ + ok: false, + error: err.message, + })), + }, + targets: await Promise.all(targets.map((target) => checkTarget(target).catch((err) => ({ + target, + ok: false, + error: err.message, + })))), +}; + +result.ok = Boolean( + result.status.ok && + result.runtime.ok && + result.sse.sessions.ok && + result.sse.agentRuns.ok && + result.targets.every((target) => target.ok), +); + +console.log(JSON.stringify(result, null, 2)); +process.exit(result.ok ? 0 : 1); diff --git a/.runtime/portal/scripts/runtime-worker-drain.mjs b/.runtime/portal/scripts/runtime-worker-drain.mjs new file mode 100755 index 0000000..e9dbbff --- /dev/null +++ b/.runtime/portal/scripts/runtime-worker-drain.mjs @@ -0,0 +1,101 @@ +#!/usr/bin/env node +import fs from 'node:fs'; +import path from 'node:path'; +import { createClient } from 'redis'; + +function loadEnvFile(filePath) { + if (!fs.existsSync(filePath)) return; + for (const line of fs.readFileSync(filePath, 'utf8').split('\n')) { + const trimmed = line.trim(); + if (!trimmed || trimmed.startsWith('#')) continue; + const idx = trimmed.indexOf('='); + if (idx < 0) continue; + const key = trimmed.slice(0, idx).trim(); + const value = trimmed.slice(idx + 1).trim(); + if (!process.env[key]) process.env[key] = value; + } +} + +loadEnvFile(path.join(process.cwd(), '.env')); + +const redisUrl = process.env.MEMIND_RUNTIME_REDIS_URL || 'redis://127.0.0.1:6379/0'; +const namespace = process.env.MEMIND_RUNTIME_REDIS_NAMESPACE || 'memind:runtime'; +const configuredWorkers = (process.env.TKMIND_API_TARGETS || process.env.TKMIND_API_TARGET || '') + .split(',') + .map((value) => value.trim()) + .filter(Boolean) + .map((_, index) => `goosed-${index + 1}`); +const action = process.argv[2] || 'status'; +const workerId = process.argv[3] || null; + +function usage() { + console.error('Usage: node scripts/runtime-worker-drain.mjs [goosed-N]'); +} + +function workerKey(id, field) { + return [namespace, 'worker', id, field].join(':'); +} + +if (!['status', 'drain', 'undrain'].includes(action)) { + usage(); + process.exit(2); +} +if (['drain', 'undrain'].includes(action) && !workerId) { + usage(); + process.exit(2); +} + +const client = createClient({ url: redisUrl }); +client.on('error', (err) => { + console.error(`Redis error: ${err instanceof Error ? err.message : err}`); +}); +await client.connect(); + +if (action === 'drain') { + await client.set(workerKey(workerId, 'drain'), '1'); +} +if (action === 'undrain') { + await client.del(workerKey(workerId, 'drain')); +} + +const keys = await client.keys(workerKey('*', 'active_streams')); +const workers = [ + ...configuredWorkers, + ...keys + .map((key) => key.split(':').at(-2)) + .filter(Boolean), +]; +if (workerId && !workers.includes(workerId)) workers.push(workerId); + +const rows = []; +for (const id of [...new Set(workers)].sort()) { + const values = await client.mGet([ + workerKey(id, 'active_streams'), + workerKey(id, 'drain'), + workerKey(id, 'stream_open_count'), + workerKey(id, 'stream_abort_count'), + workerKey(id, 'stream_error_count'), + workerKey(id, 'last_stream_started_at'), + workerKey(id, 'last_stream_ended_at'), + ]); + rows.push({ + id, + activeStreams: Number(values[0] || 0), + drain: /^(1|true|yes)$/i.test(String(values[1] || '')), + streamOpenCount: Number(values[2] || 0), + streamAbortCount: Number(values[3] || 0), + streamErrorCount: Number(values[4] || 0), + lastStreamStartedAt: values[5] ? Number(values[5]) : null, + lastStreamEndedAt: values[6] ? Number(values[6]) : null, + }); +} + +await client.quit(); + +console.log(JSON.stringify({ + ok: true, + action, + workerId, + namespace, + workers: rows, +}, null, 2)); diff --git a/.runtime/portal/server.mjs b/.runtime/portal/server.mjs index 1c34779..2f7a22a 100644 --- a/.runtime/portal/server.mjs +++ b/.runtime/portal/server.mjs @@ -8620,6 +8620,7 @@ function createRuntimeRouter({ if (/^(1|true|yes)$/i.test(String(values[5] ?? ""))) return Number.POSITIVE_INFINITY; return readNumber(values[0]) * 3 + readNumber(values[1]) + readNumber(values[2]) * 0.01 + readNumber(values[3]) * 5 + readNumber(values[4]) * 2; }; + const workerScoreFromValues = (values = []) => readNumber(values[0]) * 3 + readNumber(values[1]) + readNumber(values[2]) * 0.01 + readNumber(values[3]) * 5 + readNumber(values[4]) * 2; return { async pickTarget(fallbackPick, orderedTargets = targets) { const client = await getClient(); @@ -8650,20 +8651,24 @@ function createRuntimeRouter({ if (!client || !target) return; const workerId = workerIdForTarget(target); const streamKey = sessionId ? key("stream", sessionId, "status") : null; - const multi = client.multi().incr(key("worker", workerId, "active_streams")).set(key("worker", workerId, "heartbeat"), String(Date.now()), { EX: 30 }); + const now = String(Date.now()); + const multi = client.multi().incr(key("worker", workerId, "active_streams")).incr(key("worker", workerId, "stream_open_count")).set(key("worker", workerId, "last_stream_started_at"), now).set(key("worker", workerId, "heartbeat"), now, { EX: 30 }); if (streamKey) multi.set(streamKey, "active", { EX: 60 * 60 }); await multi.exec().catch(() => null); }, - async streamEnded(sessionId, target) { + async streamEnded(sessionId, target, { status = "closed" } = {}) { const client = await getClient(); if (!client || !target) return; const workerId = workerIdForTarget(target); const streamKey = sessionId ? key("stream", sessionId, "status") : null; const activeKey = key("worker", workerId, "active_streams"); + const now = String(Date.now()); const nextValue = await client.decr(activeKey).catch(() => null); - const multi = client.multi().set(key("worker", workerId, "heartbeat"), String(Date.now()), { EX: 30 }); + const multi = client.multi().set(key("worker", workerId, "last_stream_ended_at"), now).set(key("worker", workerId, "heartbeat"), now, { EX: 30 }); if (Number(nextValue ?? 0) < 0) multi.set(activeKey, "0"); - if (streamKey) multi.set(streamKey, "closed", { EX: 600 }); + if (status === "aborted") multi.incr(key("worker", workerId, "stream_abort_count")); + if (status === "error") multi.incr(key("worker", workerId, "stream_error_count")); + if (streamKey) multi.set(streamKey, status, { EX: 600 }); await multi.exec().catch(() => null); }, async getStatus() { @@ -8678,7 +8683,12 @@ function createRuntimeRouter({ key("worker", workerId, "error_rate"), key("worker", workerId, "memory_pressure"), key("worker", workerId, "heartbeat"), - key("worker", workerId, "drain") + key("worker", workerId, "drain"), + key("worker", workerId, "stream_open_count"), + key("worker", workerId, "stream_abort_count"), + key("worker", workerId, "stream_error_count"), + key("worker", workerId, "last_stream_started_at"), + key("worker", workerId, "last_stream_ended_at") ]).catch(() => []) : []; workers.push({ id: workerId, @@ -8689,7 +8699,13 @@ function createRuntimeRouter({ errorRate: readNumber(values?.[3]), memoryPressure: readNumber(values?.[4]), heartbeat: values?.[5] ? Number(values[5]) : null, - drain: /^(1|true|yes)$/i.test(String(values?.[6] ?? "")) + drain: /^(1|true|yes)$/i.test(String(values?.[6] ?? "")), + streamOpenCount: readNumber(values?.[7]), + streamAbortCount: readNumber(values?.[8]), + streamErrorCount: readNumber(values?.[9]), + lastStreamStartedAt: values?.[10] ? Number(values[10]) : null, + lastStreamEndedAt: values?.[11] ? Number(values[11]) : null, + score: workerScoreFromValues(values) }); } return { @@ -9215,9 +9231,9 @@ function createTkmindProxy({ if (!runtimeRouter || !sessionId || !target) return; await runtimeRouter.streamStarted(sessionId, target).catch(() => null); } - async function markStreamEnded(sessionId, target) { + async function markStreamEnded(sessionId, target, options = {}) { if (!runtimeRouter || !sessionId || !target) return; - await runtimeRouter.streamEnded(sessionId, target).catch(() => null); + await runtimeRouter.streamEnded(sessionId, target, options).catch(() => null); } async function getRuntimeStatus() { const targetStatuses = []; @@ -9262,8 +9278,8 @@ function createTkmindProxy({ if (!session?.id) { throw new Error("\u521B\u5EFA\u4F1A\u8BDD\u5931\u8D25\uFF1A\u7F3A\u5C11 session id"); } - await userAuth2.registerAgentSession(userId, session.id, startTarget); await rememberSessionTarget(session.id, startTarget); + await userAuth2.registerAgentSession(userId, session.id, startTarget); if (resolvedSessionPolicy?.gooseMode) { const modeRes = await apiFetch(startTarget, apiSecret, "/agent/update_session", { method: "POST", @@ -9302,7 +9318,7 @@ function createTkmindProxy({ if (targets.length <= 1 || !sessionId) return primaryTarget; try { const routedTarget = await runtimeRouter?.resolveSessionTarget(sessionId); - if (routedTarget && targets.includes(routedTarget)) return routedTarget; + if (routedTarget) return routedTarget; const { target, node } = await userAuth2.getSessionTarget(sessionId); if (target && targets.includes(target)) return target; return targets[node] ?? primaryTarget; @@ -9493,12 +9509,12 @@ function createTkmindProxy({ } const session = JSON.parse(text); if (session?.id) { + await rememberSessionTarget(session.id, startTarget); await userAuth2.registerAgentSession( req.currentUser.id, session.id, startTarget ); - await rememberSessionTarget(session.id, startTarget); if (sessionPolicy.gooseMode) { const modeRes = await apiFetch(startTarget, apiSecret, "/agent/update_session", { method: "POST", @@ -9801,12 +9817,18 @@ function createTkmindProxy({ }, 2e4); res.on("drain", () => source.resume()); await markStreamStarted(sessionId, sessionTarget); + let streamCloseStatus = "closed"; try { - await pipeline(source, linkSanitizer, billingTransform, clientSink); + try { + await pipeline(source, linkSanitizer, billingTransform, clientSink); + } catch (err) { + streamCloseStatus = clientClosed || upstreamAbort.signal.aborted ? "aborted" : "error"; + throw err; + } } finally { clearInterval(keepalive); req.off("close", abortUpstream); - await markStreamEnded(sessionId, sessionTarget); + await markStreamEnded(sessionId, sessionTarget, { status: streamCloseStatus }); if (!res.writableEnded) res.end(); } } catch (err) { @@ -9949,9 +9971,9 @@ function createTkmindProxy({ sessionScoped, proxyFallback, proxySessionEvents, - getRuntimeStatus, resolveTarget, startSessionForUser, + getRuntimeStatus, submitSessionReplyForUser, apiFetch: async (pathname, init) => apiFetch(await pickTarget(), apiSecret, pathname, init), apiFetchTo: (target, pathname, init) => apiFetch(target, apiSecret, pathname, init) @@ -14764,7 +14786,7 @@ function buildOutputProfile(mimeType) { mimeType: "image/png", baseQuality: 90, minQuality: 70, - encode: (pipeline, quality) => pipeline.png({ + encode: (pipeline2, quality) => pipeline2.png({ compressionLevel: 9, adaptiveFiltering: true, palette: true, @@ -14779,7 +14801,7 @@ function buildOutputProfile(mimeType) { mimeType: "image/webp", baseQuality: 84, minQuality: WEBP_MIN_QUALITY, - encode: (pipeline, quality) => pipeline.webp({ + encode: (pipeline2, quality) => pipeline2.webp({ quality, alphaQuality: Math.min(100, quality + 8), effort: 5 @@ -14791,7 +14813,7 @@ function buildOutputProfile(mimeType) { mimeType: "image/jpeg", baseQuality: 84, minQuality: JPEG_MIN_QUALITY, - encode: (pipeline, quality) => pipeline.jpeg({ + encode: (pipeline2, quality) => pipeline2.jpeg({ quality, mozjpeg: true, chromaSubsampling: "4:4:4" @@ -14803,7 +14825,7 @@ function deriveOutputFilename(filename, extension) { return `${raw}${extension}`; } async function renderVariant(buffer, metadata, profile, { width, height, quality, maxPixels }) { - const pipeline = sharp(buffer, { + const pipeline2 = sharp(buffer, { failOn: "error", limitInputPixels: maxPixels, sequentialRead: true @@ -14813,7 +14835,7 @@ async function renderVariant(buffer, metadata, profile, { width, height, quality fit: "inside", withoutEnlargement: true }); - const encoded = await profile.encode(pipeline, quality).toBuffer(); + const encoded = await profile.encode(pipeline2, quality).toBuffer(); return encoded; } async function normalizeImageForStorage({ @@ -35854,20 +35876,20 @@ api.get("/status", async (_req, res, next) => { }); api.get("/runtime/status", async (_req, res) => { await userAuthReady; - if (!tkmindProxy) { - return res.status(503).json({ ok: false, message: "\u4F1A\u8BDD\u4EE3\u7406\u5C1A\u672A\u5C31\u7EEA" }); + if (!tkmindProxy?.getRuntimeStatus) { + return res.status(503).json({ ok: false, message: "runtime router unavailable" }); } try { const status = await tkmindProxy.getRuntimeStatus(); return res.json({ ok: true, - at: Date.now(), + timestamp: (/* @__PURE__ */ new Date()).toISOString(), ...status }); } catch (err) { return res.status(502).json({ ok: false, - message: err instanceof Error ? err.message : "\u8BFB\u53D6 runtime \u72B6\u6001\u5931\u8D25" + message: err instanceof Error ? err.message : "runtime status failed" }); } }); diff --git a/docs/architecture/memind-2-streaming-agent-runtime-plan.md b/docs/architecture/memind-2-streaming-agent-runtime-plan.md index 988cf74..0ddd67e 100644 --- a/docs/architecture/memind-2-streaming-agent-runtime-plan.md +++ b/docs/architecture/memind-2-streaming-agent-runtime-plan.md @@ -11,6 +11,7 @@ - P4 Tool Gateway v1: 已完成第一步,普通用户默认不再暴露 Aider/OpenHands,显式用户白名单保留。 - P5 Worker Pool 运维化: 已完成第一步,Redis Router 支持 worker drain。 - 生产同步分支: 已从远程 `origin/main` 新建干净副本和分支 `memind-streaming-runtime-20260702`,用于远程开发机后续直接拉取。 +- P3.5/P5.5 生产化补强: 进行中,先完成生产数据库和 `MindSpace` 备份,再增加 runtime metrics、健康检查脚本和 drain 运维脚本。 ## 目标 @@ -431,6 +432,32 @@ docker exec memind-runtime-redis redis-cli DEL memind:runtime:worker:goosed-3:dr - 不提交生产 `.env`、数据库、`MindSpace/`、`data/`、`users/`、`.tailscale/`、`logs/`、证书和密钥。 - 只提交代码、配置模板、部署样例和架构/运维文档。 +### 2026-07-02 P3.5/P5.5 生产数据保护与运维补强 + +生产数据保护: + +- 允许重启生产服务,但重启前必须保护好数据库和 `/Users/john/Project/Memind/MindSpace`。 +- 备份目录: `/Users/john/Project/memind_backups/20260702-065813-pre-p35-p55`。 +- `MindSpace` 使用 `rsync -a` 全量副本,文件数 `2137`,大小约 `120M`。 +- 数据库为 Aliyun RDS MySQL `goose`,已导出: + - `mysql-schema.sql` + - `mysql-manifest.json` + - `mysql-jsonl/*.jsonl` + - 表数 `85`,行数 `28296` + +改造内容: + +- Redis Router worker 状态增加: + - `stream_open_count` + - `stream_abort_count` + - `stream_error_count` + - `last_stream_started_at` + - `last_stream_ended_at` + - `score` +- 新增只读检查脚本 `scripts/check-stream-runtime.mjs`。 +- 新增 drain 运维脚本 `scripts/runtime-worker-drain.mjs`。 +- runtime 构建模板同步上述脚本,并将 RUNBOOK 中主路径更新为 `mm.tkmind.cn -> local nginx -> Portal :8081`。 + ## 回滚策略 - P0: 修改前保留 `server.mjs` 备份;如启动失败,恢复备份并 `launchctl kickstart` Portal。 diff --git a/scripts/build-portal-runtime.mjs b/scripts/build-portal-runtime.mjs index 6e036a1..3347872 100755 --- a/scripts/build-portal-runtime.mjs +++ b/scripts/build-portal-runtime.mjs @@ -267,6 +267,14 @@ async function writeMetadata() { path.join(root, 'scripts', 'wechat-mp-menu.mjs'), path.join(runtimeRoot, 'scripts', 'wechat-mp-menu.mjs'), ); + await fs.copyFile( + path.join(root, 'scripts', 'check-stream-runtime.mjs'), + path.join(runtimeRoot, 'scripts', 'check-stream-runtime.mjs'), + ); + await fs.copyFile( + path.join(root, 'scripts', 'runtime-worker-drain.mjs'), + path.join(runtimeRoot, 'scripts', 'runtime-worker-drain.mjs'), + ); await writeFile( path.join(runtimeRoot, 'RUNBOOK.txt'), [ @@ -300,11 +308,16 @@ async function writeMetadata() { ' H5_USERS_ROOT / MINDSPACE_STORAGE_ROOT / MEMIND_SHARED_PUBLISH_ROOT', '', 'Deployment and operations transport:', - ' 105 fixed IP: 120.26.184.105', - ' 103 / Studio fixed IP: 58.38.22.103', - ' m.tkmind.cn: 105 nginx -> 127.0.0.1:19081 -> reverse SSH tunnel -> Portal :8081', - ' scripts/memind-portal-tunnel.sh must stay in runtime; release restarts cn.tkmind.memind-portal-tunnel', - ' Do not switch back to 10.10.* LAN paths unless explicitly required.', + ' H5 public domain: mm.tkmind.cn', + ' Current public path: mm.tkmind.cn -> local nginx -> Portal :8081', + ' Legacy rollback-only path: m.tkmind.cn -> 105 nginx -> reverse SSH tunnel -> Portal :8081', + ' Future H5 traffic must not depend on 105 forwarding unless explicitly rolling back.', + '', + 'Streaming runtime operations:', + ' node scripts/check-stream-runtime.mjs', + ' node scripts/runtime-worker-drain.mjs status', + ' node scripts/runtime-worker-drain.mjs drain goosed-3', + ' node scripts/runtime-worker-drain.mjs undrain goosed-3', '', ].join('\n'), ); @@ -322,6 +335,8 @@ async function main() { await writeMetadata(); await fs.chmod(path.join(runtimeRoot, 'scripts', 'run-memind-portal-prod.sh'), 0o755); await fs.chmod(path.join(runtimeRoot, 'scripts', 'wechat-mp-menu.mjs'), 0o755); + await fs.chmod(path.join(runtimeRoot, 'scripts', 'check-stream-runtime.mjs'), 0o755); + await fs.chmod(path.join(runtimeRoot, 'scripts', 'runtime-worker-drain.mjs'), 0o755); await fs.chmod(path.join(runtimeRoot, 'scripts', 'memind-portal-tunnel.sh'), 0o755); console.log(''); console.log(`Portal runtime 已生成: ${runtimeRoot}`); diff --git a/scripts/check-stream-runtime.mjs b/scripts/check-stream-runtime.mjs new file mode 100755 index 0000000..4cb9074 --- /dev/null +++ b/scripts/check-stream-runtime.mjs @@ -0,0 +1,98 @@ +#!/usr/bin/env node +import fs from 'node:fs'; +import path from 'node:path'; +import { Agent, fetch as undiciFetch } from 'undici'; + +function loadEnvFile(filePath) { + if (!fs.existsSync(filePath)) return; + for (const line of fs.readFileSync(filePath, 'utf8').split('\n')) { + const trimmed = line.trim(); + if (!trimmed || trimmed.startsWith('#')) continue; + const idx = trimmed.indexOf('='); + if (idx < 0) continue; + const key = trimmed.slice(0, idx).trim(); + const value = trimmed.slice(idx + 1).trim(); + if (!process.env[key]) process.env[key] = value; + } +} + +loadEnvFile(path.join(process.cwd(), '.env')); + +const insecureDispatcher = new Agent({ connect: { rejectUnauthorized: false } }); +const publicBase = (process.env.H5_PUBLIC_BASE_URL || 'https://mm.tkmind.cn').replace(/\/$/, ''); +const targets = (process.env.TKMIND_API_TARGETS || process.env.TKMIND_API_TARGET || '') + .split(',') + .map((value) => value.trim()) + .filter(Boolean); + +async function fetchText(url, init = {}) { + const res = await undiciFetch(url, { + ...init, + dispatcher: url.startsWith('https://127.0.0.1') ? insecureDispatcher : undefined, + }); + const text = await res.text(); + return { res, text }; +} + +async function checkJson(pathname) { + const { res, text } = await fetchText(`${publicBase}${pathname}`); + let json = null; + try { + json = JSON.parse(text); + } catch { + json = null; + } + return { ok: res.ok, status: res.status, json, text: json ? undefined : text.slice(0, 160) }; +} + +async function checkSseHeaders(pathname) { + const { res } = await fetchText(`${publicBase}${pathname}`, { + headers: { Accept: 'text/event-stream' }, + }); + return { + ok: res.status === 401 || res.ok, + status: res.status, + xAccelBuffering: res.headers.get('x-accel-buffering'), + contentType: res.headers.get('content-type'), + cacheControl: res.headers.get('cache-control'), + }; +} + +async function checkTarget(target) { + const { res, text } = await fetchText(new URL('/status', target).toString()); + return { target, ok: res.ok, status: res.status, text: text.slice(0, 80) }; +} + +const result = { + ok: true, + publicBase, + checkedAt: new Date().toISOString(), + status: await checkJson('/api/status').catch((err) => ({ ok: false, error: err.message })), + runtime: await checkJson('/api/runtime/status').catch((err) => ({ ok: false, error: err.message })), + sse: { + sessions: await checkSseHeaders('/api/sessions/check-stream-runtime/events').catch((err) => ({ + ok: false, + error: err.message, + })), + agentRuns: await checkSseHeaders('/api/agent/runs/check-stream-runtime/events').catch((err) => ({ + ok: false, + error: err.message, + })), + }, + targets: await Promise.all(targets.map((target) => checkTarget(target).catch((err) => ({ + target, + ok: false, + error: err.message, + })))), +}; + +result.ok = Boolean( + result.status.ok && + result.runtime.ok && + result.sse.sessions.ok && + result.sse.agentRuns.ok && + result.targets.every((target) => target.ok), +); + +console.log(JSON.stringify(result, null, 2)); +process.exit(result.ok ? 0 : 1); diff --git a/scripts/runtime-worker-drain.mjs b/scripts/runtime-worker-drain.mjs new file mode 100755 index 0000000..e9dbbff --- /dev/null +++ b/scripts/runtime-worker-drain.mjs @@ -0,0 +1,101 @@ +#!/usr/bin/env node +import fs from 'node:fs'; +import path from 'node:path'; +import { createClient } from 'redis'; + +function loadEnvFile(filePath) { + if (!fs.existsSync(filePath)) return; + for (const line of fs.readFileSync(filePath, 'utf8').split('\n')) { + const trimmed = line.trim(); + if (!trimmed || trimmed.startsWith('#')) continue; + const idx = trimmed.indexOf('='); + if (idx < 0) continue; + const key = trimmed.slice(0, idx).trim(); + const value = trimmed.slice(idx + 1).trim(); + if (!process.env[key]) process.env[key] = value; + } +} + +loadEnvFile(path.join(process.cwd(), '.env')); + +const redisUrl = process.env.MEMIND_RUNTIME_REDIS_URL || 'redis://127.0.0.1:6379/0'; +const namespace = process.env.MEMIND_RUNTIME_REDIS_NAMESPACE || 'memind:runtime'; +const configuredWorkers = (process.env.TKMIND_API_TARGETS || process.env.TKMIND_API_TARGET || '') + .split(',') + .map((value) => value.trim()) + .filter(Boolean) + .map((_, index) => `goosed-${index + 1}`); +const action = process.argv[2] || 'status'; +const workerId = process.argv[3] || null; + +function usage() { + console.error('Usage: node scripts/runtime-worker-drain.mjs [goosed-N]'); +} + +function workerKey(id, field) { + return [namespace, 'worker', id, field].join(':'); +} + +if (!['status', 'drain', 'undrain'].includes(action)) { + usage(); + process.exit(2); +} +if (['drain', 'undrain'].includes(action) && !workerId) { + usage(); + process.exit(2); +} + +const client = createClient({ url: redisUrl }); +client.on('error', (err) => { + console.error(`Redis error: ${err instanceof Error ? err.message : err}`); +}); +await client.connect(); + +if (action === 'drain') { + await client.set(workerKey(workerId, 'drain'), '1'); +} +if (action === 'undrain') { + await client.del(workerKey(workerId, 'drain')); +} + +const keys = await client.keys(workerKey('*', 'active_streams')); +const workers = [ + ...configuredWorkers, + ...keys + .map((key) => key.split(':').at(-2)) + .filter(Boolean), +]; +if (workerId && !workers.includes(workerId)) workers.push(workerId); + +const rows = []; +for (const id of [...new Set(workers)].sort()) { + const values = await client.mGet([ + workerKey(id, 'active_streams'), + workerKey(id, 'drain'), + workerKey(id, 'stream_open_count'), + workerKey(id, 'stream_abort_count'), + workerKey(id, 'stream_error_count'), + workerKey(id, 'last_stream_started_at'), + workerKey(id, 'last_stream_ended_at'), + ]); + rows.push({ + id, + activeStreams: Number(values[0] || 0), + drain: /^(1|true|yes)$/i.test(String(values[1] || '')), + streamOpenCount: Number(values[2] || 0), + streamAbortCount: Number(values[3] || 0), + streamErrorCount: Number(values[4] || 0), + lastStreamStartedAt: values[5] ? Number(values[5]) : null, + lastStreamEndedAt: values[6] ? Number(values[6]) : null, + }); +} + +await client.quit(); + +console.log(JSON.stringify({ + ok: true, + action, + workerId, + namespace, + workers: rows, +}, null, 2)); diff --git a/server.mjs b/server.mjs index 59fc778..e3c3fec 100644 --- a/server.mjs +++ b/server.mjs @@ -1721,7 +1721,7 @@ api.use(jsonUnlessMultipart); api.use(async (req, res, next) => { await userAuthReady; - if (req.path === '/status') return next(); + if (req.path === '/status' || req.path === '/runtime/status') return next(); if (req.path.startsWith('/internal/agent/')) return next(); if (req.path === '/agent/mindspace_page_patch') return next(); if (req.path === '/agent/mindspace_asset_delete') return next(); diff --git a/tkmind-proxy.mjs b/tkmind-proxy.mjs index 24f638d..83727db 100644 --- a/tkmind-proxy.mjs +++ b/tkmind-proxy.mjs @@ -143,6 +143,13 @@ function createRuntimeRouter({ readNumber(values[4]) * 2 ); }; + const workerScoreFromValues = (values = []) => ( + readNumber(values[0]) * 3 + + readNumber(values[1]) + + readNumber(values[2]) * 0.01 + + readNumber(values[3]) * 5 + + readNumber(values[4]) * 2 + ); return { async pickTarget(fallbackPick, orderedTargets = targets) { const client = await getClient(); @@ -178,25 +185,32 @@ function createRuntimeRouter({ if (!client || !target) return; const workerId = workerIdForTarget(target); const streamKey = sessionId ? key('stream', sessionId, 'status') : null; + const now = String(Date.now()); const multi = client .multi() .incr(key('worker', workerId, 'active_streams')) - .set(key('worker', workerId, 'heartbeat'), String(Date.now()), { EX: 30 }); + .incr(key('worker', workerId, 'stream_open_count')) + .set(key('worker', workerId, 'last_stream_started_at'), now) + .set(key('worker', workerId, 'heartbeat'), now, { EX: 30 }); if (streamKey) multi.set(streamKey, 'active', { EX: 60 * 60 }); await multi.exec().catch(() => null); }, - async streamEnded(sessionId, target) { + async streamEnded(sessionId, target, { status = 'closed' } = {}) { const client = await getClient(); if (!client || !target) return; const workerId = workerIdForTarget(target); const streamKey = sessionId ? key('stream', sessionId, 'status') : null; const activeKey = key('worker', workerId, 'active_streams'); + const now = String(Date.now()); const nextValue = await client.decr(activeKey).catch(() => null); const multi = client .multi() - .set(key('worker', workerId, 'heartbeat'), String(Date.now()), { EX: 30 }); + .set(key('worker', workerId, 'last_stream_ended_at'), now) + .set(key('worker', workerId, 'heartbeat'), now, { EX: 30 }); if (Number(nextValue ?? 0) < 0) multi.set(activeKey, '0'); - if (streamKey) multi.set(streamKey, 'closed', { EX: 600 }); + if (status === 'aborted') multi.incr(key('worker', workerId, 'stream_abort_count')); + if (status === 'error') multi.incr(key('worker', workerId, 'stream_error_count')); + if (streamKey) multi.set(streamKey, status, { EX: 600 }); await multi.exec().catch(() => null); }, async getStatus() { @@ -214,6 +228,11 @@ function createRuntimeRouter({ key('worker', workerId, 'memory_pressure'), key('worker', workerId, 'heartbeat'), key('worker', workerId, 'drain'), + key('worker', workerId, 'stream_open_count'), + key('worker', workerId, 'stream_abort_count'), + key('worker', workerId, 'stream_error_count'), + key('worker', workerId, 'last_stream_started_at'), + key('worker', workerId, 'last_stream_ended_at'), ]) .catch(() => []) : []; @@ -227,6 +246,12 @@ function createRuntimeRouter({ memoryPressure: readNumber(values?.[4]), heartbeat: values?.[5] ? Number(values[5]) : null, drain: /^(1|true|yes)$/i.test(String(values?.[6] ?? '')), + streamOpenCount: readNumber(values?.[7]), + streamAbortCount: readNumber(values?.[8]), + streamErrorCount: readNumber(values?.[9]), + lastStreamStartedAt: values?.[10] ? Number(values[10]) : null, + lastStreamEndedAt: values?.[11] ? Number(values[11]) : null, + score: workerScoreFromValues(values), }); } return { @@ -861,9 +886,9 @@ export function createTkmindProxy({ await runtimeRouter.streamStarted(sessionId, target).catch(() => null); } - async function markStreamEnded(sessionId, target) { + async function markStreamEnded(sessionId, target, options = {}) { if (!runtimeRouter || !sessionId || !target) return; - await runtimeRouter.streamEnded(sessionId, target).catch(() => null); + await runtimeRouter.streamEnded(sessionId, target, options).catch(() => null); } async function getRuntimeStatus() { @@ -1526,12 +1551,18 @@ export function createTkmindProxy({ }, 20000); res.on('drain', () => source.resume()); await markStreamStarted(sessionId, sessionTarget); + let streamCloseStatus = 'closed'; try { - await pipeline(source, linkSanitizer, billingTransform, clientSink); + try { + await pipeline(source, linkSanitizer, billingTransform, clientSink); + } catch (err) { + streamCloseStatus = clientClosed || upstreamAbort.signal.aborted ? 'aborted' : 'error'; + throw err; + } } finally { clearInterval(keepalive); req.off('close', abortUpstream); - await markStreamEnded(sessionId, sessionTarget); + await markStreamEnded(sessionId, sessionTarget, { status: streamCloseStatus }); if (!res.writableEnded) res.end(); } } catch (err) {