"""招标邮件采集:IMAP 拉取 → 解析 → 落库 bid_tenders。 凭据来自 bid_mail_accounts 表(password 用 appPublic.aes AES 加密存储,明文不入库、不写死代码)。 阿里云企业邮:imap.qiye.aliyun.com:993 (SSL)。若开启「安全登录」,password 需填客户端专用密码。 IMAP 是阻塞库 → 统一放线程池执行(run_in_executor),不阻塞事件循环。 """ import asyncio import email import email.header import email.utils import imaplib import json import logging import re from datetime import datetime, timedelta from .bid_common import (get_db, new_id, rec_to_dict, to_int, TD_NEW) logger = logging.getLogger("pipeline.bidding.mail") # 默认招标关键词(bid_mail_accounts.filter_keywords 未配置时使用) DEFAULT_KEYWORDS = ["招标", "投标", "采购", "询价", "竞争性磋商", "邀标", "标书", "中标", "资格预审"] MAX_BODY = 60000 # 单封邮件正文入库上限 MAX_MAILS = 200 # 单次最多处理邮件数 def _decode_header(raw): """MIME 邮件头解码(中文主题/发件人)。""" if not raw: return "" try: parts = email.header.decode_header(raw) out = [] for txt, enc in parts: if isinstance(txt, bytes): out.append(txt.decode(enc or "utf-8", errors="ignore")) else: out.append(txt) return "".join(out).strip() except Exception: return str(raw) def _body_text(msg): """提取邮件正文纯文本(优先 text/plain,回退 text/html 去标签)。""" plain, html = "", "" if msg.is_multipart(): for part in msg.walk(): ctype = part.get_content_type() disp = str(part.get("Content-Disposition") or "") if "attachment" in disp: continue if ctype not in ("text/plain", "text/html"): continue try: payload = part.get_payload(decode=True) or b"" charset = part.get_content_charset() or "utf-8" txt = payload.decode(charset, errors="ignore") except Exception: continue if ctype == "text/plain": plain += txt + "\n" else: html += txt + "\n" else: try: payload = msg.get_payload(decode=True) or b"" charset = msg.get_content_charset() or "utf-8" txt = payload.decode(charset, errors="ignore") except Exception: txt = "" if (msg.get_content_type() or "") == "text/html": html = txt else: plain = txt body = plain.strip() if not body and html: body = re.sub(r"<(script|style)[^>]*>.*?", " ", html, flags=re.S | re.I) body = re.sub(r"<[^>]+>", " ", body) body = re.sub(r" ?", " ", body) body = re.sub(r"[ \t]+", " ", body) body = re.sub(r"\n{3,}", "\n\n", body) return body.strip()[:MAX_BODY] def _attachments(msg): """附件清单 [{name, size, ctype}](只记录元信息,文件由人工上传到 bid_tender_files)。""" out = [] if not msg.is_multipart(): return out for part in msg.walk(): disp = str(part.get("Content-Disposition") or "") fname = part.get_filename() if "attachment" not in disp and not fname: continue name = _decode_header(fname) if fname else "" try: size = len(part.get_payload(decode=True) or b"") except Exception: size = 0 out.append({"name": name, "size": size, "ctype": part.get_content_type()}) return out def _match_keywords(subject, body, keywords): text = (subject or "") + "\n" + (body or "") for k in keywords: k = (k or "").strip() if k and k in text: return k return "" def _fetch_sync(host, port, use_ssl, username, password, folder, days, keywords): """同步 IMAP 拉取(在线程池中执行)。返回 (mails, error)。""" mails = [] m = None try: if use_ssl: m = imaplib.IMAP4_SSL(host, port) else: m = imaplib.IMAP4(host, port) m.login(username, password) m.select(folder or "INBOX") since = (datetime.now() - timedelta(days=max(1, days))).strftime("%d-%b-%Y") typ, data = m.search(None, "(SINCE %s)" % since) if typ != "OK": return [], "IMAP search 失败: %s" % typ ids = (data[0] or b"").split() ids = ids[-MAX_MAILS:] for i in ids: typ, msg_data = m.fetch(i, "(RFC822)") if typ != "OK" or not msg_data: continue raw = None for chunk in msg_data: if isinstance(chunk, tuple) and len(chunk) > 1: raw = chunk[1] break if not raw: continue msg = email.message_from_bytes(raw) subject = _decode_header(msg.get("Subject")) body = _body_text(msg) hit = _match_keywords(subject, body, keywords) if not hit: continue recv = "" try: dt = email.utils.parsedate_to_datetime(msg.get("Date")) if dt: recv = dt.strftime("%Y-%m-%d %H:%M:%S") except Exception: recv = "" mails.append({ "uid": (msg.get("Message-ID") or "").strip()[:64] or ("seq-" + i.decode()), "seq": i.decode(), "from": _decode_header(msg.get("From"))[:250], "subject": subject[:480], "received_at": recv, "body": body, "attachments": _attachments(msg), "keyword": hit, }) return mails, "" except imaplib.IMAP4.error as e: return [], "IMAP 认证/协议错误: %s(阿里云企业邮开启安全登录后需使用客户端专用密码)" % str(e)[:200] except Exception as e: return [], "%s: %s" % (type(e).__name__, str(e)[:200]) finally: if m is not None: try: m.logout() except Exception: pass async def get_accounts(sor, org_id="", account_id=""): """取启用的邮箱配置(解密密码)。""" sql = "SELECT * FROM bid_mail_accounts WHERE enabled='1'" p = {} if account_id: sql += " AND id=${aid}$" p["aid"] = account_id if org_id: sql += " AND org_id=${oid}$" p["oid"] = org_id recs = await sor.sqlExe(sql, p) await sor.sqlExe("COMMIT", {}) out = [] for r in (recs or []): d = rec_to_dict(r) d["password_plain"] = _decrypt(d.get("password", "")) out.append(d) return out def _decrypt(enc): """DB 密码解密(AES;明文兼容:解密失败按明文用)。""" if not enc: return "" try: from appPublic.jsonConfig import getConfig from appPublic.aes import aes_decode_b64 key = getConfig().password_key if key: v = aes_decode_b64(key, enc) if v: return v except Exception: pass return enc def encrypt_password(plain): """页面/脚本写入密码时调用(AES 加密),解密失败场景保持明文兼容。""" if not plain: return "" try: from appPublic.jsonConfig import getConfig from appPublic.aes import aes_encode_b64 key = getConfig().password_key if key: return aes_encode_b64(key, plain) except Exception: pass return plain async def fetch_account_mails(account, days=None): """按单个邮箱配置拉取邮件(线程池执行 IMAP)。返回 (mails, error)。""" kws = [k.strip() for k in (account.get("filter_keywords") or "").split(",") if k.strip()] if not kws: kws = DEFAULT_KEYWORDS d = to_int(days if days is not None else account.get("fetch_days"), 7) loop = asyncio.get_event_loop() return await loop.run_in_executor( None, _fetch_sync, account.get("host") or "imap.qiye.aliyun.com", to_int(account.get("port"), 993), str(account.get("use_ssl") or "1") == "1", account.get("username") or "", account.get("password_plain") or "", account.get("folder") or "INBOX", d, kws) async def fetch_tenders_to_db(org_id="", account_id="", days=None): """采集入口:拉邮件 → 按 mail_uid 去重 → 落库 bid_tenders(status=new)。 返回 (ok, message)。message 含每个邮箱的新增/跳过数与错误原因(诚实报错,不静默)。 """ db, dbname = get_db() lines = [] total_new = 0 async with db.sqlorContext(dbname) as sor: accounts = await get_accounts(sor, org_id=org_id, account_id=account_id) if not accounts: return False, ("未配置启用的招标邮箱(bid_mail_accounts)。" "请在「招标邮箱配置」页面添加:host=imap.qiye.aliyun.com port=993 " "username=邮箱地址 password=客户端专用密码") for acc in accounts: mails, err = await fetch_account_mails(acc, days=days) if err: lines.append("[%s] 采集失败: %s" % (acc.get("account_name") or acc.get("username"), err)) continue new_n, skip_n = 0, 0 for mail in mails: exist = await sor.sqlExe( "SELECT id FROM bid_tenders WHERE mail_uid=${uid}$ AND account_id=${aid}$ LIMIT 1", {"uid": mail["uid"], "aid": acc.get("id") or ""}) if exist: skip_n += 1 continue data = { "id": new_id(), "org_id": acc.get("org_id") or "0", "source": "mail", "account_id": acc.get("id") or "", "mail_uid": mail["uid"], "mail_from": mail["from"], "mail_subject": mail["subject"], "raw_content": mail["body"], "attach_json": json.dumps(mail["attachments"], ensure_ascii=False), "title": mail["subject"][:250], "status": TD_NEW, "created_by": "agent.tender_scout", } if mail.get("received_at"): data["received_at"] = mail["received_at"] await sor.C("bid_tenders", data) new_n += 1 await sor.sqlExe( "UPDATE bid_mail_accounts SET last_fetch_at=NOW(), updated_at=NOW() " "WHERE id=${aid}$", {"aid": acc.get("id") or ""}) await sor.sqlExe("COMMIT", {}) total_new += new_n lines.append("[%s] 命中招标邮件 %d 封,新增 %d,已存在跳过 %d" % (acc.get("account_name") or acc.get("username"), len(mails), new_n, skip_n)) msg = "\n".join(lines) if total_new == 0 and all("采集失败" in x for x in lines): return False, msg return True, msg