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 加固)。
|
*
|
* <p>{@link AuditSpoolWriter} 与 {@link AuditSpoolReplayer} <b>共享同一实例</b>(由
|
* {@code AuditSdkAutoConfiguration} 注入单例),把彼此的 spool 文件操作串行化,从根上消除两条竞态:
|
* <ul>
|
* <li><b>Codex #1(Critical)</b>:writer 每次 open-append-close 与 replayer 的
|
* "认领 rename + readAllLines" 之间的窗口——writer 恰在 replayer 改名旧文件后经已打开的 FD
|
* 写入被改名 inode,事件随 {@code .replaying} 被删而永久丢失。</li>
|
* <li><b>Codex #2(Important)</b>:认领 rename 与 setLastModified 刷新租约是两个独立操作,
|
* 实例 B 可在实例 A rename 成功后、刷新 mtime 前把它当孤儿抢走。</li>
|
* </ul>
|
* 两条竞态的正确性从此由本锁保证(不再依赖 60 秒静默窗口这种活跃度启发式;静默窗口降级为纯优化)。
|
*
|
* <h3>两层结构</h3>
|
* <ol>
|
* <li><b>进程内层</b>:单一监视器 {@link #monitor},writer/replayer 所有临界区共用,保证同 JVM 串行。
|
* 它同时保证同 JVM 内 {@link FileChannel#tryLock()} 请求永不重叠——否则 JDK 会抛
|
* {@link OverlappingFileLockException}(FileLock 是 JVM 级而非线程级),故本类先取进程内锁、
|
* 在 {@code synchronized} 块内获/放 FileLock,天然无重叠。</li>
|
* <li><b>跨进程层</b>:spool 目录下常驻的 {@code spool.lock} 文件 + {@link FileChannel#lock()} 系
|
* 独占 {@link FileLock}。锁文件 {@link RandomAccessFile}/{@link FileChannel} 惰性打开、常驻,
|
* {@link DisposableBean#destroy()} 时释放。</li>
|
* </ol>
|
*
|
* <h3>降级与终局(绝不让业务/发送线程崩溃)</h3>
|
* <ul>
|
* <li>打开锁文件或 {@code tryLock} 抛 {@link IOException}(不支持的文件系统等)→ 置
|
* {@link #fileLockDisabled},之后<b>不再尝试跨进程锁</b>,{@link #warnDegradeOnce} 去重 warn 一次。</li>
|
* <li><b>writer 路径</b> {@link #lockAndRun(IoRunnable, IoRunnable)}:拿到跨进程 FileLock → 执行
|
* {@code lockedAction}(正常追加共享活动文件);<b>未拿到</b>(超时/中断/不可用/已销毁)→ 执行
|
* {@code degradedAction}(Codex#A:写唯一名 {@code .retry} 应急段——一事件一文件,无共享名=无
|
* rename 竞态,天然是重放协议的输入,宁多文件不丢事件)。进程内 {@link #monitor} 全程持有。</li>
|
* <li><b>replayer 路径</b> {@link #tryLockAndRun}:非阻塞 {@code tryLock} 拿不到(他进程持有)就<b>跳过本轮</b>
|
* (返回 {@code false},下轮再来);降级态仅进程内锁执行返回 {@code true}(该退化路径仅在"单实例
|
* 拓扑"下安全——多实例共享同一 spool 卷时跨进程互斥完全失效,需运维保证不支持 FileLock 的文件系统
|
* 只单实例部署,既有建议);<b>已销毁态直接返回 {@code false}、不执行 action</b>(Codex 四轮
|
* Critical#1,见下一条)。</li>
|
* <li><b>Codex#D 终止位</b> {@link #destroyed}:{@link #destroy()} 后 {@link #ensureChannelOpen} 拒绝
|
* 重开 channel——迟到的收尾/写入调用({@link #lockAndRun(IoRunnable, IoRunnable)},"始终执行"语义)
|
* 只走进程内锁执行 degradedAction({@link #warnPostDestroyOnce} 去重 warn 一次),杜绝"已 close
|
* 的 raf/channel 被重开、停机后又持有跨进程锁"。</li>
|
* <li><b>Codex 四轮 Critical#1</b>(destroyed 态认领未拒绝):{@link #tryLockAndRun}(认领/孤儿回收
|
* 路径,"可延后到下轮"语义)在 destroyed 后<b>直接返回 {@code false},不执行 action</b>——而非像
|
* 旧版那样退化为"仅进程内锁执行并返回 true"。停机窗口内不应该再对 spool 文件做任何认领:认领
|
* 延后到下一轮/下一个存活实例处理无害,而写入路径 {@link #lockAndRun(IoRunnable, IoRunnable)}
|
* 必须"始终执行"(事件不能丢),二者语义有意不同,不可混用同一套 destroyed 处理。</li>
|
* </ul>
|
*
|
* <p><b>残余声明</b>:单文件重放(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 降级隔离):<b>始终执行</b>,但按是否持有跨进程锁二选一执行。
|
*
|
* <p>顺序:进程内 {@link #monitor} → 尽力获 FileLock(带超时重试,共最多 {@link #WRITER_LOCK_WAIT_MS})。
|
* <ul>
|
* <li>拿到 FileLock → 执行 {@code lockedAction}(正常追加共享活动文件),倒序释放。</li>
|
* <li>未拿到(超时/中断/不可用/已销毁)→ 执行 {@code degradedAction}(写唯一名 {@code .retry} 应急段)。
|
* 此时<b>绝不</b>再追加共享活动名——唯一名文件无共享、无 rename 竞态,天然是重放协议的输入
|
* (Codex#A:审计事件宁可多几个文件也不丢弃)。</li>
|
* </ul>
|
*
|
* @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 语义便捷重载:降级动作与持锁动作<b>相同</b>。
|
*
|
* <p>适用于收尾等"操作对象本就是唯一名/owner-token 名"的调用方(如 {@link AuditSpoolReplayer} 的
|
* {@code finishReplay}:写 {@code .retry}/{@code .bad}、删 owner-token 的 {@code .replaying.<uuid>}
|
* 都是唯一名,仅进程内锁即安全,无需 FileLock 也可正确落地)。
|
*/
|
void lockAndRun(IoRunnable action) throws IOException {
|
lockAndRun(action, action);
|
}
|
|
/**
|
* replayer 语义:FileLock 被<b>别的进程</b>持有则跳过(返回 {@code false},不执行 {@code action},下轮再来)。
|
*
|
* <p>顺序:进程内 {@link #monitor} → 已销毁则直接拒绝 → 非阻塞 {@link FileChannel#tryLock()} → 执行 → 倒序释放。
|
* 不支持 FileLock(降级态)→ 仅进程内锁执行并返回 {@code true}(仅单实例拓扑安全,见类 Javadoc)。
|
*
|
* <p><b>Codex 四轮 Critical#1(destroyed 态认领未拒绝)</b>:旧版 destroyed 与 fileLockDisabled 走同一
|
* 分支——仅进程内锁执行 {@code action} 并返回 {@code true},即停机后仍会"认领"。认领/孤儿回收是可延后
|
* 到下一轮或下一个存活实例处理的操作,停机窗口内没有必要、也不应该再对 spool 文件动手,故 destroyed
|
* 单独判断、直接返回 {@code false}、<b>不执行 action</b>——与"必须始终执行"的写入语义
|
* {@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} 常驻句柄;打开失败即降级。<b>已销毁则拒绝重开</b>(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) {
|
// 停机释放,忽略
|
}
|
}
|
}
|
}
|