刘光辉
15 小时以前 34981c30a78e8bbd7791131059a9210f9928b62c
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
package jnpf.audit.sdk;
 
import jnpf.util.JsonUtil;
import lombok.extern.slf4j.Slf4j;
 
import java.io.File;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.NoSuchFileException;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.nio.file.StandardOpenOption;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicLong;
 
import jnpf.audit.model.AuditEventDTO;
 
/**
 * 本地兜底 spool 写入器(Task 7 Step 3):三级降级最后一级。
 *
 * <p>按日追加一行 JSON 到 {@code ${audit.spool-dir}/audit-events-yyyyMMdd.jsonl};
 * IO 失败仅 {@code log.error} + 丢弃计数,<b>绝不向业务抛异常</b>——审计落盘失败不能反噬业务。
 *
 * <p><b>并发安全(Codex 批次二复验 #1,双层锁协议终局修复)</b>:{@link #write} 的
 * "open(append)→写一行→flush→close" 整段搬进与 {@link AuditSpoolReplayer} 共享的
 * {@link AuditSpoolLock} 临界区(进程内 monitor + 跨进程 {@code spool.lock} FileLock)。因此
 * append 与重放器的"认领 rename"严格互斥——彻底消除"writer 经已打开 FD 向被改名 inode 写入、
 * 事件随 {@code .replaying} 被删而永久丢失"的窗口竞态(原 60 秒静默窗口仅活跃度启发式,不足以保证正确性,
 * 现降级为纯优化)。{@link Files#write} 单次调用内即开-写-关(写入量 &lt;0.5 QPS,每次开关流性能无虞),
 * 锁内只做这一快操作;重放器把活动文件 rename 走后,下一次 write 以 {@link StandardOpenOption#CREATE}
 * 在原名新建文件继续追加,天然不写回旧 inode(终验四轮③)。
 *
 * <p><b>降级隔离(Codex 批次二三轮 #A:消灭共享名,杜绝降级路径 rename 竞态)</b>:跨进程 FileLock
 * 一时拿不到(超时/中断/不可用/已销毁)时,{@link AuditSpoolLock#lockAndRun(AuditSpoolLock.IoRunnable, AuditSpoolLock.IoRunnable)}
 * 不再让本次写入去<b>追加共享活动名</b>(那样会与其它进程的"认领 rename+删除"重演 #1 丢事件时序:无锁
 * open 活动文件 → 他进程 rename+读 → 本进程向被改名 inode 写 → 他进程删 → 事件丢失),而是改走
 * "应急段"路径:把该行写入唯一名文件(一事件一文件,仅降级场景、低 QPS 可接受),由重放器锁内认领
 * 消费——宁多几个文件也不丢事件。进程内 monitor 全程持有(同 JVM 串行不变)。
 *
 * <p><b>写后原子发布(Codex 四轮复验 Critical#1,消灭跨进程半写窗口)</b>:上一版把应急段直接用可扫描的
 * {@code .retry} 名一步 {@link StandardOpenOption#CREATE_NEW} 建文件+写入——但"目录项创建可见"与
 * "内容写完 close"之间存在窗口,另一进程的 {@link AuditSpoolReplayer} 按文件名扫描/认领可能落进这个
 * 窗口:读到空文件 → 收尾把它当处理完删掉(token 丢失)而本进程仍持已打开 FD 向已被 unlink 的 inode
 * 追加 → 事件静默丢失;或读到半行 → 前缀被当坏行隔离进 {@code .bad},行内容本不完整却被判定"已处理"。
 * 现改为<b>先写后发布</b>两步,rename 的原子性物理消灭这个窗口:
 * <ol>
 *   <li>写入唯一名 {@code <活动名>.<UUID>.retry.writing}({@link StandardOpenOption#CREATE_NEW},
 *       {@link Files#write} 单次调用内写完 flush+close)——此时文件名不在重放器扫描集
 *       ({@code .retry}/{@code .jsonl})内,任何进程都不会认领它;</li>
 *   <li>{@link Files#move} + {@link StandardCopyOption#ATOMIC_MOVE} 原子改名发布为
 *       {@code <活动名>.<UUID>.retry}——rename 是文件系统原子操作,{@code .retry} 这个可扫描名
 *       从"不存在"直接跳到"内容完整存在",不存在中间态,跨进程认领窗口物理消失。</li>
 * </ol>
 * rename 因源文件已不存在而失败({@link NoSuchFileException})时,说明本次写入在上面第①步
 * {@link Files#write} 调用<b>内部</b>被阻塞过久(磁盘/网络存储 IO 卡顿),导致
 * {@link AuditSpoolReplayer#recoverOrphanedWriting()} 误判其超龄并把源文件隔离为 {@code .stale}——
 * 此为 Codex 五轮修复的自愈场景,见下一段。其它 rename 失败(磁盘满、权限突变等极端 IO 异常)时仍只
 * {@code log.error}(脱敏,仅记文件名与 {@code clientEventId},不记事件正文)、<b>保留 {@code .writing}
 * 文件不删、不重试</b>——内容已完整落盘,不计入 {@link #discardCount}(不算事件丢失);交由
 * {@link AuditSpoolReplayer} 的孤儿回收兜底(mtime 超阈值后锁内重验+隔离为 {@code .stale},等待人工
 * 核查补录,见其类 Javadoc)。
 *
 * <p><b>Codex 五轮修复(自愈重发:源被隔离/挪走时重试一次)</b>:{@code .writing} 源文件被
 * {@link AuditSpoolReplayer#recoverOrphanedWriting()} 隔离为 {@code .stale} 后,本次 rename 会抛
 * {@link NoSuchFileException}——此时事件内容并未真的丢失(已完整落在 {@code .stale} 里,等待人工
 * 核查),但本次投递尚未自动完成。为了不让审计事件仅仅因为一次极端的 IO 阻塞就必须依赖人工补录,
 * {@link #write} 在此情形下<b>只重试一次</b>:换一个新 {@code UUID} 重新走一遍「写 {@code .writing}
 * → move 发布为 {@code .retry}」两步协议(见 {@link #publishOnceOrRetryAfterQuarantine})。重试
 * 成功——事件已正常投递到重放通道,{@code .stale} 里的旧副本至多在人工核查补录时造成一次重复,落库按
 * {@code clientEventId} 幂等吸收,无害;重试仍以 {@code NoSuchFileException} 失败(理论上极罕见,需
 * 连续两次不同 UUID 的文件都被隔离)则不再重试,{@code log.error}(含 {@code clientEventId})交由
 * 人工核查 {@code .stale} 隔离区,原样保留其中已经写好的完整内容。
 */
