Kaynağa Gözat

Merge remote-tracking branch 'origin/master'

YPZ 1 hafta önce
ebeveyn
işleme
ff30ce7f07
26 değiştirilmiş dosya ile 863 ekleme ve 77 silme
  1. 241 67
      java/storlead-knowledge/storlead-knowledge-core/src/main/java/com/storlead/knowledge/utils/HttpService.java
  2. 1 1
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/EmailFolderRuleApiController.java
  3. 1 1
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MaiAttachmentApiController.java
  4. 1 1
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MailBlacklistRecordApiController.java
  5. 1 1
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MailTemplatesApiController.java
  6. 1 1
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MailboxAutoReplySetController.java
  7. 2 2
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/SmtpPopSettingsApiController.java
  8. 1 1
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/UserEmailFolderApiController.java
  9. 4 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/package-info.java
  10. 2 2
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/mkt/MailSignatureApiController.java
  11. 102 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/mkt/SmtpPopSettingsApiController.java
  12. 4 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/mkt/package-info.java
  13. 9 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/client/MktMailRemoteClient.java
  14. 6 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/constant/MailRemoteApiConstants.java
  15. 38 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/dto/MktMailStatusChangeRemoteDTO.java
  16. 26 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/dto/MktMailStatusSyncQueryRemoteDTO.java
  17. 23 0
      java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/support/MailRemoteHttpSupport.java
  18. 20 0
      java/storlead-mail/storlead-mail-core/src/main/java/com/storlead/mail/pojo/entity/SmtpPopSettingsEntity.java
  19. 41 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/entity/MailStatusSyncCursorEntity.java
  20. 9 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/mapper/MailStatusSyncCursorMapper.java
  21. 19 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/MailStatusSyncCursorService.java
  22. 12 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/MarketEmailsStatusSyncService.java
  23. 48 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/impl/MailStatusSyncCursorServiceImpl.java
  24. 201 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/impl/MarketEmailsStatusSyncServiceImpl.java
  25. 33 0
      java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/task/MarketEmailsStatusSyncTask.java
  26. 17 0
      java/storlead-sasa/storlead-trade/src/main/resources/sql/mail_status_sync_cursor.sql

+ 241 - 67
java/storlead-knowledge/storlead-knowledge-core/src/main/java/com/storlead/knowledge/utils/HttpService.java

@@ -144,73 +144,73 @@ public class HttpService {
         postStream(url, authorization, body, emitter, null);
     }
 
-    public void postStream(
-            String url,
-            String authorization,
-            String body,
-            SseEmitter emitter,
-            java.util.function.Consumer<String> onMessageEnd
-    ) {
-        try {
-            HttpRequest request = HttpRequest.newBuilder()
-                    .uri(URI.create(url))
-                    .timeout(Duration.ofSeconds(200))
-                    .header("Authorization", authorization)
-                    .header("Accept", "text/event-stream")
-                    .header("Content-Type", "application/json")
-                    .POST(HttpRequest.BodyPublishers.ofString(body))
-                    .build();
-
-            HttpResponse<InputStream> response =
-                    httpClient.send(request, HttpResponse.BodyHandlers.ofInputStream());
-            StringBuilder fullAnswerBuilder = new StringBuilder();
-            try (BufferedReader reader = new BufferedReader(
-                    new InputStreamReader(response.body(), StandardCharsets.UTF_8))) {
-
-                String line;
-
-                while ((line = reader.readLine()) != null) {
-
-                    if (line.startsWith("data:")) {
-                        String data = line.substring(5);
-
-                        if (!data.contains("message_end")) {
-                            try {
-                                JsonNode node = JacksonHolder.OBJECT_MAPPER.readTree(data);
-                                JsonNode answerNode = node.get("answer");
-                                if (answerNode != null && !answerNode.isNull()) {
-                                    fullAnswerBuilder.append(answerNode.asText());
-                                }
-                            } catch (Exception ignore) {
-                                // 非 JSON 或格式异常,跳过
-                            }
-                            //System.out.println(data);
-                            emitter.send(data);
-                        }
-
-                        if (data.contains("message_end")) {
-                            // 触发消息结束回调,传递完整数据
-                            ObjectNode endNode = (ObjectNode) JacksonHolder.OBJECT_MAPPER.readTree(data);
-                            endNode.put("answer", fullAnswerBuilder.toString());
-
-                            String finalData = JacksonHolder.OBJECT_MAPPER.writeValueAsString(endNode);
-                            if (onMessageEnd != null) {
-                                onMessageEnd.accept(finalData.replace("/",""));
-                            }
-                            emitter.send(finalData);
-                            emitter.complete();
-                            return;
-                        }
-                    }
-                }
-            }
-
-            emitter.complete();
-
-        } catch (Exception e) {
-            emitter.completeWithError(e);
-        }
-    }
+//    public void postStream(
+//            String url,
+//            String authorization,
+//            String body,
+//            SseEmitter emitter,
+//            java.util.function.Consumer<String> onMessageEnd
+//    ) {
+//        try {
+//            HttpRequest request = HttpRequest.newBuilder()
+//                    .uri(URI.create(url))
+//                    .timeout(Duration.ofSeconds(200))
+//                    .header("Authorization", authorization)
+//                    .header("Accept", "text/event-stream")
+//                    .header("Content-Type", "application/json")
+//                    .POST(HttpRequest.BodyPublishers.ofString(body))
+//                    .build();
+//
+//            HttpResponse<InputStream> response =
+//                    httpClient.send(request, HttpResponse.BodyHandlers.ofInputStream());
+//            StringBuilder fullAnswerBuilder = new StringBuilder();
+//            try (BufferedReader reader = new BufferedReader(
+//                    new InputStreamReader(response.body(), StandardCharsets.UTF_8))) {
+//
+//                String line;
+//
+//                while ((line = reader.readLine()) != null) {
+//
+//                    if (line.startsWith("data:")) {
+//                        String data = line.substring(5);
+//
+//                        if (!data.contains("message_end")) {
+//                            try {
+//                                JsonNode node = JacksonHolder.OBJECT_MAPPER.readTree(data);
+//                                JsonNode answerNode = node.get("answer");
+//                                if (answerNode != null && !answerNode.isNull()) {
+//                                    fullAnswerBuilder.append(answerNode.asText());
+//                                }
+//                            } catch (Exception ignore) {
+//                                // 非 JSON 或格式异常,跳过
+//                            }
+//                            //System.out.println(data);
+//                            emitter.send(data);
+//                        }
+//
+//                        if (data.contains("message_end")) {
+//                            // 触发消息结束回调,传递完整数据
+//                            ObjectNode endNode = (ObjectNode) JacksonHolder.OBJECT_MAPPER.readTree(data);
+//                            endNode.put("answer", fullAnswerBuilder.toString());
+//
+//                            String finalData = JacksonHolder.OBJECT_MAPPER.writeValueAsString(endNode);
+//                            if (onMessageEnd != null) {
+//                                onMessageEnd.accept(finalData.replace("/",""));
+//                            }
+//                            emitter.send(finalData);
+//                            emitter.complete();
+//                            return;
+//                        }
+//                    }
+//                }
+//            }
+//
+//            emitter.complete();
+//
+//        } catch (Exception e) {
+//            emitter.completeWithError(e);
+//        }
+//    }
 
     /* ================= PUT ================= */
 
