diff --git a/.gitignore b/.gitignore
index 9154f4c..b988dc2 100644
--- a/.gitignore
+++ b/.gitignore
@@ -24,3 +24,17 @@
hs_err_pid*
replay_pid*
+# Maven 构建产物
+target/
+logs/
+
+# IDE
+.idea/
+*.iml
+.vscode/
+.settings/
+.project
+.classpath
+
+# 本地配置覆盖(密钥等敏感信息)
+app-env.yml
diff --git a/pom.xml b/pom.xml
new file mode 100644
index 0000000..914c415
--- /dev/null
+++ b/pom.xml
@@ -0,0 +1,157 @@
+
+
+ * 职责:启动 Solon 容器,加载各模块(回调网关、分发工作线程、UserSig、补拉、统计等)。 + *
+ * {@code @EnableScheduling} 启用 Solon 定时任务,配合 solon-scheduling-simple 驱动, + * 用于补拉巡检、分区自建、死信重置、用量统计等定时作业。 + * + * @author imutil + */ +@SolonMain +@EnableScheduling +public class App { + + public static void main(String[] args) { + Solon.start(App.class, args); + } +} diff --git a/src/main/java/com/imutil/common/BizException.java b/src/main/java/com/imutil/common/BizException.java new file mode 100644 index 0000000..248cd8e --- /dev/null +++ b/src/main/java/com/imutil/common/BizException.java @@ -0,0 +1,35 @@ +package com.imutil.common; + +import lombok.Getter; + +/** + * 业务异常 + *
+ * 业务逻辑中主动抛出,由 {@link com.imutil.filter.GlobalExceptionFilter} 捕获后 + * 以对应 code 返回前端。默认 code=500。 + * + * @author imutil + */ +@Getter +public class BizException extends RuntimeException { + + private static final long serialVersionUID = 1L; + + /** 返回代码 */ + private final int code; + + public BizException(String message) { + super(message); + this.code = 500; + } + + public BizException(int code, String message) { + super(message); + this.code = code; + } + + public BizException(int code, String message, Throwable cause) { + super(message, cause); + this.code = code; + } +} diff --git a/src/main/java/com/imutil/common/Httpx.java b/src/main/java/com/imutil/common/Httpx.java new file mode 100644 index 0000000..7569908 --- /dev/null +++ b/src/main/java/com/imutil/common/Httpx.java @@ -0,0 +1,67 @@ +package com.imutil.common; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.time.Duration; + +/** + * HTTP 工具封装 + *
+ * 基于 JDK 21 内置 java.net.http.HttpClient,无第三方依赖。
+ * 用于分发工作线程向业务系统转发回调、调用腾讯 REST API。
+ *
+ * @author imutil
+ */
+public final class Httpx {
+
+ /** 全局复用 HttpClient(线程安全) */
+ private static final HttpClient CLIENT = HttpClient.newBuilder()
+ .connectTimeout(Duration.ofSeconds(5))
+ .build();
+
+ private Httpx() {
+ }
+
+ /**
+ * POST JSON 请求,返回 HTTP 状态码
+ *
+ * @param url 目标地址
+ * @param body JSON 请求体
+ * @return HTTP 状态码;2xx 视为成功
+ */
+ public static int postJson(String url, String body) throws Exception {
+ HttpRequest req = HttpRequest.newBuilder()
+ .uri(URI.create(url))
+ .header("Content-Type", "application/json; charset=utf-8")
+ .POST(HttpRequest.BodyPublishers.ofString(body == null ? "" : body))
+ .timeout(Duration.ofSeconds(10))
+ .build();
+ HttpResponse
+ * 基于 MyBatis-Plus 内置雪花算法({@link IdWorker})生成全局唯一 ID。
+ * 提供字符串形式以规避 Long 雪花 ID 经 JSON 传前端时精度丢失(>2^53)。
+ *
+ * @author imutil
+ */
+public final class Ids {
+
+ private Ids() {
+ }
+
+ /**
+ * 生成雪花 ID(long)
+ */
+ public static long nextId() {
+ return IdWorker.getId();
+ }
+
+ /**
+ * 生成雪花 ID(字符串)
+ */
+ public static String nextIdStr() {
+ return String.valueOf(IdWorker.getId());
+ }
+}
diff --git a/src/main/java/com/imutil/common/Jsons.java b/src/main/java/com/imutil/common/Jsons.java
new file mode 100644
index 0000000..dd9de63
--- /dev/null
+++ b/src/main/java/com/imutil/common/Jsons.java
@@ -0,0 +1,49 @@
+package com.imutil.common;
+
+import org.noear.snack4.ONode;
+
+import java.lang.reflect.Type;
+
+/**
+ * JSON 工具封装
+ *
+ * 基于 Solon 内置的 snack4({@link ONode}),不额外引入 fastjson/jackson,保持依赖精简。
+ * 用于 Redis 缓存对象序列化、回调消息体解析等。
+ *
+ * @author imutil
+ */
+public final class Jsons {
+
+ private Jsons() {
+ }
+
+ /**
+ * 对象序列化为 JSON 字符串
+ */
+ public static String stringify(Object obj) {
+ if (obj == null) {
+ return null;
+ }
+ return ONode.serialize(obj);
+ }
+
+ /**
+ * JSON 字符串反序列化为对象
+ */
+ public static
+ * 作为 Redis 的兜底热缓存,承载映射/授权等高频读、短 TTL 的数据,
+ * 减少网络往返。多节点下数据可能短暂不一致,变更时需主动失效并依赖 Redis 兜底。
+ *
+ * @author imutil
+ */
+@Component
+public class LocalCache {
+
+ private final Cache
+ * 统一回调路径与补拉路径的 msg_key 算法,确保同一条消息无论从哪条路径进入
+ * 都生成相同 key,从而跨路径幂等去重(回调已落库的消息不会被补拉重复落库)。
+ *
+ * 腾讯侧一条消息由 (from, 会话目标, MsgSeq, MsgRandom) 唯一确定。
+ * 不含 CallbackCommand:同一条消息的发送/接收回调 command 不同,含 command 反而不利于去重。
+ *
+ * @author imutil
+ */
+public final class MsgKeys {
+
+ private MsgKeys() {
+ }
+
+ /**
+ * 生成消息唯一键
+ *
+ * @param fromAccount 发送方账号,可空
+ * @param convTarget 会话目标(GROUP=群ID,C2C=对端账号)
+ * @param msgSeq 腾讯 MsgSeq
+ * @param msgRandom 腾讯 MsgRandom
+ * @return msg_key,格式 from:convTarget:msgSeq:msgRandom
+ */
+ public static String build(String fromAccount, String convTarget, long msgSeq, long msgRandom) {
+ return (fromAccount == null ? "" : fromAccount) + ":" +
+ (convTarget == null ? "" : convTarget) + ":" + msgSeq + ":" + msgRandom;
+ }
+}
diff --git a/src/main/java/com/imutil/common/PasswordUtil.java b/src/main/java/com/imutil/common/PasswordUtil.java
new file mode 100644
index 0000000..d972ddd
--- /dev/null
+++ b/src/main/java/com/imutil/common/PasswordUtil.java
@@ -0,0 +1,84 @@
+package com.imutil.common;
+
+import javax.crypto.SecretKeyFactory;
+import javax.crypto.spec.PBEKeySpec;
+import java.security.SecureRandom;
+import java.util.Base64;
+
+/**
+ * 密码哈希工具(PBKDF2WithHmacSHA256,JDK 内置,无额外依赖)
+ *
+ * 存储格式:{@code iterations:saltBase64:hashBase64}。
+ * verify 时按存储的 iterations/salt 重新推导并常量时间比对。
+ * 每次哈希用随机 salt,防彩虹表。
+ *
+ * @author imutil
+ */
+public final class PasswordUtil {
+
+ /** 迭代次数(OWASP 2023 建议 PBKDF2-HMAC-SHA256 ≥ 120000) */
+ private static final int ITERATIONS = 120_000;
+ private static final int KEY_BITS = 256;
+ private static final int SALT_BYTES = 16;
+ private static final String ALGO = "PBKDF2WithHmacSHA256";
+
+ private PasswordUtil() {
+ }
+
+ /**
+ * 对明文密码哈希,返回 "iterations:salt:hash"
+ */
+ public static String hash(String rawPassword) {
+ byte[] salt = new byte[SALT_BYTES];
+ new SecureRandom().nextBytes(salt);
+ byte[] dk = derive(rawPassword, salt, ITERATIONS);
+ return ITERATIONS + ":" + b64(salt) + ":" + b64(dk);
+ }
+
+ /**
+ * 校验明文密码与存储hash是否匹配
+ */
+ public static boolean verify(String rawPassword, String stored) {
+ if (rawPassword == null || stored == null) {
+ return false;
+ }
+ String[] parts = stored.split(":");
+ if (parts.length != 3) {
+ return false;
+ }
+ try {
+ int iter = Integer.parseInt(parts[0]);
+ byte[] salt = Base64.getDecoder().decode(parts[1]);
+ byte[] expected = Base64.getDecoder().decode(parts[2]);
+ byte[] actual = derive(rawPassword, salt, iter);
+ return constantTimeEquals(expected, actual);
+ } catch (Exception e) {
+ return false;
+ }
+ }
+
+ private static byte[] derive(String rawPassword, byte[] salt, int iter) {
+ try {
+ PBEKeySpec spec = new PBEKeySpec(rawPassword.toCharArray(), salt, iter, KEY_BITS);
+ return SecretKeyFactory.getInstance(ALGO).generateSecret(spec).getEncoded();
+ } catch (Exception e) {
+ throw new IllegalStateException("PBKDF2 推导失败", e);
+ }
+ }
+
+ private static String b64(byte[] b) {
+ return Base64.getEncoder().encodeToString(b);
+ }
+
+ /** 常量时间比较,防时序攻击 */
+ private static boolean constantTimeEquals(byte[] a, byte[] b) {
+ if (a.length != b.length) {
+ return false;
+ }
+ int r = 0;
+ for (int i = 0; i < a.length; i++) {
+ r |= a[i] ^ b[i];
+ }
+ return r == 0;
+ }
+}
diff --git a/src/main/java/com/imutil/common/PathWhitelist.java b/src/main/java/com/imutil/common/PathWhitelist.java
new file mode 100644
index 0000000..c715438
--- /dev/null
+++ b/src/main/java/com/imutil/common/PathWhitelist.java
@@ -0,0 +1,51 @@
+package com.imutil.common;
+
+import java.util.Set;
+
+/**
+ * 鉴权白名单路径
+ *
+ * 以下路径不走 app_key 租户鉴权:
+ * - /callback/** 腾讯回调入口(走签名校验,见 T5)
+ * - /health 健康检查
+ * - /favicon.ico
+ *
+ * @author imutil
+ */
+public final class PathWhitelist {
+
+ /** 白名单路径前缀(contextPath 之后的部分) */
+ public static final Set
+ * 读取 app.yml 的 imutil.redis.* 配置,构建 JedisPool 注册为 Solon Bean。
+ * 用于限流令牌桶、映射/授权缓存、补拉水位线等高频读写场景。
+ *
+ * @author imutil
+ */
+@Configuration
+public class RedisConfig {
+
+ @Inject("${imutil.redis.host:localhost}")
+ private String host;
+
+ @Inject("${imutil.redis.port:6379}")
+ private int port;
+
+ @Inject("${imutil.redis.password:}")
+ private String password;
+
+ @Inject("${imutil.redis.database:0}")
+ private int database;
+
+ @Inject("${imutil.redis.timeout:2000}")
+ private int soTimeout;
+
+ @Inject("${imutil.redis.connectTimeout:2000}")
+ private int connectTimeout;
+
+ /**
+ * 构建 Jedis 连接池
+ */
+ @Bean
+ public JedisPool jedisPool() {
+ GenericObjectPoolConfig
+ * 基于 JedisPool 的轻量封装,覆盖限流/缓存/水位线所需能力:
+ * 通用 KV、带 TTL 的 set、setNx(分布式锁/令牌)、incr 计数、过期时间。
+ *
+ * @author imutil
+ */
+@Component
+public class RedisService {
+
+ @Inject
+ private JedisPool jedisPool;
+
+ /**
+ * 设置键值(无过期)
+ */
+ public void set(String key, String value) {
+ try (Jedis j = jedisPool.getResource()) {
+ j.set(key, value);
+ }
+ }
+
+ /**
+ * 设置键值并指定 TTL(秒)
+ */
+ public void setex(String key, String value, long ttlSeconds) {
+ try (Jedis j = jedisPool.getResource()) {
+ j.setex(key, ttlSeconds, value);
+ }
+ }
+
+ /**
+ * 读取键值
+ */
+ public String get(String key) {
+ try (Jedis j = jedisPool.getResource()) {
+ return j.get(key);
+ }
+ }
+
+ /**
+ * 读取并反序列化为对象(JSON)
+ */
+ public
+ * 由 {@link com.imutil.filter.TenantAuthFilter} 在请求入口解析并设置,
+ * 业务层通过 {@link #get()} 获取当前租户,所有数据查询强制带 tenant_id 过滤,
+ * 实现"系统 A 绝对拿不到系统 B 数据"的隔离。
+ *
+ * 必须在请求结束时 {@link #clear()} 清理,避免线程复用串租户。
+ *
+ * @author imutil
+ */
+public final class TenantContext {
+
+ private static final ThreadLocal
+ * 路由前缀 /admin。鉴权由 {@link com.imutil.filter.AdminAuthFilter} 拦截(除 login/logout)。
+ * 页面用 FreeMarker 渲染(对齐 yxtech),数据 CRUD 直连对应 Mapper。
+ *
+ * @author imutil
+ */
+@Mapping("/admin")
+@Controller
+@Slf4j
+public class AdminController {
+
+ @Inject("${server.contextPath:}")
+ private String basePath;
+
+ @Inject
+ private AdminUserService adminUserService;
+
+ @Inject
+ private TenantMapper tenantMapper;
+
+ @Inject
+ private CrossTenantGrantMapper grantMapper;
+
+ @Inject
+ private DistQueueMapper distQueueMapper;
+
+ @Inject
+ private UsageStatMapper usageStatMapper;
+
+ // ==================== 登录 / 登出 ====================
+
+ @Get
+ @Mapping("")
+ public void index(Context ctx) throws Throwable {
+ ctx.redirect(basePath + "/admin/home");
+ }
+
+ @Get
+ @Mapping("/login")
+ public Object loginPage(@Param(defaultValue = "") String error) {
+ if (StpUtil.isLogin()) {
+ // 已登录不再渲染登录页(重定向由调用方处理,此处仍渲染避免死循环)
+ }
+ ModelAndView mv = new ModelAndView("login.ftl");
+ mv.put("basePath", basePath);
+ if ("1".equals(error)) {
+ mv.put("errorMsg", "用户名或密码错误");
+ }
+ return mv;
+ }
+
+ @Post
+ @Mapping("/login")
+ public void doLogin(@Param(defaultValue = "") String username,
+ @Param(defaultValue = "") String password,
+ Context ctx) throws Throwable {
+ AdminUser u = adminUserService.login(username, password);
+ if (u != null) {
+ StpUtil.login(u.getId());
+ log.info("管理后台登录成功 id={} username={}", u.getId(), username);
+ ctx.redirect(basePath + "/admin/home");
+ } else {
+ log.warn("管理后台登录失败 username={}", username);
+ ctx.redirect(basePath + "/admin/login?error=1");
+ }
+ }
+
+ @Get
+ @Mapping("/logout")
+ public void logout(Context ctx) throws Throwable {
+ StpUtil.logout();
+ ctx.redirect(basePath + "/admin/login");
+ }
+
+ // ==================== 首页(仪表盘) ====================
+
+ @Get
+ @Mapping("/home")
+ public Object home() {
+ ModelAndView mv = view("home.ftl", "仪表盘", "home");
+ mv.put("tenantCount", tenantMapper.selectCount(null));
+ mv.put("queuePending", distQueueMapper.selectCount(Wrappers.
+ * 腾讯控制台将回调 URL 配置为 http://host/imutil/callback/im 。
+ * 本控制器负责:签名校验(防伪造+防重放)→ 委托 CallbackService 落库+入队 → 5s 内返回 ActionResults。
+ *
+ * 注意:仅做"落库+写队列"即返回,HTTP 转发业务系统由分发工作线程异步进行(T6),
+ * 确保腾讯回调在 5 秒内得到响应,避免被腾讯判定失败。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Controller
+public class CallbackController {
+
+ @Inject
+ private CallbackService callbackService;
+
+ @Inject("${imutil.tencent.callbackToken:}")
+ private String callbackToken;
+
+ /**
+ * 腾讯 IM 回调统一入口
+ * URL 参数:SdkAppid / CallbackCommand / Sign / RequestTime / contenttype 等
+ */
+ @Mapping(value = "/callback/im", method = MethodType.POST)
+ public void callback(Context ctx,
+ @Param(value = "CallbackCommand", required = false) String command,
+ @Param(value = "Sign", required = false) String sign,
+ @Param(value = "RequestTime", required = false) String requestTime,
+ @Param(value = "SdkAppid", required = false) String sdkAppid) throws Throwable {
+ String body = ctx.body();
+
+ // 1. 签名校验(未配置 Token 时跳过,便于联调,生产必须配置)
+ if (callbackToken != null && !callbackToken.isEmpty()) {
+ if (!TencentCallbackSign.verify(callbackToken, requestTime, sign)) {
+ log.warn("回调签名校验失败 command={} sdkAppid={} requestTime={}", command, sdkAppid, requestTime);
+ ctx.output("{\"ActionStatus\":\"FAIL\",\"ErrorCode\":401,\"ErrorInfo\":\"sign invalid\"}");
+ return;
+ }
+ }
+
+ // 2. 落库 + 入队(同事务,快)
+ String result;
+ try {
+ result = callbackService.handleCallback(command, body);
+ } catch (Throwable e) {
+ log.error("回调处理异常 command={} sdkAppid={}", command, sdkAppid, e);
+ // 处理异常仍返回 OK,避免腾讯反复重试同一回调(消息已可能在事务中落库)
+ // 落库失败的消息由补拉服务(T10)兜底
+ result = "{\"ActionStatus\":\"OK\",\"ErrorCode\":0,\"ErrorInfo\":\"\"}";
+ }
+ ctx.output(result);
+ }
+}
diff --git a/src/main/java/com/imutil/controller/SigController.java b/src/main/java/com/imutil/controller/SigController.java
new file mode 100644
index 0000000..779a749
--- /dev/null
+++ b/src/main/java/com/imutil/controller/SigController.java
@@ -0,0 +1,68 @@
+package com.imutil.controller;
+
+import com.imutil.common.BizException;
+import com.imutil.common.TenantContext;
+import com.imutil.model.Result;
+import com.imutil.service.UserMappingService;
+import com.imutil.tencent.UserSigUtil;
+import org.noear.solon.annotation.Controller;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.annotation.Mapping;
+import org.noear.solon.annotation.Param;
+import org.noear.solon.core.handle.MethodType;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * UserSig 签发接口
+ *
+ * 业务系统后端在客户端登录时调用本接口,本工具用 SecretKey 生成 UserSig 返回。
+ * 客户端拿 im_user_id + UserSig 直连腾讯 IM SDK 收发消息(控制面走本工具,数据面走腾讯)。
+ *
+ * 鉴权:走 {@link com.imutil.filter.TenantAuthFilter},需带 X-App-Key/X-App-Secret。
+ *
+ * @author imutil
+ */
+@Controller
+public class SigController {
+
+ @Inject
+ private UserMappingService userMappingService;
+
+ @Inject("${imutil.tencent.sdkAppId:0}")
+ private long sdkAppId;
+
+ @Inject("${imutil.tencent.secretKey:}")
+ private String secretKey;
+
+ @Inject("${imutil.tencent.usersigExpireDays:7}")
+ private int expireDays;
+
+ /**
+ * 签发 UserSig(首次自动创建 IM 账号)
+ *
+ * 入参:bizUserId(必填)、nick/faceUrl(可选,首次创建时同步到 IM)
+ */
+ @Mapping(value = "/sig/generate", method = MethodType.POST)
+ public Result> generate(@Param("bizUserId") String bizUserId,
+ @Param(value = "nick", required = false) String nick,
+ @Param(value = "faceUrl", required = false) String faceUrl) {
+ if (bizUserId == null || bizUserId.isEmpty()) {
+ throw new BizException(400, "bizUserId 必填");
+ }
+ String tenantId = TenantContext.require();
+ // 获取或创建 IM 账号
+ String imUserId = userMappingService.getOrCreate(tenantId, bizUserId, nick, faceUrl);
+ // 签发 UserSig
+ long expireSec = expireDays * 86400L;
+ String userSig = UserSigUtil.genSig(sdkAppId, secretKey, imUserId, expireSec);
+
+ Map
+ * 用于 Sa-Token 登录鉴权。密码以 PBKDF2 哈希存储(格式见 {@link com.imutil.common.PasswordUtil})。
+ * role:admin=超级管理员,viewer=只读(预留)。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("admin_user")
+public class AdminUser {
+
+ @TableId(type = IdType.AUTO)
+ private Long id;
+
+ /** 登录用户名 */
+ private String username;
+
+ /** 密码hash(格式 iterations:salt:hash,见 PasswordUtil) */
+ private String passwordHash;
+
+ /** 角色:admin / viewer */
+ private String role;
+
+ /** 状态:1=启用 0=停用 */
+ private Integer status;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/ApiCallLog.java b/src/main/java/com/imutil/entity/ApiCallLog.java
new file mode 100644
index 0000000..ae53191
--- /dev/null
+++ b/src/main/java/com/imutil/entity/ApiCallLog.java
@@ -0,0 +1,36 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 腾讯管理 API 调用审计实体
+ *
+ * 所有调腾讯后台 API 的操作记录,便于配额核算与问题追溯。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("api_call_log")
+public class ApiCallLog {
+
+ @TableId(type = IdType.AUTO)
+ private Long id;
+
+ private String tenantId;
+
+ private String apiName;
+
+ /** 入参 JSON */
+ private String params;
+
+ private String result;
+
+ private String caller;
+
+ private OffsetDateTime calledAt;
+}
diff --git a/src/main/java/com/imutil/entity/CrossTenantAudit.java b/src/main/java/com/imutil/entity/CrossTenantAudit.java
new file mode 100644
index 0000000..61a2f94
--- /dev/null
+++ b/src/main/java/com/imutil/entity/CrossTenantAudit.java
@@ -0,0 +1,33 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 跨租户通讯审计实体
+ *
+ * 所有经授权放行的跨租户消息单独审计,便于追溯。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("cross_tenant_audit")
+public class CrossTenantAudit {
+
+ @TableId(type = IdType.AUTO)
+ private Long id;
+
+ private Long grantId;
+
+ private String msgKey;
+
+ private String fromImUserId;
+
+ private String toImUserId;
+
+ private OffsetDateTime actionTime;
+}
diff --git a/src/main/java/com/imutil/entity/CrossTenantGrant.java b/src/main/java/com/imutil/entity/CrossTenantGrant.java
new file mode 100644
index 0000000..a7d63c8
--- /dev/null
+++ b/src/main/java/com/imutil/entity/CrossTenantGrant.java
@@ -0,0 +1,51 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 跨租户通讯授权实体
+ *
+ * 用于 3.4 跨租户拦截的放行依据。from_im_user_id / to_im_user_id 为 NULL
+ * 分别表示授权方任意账户 / 目标租户全员。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("cross_tenant_grant")
+public class CrossTenantGrant {
+
+ @TableId(type = IdType.ASSIGN_ID)
+ private Long grantId;
+
+ private String fromTenantId;
+
+ private String fromImUserId;
+
+ private String toTenantId;
+
+ private String toImUserId;
+
+ /** 权限:send_msg,add_friend,join_group */
+ private String permissions;
+
+ /** 方向:0=单向 1=双向 */
+ private Integer direction;
+
+ private OffsetDateTime startAt;
+
+ private OffsetDateTime endAt;
+
+ /** 1=active 0=revoked 2=expired */
+ private Integer status;
+
+ private String approvedByFrom;
+
+ private String approvedByTo;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/DistQueue.java b/src/main/java/com/imutil/entity/DistQueue.java
new file mode 100644
index 0000000..2f4b7a5
--- /dev/null
+++ b/src/main/java/com/imutil/entity/DistQueue.java
@@ -0,0 +1,52 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 分发队列实体(替代 MQ)
+ *
+ * 回调网关同事务写入 pending,分发工作线程 FOR UPDATE SKIP LOCKED 抢占消费。
+ * id 为 bigserial 严格递增,保证同会话消息按 id 顺序消费。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("dist_queue")
+public class DistQueue {
+
+ /** 自增ID(消费顺序依据) */
+ @TableId(type = IdType.AUTO)
+ private Long id;
+
+ private String msgKey;
+
+ private String tenantId;
+
+ private String convId;
+
+ /** 业务系统回调地址 */
+ private String targetUrl;
+
+ /** 分发给业务系统的回调快照 JSON */
+ private String payload;
+
+ /** 状态:0=pending 1=processing 2=done 3=dead */
+ private Integer status;
+
+ private Integer retryCount;
+
+ private OffsetDateTime nextRetryAt;
+
+ private String lockedBy;
+
+ private OffsetDateTime lockedAt;
+
+ private OffsetDateTime createdAt;
+
+ private OffsetDateTime updatedAt;
+}
diff --git a/src/main/java/com/imutil/entity/GroupMapping.java b/src/main/java/com/imutil/entity/GroupMapping.java
new file mode 100644
index 0000000..c41ef9b
--- /dev/null
+++ b/src/main/java/com/imutil/entity/GroupMapping.java
@@ -0,0 +1,34 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 群映射实体
+ *
+ * im_group_id 由本工具统一分配,避免各系统自建群号撞号。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("group_mapping")
+public class GroupMapping {
+
+ @TableId(type = IdType.ASSIGN_ID)
+ private Long id;
+
+ private String tenantId;
+
+ private String bizGroupId;
+
+ private String imGroupId;
+
+ /** 群类型:Public/Private/ChatRoom/AVChatRoom */
+ private String groupType;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/ImMessage.java b/src/main/java/com/imutil/entity/ImMessage.java
new file mode 100644
index 0000000..db984b0
--- /dev/null
+++ b/src/main/java/com/imutil/entity/ImMessage.java
@@ -0,0 +1,60 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 消息主表实体(按月 RANGE 分区)
+ *
+ * 复合主键 (msg_key, msg_time),msg_time 同时为分区键。
+ * 注意:msg_time 为分区键,禁止更新(更新分区键会触发行迁移错误)。
+ * msg_body 存原始回调 JSON 字符串。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("im_message")
+public class ImMessage {
+
+ /** 腾讯 MsgKey,去重用 */
+ @TableId(type = IdType.INPUT)
+ private String msgKey;
+
+ /** 租户ID(取发送方所属租户) */
+ private String tenantId;
+
+ /** 消息时间(分区键) */
+ private OffsetDateTime msgTime;
+
+ /** 会话类型:1=C2C 2=GROUP */
+ private Integer convType;
+
+ /** 会话ID:C2C=对端账号 GROUP=群ID */
+ private String convId;
+
+ private String fromAccount;
+
+ private String toAccount;
+
+ private String groupId;
+
+ private String msgType;
+
+ /** 消息体原始 JSON */
+ private String msgBody;
+
+ /** 来源:CALLBACK/PULL_BACK/IMPORT */
+ private String source;
+
+ /** 是否跨租户授权通讯 */
+ private Boolean isCrossTenant;
+
+ /** 分发状态:0=待分发 1=已分发 2=失败 */
+ private Integer distStatus;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/PullWatermark.java b/src/main/java/com/imutil/entity/PullWatermark.java
new file mode 100644
index 0000000..ea4daa1
--- /dev/null
+++ b/src/main/java/com/imutil/entity/PullWatermark.java
@@ -0,0 +1,35 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 补拉水位线实体
+ *
+ * 记录每个会话最后拉取的消息 Seq/时间,补拉服务增量拉取游标。
+ * Redis 兜底,PG 持久化。复合主键 (tenant_id, conv_id)。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("pull_watermark")
+public class PullWatermark {
+
+ @TableId(type = IdType.INPUT)
+ private String tenantId;
+
+ private String convId;
+
+ /** 会话类型:1=C2C 2=GROUP */
+ private Integer convType;
+
+ private Long lastSeq;
+
+ private OffsetDateTime lastTime;
+
+ private OffsetDateTime updatedAt;
+}
diff --git a/src/main/java/com/imutil/entity/Recording.java b/src/main/java/com/imutil/entity/Recording.java
new file mode 100644
index 0000000..b3eb79e
--- /dev/null
+++ b/src/main/java/com/imutil/entity/Recording.java
@@ -0,0 +1,34 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 录制文件实体
+ *
+ * cos_path 按 tenant_id 目录隔离,各系统只能查/下载本租户录制。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("recording")
+public class Recording {
+
+ @TableId(type = IdType.INPUT)
+ private String fileId;
+
+ private String tenantId;
+
+ private Long roomId;
+
+ private String cosPath;
+
+ /** 时长(秒) */
+ private Integer duration;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/Tenant.java b/src/main/java/com/imutil/entity/Tenant.java
new file mode 100644
index 0000000..bfaf84d
--- /dev/null
+++ b/src/main/java/com/imutil/entity/Tenant.java
@@ -0,0 +1,52 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 租户表实体
+ *
+ * 每个业务系统对应一个租户,tenant_id 即前缀码(如 sa/sb),用于多租户隔离。
+ * app_key/app_secret 为业务系统调用本工具 REST API 的凭证。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("tenant")
+public class Tenant {
+
+ /** 租户ID(前缀码,如 sa) */
+ @TableId(type = IdType.INPUT)
+ private String tenantId;
+
+ /** 租户名称 */
+ private String tenantName;
+
+ /** 业务系统调用凭证 key */
+ private String appKey;
+
+ /** 业务系统调用凭证 secret */
+ private String appSecret;
+
+ /** IM UserID 前缀,与 tenantId 一致 */
+ private String prefixCode;
+
+ /** 该租户的回调分发地址 */
+ private String callbackUrl;
+
+ /** IM API QPS 配额 */
+ private Integer quotaImQps;
+
+ /** TRTC 并发房间配额 */
+ private Integer quotaTrtcConcurrent;
+
+ /** 状态:1=启用 0=停用 */
+ private Integer status;
+
+ /** 创建时间 */
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/TrtcRoom.java b/src/main/java/com/imutil/entity/TrtcRoom.java
new file mode 100644
index 0000000..1edea1c
--- /dev/null
+++ b/src/main/java/com/imutil/entity/TrtcRoom.java
@@ -0,0 +1,35 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * TRTC 音视频房间实体
+ *
+ * room_id 由本工具统一分配(雪花),避免各系统自建房号撞号。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("trtc_room")
+public class TrtcRoom {
+
+ @TableId(type = IdType.ASSIGN_ID)
+ private Long roomId;
+
+ private String tenantId;
+
+ private String bizRoomId;
+
+ /** 关联的 IM 群ID(可选) */
+ private String imGroupId;
+
+ /** 状态:1=进行中 0=已结束 */
+ private Integer status;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/UsageStat.java b/src/main/java/com/imutil/entity/UsageStat.java
new file mode 100644
index 0000000..1b801cf
--- /dev/null
+++ b/src/main/java/com/imutil/entity/UsageStat.java
@@ -0,0 +1,42 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 用量统计实体(计费拆账依据)
+ *
+ * stat_level:1=小时 2=天。按 tenant_id 聚合 IM/TRTC 用量。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("usage_stat")
+public class UsageStat {
+
+ @TableId(type = IdType.ASSIGN_ID)
+ private Long id;
+
+ private String tenantId;
+
+ private OffsetDateTime statTime;
+
+ /** 1=小时 2=天 */
+ private Integer statLevel;
+
+ private Long imMsgCount;
+
+ private Long imDau;
+
+ private Long trtcDurationSec;
+
+ private Integer trtcMaxConcurrentRoom;
+
+ private Long apiCallCount;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/entity/UserMapping.java b/src/main/java/com/imutil/entity/UserMapping.java
new file mode 100644
index 0000000..7b76776
--- /dev/null
+++ b/src/main/java/com/imutil/entity/UserMapping.java
@@ -0,0 +1,44 @@
+package com.imutil.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.Data;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 用户映射实体(业务用户 ↔ IM 用户)
+ *
+ * im_user_id = prefix_code + '_' + biz_user_id,前缀法保证跨租户不撞号。
+ * is_default 标记默认账户(admin/客服等),is_global 标记跨租户通行账户。
+ *
+ * @author imutil
+ */
+@Data
+@TableName("user_mapping")
+public class UserMapping {
+
+ @TableId(type = IdType.ASSIGN_ID)
+ private Long id;
+
+ /** 租户ID */
+ private String tenantId;
+
+ /** 业务系统用户ID */
+ private String bizUserId;
+
+ /** IM 用户ID(带前缀) */
+ private String imUserId;
+
+ /** 是否默认账户 */
+ private Boolean isDefault;
+
+ /** 是否全局跨租户账户 */
+ private Boolean isGlobal;
+
+ /** 状态:1=正常 0=封禁 */
+ private Integer status;
+
+ private OffsetDateTime createdAt;
+}
diff --git a/src/main/java/com/imutil/filter/AdminAuthFilter.java b/src/main/java/com/imutil/filter/AdminAuthFilter.java
new file mode 100644
index 0000000..c906e6a
--- /dev/null
+++ b/src/main/java/com/imutil/filter/AdminAuthFilter.java
@@ -0,0 +1,52 @@
+package com.imutil.filter;
+
+import cn.dev33.satoken.stp.StpUtil;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.core.handle.Context;
+import org.noear.solon.core.handle.Filter;
+import org.noear.solon.core.handle.FilterChain;
+
+/**
+ * 管理后台鉴权过滤器
+ *
+ * 拦截 {@code /admin/**},未登录 Sa-Token 跳转登录页。
+ * {@code /admin/login}、{@code /admin/logout} 放行;
+ * {@code /admin} 已在 {@link com.imutil.common.PathWhitelist},不走路租户鉴权({@link TenantAuthFilter})。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class AdminAuthFilter implements Filter {
+
+ @Inject("${server.contextPath:}")
+ private String basePath;
+
+ @Override
+ public void doFilter(Context ctx, FilterChain chain) throws Throwable {
+ String path = ctx.pathNew();
+ // 仅拦截 /admin/**
+ if (!path.startsWith("/admin")) {
+ chain.doFilter(ctx);
+ return;
+ }
+ // 登录/登出页放行
+ if (path.equals("/admin/login") || path.equals("/admin/logout")) {
+ chain.doFilter(ctx);
+ return;
+ }
+ // 已登录放行
+ if (StpUtil.isLogin()) {
+ chain.doFilter(ctx);
+ return;
+ }
+ // 未登录:GET 跳登录页,其他返回 401 JSON
+ if ("GET".equalsIgnoreCase(ctx.method())) {
+ ctx.redirect(basePath + "/admin/login");
+ } else {
+ ctx.output("{\"success\":false,\"code\":401,\"message\":\"未登录或会话已过期\"}");
+ }
+ }
+}
diff --git a/src/main/java/com/imutil/filter/GlobalExceptionFilter.java b/src/main/java/com/imutil/filter/GlobalExceptionFilter.java
new file mode 100644
index 0000000..979188d
--- /dev/null
+++ b/src/main/java/com/imutil/filter/GlobalExceptionFilter.java
@@ -0,0 +1,73 @@
+package com.imutil.filter;
+
+import com.imutil.common.BizException;
+import com.imutil.model.Result;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.core.exception.StatusException;
+import org.noear.solon.core.handle.Context;
+import org.noear.solon.core.handle.Filter;
+import org.noear.solon.core.handle.FilterChain;
+
+/**
+ * 全局异常拦截器
+ *
+ * 捕获控制器抛出的所有异常,统一以 {@link Result#error} 格式返回,
+ * 避免直接暴露框架异常信息。HTTP 状态码始终 200,错误经 Result.code 区分。
+ *
+ * 处理规则:
+ * - {@link BizException}:业务异常,按其 code/message 返回
+ * - {@link StatusException}:框架状态异常(401/403/404/405 等)
+ * - 其他 Throwable:服务器内部错误
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class GlobalExceptionFilter implements Filter {
+
+ @Override
+ public void doFilter(Context ctx, FilterChain chain) throws Throwable {
+ try {
+ chain.doFilter(ctx);
+ } catch (BizException e) {
+ log.warn("业务异常 [{} {}] code={} : {}", ctx.method(), ctx.path(), e.getCode(), e.getMessage());
+ renderJson(ctx, Result.error(e.getCode(), e.getMessage()));
+ } catch (StatusException e) {
+ int code = e.getCode();
+ // 浏览器自动发起的 favicon.ico,静默 204
+ if (code == 404 && "/favicon.ico".equals(ctx.path())) {
+ ctx.status(204);
+ return;
+ }
+ String msg = resolveMessage(code);
+ log.warn("请求异常 [{} {}] status={} : {}", ctx.method(), ctx.path(), code, e.getMessage());
+ renderJson(ctx, Result.error(code, msg));
+ } catch (Throwable e) {
+ log.error("系统未知异常 [{} {}]", ctx.method(), ctx.path(), e);
+ renderJson(ctx, Result.error("服务器内部错误"));
+ }
+ }
+
+ /**
+ * 按 HTTP 状态码返回中文提示
+ */
+ private String resolveMessage(int code) {
+ return switch (code) {
+ case 400 -> "请求参数错误,请检查入参是否完整";
+ case 401 -> "未登录或登录已过期,请重新登录";
+ case 403 -> "无权限访问该资源";
+ case 404 -> "请求的接口不存在";
+ case 405 -> "请求方法不允许,请检查请求方式(GET/POST 等)";
+ default -> code >= 500 ? "服务处理失败,请稍后重试" : "请求异常(" + code + ")";
+ };
+ }
+
+ /**
+ * 序列化为 JSON 写入响应,HTTP 200
+ */
+ private void renderJson(Context ctx, Result> result) throws Throwable {
+ ctx.status(200);
+ ctx.render(result);
+ }
+}
diff --git a/src/main/java/com/imutil/filter/TenantAuthFilter.java b/src/main/java/com/imutil/filter/TenantAuthFilter.java
new file mode 100644
index 0000000..f968561
--- /dev/null
+++ b/src/main/java/com/imutil/filter/TenantAuthFilter.java
@@ -0,0 +1,141 @@
+package com.imutil.filter;
+
+import com.imutil.common.BizException;
+import com.imutil.common.PathWhitelist;
+import com.imutil.common.RedisService;
+import com.imutil.common.TenantContext;
+import com.imutil.entity.Tenant;
+import com.imutil.model.Result;
+import com.imutil.service.TenantService;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.core.handle.Context;
+import org.noear.solon.core.handle.Filter;
+import org.noear.solon.core.handle.FilterChain;
+
+import java.nio.charset.StandardCharsets;
+import java.util.Base64;
+
+/**
+ * 租户鉴权过滤器
+ *
+ * 拦截所有非白名单路径,校验 app_key/app_secret 凭证,解析租户写入 {@link TenantContext}。
+ * 凭证传递方式(二选一):
+ * - Header {@code X-App-Key} + {@code X-App-Secret}
+ * - HTTP Basic:{@code Authorization: Basic base64(app_key:app_secret)}
+ *
+ * 多租户凭证各自独立,不能用单一 Sa-Token http-basic,故自实现校验。
+ * 白名单(/callback、/health)不走本过滤,回调走签名校验。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class TenantAuthFilter implements Filter {
+
+ private static final String HDR_KEY = "X-App-Key";
+ private static final String HDR_SECRET = "X-App-Secret";
+ private static final String HDR_AUTH = "Authorization";
+
+ @Inject
+ private TenantService tenantService;
+
+ @Inject
+ private RedisService redisService;
+
+ @Override
+ public void doFilter(Context ctx, FilterChain chain) throws Throwable {
+ // 去除 contextPath 后的路径
+ String path = ctx.pathNew();
+ if (PathWhitelist.isWhitelisted(path)) {
+ chain.doFilter(ctx);
+ return;
+ }
+
+ try {
+ String[] cred = extractCredential(ctx);
+ if (cred == null) {
+ deny(ctx, 401, "缺少鉴权凭证(X-App-Key/X-App-Secret 或 Basic)");
+ return;
+ }
+ String appKey = cred[0];
+ String appSecret = cred[1];
+
+ Tenant tenant = tenantService.getByAppKey(appKey);
+ if (tenant == null) {
+ deny(ctx, 401, "app_key 无效");
+ return;
+ }
+ if (tenant.getStatus() != null && tenant.getStatus() != 1) {
+ deny(ctx, 403, "租户已停用");
+ return;
+ }
+ // 密钥校验(常量时间比较,防时序攻击)
+ if (tenant.getAppSecret() == null || !constantTimeEquals(tenant.getAppSecret(), appSecret)) {
+ deny(ctx, 401, "app_secret 错误");
+ return;
+ }
+
+ // 防爆破:记录失败计数已清,此处鉴权成功
+ TenantContext.set(tenant.getTenantId());
+ chain.doFilter(ctx);
+ } finally {
+ // 请求结束清理,防线程复用串租户
+ TenantContext.clear();
+ }
+ }
+
+ /**
+ * 提取凭证:优先自定义 Header,其次 HTTP Basic
+ *
+ * @return [appKey, appSecret],提取失败返回 null
+ */
+ private String[] extractCredential(Context ctx) {
+ String key = ctx.header(HDR_KEY);
+ String secret = ctx.header(HDR_SECRET);
+ if (key != null && !key.isEmpty() && secret != null && !secret.isEmpty()) {
+ return new String[]{key, secret};
+ }
+ // HTTP Basic
+ String auth = ctx.header(HDR_AUTH);
+ if (auth != null && auth.startsWith("Basic ")) {
+ try {
+ String decoded = new String(Base64.getDecoder().decode(auth.substring(6)), StandardCharsets.UTF_8);
+ int idx = decoded.indexOf(':');
+ if (idx > 0) {
+ return new String[]{decoded.substring(0, idx), decoded.substring(idx + 1)};
+ }
+ } catch (IllegalArgumentException e) {
+ return null;
+ }
+ }
+ return null;
+ }
+
+ /**
+ * 常量时间字符串比较,防止时序攻击
+ */
+ private boolean constantTimeEquals(String a, String b) {
+ if (a == null || b == null) {
+ return false;
+ }
+ if (a.length() != b.length()) {
+ return false;
+ }
+ int r = 0;
+ for (int i = 0; i < a.length(); i++) {
+ r |= a.charAt(i) ^ b.charAt(i);
+ }
+ return r == 0;
+ }
+
+ /**
+ * 返回鉴权失败响应(HTTP 200 + Result 401/403)
+ */
+ private void deny(Context ctx, int code, String msg) throws Throwable {
+ log.warn("鉴权失败 [{} {}] code={} : {}", ctx.method(), ctx.path(), code, msg);
+ ctx.status(200);
+ ctx.render(Result.error(code, msg));
+ }
+}
diff --git a/src/main/java/com/imutil/mapper/AdminUserMapper.java b/src/main/java/com/imutil/mapper/AdminUserMapper.java
new file mode 100644
index 0000000..d6892e3
--- /dev/null
+++ b/src/main/java/com/imutil/mapper/AdminUserMapper.java
@@ -0,0 +1,12 @@
+package com.imutil.mapper;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.imutil.entity.AdminUser;
+
+/**
+ * 管理后台用户 Mapper
+ *
+ * @author imutil
+ */
+public interface AdminUserMapper extends BaseMapper
+ * 提供 FOR UPDATE SKIP LOCKED 抢占消费的 SQL。
+ *
+ * @author imutil
+ */
+public interface DistQueueMapper extends BaseMapper
+ * upsert 用 PG ON CONFLICT,selectNeedCheck 按 updated_at 升序取最久未补拉的会话,
+ * 配合补拉成功后推进 updated_at,实现会话轮询公平覆盖。
+ *
+ * @author imutil
+ */
+public interface PullWatermarkMapper extends BaseMapper
+ * 回调落库 / 补拉推进时调用。last_seq、last_time 取 GREATEST 不回退(防乱序回调回退游标),
+ * updated_at 始终刷新为 now(驱动 selectNeedCheck 的轮询顺序)。
+ *
+ * @param tenantId 租户ID
+ * @param convId 会话ID
+ * @param convType 会话类型 1=C2C 2=GROUP
+ * @param lastSeq 最新消息 Seq
+ * @param lastTime 最新消息时间
+ * @param now 当前时间(updated_at)
+ * @return 影响行数(1=新增或更新)
+ */
+ @Insert("INSERT INTO pull_watermark (tenant_id, conv_id, conv_type, last_seq, last_time, updated_at) " +
+ "VALUES (#{tenantId}, #{convId}, #{convType}, #{lastSeq}, #{lastTime}, #{now}) " +
+ "ON CONFLICT (tenant_id, conv_id) DO UPDATE SET " +
+ "last_seq = GREATEST(pull_watermark.last_seq, EXCLUDED.last_seq), " +
+ "last_time = GREATEST(pull_watermark.last_time, EXCLUDED.last_time), " +
+ "updated_at = EXCLUDED.updated_at")
+ int upsert(@Param("tenantId") String tenantId,
+ @Param("convId") String convId,
+ @Param("convType") int convType,
+ @Param("lastSeq") long lastSeq,
+ @Param("lastTime") OffsetDateTime lastTime,
+ @Param("now") OffsetDateTime now);
+
+ /**
+ * 取最久未补拉的 N 个会话(按 updated_at 升序)
+ *
+ * 补拉任务每轮调用,updated_at 最老的优先;补拉后 updated_at 推进到 now,
+ * 该会话自然排到队尾,实现轮询式公平覆盖。
+ *
+ * @param limit 每轮会话数
+ * @return 待补拉会话列表
+ */
+ @Select("SELECT tenant_id, conv_id, conv_type, last_seq, last_time, updated_at " +
+ "FROM pull_watermark ORDER BY updated_at ASC LIMIT #{limit}")
+ List
+ * 提供按小时窗口的聚合 upsert(PG ON CONFLICT 幂等,可重复跑)。
+ * 两个数据源各自 upsert 指定列,避免互相覆盖:
+ * - im_message → im_msg_count / im_dau
+ * - api_call_log → api_call_count
+ *
+ * @author imutil
+ */
+public interface UsageStatMapper extends BaseMapper
+ * 仅更新 im_msg_count / im_dau,不触碰已有 api_call_count。
+ *
+ * @param start 窗口起点(含)
+ * @param end 窗口终点(不含)
+ * @param statTime 统计时刻(=窗口起点整点,作为 usage_stat.stat_time)
+ * @return 受影响租户数
+ */
+ @Insert("INSERT INTO usage_stat (tenant_id, stat_time, stat_level, im_msg_count, im_dau, " +
+ "api_call_count, trtc_duration_sec, trtc_max_concurrent_room, created_at) " +
+ "SELECT tenant_id, #{statTime}, 1, count(*), count(DISTINCT from_account), 0, 0, 0, now() " +
+ "FROM im_message WHERE created_at >= #{start} AND created_at < #{end} " +
+ "GROUP BY tenant_id " +
+ "ON CONFLICT (tenant_id, stat_time, stat_level) DO UPDATE SET " +
+ "im_msg_count = EXCLUDED.im_msg_count, im_dau = EXCLUDED.im_dau")
+ int aggregateMsgHour(@Param("start") OffsetDateTime start,
+ @Param("end") OffsetDateTime end,
+ @Param("statTime") OffsetDateTime statTime);
+
+ /**
+ * 从 api_call_log 聚合某小时窗口的腾讯 API 调用次数,upsert 到 usage_stat(stat_level=1 小时)
+ *
+ * 仅更新 api_call_count,不触碰已有 im_msg_count / im_dau。
+ *
+ * @return 受影响租户数
+ */
+ @Insert("INSERT INTO usage_stat (tenant_id, stat_time, stat_level, im_msg_count, im_dau, " +
+ "api_call_count, trtc_duration_sec, trtc_max_concurrent_room, created_at) " +
+ "SELECT tenant_id, #{statTime}, 1, 0, 0, count(*), 0, 0, now() " +
+ "FROM api_call_log WHERE called_at >= #{start} AND called_at < #{end} " +
+ "GROUP BY tenant_id " +
+ "ON CONFLICT (tenant_id, stat_time, stat_level) DO UPDATE SET " +
+ "api_call_count = EXCLUDED.api_call_count")
+ int aggregateApiCallHour(@Param("start") OffsetDateTime start,
+ @Param("end") OffsetDateTime end,
+ @Param("statTime") OffsetDateTime statTime);
+}
diff --git a/src/main/java/com/imutil/mapper/UserMappingMapper.java b/src/main/java/com/imutil/mapper/UserMappingMapper.java
new file mode 100644
index 0000000..9a77e63
--- /dev/null
+++ b/src/main/java/com/imutil/mapper/UserMappingMapper.java
@@ -0,0 +1,12 @@
+package com.imutil.mapper;
+
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import com.imutil.entity.UserMapping;
+
+/**
+ * 用户映射 Mapper
+ *
+ * @author imutil
+ */
+public interface UserMappingMapper extends BaseMapper
+ * 所有 REST 接口与回调内部响应统一使用本类。HTTP 状态码始终为 200,
+ * 成功/失败通过 {@link #code} 区分(200 成功,其余失败)。
+ *
+ * @author imutil
+ */
+@Data
+public class Result
+ * 登录校验(PBKDF2)+ 首次启动默认管理员初始化。
+ *
+ * @author imutil
+ */
+public interface AdminUserService {
+
+ /**
+ * 登录校验
+ *
+ * @return 凭证正确且账号启用返回 AdminUser,否则 null
+ */
+ AdminUser login(String username, String password);
+
+ /**
+ * 应用启动时确保存在至少一个管理员(表空则按 app.yml imutil.admin 初始化)
+ */
+ void ensureDefaultAdmin();
+}
diff --git a/src/main/java/com/imutil/service/CallbackService.java b/src/main/java/com/imutil/service/CallbackService.java
new file mode 100644
index 0000000..e135c0b
--- /dev/null
+++ b/src/main/java/com/imutil/service/CallbackService.java
@@ -0,0 +1,18 @@
+package com.imutil.service;
+
+/**
+ * 腾讯回调处理服务
+ *
+ * @author imutil
+ */
+public interface CallbackService {
+
+ /**
+ * 处理腾讯回调:落库消息(消息类回调)+ 写分发队列
+ *
+ * @param callbackCommand 回调命令,如 C2C.CallbackAfterRecvMsg
+ * @param body 回调请求体 JSON
+ * @return 处理结果,OK/FAIL
+ */
+ String handleCallback(String callbackCommand, String body);
+}
diff --git a/src/main/java/com/imutil/service/DispatchService.java b/src/main/java/com/imutil/service/DispatchService.java
new file mode 100644
index 0000000..cd3b96e
--- /dev/null
+++ b/src/main/java/com/imutil/service/DispatchService.java
@@ -0,0 +1,35 @@
+package com.imutil.service;
+
+import com.imutil.entity.DistQueue;
+
+import java.util.List;
+
+/**
+ * 回调分发服务
+ *
+ * 从 dist_queue 抢占待分发记录,HTTP 转发业务系统,按回执更新状态。
+ *
+ * @author imutil
+ */
+public interface DispatchService {
+
+ /**
+ * 抢占一批 pending 记录置为 processing(事务内 FOR UPDATE SKIP LOCKED + lock)
+ *
+ * @param workerName 工作线程标识
+ * @return 抢占到的记录列表
+ */
+ List
+ * 回调链路的兜底:腾讯回调丢失或本工具落库失败时,定期主动从腾讯拉取最新消息补全本地,
+ * 保证 im_message 不丢、业务系统不漏收(对应「已知问题」第 3 条)。
+ *
+ * @author imutil
+ */
+public interface PullService {
+
+ /**
+ * 执行一轮补拉:取最久未补拉的若干会话,逐个调腾讯 API 拉最新消息,
+ * msg_key 幂等落库(source=PULL_BACK)并入队分发,最后推进水位线
+ */
+ void pullRound();
+}
diff --git a/src/main/java/com/imutil/service/TenantService.java b/src/main/java/com/imutil/service/TenantService.java
new file mode 100644
index 0000000..db7cce7
--- /dev/null
+++ b/src/main/java/com/imutil/service/TenantService.java
@@ -0,0 +1,29 @@
+package com.imutil.service;
+
+import com.imutil.entity.Tenant;
+
+/**
+ * 租户服务
+ *
+ * @author imutil
+ */
+public interface TenantService {
+
+ /**
+ * 根据 app_key 查询租户(带缓存)
+ *
+ * @param appKey 业务系统凭证 key
+ * @return 租户实体,不存在返回 null
+ */
+ Tenant getByAppKey(String appKey);
+
+ /**
+ * 根据 tenant_id 查询租户(带缓存)
+ */
+ Tenant getById(String tenantId);
+
+ /**
+ * 失效租户缓存(租户变更时调用)
+ */
+ void evictCache(String tenantId);
+}
diff --git a/src/main/java/com/imutil/service/UsageStatService.java b/src/main/java/com/imutil/service/UsageStatService.java
new file mode 100644
index 0000000..87b8c09
--- /dev/null
+++ b/src/main/java/com/imutil/service/UsageStatService.java
@@ -0,0 +1,23 @@
+package com.imutil.service;
+
+import java.time.OffsetDateTime;
+
+/**
+ * 用量统计服务
+ *
+ * 按小时窗口从 im_message / api_call_log 聚合用量到 usage_stat(计费拆账依据)。
+ *
+ * @author imutil
+ */
+public interface UsageStatService {
+
+ /**
+ * 聚合某小时窗口 [start, end) 的用量,upsert 到 usage_stat(stat_level=1)
+ *
+ * 幂等:ON CONFLICT,重复跑不产生重复数据,取最新聚合值。
+ *
+ * @param start 窗口起点(含,整点)
+ * @param end 窗口终点(不含)
+ */
+ void aggregateHour(OffsetDateTime start, OffsetDateTime end);
+}
diff --git a/src/main/java/com/imutil/service/UserMappingService.java b/src/main/java/com/imutil/service/UserMappingService.java
new file mode 100644
index 0000000..a035908
--- /dev/null
+++ b/src/main/java/com/imutil/service/UserMappingService.java
@@ -0,0 +1,25 @@
+package com.imutil.service;
+
+/**
+ * 账号映射服务
+ *
+ * 维护业务用户ID ↔ IM UserID(带租户前缀)映射,首次使用时调腾讯 account_import 创建 IM 账号。
+ *
+ * @author imutil
+ */
+public interface UserMappingService {
+
+ /**
+ * 获取或创建 IM 用户ID
+ *
+ * 不存在则调腾讯 account_import 创建 IM 账号并写映射;存在则直接返回。
+ * im_user_id = tenant_id + '_' + biz_user_id(前缀法保证跨租户不撞号)。
+ *
+ * @param tenantId 租户ID(即前缀码)
+ * @param bizUserId 业务用户ID
+ * @param nick 昵称(首次创建时同步到 IM,可选)
+ * @param faceUrl 头像(可选)
+ * @return IM 用户ID
+ */
+ String getOrCreate(String tenantId, String bizUserId, String nick, String faceUrl);
+}
diff --git a/src/main/java/com/imutil/service/impl/AdminUserServiceImpl.java b/src/main/java/com/imutil/service/impl/AdminUserServiceImpl.java
new file mode 100644
index 0000000..fcbe454
--- /dev/null
+++ b/src/main/java/com/imutil/service/impl/AdminUserServiceImpl.java
@@ -0,0 +1,77 @@
+package com.imutil.service.impl;
+
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import com.imutil.common.PasswordUtil;
+import com.imutil.entity.AdminUser;
+import com.imutil.mapper.AdminUserMapper;
+import com.imutil.service.AdminUserService;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Init;
+import org.noear.solon.annotation.Inject;
+
+/**
+ * 管理后台用户服务实现
+ *
+ * 密码用 {@link PasswordUtil}(PBKDF2)哈希存储与校验。
+ * {@code @Init} 在容器启动后检查 admin_user 表,为空则按 app.yml imutil.admin 初始化默认管理员。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class AdminUserServiceImpl implements AdminUserService {
+
+ @Inject
+ private AdminUserMapper adminUserMapper;
+
+ @Inject("${imutil.admin.defaultUsername:admin}")
+ private String defaultUsername;
+
+ @Inject("${imutil.admin.defaultPassword:admin123}")
+ private String defaultPassword;
+
+ /**
+ * 应用启动后初始化默认管理员(仅当表为空)
+ */
+ @Init
+ public void init() {
+ ensureDefaultAdmin();
+ }
+
+ @Override
+ public AdminUser login(String username, String password) {
+ if (username == null || username.isEmpty() || password == null || password.isEmpty()) {
+ return null;
+ }
+ AdminUser u = adminUserMapper.selectOne(Wrappers.
+ * 流程:识别租户 → 消息类回调落 im_message(幂等)→ 所有回调写 dist_queue 分发业务系统。
+ * {@code @Tran} 保证消息落库与分发入队同事务:要么同时成功,要么都不入库(避免半写)。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class CallbackServiceImpl implements CallbackService {
+
+ /** 消息类回调命令关键字(命中则额外落 im_message) */
+ private static final String[] MSG_COMMAND_KEYWORDS = {"SendMsg", "RecvMsg"};
+
+ @Inject
+ private ImMessageMapper imMessageMapper;
+
+ @Inject
+ private DistQueueMapper distQueueMapper;
+
+ @Inject
+ private GroupMappingMapper groupMappingMapper;
+
+ @Inject
+ private PullWatermarkMapper pullWatermarkMapper;
+
+ @Inject
+ private TenantService tenantService;
+
+ @Override
+ @Tran
+ public String handleCallback(String callbackCommand, String body) {
+ if (callbackCommand == null || callbackCommand.isEmpty()) {
+ return ok();
+ }
+ ONode node;
+ try {
+ node = ONode.ofJson(body == null ? "{}" : body);
+ } catch (Exception e) {
+ log.warn("回调 body 解析失败 command={} : {}", callbackCommand, e.getMessage());
+ return ok();
+ }
+
+ // 1. 识别租户
+ String tenantId = identifyTenant(callbackCommand, node);
+ if (tenantId == null) {
+ // 无法识别租户(如腾讯系统消息 administrator),不落库不分发
+ log.debug("回调无法识别租户,跳过 command={} from={}", callbackCommand, node.get("FromAccount").getString());
+ return ok();
+ }
+
+ // 2. 消息类回调落库 im_message(幂等:msg_key 存在则跳过)
+ // msg_key = command:from:convTarget:msgSeq:msgRandom,含 MsgSeq+MsgRandom 全局唯一,单独作幂等键;
+ // 不依赖 MsgTimeStamp(避免腾讯回调时间戳偏差导致漏判)。
+ // DB 主键 (msg_key, msg_time) 因分区表约束保留 msg_time,作兜底防护。
+ String msgKey = null;
+ if (isMessageCallback(callbackCommand)) {
+ ImMessage msg = parseMessage(callbackCommand, node, tenantId);
+ if (msg != null) {
+ msgKey = msg.getMsgKey();
+ long exists = imMessageMapper.selectCount(Wrappers.
+ * 抢占消费:{@code @Tran} 内 fetchPending(FOR UPDATE SKIP LOCKED) + lock,提交后释放行锁,
+ * 记录置 processing;随后 HTTP 转发业务系统,按回执更新 done/retry/dead。
+ * 失败采用指数退避:backoff = retryBaseMs * 2^min(retryCount, 6)。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class DispatchServiceImpl implements DispatchService {
+
+ @Inject
+ private DistQueueMapper distQueueMapper;
+
+ @Inject("${imutil.dispatch.fetchBatch:50}")
+ private int fetchBatch;
+
+ @Inject("${imutil.dispatch.maxRetry:5}")
+ private int maxRetry;
+
+ @Inject("${imutil.dispatch.retryBaseMs:2000}")
+ private long retryBaseMs;
+
+ @Inject("${imutil.dispatch.lockTimeoutMin:3}")
+ private int lockTimeoutMin;
+
+ @Override
+ @Tran
+ public List
+ * 流程:取最久未补拉的 N 个会话 → 按会话类型调腾讯 API 拉最新消息 →
+ * 逐条 msg_key 幂等(与回调路径共享 {@link MsgKeys} 算法)→ 本地不存在则落库(source=PULL_BACK)+入队分发 → 推进水位线。
+ *
+ * 仅"本次新插入"的消息入队分发,避免对回调已正常落库的消息重复推送业务系统。
+ * 每轮串行执行 + 会话数/条数上限,受腾讯 API QPS 约束。
+ *
+ * C2C 会话需 from+to 配对,水位线仅存对端(conv_id),本端取该会话最近一条消息的 from_account;
+ * 多发送方场景仅覆盖最近一个 from(已知限制,由腾讯 get_roam_msg 成对特性决定)。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class PullServiceImpl implements PullService {
+
+ @Inject
+ private PullWatermarkMapper pullWatermarkMapper;
+
+ @Inject
+ private ImMessageMapper imMessageMapper;
+
+ @Inject
+ private DistQueueMapper distQueueMapper;
+
+ @Inject
+ private TencentImClient tencentImClient;
+
+ @Inject
+ private TenantService tenantService;
+
+ @Inject("${imutil.pull.convsPerRound:20}")
+ private int convsPerRound;
+
+ @Inject("${imutil.pull.maxMsgPerConv:20}")
+ private int maxMsgPerConv;
+
+ @Inject("${imutil.pull.lookbackMinutes:30}")
+ private int lookbackMinutes;
+
+ @Override
+ public void pullRound() {
+ List
+ * app_key → tenant 查询走 Redis 主缓存 + Caffeine 本地兜底,降低 DB 压力。
+ * 缓存 TTL 5 分钟,租户变更需主动 evict。
+ *
+ * @author imutil
+ */
+@Component
+public class TenantServiceImpl implements TenantService {
+
+ private static final String CACHE_KEY_BY_KEY = "imutil:tenant:bykey:";
+ private static final String CACHE_KEY_BY_ID = "imutil:tenant:byid:";
+ private static final long CACHE_TTL_SEC = 300;
+
+ @Inject
+ private TenantMapper tenantMapper;
+
+ @Inject
+ private RedisService redisService;
+
+ @Inject
+ private LocalCache localCache;
+
+ @Override
+ public Tenant getByAppKey(String appKey) {
+ if (appKey == null || appKey.isEmpty()) {
+ return null;
+ }
+ String key = CACHE_KEY_BY_KEY + appKey;
+ // 1. 本地缓存
+ Tenant t = localCache.get(key, Tenant.class);
+ if (t != null) {
+ return t;
+ }
+ // 2. Redis
+ t = redisService.getJson(key, Tenant.class);
+ if (t == null) {
+ // 3. DB
+ t = tenantMapper.selectOne(Wrappers.
+ * 两个数据源各自 upsert 指定列(见 {@link UsageStatMapper}),互不覆盖:
+ * - im_message → im_msg_count / im_dau
+ * - api_call_log → api_call_count
+ * trtc_duration_sec / trtc_max_concurrent_room 待音视频模块实现后补充。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class UsageStatServiceImpl implements UsageStatService {
+
+ @Inject
+ private UsageStatMapper usageStatMapper;
+
+ @Override
+ public void aggregateHour(OffsetDateTime start, OffsetDateTime end) {
+ // statTime 用窗口起点(整点),作为 usage_stat 的统计时刻
+ OffsetDateTime statTime = start;
+ int msgTenants = usageStatMapper.aggregateMsgHour(start, end, statTime);
+ int apiTenants = usageStatMapper.aggregateApiCallHour(start, end, statTime);
+ log.info("用量小时聚合完成 window=[{}, {}) msg租户数={} api租户数={}", start, end, msgTenants, apiTenants);
+ }
+}
diff --git a/src/main/java/com/imutil/service/impl/UserMappingServiceImpl.java b/src/main/java/com/imutil/service/impl/UserMappingServiceImpl.java
new file mode 100644
index 0000000..415434c
--- /dev/null
+++ b/src/main/java/com/imutil/service/impl/UserMappingServiceImpl.java
@@ -0,0 +1,118 @@
+package com.imutil.service.impl;
+
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import com.imutil.common.BizException;
+import com.imutil.entity.UserMapping;
+import com.imutil.mapper.UserMappingMapper;
+import com.imutil.service.UserMappingService;
+import com.imutil.tencent.TencentImClient;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.data.annotation.Tran;
+
+import java.nio.charset.StandardCharsets;
+
+/**
+ * 账号映射服务实现
+ *
+ * im_user_id = tenant_id + '_' + biz_user_id(前缀法)。
+ * 首次创建:调腾讯 account_import 导入 IM 账号 → 写 user_mapping(同事务)。
+ * 兼容账号已存在场景(account_import 失败但 account_check 命中则视为成功)。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class UserMappingServiceImpl implements UserMappingService {
+
+ @Inject
+ private UserMappingMapper userMappingMapper;
+
+ @Inject
+ private TencentImClient tencentImClient;
+
+ @Override
+ @Tran
+ public String getOrCreate(String tenantId, String bizUserId, String nick, String faceUrl) {
+ // 0. 拼接 IM UserID 并校验合法性(腾讯约束:UTF-8 ≤32字节,仅字母/数字/下划线/横线)
+ String imUserId = tenantId + "_" + bizUserId;
+ validateImUserId(imUserId);
+
+ // 1. 查映射是否存在
+ UserMapping exist = userMappingMapper.selectOne(Wrappers.
+ * 腾讯 IM 约束:长度 ≤ 32 字节(UTF-8),允许字母/数字/下划线/横线。
+ * bizUserId 由业务系统传入,可能含中文或特殊字符,需在拼出 imUserId 后前置校验,
+ * 避免透传腾讯后台的错误码(业务侧语意不清晰)。
+ *
+ * @param imUserId 待校验的 IM 用户ID
+ */
+ private void validateImUserId(String imUserId) {
+ if (imUserId == null || imUserId.isEmpty()) {
+ throw new BizException(400, "IM UserID 不能为空");
+ }
+ if (imUserId.getBytes(StandardCharsets.UTF_8).length > 32) {
+ throw new BizException(400, "IM UserID 过长(>32字节),请缩短 bizUserId");
+ }
+ for (int i = 0; i < imUserId.length(); i++) {
+ char c = imUserId.charAt(i);
+ boolean legal = (c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z')
+ || (c >= '0' && c <= '9') || c == '_' || c == '-';
+ if (!legal) {
+ throw new BizException(400, "IM UserID 含非法字符 '" + c + "',仅允许字母/数字/下划线/横线");
+ }
+ }
+ }
+}
diff --git a/src/main/java/com/imutil/task/DispatchRecoverTask.java b/src/main/java/com/imutil/task/DispatchRecoverTask.java
new file mode 100644
index 0000000..2442325
--- /dev/null
+++ b/src/main/java/com/imutil/task/DispatchRecoverTask.java
@@ -0,0 +1,38 @@
+package com.imutil.task;
+
+import com.imutil.service.DispatchService;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.scheduling.annotation.Scheduled;
+
+/**
+ * 分发卡死巡检任务
+ *
+ * 定时重置超时未回执的 processing 记录回 pending,避免工作线程宕机导致记录永久卡在 processing。
+ * 对应 app.yml 的 solon.scheduling.job.dispatchRecoverJob。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class DispatchRecoverTask {
+
+ @Inject
+ private DispatchService dispatchService;
+
+ /**
+ * 由 app.yml dispatchRecoverJob 驱动(默认 fixedDelay=60s)
+ */
+ @Scheduled(name = "dispatchRecoverJob")
+ public void run() {
+ try {
+ int n = dispatchService.recoverStuck();
+ if (n > 0) {
+ log.info("分发卡死巡检:重置 {} 条 processing 记录为 pending", n);
+ }
+ } catch (Throwable e) {
+ log.error("分发卡死巡检异常", e);
+ }
+ }
+}
diff --git a/src/main/java/com/imutil/task/PullCheckTask.java b/src/main/java/com/imutil/task/PullCheckTask.java
new file mode 100644
index 0000000..7ef9b29
--- /dev/null
+++ b/src/main/java/com/imutil/task/PullCheckTask.java
@@ -0,0 +1,36 @@
+package com.imutil.task;
+
+import com.imutil.service.PullService;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.scheduling.annotation.Scheduled;
+
+/**
+ * 历史消息补拉巡检任务
+ *
+ * 定时触发补拉服务,兜底腾讯回调丢失 / 本工具落库失败的消息,
+ * 保证 im_message 不丢、业务系统不漏收。
+ * 对应 app.yml 的 solon.scheduling.job.pullCheckJob。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class PullCheckTask {
+
+ @Inject
+ private PullService pullService;
+
+ /**
+ * 由 app.yml pullCheckJob 驱动(默认 fixedDelay=5分钟)
+ */
+ @Scheduled(name = "pullCheckJob")
+ public void run() {
+ try {
+ pullService.pullRound();
+ } catch (Throwable e) {
+ log.error("历史消息补拉巡检异常", e);
+ }
+ }
+}
diff --git a/src/main/java/com/imutil/task/UsageStatTask.java b/src/main/java/com/imutil/task/UsageStatTask.java
new file mode 100644
index 0000000..a8b2d46
--- /dev/null
+++ b/src/main/java/com/imutil/task/UsageStatTask.java
@@ -0,0 +1,43 @@
+package com.imutil.task;
+
+import com.imutil.service.UsageStatService;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+import org.noear.solon.scheduling.annotation.Scheduled;
+
+import java.time.OffsetDateTime;
+import java.time.temporal.ChronoUnit;
+
+/**
+ * 用量统计定时聚合任务
+ *
+ * 每小时聚合「上一整点小时」窗口的用量到 usage_stat(stat_level=1)。
+ * 取上一小时(而非当前小时)确保窗口数据已全部落库完整。
+ * 对应 app.yml 的 solon.scheduling.job.usageStatJob。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class UsageStatTask {
+
+ @Inject
+ private UsageStatService usageStatService;
+
+ /**
+ * 由 app.yml usageStatJob 驱动(默认每小时 03 分,避开整点边界)
+ */
+ @Scheduled(name = "usageStatJob")
+ public void run() {
+ OffsetDateTime now = OffsetDateTime.now();
+ // 当前整点 = 窗口终点,上一整点 = 窗口起点
+ OffsetDateTime end = now.truncatedTo(ChronoUnit.HOURS);
+ OffsetDateTime start = end.minusHours(1);
+ try {
+ usageStatService.aggregateHour(start, end);
+ } catch (Throwable e) {
+ log.error("用量小时聚合异常 window=[{}, {})", start, end, e);
+ }
+ }
+}
diff --git a/src/main/java/com/imutil/tencent/TencentCallbackSign.java b/src/main/java/com/imutil/tencent/TencentCallbackSign.java
new file mode 100644
index 0000000..c415ec5
--- /dev/null
+++ b/src/main/java/com/imutil/tencent/TencentCallbackSign.java
@@ -0,0 +1,93 @@
+package com.imutil.tencent;
+
+import java.nio.charset.StandardCharsets;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
+
+/**
+ * 腾讯 IM 第三方回调签名校验
+ *
+ * 算法(参考腾讯文档 269/1522):
+ * Sign = sha256(Token + RequestTime)
+ * - Token:控制台回调 URL 配置的鉴权 Token(非 SecretKey)
+ * - RequestTime:回调请求 URL 参数携带的时间戳(秒)
+ * - RequestTime 与当前时间相差超过 1 分钟视为无效(防重放)
+ *
+ * @author imutil
+ */
+public final class TencentCallbackSign {
+
+ /** 签名允许的最大时间偏差(秒),文档建议 60 秒 */
+ public static final long MAX_TIME_DRIFT_SEC = 60L;
+
+ private TencentCallbackSign() {
+ }
+
+ /**
+ * 校验回调签名
+ *
+ * @param token 控制台配置的鉴权 Token
+ * @param requestTime 请求时间戳(秒,字符串形式)
+ * @param sign URL 中的 Sign 参数
+ * @return true=校验通过
+ */
+ public static boolean verify(String token, String requestTime, String sign) {
+ if (token == null || token.isEmpty() || requestTime == null || requestTime.isEmpty() || sign == null) {
+ return false;
+ }
+ long ts;
+ try {
+ ts = Long.parseLong(requestTime);
+ } catch (NumberFormatException e) {
+ return false;
+ }
+ // 时效校验(防重放)
+ long now = System.currentTimeMillis() / 1000L;
+ if (Math.abs(now - ts) > MAX_TIME_DRIFT_SEC) {
+ return false;
+ }
+ String expected = sha256Hex(token + requestTime);
+ return constantTimeEquals(expected, sign == null ? "" : sign.toLowerCase());
+ }
+
+ /**
+ * 计算 SHA-256 十六进制摘要(小写)
+ */
+ public static String sha256Hex(String input) {
+ try {
+ MessageDigest md = MessageDigest.getInstance("SHA-256");
+ byte[] digest = md.digest(input.getBytes(StandardCharsets.UTF_8));
+ return toHexLower(digest);
+ } catch (NoSuchAlgorithmException e) {
+ throw new IllegalStateException("SHA-256 不可用", e);
+ }
+ }
+
+ private static String toHexLower(byte[] bytes) {
+ char[] hex = new char[bytes.length * 2];
+ String digits = "0123456789abcdef";
+ for (int i = 0; i < bytes.length; i++) {
+ int v = bytes[i] & 0xFF;
+ hex[i * 2] = digits.charAt(v >>> 4);
+ hex[i * 2 + 1] = digits.charAt(v & 0x0F);
+ }
+ return new String(hex);
+ }
+
+ /**
+ * 常量时间比较,防时序攻击
+ */
+ private static boolean constantTimeEquals(String a, String b) {
+ if (a == null || b == null) {
+ return false;
+ }
+ if (a.length() != b.length()) {
+ return false;
+ }
+ int r = 0;
+ for (int i = 0; i < a.length(); i++) {
+ r |= a.charAt(i) ^ b.charAt(i);
+ }
+ return r == 0;
+ }
+}
diff --git a/src/main/java/com/imutil/tencent/TencentImClient.java b/src/main/java/com/imutil/tencent/TencentImClient.java
new file mode 100644
index 0000000..01af1f6
--- /dev/null
+++ b/src/main/java/com/imutil/tencent/TencentImClient.java
@@ -0,0 +1,198 @@
+package com.imutil.tencent;
+
+import com.imutil.common.Httpx;
+import com.imutil.common.Jsons;
+import com.imutil.common.TenantContext;
+import com.imutil.entity.ApiCallLog;
+import com.imutil.mapper.ApiCallLogMapper;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.snack4.ONode;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Inject;
+
+import java.net.URLEncoder;
+import java.nio.charset.StandardCharsets;
+import java.time.OffsetDateTime;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 腾讯 IM REST API 客户端(管理 API 收口)
+ *
+ * 鉴权方式:请求 URL 带 admin 的 UserSig + identifier + sdkappid query 参数(非 TC3 签名)。
+ * 所有调腾讯后台 API 的入口集中在此,密钥仅本类持有,业务系统不得直接调用。
+ *
+ * 每次调用同步写 {@code api_call_log} 审计(try-catch,写入失败不影响主流程),
+ * 作为用量统计(api_call_count)与问题追溯的数据源。
+ *
+ * 返回体统一含 ActionStatus(OK/FAIL)、ErrorCode、ErrorInfo,本类透传状态。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class TencentImClient {
+
+ @Inject("${imutil.tencent.sdkAppId:0}")
+ private long sdkAppId;
+
+ @Inject("${imutil.tencent.secretKey:}")
+ private String secretKey;
+
+ @Inject("${imutil.tencent.adminUserId:administrator}")
+ private String adminUserId;
+
+ @Inject("${imutil.tencent.apiHost:console.tim.qq.com}")
+ private String apiHost;
+
+ @Inject
+ private ApiCallLogMapper apiCallLogMapper;
+
+ /**
+ * 生成管理员 UserSig(长效,用于调后台 API)
+ */
+ private String genAdminSig() {
+ return UserSigUtil.genSig(sdkAppId, secretKey, adminUserId, 30L * 86400);
+ }
+
+ /**
+ * 调用 IM REST API
+ *
+ * @param command 命令路径,如 im_open_login_svc/account_import
+ * @param bodyJson 请求体 JSON
+ * @return 响应 JSON 字符串
+ */
+ public String callApi(String command, String bodyJson) {
+ // 租户来源:当前请求上下文,无则记 system(admin 后台调用等无租户上下文场景)
+ String tenantId = TenantContext.get();
+ String tid = (tenantId == null || tenantId.isEmpty()) ? "system" : tenantId;
+
+ String result;
+ try {
+ String adminSig = genAdminSig();
+ String url = "https://" + apiHost + "/v4/" + command
+ + "?sdkappid=" + sdkAppId
+ + "&identifier=" + URLEncoder.encode(adminUserId, StandardCharsets.UTF_8)
+ + "&usersig=" + URLEncoder.encode(adminSig, StandardCharsets.UTF_8)
+ + "&contenttype=json&platform=10&apn=1";
+ Httpx.Response resp = Httpx.postJsonDetail(url, bodyJson);
+ log.info("调腾讯API {} code={} body={}", command, resp.statusCode(), resp.body());
+ result = resp.body();
+ } catch (Exception e) {
+ log.error("调腾讯API异常 {} : {}", command, e.toString());
+ result = "{\"ActionStatus\":\"FAIL\",\"ErrorCode\":-1,\"ErrorInfo\":\""
+ + e.getClass().getSimpleName() + "\"}";
+ }
+ // 审计落库(失败不影响主流程)
+ recordCallLog(tid, command, bodyJson, result);
+ return result;
+ }
+
+ /**
+ * 写入 API 调用审计
+ *
+ * 同步写 + try-catch:审计写入异常不阻塞业务调用。
+ * params/result 完整存(审计完整性优先,text 字段无长度限制)。
+ */
+ private void recordCallLog(String tenantId, String apiName, String params, String result) {
+ try {
+ ApiCallLog entry = new ApiCallLog();
+ entry.setTenantId(tenantId);
+ entry.setApiName(apiName);
+ entry.setParams(params);
+ entry.setResult(result);
+ entry.setCaller("system");
+ entry.setCalledAt(OffsetDateTime.now());
+ apiCallLogMapper.insert(entry);
+ } catch (Exception e) {
+ log.warn("写入 api_call_log 失败 api={} : {}", apiName, e.getMessage());
+ }
+ }
+
+ /**
+ * 导入账号(创建 IM 用户,幂等:已存在亦返回 OK)
+ *
+ * @return ActionStatus 是否 OK
+ */
+ public boolean accountImport(String imUserId, String nick, String faceUrl) {
+ Map
+ * 用于验证 admin UserSig 有效性:返回 ActionStatus=OK 即签名校验通过(与账号是否存在无关)。
+ *
+ * @return true=API 调用成功(签名有效)
+ */
+ public boolean accountCheck(String imUserId) {
+ Map
+ * 命令字 openim_admin/get_roam_msg,按时间窗返回指定两人之间最近的 MaxCnt 条消息。
+ * 补拉服务据此增量补全回调丢失的单聊消息。
+ *
+ * @param fromAccount 发送方 IM 账号
+ * @param toAccount 接收方 IM 账号
+ * @param maxCnt 单次拉取条数上限
+ * @param minTime 拉取时间窗起点(秒级 epoch)
+ * @param maxInterval 拉取时间窗跨度(秒),从 minTime 起算
+ * @return 腾讯响应原始 JSON(含 MsgList/Complete/LastMsgTime/LastMsgSeq),由调用方解析
+ */
+ public String getRoamMsg(String fromAccount, String toAccount, int maxCnt, long minTime, long maxInterval) {
+ Map
+ * 命令字 group_open_http_svc/group_msg_get_simple,不传 ReqMsgSeq 时返回最新 ReqMsgNumber 条。
+ * 补拉服务据此增量补全回调丢失的群消息。
+ *
+ * @param groupId 群 ID
+ * @param reqMsgNumber 单次拉取条数上限
+ * @return 腾讯响应原始 JSON(含 RspMsgList/IsFinished),由调用方解析
+ */
+ public String getGroupMsg(String groupId, int reqMsgNumber) {
+ Map
+ * 使用官方 {@code com.github.tencentyun:tls-sig-api-v2} 实现,确保 HMAC-SHA256 + zlib raw deflate
+ * 算法与腾讯服务端校验完全一致,避免手写算法细节偏差导致的 70003(UserSig illegal)错误。
+ *
+ * 官方算法流水线:SigDict(固定字段顺序)→ HMAC-SHA256(key, json) → zlib raw deflate → base64。
+ * 密钥、字段顺序、压缩参数等易错点全部由官方 SDK 处理,本类仅做静态薄封装。
+ *
+ * @author imutil
+ */
+public final class UserSigUtil {
+
+ private UserSigUtil() {
+ }
+
+ /**
+ * 生成 UserSig
+ *
+ * 线程安全说明:每次调用 new 一个 {@link TLSSigAPIv2}(官方实例非线程安全),
+ * 构造开销极低,无需缓存。
+ *
+ * @param sdkAppId 应用 SDKAppID
+ * @param secretKey 应用 SecretKey(控制台「基本配置」中的密钥,必须与 sdkAppId 同一应用)
+ * @param identifier 用户标识(IM UserID,本项目为带租户前缀的隔离 ID)
+ * @param expireSec 有效期(秒)
+ * @return UserSig 字符串
+ */
+ public static String genSig(long sdkAppId, String secretKey, String identifier, long expireSec) {
+ TLSSigAPIv2 api = new TLSSigAPIv2(sdkAppId, secretKey);
+ return api.genUserSig(identifier, expireSec);
+ }
+}
diff --git a/src/main/java/com/imutil/worker/DispatchWorker.java b/src/main/java/com/imutil/worker/DispatchWorker.java
new file mode 100644
index 0000000..1b971e4
--- /dev/null
+++ b/src/main/java/com/imutil/worker/DispatchWorker.java
@@ -0,0 +1,92 @@
+package com.imutil.worker;
+
+import com.imutil.entity.DistQueue;
+import com.imutil.service.DispatchService;
+import lombok.extern.slf4j.Slf4j;
+import org.noear.solon.annotation.Component;
+import org.noear.solon.annotation.Init;
+import org.noear.solon.annotation.Inject;
+
+import java.util.List;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * 分发工作线程
+ *
+ * 应用启动后({@code @Init})创建 daemon 线程池,每个线程循环:
+ * 抢占一批 pending → 逐条 HTTP 转发业务系统 → 回执更新状态。
+ * 空闲时按 pollIntervalMs 休眠,避免空转。
+ *
+ * daemon 线程不阻止 JVM 退出;业务系统宕机时由回执失败+重试+死信兜底(见 DispatchService)。
+ *
+ * @author imutil
+ */
+@Slf4j
+@Component
+public class DispatchWorker {
+
+ @Inject
+ private DispatchService dispatchService;
+
+ @Inject("${imutil.dispatch.workerCount:4}")
+ private int workerCount;
+
+ @Inject("${imutil.dispatch.pollIntervalMs:1000}")
+ private long pollIntervalMs;
+
+ private final AtomicInteger seq = new AtomicInteger(0);
+
+ private ExecutorService pool;
+
+ /**
+ * 应用启动后启动工作线程
+ */
+ @Init
+ public void start() {
+ pool = Executors.newFixedThreadPool(workerCount, r -> {
+ Thread t = new Thread(r, "dispatch-worker-" + seq.incrementAndGet());
+ t.setDaemon(true);
+ return t;
+ });
+ for (int i = 0; i < workerCount; i++) {
+ pool.submit(this::loop);
+ }
+ log.info("分发工作线程已启动 count={}", workerCount);
+ }
+
+ /**
+ * 工作循环
+ */
+ private void loop() {
+ String workerName = "w-" + seq.get();
+ while (true) {
+ try {
+ List 默认账号 admin / admin123(首次启动初始化,请尽快修改密码)。死信 > 0 时请到「队列监控」处理。
+
+
+ <#else>
+
+
+
+ <#list grants as g>
+ 授权ID 源租户 源用户 目标租户 目标用户 权限 方向 状态 操作
+
+ #list>
+
+ ${g.grantId?c}
+ ${g.fromTenantId!}
+ ${g.fromImUserId!'(任意)'}
+ ${g.toTenantId!}
+ ${g.toImUserId!'(全员)'}
+ ${g.permissions!}
+ <#if g.direction?? && g.direction == 1>双向<#else>单向#if>
+
+ <#if g.status?? && g.status == 1>生效
+ <#else>撤销#if>
+
+
+ <#if g.status?? && g.status == 1>
+
+ #if>
+
+ 待分发 (pending,最近 50 条)
+ <#if pendings?has_content>
+
+
+ <#else>
+
+
+ <#list pendings as q>
+ ID 租户 会话 重试次数 下次重试时间
+
+ #list>
+
+ ${q.id?c}
+ ${q.tenantId!}
+ ${q.convId!}
+ ${q.retryCount!0}
+ <#if q.nextRetryAt??>${q.nextRetryAt?string('yyyy-MM-dd HH:mm:ss')}<#else>-#if>
+ 死信 (dead,共 ${deadCount!'0'} 条,最近 50 条)
+ <#if deads?has_content>
+
+
+ <#else>
+
+
+ <#list deads as q>
+ ID 租户 会话 重试次数 目标URL 操作
+
+ #list>
+
+ ${q.id?c}
+ ${q.tenantId!}
+ ${q.convId!}
+ ${q.retryCount!0}
+ ${q.targetUrl!}
+
+
+
+
+
+
+ <#else>
+
+
+
+
+ <#list tenants as t>
+ 租户ID 名称 前缀 AppKey AppSecret
+ 回调URL IM QPS 状态 操作
+
+
+ #list>
+
+ ${t.tenantId!}
+ ${t.tenantName!}
+ ${t.prefixCode!}
+ ${t.appKey!}
+ ${t.appSecret!}
+ ${t.callbackUrl!'-'}
+ ${t.quotaImQps!'-'}
+
+ <#if t.status?? && t.status == 1>
+ 启用
+ <#else>
+ 停用
+ #if>
+
+
+
+
+
+
+
+ <#else>
+
+
+
+ <#list stats as s>
+ 租户 统计时刻 粒度 IM消息数 IM DAU API调用数 TRTC时长(秒)
+
+ #list>
+
+ ${s.tenantId!}
+ <#if s.statTime??>${s.statTime?string('yyyy-MM-dd HH:mm')}<#else>-#if>
+ <#if s.statLevel?? && s.statLevel == 1>小时<#else>天#if>
+ ${s.imMsgCount!0}
+ ${s.imDau!0}
+ ${s.apiCallCount!0}
+ ${s.trtcDurationSec!0}
+