刘光辉
昨天 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
package jnpf.audit.sdk;
 
import jnpf.audit.AuditDeliveryException;
import jnpf.audit.AuditEventApi;
import jnpf.audit.model.AuditEventDTO;
import jnpf.base.ActionResult;
import jnpf.base.ActionResultCode;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.ObjectProvider;
 
import java.util.List;
import java.util.Set;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
 
/**
 * 两级降级异步发送器:Feign(AuditEventApi 同步)→ 本地 spool 落盘。
 *
 * <p><b>2026-08-11 变更</b>:原第一级 MQ({@code streamBridge.send("AUDIT_EVENT", …)})已移除。
 * 事件总线于 2026-08-10 切 Redis 后栈内不再有 RocketMQ binder,该级**每次必然失败**再由 Feign 接管——
 * 它已不是降级级次,只是每条事件多一次异常构造。移除后行为等价(此前本就每条都走 Feign)。
 * 若将来要恢复异步投递,正解是提供一个写 Redis 的 {@link AuditEventApi} 实现 bean 替代 Feign,
 * 本类一行不用改(见 tasks/todo.md 阶段 4,当前因 Redis 无持久化 + noeviction 暂不做)。
 *
 * <p>本类的并发/时序注释对应降级链路上经过验证的并发安全设计——重构时不要随手删改,
 * 每一处都对应一个真实存在过的竞态场景。
 * deliverOnce 只做 Feign 一级并给出明确成败信号(不写 spool),供普通路径与 replayer 各自决定失败处理。
 */
@Slf4j
public class AuditEventPublisher implements DisposableBean {
    private final ObjectProvider<AuditEventApi> feignApi;
    private final AuditSpoolWriter spool;
    private final ThreadPoolExecutor pool = new ThreadPoolExecutor(1, 2, 60, TimeUnit.SECONDS,
            new ArrayBlockingQueue<>(1000),
            r -> { Thread t = new Thread(r, "audit-sender"); t.setDaemon(true); return t; },
            new ThreadPoolExecutor.AbortPolicy());   // 队列满走 catch → spool,不阻塞业务
 
    /** pending 集合:提交前在调用线程登记、run() finally 移除——统一覆盖"排队中/执行中/已出队未开跑"三态,
        无登记窗口(终验三轮④:原 run() 内 add 的写法在出队与登记之间有窗口,shutdownNow 撞上会双漏) */
    private final Set<AuditSendTask> pending = ConcurrentHashMap.newKeySet();
 
    public AuditEventPublisher(ObjectProvider<AuditEventApi> feignApi,
                               AuditSpoolWriter spool) {
        this.feignApi = feignApi;
        this.spool = spool;
    }
 
    /** 任务自携事件对象:异常/停机 drain 时能恢复出具体事件落 spool(评审#2,裸 lambda 恢复不出事件) */
    class AuditSendTask implements Runnable {
        final AuditEventDTO event;
        AuditSendTask(AuditEventDTO event) { this.event = event; }
        @Override public void run() {
            try {
                deliverOnce(event);
            } catch (Throwable t) {   // 最外层兜底:任何异常(含 JSON 序列化)都落 spool,绝不静默丢(评审#2)
                spool.write(event);
                log.error("AUDIT-LOST clientEventId={} type={} target={} 已落 spool",
                        event.getClientEventId(), event.getEventType(), event.getTargetTable(), t);
            } finally {
                pending.remove(this);
            }
        }
    }
 
    public void publishAsync(List<AuditEventDTO> events) {
        for (AuditEventDTO e : events) {
            AuditSendTask task = new AuditSendTask(e);
            pending.add(task);                    // 先登记再入池(终验三轮④)
            try {
                pool.execute(task);
            } catch (RejectedExecutionException full) {
                // 先写 spool 再移除登记(终验四轮②:反序会在 remove 与 write 之间留下双侧不可见窗口)
                try {
                    spool.write(e);   // 队列满直接落 spool(spec §5.5)
                } finally {
                    pending.remove(task);
                }
            }
        }
    }
 
    /**
     * 单次投递:Feign 同步,成功返回、失败抛异常——本方法不写 spool。
     * 终验#2/#3:普通路径与 replayer 都需要明确的成败信号各自决定失败处理,
     * 原 send() 失败时自己写 spool 正常返回,replayer 无法判断该不该保留待重试行。
     */
    void deliverOnce(AuditEventDTO e) {
        ActionResult<Integer> result = feignApi.getObject().record(e);
        if (result == null || !ActionResultCode.SUCCESS.getCode().equals(result.getCode())) {
            // HTTP 层失败由 fallbackFactory 抛出;业务层失败(HTTP 200 + fail body)靠本校验捕获(Codex 批次一 Critical#1)
            throw new AuditDeliveryException("Feign 投递返回非成功: " + (result == null ? "null" : result.getCode()), null);
        }
    }
 
    /**
     * Codex 批次二 #5:以 {@link DisposableBean#destroy()} 替代 jakarta {@code @PreDestroy},
     * 跨 boot2(JDK8)/boot3 中立——{@code @PreDestroy} 在 boot2 走 {@code javax.annotation}、
     * boot3 走 {@code jakarta.annotation},SDK 二进制无法同时兼容;DisposableBean 是 Spring 原生接口,
     * 两版一致。Spring 对 {@code @Bean} 会自动识别 DisposableBean,{@code AuditSdkAutoConfiguration} 无需改。
     */
    @Override
    public void destroy() {
        drain();
    }
 
    public void drain() {   // 停机 drain(spec §5.5);不向外抛异常(终验三轮④:异常逃出会绕过全部兜底)
        pool.shutdown();
        try {
            if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
                pool.shutdownNow();
                spoolPending();
            }
        } catch (InterruptedException ie) {
            pool.shutdownNow();
            spoolPending();
            Thread.currentThread().interrupt();   // 兜底做完再恢复中断标记
        }
    }
 
    private void spoolPending() {
        // pending 统一覆盖排队中/执行中/出队未开跑三态;与任务自身 catch 落盘可能重复——
        // spool 重放与落库均按 client_event_id 幂等,重复无害
        for (AuditSendTask task : pending) {
            spool.write(task.event);
        }
    }
}