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 文件操作串行化,从根上消除两条竞态: *
残余声明:单文件重放(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})。 *
适用于收尾等"操作对象本就是唯一名/owner-token 名"的调用方(如 {@link AuditSpoolReplayer} 的
* {@code finishReplay}:写 {@code .retry}/{@code .bad}、删 owner-token 的 {@code .replaying. 顺序:进程内 {@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) {
// 停机释放,忽略
}
}
}
}