@@ -384,4 +384,178 @@ public class HttpService {
 
         return out.toByteArray();
     }
+
+
+    /**
+     * 流式请求 Dify 工作流,并实时过滤 <think> 标签内容
+     *
+     * @param url            Dify 工作流 API 地址
+     * @param authorization  认证 Token(可为 null)
+     * @param body           请求体(JSON)
+     * @param emitter        SseEmitter 用于推送实时过滤后的文本
+     * @param onMessageEnd   回调,传递最终过滤后的完整文本
+     */
+    public void postStream(
+            String url,
+            String authorization,
+            String body,
+            SseEmitter emitter,
+            Consumer<String> onMessageEnd
+    ) {
+        try {
+            // 构建 HTTP 请求
+            HttpRequest.Builder builder = HttpRequest.newBuilder()
+                    .uri(URI.create(url))
+                    .timeout(Duration.ofSeconds(300))
+                    .header("Accept", "text/event-stream")
+                    .header("Content-Type", "application/json")
+                    .POST(HttpRequest.BodyPublishers.ofString(body));
+            if (authorization != null && !authorization.isEmpty()) {
+                builder.header("Authorization", authorization);
+            }
+            HttpRequest request = builder.build();
+
+            // 发送请求并获取流式响应
+            HttpResponse<InputStream> response =
+                    httpClient.send(request, HttpResponse.BodyHandlers.ofInputStream());
+
+            // 创建 Think 过滤器
+            ThinkFilter filter = new ThinkFilter();
+
+            try (BufferedReader reader = new BufferedReader(
+                    new InputStreamReader(response.body(), StandardCharsets.UTF_8))) {
+                String line;
+                while ((line = reader.readLine()) != null) {
+                    if (line.isBlank()) continue;
+
+                    // 解析 SSE 数据行(格式:data: {...})
+                    if (line.startsWith("data:")) {
+                        String jsonStr = line.substring(5).trim();
+                        if (jsonStr.isEmpty()) continue;
+
+                        try {
+                            JsonNode root = objectMapper.readTree(jsonStr);
+                            String event = root.path("event").asText();
+                            JsonNode dataNode = root.path("data");
+
+                            // 处理文本片段
+                            if ("text_chunk".equals(event)) {
+                                String chunk = dataNode.path("text").asText();
+                                if (chunk != null && !chunk.isEmpty()) {
+                                    // 实时过滤 chunk,返回新增的干净文本
+                                    String cleanPart = filter.processChunk(chunk);
+                                    if (!cleanPart.isEmpty()) {
+                                        emitter.send(cleanPart); // 推送实时增量
+                                    }
+                                }
+                            }
+
+                            // 处理工作流结束
+                            if ("workflow_finished".equals(event)) {
+                                // 获取最终完整过滤结果(处理缓冲区残留)
+                                String finalClean = filter.getFinalText();
+                                // 触发回调
+                                if (onMessageEnd != null) {
+                                    onMessageEnd.accept(finalClean);
+                                }
+                                // 推送最终完整结果(可选)
+                                //emitter.send(finalClean);
+                                emitter.complete();
+                                return;
+                            }
+                        } catch (Exception ignore) {
+                            // 单条 JSON 解析失败,跳过
+                        }
+                    }
+                }
+                // 如果流意外结束,兜底完成
+                emitter.complete();
+            }
+
+        } catch (Exception e) {
+            emitter.completeWithError(e);
+        }
+    }
+
+    /**
+     * 内部类:实时过滤 <think>...</think> 标签,支持跨 Chunk 处理
+     */
+    private static class ThinkFilter {
+        private final StringBuilder buffer = new StringBuilder();   // 用于累积未闭合的标签片段
+        private boolean inThink = false;                           // 是否当前处于 think 块内
+        private final StringBuilder result = new StringBuilder();   // 累积已输出的干净文本
+        private int lastSentLength = 0;                            // 已推送长度(用于增量)
+
+        /**
+         * 处理一个新的文本块,返回新增的过滤后文本(用于实时推送)
+         */
+        public String processChunk(String chunk) {
+            // 将新块追加到缓冲区,然后按字符处理
+            buffer.append(chunk);
+            StringBuilder output = new StringBuilder();
+            int i = 0;
+            while (i < buffer.length()) {
+                // 检查是否匹配 <think> 开始标签
+                if (!inThink && startsWithTag(buffer, i, "<think>")) {
+                    inThink = true;
+                    i += 6; // 跳过 "<think>"
+                    // 清除缓冲区中已消费的部分
+                    buffer.delete(0, i);
+                    i = 0;
+                    continue;
+                }
+                // 检查是否匹配 </think> 结束标签
+                if (inThink && startsWithTag(buffer, i, "</think>")) {
+                    inThink = false;
+                    i += 8; // 跳过 "</think>"
+                    buffer.delete(0, i);
+                    i = 0;
+                    continue;
+                }
+
+                // 如果不在 think 内,则输出字符
+                if (!inThink) {
+                    char c = buffer.charAt(i);
+                    output.append(c);
+                    result.append(c); // 同时保存到完整结果
+                }
+                i++;
+            }
+            // 清空缓冲区(所有字符已处理)
+            buffer.setLength(0);
+
+            // 获取新增部分(上次发送位置之后)
+            String newPart = output.toString();
+            return newPart;
+        }
+
+        /**
+         * 获取完整过滤后的文本(在流结束时调用,处理可能的剩余缓冲区)
+         */
+        public String getFinalText() {
+            // 如果缓冲区还有未处理的内容(通常是未闭合标签,直接丢弃)
+            // 但为了安全,也可尝试补全,此处简单丢弃
+            if (buffer.length() > 0) {
+                // 如果不在 think 内,且缓冲区有内容,则可能是标签的一部分,但已无法匹配,直接丢弃
+                // 实际很少出现,保留以防万一
+            }
+            return result.toString();
+        }
+
+        /**
+         * 辅助方法:检查 buffer 从 offset 开始是否匹配指定的 tag
+         */
+        private boolean startsWithTag(StringBuilder buffer, int offset, String tag) {
+            int tagLen = tag.length();
+            if (offset + tagLen > buffer.length()) {
+                return false;
+            }
+            for (int i = 0; i < tagLen; i++) {
+                if (buffer.charAt(offset + i) != tag.charAt(i)) {
+                    return false;
+                }
+            }
+            return true;
+        }
+    }
 }

