Sfoglia il codice sorgente

邮件服务接口增加和调整,业务拆分

chenkq 2 settimane fa
parent
commit
50267da8e5
13 ha cambiato i file con 1085 aggiunte e 569 eliminazioni
  1. 1 1
      storlead-mail/README.md
  2. 215 0
      storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailMessageParser.java
  3. 17 0
      storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceiveBodyTask.java
  4. 171 0
      storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceiveEngine.java
  5. 31 0
      storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceiveHeadDraft.java
  6. 62 0
      storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceivePersistSpi.java
  7. 183 0
      storlead-mail/storlead-mail-crm/src/main/java/com/storlead/sales/mail/crm/receive/CrmMailReceivePersist.java
  8. 13 20
      storlead-mail/storlead-mail-crm/src/main/java/com/storlead/sales/mail/crm/service/impl/EmailsServiceImpl.java
  9. 8 65
      storlead-mail/storlead-mail-crm/src/main/java/com/storlead/sales/mail/crm/util/ReceiveMailQueueThreadPool.java
  10. 357 0
      storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/receive/MktMailReceivePersist.java
  11. 5 0
      storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/service/impl/MktMailMessageServiceImpl.java
  12. 7 483
      storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/service/impl/MktMailReceiveServiceImpl.java
  13. 15 0
      storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/support/MktMailHeaders.java

+ 1 - 1
storlead-mail/README.md

@@ -6,7 +6,7 @@
 
 | 模块 | Artifact | 职责 |
 |------|----------|------|
-| common | `storlead-mail-common` | 账户(`smtp_pop_settings`)、SMTP/IMAP 连接、公共投递、`MailSceneEnum` |
+| common | `storlead-mail-common` | 账户、SMTP/IMAP、公共投递、`MailReceiveEngine` 收信引擎、`MailSceneEnum` |
 | crm | `storlead-mail-crm` | CRM 个人邮箱:`emails` / `client_sent_emails`、收信、文件夹、CRM 集成任务 |
 | mkt | `storlead-mail-mkt` | 营销投递:`mkt_mail_message` / `mkt_mail_attachment`、草稿、附件、延时、打开追踪 |
 | dispatch | `storlead-mail-dispatch` | **统一调度**:定时发送(CRM+MKT)、CRM 收信拉取、服务器邮件延迟删除 |

+ 215 - 0
storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailMessageParser.java

@@ -0,0 +1,215 @@
+package com.storlead.sales.mail.common.receive;
+
+import com.storlead.framework.common.util.DateUtil;
+import com.sun.mail.imap.IMAPFolder;
+import com.sun.mail.pop3.POP3Folder;
+import lombok.extern.log4j.Log4j2;
+import org.springframework.util.StringUtils;
+
+import javax.mail.Address;
+import javax.mail.Flags;
+import javax.mail.Folder;
+import javax.mail.Message;
+import javax.mail.Multipart;
+import javax.mail.internet.InternetAddress;
+import javax.mail.internet.MimeUtility;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+/**
+ * 收信 Message 解析(与 CRM/MKT 表无关)。
+ */
+@Log4j2
+public final class MailMessageParser {
+
+    private MailMessageParser() {
+    }
+
+    public static MailReceiveHeadDraft parseHead(Message message, String folderCode, boolean isPullHead) {
+        try {
+            MailReceiveHeadDraft draft = new MailReceiveHeadDraft();
+            String subject = message.getSubject() == null ? "" : MimeUtility.decodeText(message.getSubject());
+            draft.setSubject(subject);
+            draft.setFolder(folderCode);
+            draft.setInOutMark("INBOX".equalsIgnoreCase(folderCode) ? 1 : 2);
+
+            Map<String, String> fromMap = decodeAddresses(message.getFrom());
+            draft.setFromAddr(joinKeys(fromMap));
+            draft.setFromName(joinValues(fromMap));
+
+            Map<String, String> toMap = decodeAddresses(message.getRecipients(Message.RecipientType.TO));
+            draft.setRecipient(joinKeys(toMap));
+            draft.setRecipientName(joinValues(toMap));
+
+            Map<String, String> ccMap = decodeAddresses(message.getRecipients(Message.RecipientType.CC));
+            draft.setRecipientCc(joinKeys(ccMap));
+            draft.setRecipientCcName(joinValues(ccMap));
+
+            draft.setEmailSize((long) Math.max(message.getSize(), 0));
+            if (message.getSentDate() != null) {
+                draft.setSentDate(message.getSentDate());
+            }
+            if (message.getReceivedDate() != null) {
+                draft.setRecipientDate(message.getReceivedDate());
+            } else if (message.getSentDate() != null) {
+                draft.setRecipientDate(message.getSentDate());
+            }
+            draft.setIsRead(message.isSet(Flags.Flag.SEEN) ? 1 : 0);
+
+            if (isPullHead) {
+                try {
+                    Object content = message.getContent();
+                    if (content instanceof String) {
+                        draft.setIsOnlyHead(0);
+                        draft.setContent(content.toString());
+                    } else {
+                        draft.setIsOnlyHead(1);
+                    }
+                } catch (Exception e) {
+                    draft.setIsOnlyHead(1);
+                }
+            }
+            draft.setInReplyTo(resolveInReplyTo(message));
+            return draft;
+        } catch (Exception e) {
+            log.error("parseHead error", e);
+            return null;
+        }
+    }
+
+    public static String extractBody(Message message) {
+        try {
+            Object content = message.getContent();
+            if (content instanceof String) {
+                return content.toString();
+            }
+            if (content instanceof Multipart) {
+                return parseMultipartText((Multipart) content);
+            }
+        } catch (Exception e) {
+            log.error("extractBody error", e);
+        }
+        return "";
+    }
+
+    public static String resolveMsgUid(Message message, String protocolType) {
+        try {
+            Folder folder = message.getFolder();
+            if ("IMAP".equalsIgnoreCase(protocolType) && folder instanceof IMAPFolder) {
+                return String.valueOf(((IMAPFolder) folder).getUID(message));
+            }
+            if ("POP3".equalsIgnoreCase(protocolType) && folder instanceof POP3Folder) {
+                return ((POP3Folder) folder).getUID(message);
+            }
+        } catch (Exception e) {
+            log.error("resolveMsgUid error", e);
+        }
+        return null;
+    }
+
+    public static String firstHeader(Message message, String name) {
+        try {
+            String[] values = message.getHeader(name);
+            if (values == null || values.length == 0) {
+                return null;
+            }
+            return values[0];
+        } catch (Exception e) {
+            return null;
+        }
+    }
+
+    public static String resolveInReplyTo(Message message) {
+        String inReplyTo = firstHeader(message, "In-Reply-To");
+        if (!StringUtils.hasText(inReplyTo)) {
+            String references = firstHeader(message, "References");
+            if (StringUtils.hasText(references)) {
+                String[] parts = references.trim().split("\\s+");
+                inReplyTo = parts[parts.length - 1];
+            }
+        }
+        return normalizeMessageId(inReplyTo);
+    }
+
+    public static String normalizeMessageId(String messageId) {
+        if (!StringUtils.hasText(messageId)) {
+            return messageId;
+        }
+        return messageId.trim();
+    }
+
+    public static String stripBrackets(String messageId) {
+        if (!StringUtils.hasText(messageId)) {
+            return messageId;
+        }
+        String v = messageId.trim();
+        if (v.startsWith("<") && v.endsWith(">") && v.length() > 2) {
+            return v.substring(1, v.length() - 1);
+        }
+        return v;
+    }
+
+    public static java.util.Date resolveBeginDate(java.util.Date maxRecipientDate, Integer pullDay) {
+        if (maxRecipientDate != null) {
+            return maxRecipientDate;
+        }
+        java.time.LocalDate currentDate = java.time.LocalDate.now();
+        if (pullDay == null || Integer.valueOf(0).equals(pullDay)) {
+            return DateUtil.getDateByFormat(currentDate.minusYears(5).toString(), "yyyy-MM-dd");
+        }
+        return DateUtil.getDateByFormat(currentDate.minusDays(pullDay).toString(), "yyyy-MM-dd");
+    }
+
+    private static String parseMultipartText(Multipart multipart) throws Exception {
+        StringBuilder html = new StringBuilder();
+        StringBuilder plain = new StringBuilder();
+        for (int i = 0; i < multipart.getCount(); i++) {
+            javax.mail.BodyPart part = multipart.getBodyPart(i);
+            Object content = part.getContent();
+            if (content instanceof Multipart) {
+                String nested = parseMultipartText((Multipart) content);
+                if (StringUtils.hasText(nested)) {
+                    html.append(nested);
+                }
+            } else if (part.isMimeType("text/html")) {
+                html.append(content == null ? "" : content.toString());
+            } else if (part.isMimeType("text/plain")) {
+                plain.append(content == null ? "" : content.toString());
+            }
+        }
+        return html.length() > 0 ? html.toString() : plain.toString();
+    }
+
+    private static Map<String, String> decodeAddresses(Address[] addresses) throws Exception {
+        Map<String, String> map = new HashMap<>();
+        if (addresses == null) {
+            return map;
+        }
+        for (Address address : addresses) {
+            InternetAddress internetAddress = (InternetAddress) address;
+            String personal = internetAddress.getPersonal();
+            if (personal != null) {
+                personal = MimeUtility.decodeText(personal);
+            }
+            map.put(internetAddress.getAddress(), personal != null ? personal : "N/A");
+        }
+        return map;
+    }
+
+    private static String joinKeys(Map<String, String> map) {
+        if (map == null || map.isEmpty()) {
+            return "";
+        }
+        return String.join(",", map.keySet());
+    }
+
+    private static String joinValues(Map<String, String> map) {
+        if (map == null || map.isEmpty()) {
+            return "";
+        }
+        return map.values().stream()
+                .filter(v -> v != null && !"N/A".equals(v))
+                .collect(Collectors.joining(","));
+    }
+}