@Slf4j
public class AuditSpoolWriter {
 
    private static final DateTimeFormatter DAY = DateTimeFormatter.ofPattern("yyyyMMdd");
 
    private final File spoolDir;
    private final AuditSpoolLock spoolLock;
    private final AtomicLong discardCount = new AtomicLong();
 
    public AuditSpoolWriter(String spoolDir, AuditSpoolLock spoolLock) {
        this.spoolDir = new File(spoolDir);
        this.spoolLock = spoolLock;
    }
 
    public void write(AuditEventDTO event) {
        try {
            String fileName = "audit-events-" + LocalDate.now().format(DAY) + ".jsonl";
            Path path = new File(spoolDir, fileName).toPath();
            byte[] bytes = (JsonUtil.getObjectToString(event) + "\n")   // jsonl 规范用 \n
                    .getBytes(StandardCharsets.UTF_8);
            // Codex 批次二复验 #1:持锁路径 open-append-flush-close 搬进共享临界区,与 replayer 的认领 rename 互斥
            spoolLock.lockAndRun(
                    // lockedAction:持有跨进程锁 → 正常追加共享活动文件
                    () -> {
                        if (!spoolDir.exists()) {
                            Files.createDirectories(spoolDir.toPath());
                        }
                        Files.write(path, bytes,
                                StandardOpenOption.CREATE, StandardOpenOption.APPEND);
                    },
                    // degradedAction(Codex 批次二三轮 #A + 四轮 Critical#1 + 五轮自愈):未持锁 → 写后
                    // 原子发布唯一名 .retry 应急段,消灭共享名/杜绝 rename 竞态、消灭跨进程半写窗口(见类
                    // Javadoc「写后原子发布」段);发布源被 AuditSpoolReplayer 孤儿隔离时自愈重试一次
                    // (见类 Javadoc「Codex 五轮修复(自愈重发)」段)。
                    () -> {
                        if (!spoolDir.exists()) {
                            Files.createDirectories(spoolDir.toPath());
                        }
                        publishOnceOrRetryAfterQuarantine(fileName, bytes, event);
                    });
        } catch (Throwable t) {
            long n = discardCount.incrementAndGet();
            log.error("AUDIT-SPOOL-LOST clientEventId={} 落盘失败,累计丢弃={}",
                    event == null ? null : event.getClientEventId(), n, t);
        }
    }
 