+ 1 - 1
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/EmailFolderRuleApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/EmailFolderRuleApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.pojo.dto.EmailTemplatesDTO;

+ 1 - 1
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/MaiAttachmentApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MaiAttachmentApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.webclient.client.CrmMailFileRemoteClient;

+ 1 - 1
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/MailBlacklistRecordApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MailBlacklistRecordApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.pojo.dto.EmailTemplatesDTO;

+ 1 - 1
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/MailTemplatesApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MailTemplatesApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.pojo.dto.EmailTemplatesDTO;

+ 1 - 1
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/MailboxAutoReplySetController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/MailboxAutoReplySetController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.pojo.dto.EmailTemplatesDTO;

+ 2 - 2
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/SmtpPopSettingsApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/SmtpPopSettingsApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.dto.page.PageDTO;
 import com.storlead.framework.common.result.Result;
@@ -20,7 +20,7 @@ import javax.annotation.Resource;
 @Log4j2
 @RestController
 @RequestMapping("/smtp/pop/setting")
-@Api(tags = "邮件: 邮件管理")
+@Api(tags = "CRM邮件: 邮箱账户")
 public class SmtpPopSettingsApiController {
 
     @Resource

+ 1 - 1
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/UserEmailFolderApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/UserEmailFolderApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.crm;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.pojo.dto.EmailTemplatesDTO;

+ 4 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/crm/package-info.java

@@ -0,0 +1,4 @@
+/**
+ * CRM 邮件对外接口(路径不变,转发独立邮件服务)。
+ */
+package com.storlead.mail.controller.crm;

+ 2 - 2
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/MailSignatureApiController.java → java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/mkt/MailSignatureApiController.java

@@ -1,4 +1,4 @@
-package com.storlead.mail.controller;
+package com.storlead.mail.controller.mkt;
 
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.pojo.dto.EmailSignatureDTO;
@@ -15,7 +15,7 @@ import org.springframework.web.bind.annotation.RestController;
 import javax.annotation.Resource;
 
 @RestController
-@RequestMapping("/email/signature")
+@RequestMapping("/mkt/email/signature")
 @Api(tags = "邮件: 邮件签名")
 @Log4j2
 public class MailSignatureApiController {

+ 102 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/mkt/SmtpPopSettingsApiController.java

@@ -0,0 +1,102 @@
+package com.storlead.mail.controller.mkt;
+
+import com.storlead.framework.common.dto.page.PageDTO;
+import com.storlead.framework.common.result.Result;
+import com.storlead.mail.pojo.entity.SmtpPopSettingsEntity;
+import com.storlead.mail.webclient.client.MktAccountRemoteClient;
+import io.swagger.annotations.Api;
+import io.swagger.annotations.ApiOperation;
+import lombok.extern.log4j.Log4j2;
+import org.springframework.util.StringUtils;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestBody;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import javax.annotation.Resource;
+
+/**
+ * 营销邮箱账户(转发独立邮件服务 /mkt/account/**)。
+ */
+@Log4j2
+@RestController
+@RequestMapping("/mkt/email/account")
+@Api(tags = "营销邮件: 邮箱账户")
+public class SmtpPopSettingsApiController {
+
+    private static final String SCENE_MKT = "MKT";
+    private static final String SCOPE_ORG = "ORG";
+    private static final String SCOPE_PERSONAL = "PERSONAL";
+
+    @Resource
+    private MktAccountRemoteClient mktAccountRemoteClient;
+
+    @PostMapping("/page_list")
+    @ApiOperation("获取营销邮箱配置信息")
+    public Result<?> pageList(PageDTO dto) {
+        return mktAccountRemoteClient.pageList(dto);
+    }
+
+    @PostMapping("/list")
+    @ApiOperation("营销邮箱账户列表")
+    public Result<?> list(Long orgId) {
+        return mktAccountRemoteClient.list(orgId);
+    }
+
+    @PostMapping("/save")
+    @ApiOperation("保存/修改营销邮箱")
+    public Result<?> save(@RequestBody SmtpPopSettingsEntity entity) {
+        applyMktDefaults(entity);
+        return mktAccountRemoteClient.save(entity);
+    }
+
+    @PostMapping("/delete")
+    @ApiOperation("删除营销邮箱")
+    public Result<?> delete(Long id) {
+        return mktAccountRemoteClient.delete(id);
+    }
+
+    @PostMapping("/default_use")
+    @ApiOperation("设置默认营销发信账户")
+    public Result<?> defaultUse(Long id) {
+        return mktAccountRemoteClient.defaultUse(id);
+    }
+
+    @PostMapping("/test_connect")
+    @ApiOperation("测试营销邮箱连通性")
+    public Result<?> testConnect(@RequestBody SmtpPopSettingsEntity mailSet) {
+        applyMktDefaults(mailSet);
+        return mktAccountRemoteClient.testConnect(mailSet);
+    }
+
+    /**
+     * 营销账户默认必填:scene 固定 MKT,避免落入 CRM;其余空值补默认。
+     */
+    private void applyMktDefaults(SmtpPopSettingsEntity entity) {
+        if (entity == null) {
+            return;
+        }
+        entity.setScene(SCENE_MKT);
+        if (!StringUtils.hasText(entity.getAccountScope())) {
+            entity.setAccountScope(entity.getOrgId() != null ? SCOPE_ORG : SCOPE_PERSONAL);
+        }
+        if (entity.getEnableReceive() == null) {
+            entity.setEnableReceive(0);
+        }
+        if (entity.getUseDefault() == null) {
+            entity.setUseDefault(0);
+        }
+        if (entity.getIsAutoReply() == null) {
+            entity.setIsAutoReply(0);
+        }
+        if (entity.getIsDelete() == null) {
+            entity.setIsDelete(0);
+        }
+        if (entity.getEnabled() == null) {
+            entity.setEnabled(Boolean.TRUE);
+        }
+        if (entity.getIsDeleteMail() == null) {
+            entity.setIsDeleteMail(0);
+        }
+    }
+}

+ 4 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/controller/mkt/package-info.java

@@ -0,0 +1,4 @@
+/**
+ * 营销邮件对外接口(转发独立邮件服务 /mkt/**)。
+ */
+package com.storlead.mail.controller.mkt;

+ 9 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/client/MktMailRemoteClient.java

@@ -3,6 +3,8 @@ package com.storlead.mail.webclient.client;
 import com.storlead.framework.common.result.Result;
 import com.storlead.mail.webclient.constant.MailRemoteApiConstants;
 import com.storlead.mail.webclient.dto.MktMailQueryRemoteDTO;
+import com.storlead.mail.webclient.dto.MktMailStatusChangeRemoteDTO;
+import com.storlead.mail.webclient.dto.MktMailStatusSyncQueryRemoteDTO;
 import com.storlead.mail.webclient.dto.MktSendMailRemoteDTO;
 import com.storlead.mail.webclient.support.MailRemoteHttpSupport;
 import lombok.RequiredArgsConstructor;
@@ -64,6 +66,13 @@ public class MktMailRemoteClient {
         return http.postForm(MailRemoteApiConstants.MKT_OPEN_STATS, http.formOf("bizRef", bizRef));
     }
 
+    /**
+     * 按 status_updated_at 增量拉取营销邮件状态变更。
+     */
+    public Result<List<MktMailStatusChangeRemoteDTO>> statusChanges(MktMailStatusSyncQueryRemoteDTO query) {
+        return http.postJsonForList(MailRemoteApiConstants.MKT_STATUS_CHANGES, query, MktMailStatusChangeRemoteDTO.class);
+    }
+
     /** 打开追踪像素(通常由收件端浏览器直接访问 B 服务)。 */
     public byte[] track(String token) {
         return http.getForBytes(MailRemoteApiConstants.MKT_TRACK_PREFIX + token, null);

+ 6 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/constant/MailRemoteApiConstants.java

@@ -43,6 +43,12 @@ public final class MailRemoteApiConstants {
     /** 按业务单号汇总打开追踪(opened / open_count),参数:bizRef */
     public static final String MKT_OPEN_STATS = "/mkt/open_stats";
 
+    /**
+     * 营销邮件状态增量变更(供 trade 同步 market_emails)。
+     * JSON:MktMailStatusSyncQueryRemoteDTO;返回 List&lt;MktMailStatusChangeRemoteDTO&gt;
+     */
+    public static final String MKT_STATUS_CHANGES = "/mkt/status/changes";
+
     /**
      * 打开追踪像素(免登录,返回 1x1 GIF)。
      * 完整路径:{@value #MKT_TRACK_PREFIX}{token}

+ 38 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/dto/MktMailStatusChangeRemoteDTO.java

@@ -0,0 +1,38 @@
+package com.storlead.mail.webclient.dto;
+
+import com.fasterxml.jackson.annotation.JsonFormat;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.Data;
+import org.springframework.format.annotation.DateTimeFormat;
+
+import java.time.LocalDateTime;
+
+/**
+ * 营销邮件状态变更行(对齐 B 服务 mkt_mail_status 增量字段)。
+ */
+@Data
+public class MktMailStatusChangeRemoteDTO {
+
+    @ApiModelProperty("mkt_mail_status 主键")
+    private Long id;
+
+    @ApiModelProperty("邮件ID,对应 market_emails.email_id")
+    private Long mailId;
+
+    @ApiModelProperty("发送状态:-2草稿 -1定时 0发送中 1成功 2失败")
+    private Integer sendStatus;
+
+    @ApiModelProperty("是否已打开(像素):0否 1是")
+    private Integer opened;
+
+    @ApiModelProperty("是否已读:0否 1是")
+    private Integer isRead;
+
+    @ApiModelProperty("是否已回复:0否 1是")
+    private Integer replied;
+
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @ApiModelProperty("状态最后更新时间(增量游标)")
+    private LocalDateTime statusUpdatedAt;
+}

+ 26 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/dto/MktMailStatusSyncQueryRemoteDTO.java

@@ -0,0 +1,26 @@
+package com.storlead.mail.webclient.dto;
+
+import com.fasterxml.jackson.annotation.JsonFormat;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.Data;
+import org.springframework.format.annotation.DateTimeFormat;
+
+import java.time.LocalDateTime;
+
+/**
+ * 营销邮件状态增量查询(对齐 B 服务按 status_updated_at 游标拉取)。
+ */
+@Data
+public class MktMailStatusSyncQueryRemoteDTO {
+
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @ApiModelProperty("游标时间:status_updated_at > since;首次为空表示从头拉")
+    private LocalDateTime since;
+
+    @ApiModelProperty("同秒防漏:status_updated_at = since 时 id > lastId")
+    private Long lastId;
+
+    @ApiModelProperty("每批条数,默认 200")
+    private Integer limit = 200;
+}

+ 23 - 0
java/storlead-mail/storlead-mail-api/src/main/java/com/storlead/mail/webclient/support/MailRemoteHttpSupport.java

@@ -122,6 +122,29 @@ public class MailRemoteHttpSupport {
         return convertResult(raw, dataType);
     }
 
+    /**
+     * POST JSON,将 Result.result 转为 List&lt;T&gt;。
+     */
+    public <T> Result<List<T>> postJsonForList(String path, Object body, Class<T> elementType) {
+        Result<?> raw = postJson(path, body);
+        if (raw == null) {
+            return null;
+        }
+        Result<List<T>> typed = new Result<>();
+        typed.setSuccess(raw.isSuccess());
+        typed.setMessage(raw.getMessage());
+        typed.setCode(raw.getCode());
+        typed.setTimestamp(raw.getTimestamp());
+        Object data = raw.getResult();
+        if (data == null) {
+            typed.setResult(null);
+            return typed;
+        }
+        JavaType listType = objectMapper.getTypeFactory().constructCollectionType(List.class, elementType);
+        typed.setResult(objectMapper.convertValue(data, listType));
+        return typed;
+    }
+
     public Result<?> postForm(String path, Object formObject) {
         return postForm(path, toForm(formObject));
     }

+ 20 - 0
java/storlead-mail/storlead-mail-core/src/main/java/com/storlead/mail/pojo/entity/SmtpPopSettingsEntity.java

@@ -113,4 +113,24 @@ public class SmtpPopSettingsEntity extends SysBaseField {
     @ApiModelProperty(value = "错误次数")
     @TableField("connect_error_count")
     private Integer connectErrorCount;
+
+    @ApiModelProperty(value = "场景:CRM / MKT")
+    @TableField("scene")
+    private String scene;
+
+    @ApiModelProperty(value = "账户范围:PERSONAL / ORG")
+    @TableField("account_scope")
+    private String accountScope;
+
+    @ApiModelProperty(value = "组织ID(组织级营销账号)")
+    @TableField("org_id")
+    private Long orgId;
+
+    @ApiModelProperty(value = "日发送上限,空表示不限制")
+    @TableField("daily_send_limit")
+    private Integer dailySendLimit;
+
+    @ApiModelProperty(value = "是否允许收信:1是 0否")
+    @TableField("enable_receive")
+    private Integer enableReceive;
 }

+ 41 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/entity/MailStatusSyncCursorEntity.java

@@ -0,0 +1,41 @@
+package com.storlead.trade.entity;
+
+import com.baomidou.mybatisplus.annotation.IdType;
+import com.baomidou.mybatisplus.annotation.TableField;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import com.fasterxml.jackson.annotation.JsonFormat;
+import com.storlead.framework.mybatis.entity.SysBaseField;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import org.springframework.format.annotation.DateTimeFormat;
+
+import java.time.LocalDateTime;
+
+/**
+ * 邮件状态同步游标(trade 库)。
+ */
+@Data
+@EqualsAndHashCode(callSuper = true)
+@TableName("mail_status_sync_cursor")
+public class MailStatusSyncCursorEntity extends SysBaseField {
+
+    @TableId(type = IdType.AUTO)
+    @ApiModelProperty("主键")
+    private Long id;
+
+    @ApiModelProperty("同步源标识")
+    @TableField("source")
+    private String source;
+
+    @JsonFormat(timezone = "GMT+8", pattern = "yyyy-MM-dd HH:mm:ss")
+    @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss")
+    @ApiModelProperty("已同步到的 status_updated_at")
+    @TableField("last_status_updated_at")
+    private LocalDateTime lastStatusUpdatedAt;
+
+    @ApiModelProperty("已同步到的 mkt_mail_status.id")
+    @TableField("last_id")
+    private Long lastId;
+}

+ 9 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/mapper/MailStatusSyncCursorMapper.java

@@ -0,0 +1,9 @@
+package com.storlead.trade.mapper;
+
+import com.storlead.trade.entity.MailStatusSyncCursorEntity;
+import com.storlead.trade.mapper.support.TradeBaseMapper;
+import org.apache.ibatis.annotations.Mapper;
+
+@Mapper
+public interface MailStatusSyncCursorMapper extends TradeBaseMapper<MailStatusSyncCursorEntity> {
+}

+ 19 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/MailStatusSyncCursorService.java

@@ -0,0 +1,19 @@
+package com.storlead.trade.service;
+
+import com.storlead.framework.mybatis.service.MyBaseService;
+import com.storlead.trade.entity.MailStatusSyncCursorEntity;
+
+import java.time.LocalDateTime;
+
+public interface MailStatusSyncCursorService extends MyBaseService<MailStatusSyncCursorEntity> {
+
+    /**
+     * 获取或初始化指定 source 的游标。
+     */
+    MailStatusSyncCursorEntity getOrInit(String source);
+
+    /**
+     * 推进游标(本批成功写完 market_emails 后调用)。
+     */
+    void advance(String source, LocalDateTime lastStatusUpdatedAt, Long lastId);
+}

+ 12 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/MarketEmailsStatusSyncService.java

@@ -0,0 +1,12 @@
+package com.storlead.trade.service;
+
+/**
+ * 将邮件服务 mkt_mail_status 增量同步到 trade.market_emails。
+ */
+public interface MarketEmailsStatusSyncService {
+
+    /**
+     * 执行一轮(或多页)增量同步,返回本轮成功回写条数。
+     */
+    int syncOnce();
+}

+ 48 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/impl/MailStatusSyncCursorServiceImpl.java

@@ -0,0 +1,48 @@
+package com.storlead.trade.service.impl;
+
+import com.baomidou.dynamic.datasource.annotation.DS;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.storlead.framework.common.constant.CommonConstant;
+import com.storlead.framework.common.constant.DSConstants;
+import com.storlead.trade.entity.MailStatusSyncCursorEntity;
+import com.storlead.trade.mapper.MailStatusSyncCursorMapper;
+import com.storlead.trade.service.MailStatusSyncCursorService;
+import com.storlead.trade.service.impl.support.TradeDataSourceServiceImpl;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.time.LocalDateTime;
+
+@Service
+@DS(DSConstants.DATASOURCE_TRADE)
+public class MailStatusSyncCursorServiceImpl
+        extends TradeDataSourceServiceImpl<MailStatusSyncCursorMapper, MailStatusSyncCursorEntity>
+        implements MailStatusSyncCursorService {
+
+    @Override
+    public MailStatusSyncCursorEntity getOrInit(String source) {
+        MailStatusSyncCursorEntity cursor = this.getOne(new LambdaQueryWrapper<MailStatusSyncCursorEntity>()
+                .eq(MailStatusSyncCursorEntity::getSource, source)
+                .eq(MailStatusSyncCursorEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .last("limit 1"));
+        if (cursor != null) {
+            return cursor;
+        }
+        cursor = new MailStatusSyncCursorEntity();
+        cursor.setSource(source);
+        cursor.setLastId(0L);
+        cursor.setIsDelete(CommonConstant.DEL_FLAG_0);
+        cursor.setEnabled(Boolean.TRUE);
+        this.save(cursor);
+        return cursor;
+    }
+
+    @Override
+    @Transactional(rollbackFor = Exception.class)
+    public void advance(String source, LocalDateTime lastStatusUpdatedAt, Long lastId) {
+        MailStatusSyncCursorEntity cursor = getOrInit(source);
+        cursor.setLastStatusUpdatedAt(lastStatusUpdatedAt);
+        cursor.setLastId(lastId);
+        this.updateById(cursor);
+    }
+}

+ 201 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/service/impl/MarketEmailsStatusSyncServiceImpl.java

@@ -0,0 +1,201 @@
+package com.storlead.trade.service.impl;
+
+import com.baomidou.dynamic.datasource.annotation.DS;
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
+import com.storlead.framework.common.constant.CommonConstant;
+import com.storlead.framework.common.constant.DSConstants;
+import com.storlead.framework.common.result.Result;
+import com.storlead.mail.webclient.client.MktMailRemoteClient;
+import com.storlead.mail.webclient.dto.MktMailStatusChangeRemoteDTO;
+import com.storlead.mail.webclient.dto.MktMailStatusSyncQueryRemoteDTO;
+import com.storlead.trade.entity.MailStatusSyncCursorEntity;
+import com.storlead.trade.entity.MarketEmailsEntity;
+import com.storlead.trade.service.MailStatusSyncCursorService;
+import com.storlead.trade.service.MarketEmailsService;
+import com.storlead.trade.service.MarketEmailsStatusSyncService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+import org.springframework.util.CollectionUtils;
+
+import javax.annotation.Resource;
+import java.time.LocalDateTime;
+import java.util.List;
+
+/**
+ * 从邮件服务增量拉取 mkt_mail_status,按 email_id 回写 market_emails。
+ */
+@Slf4j
+@Service
+@DS(DSConstants.DATASOURCE_TRADE)
+public class MarketEmailsStatusSyncServiceImpl implements MarketEmailsStatusSyncService {
+
+    public static final String SOURCE_MKT_MAIL_STATUS = "mkt_mail_status";
+
+    private static final int BATCH_LIMIT = 200;
+    private static final int MAX_PAGES = 50;
+
+    /** trade.status:0草稿/待发,1已发送,2发送失败 */
+    private static final int TRADE_STATUS_DRAFT = 0;
+    private static final int TRADE_STATUS_SENT = 1;
+    private static final int TRADE_STATUS_FAIL = 2;
+
+    @Resource
+    private MktMailRemoteClient mktMailRemoteClient;
+
+    @Resource
+    private MailStatusSyncCursorService mailStatusSyncCursorService;
+
+    @Resource
+    private MarketEmailsService marketEmailsService;
+
+    @Override
+    public int syncOnce() {
+        MailStatusSyncCursorEntity cursor = mailStatusSyncCursorService.getOrInit(SOURCE_MKT_MAIL_STATUS);
+        LocalDateTime since = cursor.getLastStatusUpdatedAt();
+        Long lastId = cursor.getLastId() == null ? 0L : cursor.getLastId();
+
+        int totalUpdated = 0;
+        for (int page = 0; page < MAX_PAGES; page++) {
+            MktMailStatusSyncQueryRemoteDTO query = new MktMailStatusSyncQueryRemoteDTO();
+            query.setSince(since);
+            query.setLastId(lastId);
+            query.setLimit(BATCH_LIMIT);
+
+            Result<List<MktMailStatusChangeRemoteDTO>> result;
+            try {
+                result = mktMailRemoteClient.statusChanges(query);
+            } catch (Exception e) {
+                log.error("拉取营销邮件状态失败 since={} lastId={}", since, lastId, e);
+                break;
+            }
+            if (result == null || !result.isSuccess()) {
+                log.warn("拉取营销邮件状态返回失败 message={}", result == null ? null : result.getMessage());
+                break;
+            }
+            List<MktMailStatusChangeRemoteDTO> changes = result.getResult();
+            if (CollectionUtils.isEmpty(changes)) {
+                break;
+            }
+
+            LocalDateTime batchMaxTime = since;
+            Long batchMaxId = lastId;
+            int pageUpdated = 0;
+            for (MktMailStatusChangeRemoteDTO change : changes) {
+                if (change == null || change.getMailId() == null) {
+                    continue;
+                }
+                if (applyChange(change)) {
+                    pageUpdated++;
+                }
+                if (change.getStatusUpdatedAt() != null) {
+                    if (batchMaxTime == null || change.getStatusUpdatedAt().isAfter(batchMaxTime)
+                            || (change.getStatusUpdatedAt().equals(batchMaxTime)
+                            && change.getId() != null && change.getId() > batchMaxId)) {
+                        batchMaxTime = change.getStatusUpdatedAt();
+                        batchMaxId = change.getId() == null ? batchMaxId : change.getId();
+                    }
+                } else if (change.getId() != null && change.getId() > batchMaxId) {
+                    batchMaxId = change.getId();
+                }
+            }
+
+            // 本批处理完再推进游标(失败不推进则下次重拉,更新需幂等)
+            since = batchMaxTime;
+            lastId = batchMaxId;
+            mailStatusSyncCursorService.advance(SOURCE_MKT_MAIL_STATUS, since, lastId);
+            totalUpdated += pageUpdated;
+
+            if (changes.size() < BATCH_LIMIT) {
+                break;
+            }
+        }
+        return totalUpdated;
+    }
+
+    /**
+     * 按 email_id 回写;发送状态不降级;已读/已回复只升不降。
+     */
+    private boolean applyChange(MktMailStatusChangeRemoteDTO change) {
+        MarketEmailsEntity existing = marketEmailsService.getOne(new LambdaQueryWrapper<MarketEmailsEntity>()
+                .eq(MarketEmailsEntity::getEmailId, change.getMailId())
+                .eq(MarketEmailsEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .last("limit 1"));
+        if (existing == null) {
+            log.debug("market_emails 无 email_id={} 记录,跳过", change.getMailId());
+            return false;
+        }
+
+        Integer mappedStatus = mapSendStatus(change.getSendStatus());
+        Integer targetReader = resolveReader(change);
+        Integer targetReply = Integer.valueOf(1).equals(change.getReplied()) ? 1 : null;
+
+        LambdaUpdateWrapper<MarketEmailsEntity> update = new LambdaUpdateWrapper<>();
+        update.eq(MarketEmailsEntity::getId, existing.getId());
+        boolean changed = false;
+
+        if (mappedStatus != null && shouldUpgradeSendStatus(existing.getStatus(), mappedStatus)) {
+            update.set(MarketEmailsEntity::getStatus, mappedStatus);
+            changed = true;
+        }
+        if (Integer.valueOf(1).equals(targetReader)
+                && !Integer.valueOf(1).equals(existing.getHasReader())) {
+            update.set(MarketEmailsEntity::getHasReader, 1);
+            changed = true;
+        }
+        if (Integer.valueOf(1).equals(targetReply)
+                && !Integer.valueOf(1).equals(existing.getHasReply())) {
+            update.set(MarketEmailsEntity::getHasReply, 1);
+            changed = true;
+        }
+        if (!changed) {
+            return false;
+        }
+        return marketEmailsService.update(update);
+    }
+
+    private static Integer mapSendStatus(Integer sendStatus) {
+        if (sendStatus == null) {
+            return null;
+        }
+        if (sendStatus == 1) {
+            return TRADE_STATUS_SENT;
+        }
+        if (sendStatus == 2) {
+            return TRADE_STATUS_FAIL;
+        }
+        // -2草稿 -1定时 0发送中 → trade 0
+        return TRADE_STATUS_DRAFT;
+    }
+
+    /** 已打开或已读都视为 has_reader=1 */
+    private static Integer resolveReader(MktMailStatusChangeRemoteDTO change) {
+        if (Integer.valueOf(1).equals(change.getOpened()) || Integer.valueOf(1).equals(change.getIsRead())) {
+            return 1;
+        }
+        return null;
+    }
+
+    /**
+     * 发送状态:成功不降级;失败可升为成功;草稿可变为成功/失败。
+     */
+    private static boolean shouldUpgradeSendStatus(Integer current, Integer mapped) {
+        if (mapped == null) {
+            return false;
+        }
+        if (current == null) {
+            return true;
+        }
+        if (current.equals(mapped)) {
+            return false;
+        }
+        if (TRADE_STATUS_SENT == current) {
+            return false;
+        }
+        if (TRADE_STATUS_FAIL == current) {
+            return TRADE_STATUS_SENT == mapped;
+        }
+        // 当前为草稿/待发
+        return true;
+    }
+}

+ 33 - 0
java/storlead-sasa/storlead-trade/src/main/java/com/storlead/trade/task/MarketEmailsStatusSyncTask.java

@@ -0,0 +1,33 @@
+package com.storlead.trade.task;
+
+import com.storlead.trade.service.MarketEmailsStatusSyncService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import javax.annotation.Resource;
+
+/**
+ * 定时增量同步邮件库 mkt_mail_status → trade.market_emails。
+ */
+@Slf4j
+@Component
+public class MarketEmailsStatusSyncTask {
+
+    @Resource
+    private MarketEmailsStatusSyncService marketEmailsStatusSyncService;
+
+    /** 每 2 分钟跑一轮 */
+    // @Scheduled(cron = "0 */2 * * * ?")
+    public void sync() {
+        long start = System.currentTimeMillis();
+        try {
+            int updated = marketEmailsStatusSyncService.syncOnce();
+            if (updated > 0) {
+                log.info("营销邮件状态同步完成 updated={} cost={}ms", updated, System.currentTimeMillis() - start);
+            }
+        } catch (Exception e) {
+            log.error("营销邮件状态同步失败", e);
+        }
+    }
+}

+ 17 - 0
java/storlead-sasa/storlead-trade/src/main/resources/sql/mail_status_sync_cursor.sql

@@ -0,0 +1,17 @@
+-- trade 库:邮件状态同步游标(消费 mkt_mail_status 增量)
+CREATE TABLE IF NOT EXISTS `mail_status_sync_cursor` (
+  `id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键',
+  `source` varchar(64) NOT NULL COMMENT '同步源,如 mkt_mail_status',
+  `last_status_updated_at` datetime DEFAULT NULL COMMENT '已同步到的 status_updated_at',
+  `last_id` bigint DEFAULT NULL COMMENT '已同步到的 mkt_mail_status.id',
+  `owner_by` bigint DEFAULT NULL COMMENT '数据归属人',
+  `create_by` bigint DEFAULT NULL COMMENT '创建人',
+  `update_by` bigint DEFAULT NULL COMMENT '最后更新人',
+  `create_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
+  `update_time` datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
+  `is_delete` tinyint NOT NULL DEFAULT '0' COMMENT '逻辑删除:0未删除;1已删除',
+  `enabled` tinyint NOT NULL DEFAULT '1' COMMENT '是否启用:1启用;0禁用',
+  `sort` int DEFAULT '0' COMMENT '排序号',
+  PRIMARY KEY (`id`),
+  UNIQUE KEY `uk_source` (`source`)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin COMMENT='邮件状态同步游标';