+ 17 - 0
storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceiveBodyTask.java

@@ -0,0 +1,17 @@
+package com.storlead.sales.mail.common.receive;
+
+import lombok.Data;
+
+import java.util.Date;
+
+/**
+ * 待补正文的收信任务。
+ */
+@Data
+public class MailReceiveBodyTask {
+
+    private Long mailId;
+    private String msgUid;
+    private String folder;
+    private Date recipientDate;
+}

+ 171 - 0
storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceiveEngine.java

@@ -0,0 +1,171 @@
+package com.storlead.sales.mail.common.receive;
+
+import com.storlead.sales.mail.common.connection.MailConnection;
+import com.storlead.sales.mail.common.connection.client.MailConnectionUtil;
+import com.storlead.sales.mail.common.connection.config.MailProperties;
+import com.storlead.sales.mail.common.entity.SmtpPopSettingsEntity;
+import com.storlead.sales.mail.common.enums.MailSceneEnum;
+import com.storlead.sales.mail.common.util.MailReceiveQueueThreadPool;
+import lombok.extern.log4j.Log4j2;
+import org.springframework.stereotype.Component;
+import org.springframework.util.CollectionUtils;
+import org.springframework.util.StringUtils;
+
+import javax.mail.Folder;
+import javax.mail.Message;
+import javax.mail.search.ComparisonTerm;
+import javax.mail.search.ReceivedDateTerm;
+import javax.mail.search.SearchTerm;
+import java.util.Arrays;
+import java.util.Date;
+import java.util.List;
+
+/**
+ * 公共收信引擎:对齐 CRM receiveEmails 三阶段(拉头 → 正文 → 附件),
+ * 落库由 {@link MailReceivePersistSpi} 实现。
+ * <p>同账户串行:拉头完成后再补正文/附件,避免并发乱序。
+ */
+@Log4j2
+@Component
+public class MailReceiveEngine {
+
+    public static final String FOLDER_INBOX = "INBOX";
+    public static final String FOLDER_SENT = "SENT";
+    private static final String SERVER_INBOX = "INBOX";
+    private static final String SERVER_SENT = "Sent Items";
+
+    /**
+     * 异步收信入口(与 CRM receiveEmails 形态一致)。
+     */
+    public void receive(SmtpPopSettingsEntity account,
+                        MailSceneEnum scene,
+                        MailReceivePersistSpi persist,
+                        boolean isBackTask) {
+        if (account == null || account.getId() == null || persist == null || scene == null) {
+            return;
+        }
+        String taskKey = buildTaskKey(scene, account);
+        MailReceiveQueueThreadPool pool = MailReceiveQueueThreadPool.getInstance();
+        if (Boolean.TRUE.equals(pool.getTaskInProgressByMail(taskKey))) {
+            log.warn("receive already running: {}", taskKey);
+            return;
+        }
+        pool.addTask(taskKey, () -> {
+            try {
+                receiveHead(account, persist, isBackTask);
+                loadBodies(account, persist);
+                persist.loadPendingAttachmentFiles(account);
+            } catch (Exception e) {
+                log.error("MailReceiveEngine receive error, scene={}, accountId={}",
+                        scene.getCode(), account.getId(), e);
+            }
+        });
+    }
+
+    public static String buildTaskKey(MailSceneEnum scene, SmtpPopSettingsEntity account) {
+        return scene.getCode() + ":" + account.getId() + ":" + account.getEmailAddress();
+    }
+
+    private void receiveHead(SmtpPopSettingsEntity account, MailReceivePersistSpi persist, boolean isBackTask) {
+        pullOneFolder(account, persist, SERVER_INBOX, FOLDER_INBOX, isBackTask);
+        pullOneFolder(account, persist, SERVER_SENT, FOLDER_SENT, isBackTask);
+    }
+
+    private void pullOneFolder(SmtpPopSettingsEntity account,
+                               MailReceivePersistSpi persist,
+                               String serverFolder,
+                               String folderCode,
+                               boolean isBackTask) {
+        try {
+            List<Message> messages = searchMessages(account, persist, serverFolder, folderCode);
+            if (CollectionUtils.isEmpty(messages)) {
+                return;
+            }
+            log.info("==========邮件入库开始=========={} accountId={}", folderCode, account.getId());
+            List<Long> savedIds = new java.util.ArrayList<>();
+            for (Message message : messages) {
+                try {
+                    String[] headers = message.getHeader("Message-ID");
+                    if (headers == null || headers.length == 0 || !StringUtils.hasText(headers[0])) {
+                        continue;
+                    }
+                    // 是否新建/合并由 Persist 决定(MKT:同 message_id 的本地发信记录保留 id 并更新服务器字段)
+                    Long mailId = persist.saveHead(account, message, folderCode, isBackTask);
+                    if (mailId != null) {
+                        savedIds.add(mailId);
+                    }
+                } catch (Exception ex) {
+                    log.error("pull one mail error, accountId={}, folder={}", account.getId(), folderCode, ex);
+                }
+            }
+            log.info("==========邮件入库结束=========={} size={}", folderCode, savedIds.size());
+            if (!CollectionUtils.isEmpty(savedIds)) {
+                persist.afterHeadsSaved(account, folderCode, savedIds, isBackTask);
+            }
+        } catch (Exception e) {
+            log.error("pullOneFolder error, folder={}, accountId={}", folderCode, account.getId(), e);
+        }
+    }
+
+    private List<Message> searchMessages(SmtpPopSettingsEntity account,
+                                         MailReceivePersistSpi persist,
+                                         String serverFolder,
+                                         String folderCode) {
+        try {
+            Date beginDate = MailMessageParser.resolveBeginDate(
+                    persist.maxRecipientDate(account.getId(), folderCode), account.getPullDay());
+            MailProperties properties = new MailProperties(account);
+            MailConnection connection = MailConnectionUtil.receiveEmailsConnection(properties);
+            if (connection == null) {
+                return null;
+            }
+            Folder inbox = connection.getFolder(serverFolder);
+            if (inbox == null) {
+                return null;
+            }
+            inbox.open(Folder.READ_ONLY);
+            SearchTerm searchTerm = new ReceivedDateTerm(ComparisonTerm.GT, beginDate);
+            Message[] messages = inbox.search(searchTerm);
+            if (messages == null || messages.length == 0) {
+                return null;
+            }
+            return Arrays.asList(messages);
+        } catch (Exception e) {
+            log.error("searchMessages error, folder={}", serverFolder, e);
+            return null;
+        }
+    }
+
+    private void loadBodies(SmtpPopSettingsEntity account, MailReceivePersistSpi persist) {
+        List<MailReceiveBodyTask> tasks = persist.listNeedBody(account.getId());
+        if (CollectionUtils.isEmpty(tasks)) {
+            return;
+        }
+        try {
+            MailProperties properties = new MailProperties(account);
+            MailConnection connection = MailConnectionUtil.receiveEmailsConnection(properties);
+            if (connection == null) {
+                return;
+            }
+            for (MailReceiveBodyTask task : tasks) {
+                try {
+                    if (!StringUtils.hasText(task.getMsgUid())) {
+                        continue;
+                    }
+                    String serverFolder = FOLDER_INBOX.equals(task.getFolder()) ? SERVER_INBOX : SERVER_SENT;
+                    Message message = connection.getMessageByUid(task.getMsgUid(), serverFolder);
+                    if (message == null) {
+                        continue;
+                    }
+                    String body = persist.resolveBody(message);
+                    persist.updateBody(task.getMailId(), body);
+                    persist.saveAttachmentsFromMessage(account, task, message);
+                } catch (Exception ex) {
+                    log.error("load body error, mailId={}", task.getMailId(), ex);
+                }
+            }
+        } catch (Exception e) {
+            log.error("loadBodies error, accountId={}", account.getId(), e);
+        }
+    }
+}

+ 31 - 0
storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceiveHeadDraft.java

@@ -0,0 +1,31 @@
+package com.storlead.sales.mail.common.receive;
+
+import lombok.Data;
+
+import java.util.Date;
+
+/**
+ * 收信拉头解析结果(与业务表解耦)。
+ */
+@Data
+public class MailReceiveHeadDraft {
+
+    private String messageId;
+    private String subject;
+    private String fromAddr;
+    private String fromName;
+    private String recipient;
+    private String recipientName;
+    private String recipientCc;
+    private String recipientCcName;
+    private String folder;
+    private Integer inOutMark;
+    private Date sentDate;
+    private Date recipientDate;
+    private Long emailSize;
+    private Integer isRead;
+    private Integer isOnlyHead;
+    private String content;
+    private String msgUid;
+    private String inReplyTo;
+}

+ 62 - 0
storlead-mail/storlead-mail-common/src/main/java/com/storlead/sales/mail/common/receive/MailReceivePersistSpi.java

