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 落盘。 * *

2026-08-11 变更:原第一级 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 暂不做)。 * *

本类的并发/时序注释对应降级链路上经过验证的并发安全设计——重构时不要随手删改, * 每一处都对应一个真实存在过的竞态场景。 * deliverOnce 只做 Feign 一级并给出明确成败信号(不写 spool),供普通路径与 replayer 各自决定失败处理。 */ @Slf4j public class AuditEventPublisher implements DisposableBean { private final ObjectProvider 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 pending = ConcurrentHashMap.newKeySet(); public AuditEventPublisher(ObjectProvider 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 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 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); } } }