Files
dsh_shenxian/scripts/dshlog.mjs
T
admin d2ef362a98 feat(overlay): 覆盖网络线 序㊾ —— 探针观测面改「两台中继并集」(附 序㊽ 源码/文档补提交)
序㊾(本棒):
- scripts/overlay-probe.cjs:OBS-01 / OBS-08 / OBS-09 的数据源由「只读 47 中继」
  改为「按两台中继取并集」,消除 worker 归属漂移时的假红 / 假 SKIP
  · endpoints 以 network:hostId:port 为键合并,online 取「或」、localPort 取在线那一侧
  · used 按 network/hostId 去重计数(不求和,避免凭空放大在册数)
  · localPort 属中继机回环落点 ⇒ 按归属分机探活(106 侧落点由 106 机上探)
  · derived(OBS-11)保持 47 视角;阈值与判据一律未放宽
  · OBS-16 计数约束:对 47 /status 的读取仍为三次、Δ 只取 47 的 counters;
    对端 106 的采样为独立一次,落在第三次采样之后,不进 (status2, status3] 门窗口
  · 新增 --peer-status-fixture(并集的对端那一半)与「并集不可取证」强制留痕
- 交接单《覆盖网络-序45-低熵块治理-测熵与实现》§16 全节(§8 前前缀逐字未变)
- 参数表 §11.16 补记(§10 现算指纹未变,值格未动)

附(前几棒已完成并已部署、但尚未入仓的源码 / 文档):
- src/net/relay/content/*.ts、src/net/relay/index.ts、main.ts:块级寻址 C 域分离
- src/supervisor/orchestrator.ts、src/worker/agent.ts:日志采集与巡检(方案 C)
- test/overlay-content.test.mjs:随附用例(npm test = 200 pass / 0 fail / 1 skipped,Node 22)
- scripts/dshlog.mjs(跨机日志取证)、scripts/overlay-entropy.cjs(熵探针)
- dsh-server-docs/04-调整方案/129、133;INDEX.md / docs-manifest.json / 交接单 README 登记
2026-09-19 05:35:35 +08:00

570 lines
28 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env node
/**
* dshlog.mjs — DSH 平台「多节点日志采集 / 实时巡检 / 事后溯源」单一入口(方案 C:零新增常驻服务)
*
* 设计边界(2026-09-18 定)
* · 服务器侧**零安装**:远端只用系统自带的 `journalctl`,采集靠 ssh 一次性拉取,不装 agent、不开监听口
* · 归档在**本机**:`E:/dsh-logs/<host>/<YYYY-MM-DD>.ndjson.gz`(可按天清理)
* · 日志原文**保真**:`msg` 字段逐字节原样落盘(本项目大量判据依赖日志原文的行数/字节数)
* · 失败**可分**:远端元信息走 stderr 哨兵 `__DSHLOG_EOF__`,与 stdout 的纯数据分离
* ⇒ "拉取失败" 与 "确无日志" 永不混淆(这是本项目反复踩过的坑)
*
* 子命令
* collect 拉取归档(增量,默认只拉「上次游标之后」)
* ls 列出归档分片
* q 关键词/正则查询(跨机跨单元)
* timeline 时间线重建(多机日志按时间戳合并,用于事后溯源)
* watch 实时巡检(规则命中 → 判据表 + 非零退出码,可挂 automation)
* stats 归档统计
*
* 约定
* · 时间一律走 epoch(`--since @<epoch>`;⛔ 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");
const DEFAULT_HOSTS = {
"47": { ssh: ["-p", "22", "[email protected]"], label: "Manager + w-47" },
"106": { ssh: ["-p", "22", "test106"], label: "w-106" },
};
// 巡检规则:只认**指向本项目自身故障**的信号。
// ⚠️ 每条规则都是误报与漏报的取舍 —— 下面两组是实测调过的:
// · 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);