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);
|
}
|
}
|
}
|