merge: 并入同事的场外申购/推广/行情/NL2SQL 线(11 提交、334 文件)

冲突仅 3 个文件,全部取并集(双方都没有需要丢弃的改动):
- app/main.py:import 双方路由(我方 knowledge_management + 同事的 offsite_fund/
  promotion_material);include_router 段本已自动合并
- app/service/agent/bootstrap.py:import 与工具注册均取并集
  (query_customer_profile + query_financial_data 都注册)
- tests/integration/test_config_release_mysql.py:outbox 清理同时保留
  架构师的 event_type 限定(防误删其它域 outbox 行)与同事新增的 peer_release_id

同事这轮带入:11 个 alembic 迁移(建 offsite_* / promotion_* 等表)、
场外申购与推广素材 Agent、financial NL2SQL 工具。
注意:本库尚无 offsite_*/promotion_* 表,跑相关测试前需要执行 alembic upgrade。

边界核对:同事的场外代码未写入场内交易表(fin_sim_order/fin_capital_flow/fin_cash_ledger),
符合 AGENTS.md 规则 8。
This commit is contained in:
qyqy
2026-09-11 18:56:26 +08:00
118 changed files with 18967 additions and 191 deletions
@@ -0,0 +1,551 @@
from __future__ import annotations
from datetime import datetime
from decimal import Decimal
from pathlib import Path
from typing import Any
import pytest
from sqlalchemy import select, text, update
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
from sqlalchemy.pool import StaticPool
from app.core.config import Settings
from app.core.contracts import RequestContext
from app.core.offsite_fund_contracts import ReceiveRecognizedMailRequest
from app.model.audit import InteractionAudit
from app.model.base import Base
from app.model.offsite_fund import (
OffsiteFundMail,
OffsiteMailCursor,
OffsiteNotification,
OffsiteRecognitionAttempt,
)
from app.service.offsite_document_recognition_adapter import StructuredRecognitionResult
from app.service.offsite_mail_adapter import (
RawMailAttachment,
RawMailMessage,
SavedMailAttachment,
SavedMailMessage,
)
from app.worker.offsite_mail_worker import OffsiteMailWorker
@pytest.mark.asyncio
async def test_worker_advances_uid_only_after_each_mail_succeeds(tmp_path: Path) -> None:
maker, engine = await _database()
receiver = FakeReceiver(_raw_mail("5"), _raw_mail("6"), _raw_mail("7"))
service = FakeService(fail_uid="6")
worker = OffsiteMailWorker(
_settings(),
receiver=receiver,
storage=FakeStorage(tmp_path),
recognizer=FakeRecognizer(),
identity_resolver=_identity,
session_factory=maker,
service_factory=lambda _session: service,
)
try:
assert await worker.run_once()
cursor = await _cursor(maker)
assert cursor.last_uid == "5"
assert cursor.blocked_uid == "6"
assert cursor.status == "failed"
assert receiver.fetch_last_uid == "0"
assert service.calls == ["5", "6"]
async with maker() as session, session.begin():
await session.execute(
update(OffsiteMailCursor).values(next_retry_at=None)
)
service.fail_uid = None
assert await worker.run_once()
cursor = await _cursor(maker)
assert cursor.last_uid == "7"
assert cursor.status == "idle"
assert cursor.blocked_uid is None
assert service.calls == ["5", "6", "6", "7"]
assert receiver.fetch_last_uid == "5"
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_worker_writes_alert_audit_when_cursor_becomes_blocked(tmp_path: Path) -> None:
"""游标锁死必须留下可见告警:前端"收件箱已停止收信"提示读的就是这条记录。"""
maker, engine = await _database()
service = FakeService(fail_uid="6")
worker = OffsiteMailWorker(
_settings(),
receiver=FakeReceiver(_raw_mail("6")),
storage=FakeStorage(tmp_path),
recognizer=FakeRecognizer(),
identity_resolver=_identity,
session_factory=maker,
service_factory=lambda _session: service,
)
try:
# 测试环境不需要真实退避等待,每轮清空重试时间以便连续推进失败计数。
for _ in range(3):
await worker.run_once()
async with maker() as session, session.begin():
await session.execute(
update(OffsiteMailCursor).values(next_retry_at=None)
)
cursor = await _cursor(maker)
assert cursor.status == "blocked"
assert cursor.retry_count == 3
async with maker() as session:
audits = (
await session.execute(
select(InteractionAudit).where(
InteractionAudit.action_type == "offsite.mail_cursor_blocked"
)
)
).scalars().all()
assert len(audits) == 1
assert audits[0].portal == "worker"
assert audits[0].detail["imap_uid"] == "6"
assert audits[0].detail["retry_count"] == 3
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_worker_does_not_write_without_configured_worker_identity(tmp_path: Path) -> None:
maker, engine = await _database()
receiver = FakeReceiver(_raw_mail("5"))
worker = OffsiteMailWorker(
_settings(offsite_worker_user_id=""),
receiver=receiver,
storage=FakeStorage(tmp_path),
recognizer=FakeRecognizer(),
identity_resolver=_identity,
session_factory=maker,
service_factory=lambda _session: FakeService(),
)
try:
assert not await worker.run_once()
assert receiver.health_calls == 0
async with maker() as session:
assert await session.scalar(select(OffsiteMailCursor.id)) is None
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_worker_rejects_real_imap_processing_with_mock_recognition() -> None:
receiver = FakeReceiver(_raw_mail("5"))
worker = OffsiteMailWorker(
_settings(offsite_worker_user_id="1"),
receiver=receiver,
)
assert not await worker.run_once()
assert receiver.health_calls == 0
@pytest.mark.asyncio
async def test_worker_wires_saved_attachment_recognition_into_business_service(
tmp_path: Path,
) -> None:
maker, engine = await _database()
service = FakeService()
worker = OffsiteMailWorker(
_settings(),
receiver=FakeReceiver(_raw_mail_with_attachment("5")),
storage=FakeStorage(tmp_path),
recognizer=SuccessfulRecognizer(),
identity_resolver=_identity,
session_factory=maker,
service_factory=lambda _session: service,
)
try:
assert await worker.run_once()
assert len(service.payloads) == 1
attachment = service.payloads[0].attachments[0]
assert attachment.document_type == "subscription"
assert attachment.original_file_path.endswith("5.eml")
assert attachment.extracted_fields["申请编号"] == "SUB-005"
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_worker_retries_incomplete_recognition_once_and_records_attempts(
tmp_path: Path,
) -> None:
maker, engine = await _database()
service = FakeService()
recognizer = FlakyRecognizer()
worker = OffsiteMailWorker(
_settings(),
receiver=FakeReceiver(_raw_mail_with_attachment("9")),
storage=FakeStorage(tmp_path),
recognizer=recognizer,
identity_resolver=_identity,
session_factory=maker,
service_factory=lambda _session: service,
)
try:
assert await worker.run_once()
assert recognizer.calls == 2
assert len(service.payloads) == 1
async with maker() as session:
attempts = (
await session.execute(
select(OffsiteRecognitionAttempt).order_by(
OffsiteRecognitionAttempt.attempt_no
)
)
).scalars().all()
assert len(attempts) == 2
assert [attempt.attempt_no for attempt in attempts] == [1, 2]
assert all(attempt.source == "automatic" for attempt in attempts)
assert attempts[0].status == "recognition_exception"
assert attempts[1].status == "success"
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_worker_recovers_stale_sending_notification_without_marking_success(
tmp_path: Path,
) -> None:
maker, engine = await _database()
worker = OffsiteMailWorker(
_settings(),
receiver=FakeReceiver(),
storage=FakeStorage(tmp_path),
recognizer=FakeRecognizer(),
session_factory=maker,
)
old = datetime(2026, 9, 1, 0, 0, 0)
try:
async with maker() as session, session.begin():
session.add(
OffsiteNotification(
id=1,
notification_type="mail_return",
business_key="20260901-001-A01",
receiver_id="15008108550@163.com",
operator_id="operator-001",
agent_draft="draft",
final_content="final",
payload={},
status="发送中",
retry_count=0,
created_at=old,
updated_at=old,
)
)
assert await worker.recover_stale_notifications()
async with maker() as session:
notification = await session.get(OffsiteNotification, 1)
assert notification is not None
assert notification.status == "发送失败"
assert notification.provider_message_id is None
assert notification.failure_reason is not None
assert "人工核验" in notification.failure_reason
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_worker_waits_in_idle_then_scans_new_uid_without_duplicate_fetch() -> None:
receiver = IdleCompensationReceiver(_raw_mail("8"))
worker = OffsiteMailWorker(
_settings(offsite_imap_idle_enabled=True),
receiver=receiver,
)
messages = await worker._fetch_with_idle_compensation("7")
assert [item.imap_uid for item in messages] == ["8"]
assert receiver.fetch_last_uids == ["7", "7"]
assert receiver.idle_wait_calls == [120]
@pytest.mark.asyncio
async def test_worker_idle_timeout_returns_without_second_scan() -> None:
receiver = IdleCompensationReceiver()
receiver.idle_result = False
worker = OffsiteMailWorker(
_settings(offsite_imap_idle_enabled=True),
receiver=receiver,
)
messages = await worker._fetch_with_idle_compensation("7")
assert messages == ()
assert receiver.fetch_last_uids == ["7"]
assert receiver.idle_wait_calls == [120]
class FakeReceiver:
def __init__(self, *messages: RawMailMessage) -> None:
self.messages = messages
self.last_scanned_uid: str | None = "0"
self.fetch_last_uid: str | None = None
self.health_calls = 0
def health_check(self) -> dict[str, object]:
self.health_calls += 1
return {"status": "ok"}
def fetch_since(self, last_uid: str | None, *, limit: int) -> tuple[RawMailMessage, ...]:
del limit
self.fetch_last_uid = last_uid
self.last_scanned_uid = self.messages[-1].imap_uid if self.messages else last_uid
return self.messages
def close(self) -> None:
return None
class IdleCompensationReceiver:
def __init__(self, *messages: RawMailMessage) -> None:
self.messages = messages
self.last_scanned_uid: str | None = None
self.fetch_last_uids: list[str | None] = []
self.idle_wait_calls: list[float] = []
self.idle_result = True
self._fetch_count = 0
def health_check(self) -> dict[str, object]:
return {"status": "ok"}
def fetch_since(self, last_uid: str | None, *, limit: int) -> tuple[RawMailMessage, ...]:
del limit
self._fetch_count += 1
self.fetch_last_uids.append(last_uid)
if self._fetch_count == 1:
self.last_scanned_uid = None
return ()
self.last_scanned_uid = self.messages[-1].imap_uid if self.messages else last_uid
return self.messages
def wait_for_new_mail(self, timeout_seconds: float) -> bool:
self.idle_wait_calls.append(timeout_seconds)
return self.idle_result
def close(self) -> None:
return None
class FakeStorage:
def __init__(self, root: Path) -> None:
self.root = root
def save(self, mail: RawMailMessage) -> SavedMailMessage:
path = self.root / f"{mail.imap_uid}.eml"
path.write_bytes(mail.raw_message)
return SavedMailMessage(
imap_uid=mail.imap_uid,
message_id=mail.message_id,
sender=mail.sender,
return_path=mail.return_path,
auth_result=mail.auth_result,
eml_path=str(path),
attachments=(
SavedMailAttachment(
filename="empty.txt",
file_hash="a" * 64,
media_type="text/plain",
size_bytes=0,
original_file_path=str(path),
),
) if mail.attachments else (),
)
class FakeRecognizer:
async def recognize(self, source: Any) -> Any:
del source
raise AssertionError("本测试邮件没有附件,不应调用识别器")
class SuccessfulRecognizer:
async def recognize(self, source: Any) -> StructuredRecognitionResult:
del source
return StructuredRecognitionResult(
document_type="subscription",
extracted_fields={
"基金代码": "000001",
"基金名称": "测试基金",
"账户标识": "ACCT-001",
"投资者名称": "测试客户",
"申请编号": "SUB-005",
"申请日期": "2026-09-10",
"代销机构": "测试代销",
"申购金额": "10000",
"金额单位": "元",
},
field_confidence={"申请编号": Decimal("0.99")},
missing_fields=(),
low_confidence_fields=(),
page_evidence={"申请编号": [{"page": 1}]},
ocr_text="申请编号:SUB-005",
ocr_status="mock",
llm_status="mock",
)
class FlakyRecognizer:
def __init__(self) -> None:
self.calls = 0
async def recognize(self, source: Any) -> StructuredRecognitionResult:
del source
self.calls += 1
if self.calls == 1:
return StructuredRecognitionResult(
document_type="subscription",
extracted_fields={"基金代码": "000001"},
field_confidence={"基金代码": Decimal("0.99")},
missing_fields=("申购金额",),
low_confidence_fields=(),
page_evidence={},
ocr_text="申购",
ocr_status="success",
llm_status="success",
)
return await SuccessfulRecognizer().recognize(None)
class FakeService:
def __init__(self, fail_uid: str | None = None) -> None:
self.fail_uid = fail_uid
self.calls: list[str] = []
self.payloads: list[ReceiveRecognizedMailRequest] = []
async def receive_recognized_mail(
self, payload: ReceiveRecognizedMailRequest, context: RequestContext
) -> dict[str, object]:
del context
self.calls.append(payload.imap_uid)
self.payloads.append(payload)
if payload.imap_uid == self.fail_uid:
return {"code": 500, "message": "模拟入库失败", "data": {}}
return {"code": 0, "message": "ok", "data": {"business": False}}
async def _database() -> tuple[async_sessionmaker[AsyncSession], Any]:
engine = create_async_engine(
"sqlite+aiosqlite://",
connect_args={"check_same_thread": False},
poolclass=StaticPool,
)
async with engine.begin() as connection:
await connection.run_sync(
lambda sync: Base.metadata.create_all(
sync,
tables=[
OffsiteFundMail.__table__,
OffsiteMailCursor.__table__,
OffsiteNotification.__table__,
OffsiteRecognitionAttempt.__table__,
],
)
)
await connection.execute(
text(
"""
CREATE TABLE interaction_audit (
id INTEGER PRIMARY KEY AUTOINCREMENT,
actor_type VARCHAR(16) NOT NULL,
actor_id BIGINT NULL,
target_customer_id BIGINT NULL,
session_id VARCHAR(64) NULL,
portal VARCHAR(32) NULL,
action_type VARCHAR(64) NOT NULL,
detail JSON NOT NULL,
created_at DATETIME NOT NULL
)
"""
)
)
return async_sessionmaker(engine, expire_on_commit=False), engine
async def _identity(identity: RequestContext) -> RequestContext:
return identity.model_copy(
update={
"roles": ("operator",),
"permissions": ("offsite:write",),
"data_scope": "all",
}
)
def _settings(**updates: object) -> Settings:
# 外部识别开关必须显式给值:`app/core/config.py` 导入时用 load_dotenv 把本机 `.env`
# 灌进了 os.environ,仅靠 `_env_file=None` 挡不住环境变量;构造参数优先级最高,
# 否则"识别未就绪应拒绝处理"的用例会被本机开关带成"已就绪"。
values: dict[str, object] = {
"jwt_issuer": "jr-local",
"jwt_audience": "jr-agent-platform",
"mysql_dsn": "sqlite+aiosqlite:///test.db",
"redis_url": "redis://127.0.0.1:6379/0",
"milvus_uri": "http://127.0.0.1:19530",
"neo4j_uri": "bolt://127.0.0.1:7687",
"offsite_mail_worker_enabled": True,
"offsite_imap_enabled": True,
"offsite_imap_idle_timeout_seconds": 120,
"offsite_worker_user_id": "1",
"offsite_max_retry_count": 3,
"offsite_ocr_enabled": False,
"offsite_deepseek_enabled": False,
"offsite_deepseek_api_key": "test-deepseek-key",
}
values.update(updates)
return Settings(_env_file=None, **values)
def _raw_mail(uid: str) -> RawMailMessage:
return RawMailMessage(
imap_uid=uid,
message_id=f"<{uid}@worker.test>",
sender="15008108550@163.com",
return_path="15008108550@163.com",
auth_result={"spf": "pass"},
raw_message=f"mail-{uid}".encode(),
received_at=datetime(2026, 9, 10, 9, 30, 0),
attachments=(),
)
def _raw_mail_with_attachment(uid: str) -> RawMailMessage:
mail = _raw_mail(uid)
return RawMailMessage(
imap_uid=mail.imap_uid,
message_id=mail.message_id,
sender=mail.sender,
return_path=mail.return_path,
auth_result=mail.auth_result,
raw_message=mail.raw_message,
received_at=mail.received_at,
attachments=(
RawMailAttachment(
filename="申购申请单.pdf",
media_type="application/pdf",
payload=b"pdf",
),
),
)
async def _cursor(maker: async_sessionmaker[AsyncSession]) -> OffsiteMailCursor:
async with maker() as session:
cursor = await session.scalar(select(OffsiteMailCursor))
assert cursor is not None
return cursor
+24 -1
View File
@@ -4,7 +4,7 @@ from unittest.mock import AsyncMock, Mock
import pytest
from app.model.platform import DomainEventOutbox
from app.worker.outbox_worker import OutboxWorker
from app.worker.outbox_worker import OutboxHandlerError, OutboxWorker
def event() -> DomainEventOutbox:
@@ -66,9 +66,32 @@ async def test_handler_failure_records_attempt(attempts: int, expected: str) ->
assert row.status == expected
assert row.retry_count == attempts + 1
assert "do-not-log" not in (row.last_error or "")
# 其他异常一律只记类名:消息可能含凭据/SQL/客户标识,不落库。
assert row.last_error == "RuntimeError"
session.commit.assert_awaited_once()
@pytest.mark.asyncio
async def test_handler_error_reason_is_recorded() -> None:
"""`OutboxHandlerError` 的固定文案可以落库 —— 否则失败原因根本无从分辨。
实测库里 373 条死信的 `last_error` 全是裸的 `"ValueError"`:`dispatch`(run not
found)、`dispatch_run_completed`、`dispatch_memory_extraction`、
`dispatch_profile_rebuild` 抛的都是 ValueError,只记类名等于把"哪一处失败"也丢了。
这些 handler 现在改抛 `OutboxHandlerError`,它的 reason 由代码写死、不含请求数据。
"""
row = event()
session = AsyncMock()
session.add = Mock()
session.scalar.side_effect = [row, None]
handler = AsyncMock(side_effect=OutboxHandlerError("run not found"))
await OutboxWorker(session, {"test.event": handler}).publish_one()
assert row.last_error == "OutboxHandlerError: run not found"
assert row.status == "failed"
@pytest.mark.asyncio
async def test_empty_matching_queue_releases_transaction() -> None:
session = AsyncMock()