Files
group_fqcd_jr/app/service/offsite_fund_service.py
T
lzf_0626 b5b068065f 修掉合并带入的 3 个 mypy 错(全在袁聪的场外邮件模块,各一行、不动逻辑)
1. offsite_document_recognition_adapter.py:788 —— 冗余 cast。
   `value in ("summary", ...)` 已经把类型收窄到那个字面量联合,cast 多余,删掉即可。
   该文件另有 13 处 cast,删这一个不影响 import。

2. offsite_fund_service.py:1203 —— dict 不变型。
   字面量里只有一个 bool,mypy 推断成 dict[str, bool | None],而 dict 是**不变型**,
   不是 dict[str, object] 的子类型。运行时本来就是合法值,加一行显式标注即可。

3. offsite_fund_service.py:1240 —— 标注宽于实际。
   `get_content()` 只定义在 EmailMessage 上,而调用方传的是
   BytesParser(policy=policy.default).parsebytes(...) 的返回值(本来就是 EmailMessage)。
   形参从基类 Message 改成 EmailMessage;import 同步替换(Message 仅此一处使用)。

三条都不影响运行:tests/unit/worker/test_offsite_mail_worker.py 的 6 个用例覆盖的正是
这条路径,修前修后都通过。修的意义在于——strict=true 下留着这 3 个错,mypy 在这个文件上
就失去价值,而 offsite_* 是刚合并进来、最需要类型检查兜底的新域。

验证:mypy 163 文件 0 错 / ruff 干净 / unit+contract 810 passed / integration 76 passed。

(另:本机 pytest 默认 basetemp 被权限占住的问题仍在,用重定向 TEMP 绕开,注意先建目录;
NL_develop 的 conftest 修复合并后可根治。)
2026-09-11 20:15:17 +08:00

