package jnpf.audit.sdk; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.DisposableBean; import java.io.Closeable; import java.io.File; import java.io.IOException; import java.io.RandomAccessFile; import java.nio.channels.FileChannel; import java.nio.channels.FileLock; import java.nio.channels.OverlappingFileLockException; import java.nio.file.Files; /** * spool 双层锁协议(Codex 批次二复验 #1/#2 终局修复;批次二三轮 #A/#D 加固)。 * *

{@link AuditSpoolWriter} 与 {@link AuditSpoolReplayer} 共享同一实例(由 * {@code AuditSdkAutoConfiguration} 注入单例),把彼此的 spool 文件操作串行化,从根上消除两条竞态: *

* 两条竞态的正确性从此由本锁保证(不再依赖 60 秒静默窗口这种活跃度启发式;静默窗口降级为纯优化)。 * *

两层结构

*
    *
  1. 进程内层:单一监视器 {@link #monitor},writer/replayer 所有临界区共用,保证同 JVM 串行。 * 它同时保证同 JVM 内 {@link FileChannel#tryLock()} 请求永不重叠——否则 JDK 会抛 * {@link OverlappingFileLockException}(FileLock 是 JVM 级而非线程级),故本类先取进程内锁、 * 在 {@code synchronized} 块内获/放 FileLock,天然无重叠。
  2. *
  3. 跨进程层:spool 目录下常驻的 {@code spool.lock} 文件 + {@link FileChannel#lock()} 系 * 独占 {@link FileLock}。锁文件 {@link RandomAccessFile}/{@link FileChannel} 惰性打开、常驻, * {@link DisposableBean#destroy()} 时释放。
  4. *
* *

降级与终局(绝不让业务/发送线程崩溃)

* * *

残余声明:单文件重放(readAllLines→网络 deliver→收尾)超 10 分钟仍可能被其它进程当孤儿 * 抢走 → 后果=重复重放,落库按 client_event_id 幂等保证不丢不重入库,可接受。{@code spool.lock} * 自身永不参与扫描/重放/删除({@link AuditSpoolReplayer} 扫描过滤已按 {@link #LOCK_FILE_NAME} 排除)。 */ @Slf4j public class AuditSpoolLock implements DisposableBean { /** 跨进程锁文件名——常驻 spool 目录,永不参与扫描/重放/删除。 */ static final String LOCK_FILE_NAME = "spool.lock"; /** writer 路径 FileLock 获取的最长等待(带重试);超时后降级为唯一名应急段。 */ private static final long WRITER_LOCK_WAIT_MS = 5_000L; /** 重试轮询间隔。 */ private static final long RETRY_SLEEP_MS = 50L; /** * 进程内层:writer/replayer 所有临界区共用的单一监视器,保证同 JVM 串行; * 亦保证同 JVM 内 FileLock 请求永不重叠(否则抛 {@link OverlappingFileLockException})。 */ private final Object monitor = new Object(); private final File spoolDir; private final File lockFile; /** 跨进程层常驻句柄,destroy 时释放;仅在持有 {@link #monitor} 时读写。 */ private RandomAccessFile raf; private FileChannel channel; /** 降级标志:FileLock 不支持/锁文件打不开时置位,之后不再尝试跨进程锁;仅在持有 {@link #monitor} 时读写。 */ private boolean fileLockDisabled = false; /** warn 去重标志;仅在持有 {@link #monitor} 时读写。 */ private boolean degradeWarned = false; /** Codex#D:终止位——destroy 后拒绝重开 channel,迟到调用只走进程内锁。仅在持有 {@link #monitor} 时写。 */ private volatile boolean destroyed = false; /** * 已销毁后迟到调用的 warn 去重标志;仅在持有 {@link #monitor} 时读写。两条调用路径共用一个去重开关: * {@link #acquireFileLockWithTimeout()}(writer/{@code lockAndRun} 路径,destroyed → 返回 {@code null} * → 调用方走 degradedAction,仅进程内锁执行)与 {@link #tryLockAndRun}(replayer/认领路径,destroyed → * 直接返回 {@code false}、不执行 action——Codex 四轮 Critical#1)。 */ private boolean postDestroyWarned = false; public AuditSpoolLock(String spoolDir) { this.spoolDir = new File(spoolDir); this.lockFile = new File(this.spoolDir, LOCK_FILE_NAME); } /** 可抛 {@link IOException} 的无返回回调(JDK8 函数式接口)。 */ @FunctionalInterface public interface IoRunnable { void run() throws IOException; } /** * writer 语义(Codex#A 降级隔离):始终执行,但按是否持有跨进程锁二选一执行。 * *

顺序:进程内 {@link #monitor} → 尽力获 FileLock(带超时重试,共最多 {@link #WRITER_LOCK_WAIT_MS})。 *

* * @param lockedAction 持有跨进程锁时执行的动作(共享活动文件追加) * @param degradedAction 未持锁时执行的降级动作(唯一名应急段写入) * @throws IOException 所选动作抛出的 IO 异常,由调用方按原有兜底处理 */ void lockAndRun(IoRunnable lockedAction, IoRunnable degradedAction) throws IOException { synchronized (monitor) { FileLock fileLock = acquireFileLockWithTimeout(); if (fileLock == null) { // Codex#A:未持有跨进程锁(超时/中断/不可用/已销毁)→ 走降级动作(写唯一名,消灭共享) degradedAction.run(); return; } try { lockedAction.run(); } finally { releaseQuietly(fileLock); } } } /** * writer 语义便捷重载:降级动作与持锁动作相同。 * *

适用于收尾等"操作对象本就是唯一名/owner-token 名"的调用方(如 {@link AuditSpoolReplayer} 的 * {@code finishReplay}:写 {@code .retry}/{@code .bad}、删 owner-token 的 {@code .replaying.} * 都是唯一名,仅进程内锁即安全,无需 FileLock 也可正确落地)。 */ void lockAndRun(IoRunnable action) throws IOException { lockAndRun(action, action); } /** * replayer 语义:FileLock 被别的进程持有则跳过(返回 {@code false},不执行 {@code action},下轮再来)。 * *

顺序:进程内 {@link #monitor} → 已销毁则直接拒绝 → 非阻塞 {@link FileChannel#tryLock()} → 执行 → 倒序释放。 * 不支持 FileLock(降级态)→ 仅进程内锁执行并返回 {@code true}(仅单实例拓扑安全,见类 Javadoc)。 * *

Codex 四轮 Critical#1(destroyed 态认领未拒绝):旧版 destroyed 与 fileLockDisabled 走同一 * 分支——仅进程内锁执行 {@code action} 并返回 {@code true},即停机后仍会"认领"。认领/孤儿回收是可延后 * 到下一轮或下一个存活实例处理的操作,停机窗口内没有必要、也不应该再对 spool 文件动手,故 destroyed * 单独判断、直接返回 {@code false}、不执行 action——与"必须始终执行"的写入语义 * {@link #lockAndRun(IoRunnable, IoRunnable)} 明确区分开。 * * @return {@code true}=已执行(拿到 FileLock 或已降级);{@code false}=FileLock 被他进程持有本轮跳过, * 或已销毁直接拒绝 * @throws IOException {@code action} 抛出的 IO 异常 */ boolean tryLockAndRun(IoRunnable action) throws IOException { synchronized (monitor) { if (destroyed) { // Codex 四轮 Critical#1:destroyed 后直接拒绝认领,不执行 action、不退化为仅进程内锁执行。 warnPostDestroyOnce(); return false; } ensureChannelOpen(); if (fileLockDisabled || channel == null) { // fileLockDisabled(不支持 FileLock 的文件系统)保持原有退化行为:仅进程内锁执行 action。 // 该退化路径仅在"单实例拓扑"下安全——多实例共享同一 spool 卷时跨进程互斥完全失效, // 需运维保证不支持 FileLock 的文件系统只单实例部署(既有建议,见类 Javadoc)。 action.run(); return true; } FileLock fileLock; try { fileLock = channel.tryLock(); // 非阻塞,整文件独占 } catch (OverlappingFileLockException overlap) { // 理论不可达:monitor 已串行化同 JVM,FileLock 请求不会重叠。防御性当作"被占"跳过本轮。 return false; } catch (IOException unsupported) { markDegraded(unsupported); // 不支持的 FS 等 → 降级,仅进程内锁执行 action.run(); return true; } if (fileLock == null) { return false; // 别的进程持有 → 跳过本轮 } try { action.run(); return true; } finally { releaseQuietly(fileLock); } } } /** * 带超时重试获取 FileLock;拿不到返回 {@code null}(调用方走降级动作)。调用方必须持有 {@link #monitor}。 */ private FileLock acquireFileLockWithTimeout() { ensureChannelOpen(); if (fileLockDisabled || destroyed || channel == null) { if (destroyed) { warnPostDestroyOnce(); } return null; // 已降级/已销毁/未开 channel:调用方走降级动作 } long deadline = System.currentTimeMillis() + WRITER_LOCK_WAIT_MS; while (true) { try { FileLock lock = channel.tryLock(); if (lock != null) { return lock; } } catch (OverlappingFileLockException overlap) { // 理论不可达(monitor 串行化);当作瞬时,落入重试/超时逻辑 } catch (IOException unsupported) { markDegraded(unsupported); // 不支持 → 降级 return null; } if (System.currentTimeMillis() >= deadline) { // 超时 → 走降级动作(唯一名应急段,审计事件宁可多文件也不丢);不置 fileLockDisabled(他进程本次持有,属瞬时) warnDegradeOnce("FileLock 获取超时(" + WRITER_LOCK_WAIT_MS + "ms),本次写入降级为唯一名应急段"); return null; } try { Thread.sleep(RETRY_SLEEP_MS); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); return null; // 被中断 → 走降级动作(唯一名应急段),尽快完成写入 } } } /** * 惰性打开 {@code spool.lock} 常驻句柄;打开失败即降级。已销毁则拒绝重开(Codex#D)。 * 调用方必须持有 {@link #monitor}。 */ private void ensureChannelOpen() { if (channel != null || fileLockDisabled || destroyed) { return; // Codex#D:destroyed 时不重开 channel,迟到调用走进程内锁 } try { if (!spoolDir.exists()) { Files.createDirectories(spoolDir.toPath()); } raf = new RandomAccessFile(lockFile, "rw"); channel = raf.getChannel(); } catch (IOException openFail) { markDegraded(openFail); // 无法打开锁文件 → 仅进程内锁 } } /** 置永久降级标志并 warn 一次。调用方必须持有 {@link #monitor}。 */ private void markDegraded(Throwable cause) { fileLockDisabled = true; warnDegradeOnce("跨进程 FileLock 不可用,降级(写入走唯一名应急段/重放走仅进程内锁): " + cause); } /** warn 去重。调用方必须持有 {@link #monitor}。 */ private void warnDegradeOnce(String msg) { if (!degradeWarned) { degradeWarned = true; log.warn("[audit] spool 锁降级: {}", msg); } } /** Codex#D:已销毁后迟到调用的 warn 去重。调用方必须持有 {@link #monitor}。 */ private void warnPostDestroyOnce() { if (!postDestroyWarned) { postDestroyWarned = true; log.warn("[audit] spool 锁已销毁,迟到调用按路径降级(写入路径=仅进程内锁执行," + "认领路径=直接拒绝,不再重开跨进程 channel)"); } } private void releaseQuietly(FileLock fileLock) { if (fileLock == null) { return; } try { fileLock.release(); } catch (IOException releaseFail) { log.warn("[audit] spool FileLock 释放失败: {}", releaseFail.toString()); } } /** * Codex#D:置终止位、释放跨进程句柄。置 {@link #destroyed} 后 {@link #ensureChannelOpen} 永不再重开 * channel,迟到的收尾/写入调用只走进程内锁(不会因 {@code ensureChannelOpen} 复活一个停机后仍持有的跨进程锁)。 */ @Override public void destroy() { synchronized (monitor) { destroyed = true; closeQuietly(channel); closeQuietly(raf); channel = null; raf = null; } } private static void closeQuietly(Closeable c) { if (c != null) { try { c.close(); } catch (IOException ignore) { // 停机释放,忽略 } } } }