package jnpf.audit.sdk;
|
|
import jnpf.audit.model.AuditEventDTO;
|
import jnpf.util.JsonUtil;
|
import lombok.extern.slf4j.Slf4j;
|
import org.springframework.beans.factory.DisposableBean;
|
import org.springframework.beans.factory.InitializingBean;
|
|
import java.io.File;
|
import java.io.IOException;
|
import java.nio.charset.StandardCharsets;
|
import java.nio.file.Files;
|
import java.nio.file.StandardCopyOption;
|
import java.nio.file.StandardOpenOption;
|
import java.security.MessageDigest;
|
import java.security.NoSuchAlgorithmException;
|
import java.util.ArrayList;
|
import java.util.List;
|
import java.util.UUID;
|
import java.util.concurrent.Executors;
|
import java.util.concurrent.ScheduledExecutorService;
|
import java.util.concurrent.TimeUnit;
|
import java.util.regex.Matcher;
|
import java.util.regex.Pattern;
|
|
/**
|
* spool 重放器(Task 7 Step 3;评审#3:spool 必须可恢复才算兜底)。
|
*
|
* <p>周期性扫描 spool 目录的 {@code .jsonl} / {@code .retry} 文件 → 原子重命名加
|
* {@code .replaying.<owner-token>} 后缀(认领,token 每次新生成,见"#B+#C"段)→ 逐行反序列化走
|
* {@link AuditEventPublisher#deliverOnce(AuditEventDTO)}
|
* (成功丢弃该行、抛异常保留)→ 全部成功删 {@code .replaying};部分失败把剩余行写入
|
* {@code <稳定 base 文件名>.<UUID>.retry} 新文件(base+UUID 保证同轮多文件失败不撞名,终验五轮⑤;
|
* base 名的取法见下文"内审 Important 修复"段 {@link #stableBaseName(String)}),
|
* {@code .retry} 写成功后删对应 {@code .replaying}(不留孤儿);<b>绝不写回活动原名</b>(终验四轮③:
|
* 重命名期间写入端可能已在原名新建文件追加新事件,写回原名会覆盖丢数据)。落库幂等(client_event_id)
|
* 使重放天然安全。
|
*
|
* <p>调度手段:自建单线程 daemon {@link ScheduledExecutorService}(非 @Scheduled),不依赖宿主
|
* 是否开启 @EnableScheduling,生命周期由本 bean 自控(见 task-7-report 决策说明)。
|
* 由 {@code audit.spool-replay-enabled}(默认 true)开关。
|
*
|
* <p><b>双层锁协议(Codex 批次二复验 #1/#2 终局修复)</b>:与 {@link AuditSpoolWriter} 共享同一
|
* {@link AuditSpoolLock}(进程内 monitor + 跨进程 {@code spool.lock} FileLock)。快速文件操作在锁内、
|
* 慢速网络重放在锁外,临界区划分:
|
* <ol>
|
* <li><b>认领</b>({@link #replayFile}):{@code rename→.replaying.<owner-token>} + {@code setLastModified}
|
* 刷新租约在<b>同一临界区</b>——#1 与 writer 的 append 互斥(消除"取走文件瞬间"与"正在 append"的窗口
|
* 竞态);#2 rename 与刷新原子成对(消除"rename 成功后、刷新 mtime 前被孤儿回收误抢"的窗口)。
|
* 临界区内<b>先锁内重验</b>(#B)再 rename。replayer 路径 {@code tryLock} 拿不到就跳过本轮,下轮再来。</li>
|
* <li><b>readAllLines</b>:锁外——认领后 {@code .replaying} 已不在任何写入端/扫描集视线内(孤儿回收
|
* 10 分钟阈值 + 刚刷新的 mtime 保护),读取安全且不阻塞 writer。</li>
|
* <li><b>逐行 deliverOnce</b>:锁外(网络慢操作)。</li>
|
* <li><b>收尾</b>({@link #finishReplay}):写 {@code .bad} / 写 {@code .retry} / 删 {@code .replaying}
|
* 在锁内(各自快操作)。</li>
|
* <li><b>孤儿回收 rename</b>({@link #recoverOrphanedReplaying}):检查+rename 在锁内。</li>
|
* </ol>
|
* 60 秒静默窗口({@link #QUIET_WINDOW_MS})<b>保留为纯优化</b>(减少活跃文件搅动),非正确性依赖——
|
* 正确性已由锁协议保证。<b>残余</b>:单文件重放超 10 分钟仍可能被其它进程当孤儿抢走 → 后果=重复重放,
|
* 落库按 client_event_id 幂等保证不丢不重入库,可接受;{@code spool.lock} 永不参与扫描/重放/删除
|
* (扫描过滤已按 {@link AuditSpoolLock#LOCK_FILE_NAME} 排除)。
|
*
|
* <p><b>裁定修复 D(孤儿 .replaying 恢复)</b>:进程若在"改名成 .replaying"之后、"处理完"之前崩溃,
|
* 会留下孤儿——扫描集固定为 {@code {.jsonl,.retry}},永远不会认领它,事件永久滞留。每轮扫描时,
|
* 额外把<b>修改时间早于(当前时刻 - 10 分钟,即两个调度周期)</b>的 .replaying 文件原子改名为
|
* {@code <稳定 base 文件名>.<UUID>.retry}(锁内),回归 retry 通道由下一轮自然消化;不在回收动作
|
* 内直接处理文件内容,避免与"本实例/并发实例正在重放中的文件"混淆。年龄阈值即为防止误抢:正常一轮
|
* 处理远快于一个调度周期,若 .replaying 存活跨越两个周期才可判定为孤儿。落库按 client_event_id
|
* 幂等,重复重放安全。
|
*
|
* <p><b>Codex 批次二三轮 #B+#C(认领名 owner-token 化 + 锁内重验)</b>:认领 rename 目标由固定后缀
|
* {@code .replaying} 改为 owner-token 化的 {@code .replaying.<ownerUUID>}(token 每次认领新生成,
|
* {@link #REPLAYING_TOKEN_SUFFIX} 识别)。收尾/删除只操作本 owner token 的路径——陈旧 owner(被抢后收尾)
|
* 与新 owner 永不同名,{@code deleteIfExists} 物理上不可能删掉新 owner 刚认领的文件(消除 #C 固定名跨
|
* owner 冲突)。同时认领({@link #replayFile})与孤儿回收({@link #recoverOrphanedReplaying})都改为
|
* <b>锁内先重验</b>(exists + mtime 仍超阈/仍静默)再 rename——锁外 {@code listFiles} 预筛只作候选收集,
|
* 消除"预筛后、锁内 rename 前,路径被删除并复用为新文件 → 旧检查结果误搬新文件"的窗口(#B)。
|
* {@code setLastModified} 失败保留 warn+继续:owner-token 唯一名下租约失败的最坏后果=被孤儿回收→重复重放
|
* (幂等吸收),不再有丢失路径。
|
*
|
* <p><b>Codex 批次二三轮 #D(优雅退出)</b>:{@link #destroy()} 在 {@code shutdownNow} 后
|
* {@code awaitTermination}({@link #SHUTDOWN_AWAIT_MS}) 等待重放线程真正退出(超时 warn),与
|
* {@link AuditSpoolLock} 的终止位配合,避免"停机返回但重放线程仍在锁内收尾"。
|
*
|
* <p><b>内审 Important 修复(retry 文件名收敛,防 NAME_MAX 无界增长)</b>:部分失败写 {@code .retry}、
|
* 孤儿回收改名、坏行隔离写 {@code .bad},前缀都先经 {@link #stableBaseName(String)} 从文件名剥离<b>所有</b>
|
* {@code .<UUID>.retry} 层与 {@code .replaying.<UUID>} token 后缀还原出稳定 base(如 {@code audit-events-20260723.jsonl}),
|
* 再拼<b>恰好一层</b>新 UUID,文件名长度从此有界,不随失败轮次增长。UUID 段用 {@link #RETRY_LAYER_SUFFIX}
|
* 严格匹配标准 8-4-4-4-12 十六进制格式并锚定字符串末尾,不会误剥业务文件名自身携带的点段(如 {@code .jsonl})。
|
* (Codex 五轮起:{@link #recoverOrphanedWriting()} 的隔离改名不再走本机制——直接在原名后追加
|
* {@code .stale},见其专属段;{@code .retry.writing} 层剥离分支保留作通用防御,当前无调用方会触发。)
|
*
|
* <p><b>坏行隔离与脱敏(Codex 批次二 #3/#4)</b>:反序列化失败的坏行不静默丢、不进 remaining,收尾时批量
|
* 隔离进 {@code <稳定 base>.<UUID>.bad}(不参与任何自动重放,留人工修复补录,见 {@link #finishReplay});
|
* 坏行 ERROR 日志只记文件名/行号/行长/内容 SHA-256 前 12 位({@link #sha256Prefix12}),绝不输出行内容正文。
|
*
|
* <p><b>写后原子发布孤儿回收(Codex 四轮 Critical#1)</b>:{@link AuditSpoolWriter} 降级路径改为「先写
|
* 唯一名 {@code .writing} 临时文件、再原子改名发布为 {@code .retry}」(消灭跨进程半写窗口,见其类
|
* Javadoc)。若 writer 恰好崩溃于两步之间(或改名本身因极端 IO 异常失败),会遗留
|
* {@code <稳定base>.<UUID>.retry.writing} 孤儿——扫描集固定排除 {@code .writing}(见
|
* {@link #replayOnce()}),主流程永远不会认领它。每轮额外调用 {@link #recoverOrphanedWriting()}
|
* 处理这类孤儿;具体动作见下一段——Codex 五轮复验推翻了本段最初"回归 .retry 通道当正常事件重放"的方案。
|
*
|
* <p><b>Codex 五轮修复(.writing 隔离区,不能仅凭年龄自动发布)</b>:四轮版本曾把超龄
|
* {@code .retry.writing} 直接改名回归 {@code .retry} 通道当正常事件重放——Codex 五轮复验指出这不安全:
|
* degraded writer(未持 FileLock 的降级路径)若恰好在其唯一一次 {@link Files#write} 调用<b>内部</b>
|
* 被操作系统/磁盘 IO 阻塞超过 10 分钟,该 writer 线程仍然存活、仍持有该文件的打开 FD,只是尚未返回;
|
* 此时它的 {@code .writing} 文件在孤儿回收眼中已"超龄",若照旧改名发布为 {@code .retry},重放器会把它
|
* 当正常事件认领(rename→{@code .replaying})、读到空文件(内容还没真正落盘)、判定"处理完毕"并删除
|
* {@code .replaying};随后卡住的 write 调用恢复,向着已被 unlink 的 inode 追加——数据静默消失,不留
|
* 任何错误痕迹。半行截断时同理:坏行隔离机制只能吸收"取走时刻"已经写入的前缀,写调用尚未完成的后半段
|
* 一样会丢。
|
*
|
* <p>结论:<b>超龄不足以断言 writer 已死</b>,{@link #recoverOrphanedWriting()} 因此只做<b>隔离</b>、
|
* 绝不发布、绝不删除——锁内重验通过后把 {@code .retry.writing} 原子改名为 {@code <原名>.stale}(同
|
* 目录),{@code log.error} 告警(脱敏:仅记文件名/大小/mtime,不含内容)后即止步:不进 {@code .retry}
|
* 通道(不会被重放器认领)、不进 {@link #finishReplay} 收尾(不会被删除)。{@code .stale} 排除在一切
|
* 扫描/认领之外——主扫描({@link #replayOnce()})只认 {@code .jsonl}/{@code .retry},
|
* {@link #recoverOrphanedReplaying()} 只认 {@code .replaying.<uuid>},{@link #recoverOrphanedWriting()}
|
* 自身只认 {@code .retry.writing},三者的后缀判定都不会匹配以 {@code .stale} 收尾的文件名,故隔离后的
|
* 文件物理上不可能被本类任何逻辑再次触碰;也无需纳入 {@link #stableBaseName(String)} 的剥层范围——
|
* {@code .stale} 从不参与任何"新建 retry/bad 文件"的前缀推导。
|
*
|
* <p>安全性论证的核心是 POSIX {@code rename(2)} 的语义:<b>rename 不 unlink</b>——它只改变目录项指向的
|
* 名字,被卡住的 write 调用持有的是对 inode 的 FD,与路径名无关,恢复后仍会把内容完整写入<b>同一个
|
* inode</b>(现在挂在 {@code .stale} 这个名字下),不会因改名而丢失或写偏。事件因此保全在隔离区,等待
|
* 人工核查后手工补录(对照 {@code clientEventId} 去重);隔离区本身永不被任何自动化流程删除,因此
|
* "重放器把仍在写的文件当孤儿抢走→空读→删除→静默丢失"这条路径被<b>物理消除</b>(不是靠时序窗口收窄,
|
* 而是代码上根本不存在能删除 {@code .stale} 的路径)。writer 侧针对"发布源被隔离"的自愈见
|
* {@link AuditSpoolWriter} 类 Javadoc"Codex 五轮修复(自愈重发)"段——move 目标源被隔离后不会让事件
|
* 真的丢,只会在 {@code .stale} 里留一份至多被人工补录时重复的副本(落库按 {@code clientEventId} 幂等
|
* 吸收)。
|
*/
|
@Slf4j
|
public class AuditSpoolReplayer implements InitializingBean, DisposableBean {
|
|
private static final long DEFAULT_REPLAY_INTERVAL_MS = 300_000L; // spec: fixedDelay 5 分钟(默认,可经 audit.spool-replay-interval-ms 覆盖)
|
private static final String REPLAYING_SUFFIX = ".replaying";
|
// 裁定修复 D:孤儿 .replaying.<uuid> 年龄阈值 = 两个默认调度周期(10 分钟),防止误抢并发实例正在处理的文件。
|
// 刻意锚定 DEFAULT_REPLAY_INTERVAL_MS 而非实例调度周期:孤儿年龄阈值是"多久才敢判定文件被崩溃遗弃"的
|
// 崩溃恢复安全底线,与"多久轮询一次"是两个关注点——把轮询周期调短(如测试环境 15s 加速验证)不应连带
|
// 缩小这条安全窗口,否则一次略慢的正常重放就可能被误判为孤儿并被并发实例抢走(幂等吸收但徒增搅动)。
|
private static final long ORPHAN_REPLAYING_AGE_MS = 2 * DEFAULT_REPLAY_INTERVAL_MS;
|
// Codex 批次二 #1:活动 .jsonl 静默窗口——距今不足 60 秒本轮跳过。锁协议落地后此窗口降级为纯优化
|
// (减少活跃文件搅动),非正确性依赖;正确性由 AuditSpoolLock 双层锁保证
|
private static final long QUIET_WINDOW_MS = 60_000L;
|
// 内审 Important 修复:精确匹配"尾部恰好一层 <UUID>.retry",UUID 严格按 8-4-4-4-12 十六进制格式
|
// 锚定字符串末尾($),避免误剥业务文件名自身携带的点段(如 audit-events-20260723.jsonl 的 .jsonl)
|
private static final Pattern RETRY_LAYER_SUFFIX = Pattern.compile(
|
"\\.[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\\.retry$");
|
// Codex 批次二三轮 #B+#C:认领名 owner-token 化——rename 目标 <src原名>.replaying.<ownerUUID>,
|
// 每次认领新生成 owner token。孤儿匹配/stableBaseName 剥离都以本模式识别 .replaying.<uuid> 层,
|
// 锚定字符串末尾($),UUID 严格 8-4-4-4-12 十六进制格式
|
private static final Pattern REPLAYING_TOKEN_SUFFIX = Pattern.compile(
|
"\\.replaying\\.[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$");
|
// Codex 四轮 Critical#1:AuditSpoolWriter 降级路径写后原子发布的临时名 <稳定base>.<UUID>.retry.writing——
|
// 崩溃/rename 失败会遗留孤儿,本模式供扫描排除与 recoverOrphanedWriting() 孤儿回收识别,
|
// 锚定字符串末尾($),UUID 严格 8-4-4-4-12 十六进制格式
|
private static final Pattern RETRY_WRITING_SUFFIX = Pattern.compile(
|
"\\.[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\\.retry\\.writing$");
|
// Codex 五轮:recoverOrphanedWriting() 隔离目标后缀——只隔离不发布不删除。任何以此结尾的文件名都不再
|
// 匹配 RETRY_WRITING_SUFFIX/REPLAYING_TOKEN_SUFFIX(二者均要求恰好以 .retry.writing / .replaying.<uuid>
|
// 收尾),也不以 .jsonl/.retry 结尾,故隔离后天然不会被本类任何扫描逻辑再次命中,见类 Javadoc"Codex 五轮修复"段
|
private static final String STALE_SUFFIX = ".stale";
|
// Codex 批次二三轮 #D:destroy 时等待重放线程退出的上限
|
private static final long SHUTDOWN_AWAIT_MS = 5_000L;
|
|
private final AuditEventPublisher publisher;
|
private final File spoolDir;
|
private final boolean enabled;
|
private final AuditSpoolLock spoolLock;
|
// 调度周期:默认 DEFAULT_REPLAY_INTERVAL_MS(5 分钟),可经 audit.spool-replay-interval-ms 覆盖。
|
// 仅影响"多久轮询一次 spool",不影响孤儿年龄阈值(ORPHAN_REPLAYING_AGE_MS,固定安全底线,见其注释)。
|
// 生产默认不动;测试环境可调短(如 15s)加速降级链恢复重放验证。
|
private final long replayIntervalMs;
|
private ScheduledExecutorService scheduler;
|
|
public AuditSpoolReplayer(AuditEventPublisher publisher, String spoolDir, boolean enabled,
|
AuditSpoolLock spoolLock) {
|
this(publisher, spoolDir, enabled, spoolLock, DEFAULT_REPLAY_INTERVAL_MS);
|
}
|
|
public AuditSpoolReplayer(AuditEventPublisher publisher, String spoolDir, boolean enabled,
|
AuditSpoolLock spoolLock, long replayIntervalMs) {
|
this.publisher = publisher;
|
this.spoolDir = new File(spoolDir);
|
this.enabled = enabled;
|
this.spoolLock = spoolLock;
|
// 兜底:非正数(配置笔误/0)回落默认周期,绝不让 scheduleWithFixedDelay 收到非法周期
|
this.replayIntervalMs = replayIntervalMs > 0 ? replayIntervalMs : DEFAULT_REPLAY_INTERVAL_MS;
|
}
|
|
@Override
|
public void afterPropertiesSet() {
|
if (!enabled) {
|
log.info("[audit] spool replayer 已禁用 (audit.spool-replay-enabled=false)");
|
return;
|
}
|
scheduler = Executors.newSingleThreadScheduledExecutor(r -> {
|
Thread t = new Thread(r, "audit-spool-replayer");
|
t.setDaemon(true);
|
return t;
|
});
|
// initialDelay 取与周期一致:避免启动瞬间与业务争 IO;遗留 spool 至多延迟一个周期被重放
|
scheduler.scheduleWithFixedDelay(this::replaySafely,
|
replayIntervalMs, replayIntervalMs, TimeUnit.MILLISECONDS);
|
}
|
|
@Override
|
public void destroy() {
|
if (scheduler == null) {
|
return;
|
}
|
// Codex 批次二三轮 #D:shutdownNow 后等待重放线程真正退出,避免"停机返回但重放线程仍在锁内/
|
// 正调 AuditSpoolLock 收尾"的窗口——与 AuditSpoolLock.destroy 的终止位配合优雅退出。
|
scheduler.shutdownNow();
|
try {
|
if (!scheduler.awaitTermination(SHUTDOWN_AWAIT_MS, TimeUnit.MILLISECONDS)) {
|
log.warn("[audit] spool 重放线程 {}ms 内未退出(可能卡在 IO/网络重放中)", SHUTDOWN_AWAIT_MS);
|
}
|
} catch (InterruptedException ie) {
|
Thread.currentThread().interrupt(); // 恢复中断标记后返回
|
}
|
}
|
|
private void replaySafely() {
|
try {
|
replayOnce();
|
} catch (Throwable t) {
|
// 调度线程抛异常会导致 scheduleWithFixedDelay 停摆,后续轮次不再执行——必须吞住
|
log.error("[audit] spool 重放轮次异常", t);
|
}
|
}
|
|
void replayOnce() {
|
if (!spoolDir.isDirectory()) {
|
return;
|
}
|
recoverOrphanedReplaying();
|
// Codex 四轮 Critical#1:孤儿 .retry.writing 回收——AuditSpoolWriter 降级路径写后原子发布中断
|
// (崩溃/rename 失败)遗留的临时文件,扫描集固定排除 .writing,永远不会被下面的主流程认领
|
recoverOrphanedWriting();
|
// 扫描集只含 .jsonl/.retry,显式排除 spool.lock(锁文件永不参与扫描/重放)与 .writing
|
// (Codex 四轮 Critical#1:写后原子发布的临时文件,内容可能不完整,绝不能被本扫描直接命中;
|
// Codex 五轮起,超龄 .writing 孤儿只会被 recoverOrphanedWriting() 隔离为 .stale,不再回归本扫描集)。
|
// .stale(隔离区,见类 Javadoc"Codex 五轮修复"段)天然不以 .jsonl/.retry 结尾,被本谓词排除,
|
// 无需额外条件
|
File[] targets = spoolDir.listFiles((dir, name) ->
|
!name.equals(AuditSpoolLock.LOCK_FILE_NAME)
|
&& !name.endsWith(".writing")
|
&& (name.endsWith(".jsonl") || name.endsWith(".retry")));
|
if (targets == null) {
|
return;
|
}
|
long now = System.currentTimeMillis();
|
for (File src : targets) {
|
// Codex 批次二 #1:60 秒静默窗口——距今不足 60 秒的 .jsonl 视为写入端可能活跃,本轮跳过,
|
// 下一轮再认领(<0.5 QPS,延迟一周期无害)。锁协议落地后此跳过仅为优化(减少活跃文件搅动),
|
// 即便不跳过、正确性也由 AuditSpoolLock 的 append↔rename 互斥保证。.retry 由本重放器自建、
|
// 无外部写入端,不受此限;.bad 本就不在扫描集(且从不重放)。
|
if (src.getName().endsWith(".jsonl")
|
&& now - src.lastModified() < QUIET_WINDOW_MS) {
|
continue;
|
}
|
replayFile(src);
|
}
|
}
|
|
/**
|
* 内审 Important 修复:从文件名剥离<b>所有</b> {@code .<UUID>.retry} 层与 {@code .replaying.<UUID>} token 后缀,
|
* 还原出稳定 base name(如 {@code audit-events-20260723.jsonl}),见类注释"内审 Important 修复"段。
|
*
|
* <p>剥离顺序:先去掉末尾 owner-token 化的 {@code .replaying.<UUID>} 后缀(认领时最后追加,
|
* 是最外层,最多一层,用 {@link #REPLAYING_TOKEN_SUFFIX} 匹配),再去掉 {@link AuditSpoolWriter}
|
* 降级路径写后原子发布的临时层 {@code .<UUID>.retry.writing}(Codex 四轮 Critical#1,
|
* {@link #RETRY_WRITING_SUFFIX} 匹配,同样最多一层、与 replaying token 互斥不共存),最后循环剥离
|
* 末尾的 {@code .<UUID>.retry} 层,直到不再匹配为止(病态场景下可能有多层,历史累积的每一层都会被剥掉)。
|
* UUID 段用 {@link #RETRY_LAYER_SUFFIX}/{@link #REPLAYING_TOKEN_SUFFIX}/{@link #RETRY_WRITING_SUFFIX}
|
* 严格匹配标准 8-4-4-4-12 十六进制格式并锚定字符串末尾,不会误剥业务文件名本身自带的点段。
|
*
|
* @param name 原始文件名(可能已带若干层 retry 后缀 + 一层 replaying token 或一层 writing 临时层,
|
* 也可能是从未重放过的原始文件名)
|
* @return 剥离干净的稳定 base name
|
*/
|
private String stableBaseName(String name) {
|
String base = name;
|
Matcher replayingMatcher = REPLAYING_TOKEN_SUFFIX.matcher(base);
|
if (replayingMatcher.find()) {
|
base = base.substring(0, replayingMatcher.start());
|
}
|
// Codex 四轮 Critical#1:.retry.writing 只会是 AuditSpoolWriter 降级路径追加的最外层(写后原子
|
// 发布的临时名,从不嵌套),与 .replaying.<uuid> 互斥同时出现,故也只剥最多一层
|
Matcher writingMatcher = RETRY_WRITING_SUFFIX.matcher(base);
|
if (writingMatcher.find()) {
|
base = base.substring(0, writingMatcher.start());
|
}
|
Matcher matcher;
|
while ((matcher = RETRY_LAYER_SUFFIX.matcher(base)).find()) {
|
base = base.substring(0, matcher.start());
|
}
|
return base;
|
}
|
|
/**
|
* 裁定修复 D:回收崩溃遗留的 .replaying 孤儿文件,见类注释。
|
*
|
* <p>只对"修改时间早于 now - {@link #ORPHAN_REPLAYING_AGE_MS}"的 {@code .replaying.<uuid>} 生效——
|
* 年龄阈值防止误抢并发实例(或本实例本轮稍早)正在处理中的文件;检查+rename 在
|
* {@link AuditSpoolLock} 临界区(锁协议临界区 5)完成,交回 {@code .retry} 通道由本轮或下一轮
|
* {@link #replayFile(File)} 正常消化,本方法不直接读取/处理文件内容。
|
*
|
* <p><b>Codex 批次二三轮 #B(锁内重验)</b>:锁外 {@code listFiles} 预筛只作候选收集;预筛与锁内
|
* rename 之间,该路径可能已被删除并复用为一个<b>新</b>文件(同名巧合),若沿用预筛结果直接搬走会把
|
* 新文件误当旧孤儿。故锁内 rename 前<b>重新确认</b> orphan 仍存在且 mtime 仍超阈值,才执行 rename。
|
*/
|
private void recoverOrphanedReplaying() {
|
File[] orphans = spoolDir.listFiles((dir, name) ->
|
!name.equals(AuditSpoolLock.LOCK_FILE_NAME)
|
&& REPLAYING_TOKEN_SUFFIX.matcher(name).find());
|
if (orphans == null || orphans.length == 0) {
|
return;
|
}
|
long now = System.currentTimeMillis();
|
for (final File orphan : orphans) {
|
if (now - orphan.lastModified() < ORPHAN_REPLAYING_AGE_MS) {
|
continue; // 锁外预筛:未超龄,可能正被处理中,本轮不碰
|
}
|
String orphanName = orphan.getName();
|
// 内审 Important 修复:先归一到稳定 base 再拼恰好一层新 UUID,防止孤儿在已累积多层
|
// 后缀的名字上继续追加、文件名无界增长
|
final File retry = new File(orphan.getParentFile(),
|
stableBaseName(orphanName) + "." + UUID.randomUUID() + ".retry");
|
final boolean[] moved = {false};
|
try {
|
// 锁协议临界区 5:检查+rename 在锁内,与并发实例/本进程认领互斥;tryLock 拿不到则本轮跳过该孤儿
|
boolean lockRun = spoolLock.tryLockAndRun(() -> {
|
// Codex#B 锁内重验:预筛后路径可能已被删除并复用为新文件,重新确认仍存在且仍超龄才搬走
|
if (!orphan.exists()) {
|
return;
|
}
|
if (System.currentTimeMillis() - orphan.lastModified() < ORPHAN_REPLAYING_AGE_MS) {
|
return;
|
}
|
Files.move(orphan.toPath(), retry.toPath(), StandardCopyOption.ATOMIC_MOVE);
|
moved[0] = true;
|
});
|
if (lockRun && moved[0]) {
|
log.warn("[audit] 回收孤儿 .replaying 文件: {} -> {}", orphanName, retry.getName());
|
}
|
// lockRun=false(他进程持锁)或 moved=false(锁内重验未过):本轮跳过该孤儿,下轮再判
|
} catch (IOException moveFail) {
|
// 改名失败:可能已被其他实例/本轮稍早处理,下一轮再判一次
|
log.warn("[audit] 孤儿 .replaying 回收改名失败,下轮再试: {}", orphanName);
|
}
|
}
|
}
|
|
/**
|
* Codex 五轮修复:回收 {@link AuditSpoolWriter} 降级路径「写后原子发布」中断遗留的
|
* {@code <稳定base>.<UUID>.retry.writing} 孤儿——<b>只隔离,永不发布为 .retry,永不删除</b>。
|
*
|
* <p>四轮版本曾把超龄孤儿直接改名回归 {@code .retry} 通道当正常事件重放,Codex 五轮复验指出这不安全:
|
* 超龄(mtime 早于 {@link #ORPHAN_REPLAYING_AGE_MS})只能说明"这个文件很久没有被更新",不能证明
|
* writer 已经死亡——它可能仍卡在那唯一一次 {@link java.nio.file.Files#write} 调用内部(磁盘/网络
|
* 存储 IO 阻塞),FD 仍打开、进程仍存活。若照旧回归 {@code .retry} 通道,重放器会认领→读到空文件
|
* (或崩溃语义下的半行前缀)→判定处理完毕→删除;卡住的 write 恢复后写向已 unlink 的 inode,数据
|
* 静默消失,且不留任何错误日志。完整论证见类 Javadoc"Codex 五轮修复(.writing 隔离区)"段。
|
*
|
* <p>修复后动作:锁内重验(exists + 仍超龄,防止误抢正在写入中的文件)通过后,把 {@code .retry.writing}
|
* <b>原子改名为 {@code <原名>.stale}</b>(同目录)并 {@code log.error} 告警(脱敏:仅文件名/大小/mtime,
|
* 不含内容)——不解析内容、不写回任何可被扫描/认领的名字、不删除。{@code rename} 不 unlink:即便该
|
* writer 随后真的恢复并把内容写完,也是写进同一个 inode(现名 {@code .stale}),内容不会因改名而
|
* 丢失,只是滞留隔离区待人工核查补录(对照 {@code clientEventId} 去重,落库天然幂等)。writer 侧针对
|
* "发布源被隔离"的自愈见 {@link AuditSpoolWriter} 类 Javadoc。
|
*/
|
private void recoverOrphanedWriting() {
|
// RETRY_WRITING_SUFFIX 只匹配以 ".<uuid>.retry.writing" 结尾的文件名——本方法产出的 .stale
|
// 文件不再以此结尾,天然不会被下一轮扫描重新命中(见类 Javadoc"Codex 五轮修复"段)
|
File[] orphans = spoolDir.listFiles((dir, name) ->
|
!name.equals(AuditSpoolLock.LOCK_FILE_NAME)
|
&& RETRY_WRITING_SUFFIX.matcher(name).find());
|
if (orphans == null || orphans.length == 0) {
|
return;
|
}
|
long now = System.currentTimeMillis();
|
for (final File orphan : orphans) {
|
if (now - orphan.lastModified() < ORPHAN_REPLAYING_AGE_MS) {
|
continue; // 锁外预筛:未超龄,可能正被 writer 写入中(或卡在 Files.write 内部),本轮不碰
|
}
|
String orphanName = orphan.getName();
|
// Codex 五轮:只隔离,原名后追加 .stale——不发布为 .retry(不会被重放器认领)、不删除
|
// (.stale 排除在一切扫描之外,见类 Javadoc"Codex 五轮修复"段的完整论证)
|
final File stale = new File(orphan.getParentFile(), orphanName + STALE_SUFFIX);
|
final long[] quarantinedSize = {-1L};
|
final long[] quarantinedMtime = {-1L};
|
final boolean[] moved = {false};
|
try {
|
// 锁内重验(同 recoverOrphanedReplaying 的 #B 模式):预筛后、锁内 rename 前,路径可能
|
// 已被 writer 完成发布并删除、或被别的实例先一步隔离,重新确认仍存在且仍超龄才动手
|
boolean lockRun = spoolLock.tryLockAndRun(() -> {
|
if (!orphan.exists()) {
|
return;
|
}
|
if (System.currentTimeMillis() - orphan.lastModified() < ORPHAN_REPLAYING_AGE_MS) {
|
return;
|
}
|
quarantinedSize[0] = orphan.length();
|
quarantinedMtime[0] = orphan.lastModified();
|
Files.move(orphan.toPath(), stale.toPath(), StandardCopyOption.ATOMIC_MOVE);
|
moved[0] = true;
|
});
|
if (lockRun && moved[0]) {
|
// log.error 而非 warn:这是需要人工核查补录的隔离事件,不是可自愈的瞬时状况
|
log.error("[audit] 超龄 .retry.writing 隔离至 .stale(不自动发布/不自动删除,"
|
+ "需人工核查补录,见 AuditSpoolReplayer 类 Javadoc"
|
+ "\".writing 隔离区处置\"说明段): file={} size={}bytes mtime={}",
|
stale.getName(), quarantinedSize[0], quarantinedMtime[0]);
|
}
|
// lockRun=false(他进程持锁)或 moved=false(锁内重验未过):本轮跳过该孤儿,下轮再判
|
} catch (IOException moveFail) {
|
// 改名失败:可能已被其他实例/本轮稍早处理,下一轮再判一次
|
log.warn("[audit] 孤儿 .retry.writing 隔离改名失败,下轮再试: {}", orphanName);
|
}
|
}
|
}
|
|
private void replayFile(final File src) {
|
// 锁协议临界区 2(认领):Codex 批次二 #1/#2/#B/#C。
|
// #B+#C 认领名 owner-token 化——每次认领新生成 ownerToken,rename 目标 <src原名>.replaying.<ownerToken>。
|
// 收尾/删除只操作本 owner token 的路径:陈旧 owner(被抢后收尾)与新 owner 永不同名,
|
// deleteIfExists 物理上不可能误删他 owner 刚认领的同名文件(消除 Codex#C 固定名跨 owner 冲突)。
|
final File replaying = new File(src.getParentFile(),
|
src.getName() + REPLAYING_SUFFIX + "." + UUID.randomUUID());
|
final boolean isJsonl = src.getName().endsWith(".jsonl");
|
final boolean[] renamed = {false};
|
// rename→.replaying.<owner> 与 setLastModified 刷新租约在同一临界区(进程内 monitor + 跨进程 FileLock):
|
// #1:与 writer 的 append 互斥,消除"取走文件瞬间"与"正在 append"的窗口竞态;
|
// #2:rename 与刷新原子成对,消除"rename 成功后、刷新 mtime 前被孤儿回收误抢"的窗口。
|
try {
|
boolean lockRun = spoolLock.tryLockAndRun(() -> {
|
// Codex#B 锁内重验:锁外预筛(quiet window)与锁内之间,src 可能已被他 owner 取走/删除。
|
// 锁内重新确认 src 仍存在(.jsonl 且仍在静默窗口内说明写入端可能活跃,让给下轮)才认领。
|
if (!src.exists()) {
|
return;
|
}
|
if (isJsonl && System.currentTimeMillis() - src.lastModified() < QUIET_WINDOW_MS) {
|
return;
|
}
|
Files.move(src.toPath(), replaying.toPath(), StandardCopyOption.ATOMIC_MOVE);
|
// setLastModified 返回 false 仅告警不中断。Codex#C:owner-token 唯一名下,租约建立失败的最坏
|
// 后果=本文件超 10 分钟后被孤儿回收→重复重放(落库按 client_event_id 幂等吸收),不再有丢失路径。
|
if (!replaying.setLastModified(System.currentTimeMillis())) {
|
log.warn("[audit] 认领后刷新 spool 租约 mtime 失败(最坏后果=重复重放,幂等吸收): {}",
|
replaying.getName());
|
}
|
renamed[0] = true;
|
});
|
if (!lockRun) {
|
// FileLock 被别的进程持有 → 本轮跳过该文件,下轮再认领
|
return;
|
}
|
if (!renamed[0]) {
|
// 锁内重验未过(已被他 owner 取走/删除,或 .jsonl 仍活跃)→ 本轮跳过,下轮再判
|
return;
|
}
|
} catch (IOException moveFail) {
|
// 改名失败(正被他进程/他轮取走)→ 跳过,下轮再试
|
log.warn("[audit] spool 文件占用改名失败,跳过: {}", src.getName());
|
return;
|
}
|
|
// 锁协议临界区外(readAllLines):认领后 .replaying 已被本进程独占,读取不阻塞 writer
|
List<String> lines;
|
try {
|
lines = Files.readAllLines(replaying.toPath(), StandardCharsets.UTF_8);
|
} catch (IOException readFail) {
|
// 读失败不删 .replaying(避免读失败即丢数据),留待人工/下次进程处理
|
log.error("[audit] spool 读取失败,保留 .replaying: {}", replaying.getName(), readFail);
|
return;
|
}
|
|
// 锁协议临界区外(逐行 deliverOnce=网络慢操作):把行分类为 remaining(投递失败待重试)
|
// 与 badLines(反序列化失败待隔离),此循环不做任何 spool 文件写操作
|
List<String> remaining = new ArrayList<>();
|
List<String> badLines = new ArrayList<>();
|
for (int i = 0; i < lines.size(); i++) {
|
String line = lines.get(i);
|
if (line == null || line.trim().isEmpty()) {
|
continue;
|
}
|
AuditEventDTO event;
|
try {
|
event = JsonUtil.getJsonToBean(line, AuditEventDTO.class);
|
} catch (Exception parseFail) {
|
// Codex 批次二 #3:坏行留待收尾批量隔离到 .bad(不进 remaining、不静默丢)。
|
// Codex 批次二 #4:日志只记文件名/行号/行长/内容 SHA-256 前 12 位 + 异常类名,绝不输出行正文
|
// (parseFail 只取类名——Jackson 异常正文可能内嵌源片段,故不整体打印)。
|
badLines.add(line);
|
log.error("[audit] spool 坏行检出待隔离: file={} lineNo={} lineLen={} sha256_12={} cause={}",
|
replaying.getName(), i + 1, line.length(), sha256Prefix12(line),
|
parseFail.getClass().getName());
|
continue;
|
}
|
try {
|
publisher.deliverOnce(event); // 成功→丢弃该行;失败→保留待下轮
|
} catch (Exception deliverFail) {
|
// 好行投递失败(AuditDeliveryException)→ 保留待下轮,行为不变
|
remaining.add(line);
|
}
|
}
|
|
// 锁协议临界区 4(收尾):写 .bad / 写 .retry / 删 .replaying 在锁内(各自快操作)。
|
// 用 writer 语义 lockAndRun(始终执行):投递已完成,收尾必须落地,避免留 .replaying 被重复重放
|
try {
|
spoolLock.lockAndRun(() -> finishReplay(src, replaying, remaining, badLines));
|
} catch (IOException writeFail) {
|
log.error("[audit] spool 重放收尾写盘失败,保留 .replaying: {}", replaying.getName(), writeFail);
|
}
|
}
|
|
/**
|
* 锁协议临界区 4(收尾):批量写坏行 {@code .bad}、写剩余行 {@code .retry}、删 {@code .replaying}。
|
* 全程由 {@link AuditSpoolLock#lockAndRun} 保护,仅做快速文件操作。
|
*
|
* <p>Codex 批次二 #3:坏行批量隔离到 {@code <稳定 base>.<UUID>.bad}(不重放、不静默丢);写 {@code .bad}
|
* 失败则该批坏行退回 remaining(宁可重试也不丢)。剩余行(含投递失败 + 隔离失败退回的坏行)非空则写
|
* {@code <稳定 base>.<UUID>.retry}(终验四轮③:绝不写回活动原名),否则删 {@code .replaying};有 remaining
|
* 时也在 {@code .retry} 写成功后删 {@code .replaying}(不留孤儿)。
|
*
|
* @throws IOException 写 {@code .retry} 或删 {@code .replaying} 失败——调用方据此保留 {@code .replaying} 待下轮
|
*/
|
private void finishReplay(File src, File replaying, List<String> remaining, List<String> badLines)
|
throws IOException {
|
List<String> retryLines = new ArrayList<>(remaining);
|
if (!badLines.isEmpty()) {
|
try {
|
writeNewSpoolFile(src, ".bad", badLines);
|
} catch (IOException badFail) {
|
// 写 .bad 失败:坏行退回重试(宁可重试也不丢);日志不含正文,只记文件名/条数
|
retryLines.addAll(badLines);
|
log.warn("[audit] spool 坏行隔离写盘失败,退回重试: file={} count={}",
|
replaying.getName(), badLines.size(), badFail);
|
}
|
}
|
if (retryLines.isEmpty()) {
|
Files.deleteIfExists(replaying.toPath()); // 全部成功:删 .replaying
|
} else {
|
// 部分失败:剩余行写 <稳定 base>.<UUID>.retry 新文件,绝不写回活动原名(终验四轮③)
|
writeNewSpoolFile(src, ".retry", retryLines);
|
Files.deleteIfExists(replaying.toPath()); // .retry 写成功后删 .replaying,不留孤儿
|
}
|
}
|
|
/**
|
* 把 {@code lines} 写入一个新建的 {@code <稳定 base>.<UUID><suffix>} 文件({@code .retry} 或 {@code .bad})。
|
* 前缀先经 {@link #stableBaseName(String)} 归一,再拼恰好一层新 UUID,文件名长度有界(内审 Important 修复);
|
* {@link StandardOpenOption#CREATE_NEW} 保证不覆盖既有文件,UUID 使撞名概率可忽略。
|
*/
|
private void writeNewSpoolFile(File src, String suffix, List<String> lines) throws IOException {
|
File out = new File(src.getParentFile(),
|
stableBaseName(src.getName()) + "." + UUID.randomUUID() + suffix);
|
StringBuilder sb = new StringBuilder();
|
for (String line : lines) {
|
sb.append(line).append('\n');
|
}
|
Files.write(out.toPath(), sb.toString().getBytes(StandardCharsets.UTF_8),
|
StandardOpenOption.CREATE_NEW, StandardOpenOption.WRITE);
|
}
|
|
/**
|
* Codex 批次二 #4:坏行内容的 SHA-256 十六进制<b>前 12 位</b>,供日志定位坏行而不泄露正文
|
* (spec "ERROR 日志不含字段值正文")。SHA-256 为 JDK 标配算法,{@link NoSuchAlgorithmException}
|
* 理论不可达;万一不可达则回落到固定占位串,绝不回退成打印正文。
|
*/
|
private static String sha256Prefix12(String line) {
|
try {
|
MessageDigest md = MessageDigest.getInstance("SHA-256");
|
byte[] digest = md.digest(line.getBytes(StandardCharsets.UTF_8));
|
StringBuilder hex = new StringBuilder(12);
|
for (byte b : digest) {
|
hex.append(Character.forDigit((b >> 4) & 0xF, 16));
|
hex.append(Character.forDigit(b & 0xF, 16));
|
if (hex.length() >= 12) {
|
break;
|
}
|
}
|
return hex.substring(0, 12);
|
} catch (NoSuchAlgorithmException e) {
|
return "sha256-unavailable";
|
}
|
}
|
}
|