    /**
     * degradedAction 主体:写后原子发布两步协议(写唯一名 {@code .writing} → move 发布为
     * {@code .retry}),源被孤儿隔离({@link NoSuchFileException})时自愈<b>只重试一次</b>——见类
     * Javadoc「写后原子发布」与「Codex 五轮修复(自愈重发)」两段的完整论证。
     *
     * @throws IOException 非隔离原因的 rename 失败(磁盘满等)已在方法内 log.error 吸收并返回,
     *                      仅 {@code Files.write} 本身失败(如磁盘已满导致写不进去)会外抛,
     *                      交由 {@link #write} 的外层 catch 计入 {@link #discardCount}
     */
    private void publishOnceOrRetryAfterQuarantine(String fileName, byte[] bytes, AuditEventDTO event)
            throws IOException {
        boolean retried = false;
        while (true) {
            String uuid = UUID.randomUUID().toString();
            Path writing = new File(spoolDir, fileName + "." + uuid + ".retry.writing").toPath();
            Path published = new File(spoolDir, fileName + "." + uuid + ".retry").toPath();
            // 第一步:写唯一名 .writing 临时文件,CREATE_NEW+单次 write 调用内写完 flush+close
            Files.write(writing, bytes, StandardOpenOption.CREATE_NEW, StandardOpenOption.WRITE);
            try {
                // 第二步:原子改名发布为 .retry——rename 原子性使可扫描名下只可能出现完整内容
                Files.move(writing, published, StandardCopyOption.ATOMIC_MOVE);
                return;   // 发布成功
            } catch (NoSuchFileException sourceGone) {
                // 源已被 AuditSpoolReplayer#recoverOrphanedWriting() 隔离为 .stale(本次写入曾在
                // Files.write 内被阻塞超过孤儿回收阈值)——内容未丢,只是滞留隔离区。只重试一次:
                // 换新 UUID 重写一份再发布;.stale 里的旧副本至多造成人工补录时的重复(幂等吸收)
                if (retried) {
                    log.error("AUDIT-SPOOL-RETRY-PUBLISH-FAIL file={} clientEventId={} "
                                    + "源文件已被隔离且重试后仍失败,需人工核查 .stale 隔离区补录",
                            writing.getFileName(),
                            event == null ? null : event.getClientEventId(), sourceGone);
                    return;
                }
                retried = true;
                log.warn("[audit] degraded 发布源文件被隔离(写入曾阻塞超过孤儿回收阈值),自愈重写一次: "
                                + "file={} clientEventId={}",
                        writing.getFileName(), event == null ? null : event.getClientEventId());
                // 继续循环:新 UUID 重写一份再尝试发布
            } catch (IOException renameFail) {
                // 其它极端 IO 异常(磁盘满/权限突变等):内容已完整落在 .writing(上一步写入已成功),
                // 不算事件丢失——不计入 discardCount,保留 .writing 交孤儿回收兜底
                log.error("AUDIT-SPOOL-RETRY-PUBLISH-FAIL file={} clientEventId={} "
                                + "原子发布失败,内容已保留于 .writing 待孤儿回收",
                        writing.getFileName(),
                        event == null ? null : event.getClientEventId(), renameFail);
                return;
            }
        }
    }
 
    /** 累计因落盘 IO 失败而丢弃的事件数(观测用)。 */
    public long discardCount() {
        return discardCount.get();
    }
}