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 必须可恢复才算兜底)。 * *

周期性扫描 spool 目录的 {@code .jsonl} / {@code .retry} 文件 → 原子重命名加 * {@code .replaying.} 后缀(认领,token 每次新生成,见"#B+#C"段)→ 逐行反序列化走 * {@link AuditEventPublisher#deliverOnce(AuditEventDTO)} * (成功丢弃该行、抛异常保留)→ 全部成功删 {@code .replaying};部分失败把剩余行写入 * {@code <稳定 base 文件名>..retry} 新文件(base+UUID 保证同轮多文件失败不撞名,终验五轮⑤; * base 名的取法见下文"内审 Important 修复"段 {@link #stableBaseName(String)}), * {@code .retry} 写成功后删对应 {@code .replaying}(不留孤儿);绝不写回活动原名(终验四轮③: * 重命名期间写入端可能已在原名新建文件追加新事件,写回原名会覆盖丢数据)。落库幂等(client_event_id) * 使重放天然安全。 * *

调度手段:自建单线程 daemon {@link ScheduledExecutorService}(非 @Scheduled),不依赖宿主 * 是否开启 @EnableScheduling,生命周期由本 bean 自控(见 task-7-report 决策说明)。 * 由 {@code audit.spool-replay-enabled}(默认 true)开关。 * *

双层锁协议(Codex 批次二复验 #1/#2 终局修复):与 {@link AuditSpoolWriter} 共享同一 * {@link AuditSpoolLock}(进程内 monitor + 跨进程 {@code spool.lock} FileLock)。快速文件操作在锁内、 * 慢速网络重放在锁外,临界区划分: *

    *
  1. 认领({@link #replayFile}):{@code rename→.replaying.} + {@code setLastModified} * 刷新租约在同一临界区——#1 与 writer 的 append 互斥(消除"取走文件瞬间"与"正在 append"的窗口 * 竞态);#2 rename 与刷新原子成对(消除"rename 成功后、刷新 mtime 前被孤儿回收误抢"的窗口)。 * 临界区内先锁内重验(#B)再 rename。replayer 路径 {@code tryLock} 拿不到就跳过本轮,下轮再来。
  2. *
  3. readAllLines:锁外——认领后 {@code .replaying} 已不在任何写入端/扫描集视线内(孤儿回收 * 10 分钟阈值 + 刚刷新的 mtime 保护),读取安全且不阻塞 writer。
  4. *
  5. 逐行 deliverOnce:锁外(网络慢操作)。
  6. *
  7. 收尾({@link #finishReplay}):写 {@code .bad} / 写 {@code .retry} / 删 {@code .replaying} * 在锁内(各自快操作)。
  8. *
  9. 孤儿回收 rename({@link #recoverOrphanedReplaying}):检查+rename 在锁内。
  10. *
* 60 秒静默窗口({@link #QUIET_WINDOW_MS})保留为纯优化(减少活跃文件搅动),非正确性依赖—— * 正确性已由锁协议保证。残余:单文件重放超 10 分钟仍可能被其它进程当孤儿抢走 → 后果=重复重放, * 落库按 client_event_id 幂等保证不丢不重入库,可接受;{@code spool.lock} 永不参与扫描/重放/删除 * (扫描过滤已按 {@link AuditSpoolLock#LOCK_FILE_NAME} 排除)。 * *

裁定修复 D(孤儿 .replaying 恢复):进程若在"改名成 .replaying"之后、"处理完"之前崩溃, * 会留下孤儿——扫描集固定为 {@code {.jsonl,.retry}},永远不会认领它,事件永久滞留。每轮扫描时, * 额外把修改时间早于(当前时刻 - 10 分钟,即两个调度周期)的 .replaying 文件原子改名为 * {@code <稳定 base 文件名>..retry}(锁内),回归 retry 通道由下一轮自然消化;不在回收动作 * 内直接处理文件内容,避免与"本实例/并发实例正在重放中的文件"混淆。年龄阈值即为防止误抢:正常一轮 * 处理远快于一个调度周期,若 .replaying 存活跨越两个周期才可判定为孤儿。落库按 client_event_id * 幂等,重复重放安全。 * *

Codex 批次二三轮 #B+#C(认领名 owner-token 化 + 锁内重验):认领 rename 目标由固定后缀 * {@code .replaying} 改为 owner-token 化的 {@code .replaying.}(token 每次认领新生成, * {@link #REPLAYING_TOKEN_SUFFIX} 识别)。收尾/删除只操作本 owner token 的路径——陈旧 owner(被抢后收尾) * 与新 owner 永不同名,{@code deleteIfExists} 物理上不可能删掉新 owner 刚认领的文件(消除 #C 固定名跨 * owner 冲突)。同时认领({@link #replayFile})与孤儿回收({@link #recoverOrphanedReplaying})都改为 * 锁内先重验(exists + mtime 仍超阈/仍静默)再 rename——锁外 {@code listFiles} 预筛只作候选收集, * 消除"预筛后、锁内 rename 前,路径被删除并复用为新文件 → 旧检查结果误搬新文件"的窗口(#B)。 * {@code setLastModified} 失败保留 warn+继续:owner-token 唯一名下租约失败的最坏后果=被孤儿回收→重复重放 * (幂等吸收),不再有丢失路径。 * *

Codex 批次二三轮 #D(优雅退出):{@link #destroy()} 在 {@code shutdownNow} 后 * {@code awaitTermination}({@link #SHUTDOWN_AWAIT_MS}) 等待重放线程真正退出(超时 warn),与 * {@link AuditSpoolLock} 的终止位配合,避免"停机返回但重放线程仍在锁内收尾"。 * *

内审 Important 修复(retry 文件名收敛,防 NAME_MAX 无界增长):部分失败写 {@code .retry}、 * 孤儿回收改名、坏行隔离写 {@code .bad},前缀都先经 {@link #stableBaseName(String)} 从文件名剥离所有 * {@code ..retry} 层与 {@code .replaying.} token 后缀还原出稳定 base(如 {@code audit-events-20260723.jsonl}), * 再拼恰好一层新 UUID,文件名长度从此有界,不随失败轮次增长。UUID 段用 {@link #RETRY_LAYER_SUFFIX} * 严格匹配标准 8-4-4-4-12 十六进制格式并锚定字符串末尾,不会误剥业务文件名自身携带的点段(如 {@code .jsonl})。 * (Codex 五轮起:{@link #recoverOrphanedWriting()} 的隔离改名不再走本机制——直接在原名后追加 * {@code .stale},见其专属段;{@code .retry.writing} 层剥离分支保留作通用防御,当前无调用方会触发。) * *

坏行隔离与脱敏(Codex 批次二 #3/#4):反序列化失败的坏行不静默丢、不进 remaining,收尾时批量 * 隔离进 {@code <稳定 base>..bad}(不参与任何自动重放,留人工修复补录,见 {@link #finishReplay}); * 坏行 ERROR 日志只记文件名/行号/行长/内容 SHA-256 前 12 位({@link #sha256Prefix12}),绝不输出行内容正文。 * *

写后原子发布孤儿回收(Codex 四轮 Critical#1):{@link AuditSpoolWriter} 降级路径改为「先写 * 唯一名 {@code .writing} 临时文件、再原子改名发布为 {@code .retry}」(消灭跨进程半写窗口,见其类 * Javadoc)。若 writer 恰好崩溃于两步之间(或改名本身因极端 IO 异常失败),会遗留 * {@code <稳定base>..retry.writing} 孤儿——扫描集固定排除 {@code .writing}(见 * {@link #replayOnce()}),主流程永远不会认领它。每轮额外调用 {@link #recoverOrphanedWriting()} * 处理这类孤儿;具体动作见下一段——Codex 五轮复验推翻了本段最初"回归 .retry 通道当正常事件重放"的方案。 * *

Codex 五轮修复(.writing 隔离区,不能仅凭年龄自动发布):四轮版本曾把超龄 * {@code .retry.writing} 直接改名回归 {@code .retry} 通道当正常事件重放——Codex 五轮复验指出这不安全: * degraded writer(未持 FileLock 的降级路径)若恰好在其唯一一次 {@link Files#write} 调用内部 * 被操作系统/磁盘 IO 阻塞超过 10 分钟,该 writer 线程仍然存活、仍持有该文件的打开 FD,只是尚未返回; * 此时它的 {@code .writing} 文件在孤儿回收眼中已"超龄",若照旧改名发布为 {@code .retry},重放器会把它 * 当正常事件认领(rename→{@code .replaying})、读到空文件(内容还没真正落盘)、判定"处理完毕"并删除 * {@code .replaying};随后卡住的 write 调用恢复,向着已被 unlink 的 inode 追加——数据静默消失,不留 * 任何错误痕迹。半行截断时同理:坏行隔离机制只能吸收"取走时刻"已经写入的前缀,写调用尚未完成的后半段 * 一样会丢。 * *

结论:超龄不足以断言 writer 已死,{@link #recoverOrphanedWriting()} 因此只做隔离、 * 绝不发布、绝不删除——锁内重验通过后把 {@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.},{@link #recoverOrphanedWriting()} * 自身只认 {@code .retry.writing},三者的后缀判定都不会匹配以 {@code .stale} 收尾的文件名,故隔离后的 * 文件物理上不可能被本类任何逻辑再次触碰;也无需纳入 {@link #stableBaseName(String)} 的剥层范围—— * {@code .stale} 从不参与任何"新建 retry/bad 文件"的前缀推导。 * *

安全性论证的核心是 POSIX {@code rename(2)} 的语义:rename 不 unlink——它只改变目录项指向的 * 名字,被卡住的 write 调用持有的是对 inode 的 FD,与路径名无关,恢复后仍会把内容完整写入同一个 * inode(现在挂在 {@code .stale} 这个名字下),不会因改名而丢失或写偏。事件因此保全在隔离区,等待 * 人工核查后手工补录(对照 {@code clientEventId} 去重);隔离区本身永不被任何自动化流程删除,因此 * "重放器把仍在写的文件当孤儿抢走→空读→删除→静默丢失"这条路径被物理消除(不是靠时序窗口收窄, * 而是代码上根本不存在能删除 {@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. 年龄阈值 = 两个默认调度周期(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 修复:精确匹配"尾部恰好一层 .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 目标 .replaying., // 每次认领新生成 owner token。孤儿匹配/stableBaseName 剥离都以本模式识别 .replaying. 层, // 锚定字符串末尾($),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>..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. // 收尾),也不以 .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 修复:从文件名剥离所有 {@code ..retry} 层与 {@code .replaying.} token 后缀, * 还原出稳定 base name(如 {@code audit-events-20260723.jsonl}),见类注释"内审 Important 修复"段。 * *

剥离顺序:先去掉末尾 owner-token 化的 {@code .replaying.} 后缀(认领时最后追加, * 是最外层,最多一层,用 {@link #REPLAYING_TOKEN_SUFFIX} 匹配),再去掉 {@link AuditSpoolWriter} * 降级路径写后原子发布的临时层 {@code ..retry.writing}(Codex 四轮 Critical#1, * {@link #RETRY_WRITING_SUFFIX} 匹配,同样最多一层、与 replaying token 互斥不共存),最后循环剥离 * 末尾的 {@code ..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. 互斥同时出现,故也只剥最多一层 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 孤儿文件,见类注释。 * *

只对"修改时间早于 now - {@link #ORPHAN_REPLAYING_AGE_MS}"的 {@code .replaying.} 生效—— * 年龄阈值防止误抢并发实例(或本实例本轮稍早)正在处理中的文件;检查+rename 在 * {@link AuditSpoolLock} 临界区(锁协议临界区 5)完成,交回 {@code .retry} 通道由本轮或下一轮 * {@link #replayFile(File)} 正常消化,本方法不直接读取/处理文件内容。 * *

Codex 批次二三轮 #B(锁内重验):锁外 {@code listFiles} 预筛只作候选收集;预筛与锁内 * rename 之间,该路径可能已被删除并复用为一个文件(同名巧合),若沿用预筛结果直接搬走会把 * 新文件误当旧孤儿。故锁内 rename 前重新确认 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>..retry.writing} 孤儿——只隔离,永不发布为 .retry,永不删除。 * *

四轮版本曾把超龄孤儿直接改名回归 {@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 隔离区)"段。 * *

修复后动作:锁内重验(exists + 仍超龄,防止误抢正在写入中的文件)通过后,把 {@code .retry.writing} * 原子改名为 {@code <原名>.stale}(同目录)并 {@code log.error} 告警(脱敏:仅文件名/大小/mtime, * 不含内容)——不解析内容、不写回任何可被扫描/认领的名字、不删除。{@code rename} 不 unlink:即便该 * writer 随后真的恢复并把内容写完,也是写进同一个 inode(现名 {@code .stale}),内容不会因改名而 * 丢失,只是滞留隔离区待人工核查补录(对照 {@code clientEventId} 去重,落库天然幂等)。writer 侧针对 * "发布源被隔离"的自愈见 {@link AuditSpoolWriter} 类 Javadoc。 */ private void recoverOrphanedWriting() { // RETRY_WRITING_SUFFIX 只匹配以 "..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 目标 .replaying.。 // 收尾/删除只操作本 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. 与 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 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 remaining = new ArrayList<>(); List 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} 保护,仅做快速文件操作。 * *

Codex 批次二 #3:坏行批量隔离到 {@code <稳定 base>..bad}(不重放、不静默丢);写 {@code .bad} * 失败则该批坏行退回 remaining(宁可重试也不丢)。剩余行(含投递失败 + 隔离失败退回的坏行)非空则写 * {@code <稳定 base>..retry}(终验四轮③:绝不写回活动原名),否则删 {@code .replaying};有 remaining * 时也在 {@code .retry} 写成功后删 {@code .replaying}(不留孤儿)。 * * @throws IOException 写 {@code .retry} 或删 {@code .replaying} 失败——调用方据此保留 {@code .replaying} 待下轮 */ private void finishReplay(File src, File replaying, List remaining, List badLines) throws IOException { List 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>..retry 新文件,绝不写回活动原名(终验四轮③) writeNewSpoolFile(src, ".retry", retryLines); Files.deleteIfExists(replaying.toPath()); // .retry 写成功后删 .replaying,不留孤儿 } } /** * 把 {@code lines} 写入一个新建的 {@code <稳定 base>.} 文件({@code .retry} 或 {@code .bad})。 * 前缀先经 {@link #stableBaseName(String)} 归一,再拼恰好一层新 UUID,文件名长度有界(内审 Important 修复); * {@link StandardOpenOption#CREATE_NEW} 保证不覆盖既有文件,UUID 使撞名概率可忽略。 */ private void writeNewSpoolFile(File src, String suffix, List 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 十六进制前 12 位,供日志定位坏行而不泄露正文 * (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"; } } }