刘光辉
昨天 bb638871a7fb692d80f1b7a758f991dc0879002c
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
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) {
                // 停机释放,忽略
            }
        }
    }
}