Files
admin 452924d89c feat(config): 涉密内容外置到配置目录(档案 140)
把散落在代码里的真实部署值统一收进 config/,代码改为引用配置,
使仓库副本/开源导出不再带出生产域名、IP、内网路径与凭据。

新增 config/:platform.env.example(模板)· load.sh(shell 加载器)·
index.cjs(node 加载器)· README.md(键一览与优先级)。
真实值放 config/platform.env —— 已 .gitignore 排除,不入库、不进导出。

TS 侧新增 src/platform-paths.ts 作部署路径的唯一解析处(零副作用):
platformDir/stateDir/backupDir/artifactDir/installDir/scriptPath。
config.ts 接入这些字段;内置中继种子由生产 URL 改为空(改由
DSHS_OVERLAY_BOOTSTRAP_SEEDS 提供)。修掉 5 处硬编码绝对路径,
src/** 注释中性化 116 行/53 文件。

scripts/** 36 个内部运维脚本:真令牌/PG 口令/隧道目标/主机号/路径
一律改从配置取;web/wake.html 的注册域白名单改为运行时从
location.hostname 推导;test/** 夹具 119 行/13 文件改 RFC 2606/5737
保留值,并把「内置种子必须为空」固化为回归断言。

取证:tsc 0 错;npm test 373/375(唯一失败 lease 属既有);
全仓扫描(大小写不敏感)代码面涉密标识 = 0;已部署 47 并零回归
(/opt/dsh/* 未搬家,/var/lib/dshs/platform 未被误建)。
2026-09-19 15:12:19 +08:00

579 lines
29 KiB
JavaScript
Raw Permalink 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");
// 主机表从配置读取(⛔ 不在代码里写死地址 / ssh 别名):
// DSHLOG_HOSTS='{"mgr":{"ssh":["-p","22","[email protected]"],"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);