@@ -0,0 +1,62 @@
+package com.storlead.sales.mail.common.receive;
+
+import com.storlead.sales.mail.common.entity.SmtpPopSettingsEntity;
+
+import javax.mail.Message;
+import java.util.Date;
+import java.util.List;
+
+/**
+ * 收信落库 SPI:CRM / MKT 各自实现,引擎只编排协议流程。
+ */
+public interface MailReceivePersistSpi {
+
+    /**
+     * 增量游标:该账户该 folder 下最大收信时间。
+     */
+    Date maxRecipientDate(Long accountId, String folder);
+
+    /**
+     * 是否已存在同 Message-ID(可选;引擎不再据此短路,由 {@link #saveHead} 自行去重/合并)。
+     */
+    boolean existsMessageId(Long accountId, String messageId);
+
+    /**
+     * 保存拉头结果:新建返回主键;合并更新本地已有记录可返回 null(不触发 afterHeadsSaved)或返回原 id。
+     */
+    Long saveHead(SmtpPopSettingsEntity account, Message message, String folder, boolean isBackTask);
+
+    /**
+     * 一批拉头完成后的业务钩子(CRM:集成任务/自动回复;MKT:可空)。
+     */
+    void afterHeadsSaved(SmtpPopSettingsEntity account, String folder, List<Long> mailIds, boolean isBackTask);
+
+    /**
+     * 待补正文列表。
+     */
+    List<MailReceiveBodyTask> listNeedBody(Long accountId);
+
+    /**
+     * 回写正文并标记非仅头。
+     */
+    void updateBody(Long mailId, String content);
+
+    /**
+     * 从 Message 解析并保存附件(含下载)。
+     */
+    void saveAttachmentsFromMessage(SmtpPopSettingsEntity account, MailReceiveBodyTask task, Message message);
+
+    /**
+     * 解析正文;CRM 可覆盖为完整 HTML/CID 处理。
+     */
+    default String resolveBody(Message message) {
+        return MailMessageParser.extractBody(message);
+    }
+
+    /**
+     * 第三阶段:下载仍未落地的附件文件(CRM 有 download=0;MKT 可空实现)。
+     */
+    default void loadPendingAttachmentFiles(SmtpPopSettingsEntity account) {
+        // optional
+    }
+}

+ 183 - 0
storlead-mail/storlead-mail-crm/src/main/java/com/storlead/sales/mail/crm/receive/CrmMailReceivePersist.java

