#!/usr/bin/env node /** * dshlog.mjs — DSH 平台「多节点日志采集 / 实时巡检 / 事后溯源」单一入口(方案 C:零新增常驻服务) * * 设计边界(2026-09-18 定) * · 服务器侧**零安装**:远端只用系统自带的 `journalctl`,采集靠 ssh 一次性拉取,不装 agent、不开监听口 * · 归档在**本机**:`E:/dsh-logs//.ndjson.gz`(可按天清理) * · 日志原文**保真**:`msg` 字段逐字节原样落盘(本项目大量判据依赖日志原文的行数/字节数) * · 失败**可分**:远端元信息走 stderr 哨兵 `__DSHLOG_EOF__`,与 stdout 的纯数据分离 * ⇒ "拉取失败" 与 "确无日志" 永不混淆(这是本项目反复踩过的坑) * * 子命令 * collect 拉取归档(增量,默认只拉「上次游标之后」) * ls 列出归档分片 * q 关键词/正则查询(跨机跨单元) * timeline 时间线重建(多机日志按时间戳合并,用于事后溯源) * watch 实时巡检(规则命中 → 判据表 + 非零退出码,可挂 automation) * stats 归档统计 * * 约定 * · 时间一律走 epoch(`--since @`;⛔ journalctl 不吃 date -Is 的时区偏移) * · 归档根可用 `DSHLOG_ROOT` 覆盖,默认 `E:/dsh-logs` */ import { spawnSync } from "node:child_process"; import { appendFileSync, existsSync, mkdirSync, readdirSync, readFileSync, statSync, unlinkSync, writeFileSync, } from "node:fs"; import { createGunzip, gzipSync } from "node:zlib"; import { createReadStream } from "node:fs"; import { createInterface } from "node:readline"; import path from "node:path"; // ───────────────────────────── 配置 ───────────────────────────── const ROOT = process.env.DSHLOG_ROOT || "E:/dsh-logs"; const STATE = path.join(ROOT, "state.json"); const HOSTS_FILE = path.join(ROOT, "hosts.json"); // 主机表从配置读取(⛔ 不在代码里写死地址 / ssh 别名): // DSHLOG_HOSTS='{"mgr":{"ssh":["-p","22","root@10.0.0.1"],"label":"manager"}}' // 未配置 ⇒ 空表,`hosts` 子命令会提示先配置。 const DEFAULT_HOSTS = (() => { const raw = process.env.DSHLOG_HOSTS ?? ""; if (raw.trim() === "") return {}; try { return JSON.parse(raw); } catch (e) { console.error("DSHLOG_HOSTS 不是合法 JSON:", String(e)); return {}; } })(); // 巡检规则:只认**指向本项目自身故障**的信号。 // ⚠️ 每条规则都是误报与漏报的取舍 —— 下面两组是实测调过的: // · WATCH-01 不能只写 `fatal`:sshd 的 `ssh_dispatch_run_fatal`(客户端网络断)会天天命中 // · WATCH-08 不能写 `Stopped .*`:实例正常退出、用户主动停会话都会打 `Stopped /usr/bin/bwrap …` const SEV_RULES = [ ["WATCH-01", "进程级致命", /PANIC|FATAL ERROR|unhandledRejection|uncaughtException|SIGSEGV|core dumped|segfault/i], ["WATCH-02", "内存/被杀", /out of memory|oom-kill|oom_reaper|Killed process|ENOMEM/i], ["WATCH-03", "端口/连接失败", /EADDRINUSE|ECONNREFUSED|ECONNRESET|ETIMEDOUT|EHOSTUNREACH|EPIPE/i], ["WATCH-04", "权限/属主异常", /EACCES|EPERM|permission denied/i], ["WATCH-05", "磁盘/写入失败", /ENOSPC|no space left|read-only file system|EROFS/i], ["WATCH-06", "覆盖网络信任链被拒", /no-trusted-keys|取目录全部失败|信任链被拒|reject.{0,12}trust/i], ["WATCH-07", "HTTP 5xx", /"statusCode":5\d\d|status=5\d\d/i], ["WATCH-08", "服务异常终止", /Failed with result|start request repeated|Main process exited, code=(exited|killed)|crash|panic exit/i], ]; // 登录/编排类噪声单元:命中规则也不计(否则 sshd 连接超时会天天冒充"故障") const NOISE_UNITS = /^sshd\.service$|^crond\.service$|^systemd-logind\.service$/; // ───────────────────────────── 小工具 ───────────────────────────── const nowSec = () => Math.floor(Date.now() / 1000); const pad = (n) => String(n).padStart(2, "0"); const dayStr = (d = new Date()) => `${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())}`; const tsLocal = (ms) => { const d = new Date(ms); return `${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())} ${pad(d.getHours())}:${pad(d.getMinutes())}:${pad(d.getSeconds())}`; }; /** `--k v` / `--flag` 两种形态都吃 */ function parseArgs(argv) { const o = { _: [] }; for (let i = 0; i < argv.length; i++) { const a = argv[i]; if (a.startsWith("--")) { const k = a.slice(2); const nxt = argv[i + 1]; if (nxt === undefined || nxt.startsWith("--")) o[k] = true; else { o[k] = nxt; i++; } } else o._.push(a); } return o; } /** `1h` / `30m` / `2d` / 纯秒 ⇒ 秒数 */ function dur2sec(s) { if (s === undefined || s === true) return null; const m = String(s).match(/^(\d+(?:\.\d+)?)([smhd]?)$/); if (!m) return null; const n = parseFloat(m[1]); return Math.round(n * ({ s: 1, m: 60, h: 3600, d: 86400 }[m[2] || "s"])); } function readJson(f, dflt) { try { return JSON.parse(readFileSync(f, "utf8")); } catch { return dflt; } } function loadHosts() { const h = readJson(HOSTS_FILE, null); return h && Object.keys(h).length ? h : DEFAULT_HOSTS; } // ───────────────────────── ssh(数据走 stdout,元信息走 stderr 哨兵) ───────────────────────── /** * 远端执行的脚本:⛔ 只用 journalctl;⛔ 不写服务器任何文件(数据直出 stdout)。 * 数据走 stdout(纯),元信息走 stderr(哨兵 + 行数 + 错误原文)⇒ 两条通道物理分离。 * ⚠️ `awk` 既转发 stdout 又把行数写 stderr:这样"拉取失败"与"确无日志"永远可分。 */ function remoteCollectScript({ sinceEpoch, untilEpoch, afterCursor, units, priority }) { const parts = [ "set -u", "set -o pipefail", 'J="journalctl --no-pager -o json"', ]; // ⚠️ 239 也支持 --output-fields(实测:能砍掉约一半体积) // 版本号解析必须 `NR==1` —— `journalctl --version` 是**多行**输出, // 不加限定会让 V 变成 "239\n0" ⇒ `[: integer expression expected` ⇒ 静默退回全字段(体积翻倍) parts.push('V=$(journalctl --version 2>/dev/null | awk \'NR==1{print $2+0}\')'); parts.push('if [ "${V:-0}" -ge 236 ] 2>/dev/null; then J="$J --output-fields=__REALTIME_TIMESTAMP,__CURSOR,_SYSTEMD_UNIT,_HOSTNAME,SYSLOG_IDENTIFIER,_PID,PRIORITY,MESSAGE"; fi'); if (sinceEpoch) parts.push(`J="$J --since @${sinceEpoch}"`); if (untilEpoch) parts.push(`J="$J --until @${untilEpoch}"`); if (afterCursor) parts.push(`J="$J --after-cursor='${afterCursor}'"`); if (priority) parts.push(`J="$J -p ${priority}"`); if (units && units.length) parts.push(`J="$J ${units.map((u) => `-u ${u}`).join(" ")}"`); parts.push( 'eval "$J" 2>/tmp/.dshlog_err.$$ | awk \'{n++} {print} END{printf "__DSHLOG_LINES__ %d\\n", n > "/dev/stderr"}\'', "RC=$?", 'EB=$(wc -c < /tmp/.dshlog_err.$$ 2>/dev/null || echo 0)', 'head -c 2000 /tmp/.dshlog_err.$$ >&2', "rm -f /tmp/.dshlog_err.$$", 'echo "__DSHLOG_EOF__ rc=$RC errbytes=$EB" >&2', ); return parts.join("\n"); } function sshRun(hostCfg, script, { timeoutMs = 120000 } = {}) { // 🔴 `-C` 不可省:47 的出方向未压缩带宽实测 ~20 KB/s(1.6 MB 要 84 s), // 开压缩后同样的数据 11 s(7.6×)。日志是 JSON ⇒ 压缩率极高,这是最大的单一杠杆。 const r = spawnSync("ssh", ["-C", ...hostCfg.ssh, "bash -s"], { input: script, encoding: "utf8", timeout: timeoutMs, maxBuffer: 256 * 1024 * 1024, }); const stderr = r.stderr || ""; const sentinel = stderr.split("\n").filter((l) => l.startsWith("__DSHLOG_EOF__")).pop() || ""; const m = sentinel.match(/rc=(\d+)\s+errbytes=(\d+)/); const rc = m ? parseInt(m[1], 10) : null; const errbytes = m ? parseInt(m[2], 10) : null; // 远端 journalctl 自己报告的输出行数(权威计数,用于与本地解析数对账) const lm = stderr.split("\n").filter((l) => l.startsWith("__DSHLOG_LINES__")).pop() || ""; const lm2 = lm.match(/__DSHLOG_LINES__ (\d+)/); const remoteLines = lm2 ? parseInt(lm2[1], 10) : null; // 远端 journalctl 自身的错误(剔除 ssh 噪声与自家哨兵行) const remoteErr = stderr.split("\n") .filter((l) => l.trim() && !l.startsWith("__DSHLOG_EOF__") && !l.startsWith("__DSHLOG_LINES__") && !/post-quantum|store now, decrypt later|openssh.com\/pq|WARNING: connection is not using/i.test(l)) .join("\n").trim(); return { stdout: r.stdout || "", rc, errbytes, remoteLines, remoteErr, sshStatus: r.status, sshError: r.error ? String(r.error.message || r.error) : null, }; } // ───────────────────────────── 归档读写 ───────────────────────────── function hostDir(host) { return path.join(ROOT, host); } function shardPath(host, day) { return path.join(hostDir(host), `${day}.ndjson.gz`); } function appendShard(host, day, lines) { if (!lines.length) return 0; mkdirSync(hostDir(host), { recursive: true }); const f = shardPath(host, day); // ⚠️ 必须同步写:异步 stream 在进程收尾时可能未 flush(会静默丢最后一批) // gzip 多成员文件可被 `gunzip` / `createGunzip` 顺序读回 ⇒ 追加安全 appendFileSync(f, gzipSync(Buffer.from(lines.join("\n") + "\n", "utf8"), { level: 6 })); return lines.length; } async function* readShard(host, day) { const f = shardPath(host, day); if (!existsSync(f)) return; const rl = createInterface({ input: createReadStream(f).pipe(createGunzip()), crlfDelay: Infinity }); for await (const line of rl) { if (!line.trim()) continue; try { yield JSON.parse(line); } catch { /* 损坏行跳过(不静默吞:由 stats 统计) */ } } } function listDays(host) { const d = hostDir(host); if (!existsSync(d)) return []; return readdirSync(d).filter((f) => f.endsWith(".ndjson.gz")).map((f) => f.replace(".ndjson.gz", "")).sort(); } /** 归一:journald JSON ⇒ 统一记录(msg 保真) */ function normalize(rec, host) { const tsUs = rec.__REALTIME_TIMESTAMP ? Number(rec.__REALTIME_TIMESTAMP) : 0; let msg = rec.MESSAGE; if (Array.isArray(msg)) { // journald 对非 UTF-8 字节会给数组 ⇒ 转义回可读文本(保真:字节以 \xNN 形式保留) msg = msg.map((b) => (b >= 32 && b < 127 ? String.fromCharCode(b) : `\\x${b.toString(16).padStart(2, "0")}`)).join(""); } return { ts: Math.floor(tsUs / 1000), host, unit: rec._SYSTEMD_UNIT || "", ident: rec.SYSLOG_IDENTIFIER || "", pid: rec._PID ? Number(rec._PID) : 0, pri: rec.PRIORITY ? Number(rec.PRIORITY) : 6, cursor: rec.__CURSOR || "", msg: typeof msg === "string" ? msg : JSON.stringify(msg ?? ""), }; } // ───────────────────────────── collect ───────────────────────────── /** * 校时:**以远端 NTP 同步状态为权威判据**,RTT 估算只作参考。 * ⚠️ 别用"ssh 往返估算"当权威 —— 实测 47 的 ssh RTT 达 3.7s 且往返不对称, * 单次估算能给出 1.3s 的假偏移(比真实时钟差还大)⇒ 那样的校正会**制造**错序。 */ function probeClock(cfg) { const r = sshRun(cfg, [ 'echo "__DSHLOG_NTP__ $(timedatectl show -p NTPSynchronized --value 2>/dev/null || echo unknown)"', 'echo "__DSHLOG_OFF__ $(timedatectl show -p Offset --value 2>/dev/null || echo unknown)"', 'echo "__DSHLOG_TM__ $(date +%s%N)"', ].join("\n"), { timeoutMs: 20000 }); const all = `${r.stdout}\n${r.stderr}`; const synced = (all.match(/__DSHLOG_NTP__ (\S+)/) || [])[1] || "unknown"; const offRaw = (all.match(/__DSHLOG_OFF__ (\S+)/) || [])[1] || "unknown"; const tms = (all.match(/__DSHLOG_TM__ (\d+)/) || [])[1]; return { ntpSynced: synced, ntpOffsetMs: offRaw !== "unknown" ? Math.round(Number(offRaw) / 1000) : null, remoteNowMs: tms ? Number(tms) / 1e6 : null, measuredAt: Date.now(), }; } async function cmdCollect(args) { mkdirSync(ROOT, { recursive: true }); const hosts = loadHosts(); const want = args.hosts && args.hosts !== true ? String(args.hosts).split(",").map((s) => s.trim()) : Object.keys(hosts); const st = readJson(STATE, { hosts: {} }); st.hosts = st.hosts || {}; const since = dur2sec(args.since); const full = !!args.full; const backfill = !!since; // 回填 = 显式给 --since(与"续拉"是两种语义) const chunk = args.chunk ? dur2sec(args.chunk) : Math.max(3600, Math.ceil((since || 3600) / 8)); const rows = []; for (const h of want) { const cfg = hosts[h]; if (!cfg) { rows.push([h, "SKIP", "hosts.json 无此节点", ""]); continue; } const prev = st.hosts[h] || {}; // 回填必须去重(同一 cursor 已在归档里则跳过),否则重复落盘会污染行数类判据 const seen = new Set(); if (backfill) { for (const d of listDays(h)) for await (const r of readShard(h, d)) if (r.cursor) seen.add(r.cursor); } // 时间窗切块:单次 ssh 只搬一块 ⇒ 内存/传输受控,且任何一块失败可单独重试 let windows; if (backfill) { const end = nowSec(); windows = []; for (let t = end - since; t < end; t += chunk) windows.push([t, Math.min(t + chunk, end)]); if (!windows.length) windows = [[end - 60, end]]; } else if (!full && prev.cursor) { windows = [[null, null]]; // cursor 续拉(精确、不重不漏) } else { windows = [[nowSec() - 3600, nowSec()]]; } const clock = probeClock(cfg); let ok = 0, bad = 0, dup = 0, written = 0, firstErr = "", jrcBad = null, mismatch = false, remoteReported = 0; let lastCursor = prev.cursor || "", lastTs = prev.lastTs || 0; const dayMap = {}; for (const [from, to] of windows) { if (process.env.DSHLOG_VERBOSE) console.error(` [${h}] 拉块 ${from === null ? "cursor" : tsLocal(from * 1000)} → ${to === null ? "now" : tsLocal(to * 1000)} …`); const res = sshRun(cfg, remoteCollectScript({ sinceEpoch: from, untilEpoch: to, afterCursor: from === null ? prev.cursor : null, }), { timeoutMs: 600000 }); const sshFailed = res.rc === null && res.sshStatus !== 0; if (sshFailed) { if (!firstErr) firstErr = res.sshError || "ssh 失败"; continue; } if (res.rc !== null && res.rc !== 0 && res.rc !== 1) { jrcBad = res.rc; } if (res.remoteLines !== null) remoteReported += res.remoteLines; if (res.remoteErr && !firstErr) firstErr = res.remoteErr.slice(0, 160); let chunkOk = 0, chunkBad = 0; for (const l of res.stdout.split("\n")) { if (!l.trim()) continue; let rec; try { rec = JSON.parse(l); } catch { bad++; chunkBad++; continue; } const n = normalize(rec, h); ok++; chunkOk++; if (backfill && n.cursor && seen.has(n.cursor)) { dup++; continue; } if (n.cursor) seen.add(n.cursor); if (n.cursor) lastCursor = n.cursor; if (n.ts > lastTs) lastTs = n.ts; const d = dayStr(new Date(n.ts)); (dayMap[d] = dayMap[d] || []).push(JSON.stringify(n)); } if (res.remoteLines !== null && res.remoteLines !== chunkOk + chunkBad) mismatch = true; } for (const [d, ls] of Object.entries(dayMap)) written += appendShard(h, d, ls); const state = firstErr && ok === 0 ? "FAIL" : jrcBad !== null ? `JOURNALCTL_RC${jrcBad}` : mismatch ? "MISMATCH" : ok === 0 ? "EMPTY" : "OK"; st.hosts[h] = { cursor: lastCursor, lastTs, lastRun: Date.now(), ntpSynced: clock.ntpSynced, ntpOffsetMs: clock.ntpOffsetMs, lastState: state, lastLines: ok, lastWritten: written, lastDup: dup, lastRemoteLines: remoteReported, }; let detail = `拉取 ${ok} 行`; if (remoteReported) detail += ` / 远端报 ${remoteReported}`; if (dup) detail += ` · 去重跳 ${dup}`; if (written) detail += ` · 落盘 ${written}`; if (bad) detail += ` · 解析失败 ${bad}`; if (firstErr) detail += ` · err: ${firstErr.slice(0, 120)}`; rows.push([h, state, detail, `${windows.length} 块` + (backfill ? ` since @${nowSec() - since}` : " cursor 续拉") + ` · NTP ${clock.ntpSynced}` + (clock.ntpOffsetMs === null ? "" : `(${clock.ntpOffsetMs}ms)`)]); } writeFileSync(STATE, JSON.stringify(st, null, 2), "utf8"); if (!args.quiet) printTable(["节点", "状态", "结果", "方式"], rows, args); const badRows = rows.filter((r) => !["OK", "EMPTY", "SKIP"].includes(r[1])); if (badRows.length) process.exitCode = 2; return rows; } // ───────────────────────────── prune(保留策略) ───────────────────────────── function cmdPrune(args) { const keep = args.keep ? Number(args.keep) : 3; // 与服务器侧口径一致:3 天 const cut = dayStr(new Date(Date.now() - keep * 86400e3)); const hosts = loadHosts(); const rows = []; let total = 0; for (const h of Object.keys(hosts)) { for (const d of listDays(h)) { if (d >= cut) continue; const p = shardPath(h, d); const sz = statSync(p).size; total += sz; if (args.apply) { try { unlinkSync(p); rows.push([h, d, "已删除", `${(sz / 1e6).toFixed(2)} MB`]); } catch (e) { rows.push([h, d, "删除失败", String(e.message).slice(0, 60)]); } } else rows.push([h, d, "待删(干跑)", `${(sz / 1e6).toFixed(2)} MB`]); } } printTable(["节点", "日期", "动作", "大小"], rows, args); console.log(`保留 ${keep} 天 ⇒ 截止 ${cut}${args.apply ? " · 已执行" : " · 干跑(加 --apply 才真删)"}|涉及 ${(total / 1e6).toFixed(2)} MB`); } // ───────────────────────────── ls / stats ───────────────────────────── function cmdLs(args) { const hosts = loadHosts(); const rows = []; for (const h of Object.keys(hosts)) { for (const d of listDays(h)) { const p = shardPath(h, d); rows.push([h, d, `${(statSync(p).size / 1024 / 1024).toFixed(2)} MB`]); } } printTable(["节点", "日期", "大小"], rows, args); } async function cmdStats(args) { const hosts = loadHosts(); const st = readJson(STATE, { hosts: {} }); const rows = []; for (const h of Object.keys(hosts)) { const s = (st.hosts || {})[h] || {}; rows.push([h, s.lastState || "-", String(s.lastLines ?? "-"), s.lastRun ? tsLocal(s.lastRun) : "-", s.ntpSynced ? `${s.ntpSynced}${s.ntpOffsetMs === null || s.ntpOffsetMs === undefined ? "" : ` (${s.ntpOffsetMs}ms)`}` : "-", s.cursor ? s.cursor.slice(0, 24) + "…" : "-"]); } printTable(["节点", "末次状态", "行数", "末次拉取", "NTP", "游标"], rows, args); if (args.detail) { for (const h of Object.keys(hosts)) { for (const d of listDays(h)) { let n = 0; const units = new Map(); for await (const r of readShard(h, d)) { n++; units.set(r.unit, (units.get(r.unit) || 0) + 1); } const top = [...units.entries()].sort((a, b) => b[1] - a[1]).slice(0, 4) .map(([u, c]) => `${u.replace(/\.service|\.scope/g, "")}=${c}`).join(" "); console.log(` ${h} ${d}: ${n} 行 · ${top}`); } } } } // ───────────────────────────── q(查询) ───────────────────────────── async function cmdQ(args) { const pat = args._[0]; if (!pat) { console.error("用法: dshlog q <关键词|/正则/> [--since 6h] [--host 47] [--unit dshs] [--limit 100] [--json]"); process.exitCode = 2; return; } const re = pat.startsWith("/") && pat.endsWith("/") ? new RegExp(pat.slice(1, -1), "i") : new RegExp(pat.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"), "i"); const hosts = args.host && args.host !== true ? String(args.host).split(",") : Object.keys(loadHosts()); const sinceMs = args.since ? Date.now() - dur2sec(args.since) * 1000 : 0; const untilMs = args.until ? Date.now() - dur2sec(args.until) * 1000 : Infinity; const unitRe = args.unit && args.unit !== true ? new RegExp(String(args.unit), "i") : null; const limit = args.limit ? Number(args.limit) : 100; const hits = []; let scanned = 0; for (const h of hosts) { for (const d of listDays(h)) { for await (const r of readShard(h, d)) { scanned++; if (r.ts < sinceMs || r.ts > untilMs) continue; if (unitRe && !unitRe.test(r.unit)) continue; if (!re.test(r.msg)) continue; hits.push(r); if (hits.length >= limit * 20) break; } } } hits.sort((a, b) => a.ts - b.ts); const show = hits.slice(-limit); if (args.json) console.log(JSON.stringify(show, null, 1)); else { console.log(`命中 ${hits.length} 条(扫描 ${scanned} 行)⇒ 显示末 ${show.length} 条`); for (const r of show) console.log(`${tsLocal(r.ts)} [${r.host}/${r.unit.replace(/\.service|\.scope/g, "")}] ${r.msg.slice(0, 300)}`); } if (!show.length) process.exitCode = 1; } // ───────────────────────────── timeline(事后溯源) ───────────────────────────── async function cmdTimeline(args) { const hosts = args.host && args.host !== true ? String(args.host).split(",") : Object.keys(loadHosts()); const st = readJson(STATE, { hosts: {} }); const from = args.from ? new Date(args.from).getTime() : Date.now() - 3600e3; const to = args.to ? new Date(args.to).getTime() : Date.now() + 60e3; const grep = args.grep && args.grep !== true ? new RegExp(String(args.grep), "i") : null; const priMax = args.pri ? Number(args.pri) : 7; const fix = !!args["fix-clock"]; // 默认不校正:NTP 偏差是"系统事实",掩盖它反而会误导溯源 const all = []; for (const h of hosts) { const hs = (st.hosts || {})[h] || {}; const off = fix ? (hs.ntpOffsetMs || 0) : 0; for (const d of listDays(h)) { for await (const r of readShard(h, d)) { const t = r.ts - off; if (t < from || t > to) continue; if (r.pri > priMax) continue; if (grep && !grep.test(r.msg)) continue; all.push({ ...r, ts: t }); } } } all.sort((a, b) => a.ts - b.ts || a.host.localeCompare(b.host)); const out = args.out && args.out !== true ? args.out : null; const lines = all.map((r) => `${tsLocal(r.ts)}.${String(r.ts % 1000).padStart(3, "0")} [${r.host.padEnd(3)}/${(r.unit || "-").replace(/\.service|\.scope/g, "").padEnd(16)}] pri=${r.pri} ${r.msg}`); if (out) { writeFileSync(out, lines.join("\n") + "\n", "utf8"); console.log(`已写 ${all.length} 行 ⇒ ${out}`); } else { console.log(`时间线 ${all.length} 行(${tsLocal(from)} → ${tsLocal(to)})`); for (const l of lines) console.log(l); } const ntpNote = hosts.map((h) => { const hs = (st.hosts || {})[h] || {}; return `${h}: NTP=${hs.ntpSynced || "?"}${hs.ntpOffsetMs === undefined || hs.ntpOffsetMs === null ? "" : `(偏差${hs.ntpOffsetMs}ms)`}`; }).join(" · "); console.log(`# 时钟状态 ${ntpNote}${fix ? " · 已按偏差校正(--fix-clock)" : " · 未校正(原始系统时间戳)"}`); } // ───────────────────────────── watch(实时巡检) ───────────────────────────── async function cmdWatch(args) { const since = args.since ? dur2sec(args.since) : 900; // 默认 15 分钟窗口 const hosts = args.host && args.host !== true ? String(args.host).split(",") : Object.keys(loadHosts()); // 默认先增量续拉(cursor 只搬新行,代价极小)⇒ 挂 automation 时只需这一条命令 if (!args["no-collect"]) { try { await cmdCollect({ hosts: args.host, quiet: !args.verbose, json: args.json }); } catch (e) { console.error(`⚠ 预拉取失败(继续按现有归档巡检): ${e.message}`); } } const from = Date.now() - since * 1000; const counts = new Map(); // ruleId -> [{host,ts,msg}] const byUnit = new Map(); let total = 0, priErr = 0; for (const h of hosts) { for (const d of listDays(h)) { for await (const r of readShard(h, d)) { if (r.ts < from) continue; total++; if (r.pri <= 3) priErr++; byUnit.set(`${h}/${r.unit}`, (byUnit.get(`${h}/${r.unit}`) || 0) + 1); if (NOISE_UNITS.test(r.unit)) continue; for (const [id, , re] of SEV_RULES) { if (re.test(r.msg)) { if (!counts.has(id)) counts.set(id, []); const arr = counts.get(id); if (arr.length < 200) arr.push({ host: h, ts: r.ts, msg: r.msg, unit: r.unit }); } } } } } const maxHit = args.threshold ? Number(args.threshold) : 0; // 0 = 命中即 FAIL(告警模式) const rows = []; for (const [id, name] of SEV_RULES) { const hits = counts.get(id) || []; const verdict = hits.length > maxHit ? "FAIL" : "PASS"; const sample = hits.slice(-1)[0]; rows.push([id, name, verdict, String(hits.length), sample ? `${tsLocal(sample.ts)} ${sample.host} ${sample.msg.slice(0, 90)}` : "—"]); } console.log(`巡检窗口 = 近 ${Math.round(since / 60)} 分钟 · 归档内 ${total} 行 · pri≤3 计 ${priErr} 行`); printTable(["判据", "含义", "结果", "命中", "最近一条"], rows, args); const fails = rows.filter((r) => r[2] === "FAIL"); if (args.report) { const f = path.join(ROOT, "reports", `watch-${dayStr()}-${new Date().getHours()}${pad(new Date().getMinutes())}.md`); mkdirSync(path.dirname(f), { recursive: true }); writeFileSync(f, [ `# 日志巡检报告 ${tsLocal(Date.now())}`, "", `窗口 = 近 ${Math.round(since / 60)} 分钟 | 归档行数 ${total} | pri≤3 = ${priErr}`, "", "| 判据 | 含义 | 结果 | 命中 | 最近一条 |", "|---|---|---|---|---|", ...rows.map((r) => `| ${r[0]} | ${r[1]} | ${r[2]} | ${r[3]} | ${String(r[4]).replace(/\|/g, "\\|")} |`), "", ].join("\n"), "utf8"); console.log(`报告 ⇒ ${f}`); } if (fails.length) process.exitCode = 2; } // ───────────────────────────── 输出 ───────────────────────────── function printTable(head, rows, args) { if (args.json) { console.log(JSON.stringify(rows.map((r) => Object.fromEntries(head.map((h, i) => [h, r[i]]))), null, 1)); return; } const w = head.map((h, i) => Math.max(strw(h), ...rows.map((r) => strw(String(r[i]))))); const line = (cells) => cells.map((c, i) => String(c) + " ".repeat(Math.max(0, w[i] - strw(String(c))))).join(" "); console.log(line(head)); console.log(w.map((n) => "─".repeat(n)).join(" ")); for (const r of rows) console.log(line(r)); } /** 中文按 2 列宽计 */ function strw(s) { let n = 0; for (const ch of String(s)) n += ch.charCodeAt(0) > 0x2e7f ? 2 : 1; return n; } // ───────────────────────────── main ───────────────────────────── const argv = process.argv.slice(2); const cmd = argv[0]; const args = parseArgs(argv.slice(1)); const cmds = { collect: cmdCollect, ls: cmdLs, stats: cmdStats, q: cmdQ, timeline: cmdTimeline, watch: cmdWatch, prune: cmdPrune }; if (!cmd || cmd === "help" || args.help) { console.log(`dshlog — DSH 多节点日志采集 / 巡检 / 溯源(归档根 ${ROOT}) collect [--hosts 47,106] [--since 1h [--chunk 30m]|--full] [--json] 拉取归档:无 --since = cursor 续拉(增量、不重不漏);有 --since = 回填(自动去重) ls 列出归档分片 stats [--detail] 归档与末次拉取状态 q <词|/正则/> [--since 6h] [--host 47] [--unit dshs] [--limit 100] [--json] timeline [--from ISO] [--to ISO] [--host ..] [--grep ..] [--pri 3] [--out f] [--fix-clock] watch [--since 15m] [--threshold N] [--report] [--no-collect] [--json] 巡检:默认先增量续拉再判据;有 FAIL ⇒ 退出码 2(可被自动化当判据) prune [--keep 3] [--apply] 归档保留策略(默认 3 天 · 与服务器同口径 · 干跑) 环境变量:DSHLOG_ROOT(归档根,默认 E:/dsh-logs)· DSHLOG_VERBOSE=1(打印每块进度) `); process.exit(0); } if (!cmds[cmd]) { console.error(`未知子命令: ${cmd}`); process.exit(2); } await cmds[cmd](args);