2590 lines
113 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""场外基金申购赎回业务编排服务。"""
import asyncio
from collections import defaultdict
from collections.abc import Mapping, Sequence
from datetime import UTC, date, datetime
from decimal import Decimal
from email import policy
from email.header import decode_header, make_header
from email.message import EmailMessage
from email.parser import BytesParser
from email.utils import parsedate_to_datetime
from pathlib import Path
from typing import Literal, cast
from zoneinfo import ZoneInfo
from sqlalchemy import func, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.config import get_settings
from app.core.contracts import RequestContext
from app.core.offsite_fund_contracts import (
AttachmentFileContent,
Disposition,
DocumentType,
OffsiteDocumentSummary,
OperationDecision,
ReceiveRecognizedMailRequest,
RecognizedAttachment,
)
from app.model.audit import InteractionAudit
from app.model.offsite_fund import (
OffsiteExecutionPlanTask,
OffsiteFieldCorrection,
OffsiteFundAttachment,
OffsiteFundDocument,
OffsiteFundMail,
OffsiteMailCursor,
OffsiteNotification,
OffsiteQueryRecord,
OffsiteRecognitionAttempt,
OffsiteRuleResult,
)
from app.service.offsite_document_recognition_adapter import (
EXTRACTED_FIELD_NAMES,
OffsiteDocumentRecognitionAdapter,
RecognitionSourceFile,
StructuredRecognitionResult,
)
from app.service.offsite_fund_rules import (
OffsiteFundRuleEngine,
decimal_from,
normalize_amount_yuan,
parse_application_date,
)
from app.service.offsite_nl2sql_adapter import OffsiteNl2SqlAdapter
from app.service.offsite_smtp_adapter import (
OffsiteMailReplyRequest,
OffsiteSmtpSender,
SmtpAttachment,
SmtpSendResult,
)
SHANGHAI = ZoneInfo("Asia/Shanghai")
# 浏览器可以安全内联渲染的附件类型;其余类型一律走下载,避免渲染失败或引入 XSS。
INLINE_PREVIEW_MEDIA_TYPES = {
"application/pdf": "application/pdf",
"image/png": "image/png",
"image/jpeg": "image/jpeg",
"image/jpg": "image/jpeg",
"image/pjpeg": "image/jpeg",
"image/gif": "image/gif",
"image/webp": "image/webp",
"image/bmp": "image/bmp",
"image/x-ms-bmp": "image/bmp",
"image/avif": "image/avif",
}
# 入库的 media_type 可能不规范(例如统一写成 application/octet-stream),
# 此时按文件后缀兜底识别,保证 PDF 申购单/赎回单仍能直接预览。
PREVIEW_SUFFIX_MEDIA_TYPES = {
".pdf": "application/pdf",
".png": "image/png",
".jpg": "image/jpeg",
".jpeg": "image/jpeg",
".jpe": "image/jpeg",
".gif": "image/gif",
".webp": "image/webp",
".bmp": "image/bmp",
".avif": "image/avif",
}
DEFAULT_DOWNLOAD_MEDIA_TYPE = "application/octet-stream"
# NL2SQL 核对查询侧字段:与规则引擎 database_value 的键名保持一致。
NL2SQL_FIELD_NAMES: tuple[str, ...] = (
"最新净值",
"基金最新总份额",
"申请前持有份额",
"当前最新可用份额",
)
# 不同单据类型实际需要查询的字段:赎回单不查询净值与申请前持有份额。
NL2SQL_DOCUMENT_FIELDS: dict[str, tuple[str, ...]] = {
"subscription": ("最新净值", "基金最新总份额", "申请前持有份额"),
"redemption": ("基金最新总份额", "当前最新可用份额"),
}
# 查询被阻断时规则结果的固定标记,避免把非查询类规则误判成查询失败。
NL2SQL_BLOCKED_REASON = "依赖的 NL2SQL 查询失败或无可用数据"
# 人工修正字段的目标类型:区分 OCR 识别字段与 NL2SQL 查询字段。
CORRECTION_TARGET_RECOGNITION = "recognition"
CORRECTION_TARGET_NL2SQL = "nl2sql"
# 面板默认状态:识别全部无异常时为只读的"已保存",存在缺失或低置信字段时为"修改"。
FIELD_STATE_SAVED = "saved"
FIELD_STATE_EDITING = "editing"
class OffsiteFundService:
def __init__(
self,
session: AsyncSession,
smtp_sender: OffsiteSmtpSender | None = None,
recognizer: OffsiteDocumentRecognitionAdapter | None = None,
) -> None:
self.session = session
self.rules = OffsiteFundRuleEngine()
self.nl2sql = OffsiteNl2SqlAdapter()
self.smtp_sender = smtp_sender or OffsiteSmtpSender(get_settings())
self.recognizer = recognizer
async def list_mails(
self,
context: RequestContext,
page: int = 1,
page_size: int = 20,
sender: str | None = None,
status: str | None = None,
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
filters = []
if sender:
filters.append(OffsiteFundMail.sender == sender)
if status:
filters.append(OffsiteFundMail.status == status)
total = await self.session.scalar(
select(func.count(OffsiteFundMail.id)).where(*filters)
) or 0
mails = (
await self.session.execute(
select(OffsiteFundMail)
.where(*filters)
.order_by(OffsiteFundMail.created_at.desc(), OffsiteFundMail.id.desc())
.offset((page - 1) * page_size)
.limit(page_size)
)
).scalars().all()
mail_ids = [mail.mail_id for mail in mails]
attachments = await self._mail_attachments(mail_ids)
documents = await self._mail_documents(mail_ids)
return {
"code": 0,
"message": "ok",
"data": {
"items": [
self._mail_summary(
mail,
attachments.get(mail.mail_id, ()),
documents.get(mail.mail_id, ()),
)
for mail in mails
],
"page": page,
"page_size": page_size,
"total": int(total),
},
}
async def get_mail(
self, mail_id: str, context: RequestContext
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
async with self.session.begin():
mail = await self.session.scalar(
select(OffsiteFundMail).where(OffsiteFundMail.mail_id == mail_id)
)
if mail is None:
return {"code": 404, "message": "邮件不存在", "data": {}}
attachments = await self._mail_attachments((mail.mail_id,))
documents = await self._mail_documents((mail.mail_id,))
detail = self._mail_detail(
mail,
attachments.get(mail.mail_id, ()),
documents.get(mail.mail_id, ()),
)
self._add_audit(context, "offsite.mail_viewed", {"mail_id": mail_id})
return {"code": 0, "message": "ok", "data": detail}
async def mail_recognition_fields(
self, mail_id: str, context: RequestContext
) -> dict[str, object]:
"""查询一封邮件下全部附件的 OCR 识别字段。
只读接口:返回识别原文、结构化字段、置信度和最近一次识别尝试,
不触发重新识别,也不改动附件、单据或规则结果。
缺失字段和低置信字段只记录在识别尝试表,因此按附件取最近一次尝试补齐。
"""
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
async with self.session.begin():
mail = await self.session.scalar(
select(OffsiteFundMail).where(OffsiteFundMail.mail_id == mail_id)
)
if mail is None:
return {"code": 404, "message": "邮件不存在", "data": {}}
attachments = (await self._mail_attachments((mail.mail_id,))).get(
mail.mail_id, ()
)
documents = (await self._mail_documents((mail.mail_id,))).get(
mail.mail_id, ()
)
attempts = await self._latest_recognition_attempts(
tuple(item.attachment_id for item in attachments)
)
corrections = await self._latest_recognition_corrections(
tuple(item.attachment_id for item in attachments)
)
payload = self._recognition_payload(
mail, attachments, documents, attempts, corrections
)
self._add_audit(
context, "offsite.mail_recognition_viewed", {"mail_id": mail_id}
)
return {"code": 0, "message": "ok", "data": payload}
async def _latest_recognition_attempts(
self, attachment_ids: Sequence[str]
) -> dict[str, OffsiteRecognitionAttempt]:
"""取每个附件最近一次识别尝试,按尝试序号升序覆盖为最新一条。"""
if not attachment_ids:
return {}
rows = (
await self.session.execute(
select(OffsiteRecognitionAttempt)
.where(OffsiteRecognitionAttempt.attachment_id.in_(attachment_ids))
.order_by(
OffsiteRecognitionAttempt.attempt_no.asc(),
OffsiteRecognitionAttempt.id.asc(),
)
)
).scalars().all()
latest: dict[str, OffsiteRecognitionAttempt] = {}
for row in rows:
if row.attachment_id:
latest[row.attachment_id] = row
return latest
async def _latest_recognition_corrections(
self, attachment_ids: Sequence[str]
) -> dict[str, OffsiteFieldCorrection]:
"""取每个附件最近一次人工修正,按时间与主键升序覆盖成最新一条。"""
if not attachment_ids:
return {}
rows = (
await self.session.execute(
select(OffsiteFieldCorrection)
.where(
OffsiteFieldCorrection.target_type
== CORRECTION_TARGET_RECOGNITION,
OffsiteFieldCorrection.attachment_id.in_(attachment_ids),
)
.order_by(
OffsiteFieldCorrection.created_at.asc(),
OffsiteFieldCorrection.id.asc(),
)
)
).scalars().all()
latest: dict[str, OffsiteFieldCorrection] = {}
for row in rows:
if row.attachment_id:
latest[row.attachment_id] = row
return latest
async def _latest_nl2sql_correction(
self, task_id: str
) -> OffsiteFieldCorrection | None:
"""取单据最近一次 NL2SQL 字段人工修正;多条时以最新一条为准。"""
row: OffsiteFieldCorrection | None = await self.session.scalar(
select(OffsiteFieldCorrection)
.where(
OffsiteFieldCorrection.target_type == CORRECTION_TARGET_NL2SQL,
OffsiteFieldCorrection.task_id == task_id,
)
.order_by(
OffsiteFieldCorrection.created_at.desc(),
OffsiteFieldCorrection.id.desc(),
)
.limit(1)
)
return row
async def _effective_recognition_fields(
self, attachment: OffsiteFundAttachment
) -> dict[str, object]:
"""人工修正值优先的有效识别字段;附件上的原始识别值保持落库不变。"""
corrections = await self._latest_recognition_corrections(
(attachment.attachment_id,)
)
correction = corrections.get(attachment.attachment_id)
correction_fields = (
dict(correction.corrected_fields) if correction is not None else {}
)
return self._merge_corrections(attachment.extracted_fields, correction_fields)
@staticmethod
def _merge_corrections(
original: Mapping[str, object], corrections: Mapping[str, object]
) -> dict[str, object]:
"""把人工修正值叠加到原值上,空值与未修正字段一律回落到原值。"""
fields = dict(original)
for name, value in corrections.items():
if value is None or str(value).strip() == "":
continue
fields[name] = value
return fields
@staticmethod
def _cleaned_corrections(
submitted: Mapping[str, object],
allowed: Sequence[str],
original: Mapping[str, object],
) -> dict[str, str]:
"""只保留白名单内、非空且确实与原值不同的修正值,过滤未知字段。"""
cleaned: dict[str, str] = {}
for name in allowed:
if name not in submitted:
continue
value = submitted.get(name)
text = "" if value is None else str(value).strip()
if not text:
continue
if str(original.get(name) or "").strip() == text:
continue
cleaned[name] = text
return cleaned
@classmethod
def _correction_payload(
cls, correction: OffsiteFieldCorrection | None
) -> dict[str, object] | None:
if correction is None:
return None
return {
"operator_id": correction.operator_id,
"fields": dict(correction.corrected_fields),
"changed_fields": list(correction.changed_fields),
"corrected_at": cls._iso_datetime(correction.created_at),
}
@classmethod
def _recognition_payload(
cls,
mail: OffsiteFundMail,
attachments: Sequence[OffsiteFundAttachment],
documents: Sequence[OffsiteFundDocument],
attempts: Mapping[str, OffsiteRecognitionAttempt],
corrections: Mapping[str, OffsiteFieldCorrection],
) -> dict[str, object]:
documents_by_attachment: dict[str, list[OffsiteFundDocument]] = defaultdict(list)
for document in documents:
documents_by_attachment[document.attachment_id].append(document)
items = []
for attachment in attachments:
attempt = attempts.get(attachment.attachment_id)
correction = corrections.get(attachment.attachment_id)
correction_fields = (
dict(correction.corrected_fields) if correction is not None else {}
)
missing_fields = list(attempt.missing_fields) if attempt else []
low_confidence_fields = (
list(attempt.low_confidence_fields) if attempt else []
)
items.append({
"attachment_id": attachment.attachment_id,
"filename": attachment.filename,
"document_type": attachment.document_type,
"media_type": attachment.media_type,
"size_bytes": attachment.size_bytes,
"status": attachment.status,
"ocr_text": attachment.ocr_text,
"extracted_fields": attachment.extracted_fields,
"effective_fields": cls._merge_corrections(
attachment.extracted_fields, correction_fields
),
"corrections": cls._correction_payload(correction),
"has_correction": correction is not None,
# 有缺失字段或低置信字段时默认进入修改态,否则默认只读的已保存态。
"default_state": (
FIELD_STATE_EDITING
if missing_fields or low_confidence_fields
else FIELD_STATE_SAVED
),
"field_confidence": attachment.field_confidence,
"page_evidence": attachment.page_evidence,
"missing_fields": missing_fields,
"low_confidence_fields": low_confidence_fields,
"ocr_status": attempt.ocr_status if attempt else None,
"llm_status": attempt.llm_status if attempt else None,
"latest_attempt": (
{
"attempt_no": attempt.attempt_no,
"source": attempt.source,
"status": attempt.status,
"started_at": cls._iso_datetime(attempt.started_at),
"finished_at": cls._iso_datetime(attempt.finished_at),
"error_message": attempt.error_message,
}
if attempt
else None
),
"documents": [
{
"task_id": document.task_id,
"document_type": document.document_type,
"status": document.status,
"operator_decision": document.operator_decision,
"fund_code": document.fund_code,
"fund_name": document.fund_name,
"application_no": document.application_no,
"application_date": (
document.application_date.isoformat()
if document.application_date
else None
),
"raw_application_date": document.raw_application_date,
"agency": document.agency,
"subscription_amount_yuan": (
str(document.subscription_amount_yuan)
if document.subscription_amount_yuan is not None
else None
),
"redemption_shares": (
str(document.redemption_shares)
if document.redemption_shares is not None
else None
),
}
for document in documents_by_attachment.get(
attachment.attachment_id, []
)
],
})
return {
"mail_id": mail.mail_id,
"sender": mail.sender,
"status": mail.status,
"received_date": mail.received_date.isoformat(),
"attachments": items,
}
async def nl2sql_fields(
self, task_id: str, context: RequestContext
) -> dict[str, object]:
"""查询单据核对使用的 NL2SQL 返回字段。
只读接口:返回最近一次核对落库的查询侧字段(最新净值、基金最新总份额、
申请前持有份额、当前最新可用份额),不触发新的 NL2SQL 调用,也不改写规则结果。
空值必须区分"当前单据类型不查询"、"查询失败"和"尚未核对",否则运营会误判数据源。
"""
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
async with self.session.begin():
document = await self.session.scalar(
select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id
)
)
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
rule_results = (
await self.session.execute(
select(OffsiteRuleResult).where(
OffsiteRuleResult.task_id == task_id
)
)
).scalars().all()
query_records = (
await self.session.execute(
select(OffsiteQueryRecord)
.where(OffsiteQueryRecord.task_id == task_id)
.order_by(OffsiteQueryRecord.id.asc())
)
).scalars().all()
correction = await self._latest_nl2sql_correction(task_id)
payload = self._nl2sql_fields_payload(
document, rule_results, query_records, correction
)
self._add_audit(
context, "offsite.nl2sql_fields_viewed", {"task_id": task_id}
)
return {"code": 0, "message": "ok", "data": payload}
@classmethod
def _nl2sql_fields_payload(
cls,
document: OffsiteFundDocument,
rule_results: Sequence[OffsiteRuleResult],
query_records: Sequence[OffsiteQueryRecord],
correction: OffsiteFieldCorrection | None = None,
) -> dict[str, object]:
expected = NL2SQL_DOCUMENT_FIELDS.get(document.document_type, ())
values = cls._nl2sql_original_fields(rule_results)
correction_fields = (
dict(correction.corrected_fields) if correction is not None else {}
)
query_failed = any(
record.status == "query_failed" for record in query_records
) or any(
isinstance(result.calculation, Mapping)
and result.calculation.get("原因") == NL2SQL_BLOCKED_REASON
for result in rule_results
)
fields: dict[str, str | None] = {}
field_status: dict[str, str] = {}
missing_fields: list[str] = []
for field_name in NL2SQL_FIELD_NAMES:
if field_name in values:
fields[field_name] = values[field_name]
field_status[field_name] = "success"
continue
fields[field_name] = None
missing_fields.append(field_name)
if field_name not in expected:
field_status[field_name] = "not_queried"
elif query_failed:
field_status[field_name] = "query_failed"
else:
field_status[field_name] = "pending"
for field_name, value in correction_fields.items():
text = "" if value is None else str(value).strip()
if not text:
continue
field_status[field_name] = "corrected"
if field_name in missing_fields:
missing_fields.remove(field_name)
return {
"task_id": document.task_id,
"mail_id": document.mail_id,
"document_type": document.document_type,
"fund_code": document.fund_code,
"account_identifier": document.account_identifier,
"fields": fields,
"effective_fields": cls._merge_corrections(fields, correction_fields),
"corrections": cls._correction_payload(correction),
"field_status": field_status,
"missing_fields": missing_fields,
# NL2SQL 侧允许人工修正,面板默认进入可修改态。
"default_state": FIELD_STATE_EDITING,
"queries": [
{
"rule_code": record.rule_code,
"status": record.status,
"row_count": cls._query_row_count(record.result_summary),
"queried_at": cls._iso_datetime(record.created_at),
"error_message": record.error_message,
}
for record in query_records
],
"updated_at": cls._iso_datetime(document.updated_at),
}
@classmethod
def _nl2sql_original_fields(
cls, rule_results: Sequence[OffsiteRuleResult]
) -> dict[str, str]:
"""按既有口径聚合 NL2SQL 查询侧原始字段:同名字段先到先得。"""
values: dict[str, str] = {}
for result in rule_results:
database_value = result.database_value or {}
for field_name in NL2SQL_FIELD_NAMES:
value = database_value.get(field_name)
if value is None or value == "" or field_name in values:
continue
values[field_name] = str(value)
return values
@staticmethod
def _query_row_count(result_summary: object) -> int:
"""从落库的 NL2SQL 摘要里取返回行数;结构异常时按 0 行处理。"""
if not isinstance(result_summary, Mapping):
return 0
data = result_summary.get("data")
if not isinstance(data, Mapping):
return 0
rows = data.get("rows")
if isinstance(rows, list):
return len(rows)
total = data.get("total")
return int(total) if isinstance(total, (int, float)) else 0
async def rule_results(
self, task_id: str, context: RequestContext
) -> dict[str, object]:
"""查询单据的规则判定结果。
只读接口:返回最近一次"计算 + 核对"落库的结论、单据侧值、查询侧值和计算过程,
不触发 NL2SQL 查询,也不重新判定。未核对过的单据返回空规则列表,
调用方需要先执行核对或重新判定。
"""
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
async with self.session.begin():
document = await self.session.scalar(
select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id
)
)
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
results = await self._document_rule_results(task_id)
payload = self._rule_results_payload(document, results)
self._add_audit(
context, "offsite.rule_results_viewed", {"task_id": task_id}
)
return {"code": 0, "message": "ok", "data": payload}
async def recalculate_rule_results(
self, task_id: str, operator_id: str, context: RequestContext
) -> dict[str, object]:
"""重新判定单据规则。
重新判定只重跑"计算 + 核对"两个阶段:输入是已落库的识别字段和已经查到的
NL2SQL 结果,不重新识别附件、不重新调用 NL2SQL,因此不会产生新的外部调用。
需要换一批查询数据时,先调用核对触发接口刷新查询记录,再调用本接口重新判定。
查询缺失或失败的规则按"无法判断"处理并保留原因,不做任何猜测。
"""
denied = self._permission_error(
context, ("offsite:write", "offsite:confirm", "offsite:nl2sql")
)
if denied is not None:
return denied
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
document = await self.session.scalar(
select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id
)
)
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
attachment = await self.session.scalar(
select(OffsiteFundAttachment).where(
OffsiteFundAttachment.attachment_id == document.attachment_id
)
)
if attachment is None:
return {"code": 404, "message": "原始附件不存在", "data": {}}
# 人工修正值优先:修正过的识别字段按修正后口径重算,原始识别值不变。
fields = await self._effective_recognition_fields(attachment)
blocked_rules, query_snapshot = await self._replay_query_values(
document, fields
)
await self._refresh_rule_results(
document, fields, now, blocked_rule_codes=blocked_rules
)
await self._finish_execution_plan(task_id, fields, blocked_rules, now)
# 只在规则执行态之间流转:识别异常、识别复核和人工已确认属于其它阶段的状态,
# 重新判定不得把它们覆盖掉,否则会把识别结果丢失伪装成查询失败。
if document.status in {"planned", "query_failed"}:
document.status = "query_failed" if blocked_rules else "planned"
document.updated_at = now
results = await self._document_rule_results(task_id)
payload = self._rule_results_payload(document, results)
payload["query_snapshot"] = query_snapshot
self._add_audit(context, "offsite.rule_results_recalculated", {
"task_id": task_id,
"operator_id": operator_id,
"blocked_rule_count": len(blocked_rules),
})
return {"code": 0, "message": "重新判定完成", "data": payload}
async def save_recognition_corrections(
self,
mail_id: str,
submitted: Sequence[tuple[str, dict[str, str]]],
operator_id: str,
context: RequestContext,
) -> dict[str, object]:
"""保存邮件内附件的 OCR 识别字段人工修正值。
只新增修正记录,不覆盖附件上的 Agent 原始识别值;单据的标准化展示字段
跟随修正后的有效值刷新,重新判定再按有效值重跑计算与核对。
"""
denied = self._permission_error(context, ("offsite:write",))
if denied is not None:
return denied
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
mail = await self.session.scalar(
select(OffsiteFundMail).where(OffsiteFundMail.mail_id == mail_id)
)
if mail is None:
return {"code": 404, "message": "邮件不存在", "data": {}}
attachments = (await self._mail_attachments((mail.mail_id,))).get(
mail.mail_id, ()
)
documents = (await self._mail_documents((mail.mail_id,))).get(
mail.mail_id, ()
)
by_attachment = {item.attachment_id: item for item in attachments}
documents_by_attachment: dict[str, list[OffsiteFundDocument]] = defaultdict(
list
)
for document in documents:
documents_by_attachment[document.attachment_id].append(document)
# 同一附件重复提交时以最后一次为准,避免一次请求写出多条相互覆盖的记录。
merged: dict[str, dict[str, str]] = {}
for attachment_id, fields in submitted:
merged[attachment_id] = dict(fields)
if not merged:
return {"code": 422, "message": "请至少提交一个附件的修正字段", "data": {}}
saved: list[str] = []
for attachment_id, fields in merged.items():
attachment = by_attachment.get(attachment_id)
if attachment is None:
return {
"code": 404,
"message": f"附件 {attachment_id} 不属于该邮件",
"data": {},
}
cleaned = self._cleaned_corrections(
fields, EXTRACTED_FIELD_NAMES, attachment.extracted_fields
)
self.session.add(
OffsiteFieldCorrection(
target_type=CORRECTION_TARGET_RECOGNITION,
task_id=None,
mail_id=mail.mail_id,
attachment_id=attachment_id,
operator_id=operator_id,
original_fields=dict(attachment.extracted_fields),
corrected_fields=cleaned,
changed_fields=sorted(cleaned),
created_at=now,
)
)
effective = self._merge_corrections(
attachment.extracted_fields, cleaned
)
for document in documents_by_attachment.get(attachment_id, []):
self._apply_recognition_fields(document, effective)
document.updated_at = now
saved.append(attachment_id)
attempt_map = await self._latest_recognition_attempts(tuple(by_attachment))
correction_map = await self._latest_recognition_corrections(
tuple(by_attachment)
)
payload = self._recognition_payload(
mail, attachments, documents, attempt_map, correction_map
)
self._add_audit(context, "offsite.mail_recognition_corrected", {
"mail_id": mail_id,
"operator_id": operator_id,
"attachment_ids": sorted(saved),
})
return {"code": 0, "message": "识别字段修正已保存", "data": payload}
async def save_nl2sql_corrections(
self,
task_id: str,
fields: Mapping[str, str],
operator_id: str,
context: RequestContext,
) -> dict[str, object]:
"""保存单据 NL2SQL 返回字段的人工修正值。
只新增修正记录,不改写查询记录与查询摘要;重新判定时修正值优先于查询原值。
"""
denied = self._permission_error(context, ("offsite:write",))
if denied is not None:
return denied
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
document = await self.session.scalar(
select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id
)
)
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
rule_results = (
await self.session.execute(
select(OffsiteRuleResult).where(
OffsiteRuleResult.task_id == task_id
)
)
).scalars().all()
query_records = (
await self.session.execute(
select(OffsiteQueryRecord)
.where(OffsiteQueryRecord.task_id == task_id)
.order_by(OffsiteQueryRecord.id.asc())
)
).scalars().all()
cleaned = self._cleaned_corrections(
fields, NL2SQL_FIELD_NAMES, self._nl2sql_original_fields(rule_results)
)
self.session.add(
OffsiteFieldCorrection(
target_type=CORRECTION_TARGET_NL2SQL,
task_id=task_id,
mail_id=document.mail_id,
attachment_id=None,
operator_id=operator_id,
original_fields=self._nl2sql_original_fields(rule_results),
corrected_fields=cleaned,
changed_fields=sorted(cleaned),
created_at=now,
)
)
correction = await self._latest_nl2sql_correction(task_id)
payload = self._nl2sql_fields_payload(
document, rule_results, query_records, correction
)
self._add_audit(context, "offsite.nl2sql_fields_corrected", {
"task_id": task_id,
"operator_id": operator_id,
"changed_fields": sorted(cleaned),
})
return {"code": 0, "message": "NL2SQL 字段修正已保存", "data": payload}
async def _document_rule_results(
self, task_id: str
) -> list[OffsiteRuleResult]:
return list((
await self.session.execute(
select(OffsiteRuleResult)
.where(OffsiteRuleResult.task_id == task_id)
.order_by(OffsiteRuleResult.id.asc())
)
).scalars().all())
async def _replay_query_values(
self, document: OffsiteFundDocument, fields: dict[str, object]
) -> tuple[set[str], list[dict[str, object]]]:
"""按已落库查询记录重放查询阶段结果,返回被阻断规则和查询快照。
人工修正过的 NL2SQL 字段优先于查询原值;修正补齐字段后规则不再阻断。
"""
nl2sql_corrections = await self._latest_nl2sql_correction(document.task_id)
corrected_fields = (
dict(nl2sql_corrections.corrected_fields)
if nl2sql_corrections is not None
else {}
)
blocked_rules: set[str] = set()
snapshot: list[dict[str, object]] = []
for rule_code, _question in self._nl2sql_questions(document):
record = await self.session.scalar(
select(OffsiteQueryRecord)
.where(
OffsiteQueryRecord.task_id == document.task_id,
OffsiteQueryRecord.rule_code == rule_code,
)
.order_by(OffsiteQueryRecord.id.desc())
.limit(1)
)
row = (
self._first_query_row(record.result_summary)
if record is not None and record.status == "success"
else None
)
values, _missing = self._query_values(rule_code, row)
for field_name in self._query_field_mapping(rule_code):
corrected = corrected_fields.get(field_name)
if corrected is None or str(corrected).strip() == "":
continue
values[field_name] = str(corrected)
# 修正值补齐后按规则实际需要的字段重新计算缺失,避免误判为查询失败。
missing = tuple(
name
for name in self._query_field_mapping(rule_code)
if name not in values
)
reason: str | None = None
if missing:
if record is None:
reason = "尚未执行该规则的 NL2SQL 查询"
elif record.status != "success":
reason = record.error_message or "NL2SQL 查询未成功"
else:
reason = f"查询结果缺少字段:{', '.join(missing)}"
if reason is None:
fields.update(values)
else:
blocked_rules.add(rule_code)
snapshot.append({
"rule_code": rule_code,
"status": "success" if reason is None else "query_failed",
"reason": reason,
"values": values,
})
return blocked_rules, snapshot
@classmethod
def _rule_results_payload(
cls,
document: OffsiteFundDocument,
results: Sequence[OffsiteRuleResult],
) -> dict[str, object]:
return {
"task_id": document.task_id,
"mail_id": document.mail_id,
"document_type": document.document_type,
"fund_code": document.fund_code,
"fund_name": document.fund_name,
"document_status": document.status,
"operator_decision": document.operator_decision,
"summary": {
"total": len(results),
"normal": sum(1 for item in results if item.result == "正常"),
"abnormal": sum(1 for item in results if item.result == "异常"),
"unknown": sum(1 for item in results if item.result == "无法判断"),
},
"rules": [
{
"rule_code": item.rule_code,
"rule_name": item.rule_name,
"result": item.result,
"document_value": item.document_value,
"database_value": item.database_value,
"calculation": item.calculation,
# 规则结果表不可覆盖重写:这里返回首次落库时间,结论刷新时间取外层 updated_at。
"created_at": cls._iso_datetime(item.created_at),
}
for item in results
],
"updated_at": cls._iso_datetime(document.updated_at),
}
async def open_attachment_file(
self, attachment_id: str, disposition: str, context: RequestContext
) -> dict[str, object] | AttachmentFileContent:
"""打开单个附件的原始文件,供浏览器内联预览或下载。
只读接口:不改动附件、单据与识别结果。文件路径始终限制在场外收件存储根
目录内,越界或文件缺失统一按"文件不存在"处理,既不泄露服务器真实路径,
也不把内部异常暴露给调用方。
"""
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
async with self.session.begin():
attachment = await self.session.scalar(
select(OffsiteFundAttachment).where(
OffsiteFundAttachment.attachment_id == attachment_id
)
)
if attachment is None:
return {"code": 404, "message": "附件不存在", "data": {}}
path = self._resolve_storage_path(attachment.original_file_path)
if path is None:
return {"code": 404, "message": "附件原始文件不存在", "data": {}}
filename = attachment.filename
mail_id = attachment.mail_id
media_type = self._preview_media_type(filename, attachment.media_type)
resolved_disposition = self._preview_disposition(media_type, disposition)
self._add_audit(context, "offsite.attachment_viewed", {
"attachment_id": attachment_id,
"mail_id": mail_id,
"disposition": resolved_disposition,
})
return AttachmentFileContent(
path=str(path),
filename=filename,
media_type=media_type,
disposition=resolved_disposition,
)
@staticmethod
def _resolve_storage_path(original_file_path: str) -> Path | None:
"""把附件路径限制在场外收件存储根目录内,阻断路径穿越与软链接逃逸。"""
if not original_file_path:
return None
root = Path(get_settings().offsite_mail_storage_dir).resolve()
try:
candidate = Path(original_file_path).resolve()
except OSError:
return None
if not candidate.is_relative_to(root) or not candidate.is_file():
return None
return candidate
@staticmethod
def _preview_media_type(filename: str, media_type: str) -> str:
"""确定预览用的媒体类型:优先可信白名单,其次按文件后缀兜底。"""
normalized = (media_type or "").strip().lower()
if normalized in INLINE_PREVIEW_MEDIA_TYPES:
return INLINE_PREVIEW_MEDIA_TYPES[normalized]
suffix = Path(filename or "").suffix.lower()
return PREVIEW_SUFFIX_MEDIA_TYPES.get(suffix, DEFAULT_DOWNLOAD_MEDIA_TYPE)
@staticmethod
def _preview_disposition(media_type: str, requested: str) -> Disposition:
"""只有浏览器能安全内联渲染的类型才允许 inline,其余一律下载。"""
if (requested or "").strip().lower() != "inline":
return "attachment"
if media_type in INLINE_PREVIEW_MEDIA_TYPES.values():
return "inline"
return "attachment"
async def mailbox_status(self, context: RequestContext) -> dict[str, object]:
"""收件箱运行状态(只读),供前端在游标锁死时告警。"""
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
settings = get_settings()
cursor = await self.session.scalar(
select(OffsiteMailCursor)
.where(OffsiteMailCursor.mailbox == settings.offsite_mailbox)
.order_by(OffsiteMailCursor.id.desc())
.limit(1)
)
monitoring = bool(
settings.offsite_imap_enabled and settings.offsite_mail_worker_enabled
)
if cursor is None:
return {
"code": 0,
"message": "ok",
"data": {
"mailbox": settings.offsite_mailbox,
"status": "uninitialized",
"blocked": False,
"last_uid": None,
"blocked_uid": None,
"blocked_message_id": None,
"retry_count": 0,
"last_error": None,
"updated_at": None,
"monitoring": monitoring,
"alert_message": None,
},
}
blocked = cursor.status == "blocked"
alert_message = None
if blocked:
alert_message = (
f"收件邮箱已停止收取邮件:uid {cursor.blocked_uid or '未知'} 这封邮件连续 "
f"{cursor.retry_count} 次识别失败后被挂起,处理时间 "
f"{self._iso_datetime(cursor.updated_at) or '未记录'}。"
f"原因:{cursor.last_error or '未记录'}。"
)
return {
"code": 0,
"message": "ok",
"data": {
"mailbox": cursor.mailbox,
"status": cursor.status,
"blocked": blocked,
"last_uid": cursor.last_uid,
"blocked_uid": cursor.blocked_uid,
"blocked_message_id": cursor.blocked_message_id,
"retry_count": cursor.retry_count,
"last_error": cursor.last_error,
"updated_at": self._iso_datetime(cursor.updated_at),
"monitoring": monitoring,
"alert_message": alert_message,
},
}
async def _mail_attachments(
self, mail_ids: Sequence[str]
) -> dict[str, tuple[OffsiteFundAttachment, ...]]:
if not mail_ids:
return {}
rows = (
await self.session.execute(
select(OffsiteFundAttachment)
.where(OffsiteFundAttachment.mail_id.in_(mail_ids))
.order_by(OffsiteFundAttachment.id.asc())
)
).scalars().all()
grouped: dict[str, list[OffsiteFundAttachment]] = defaultdict(list)
for row in rows:
grouped[row.mail_id].append(row)
return {mail_id: tuple(items) for mail_id, items in grouped.items()}
async def _mail_documents(
self, mail_ids: Sequence[str]
) -> dict[str, tuple[OffsiteFundDocument, ...]]:
if not mail_ids:
return {}
rows = (
await self.session.execute(
select(OffsiteFundDocument)
.where(OffsiteFundDocument.mail_id.in_(mail_ids))
.order_by(OffsiteFundDocument.id.asc())
)
).scalars().all()
grouped: dict[str, list[OffsiteFundDocument]] = defaultdict(list)
for row in rows:
grouped[row.mail_id].append(row)
return {mail_id: tuple(items) for mail_id, items in grouped.items()}
@classmethod
def _mail_summary(
cls,
mail: OffsiteFundMail,
attachments: Sequence[OffsiteFundAttachment],
documents: Sequence[OffsiteFundDocument],
) -> dict[str, object]:
parsed = cls._parse_eml(mail.original_eml_path)
return {
"mail_id": mail.mail_id,
"sender": mail.sender,
"return_path": mail.return_path,
"subject": parsed["subject"],
"sent_at": parsed["sent_at"],
"received_at": cls._iso_datetime(mail.created_at),
"received_date": mail.received_date.isoformat(),
"status": mail.status,
"has_body": parsed["has_body"],
"attachment_count": len(attachments),
"business_document_count": len(documents),
"attachment_names": [item.filename for item in attachments],
}
@classmethod
def _mail_detail(
cls,
mail: OffsiteFundMail,
attachments: Sequence[OffsiteFundAttachment],
documents: Sequence[OffsiteFundDocument],
) -> dict[str, object]:
parsed = cls._parse_eml(mail.original_eml_path)
documents_by_attachment: dict[str, list[OffsiteFundDocument]] = defaultdict(list)
for document in documents:
documents_by_attachment[document.attachment_id].append(document)
attachment_items = []
for attachment in attachments:
attachment_items.append({
"attachment_id": attachment.attachment_id,
"filename": attachment.filename,
"document_type": attachment.document_type,
"media_type": attachment.media_type,
"size_bytes": attachment.size_bytes,
"file_hash": attachment.file_hash,
"status": attachment.status,
"documents": [
{
"task_id": document.task_id,
"document_type": document.document_type,
"status": document.status,
"operator_decision": document.operator_decision,
}
for document in documents_by_attachment.get(
attachment.attachment_id, []
)
],
})
return {
"mail_id": mail.mail_id,
"imap_uid": mail.imap_uid,
"message_id": mail.message_id,
"sender": mail.sender,
"return_path": mail.return_path,
"subject": parsed["subject"],
"sent_at": parsed["sent_at"],
"received_at": cls._iso_datetime(mail.created_at),
"received_date": mail.received_date.isoformat(),
"status": mail.status,
"has_body": parsed["has_body"],
"body_text": parsed["body_text"],
"body_html": parsed["body_html"],
"attachments": attachment_items,
}
@staticmethod
def _parse_eml(path: str) -> dict[str, object]:
# 显式标注:字面量里只有一个 bool,mypy 会把它推断成 `dict[str, bool | None]`,
# 而 dict 是**不变型**,`dict[str, bool | None]` 不是 `dict[str, object]` 的子类型。
# 运行时本来就是合法值,缺的只是这一个标注。
empty: dict[str, object] = {
"subject": None,
"sent_at": None,
"has_body": False,
"body_text": None,
"body_html": None,
}
try:
message = BytesParser(policy=policy.default).parsebytes(Path(path).read_bytes())
except (OSError, ValueError):
return empty
subject = OffsiteFundService._decode_header(message.get("Subject"))
sent_at = None
raw_date = message.get("Date")
if raw_date:
try:
sent_at = parsedate_to_datetime(raw_date).isoformat()
except (TypeError, ValueError, OverflowError):
sent_at = raw_date
body_text, body_html = OffsiteFundService._mail_body(message)
return {
"subject": subject,
"sent_at": sent_at,
"has_body": bool(body_text or body_html),
"body_text": body_text,
"body_html": body_html,
}
@staticmethod
def _decode_header(value: str | None) -> str | None:
if not value:
return None
try:
return str(make_header(decode_header(value)))
except (LookupError, UnicodeError, ValueError):
return value
# 形参标注用 EmailMessage 而不是基类 Message:`get_content()` 只定义在 EmailMessage 上,
# 而调用方传的是 `BytesParser(policy=policy.default).parsebytes(...)` 的返回值 ——
# 它本来就是 EmailMessage。原先标注成基类,mypy 于是在下面报"没有该属性"。
@staticmethod
def _mail_body(message: EmailMessage) -> tuple[str | None, str | None]:
text_body: str | None = None
html_body: str | None = None
parts = message.walk() if message.is_multipart() else (message,)
for part in parts:
if part.is_multipart() or part.get_content_disposition() == "attachment":
continue
content_type = part.get_content_type()
try:
content = part.get_content()
except (LookupError, UnicodeError, ValueError):
continue
if not isinstance(content, str) or not content.strip():
continue
if content_type == "text/plain" and text_body is None:
text_body = content
elif content_type == "text/html" and html_body is None:
html_body = content
return text_body, html_body
@staticmethod
def _iso_datetime(value: datetime | None) -> str | None:
return value.replace(tzinfo=UTC).isoformat() if value is not None else None
async def receive_recognized_mail(
self, payload: ReceiveRecognizedMailRequest, context: RequestContext
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:write",))
if denied is not None:
return denied
settings = get_settings()
if payload.sender not in settings.offsite_allowed_senders:
return {"code": 403, "message": "发件人不在场外业务白名单", "data": {}}
business = [item for item in payload.attachments if item.document_type != "other"]
if not business:
return {"code": 0, "message": "ok", "data": {"business": False}}
async with self.session.begin():
existing = await self.session.scalar(select(OffsiteFundMail).where(
OffsiteFundMail.imap_uid == payload.imap_uid,
OffsiteFundMail.message_id == payload.message_id,
))
if existing is not None:
return {"code": 0, "message": "ok", "data": {"mail_id": existing.mail_id}}
mail_id = await self._next_mail_id()
now = datetime.now(UTC).replace(tzinfo=None)
mail = OffsiteFundMail(
mail_id=mail_id, imap_uid=payload.imap_uid, message_id=payload.message_id,
received_date=datetime.now(SHANGHAI).date(), sender=payload.sender,
return_path=payload.return_path, auth_result=payload.auth_result,
original_eml_path=payload.eml_path, status="recognized",
created_at=now, updated_at=now,
)
self.session.add(mail)
summaries: list[OffsiteDocumentSummary] = []
for index, item in enumerate(payload.attachments, start=1):
attachment_id = f"{mail_id}-A{index:02d}"
original_attempt = await self.session.scalar(
select(OffsiteRecognitionAttempt)
.where(
OffsiteRecognitionAttempt.message_id == payload.message_id,
OffsiteRecognitionAttempt.file_hash == item.file_hash,
)
.order_by(OffsiteRecognitionAttempt.attempt_no)
)
attachment_fields = (
original_attempt.extracted_fields
if original_attempt is not None
else item.extracted_fields
)
attachment_confidence = (
original_attempt.field_confidence
if original_attempt is not None
else {key: str(value) for key, value in item.field_confidence.items()}
)
attachment_type = (
original_attempt.document_type
if original_attempt is not None
else item.document_type
)
attachment_evidence = (
original_attempt.page_evidence
if original_attempt is not None
else item.page_evidence
)
self.session.add(OffsiteFundAttachment(
attachment_id=attachment_id, mail_id=mail_id, filename=item.filename,
file_hash=item.file_hash, media_type=item.media_type,
size_bytes=item.size_bytes, document_type=attachment_type,
original_file_path=item.original_file_path, ocr_text=item.ocr_text,
extracted_fields=attachment_fields,
field_confidence=attachment_confidence,
page_evidence=attachment_evidence, status="recognized",
created_at=now,
))
await self._link_initial_recognition_attempts(
payload.message_id, item.file_hash, mail_id, attachment_id
)
if item.document_type == "other":
continue
if item.document_type == "summary":
summaries.append(await self._create_document(
mail_id, f"{attachment_id}-S", attachment_id, item, "subscription", now))
summaries.append(await self._create_document(
mail_id, f"{attachment_id}-R", attachment_id, item, "redemption", now))
else:
summaries.append(await self._create_document(
mail_id, attachment_id, attachment_id, item, item.document_type, now))
self._add_audit(context, "offsite.mail_recognized", {
"mail_id": mail_id,
"imap_uid": payload.imap_uid,
"message_id": payload.message_id,
"business_attachment_count": len(business),
})
return {"code": 0, "message": "ok", "data": {
"mail_id": mail_id,
"documents": [item.model_dump(mode="json") for item in summaries],
}}
async def retry_document_recognition(
self, task_id: str, operator_id: str, context: RequestContext
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:write", "offsite:confirm"))
if denied is not None:
return denied
settings = get_settings()
async with self.session.begin():
document = await self.session.scalar(
select(OffsiteFundDocument)
.where(OffsiteFundDocument.task_id == task_id)
.with_for_update()
)
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
if document.status == "recognition_retrying":
return {"code": 409, "message": "单据正在识别重试中", "data": {}}
if document.status not in {"recognition_exception", "recognition_review"}:
return {"code": 422, "message": "当前单据状态不允许识别重试", "data": {}}
attachment = await self.session.scalar(
select(OffsiteFundAttachment).where(
OffsiteFundAttachment.attachment_id == document.attachment_id
)
)
mail = await self.session.scalar(
select(OffsiteFundMail).where(OffsiteFundMail.mail_id == document.mail_id)
)
if attachment is None or mail is None:
return {"code": 404, "message": "识别原始附件或邮件不存在", "data": {}}
source = RecognitionSourceFile(
filename=attachment.filename,
media_type=attachment.media_type,
payload=b"",
file_hash=attachment.file_hash,
original_file_path=attachment.original_file_path,
)
document.status = "recognition_retrying"
document.updated_at = datetime.now(UTC).replace(tzinfo=None)
result: StructuredRecognitionResult | None = None
error_message: str | None = None
started_at = datetime.now(UTC).replace(tzinfo=None)
try:
payload = await asyncio.to_thread(Path(source.original_file_path).read_bytes)
source = RecognitionSourceFile(
filename=source.filename,
media_type=source.media_type,
payload=payload,
file_hash=source.file_hash,
original_file_path=source.original_file_path,
)
if self.recognizer is None and not (
settings.offsite_ocr_enabled and settings.offsite_deepseek_enabled
):
error_message = "真实识别能力未启用"
else:
recognizer = self.recognizer or OffsiteDocumentRecognitionAdapter(settings)
result = await recognizer.recognize(source)
except asyncio.CancelledError:
raise
except (OSError, RuntimeError, ValueError, TypeError) as exc:
error_message = f"{type(exc).__name__}: {str(exc)[:450]}"
finished_at = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
document = await self.session.scalar(
select(OffsiteFundDocument)
.where(OffsiteFundDocument.task_id == task_id)
.with_for_update()
)
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
mail = await self.session.scalar(
select(OffsiteFundMail).where(OffsiteFundMail.mail_id == document.mail_id)
)
if mail is None:
return {"code": 404, "message": "关联邮件不存在", "data": {}}
attempt_no = await self._next_recognition_attempt_no(
mail.message_id, source.file_hash
)
attempt_status = self._recognition_attempt_status(result, error_message)
self.session.add(
OffsiteRecognitionAttempt(
task_id=task_id,
mail_id=document.mail_id,
attachment_id=document.attachment_id,
imap_uid=mail.imap_uid,
message_id=mail.message_id,
filename=source.filename,
file_hash=source.file_hash,
original_file_path=source.original_file_path,
attempt_no=attempt_no,
source="manual",
operator_id=operator_id,
document_type=result.document_type if result is not None else "other",
ocr_status=result.ocr_status if result is not None else "error",
llm_status=result.llm_status if result is not None else "error",
extracted_fields=result.extracted_fields if result is not None else {},
field_confidence=(
{key: str(value) for key, value in result.field_confidence.items()}
if result is not None
else {}
),
missing_fields=list(result.missing_fields) if result is not None else [],
low_confidence_fields=(
list(result.low_confidence_fields) if result is not None else []
),
page_evidence=result.page_evidence if result is not None else {},
status=attempt_status,
error_message=error_message or (
result.error_message if result is not None else "识别器未返回结果"
),
started_at=started_at,
finished_at=finished_at,
created_at=finished_at,
)
)
if result is not None and self._recognition_result_matches(document, result):
fields = result.extracted_fields
status = self._recognition_status(
self._recognized_attachment_for_retry(source, result),
cast(Literal["subscription", "redemption"], document.document_type),
)
self._apply_recognition_fields(document, fields)
document.operator_decision = "未处理"
document.status = status
document.updated_at = finished_at
await self._refresh_rule_results(document, fields, finished_at)
await self._reset_execution_plan(document.task_id, fields, finished_at)
self._add_audit(context, "offsite.document_recognition_recovered", {
"task_id": task_id,
"operator_id": operator_id,
"attempt_no": attempt_no,
"status": status,
})
return {
"code": 0,
"message": "识别重试完成",
"data": {"task_id": task_id, "status": status, "attempt_no": attempt_no},
}
document.status = "recognition_exception"
document.updated_at = finished_at
self._add_audit(context, "offsite.document_recognition_retry_failed", {
"task_id": task_id,
"operator_id": operator_id,
"attempt_no": attempt_no,
"error": error_message or "识别结果仍不满足业务字段要求",
})
return {
"code": 0,
"message": "识别重试未通过,已保留异常状态",
"data": {"task_id": task_id, "status": document.status, "attempt_no": attempt_no},
}
async def _link_initial_recognition_attempts(
self, message_id: str, file_hash: str, mail_id: str, attachment_id: str
) -> None:
await self.session.execute(
update(OffsiteRecognitionAttempt)
.where(
OffsiteRecognitionAttempt.message_id == message_id,
OffsiteRecognitionAttempt.file_hash == file_hash,
OffsiteRecognitionAttempt.attachment_id.is_(None),
)
.values(mail_id=mail_id, attachment_id=attachment_id)
)
async def _next_recognition_attempt_no(self, message_id: str, file_hash: str) -> int:
current = await self.session.scalar(
select(func.max(OffsiteRecognitionAttempt.attempt_no)).where(
OffsiteRecognitionAttempt.message_id == message_id,
OffsiteRecognitionAttempt.file_hash == file_hash,
)
)
return int(current or 0) + 1
async def _refresh_rule_results(
self,
document: OffsiteFundDocument,
fields: Mapping[str, object],
now: datetime,
blocked_rule_codes: set[str] | None = None,
) -> None:
blocked = blocked_rule_codes or set()
amount_yuan = normalize_amount_yuan(
fields.get("申购金额"), fields.get("金额单位")
)
redemption_shares = decimal_from(fields.get("赎回份额"))
if document.document_type == "subscription":
decisions = self.rules.check_subscription(
amount_yuan=amount_yuan,
nav=decimal_from(fields.get("最新净值")),
total_fund_shares=decimal_from(fields.get("基金最新总份额")),
before_holding_shares=decimal_from(fields.get("申请前持有份额")),
)
else:
decisions = self.rules.check_redemption(
redemption_shares=redemption_shares,
total_fund_shares=decimal_from(fields.get("基金最新总份额")),
available_quantity=decimal_from(fields.get("当前最新可用份额")),
)
existing = {
item.rule_code: item
for item in (
await self.session.execute(
select(OffsiteRuleResult).where(
OffsiteRuleResult.task_id == document.task_id
)
)
).scalars().all()
}
for decision in decisions:
result = decision.result
document_value = decision.document_value
database_value = decision.database_value
calculation = decision.calculation
if decision.rule_code in blocked:
result = "无法判断"
document_value = {}
database_value = {}
calculation = {"原因": "依赖的 NL2SQL 查询失败或无可用数据"}
row = existing.get(decision.rule_code)
if row is None:
self.session.add(
OffsiteRuleResult(
task_id=document.task_id,
rule_code=decision.rule_code,
rule_name=decision.rule_name,
result=result,
document_value=document_value,
database_value=database_value,
calculation=calculation,
created_at=now,
)
)
continue
row.rule_name = decision.rule_name
row.result = result
row.document_value = document_value
row.database_value = database_value
row.calculation = calculation
async def _reset_execution_plan(
self, task_id: str, fields: Mapping[str, object], now: datetime
) -> None:
tasks = (
await self.session.execute(
select(OffsiteExecutionPlanTask)
.where(OffsiteExecutionPlanTask.task_id == task_id)
.with_for_update()
)
).scalars().all()
for task in tasks:
task.input_json = dict(fields)
task.output_json = None
task.error_message = None
if task.rule_code == "subscription_minimum_amount":
task.status = "不适用" if task.stage == "查询" else "已完成"
else:
task.status = "待执行"
task.updated_at = now
@staticmethod
def _apply_recognition_fields(
document: OffsiteFundDocument, fields: Mapping[str, object]
) -> None:
raw_date = fields.get("申请日期")
document.fund_code = OffsiteFundService._text(fields.get("基金代码"))
document.fund_name = OffsiteFundService._text(fields.get("基金名称"))
document.account_identifier = OffsiteFundService._text(fields.get("账户标识"))
document.investor_name = OffsiteFundService._text(
fields.get("投资者名称") or fields.get("客户标识")
)
document.application_no = OffsiteFundService._text(fields.get("申请编号"))
document.application_date = parse_application_date(raw_date)
document.raw_application_date = OffsiteFundService._text(raw_date)
document.agency = OffsiteFundService._text(fields.get("代销机构"))
document.subscription_amount_yuan = normalize_amount_yuan(
fields.get("申购金额"), fields.get("金额单位")
)
document.redemption_shares = decimal_from(fields.get("赎回份额"))
@staticmethod
def _recognized_attachment_for_retry(
source: RecognitionSourceFile, result: StructuredRecognitionResult
) -> RecognizedAttachment:
return RecognizedAttachment(
filename=source.filename,
file_hash=source.file_hash,
original_file_path=source.original_file_path,
media_type=source.media_type,
size_bytes=len(source.payload),
document_type=result.document_type,
extracted_fields=result.extracted_fields,
field_confidence=result.field_confidence,
ocr_text=result.ocr_text,
page_evidence=result.page_evidence,
)
@staticmethod
def _recognition_result_matches(
document: OffsiteFundDocument, result: StructuredRecognitionResult | None
) -> bool:
if result is None or result.document_type != document.document_type:
return False
if result.error_message and (
result.ocr_status in {"error", "misconfigured"}
or result.llm_status in {"error", "misconfigured"}
):
return False
return OffsiteFundService._recognition_status(
OffsiteFundService._recognized_attachment_for_retry(
RecognitionSourceFile(
filename="retry",
media_type="application/octet-stream",
payload=b"",
file_hash="r" * 64,
),
result,
),
cast(Literal["subscription", "redemption"], document.document_type),
) != "recognition_exception"
@staticmethod
def _recognition_attempt_status(
result: StructuredRecognitionResult | None, error_message: str | None
) -> str:
if result is None or error_message:
return "error"
if result.error_message and (
result.ocr_status in {"error", "misconfigured"}
or result.llm_status in {"error", "misconfigured"}
):
return "error"
if result.missing_fields:
return "recognition_exception"
if result.low_confidence_fields:
return "recognition_review"
return "success"
async def confirm_document(
self, task_id: str, decision: OperationDecision, operator_id: str,
context: RequestContext,
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:confirm", "offsite:write"))
if denied is not None:
return denied
async with self.session.begin():
document = await self.session.scalar(select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id).with_for_update())
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
if document.status in {
"recognition_exception",
"recognition_review",
"recognition_retrying",
"query_failed",
}:
return {
"code": 422,
"message": "识别或数据核对尚未完成,不能进行业务确认",
"data": {},
}
document.operator_decision = decision
document.status = "operator_confirmed"
document.updated_at = datetime.now(UTC).replace(tzinfo=None)
await self._refresh_mail_status(document.mail_id)
self._add_audit(context, "offsite.document_confirmed", {
"task_id": task_id,
"operator_id": operator_id,
"decision": decision,
})
return {"code": 0, "message": "ok", "data": {"task_id": task_id, "decision": decision}}
async def recalculate_statistics(
self, fund_code: str, application_date: str, context: RequestContext
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:read", "offsite:write"))
if denied is not None:
return denied
target_date = parse_application_date(application_date)
if target_date is None:
return {"code": 422, "message": "申请日期格式不正确", "data": {}}
documents = (await self.session.execute(select(OffsiteFundDocument).where(
OffsiteFundDocument.fund_code == fund_code,
OffsiteFundDocument.application_date == target_date,
OffsiteFundDocument.operator_decision == "确认正常",
))).scalars().all()
task_ids = [document.task_id for document in documents]
successful_normal_returns: set[str] = set()
if task_ids:
successful_normal_returns = set((await self.session.execute(
select(OffsiteNotification.business_key).where(
OffsiteNotification.business_key.in_(task_ids),
OffsiteNotification.notification_type.in_(
("mail_return", "normal_return")
),
OffsiteNotification.status == "发送成功",
)
)).scalars().all())
rows = [
document for document in documents
if document.task_id in successful_normal_returns
]
subscription_total = sum(
((row.subscription_amount_yuan or Decimal("0")) for row in rows),
Decimal("0"),
)
redemption_total = sum(
((row.redemption_shares or Decimal("0")) for row in rows),
Decimal("0"),
)
agency_breakdown = self._agency_breakdown(rows)
latest_nav = await self._query_latest_nav(fund_code, target_date, context)
redemption_amount_yuan = (
redemption_total * latest_nav if latest_nav is not None else None
)
net_flow_amount_yuan = (
subscription_total - redemption_amount_yuan
if redemption_amount_yuan is not None
else (subscription_total if redemption_total == 0 else None)
)
return {"code": 0, "message": "ok", "data": {
"fund_code": fund_code, "application_date": target_date.isoformat(),
"fund_name": next((row.fund_name for row in rows if row.fund_name), None),
"subscription_amount_yuan": str(subscription_total),
"subscription_count": sum(1 for row in rows if row.document_type == "subscription"),
"redemption_shares": str(redemption_total),
"redemption_count": sum(1 for row in rows if row.document_type == "redemption"),
"latest_nav": str(latest_nav) if latest_nav is not None else None,
"redemption_amount_yuan": (
str(redemption_amount_yuan) if redemption_amount_yuan is not None else None
),
"net_flow_amount_yuan": (
str(net_flow_amount_yuan) if net_flow_amount_yuan is not None else None
),
"agency_breakdown": agency_breakdown,
}}
async def trigger_agent_nl2sql(
self, task_id: str, operator_id: str, manual_confirmed: bool,
context: RequestContext,
) -> dict[str, object]:
denied = self._permission_error(
context, ("offsite:nl2sql", "offsite:write", "financial:nl2sql:read")
)
if denied is not None:
return denied
if not manual_confirmed:
return {"code": 422, "message": "必须传递人工已确认原始文件内容状态", "data": {}}
now = datetime.now(UTC).replace(tzinfo=None)
records: list[dict[str, object]] = []
async with self.session.begin():
document = await self.session.scalar(select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id))
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
attachment = await self.session.scalar(select(OffsiteFundAttachment).where(
OffsiteFundAttachment.attachment_id == document.attachment_id))
if attachment is None:
return {"code": 404, "message": "原始附件不存在", "data": {}}
questions = self._nl2sql_questions(document)
if not questions:
return {"code": 422, "message": "查询条件不足", "data": {}}
fields = dict(attachment.extracted_fields)
query_fields: dict[str, object] = {}
blocked_rules: set[str] = set()
for rule_code, question in questions:
result = await asyncio.to_thread(self.nl2sql.query, question, context)
row = self._first_query_row(result)
raw_query_status = str(result.get("status", "error"))
query_status = raw_query_status
query_result = result
if query_status == "success" and row is None:
query_status = "query_failed"
query_result = {
**result,
"status": query_status,
"message": "查询无可用数据",
"raw_status": "success",
}
values, missing = self._query_values(rule_code, row)
if query_status != "success":
blocked_rules.add(rule_code)
query_result = {
**query_result,
"status": "query_failed",
"message": (
query_result.get("message")
if isinstance(query_result.get("message"), str)
else "查询未成功"
),
"raw_status": raw_query_status,
}
query_status = "query_failed"
elif missing:
blocked_rules.add(rule_code)
query_status = "query_failed"
query_result = {
**query_result,
"status": query_status,
"message": f"查询结果缺少字段:{', '.join(missing)}",
}
query_fields.update(values)
safe_query_result = cast(
dict[str, object], self._json_safe(query_result)
)
self.session.add(OffsiteQueryRecord(
task_id=task_id, rule_code=rule_code,
natural_language_request=question, script_path=self.nl2sql.script_path,
result_summary=safe_query_result, status=query_status,
error_message=(
safe_query_result.get("message")
if isinstance(safe_query_result.get("message"), str) else None
),
created_at=now,
))
await self._update_query_plan_status(
task_id, rule_code, safe_query_result, now
)
records.append({
"rule_code": rule_code,
"status": query_status,
"row_count": 1 if row is not None else 0,
})
fields.update(query_fields)
await self._refresh_rule_results(
document, fields, now, blocked_rule_codes=blocked_rules
)
await self._finish_execution_plan(task_id, fields, blocked_rules, now)
document.status = "query_failed" if blocked_rules else "planned"
document.updated_at = now
self._add_audit(context, "offsite.nl2sql_triggered", {
"task_id": task_id,
"operator_id": operator_id,
"query_count": len(records),
"blocked_rule_count": len(blocked_rules),
})
return {
"code": 0,
"message": "ok",
"data": {
"task_id": task_id,
"status": "query_failed" if blocked_rules else "planned",
"queries": records,
},
}
async def create_notification(
self, task_id: str, notification_type: str, operator_id: str,
context: RequestContext,
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:notify", "offsite:write"))
if denied is not None:
return denied
settings = get_settings()
receiver = {
"risk": settings.offsite_risk_receiver_id,
"settlement": settings.offsite_settlement_receiver_id,
"mail_return": settings.offsite_mail_return_receiver,
"normal_return": settings.offsite_mail_return_receiver,
"exception_return": settings.offsite_mail_return_receiver,
}[notification_type]
draft = f"{task_id} 待发送{notification_type}通知"
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
document = await self.session.scalar(select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == task_id))
if document is None:
return {"code": 404, "message": "单据不存在", "data": {}}
validation = self._validate_notification(document, notification_type)
if validation is not None:
return validation
rule_results = (await self.session.execute(select(OffsiteRuleResult).where(
OffsiteRuleResult.task_id == task_id
))).scalars().all()
payload = self._notification_payload(document, rule_results, operator_id)
notice = OffsiteNotification(
notification_type=notification_type, business_key=task_id,
receiver_id=receiver, operator_id=operator_id, agent_draft=draft,
final_content=draft, payload=payload, status="待发送",
created_at=now, updated_at=now,
)
self.session.add(notice)
await self.session.flush()
await self._refresh_mail_status(document.mail_id)
notice_id = notice.id
self._add_audit(context, "offsite.notification_created", {
"task_id": task_id,
"notification_type": notification_type,
"operator_id": operator_id,
"receiver_id": receiver,
})
return {"code": 0, "message": "ok", "data": {"notification_id": str(notice_id)}}
async def send_notification(
self,
notification_id: int,
operator_id: str,
operator_confirmed: bool,
final_content: str | None,
context: RequestContext,
) -> dict[str, object]:
denied = self._permission_error(context, ("offsite:notify", "offsite:write"))
if denied is not None:
return denied
if not operator_confirmed:
return {"code": 422, "message": "邮件发送前必须完成运营确认", "data": {}}
prepared = await self._prepare_notification_send(
notification_id, operator_id, final_content, context
)
if isinstance(prepared, dict):
return prepared
notice, mail, attachment, receiver = prepared
try:
request = self._build_mail_reply_request(
notice, mail, attachment, receiver, operator_id, final_content
)
result = await asyncio.to_thread(self.smtp_sender.send_reply, request)
except OSError as exc:
result = SmtpSendResult(
status="发送失败",
dry_run=False,
provider_message_id=None,
failure_reason=type(exc).__name__,
retry_count=notice.retry_count + 1,
request_summary={"notification_id": str(notification_id)},
)
await self._finish_notification(notification_id, result, context)
return {
"code": 0,
"message": "ok",
"data": {
"notification_id": str(notification_id),
"status": result.status,
"dry_run": result.dry_run,
"provider_message_id": result.provider_message_id,
"failure_reason": result.failure_reason,
"retry_count": result.retry_count,
},
}
async def _prepare_notification_send(
self,
notification_id: int,
operator_id: str,
final_content: str | None,
context: RequestContext,
) -> tuple[
OffsiteNotification,
OffsiteFundMail,
OffsiteFundAttachment,
str,
] | dict[str, object]:
settings = get_settings()
async with self.session.begin():
notice = await self.session.scalar(select(OffsiteNotification).where(
OffsiteNotification.id == notification_id
).with_for_update())
if notice is None:
return {"code": 404, "message": "通知不存在", "data": {}}
if notice.status == "发送成功":
return {
"code": 0,
"message": "ok",
"data": {
"notification_id": str(notification_id),
"status": notice.status,
"provider_message_id": notice.provider_message_id,
},
}
if notice.status == "发送中":
return {"code": 409, "message": "通知正在发送中", "data": {}}
if notice.retry_count > 0 and notice.retry_count >= settings.offsite_max_retry_count:
return {"code": 422, "message": "通知已达到最大重试次数", "data": {}}
document = await self.session.scalar(select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == notice.business_key))
if document is None:
return {"code": 404, "message": "通知关联的单据或邮件不存在", "data": {}}
mail = await self.session.scalar(select(OffsiteFundMail).where(
OffsiteFundMail.mail_id == document.mail_id))
if mail is None:
return {"code": 404, "message": "通知关联的单据或邮件不存在", "data": {}}
if notice.notification_type in {
"mail_return", "normal_return", "exception_return"
}:
receiver = settings.offsite_mail_return_receiver
else:
receiver = notice.receiver_id if "@" in notice.receiver_id else ""
if not receiver:
return {"code": 422, "message": "通知对象未配置可发送的邮箱地址", "data": {}}
attachment = await self.session.scalar(select(OffsiteFundAttachment).where(
OffsiteFundAttachment.attachment_id == document.attachment_id))
if attachment is None:
return {"code": 404, "message": "通知关联的原始附件不存在", "data": {}}
if final_content is not None and final_content != notice.final_content:
self._add_audit(context, "offsite.notification_content_changed", {
"notification_id": str(notification_id),
"operator_id": operator_id,
"changed": True,
})
notice.final_content = final_content
notice.operator_id = operator_id
notice.status = "发送中"
return notice, mail, attachment, receiver
def _build_mail_reply_request(
self,
notice: OffsiteNotification,
mail: OffsiteFundMail,
attachment: OffsiteFundAttachment,
receiver: str,
operator_id: str,
final_content: str | None,
) -> OffsiteMailReplyRequest:
path = Path(attachment.original_file_path)
payload = path.read_bytes()
return OffsiteMailReplyRequest(
to_address=receiver,
subject="场外基金申购赎回处理结果",
body=final_content if final_content is not None else notice.final_content,
operator_id=operator_id,
operator_confirmed=True,
reply_to_message_id=mail.message_id,
attachments=(
SmtpAttachment(
filename=attachment.filename,
media_type=attachment.media_type,
payload=payload,
),
),
retry_count=notice.retry_count,
)
async def _finish_notification(
self, notification_id: int, result: SmtpSendResult, context: RequestContext
) -> None:
now = datetime.now(UTC).replace(tzinfo=None)
async with self.session.begin():
notice = await self.session.scalar(select(OffsiteNotification).where(
OffsiteNotification.id == notification_id
).with_for_update())
if notice is None:
return
notice.status = result.status
notice.provider_message_id = result.provider_message_id
notice.failure_reason = result.failure_reason
notice.retry_count = result.retry_count
notice.sent_at = now if result.status == "发送成功" else None
notice.updated_at = now
if result.status == "发送成功" and notice.notification_type in {
"mail_return", "normal_return", "exception_return"
}:
document = await self.session.scalar(select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == notice.business_key))
if document is not None:
mail = await self.session.scalar(select(OffsiteFundMail).where(
OffsiteFundMail.mail_id == document.mail_id).with_for_update())
if mail is not None:
mail.status = "normal_return_sent"
if result.status == "发送成功":
document = await self.session.scalar(select(OffsiteFundDocument).where(
OffsiteFundDocument.task_id == notice.business_key))
if document is not None:
await self._refresh_mail_status(document.mail_id)
self._add_audit(context, "offsite.notification_send_finished", {
"notification_id": str(notification_id),
"status": result.status,
"dry_run": result.dry_run,
"retry_count": result.retry_count,
})
async def _query_latest_nav(
self, fund_code: str, application_date: object, context: RequestContext
) -> Decimal | None:
question = (
f"基金代码为{fund_code},申请日期为{application_date},"
"查询该基金最新净值"
)
result = await asyncio.to_thread(self.nl2sql.query, question, context)
if result.get("status") != "success":
return None
data = result.get("data")
if not isinstance(data, dict):
return None
rows = data.get("rows")
if not isinstance(rows, list):
return None
for row in rows:
if not isinstance(row, dict):
continue
for key in ("nav", "latest_nav", "基金净值"):
value = decimal_from(row.get(key))
if value is not None and value > 0:
return value
return None
@staticmethod
def _first_query_row(result: Mapping[str, object]) -> dict[str, object] | None:
data = result.get("data")
if not isinstance(data, Mapping):
return None
rows = data.get("rows")
if not isinstance(rows, list) or not rows:
return None
first = rows[0]
return dict(first) if isinstance(first, Mapping) else None
@staticmethod
def _json_safe(value: object) -> object:
if isinstance(value, Decimal):
return str(value)
if isinstance(value, (date, datetime)):
return value.isoformat()
if isinstance(value, Mapping):
return {str(key): OffsiteFundService._json_safe(item) for key, item in value.items()}
if isinstance(value, (list, tuple)):
return [OffsiteFundService._json_safe(item) for item in value]
return value
@staticmethod
def _query_field_mapping(rule_code: str) -> dict[str, tuple[str, ...]]:
"""规则对应查询结果的中文字段与可接受列别名。"""
if rule_code == "subscription_holding_ratio":
return {
"最新净值": ("nav",),
"基金最新总份额": ("total_fund_shares",),
"申请前持有份额": ("total_quantity", "shares"),
}
if rule_code == "subscription_single_share_limit":
return {
"最新净值": ("nav",),
"基金最新总份额": ("total_fund_shares",),
}
if rule_code == "redemption_large_ratio":
return {"基金最新总份额": ("total_fund_shares",)}
return {"当前最新可用份额": ("available_quantity",)}
@staticmethod
def _query_values(
rule_code: str, row: Mapping[str, object] | None
) -> tuple[dict[str, object], tuple[str, ...]]:
if row is None:
return {}, ("查询结果",)
mappings = OffsiteFundService._query_field_mapping(rule_code)
values: dict[str, object] = {}
missing: list[str] = []
for field_name, aliases in mappings.items():
value = next(
(row.get(alias) for alias in aliases if row.get(alias) is not None),
None,
)
number = decimal_from(value)
if number is None:
missing.append(field_name)
continue
if field_name in {"最新净值", "基金最新总份额"} and number <= 0:
missing.append(field_name)
continue
values[field_name] = str(number)
return values, tuple(missing)
async def _finish_execution_plan(
self,
task_id: str,
fields: Mapping[str, object],
blocked_rule_codes: set[str],
now: datetime,
) -> None:
results = {
item.rule_code: item
for item in (
await self.session.execute(
select(OffsiteRuleResult).where(
OffsiteRuleResult.task_id == task_id
)
)
).scalars().all()
}
tasks = (
await self.session.execute(
select(OffsiteExecutionPlanTask)
.where(OffsiteExecutionPlanTask.task_id == task_id)
.with_for_update()
)
).scalars().all()
for task in tasks:
task.input_json = dict(fields)
if task.stage == "查询":
continue
decision = results.get(task.rule_code)
if task.rule_code in blocked_rule_codes:
task.status = "无法判断"
task.output_json = {"原因": "依赖的 NL2SQL 查询失败或无可用数据"}
task.error_message = "依赖的 NL2SQL 查询失败或无可用数据"
else:
task.status = "已完成"
task.output_json = {
"result": decision.result if decision is not None else "无法判断",
"document_value": (
decision.document_value if decision is not None else {}
),
"database_value": (
decision.database_value if decision is not None else {}
),
"calculation": (
decision.calculation if decision is not None else {}
),
}
task.error_message = None
task.updated_at = now
async def _refresh_mail_status(self, mail_id: str) -> None:
mail = await self.session.scalar(select(OffsiteFundMail).where(
OffsiteFundMail.mail_id == mail_id
).with_for_update())
if mail is None:
return
documents = (await self.session.execute(select(OffsiteFundDocument).where(
OffsiteFundDocument.mail_id == mail_id
))).scalars().all()
if not documents:
return
task_ids = [document.task_id for document in documents]
notices = (await self.session.execute(select(OffsiteNotification).where(
OffsiteNotification.business_key.in_(task_ids)
))).scalars().all()
for document in documents:
related = [
notice for notice in notices if notice.business_key == document.task_id
]
if document.operator_decision == "确认正常":
if not any(
notice.status == "发送成功"
and notice.notification_type in {"mail_return", "normal_return"}
for notice in related
):
mail.status = "processing"
return
elif document.operator_decision == "确认异常":
if not any(
notice.status == "发送成功"
and notice.notification_type in {"mail_return", "exception_return"}
for notice in related
):
mail.status = "processing"
return
if any(
notice.notification_type in {"risk", "settlement"}
and notice.status != "发送成功"
for notice in notices
):
mail.status = "processing"
return
mail.status = (
"normal_return_sent"
if all(document.operator_decision == "确认正常" for document in documents)
else "completed"
)
async def _next_mail_id(self) -> str:
today = datetime.now(SHANGHAI).date()
prefix = today.strftime("%Y%m%d")
current = await self.session.scalar(select(func.count()).select_from(OffsiteFundMail).where(
OffsiteFundMail.received_date == today))
return f"{prefix}-{int(current or 0) + 1:03d}"
async def _create_document(
self,
mail_id: str,
task_id: str,
attachment_id: str,
item: RecognizedAttachment,
document_type: Literal["subscription", "redemption"],
now: datetime,
) -> OffsiteDocumentSummary:
fields = item.extracted_fields
raw_date = fields.get("申请日期")
amount_yuan = normalize_amount_yuan(fields.get("申购金额"), fields.get("金额单位"))
redemption_shares = decimal_from(fields.get("赎回份额"))
status = self._recognition_status(item, document_type)
document = OffsiteFundDocument(
task_id=task_id, mail_id=mail_id, attachment_id=attachment_id,
document_type=document_type, fund_code=self._text(fields.get("基金代码")),
fund_name=self._text(fields.get("基金名称")),
account_identifier=self._text(fields.get("账户标识")),
investor_name=self._text(fields.get("投资者名称") or fields.get("客户标识")),
application_no=self._text(fields.get("申请编号")),
application_date=parse_application_date(raw_date),
raw_application_date=self._text(raw_date), agency=self._text(fields.get("代销机构")),
subscription_amount_yuan=amount_yuan, redemption_shares=redemption_shares,
status=status, created_at=now, updated_at=now,
)
self.session.add(document)
self._append_plan_tasks(task_id, document_type, fields, now)
decisions = self.rules.check_subscription(
amount_yuan=amount_yuan, nav=decimal_from(fields.get("最新净值")),
total_fund_shares=decimal_from(fields.get("基金最新总份额")),
before_holding_shares=decimal_from(fields.get("申请前持有份额")),
) if document_type == "subscription" else self.rules.check_redemption(
redemption_shares=redemption_shares,
total_fund_shares=decimal_from(fields.get("基金最新总份额")),
available_quantity=decimal_from(fields.get("当前最新可用份额")),
)
results: dict[str, Literal["正常", "异常", "无法判断"]] = {}
for decision in decisions:
results[decision.rule_code] = decision.result
self.session.add(OffsiteRuleResult(
task_id=task_id, rule_code=decision.rule_code, rule_name=decision.rule_name,
result=decision.result, document_value=decision.document_value,
database_value=decision.database_value, calculation=decision.calculation,
created_at=now,
))
return OffsiteDocumentSummary(
task_id=task_id, document_type=cast(DocumentType, document_type), status=status,
rule_results=results)
def _append_plan_tasks(
self, task_id: str, document_type: str, fields: dict[str, object], now: datetime
) -> None:
rules: tuple[tuple[str, str], ...] = (
("subscription_minimum_amount", "申购最低金额"),
("subscription_holding_ratio", "申购后单一投资者持有比例"),
("subscription_single_share_limit", "申购单笔份额上限"),
)
if document_type == "redemption":
rules = (("redemption_large_ratio", "赎回巨额比例"),
("redemption_available_quantity", "账户可用份额"))
for rule_code, title in rules:
for stage in ("查询", "计算", "核对"):
if rule_code == "subscription_minimum_amount" and stage == "查询":
status = "不适用"
else:
status = "已完成" if stage != "查询" else "待执行"
self.session.add(OffsiteExecutionPlanTask(
task_id=task_id, stage=stage, rule_code=rule_code, title=title,
depends_on=[], input_json=fields, output_json=None,
status=status,
error_message=None, created_at=now, updated_at=now,
))
@staticmethod
def _nl2sql_questions(document: OffsiteFundDocument) -> list[tuple[str, str]]:
if not document.fund_code:
return []
base = f"基金代码为{document.fund_code}"
if document.account_identifier:
base += f",账户标识为{document.account_identifier}"
if document.document_type == "subscription":
return [
(
"subscription_holding_ratio",
base + ",查询基金最新总份额、最新净值和申请前持有份额",
),
("subscription_single_share_limit", base + ",查询基金最新总份额和最新净值"),
]
return [
("redemption_large_ratio", base + ",查询产品最新总份额"),
("redemption_available_quantity", base + ",查询账户当前最新可用份额"),
]
@staticmethod
def _text(value: object) -> str | None:
if value is None:
return None
text = str(value).strip()
return text or None
async def _update_query_plan_status(
self, task_id: str, rule_code: str, result: dict[str, object], now: datetime
) -> None:
task = await self.session.scalar(select(OffsiteExecutionPlanTask).where(
OffsiteExecutionPlanTask.task_id == task_id,
OffsiteExecutionPlanTask.rule_code == rule_code,
OffsiteExecutionPlanTask.stage == "查询",
).with_for_update())
if task is None:
return
status = str(result.get("status", "error"))
task.status = "已完成" if status in {"ready", "success"} else "查询失败"
if status == "need_confirmation":
task.status = "无法判断"
task.output_json = result
message = result.get("message")
task.error_message = message if isinstance(message, str) else None
task.updated_at = now
@staticmethod
def _permission_error(
context: RequestContext, permissions: tuple[str, ...]
) -> dict[str, object] | None:
roles = {"operator", "risk_operator", "admin", "super_admin"}
if not roles.intersection(context.roles):
return {"code": 403, "message": "当前角色不能操作场外基金流程", "data": {}}
if not set(permissions).intersection(context.permissions):
return {"code": 403, "message": "缺少场外基金操作权限", "data": {}}
return None
def _add_audit(
self, context: RequestContext, action_type: str, detail: dict[str, object]
) -> None:
self.session.add(InteractionAudit(
actor_type="user",
actor_id=int(context.user_id) if context.user_id.isdigit() else None,
target_customer_id=None,
session_id=None,
portal=context.portal,
action_type=action_type,
detail={**detail, "trace_id": context.trace_id},
created_at=datetime.now(UTC).replace(tzinfo=None),
))
@staticmethod
def _recognition_status(
item: RecognizedAttachment, document_type: Literal["subscription", "redemption"]
) -> str:
fields = item.extracted_fields
required = ["基金代码", "基金名称", "账户标识", "申请编号", "申请日期", "代销机构"]
if not (fields.get("投资者名称") or fields.get("客户标识")):
return "recognition_exception"
if document_type == "subscription":
required.extend(["申购金额", "金额单位"])
else:
required.append("赎回份额")
if any(not str(fields.get(name) or "").strip() for name in required):
return "recognition_exception"
if not item.field_confidence:
return "recognition_review"
confidence_values = [
value for value in (decimal_from(value) for value in item.field_confidence.values())
if value is not None
]
if not confidence_values:
return "recognition_review"
min_confidence = min(confidence_values)
if min_confidence < Decimal("0.80"):
return "recognition_exception"
if min_confidence < Decimal("0.95"):
return "recognition_review"
return "planned"
@staticmethod
def _agency_breakdown(rows: Sequence[OffsiteFundDocument]) -> list[dict[str, object]]:
grouped: dict[str, dict[str, object]] = {}
for row in rows:
agency = row.agency or "未识别代销机构"
item = grouped.setdefault(agency, {
"agency": agency,
"subscription_amount_yuan": Decimal("0"),
"subscription_count": 0,
"redemption_shares": Decimal("0"),
"redemption_count": 0,
})
if row.document_type == "subscription":
item["subscription_amount_yuan"] = cast(
Decimal, item["subscription_amount_yuan"]
) + (row.subscription_amount_yuan or Decimal("0"))
subscription_count = item["subscription_count"]
item["subscription_count"] = (
subscription_count + 1 if isinstance(subscription_count, int) else 1
)
elif row.document_type == "redemption":
item["redemption_shares"] = cast(
Decimal, item["redemption_shares"]
) + (row.redemption_shares or Decimal("0"))
redemption_count = item["redemption_count"]
item["redemption_count"] = (
redemption_count + 1 if isinstance(redemption_count, int) else 1
)
return [
{
**item,
"subscription_amount_yuan": str(item["subscription_amount_yuan"]),
"redemption_shares": str(item["redemption_shares"]),
}
for item in grouped.values()
]
@staticmethod
def _validate_notification(
document: OffsiteFundDocument, notification_type: str
) -> dict[str, object] | None:
if notification_type == "risk" and document.operator_decision != "确认异常":
return {"code": 422, "message": "风控通知只允许发送已确认异常单据", "data": {}}
if notification_type == "settlement" and document.operator_decision != "确认正常":
return {"code": 422, "message": "资金清算通知只允许发送已确认正常单据", "data": {}}
if notification_type == "normal_return" and document.operator_decision != "确认正常":
return {"code": 422, "message": "正常返回只允许发送已确认正常单据", "data": {}}
if notification_type == "exception_return" and document.operator_decision != "确认异常":
return {"code": 422, "message": "异常返回只允许发送已确认异常单据", "data": {}}
if notification_type == "mail_return" and document.operator_decision == "未处理":
return {"code": 422, "message": "邮件返回前必须先完成人工确认", "data": {}}
return None
@staticmethod
def _notification_payload(
document: OffsiteFundDocument,
rule_results: Sequence[OffsiteRuleResult],
operator_id: str,
) -> dict[str, object]:
anomalies = [
{
"rule_code": result.rule_code,
"rule_name": result.rule_name,
"result": result.result,
"document_value": result.document_value,
"database_value": result.database_value,
"calculation": result.calculation,
}
for result in rule_results if result.result == "异常"
]
return {
"task_id": document.task_id,
"mail_id": document.mail_id,
"attachment_id": document.attachment_id,
"document_type": document.document_type,
"fund_code": document.fund_code,
"fund_name": document.fund_name,
"account_identifier": OffsiteFundService._mask_account(document.account_identifier),
"application_no": document.application_no,
"application_date": (
document.application_date.isoformat() if document.application_date else None
),
"agency": document.agency,
"operator_id": operator_id,
"operator_decision": document.operator_decision,
"confirmed_at": datetime.now(UTC).replace(tzinfo=None).isoformat(timespec="seconds"),
"anomalies": anomalies,
}
@staticmethod
def _mask_account(value: str | None) -> str | None:
if value is None or len(value) <= 4:
return value
return f"{value[:2]}***{value[-2:]}"