@@ -0,0 +1,183 @@
+package com.storlead.sales.mail.crm.receive;
+
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
+import com.storlead.framework.common.constant.CommonConstant;
+import com.storlead.sales.mail.common.entity.SmtpPopSettingsEntity;
+import com.storlead.sales.mail.common.receive.MailMessageParser;
+import com.storlead.sales.mail.common.receive.MailReceiveBodyTask;
+import com.storlead.sales.mail.common.receive.MailReceiveEngine;
+import com.storlead.sales.mail.common.receive.MailReceivePersistSpi;
+import com.storlead.sales.mail.crm.entity.EmailFolderRuleEntity;
+import com.storlead.sales.mail.crm.entity.EmailsEntity;
+import com.storlead.sales.mail.crm.enums.EmailBoxEnum;
+import com.storlead.sales.mail.crm.integration.service.MailIntegrationTaskService;
+import com.storlead.sales.mail.crm.service.EmailBlacklistRecordService;
+import com.storlead.sales.mail.crm.service.EmailFolderRuleService;
+import com.storlead.sales.mail.crm.service.EmailsService;
+import com.storlead.sales.mail.crm.service.impl.EmailsServiceImpl;
+import lombok.extern.log4j.Log4j2;
+import org.springframework.context.annotation.Lazy;
+import org.springframework.stereotype.Component;
+import org.springframework.util.CollectionUtils;
+import org.springframework.util.StringUtils;
+
+import javax.annotation.Resource;
+import javax.mail.Message;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+/**
+ * CRM 收信落库:复用 EmailsService 已验证的 convert/正文/附件逻辑。
+ */
+@Log4j2
+@Component
+public class CrmMailReceivePersist implements MailReceivePersistSpi {
+
+    @Resource
+    private EmailsService emailsService;
+
+    @Lazy
+    @Resource
+    private EmailsServiceImpl emailsServiceImpl;
+
+    @Resource
+    private EmailFolderRuleService folderRuleService;
+
+    @Resource
+    private EmailBlacklistRecordService blackListService;
+
+    @Resource
+    private MailIntegrationTaskService mailIntegrationTaskService;
+
+    @Override
+    public Date maxRecipientDate(Long accountId, String folder) {
+        QueryWrapper<EmailsEntity> qw = new QueryWrapper<>();
+        qw.eq("smtp_pop_id", accountId);
+        qw.eq("folder", folder);
+        qw.select("max(recipient_date) as recipient_date");
+        EmailsEntity email = emailsService.getOne(qw);
+        return email == null ? null : email.getRecipientDate();
+    }
+
+    @Override
+    public boolean existsMessageId(Long accountId, String messageId) {
+        if (!StringUtils.hasText(messageId)) {
+            return false;
+        }
+        // 与历史一致:按 folder 查;引擎侧已传全账户去重,这里再查任意 folder
+        Integer cnt = emailsService.count(new LambdaQueryWrapper<EmailsEntity>()
+                .eq(EmailsEntity::getSmtpPopId, accountId)
+                .eq(EmailsEntity::getMessageId, messageId)
+                .eq(EmailsEntity::getIsDelete, CommonConstant.DEL_FLAG_0));
+        return cnt != null && cnt > 0;
+    }
+
+    @Override
+    public Long saveHead(SmtpPopSettingsEntity account, Message message, String folder, boolean isBackTask) {
+        try {
+            String[] headers = message.getHeader("Message-ID");
+            if (headers == null || headers.length == 0) {
+                return null;
+            }
+            String messageId = MailMessageParser.normalizeMessageId(headers[0]);
+            String old = emailsService.getEmailMessageIdsByMessageId(messageId, folder, account.getId());
+            if (old != null) {
+                return null;
+            }
+            EmailsEntity entity = emailsService.convertMessageToEmailVo(message, true);
+            if (entity == null) {
+                return null;
+            }
+            entity.setMessageId(messageId);
+            entity.setSmtpPopId(account.getId());
+            entity.setOwnerBy(account.getOwnerBy());
+            entity.setFolder(folder);
+            entity.setInOutMark(MailReceiveEngine.FOLDER_INBOX.equals(folder) ? 1 : 2);
+            entity.setMsgUid(MailMessageParser.resolveMsgUid(message, account.getProtocolType()));
+
+            List<EmailFolderRuleEntity> folderRules = folderRuleService.getUserFolderRules(account.getOwnerBy());
+            if (!CollectionUtils.isEmpty(folderRules)) {
+                folderRuleService.carveUpMailCustomFolder(entity, folderRules);
+            }
+            Set<String> blackls = blackListService.getBlackEmaillist(account.getOwnerBy());
+            if (!CollectionUtils.isEmpty(blackls) && blackls.contains(entity.getFrom())) {
+                entity.setIsTrash(1);
+            }
+            if (emailsService.saveOrUpdate(entity)) {
+                return entity.getId();
+            }
+        } catch (Exception e) {
+            log.error("crm saveHead error, accountId={}", account.getId(), e);
+        }
+        return null;
+    }
+
+    @Override
+    public void afterHeadsSaved(SmtpPopSettingsEntity account, String folder, List<Long> mailIds, boolean isBackTask) {
+        if (CollectionUtils.isEmpty(mailIds)) {
+            return;
+        }
+        boolean needNewRemind = Boolean.TRUE.equals(isBackTask)
+                && EmailBoxEnum.INBOX.code.equals(folder);
+        mailIntegrationTaskService.createAfterPullHeadMail(mailIds, account.getOwnerBy(), needNewRemind);
+        List<EmailsEntity> entities = emailsService.listByIds(mailIds);
+        if (!CollectionUtils.isEmpty(entities)) {
+            emailsServiceImpl.autoReplyMails(entities);
+        }
+    }
+
+    @Override
+    public List<MailReceiveBodyTask> listNeedBody(Long accountId) {
+        List<EmailsEntity> list = emailsService.list(new LambdaQueryWrapper<EmailsEntity>()
+                .select(EmailsEntity::getId, EmailsEntity::getMsgUid, EmailsEntity::getFolder,
+                        EmailsEntity::getRecipientDate)
+                .eq(EmailsEntity::getSmtpPopId, accountId)
+                .eq(EmailsEntity::getIsOnlyHead, 1)
+                .eq(EmailsEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .in(EmailsEntity::getFolder, EmailBoxEnum.INBOX.code, EmailBoxEnum.SENT.code));
+        if (CollectionUtils.isEmpty(list)) {
+            return new ArrayList<>();
+        }
+        return list.stream().map(e -> {
+            MailReceiveBodyTask t = new MailReceiveBodyTask();
+            t.setMailId(e.getId());
+            t.setMsgUid(e.getMsgUid());
+            t.setFolder(e.getFolder());
+            t.setRecipientDate(e.getRecipientDate());
+            return t;
+        }).collect(Collectors.toList());
+    }
+
+    @Override
+    public String resolveBody(Message message) {
+        return emailsService.getEmailBodyContent(message);
+    }
+
+    @Override
+    public void updateBody(Long mailId, String content) {
+        LambdaUpdateWrapper<EmailsEntity> update = new LambdaUpdateWrapper<>();
+        update.eq(EmailsEntity::getId, mailId);
+        update.set(EmailsEntity::getContent, content);
+        update.set(EmailsEntity::getIsOnlyHead, 0);
+        emailsService.update(update);
+    }
+
+    @Override
+    public void saveAttachmentsFromMessage(SmtpPopSettingsEntity account, MailReceiveBodyTask task, Message message) {
+        try {
+            emailsService.pullMailAttachmentHead(task.getMailId(), task.getRecipientDate(), message.getContent(), account);
+        } catch (Exception e) {
+            log.error("crm saveAttachments error, mailId={}", task.getMailId(), e);
+        }
+    }
+
+    @Override
+    public void loadPendingAttachmentFiles(SmtpPopSettingsEntity account) {
+        emailsService.loadMailFiles(account);
+    }
+}

+ 13 - 20
storlead-mail/storlead-mail-crm/src/main/java/com/storlead/sales/mail/crm/service/impl/EmailsServiceImpl.java

@@ -40,6 +40,9 @@ import com.storlead.sales.mail.crm.integration.handler.NewMailRemindHandler;
 import com.storlead.sales.mail.crm.integration.service.MailIntegrationTaskService;
 import com.storlead.sales.mail.crm.integration.util.MailIdsParser;
 import com.storlead.sales.mail.common.util.EmailHelper;
+import com.storlead.sales.mail.common.enums.MailSceneEnum;
+import com.storlead.sales.mail.common.receive.MailReceiveEngine;
+import com.storlead.sales.mail.crm.receive.CrmMailReceivePersist;
 import com.storlead.sales.mail.crm.util.EmailSenderWithThreadLocal;
 import com.storlead.sales.mail.crm.util.ReceiveMailQueueThreadPool;
 import com.sun.mail.imap.IMAPBodyPart;
@@ -47,6 +50,7 @@ import com.sun.mail.imap.IMAPFolder;
 import com.sun.mail.pop3.POP3Folder;
 import lombok.extern.log4j.Log4j2;
 import org.apache.commons.io.FilenameUtils;
+import org.springframework.context.annotation.Lazy;
 import org.springframework.core.env.Environment;
 import org.springframework.stereotype.Service;
 import org.springframework.util.CollectionUtils;
@@ -138,6 +142,13 @@ public class EmailsServiceImpl extends MyBaseServiceImpl<EmailsMapper, EmailsEnt
     @Resource
     private IntegrationExternalInvoker integrationExternalInvoker;
 
+    @Resource
+    private MailReceiveEngine mailReceiveEngine;
+
+    @Lazy
+    @Resource
+    private CrmMailReceivePersist crmMailReceivePersist;
+
     @Override
     public void updateSelectBindCustomerMail() {
         //   this.baseMapper.updateCleanBindCustomerMail();
@@ -416,26 +427,8 @@ public class EmailsServiceImpl extends MyBaseServiceImpl<EmailsMapper, EmailsEnt
     }
     @Override
     public void receiveEmails(SmtpPopSettingsEntity smtpPop,Boolean isBackTask) {
-        String active = environment.getProperty("spring.profiles.active");
-//        if("dev".equals(active)) {
-//            return;
-//        }
-        String taskId = smtpPop.getEmailAddress()+"_"+smtpPop.getOwnerBy().toString();
-        ReceiveMailQueueThreadPool instance = ReceiveMailQueueThreadPool.getInstance();
-        if (instance.getTaskInProgressByMail(taskId)) {
-            log.error("有任务正在进行--:>>>>>>>>>>>>>"+smtpPop.getEmailAddress());
-            return;
-        }
-        instance.addTask(taskId,new Runnable() {
-            @Override
-            public void run() {
-                receiveHeadEmails(smtpPop,isBackTask);
-            }
-        });
-
-        addLoadMailContentTask(smtpPop);
-
-        addLoadMailFilesTask(smtpPop);
+        // 公共引擎编排三阶段;CRM 落库/钩子由 CrmMailReceivePersist 完成(行为对齐原逻辑)
+        mailReceiveEngine.receive(smtpPop, MailSceneEnum.CRM, crmMailReceivePersist, Boolean.TRUE.equals(isBackTask));
     }
 
     @Override

+ 8 - 65
storlead-mail/storlead-mail-crm/src/main/java/com/storlead/sales/mail/crm/util/ReceiveMailQueueThreadPool.java

@@ -1,43 +1,16 @@
 package com.storlead.sales.mail.crm.util;
 
-import com.storlead.framework.common.thread.ThreadPoolUtil;
+import com.storlead.sales.mail.common.util.MailReceiveQueueThreadPool;
 import lombok.extern.log4j.Log4j2;
 
-import java.util.*;
-import java.util.concurrent.ThreadPoolExecutor;
-
 /**
- * @program: sp-sales-platform
- * @description:
- * @author: chenkq
- * @create: 2024-08-26 16:27
+ * CRM 收信队列:委托公共 {@link MailReceiveQueueThreadPool},保留类名兼容旧引用。
  */
 @Log4j2
 public class ReceiveMailQueueThreadPool {
 
-
     private static volatile ReceiveMailQueueThreadPool instance;
 
-    // 用于存储每个人的任务队列
-    private final Map<String, Queue<Runnable>> taskMap;
-
-    public Boolean getTaskInProgressByMail(String mailAddress) {
-        return taskInProgress.getOrDefault(mailAddress, false);
-    }
-
-    // 用于存储每个人当前是否有任务在进行中
-    private final Map<String, Boolean> taskInProgress;
-
-    // 线程池,用于执行任务
-    private final ThreadPoolExecutor executorService;
-
-    public ReceiveMailQueueThreadPool() {
-        this.taskMap = new HashMap<>();
-        this.taskInProgress = new HashMap<>();
-        this.executorService = ThreadPoolUtil.getThreadExecutor();
-    }
-
-    // 双重检查锁的单例获取方法
     public static ReceiveMailQueueThreadPool getInstance() {
         if (instance == null) {
             synchronized (ReceiveMailQueueThreadPool.class) {
@@ -49,49 +22,19 @@ public class ReceiveMailQueueThreadPool {
         return instance;
     }
 
-    // 添加任务到某个人的队列
-    public synchronized void addTask(String person, Runnable task) {
-        if (Objects.nonNull(taskMap) && Objects.nonNull(taskMap.get(person))) {
-            return; // 该人有任务正在进行中
-        }
-        taskMap.putIfAbsent(person, new LinkedList<>());
-        taskMap.get(person).offer(task);
-        taskInProgress.putIfAbsent(person, false);
-        processNextTask(person);
+    public Boolean getTaskInProgressByMail(String mailAddress) {
+        return MailReceiveQueueThreadPool.getInstance().getTaskInProgressByMail(mailAddress);
     }
 
-    // 处理某个人的下一个任务
-    private synchronized void processNextTask(String person) {
-        if (taskInProgress.getOrDefault(person, false)) {
-            return; // 该人有任务正在进行中
-        }
-
-        Queue<Runnable> queue = taskMap.get(person);
-        if (queue == null || queue.isEmpty()) {
-            log.error("没有任务了--"+person);
-            return; // 无任务可处理
-        }
-
-        Runnable task = queue.poll();
-        taskInProgress.put(person, true);
-        executorService.submit(() -> {
-            try {
-                task.run(); // 执行任务
-            } finally {
-                taskCompleted(person); // 标记任务完成并处理下一个任务
-            }
-        });
+    public synchronized void addTask(String person, Runnable task) {
+        MailReceiveQueueThreadPool.getInstance().addTask(person, task);
     }
 
-    // 标记任务完成
     public synchronized void taskCompleted(String person) {
-        taskInProgress.put(person, false);
-        taskMap.remove(person);
-        processNextTask(person);
+        MailReceiveQueueThreadPool.getInstance().taskCompleted(person);
     }
 
-    // 关闭线程池
     public void shutdown() {
-        executorService.shutdown();
+        // shared pool lifecycle managed by framework executor
     }
 }

+ 357 - 0
storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/receive/MktMailReceivePersist.java

@@ -0,0 +1,357 @@
+package com.storlead.sales.mail.mkt.receive;
+
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
+import com.storlead.framework.common.constant.CommonConstant;
+import com.storlead.framework.common.util.DateUtils;
+import com.storlead.framework.common.util.RandomGenerateHelper;
+import com.storlead.sales.mail.common.entity.SmtpPopSettingsEntity;
+import com.storlead.sales.mail.common.properties.MailFileProperties;
+import com.storlead.sales.mail.common.receive.MailMessageParser;
+import com.storlead.sales.mail.common.receive.MailReceiveBodyTask;
+import com.storlead.sales.mail.common.receive.MailReceiveHeadDraft;
+import com.storlead.sales.mail.common.receive.MailReceivePersistSpi;
+import com.storlead.sales.mail.mkt.entity.MktMailAttachmentEntity;
+import com.storlead.sales.mail.mkt.entity.MktMailMessageEntity;
+import com.storlead.sales.mail.mkt.enums.MktMailFolderEnum;
+import com.storlead.sales.mail.mkt.service.MktMailAttachmentService;
+import com.storlead.sales.mail.mkt.service.MktMailMessageService;
+import com.storlead.sales.mail.mkt.support.MktMailHeaders;
+import lombok.extern.log4j.Log4j2;
+import org.apache.commons.io.FilenameUtils;
+import org.springframework.stereotype.Component;
+import org.springframework.util.CollectionUtils;
+import org.springframework.util.StringUtils;
+
+import javax.annotation.Resource;
+import javax.mail.BodyPart;
+import javax.mail.Message;
+import javax.mail.Multipart;
+import javax.mail.Part;
+import javax.mail.internet.MimeUtility;
+import java.io.File;
+import java.io.FileOutputStream;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import java.util.stream.Collectors;
+
+/**
+ * MKT 收信落库:mkt_mail_message + reply_msg_id 绑定 + mkt_mail_attachment。
+ */
+@Log4j2
+@Component
+public class MktMailReceivePersist implements MailReceivePersistSpi {
+
+    private static final int STATUS_RECEIVED = 1;
+
+    @Resource
+    private MktMailMessageService mktMailMessageService;
+
+    @Resource
+    private MktMailAttachmentService mktMailAttachmentService;
+
+    @Resource
+    private MailFileProperties mailFileProperties;
+
+    @Override
+    public Date maxRecipientDate(Long accountId, String folder) {
+        QueryWrapper<MktMailMessageEntity> qw = new QueryWrapper<>();
+        qw.eq("account_id", accountId);
+        qw.eq("folder", folder);
+        qw.select("max(recipient_date) as recipient_date");
+        MktMailMessageEntity latest = mktMailMessageService.getOne(qw);
+        return latest == null ? null : latest.getRecipientDate();
+    }
+
+    @Override
+    public boolean existsMessageId(Long accountId, String messageId) {
+        return findByMessageId(accountId, messageId) != null;
+    }
+
+    @Override
+    public Long saveHead(SmtpPopSettingsEntity account, Message message, String folder, boolean isBackTask) {
+        try {
+            String[] headers = message.getHeader("Message-ID");
+            if (headers == null || headers.length == 0 || !StringUtils.hasText(headers[0])) {
+                return null;
+            }
+            String messageId = MailMessageParser.normalizeMessageId(headers[0]);
+            MailReceiveHeadDraft draft = MailMessageParser.parseHead(message, folder, true);
+            if (draft == null) {
+                return null;
+            }
+
+            // 优先按发信时写入的主键头对齐;兼容旧邮件再按 RFC Message-ID 对齐
+            MktMailMessageEntity existing = findLocalOutboundForMerge(account.getId(), message, messageId);
+            if (existing != null) {
+                // 保留原主键 id,只回写服务器侧字段,绝不 insert 新行
+                if (isLocalOutbound(existing) && MktMailFolderEnum.SENT.code.equals(folder)) {
+                    mergeServerFields(existing, message, draft, messageId, account);
+                }
+                return null;
+            }
+
+            MktMailMessageEntity entity = new MktMailMessageEntity();
+            entity.setAccountId(account.getId());
+            entity.setOrgId(account.getOrgId());
+            entity.setOwnerBy(account.getOwnerBy());
+            entity.setCreateBy(account.getOwnerBy());
+            entity.setMessageId(messageId);
+            entity.setSubject(draft.getSubject());
+            entity.setFromAddr(draft.getFromAddr());
+            entity.setFromName(draft.getFromName());
+            entity.setRecipient(StringUtils.hasText(draft.getRecipient()) ? draft.getRecipient() : "");
+            entity.setRecipientCc(draft.getRecipientCc());
+            entity.setFolder(folder);
+            entity.setInOutMark(draft.getInOutMark());
+            entity.setSentDate(draft.getSentDate());
+            entity.setRecipientDate(draft.getRecipientDate() == null ? new Date() : draft.getRecipientDate());
+            if (draft.getSentDate() != null) {
+                entity.setSentTime(draft.getSentDate());
+            }
+            entity.setEmailSize(draft.getEmailSize());
+            entity.setIsRead(draft.getIsRead());
+            entity.setIsOnlyHead(draft.getIsOnlyHead());
+            entity.setContent(draft.getContent());
+            entity.setMsgUid(MailMessageParser.resolveMsgUid(message, account.getProtocolType()));
+            entity.setStatus(STATUS_RECEIVED);
+            entity.setIsDelete(CommonConstant.DEL_FLAG_0);
+            entity.setEnabled(Boolean.TRUE);
+            entity.setIsTrack(0);
+            entity.setOpened(0);
+            entity.setOpenCount(0);
+
+            bindReply(account.getId(), draft.getInReplyTo(), entity);
+
+            if (mktMailMessageService.save(entity)) {
+                return entity.getId();
+            }
+        } catch (Exception e) {
+            log.error("mkt saveHead error, accountId={}", account.getId(), e);
+        }
+        return null;
+    }
+
+    /**
+     * 服务器邮件 → 本地已发出营销记录:先读自定义主键头,再回退 Message-ID。
+     */
+    private MktMailMessageEntity findLocalOutboundForMerge(Long accountId, Message message, String messageId) {
+        MktMailMessageEntity byPk = findByLocalMailIdHeader(accountId, message);
+        if (byPk != null) {
+            return byPk;
+        }
+        return findByMessageId(accountId, messageId);
+    }
+
+    private MktMailMessageEntity findByLocalMailIdHeader(Long accountId, Message message) {
+        try {
+            String[] values = message.getHeader(MktMailHeaders.LOCAL_MAIL_ID);
+            if (values == null || values.length == 0 || !StringUtils.hasText(values[0])) {
+                return null;
+            }
+            Long localId = Long.parseLong(values[0].trim());
+            return mktMailMessageService.getOne(new LambdaQueryWrapper<MktMailMessageEntity>()
+                    .eq(MktMailMessageEntity::getId, localId)
+                    .eq(MktMailMessageEntity::getAccountId, accountId)
+                    .eq(MktMailMessageEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                    .last("limit 1"));
+        } catch (Exception e) {
+            log.warn("parse {} header failed", MktMailHeaders.LOCAL_MAIL_ID, e);
+            return null;
+        }
+    }
+
+    /**
+     * 本地发信记录(SENTING)或已合并的 SENT,与服务器 Sent 对齐时保留业务字段。
+     */
+    private boolean isLocalOutbound(MktMailMessageEntity existing) {
+        String folder = existing.getFolder();
+        return MktMailFolderEnum.SENTING.code.equals(folder)
+                || MktMailFolderEnum.SENT.code.equals(folder)
+                || !StringUtils.hasText(folder);
+    }
+
+    /**
+     * 按主键 id 更新:保留 biz_ref / track / 正文等,只写服务器拉取字段;不 insert。
+     */
+    private void mergeServerFields(MktMailMessageEntity existing,
+                                   Message message,
+                                   MailReceiveHeadDraft draft,
+                                   String messageId,
+                                   SmtpPopSettingsEntity account) {
+        String msgUid = MailMessageParser.resolveMsgUid(message, account.getProtocolType());
+        LambdaUpdateWrapper<MktMailMessageEntity> update = new LambdaUpdateWrapper<>();
+        update.eq(MktMailMessageEntity::getId, existing.getId());
+        update.set(MktMailMessageEntity::getMessageId, messageId);
+        update.set(MktMailMessageEntity::getFolder, MktMailFolderEnum.SENT.code);
+        update.set(MktMailMessageEntity::getInOutMark, 2);
+        if (StringUtils.hasText(msgUid)) {
+            update.set(MktMailMessageEntity::getMsgUid, msgUid);
+        }
+        if (draft.getRecipientDate() != null) {
+            update.set(MktMailMessageEntity::getRecipientDate, draft.getRecipientDate());
+        }
+        if (draft.getSentDate() != null) {
+            update.set(MktMailMessageEntity::getSentDate, draft.getSentDate());
+            update.set(MktMailMessageEntity::getSentTime, draft.getSentDate());
+        }
+        if (draft.getEmailSize() != null) {
+            update.set(MktMailMessageEntity::getEmailSize, draft.getEmailSize());
+        }
+        if (draft.getIsRead() != null) {
+            update.set(MktMailMessageEntity::getIsRead, draft.getIsRead());
+        }
+        // 主题/收发件人以本地发信为准,服务器有值且本地为空时才补
+        if (!StringUtils.hasText(existing.getSubject()) && StringUtils.hasText(draft.getSubject())) {
+            update.set(MktMailMessageEntity::getSubject, draft.getSubject());
+        }
+        if (!StringUtils.hasText(existing.getFromAddr()) && StringUtils.hasText(draft.getFromAddr())) {
+            update.set(MktMailMessageEntity::getFromAddr, draft.getFromAddr());
+            update.set(MktMailMessageEntity::getFromName, draft.getFromName());
+        }
+        if (!StringUtils.hasText(existing.getRecipient()) && StringUtils.hasText(draft.getRecipient())) {
+            update.set(MktMailMessageEntity::getRecipient, draft.getRecipient());
+        }
+        mktMailMessageService.update(update);
+        log.info("mkt merge sent mail keep primary id={}, messageId={}", existing.getId(), messageId);
+    }
+
+    @Override
+    public void afterHeadsSaved(SmtpPopSettingsEntity account, String folder, List<Long> mailIds, boolean isBackTask) {
+        // MKT 无 CRM 集成/自动回复
+    }
+
+    @Override
+    public List<MailReceiveBodyTask> listNeedBody(Long accountId) {
+        List<MktMailMessageEntity> list = mktMailMessageService.list(new LambdaQueryWrapper<MktMailMessageEntity>()
+                .select(MktMailMessageEntity::getId, MktMailMessageEntity::getMsgUid,
+                        MktMailMessageEntity::getFolder, MktMailMessageEntity::getRecipientDate)
+                .eq(MktMailMessageEntity::getAccountId, accountId)
+                .eq(MktMailMessageEntity::getIsOnlyHead, 1)
+                .eq(MktMailMessageEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .in(MktMailMessageEntity::getFolder,
+                        MktMailFolderEnum.INBOX.code, MktMailFolderEnum.SENT.code));
+        if (CollectionUtils.isEmpty(list)) {
+            return new ArrayList<>();
+        }
+        return list.stream().map(e -> {
+            MailReceiveBodyTask t = new MailReceiveBodyTask();
+            t.setMailId(e.getId());
+            t.setMsgUid(e.getMsgUid());
+            t.setFolder(e.getFolder());
+            t.setRecipientDate(e.getRecipientDate());
+            return t;
+        }).collect(Collectors.toList());
+    }
+
+    @Override
+    public void updateBody(Long mailId, String content) {
+        LambdaUpdateWrapper<MktMailMessageEntity> update = new LambdaUpdateWrapper<>();
+        update.eq(MktMailMessageEntity::getId, mailId);
+        update.set(MktMailMessageEntity::getContent, content);
+        update.set(MktMailMessageEntity::getIsOnlyHead, 0);
+        mktMailMessageService.update(update);
+    }
+
+    @Override
+    public void saveAttachmentsFromMessage(SmtpPopSettingsEntity account, MailReceiveBodyTask task, Message message) {
+        try {
+            int exist = mktMailAttachmentService.count(new LambdaQueryWrapper<MktMailAttachmentEntity>()
+                    .eq(MktMailAttachmentEntity::getMailId, task.getMailId())
+                    .eq(MktMailAttachmentEntity::getIsDelete, CommonConstant.DEL_FLAG_0));
+            if (exist > 0) {
+                return;
+            }
+            Object content = message.getContent();
+            if (!(content instanceof Multipart)) {
+                return;
+            }
+            String day = DateUtils.date2Str(
+                    task.getRecipientDate() == null ? new Date() : task.getRecipientDate(), DateUtils.yyyyMMdd);
+            String downloadDir = mailFileProperties.getPath().getPath() + File.separator
+                    + account.getEmailAddress() + File.separator + day + File.separator;
+            Multipart multipart = (Multipart) content;
+            for (int i = 0; i < multipart.getCount(); i++) {
+                BodyPart bodyPart = multipart.getBodyPart(i);
+                if (Part.ATTACHMENT.equalsIgnoreCase(bodyPart.getDisposition())
+                        || bodyPart.isMimeType("application/octet-stream")) {
+                    saveOneAttachment(bodyPart, downloadDir, task.getMailId(), account);
+                }
+            }
+        } catch (Exception e) {
+            log.error("mkt saveAttachments error, mailId={}", task.getMailId(), e);
+        }
+    }
+
+    private void bindReply(Long accountId, String inReplyTo, MktMailMessageEntity entity) {
+        if (!StringUtils.hasText(inReplyTo)) {
+            return;
+        }
+        String normalized = MailMessageParser.normalizeMessageId(inReplyTo);
+        entity.setInReplyTo(normalized);
+        MktMailMessageEntity parent = findByMessageId(accountId, normalized);
+        if (parent == null) {
+            return;
+        }
+        entity.setReplyMsgId(parent.getId());
+        if (!StringUtils.hasText(entity.getBizRef()) && StringUtils.hasText(parent.getBizRef())) {
+            entity.setBizRef(parent.getBizRef());
+        }
+        if (!StringUtils.hasText(entity.getBizBatchNo()) && StringUtils.hasText(parent.getBizBatchNo())) {
+            entity.setBizBatchNo(parent.getBizBatchNo());
+        }
+    }
+
+    private MktMailMessageEntity findByMessageId(Long accountId, String messageId) {
+        if (!StringUtils.hasText(messageId)) {
+            return null;
+        }
+        String normalized = MailMessageParser.normalizeMessageId(messageId);
+        String bare = MailMessageParser.stripBrackets(normalized);
+        List<MktMailMessageEntity> list = mktMailMessageService.list(new LambdaQueryWrapper<MktMailMessageEntity>()
+                .eq(MktMailMessageEntity::getAccountId, accountId)
+                .eq(MktMailMessageEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
+                .and(w -> w.eq(MktMailMessageEntity::getMessageId, normalized)
+                        .or().eq(MktMailMessageEntity::getMessageId, "<" + bare + ">")
+                        .or().eq(MktMailMessageEntity::getMessageId, bare))
+                .orderByDesc(MktMailMessageEntity::getId)
+                .last("limit 1"));
+        return CollectionUtils.isEmpty(list) ? null : list.get(0);
+    }
+
+    private void saveOneAttachment(BodyPart bodyPart, String downloadDir, Long mailId, SmtpPopSettingsEntity account) {
+        try {
+            String fileName = bodyPart.getFileName();
+            if (!StringUtils.hasText(fileName)) {
+                return;
+            }
+            fileName = MimeUtility.decodeText(fileName).replace("\n", "").replace("\r", "");
+            String code = RandomGenerateHelper.generateRandomString(6);
+            if (!new File(downloadDir).exists()) {
+                new File(downloadDir).mkdirs();
+            }
+            String storedName = fileName + "." + code;
+            File file = new File(downloadDir + File.separator + storedName);
+            try (FileOutputStream output = new FileOutputStream(file)) {
+                bodyPart.getInputStream().transferTo(output);
+            }
+            MktMailAttachmentEntity att = new MktMailAttachmentEntity();
+            att.setMailId(mailId);
+            att.setAccountId(account.getId());
+            att.setFileCode(code);
+            att.setFileName(FilenameUtils.getBaseName(fileName));
+            att.setFileExt(FilenameUtils.getExtension(fileName));
+            att.setFilePath(file.getAbsolutePath().replace(mailFileProperties.getPath().getPath(), ""));
+            att.setFileSize(file.length());
+            att.setOwnerBy(account.getOwnerBy());
+            att.setCreateBy(account.getOwnerBy());
+            att.setIsDelete(CommonConstant.DEL_FLAG_0);
+            att.setEnabled(Boolean.TRUE);
+            mktMailAttachmentService.save(att);
+        } catch (Exception e) {
+            log.error("mkt saveOneAttachment error", e);
+        }
+    }
+}

+ 5 - 0
storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/service/impl/MktMailMessageServiceImpl.java

@@ -25,6 +25,7 @@ import com.storlead.sales.mail.mkt.pojo.MktSendMailDTO;
 import com.storlead.sales.mail.mkt.service.MktMailAttachmentService;
 import com.storlead.sales.mail.mkt.service.MktMailMessageService;
 import com.storlead.sales.mail.mkt.enums.MktMailFolderEnum;
+import com.storlead.sales.mail.mkt.support.MktMailHeaders;
 import com.storlead.sales.mail.mkt.util.MktMailSender;
 import lombok.extern.log4j.Log4j2;
 import org.springframework.beans.factory.annotation.Value;
@@ -383,6 +384,10 @@ public class MktMailMessageServiceImpl extends MyBaseServiceImpl<MktMailMessageM
         }
 
         message.setContent(multipart);
+        // 携带本地主键,收信从服务器拉回时按此头更新原记录,不新增
+        if (entity.getId() != null) {
+            message.setHeader(MktMailHeaders.LOCAL_MAIL_ID, String.valueOf(entity.getId()));
+        }
         message.saveChanges();
         return message;
     }

+ 7 - 483
storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/service/impl/MktMailReceiveServiceImpl.java

@@ -1,74 +1,27 @@
 package com.storlead.sales.mail.mkt.service.impl;
 
-import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
-import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
-import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
-import com.storlead.framework.common.constant.CommonConstant;
-import com.storlead.framework.common.util.DateUtil;
-import com.storlead.framework.common.util.DateUtils;
-import com.storlead.framework.common.util.RandomGenerateHelper;
-import com.storlead.sales.mail.common.connection.MailConnection;
-import com.storlead.sales.mail.common.connection.client.MailConnectionUtil;
-import com.storlead.sales.mail.common.connection.config.MailProperties;
 import com.storlead.sales.mail.common.entity.SmtpPopSettingsEntity;
-import com.storlead.sales.mail.common.properties.MailFileProperties;
-import com.storlead.sales.mail.common.util.MailReceiveQueueThreadPool;
-import com.storlead.sales.mail.mkt.entity.MktMailAttachmentEntity;
-import com.storlead.sales.mail.mkt.entity.MktMailMessageEntity;
-import com.storlead.sales.mail.mkt.enums.MktMailFolderEnum;
-import com.storlead.sales.mail.mkt.service.MktMailAttachmentService;
-import com.storlead.sales.mail.mkt.service.MktMailMessageService;
+import com.storlead.sales.mail.common.enums.MailSceneEnum;
+import com.storlead.sales.mail.common.receive.MailReceiveEngine;
+import com.storlead.sales.mail.mkt.receive.MktMailReceivePersist;
 import com.storlead.sales.mail.mkt.service.MktMailReceiveService;
-import com.sun.mail.imap.IMAPFolder;
-import com.sun.mail.pop3.POP3Folder;
 import lombok.extern.log4j.Log4j2;
-import org.apache.commons.io.FilenameUtils;
 import org.springframework.stereotype.Service;
-import org.springframework.util.CollectionUtils;
-import org.springframework.util.StringUtils;
 
 import javax.annotation.Resource;
-import javax.mail.Address;
-import javax.mail.BodyPart;
-import javax.mail.Flags;
-import javax.mail.Folder;
-import javax.mail.Message;
-import javax.mail.Multipart;
-import javax.mail.Part;
-import javax.mail.internet.InternetAddress;
-import javax.mail.internet.MimeUtility;
-import javax.mail.search.ComparisonTerm;
-import javax.mail.search.ReceivedDateTerm;
-import javax.mail.search.SearchTerm;
-import java.io.File;
-import java.io.FileOutputStream;
-import java.time.LocalDate;
-import java.util.ArrayList;
-import java.util.Arrays;
-import java.util.Date;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.Objects;
-import java.util.stream.Collectors;
 
 /**
- * 营销收信:对齐 CRM receiveEmails,落 mkt_mail_message(folder 区分),并绑定 reply_msg_id
+ * 营销收信入口:复用公共 {@link MailReceiveEngine},落库由 {@link MktMailReceivePersist} 完成。
  */
 @Log4j2
 @Service
 public class MktMailReceiveServiceImpl implements MktMailReceiveService {
 
-    private static final int STATUS_RECEIVED = 1;
-
-    @Resource
-    private MktMailMessageService mktMailMessageService;
-
     @Resource
-    private MktMailAttachmentService mktMailAttachmentService;
+    private MailReceiveEngine mailReceiveEngine;
 
     @Resource
-    private MailFileProperties mailFileProperties;
+    private MktMailReceivePersist mktMailReceivePersist;
 
     @Override
     public void receiveEmails(SmtpPopSettingsEntity account) {
@@ -79,435 +32,6 @@ public class MktMailReceiveServiceImpl implements MktMailReceiveService {
             log.warn("mkt receive skipped, enable_receive!=1, accountId={}", account.getId());
             return;
         }
-        String taskKey = "mkt_" + account.getEmailAddress() + "_" + account.getId();
-        MailReceiveQueueThreadPool pool = MailReceiveQueueThreadPool.getInstance();
-        if (Boolean.TRUE.equals(pool.getTaskInProgressByMail(taskKey))) {
-            log.warn("mkt receive already running: {}", taskKey);
-            return;
-        }
-        pool.addTask(taskKey, () -> {
-            try {
-                receiveHeadEmails(account);
-                loadMailContents(account);
-            } catch (Exception e) {
-                log.error("mkt receiveEmails error, accountId={}", account.getId(), e);
-            }
-        });
-    }
-
-    private void receiveHeadEmails(SmtpPopSettingsEntity account) {
-        pullFolder(account, "INBOX", MktMailFolderEnum.INBOX);
-        pullFolder(account, "Sent Items", MktMailFolderEnum.SENT);
-    }
-
-    private void pullFolder(SmtpPopSettingsEntity account, String serverFolder, MktMailFolderEnum folderEnum) {
-        try {
-            List<Message> messages = searchMessages(account, serverFolder, folderEnum);
-            if (CollectionUtils.isEmpty(messages)) {
-                return;
-            }
-            List<String> recentIds = listRecentMessageIds(account.getId(), folderEnum.code);
-            for (Message message : messages) {
-                try {
-                    String[] headers = message.getHeader("Message-ID");
-                    if (headers == null || headers.length == 0 || !StringUtils.hasText(headers[0])) {
-                        continue;
-                    }
-                    String messageId = normalizeMessageId(headers[0]);
-                    if (recentIds.contains(messageId) || existsByMessageId(account.getId(), messageId)) {
-                        continue;
-                    }
-                    MktMailMessageEntity entity = convertMessage(message, folderEnum);
-                    if (entity == null) {
-                        continue;
-                    }
-                    entity.setMessageId(messageId);
-                    entity.setAccountId(account.getId());
-                    entity.setOrgId(account.getOrgId());
-                    entity.setOwnerBy(account.getOwnerBy());
-                    entity.setCreateBy(account.getOwnerBy());
-                    entity.setStatus(STATUS_RECEIVED);
-                    entity.setIsDelete(CommonConstant.DEL_FLAG_0);
-                    entity.setEnabled(Boolean.TRUE);
-                    entity.setIsTrack(0);
-                    entity.setOpened(0);
-                    entity.setOpenCount(0);
-                    fillMsgUid(message, account, entity);
-                    bindReply(account.getId(), message, entity);
-
-                    if (mktMailMessageService.save(entity)) {
-                        recentIds.add(messageId);
-                    }
-                } catch (Exception ex) {
-                    log.error("mkt pull one mail error, accountId={}", account.getId(), ex);
-                }
-            }
-        } catch (Exception e) {
-            log.error("mkt pullFolder error, folder={}, accountId={}", serverFolder, account.getId(), e);
-        }
-    }
-
-    private List<Message> searchMessages(SmtpPopSettingsEntity account, String folderName, MktMailFolderEnum folderEnum) {
-        try {
-            QueryWrapper<MktMailMessageEntity> qw = new QueryWrapper<>();
-            qw.eq("account_id", account.getId());
-            qw.eq("folder", folderEnum.code);
-            qw.select("max(recipient_date) as recipient_date");
-            MktMailMessageEntity latest = mktMailMessageService.getOne(qw);
-
-            MailProperties properties = new MailProperties(account);
-            MailConnection connection = MailConnectionUtil.receiveEmailsConnection(properties);
-            if (connection == null) {
-                return null;
-            }
-            Folder inbox = connection.getFolder(folderName);
-            if (inbox == null) {
-                return null;
-            }
-            inbox.open(Folder.READ_ONLY);
-
-            Date beginDate = latest == null ? null : latest.getRecipientDate();
-            if (beginDate == null) {
-                LocalDate currentDate = LocalDate.now();
-                if (account.getPullDay() == null || Integer.valueOf(0).equals(account.getPullDay())) {
-                    beginDate = DateUtil.getDateByFormat(currentDate.minusYears(1).toString(), "yyyy-MM-dd");
-                } else {
-                    beginDate = DateUtil.getDateByFormat(currentDate.minusDays(account.getPullDay()).toString(), "yyyy-MM-dd");
-                }
-            }
-            SearchTerm searchTerm = new ReceivedDateTerm(ComparisonTerm.GT, beginDate);
-            Message[] messages = inbox.search(searchTerm);
-            if (messages == null || messages.length == 0) {
-                return null;
-            }
-            return Arrays.asList(messages);
-        } catch (Exception e) {
-            log.error("mkt searchMessages error, folder={}", folderName, e);
-            return null;
-        }
-    }
-
-    private void loadMailContents(SmtpPopSettingsEntity account) {
-        List<MktMailMessageEntity> heads = mktMailMessageService.list(new LambdaQueryWrapper<MktMailMessageEntity>()
-                .eq(MktMailMessageEntity::getAccountId, account.getId())
-                .eq(MktMailMessageEntity::getIsOnlyHead, 1)
-                .eq(MktMailMessageEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
-                .in(MktMailMessageEntity::getFolder, MktMailFolderEnum.INBOX.code, MktMailFolderEnum.SENT.code));
-        if (CollectionUtils.isEmpty(heads)) {
-            return;
-        }
-        try {
-            MailProperties properties = new MailProperties(account);
-            MailConnection connection = MailConnectionUtil.receiveEmailsConnection(properties);
-            if (connection == null) {
-                return;
-            }
-            for (MktMailMessageEntity mail : heads) {
-                try {
-                    if (!StringUtils.hasText(mail.getMsgUid())) {
-                        continue;
-                    }
-                    String serverFolder = MktMailFolderEnum.INBOX.code.equals(mail.getFolder()) ? "INBOX" : "Sent Items";
-                    Message message = connection.getMessageByUid(mail.getMsgUid(), serverFolder);
-                    if (message == null) {
-                        continue;
-                    }
-                    String body = extractBody(message);
-                    LambdaUpdateWrapper<MktMailMessageEntity> update = new LambdaUpdateWrapper<>();
-                    update.eq(MktMailMessageEntity::getId, mail.getId());
-                    update.set(MktMailMessageEntity::getContent, body);
-                    update.set(MktMailMessageEntity::getIsOnlyHead, 0);
-                    mktMailMessageService.update(update);
-                    saveAttachments(mail, message, account);
-                } catch (Exception ex) {
-                    log.error("mkt load content error, mailId={}", mail.getId(), ex);
-                }
-            }
-        } catch (Exception e) {
-            log.error("mkt loadMailContents error, accountId={}", account.getId(), e);
-        }
-    }
-
-    private void bindReply(Long accountId, Message message, MktMailMessageEntity entity) {
-        try {
-            String inReplyTo = firstHeader(message, "In-Reply-To");
-            if (!StringUtils.hasText(inReplyTo)) {
-                String references = firstHeader(message, "References");
-                if (StringUtils.hasText(references)) {
-                    String[] parts = references.trim().split("\\s+");
-                    inReplyTo = parts[parts.length - 1];
-                }
-            }
-            if (!StringUtils.hasText(inReplyTo)) {
-                return;
-            }
-            String normalized = normalizeMessageId(inReplyTo);
-            entity.setInReplyTo(normalized);
-            MktMailMessageEntity parent = findByMessageId(accountId, normalized);
-            if (parent != null) {
-                entity.setReplyMsgId(parent.getId());
-                // 继承业务单号,便于营销侧回查
-                if (!StringUtils.hasText(entity.getBizRef()) && StringUtils.hasText(parent.getBizRef())) {
-                    entity.setBizRef(parent.getBizRef());
-                }
-                if (!StringUtils.hasText(entity.getBizBatchNo()) && StringUtils.hasText(parent.getBizBatchNo())) {
-                    entity.setBizBatchNo(parent.getBizBatchNo());
-                }
-            }
-        } catch (Exception e) {
-            log.error("mkt bindReply error", e);
-        }
-    }
-
-    private MktMailMessageEntity findByMessageId(Long accountId, String messageId) {
-        if (!StringUtils.hasText(messageId)) {
-            return null;
-        }
-        String normalized = normalizeMessageId(messageId);
-        List<MktMailMessageEntity> list = mktMailMessageService.list(new LambdaQueryWrapper<MktMailMessageEntity>()
-                .eq(MktMailMessageEntity::getAccountId, accountId)
-                .eq(MktMailMessageEntity::getIsDelete, CommonConstant.DEL_FLAG_0)
-                .and(w -> w.eq(MktMailMessageEntity::getMessageId, normalized)
-                        .or().eq(MktMailMessageEntity::getMessageId, "<" + stripBrackets(normalized) + ">")
-                        .or().eq(MktMailMessageEntity::getMessageId, stripBrackets(normalized)))
-                .orderByDesc(MktMailMessageEntity::getId)
-                .last("limit 1"));
-        return CollectionUtils.isEmpty(list) ? null : list.get(0);
-    }
-
-    private boolean existsByMessageId(Long accountId, String messageId) {
-        return findByMessageId(accountId, messageId) != null;
-    }
-
-    private List<String> listRecentMessageIds(Long accountId, String folder) {
-        QueryWrapper<MktMailMessageEntity> qw = new QueryWrapper<>();
-        qw.select("message_id");
-        qw.eq("folder", folder);
-        qw.eq("account_id", accountId);
-        qw.isNotNull("message_id");
-        qw.last("order by recipient_date desc limit 200");
-        List<Object> objs = mktMailMessageService.listObjs(qw);
-        if (CollectionUtils.isEmpty(objs)) {
-            return new ArrayList<>();
-        }
-        return objs.stream().filter(Objects::nonNull).map(o -> normalizeMessageId(String.valueOf(o))).collect(Collectors.toList());
-    }
-
-    private MktMailMessageEntity convertMessage(Message message, MktMailFolderEnum folderEnum) {
-        try {
-            MktMailMessageEntity entity = new MktMailMessageEntity();
-            String subject = message.getSubject() == null ? "" : MimeUtility.decodeText(message.getSubject());
-            entity.setSubject(subject);
-            entity.setFolder(folderEnum.code);
-            entity.setInOutMark(folderEnum == MktMailFolderEnum.INBOX ? 1 : 2);
-
-            Map<String, String> fromMap = decodeAddresses(message.getFrom());
-            entity.setFromAddr(joinKeys(fromMap));
-            entity.setFromName(joinValues(fromMap));
-
-            Map<String, String> toMap = decodeAddresses(message.getRecipients(Message.RecipientType.TO));
-            String recipient = joinKeys(toMap);
-            entity.setRecipient(StringUtils.hasText(recipient) ? recipient : "");
-
-            Map<String, String> ccMap = decodeAddresses(message.getRecipients(Message.RecipientType.CC));
-            entity.setRecipientCc(joinKeys(ccMap));
-
-            entity.setEmailSize((long) Math.max(message.getSize(), 0));
-            if (message.getSentDate() != null) {
-                entity.setSentDate(message.getSentDate());
-                entity.setSentTime(message.getSentDate());
-            }
-            if (message.getReceivedDate() != null) {
-                entity.setRecipientDate(message.getReceivedDate());
-            } else if (message.getSentDate() != null) {
-                entity.setRecipientDate(message.getSentDate());
-            } else {
-                entity.setRecipientDate(new Date());
-            }
-            entity.setIsRead(message.isSet(Flags.Flag.SEEN) ? 1 : 0);
-
-            try {
-                Object content = message.getContent();
-                if (content instanceof String) {
-                    entity.setIsOnlyHead(0);
-                    entity.setContent(content.toString());
-                } else {
-                    entity.setIsOnlyHead(1);
-                }
-            } catch (Exception e) {
-                entity.setIsOnlyHead(1);
-            }
-            return entity;
-        } catch (Exception e) {
-            log.error("mkt convertMessage error", e);
-            return null;
-        }
-    }
-
-    private void fillMsgUid(Message message, SmtpPopSettingsEntity account, MktMailMessageEntity entity) {
-        try {
-            if ("IMAP".equalsIgnoreCase(account.getProtocolType()) && message.getFolder() instanceof IMAPFolder) {
-                long uid = ((IMAPFolder) message.getFolder()).getUID(message);
-                entity.setMsgUid(String.valueOf(uid));
-            } else if ("POP3".equalsIgnoreCase(account.getProtocolType()) && message.getFolder() instanceof POP3Folder) {
-                entity.setMsgUid(((POP3Folder) message.getFolder()).getUID(message));
-            }
-        } catch (Exception e) {
-            log.error("mkt fillMsgUid error", e);
-        }
-    }
-
-    private void saveAttachments(MktMailMessageEntity mail, Message message, SmtpPopSettingsEntity account) {
-        try {
-            int exist = mktMailAttachmentService.count(new LambdaQueryWrapper<MktMailAttachmentEntity>()
-                    .eq(MktMailAttachmentEntity::getMailId, mail.getId())
-                    .eq(MktMailAttachmentEntity::getIsDelete, CommonConstant.DEL_FLAG_0));
-            if (exist > 0) {
-                return;
-            }
-            Object content = message.getContent();
-            if (!(content instanceof Multipart)) {
-                return;
-            }
-            String day = DateUtils.date2Str(
-                    mail.getRecipientDate() == null ? new Date() : mail.getRecipientDate(), DateUtils.yyyyMMdd);
-            String downloadDir = mailFileProperties.getPath().getPath() + File.separator
-                    + account.getEmailAddress() + File.separator + day + File.separator;
-            Multipart multipart = (Multipart) content;
-            for (int i = 0; i < multipart.getCount(); i++) {
-                BodyPart bodyPart = multipart.getBodyPart(i);
-                if (Part.ATTACHMENT.equalsIgnoreCase(bodyPart.getDisposition())
-                        || bodyPart.isMimeType("application/octet-stream")) {
-                    saveOneAttachment(bodyPart, downloadDir, mail, account);
-                }
-            }
-        } catch (Exception e) {
-            log.error("mkt saveAttachments error, mailId={}", mail.getId(), e);
-        }
-    }
-
-    private void saveOneAttachment(BodyPart bodyPart, String downloadDir, MktMailMessageEntity mail, SmtpPopSettingsEntity account) {
-        try {
-            String fileName = bodyPart.getFileName();
-            if (!StringUtils.hasText(fileName)) {
-                return;
-            }
-            fileName = MimeUtility.decodeText(fileName).replace("\n", "").replace("\r", "");
-            String code = RandomGenerateHelper.generateRandomString(6);
-            if (!new File(downloadDir).exists()) {
-                new File(downloadDir).mkdirs();
-            }
-            String storedName = fileName + "." + code;
-            File file = new File(downloadDir + File.separator + storedName);
-            try (FileOutputStream output = new FileOutputStream(file)) {
-                bodyPart.getInputStream().transferTo(output);
-            }
-            MktMailAttachmentEntity att = new MktMailAttachmentEntity();
-            att.setMailId(mail.getId());
-            att.setAccountId(account.getId());
-            att.setFileCode(code);
-            att.setFileName(FilenameUtils.getBaseName(fileName));
-            att.setFileExt(FilenameUtils.getExtension(fileName));
-            att.setFilePath(file.getAbsolutePath().replace(mailFileProperties.getPath().getPath(), ""));
-            att.setFileSize(file.length());
-            att.setOwnerBy(mail.getOwnerBy());
-            att.setCreateBy(mail.getOwnerBy());
-            att.setIsDelete(CommonConstant.DEL_FLAG_0);
-            att.setEnabled(Boolean.TRUE);
-            mktMailAttachmentService.save(att);
-        } catch (Exception e) {
-            log.error("mkt saveOneAttachment error", e);
-        }
-    }
-
-    private static String extractBody(Message message) {
-        try {
-            Object content = message.getContent();
-            if (content instanceof String) {
-                return content.toString();
-            }
-            if (content instanceof Multipart) {
-                return parseMultipart((Multipart) content);
-            }
-        } catch (Exception e) {
-            log.error("mkt extractBody error", e);
-        }
-        return "";
-    }
-
-    private static String parseMultipart(Multipart multipart) throws Exception {
-        StringBuilder html = new StringBuilder();
-        StringBuilder plain = new StringBuilder();
-        for (int i = 0; i < multipart.getCount(); i++) {
-            BodyPart part = multipart.getBodyPart(i);
-            Object content = part.getContent();
-            if (content instanceof Multipart) {
-                String nested = parseMultipart((Multipart) content);
-                html.append(nested);
-            } else if (part.isMimeType("text/html")) {
-                html.append(content == null ? "" : content.toString());
-            } else if (part.isMimeType("text/plain")) {
-                plain.append(content == null ? "" : content.toString());
-            }
-        }
-        return html.length() > 0 ? html.toString() : plain.toString();
-    }
-
-    private static Map<String, String> decodeAddresses(Address[] addresses) throws Exception {
-        Map<String, String> map = new HashMap<>();
-        if (addresses == null) {
-            return map;
-        }
-        for (Address address : addresses) {
-            InternetAddress internetAddress = (InternetAddress) address;
-            String personal = internetAddress.getPersonal();
-            if (personal != null) {
-                personal = MimeUtility.decodeText(personal);
-            }
-            map.put(internetAddress.getAddress(), personal != null ? personal : "N/A");
-        }
-        return map;
-    }
-
-    private static String joinKeys(Map<String, String> map) {
-        if (map == null || map.isEmpty()) {
-            return "";
-        }
-        return String.join(",", map.keySet());
-    }
-
-    private static String joinValues(Map<String, String> map) {
-        if (map == null || map.isEmpty()) {
-            return "";
-        }
-        return map.values().stream()
-                .filter(v -> v != null && !"N/A".equals(v))
-                .collect(Collectors.joining(","));
-    }
-
-    private static String firstHeader(Message message, String name) throws Exception {
-        String[] values = message.getHeader(name);
-        if (values == null || values.length == 0) {
-            return null;
-        }
-        return values[0];
-    }
-
-    private static String normalizeMessageId(String messageId) {
-        if (!StringUtils.hasText(messageId)) {
-            return messageId;
-        }
-        return messageId.trim();
-    }
-
-    private static String stripBrackets(String messageId) {
-        if (!StringUtils.hasText(messageId)) {
-            return messageId;
-        }
-        String v = messageId.trim();
-        if (v.startsWith("<") && v.endsWith(">") && v.length() > 2) {
-            return v.substring(1, v.length() - 1);
-        }
-        return v;
+        mailReceiveEngine.receive(account, MailSceneEnum.MKT, mktMailReceivePersist, false);
     }
 }

+ 15 - 0
storlead-mail/storlead-mail-mkt/src/main/java/com/storlead/sales/mail/mkt/support/MktMailHeaders.java

@@ -0,0 +1,15 @@
+package com.storlead.sales.mail.mkt.support;
+
+/**
+ * 营销发信与收信对齐用的自定义邮件头:携带本地主键 id。
+ */
+public final class MktMailHeaders {
+
+    /**
+     * 发出 MIME 时写入,收信从服务器拉回 Sent/INBOX 时用此头匹配本地 mkt_mail_message.id。
+     */
+    public static final String LOCAL_MAIL_ID = "X-Storlead-Mkt-Mail-Id";
+
+    private MktMailHeaders() {
+    }
+}