From 9d1a6a11b0e1eb6933dc080c3df54942bfcf1195 Mon Sep 17 00:00:00 2001 From: Hermes Agent Date: Sun, 14 Jun 2026 20:37:39 +0800 Subject: [PATCH] feat: harden delivery and deletion workflows --- README.md | 13 +++ admin/routes/orders.py | 41 +++++++ admin/tests/test_order_deletion.py | 145 ++++++++++++++++++++++++ data/notifications/dispatcher.py | 69 ++++++++++++ data/notifications/email_service.py | 44 +++++++- data/orders/deletion_service.py | 148 +++++++++++++++++++++++++ docs/ACTIVE_REMEDIATION_2026-06-13.md | 25 +++-- docs/CROWD_DB_DATA_QUALITY.md | 4 +- docs/DATA_RETENTION_AND_DELETION.md | 8 +- docs/DELIVERY_SERVICE_DESIGN.md | 19 ++-- scripts/gaokao-delivery-dispatch.py | 37 +++++++ tests/test_delivery_dispatcher.py | 152 ++++++++++++++++++++++++++ 12 files changed, 680 insertions(+), 25 deletions(-) create mode 100644 admin/tests/test_order_deletion.py create mode 100644 data/notifications/dispatcher.py create mode 100644 data/orders/deletion_service.py create mode 100644 scripts/gaokao-delivery-dispatch.py create mode 100644 tests/test_delivery_dispatcher.py diff --git a/README.md b/README.md index 7669cc8..79f917c 100644 --- a/README.md +++ b/README.md @@ -154,6 +154,19 @@ bash scripts/dev-verify.sh GAOKAO_SKIP_INSTALL=1 bash scripts/dev-verify.sh ``` +### T12 交付事件最小执行链 + +当前 `report_ready` 不再只停留在事件表: + +- 事件表:`delivery_notifications` +- 执行器:`data.notifications.dispatcher.DeliveryDispatcher` +- CLI:`python3 scripts/gaokao-delivery-dispatch.py --channel station` + +最小语义: + +- 交付物齐全(HTML/PDF 存在)→ `ready -> sent` +- 交付物缺失 → `failed`,并写入 `failure_reason` + ### T6.7 Docker Compose 一键启动 仓库根目录已提供 `Dockerfile`、`docker-compose.yml` 与 `.env.docker.example`。默认镜像会把运行数据写入容器外部卷 `/var/lib/gaokao`,避免覆盖仓库里的 Python 包 `data/`。默认 compose 只绑定 `127.0.0.1` 且以 `dev` 模式启动,适合本机自测;正式部署前请复制 `.env.docker.example` 到 `.env` 并替换密钥/密码。 diff --git a/admin/routes/orders.py b/admin/routes/orders.py index bb374e9..77ea3b2 100644 --- a/admin/routes/orders.py +++ b/admin/routes/orders.py @@ -29,6 +29,7 @@ from admin.errors import ( ) from admin.errors.exceptions import BusinessError from data.orders.dao import DuplicateOrder, OrderNotFound, OrdersDAO +from data.orders.deletion_service import OrderDeletionService from data.orders.intake_store import IntakeStore from data.orders.models import Order, generate_order_id from data.orders.state_machine import InvalidStateTransition, next_states @@ -143,6 +144,12 @@ class OrderMutationResponse(OrderDetailPayload): action: str +class OrderDeletionResponse(BaseModel): + action: str + order_id: str + files_deleted: int = 0 + + class CreateOrderRequest(BaseModel): source: OrderSource external_id: Optional[str] = None @@ -450,3 +457,37 @@ def patch_order( finally: intake_store.close() return {"action": action, **_detail_payload(dao, order, intake=intake)} + + +@router.delete( + "/{order_id}", + response_model=OrderDeletionResponse, + summary="订单删除 / 匿名化(T12/A-4)", +) +def delete_or_anonymize_order( + order_id: str = Path(..., min_length=1), + mode: Literal["delete", "anonymize"] = Query("delete"), + reason: str = Query(..., min_length=1), + settings: Settings = Depends(get_settings_dep), + current_user: AdminUser = Depends(get_current_user), +) -> dict[str, Any]: + service = OrderDeletionService.for_db(settings.orders_db_path) + try: + try: + if mode == "delete": + result = service.delete_order( + order_id, + actor=current_user.username, + reason=reason, + ) + else: + result = service.anonymize_order( + order_id, + actor=current_user.username, + reason=reason, + ) + except OrderNotFound as exc: + raise _business_error_for_lookup(order_id) from exc + finally: + service.close() + return result.__dict__ diff --git a/admin/tests/test_order_deletion.py b/admin/tests/test_order_deletion.py new file mode 100644 index 0000000..fe44b5c --- /dev/null +++ b/admin/tests/test_order_deletion.py @@ -0,0 +1,145 @@ +from __future__ import annotations + +from pathlib import Path + +from data.notifications.email_service import DeliveryNotificationService +from data.orders.dao import OrderNotFound, OrdersDAO +from data.orders.deletion_service import OrderDeletionService +from data.orders.intake_store import IntakeStore +from data.orders.models import Order +from data.payments.service import PaymentService + + +def _seed_order(db_path: str, order_id: str = "GKO-20260614-DELETE") -> Order: + order = Order( + id=order_id, + source="web", + service_version="standard", + amount_cents=9900, + status="pending", + customer_name="张家长", + customer_phone="13800138000", + customer_wechat="wx-parent-01", + candidate_name="张三", + candidate_province="湖南", + notes="需要删除的测试订单", + ) + with OrdersDAO.connect(db_path) as dao: + return dao.create(order, actor="test", reason="seed") + + +def _mark_paid(settings, order: Order) -> None: + service = PaymentService.for_db( + settings.orders_db_path, + base_url=settings.payment_base_url, + webhook_secret=settings.payment_webhook_secret, + ) + checkout = service.create_checkout(order.id, portal_token="portal-token") + payload, headers = service.provider.build_webhook_request( + payment_id=checkout.payment_id, + amount_cents=order.amount_cents, + provider_trade_no=f"MOCK-{order.id}", + ) + service.handle_webhook(payload, headers["X-Mock-Signature"]) + + +def _prepare_order_with_artifacts( + settings, tmp_path: Path, order_id: str +) -> tuple[Path, Path]: + report_path = tmp_path / f"{order_id}-report.html" + pdf_path = tmp_path / f"{order_id}-report.pdf" + report_path.write_text("

report

", encoding="utf-8") + pdf_path.write_bytes(b"%PDF-1.4\ndelete\n") + with OrdersDAO.connect(settings.orders_db_path) as dao: + dao.update( + order_id, + {"audit_report": str(report_path), "pdf_path": str(pdf_path)}, + actor="test", + reason="attach_report", + ) + dao.transition_status(order_id, "serving", actor="test", reason="processing") + dao.transition_status( + order_id, "delivered", actor="test", reason="report_ready" + ) + return report_path, pdf_path + + +def test_admin_delete_order_removes_artifacts_and_related_records( + client, auth_headers, settings, tmp_path +): + order = _seed_order(settings.orders_db_path) + _mark_paid(settings, order) + IntakeStore.for_db(settings.orders_db_path).save( + order_id=order.id, + payload={"candidate_score": 578, "guardian_notes": "to delete"}, + submit=True, + ) + report_path, pdf_path = _prepare_order_with_artifacts(settings, tmp_path, order.id) + + notification_service = DeliveryNotificationService.for_db(settings.orders_db_path) + try: + assert len(notification_service.list_events(order.id)) == 1 + finally: + notification_service.close() + + resp = client.delete( + f"/api/orders/{order.id}?mode=delete&reason=user_request", + headers=auth_headers, + ) + assert resp.status_code == 200, resp.text + body = resp.json() + assert body["action"] == "deleted" + assert body["order_id"] == order.id + assert body["files_deleted"] == 2 + + assert not report_path.exists() + assert not pdf_path.exists() + + with OrdersDAO.connect(settings.orders_db_path) as dao: + try: + dao.get(order.id) + except OrderNotFound: + pass + else: + raise AssertionError("order should be deleted") + + notification_service = DeliveryNotificationService.for_db(settings.orders_db_path) + try: + assert notification_service.list_events(order.id) == [] + finally: + notification_service.close() + + deletion_service = OrderDeletionService.for_db(settings.orders_db_path) + try: + assert deletion_service.audit_count(order.id) == 1 + finally: + deletion_service.close() + + +def test_admin_anonymize_order_masks_pii_but_keeps_order( + client, auth_headers, settings +): + order = _seed_order(settings.orders_db_path, order_id="GKO-20260614-ANON") + + resp = client.delete( + f"/api/orders/{order.id}?mode=anonymize&reason=retention_expired", + headers=auth_headers, + ) + assert resp.status_code == 200, resp.text + body = resp.json() + assert body["action"] == "anonymized" + assert body["order_id"] == order.id + + with OrdersDAO.connect(settings.orders_db_path) as dao: + anonymized = dao.get(order.id) + assert anonymized.customer_name == "已匿名化" + assert anonymized.customer_phone is None + assert anonymized.customer_wechat is None + assert anonymized.candidate_name == "匿名考生" + assert anonymized.notes == "[ANONYMIZED] retention_expired" + + deletion_service = OrderDeletionService.for_db(settings.orders_db_path) + try: + assert deletion_service.audit_count(order.id) == 1 + finally: + deletion_service.close() diff --git a/data/notifications/dispatcher.py b/data/notifications/dispatcher.py new file mode 100644 index 0000000..e74b360 --- /dev/null +++ b/data/notifications/dispatcher.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +import json +from dataclasses import dataclass +from pathlib import Path + +from data.notifications.email_service import ( + DeliveryNotificationEvent, + DeliveryNotificationService, +) + + +@dataclass +class DispatchResult: + processed: int = 0 + sent: int = 0 + failed: int = 0 + + +class DeliveryDispatcher: + def __init__(self, service: DeliveryNotificationService) -> None: + self._service = service + + @classmethod + def for_db(cls, db_path: str) -> "DeliveryDispatcher": + return cls(DeliveryNotificationService.for_db(db_path)) + + def close(self) -> None: + self._service.close() + + def dispatch_ready_events( + self, + *, + channel: str = "station", + statuses: tuple[str, ...] = ("ready", "failed"), + limit: int = 100, + ) -> DispatchResult: + result = DispatchResult() + events = self._service.list_pending_events( + channel=channel, + statuses=statuses, + limit=limit, + ) + for event in events: + result.processed += 1 + failure_reason = self._validate_event(event) + if failure_reason is not None: + self._service.mark_failed( + event.order_id, failure_reason, event_type=event.event_type + ) + result.failed += 1 + continue + self._service.mark_sent(event.order_id, event_type=event.event_type) + result.sent += 1 + return result + + @staticmethod + def _validate_event(event: DeliveryNotificationEvent) -> str | None: + try: + payload = json.loads(event.payload_json) + except json.JSONDecodeError: + return "delivery payload invalid" + report_path = payload.get("audit_report") or payload.get("plan_file") + pdf_path = payload.get("pdf_path") + if not report_path or not Path(str(report_path)).is_file(): + return "delivery artifact missing" + if not pdf_path or not Path(str(pdf_path)).is_file(): + return "delivery artifact missing" + return None diff --git a/data/notifications/email_service.py b/data/notifications/email_service.py index 930e27c..e53e500 100644 --- a/data/notifications/email_service.py +++ b/data/notifications/email_service.py @@ -123,9 +123,43 @@ class DeliveryNotificationService: ) self._conn.commit() - def list_events(self, order_id: str) -> list[DeliveryNotificationEvent]: - rows = self._conn.execute( - "SELECT order_id, event_type, channel, payload_json, status, attempt_count, last_attempt_at, failure_reason, created_at FROM delivery_notifications WHERE order_id=? ORDER BY id ASC", - (order_id,), - ).fetchall() + def list_events( + self, + order_id: str, + *, + status: str | None = None, + channel: str | None = None, + ) -> list[DeliveryNotificationEvent]: + sql = "SELECT order_id, event_type, channel, payload_json, status, attempt_count, last_attempt_at, failure_reason, created_at FROM delivery_notifications WHERE order_id=?" + params: list[object] = [order_id] + if status is not None: + sql += " AND status=?" + params.append(status) + if channel is not None: + sql += " AND channel=?" + params.append(channel) + sql += " ORDER BY id ASC" + rows = self._conn.execute(sql, tuple(params)).fetchall() + return [DeliveryNotificationEvent(**dict(row)) for row in rows] + + def list_pending_events( + self, + *, + channel: str | None = None, + statuses: tuple[str, ...] = ("ready",), + limit: int = 100, + ) -> list[DeliveryNotificationEvent]: + placeholders = ",".join("?" for _ in statuses) + sql = ( + "SELECT order_id, event_type, channel, payload_json, status, attempt_count, last_attempt_at, failure_reason, created_at FROM delivery_notifications WHERE status IN (" + + placeholders + + ")" + ) + params: list[object] = list(statuses) + if channel is not None: + sql += " AND channel=?" + params.append(channel) + sql += " ORDER BY id ASC LIMIT ?" + params.append(limit) + rows = self._conn.execute(sql, tuple(params)).fetchall() return [DeliveryNotificationEvent(**dict(row)) for row in rows] diff --git a/data/orders/deletion_service.py b/data/orders/deletion_service.py new file mode 100644 index 0000000..042f4b0 --- /dev/null +++ b/data/orders/deletion_service.py @@ -0,0 +1,148 @@ +from __future__ import annotations + +import sqlite3 +from dataclasses import dataclass +from pathlib import Path + +from data.orders.dao import OrderNotFound, OrdersDAO +from data.orders.models import utc_now_iso +from data.orders.schema import apply_schema + + +AUDIT_SCHEMA_SQL = """ +CREATE TABLE IF NOT EXISTS order_deletion_audits ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + order_id TEXT NOT NULL, + action TEXT NOT NULL CHECK(action IN ('delete','anonymize')), + actor TEXT, + reason TEXT, + files_deleted INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_order_deletion_audits_order_id ON order_deletion_audits(order_id); +""" + + +@dataclass +class DeletionResult: + order_id: str + action: str + files_deleted: int = 0 + + +class OrderDeletionService: + def __init__(self, conn: sqlite3.Connection) -> None: + self._conn = conn + self._conn.executescript(AUDIT_SCHEMA_SQL) + self._conn.commit() + + @classmethod + def for_db(cls, db_path: str | Path) -> "OrderDeletionService": + conn = apply_schema(db_path) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA foreign_keys = ON") + return cls(conn) + + def close(self) -> None: + self._conn.close() + + def delete_order(self, order_id: str, *, actor: str, reason: str) -> DeletionResult: + with OrdersDAO(self._conn) as dao: + try: + order = dao.get(order_id) + except OrderNotFound: + raise + files_deleted = self._delete_artifacts(order.audit_report, order.pdf_path) + deleted = dao.delete(order_id) + if not deleted: + raise OrderNotFound(f"订单不存在: {order_id}") + self._insert_audit( + order_id=order_id, + action="delete", + actor=actor, + reason=reason, + files_deleted=files_deleted, + ) + self._conn.commit() + return DeletionResult( + order_id=order_id, action="deleted", files_deleted=files_deleted + ) + + def anonymize_order( + self, order_id: str, *, actor: str, reason: str + ) -> DeletionResult: + with OrdersDAO(self._conn) as dao: + try: + dao.get(order_id) + except OrderNotFound: + raise + now = utc_now_iso() + self._conn.execute( + """ + UPDATE orders + SET customer_name=?, + customer_phone_enc=NULL, + customer_phone_hash=NULL, + customer_wechat=NULL, + candidate_name=?, + candidate_id_card_enc=NULL, + candidate_interests=NULL, + candidate_strong_subjects=NULL, + candidate_weak_subjects=NULL, + candidate_family=NULL, + notes=?, + status_updated_at=? + WHERE id=? + """, + ( + "已匿名化", + "匿名考生", + f"[ANONYMIZED] {reason}", + now, + order_id, + ), + ) + self._insert_audit( + order_id=order_id, + action="anonymize", + actor=actor, + reason=reason, + files_deleted=0, + ) + self._conn.commit() + return DeletionResult( + order_id=order_id, action="anonymized", files_deleted=0 + ) + + def audit_count(self, order_id: str) -> int: + row = self._conn.execute( + "SELECT COUNT(*) FROM order_deletion_audits WHERE order_id=?", + (order_id,), + ).fetchone() + return int(row[0] if row is not None else 0) + + def _insert_audit( + self, + *, + order_id: str, + action: str, + actor: str, + reason: str, + files_deleted: int, + ) -> None: + self._conn.execute( + "INSERT INTO order_deletion_audits(order_id, action, actor, reason, files_deleted, created_at) VALUES (?, ?, ?, ?, ?, ?)", + (order_id, action, actor, reason, files_deleted, utc_now_iso()), + ) + + @staticmethod + def _delete_artifacts(*paths: str | None) -> int: + deleted = 0 + for raw_path in paths: + if not raw_path: + continue + path = Path(raw_path) + if path.is_file(): + path.unlink() + deleted += 1 + return deleted diff --git a/docs/ACTIVE_REMEDIATION_2026-06-13.md b/docs/ACTIVE_REMEDIATION_2026-06-13.md index ad6afcb..50a14f2 100644 --- a/docs/ACTIVE_REMEDIATION_2026-06-13.md +++ b/docs/ACTIVE_REMEDIATION_2026-06-13.md @@ -98,11 +98,13 @@ - `report_ready` 事件自动落库 - `status / attempt_count / last_attempt_at / failure_reason` 追踪字段 - 站内查看 + PDF 下载最小交付闭环 +- `DeliveryDispatcher` + `scripts/gaokao-delivery-dispatch.py` +- `ready -> sent` / 缺文件 `-> failed` 的最小执行链 仍缺: - 真实邮件/站内通知发送执行器 -- 自动重试 worker / 告警链 +- 自动重试调度 / 告警链 - 面向用户的独立通知审计页 - 对账/退款与交付失败补偿联动 @@ -128,16 +130,19 @@ - `docs/DATA_RETENTION_AND_DELETION.md` - 服务条款草案、删除执行 SOP - 前台资料提交同意字段落库 +- 后台 `DELETE /api/orders/{id}?mode=delete|anonymize&reason=...` 最小执行入口 +- 删除时自动清理报告 HTML/PDF,并写入 `order_deletion_audits` 仍缺: - 正式法务审定版本 -- 删除/匿名化后台执行入口 +- 前台/客服删除工单流程 +- 数据保留期自动清理任务 - 合规文本上线前最终校对 说明: -- 当前不再是“完全没有合规基线”,而是“文档基线已建,执行与审定仍待完成”。 +- 当前不再是“完全没有合规基线”,而是“文档基线 + 后台最小执行入口已建,前台流程与自动清理仍待完成”。 #### A-5 业务数据备份/恢复/密钥托管已形成基线,但生产化仍未完结 @@ -168,22 +173,24 @@ #### B-1 crowd_db 27 省“结构覆盖”不等于“高质量覆盖” 严重度: P1 -当前状态: 未解决 +当前状态: 部分解决 事实: - 27 省 JSON 文件已存在 -- 但高置信、高密度推荐数据目前重点仍在湖南 +- `risk_report` 已输出 `quality_level / quality_label` +- `gaokao-data-trace` human 输出已展示质量等级 +- 低置信省份仍主要是“结构覆盖已完成,但推荐质量待增强” 风险: -- 如果对外统一宣传“27省高质量反扎堆推荐”,有误导风险 +- 若对外统一宣传“27省高质量反扎堆推荐”,仍有误导风险 建议动作: -- 在报告输出里继续强化 confidence 文案 -- 新增“数据完整度等级”字段 -- 对低置信省份降级展示,不输出强结论 +- 在 README / 产品文案中继续避免“27省高质量推荐全覆盖”表述 +- 后续补 province-level completeness/quality summary +- 继续提升非湖南省份高置信数据密度 #### B-2 本地验证环境仍未完全固化为单命令体验 diff --git a/docs/CROWD_DB_DATA_QUALITY.md b/docs/CROWD_DB_DATA_QUALITY.md index 0abef5a..fb54003 100644 --- a/docs/CROWD_DB_DATA_QUALITY.md +++ b/docs/CROWD_DB_DATA_QUALITY.md @@ -80,4 +80,6 @@ 当前状态: - **文档基线已补齐** -- **代码主链分级展示仍待继续实现** +- **`risk_report` 已输出 `quality_level / quality_label`** +- **`gaokao-data-trace --human` 已输出质量等级** +- **低置信省份已通过 confidence 分级与 warning 机制降级** diff --git a/docs/DATA_RETENTION_AND_DELETION.md b/docs/DATA_RETENTION_AND_DELETION.md index 204f956..88418f3 100644 --- a/docs/DATA_RETENTION_AND_DELETION.md +++ b/docs/DATA_RETENTION_AND_DELETION.md @@ -40,13 +40,15 @@ - `OrdersDAO.delete()` 删除订单能力 - ON DELETE CASCADE 清理部分关联表 +- 后台 `DELETE /api/orders/{id}?mode=delete|anonymize&reason=...` 最小执行入口 +- 删除时自动清理 HTML/PDF 报告文件 +- `order_deletion_audits` 最小审计表 尚缺: -- 前台/后台删除工单流程 -- 报告文件自动清理脚本 -- 删除操作审计记录 +- 前台/客服删除工单流程 - 数据保留期自动清理任务 +- 更细粒度的匿名化策略(如案例长期保留脱敏版) ## 5. MVP 最低执行要求 diff --git a/docs/DELIVERY_SERVICE_DESIGN.md b/docs/DELIVERY_SERVICE_DESIGN.md index 5a73f67..12b86c0 100644 --- a/docs/DELIVERY_SERVICE_DESIGN.md +++ b/docs/DELIVERY_SERVICE_DESIGN.md @@ -52,13 +52,16 @@ - `delivery_notifications` 事件表 - portal 状态页 / 报告页 / PDF 下载页 - `delivered` 但无交付物时不再误报 `report_ready` +- `data.notifications.dispatcher.DeliveryDispatcher` +- `scripts/gaokao-delivery-dispatch.py` +- `ready -> sent` / `缺文件 -> failed` / `attempt_count` 递增 未完成: -- 通知触发点下沉到稳定主链 - 独立 `delivery_job` / `delivery_attempt` 模型 -- 重试与失败原因追踪 -- 多通道统一投递状态 +- 邮件/渠道真实发送执行器 +- 自动重试调度与告警链 +- 多通道统一投递状态页 ## 6. 当前最短闭环 @@ -67,10 +70,12 @@ 3. portal 显示 `report_ready` 4. 状态页提供查看/下载入口 5. `report_ready` 事件落库且幂等 +6. `gaokao-delivery-dispatch.py --channel station` 可把 ready 事件推进到 sent +7. 缺失交付物时事件转 failed,并记录 `failure_reason` ## 7. 下一步实施建议 -1. 统一 `report_ready` 触发点 -2. 增加 `delivery_status` 最小字段 -3. 增加失败重试/失败原因 -4. 再决定是否扩到邮件通道 +1. 把 dispatcher 接入定时任务或 systemd timer +2. 增加邮件或站内通知真实发送器(二选一先闭环) +3. 增加失败重试阈值与告警 +4. 再决定是否扩到多通道 diff --git a/scripts/gaokao-delivery-dispatch.py b/scripts/gaokao-delivery-dispatch.py new file mode 100644 index 0000000..31e99d4 --- /dev/null +++ b/scripts/gaokao-delivery-dispatch.py @@ -0,0 +1,37 @@ +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parent.parent +if str(ROOT) not in sys.path: + sys.path.insert(0, str(ROOT)) + + +def main() -> int: + from admin.config import load_settings + from data.notifications.dispatcher import DeliveryDispatcher + + parser = argparse.ArgumentParser( + description="Dispatch delivery notification events" + ) + parser.add_argument("--channel", default="station") + parser.add_argument("--limit", type=int, default=100) + args = parser.parse_args() + + settings = load_settings() + dispatcher = DeliveryDispatcher.for_db(settings.orders_db_path) + try: + result = dispatcher.dispatch_ready_events( + channel=args.channel, limit=args.limit + ) + finally: + dispatcher.close() + print(json.dumps(result.__dict__, ensure_ascii=False, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_delivery_dispatcher.py b/tests/test_delivery_dispatcher.py new file mode 100644 index 0000000..bf23ff4 --- /dev/null +++ b/tests/test_delivery_dispatcher.py @@ -0,0 +1,152 @@ +from __future__ import annotations + +import subprocess +from pathlib import Path + +from data.notifications.email_service import DeliveryNotificationService +from data.orders.dao import OrdersDAO +from data.orders.intake_store import IntakeStore +from data.orders.models import Order +from data.payments.service import PaymentService + + +PROJECT_ROOT = Path(__file__).resolve().parents[1] + + +def _seed_order(db_path: str, order_id: str = "GKO-20260614-DISPATCH") -> Order: + order = Order( + id=order_id, + source="web", + service_version="standard", + amount_cents=9900, + status="pending", + customer_name="张家长", + customer_phone="13800138000", + candidate_name="张三", + candidate_province="湖南", + ) + with OrdersDAO.connect(db_path) as dao: + return dao.create(order, actor="test", reason="seed") + + +def _mark_paid(settings, order: Order) -> None: + service = PaymentService.for_db( + settings.orders_db_path, + base_url=settings.payment_base_url, + webhook_secret=settings.payment_webhook_secret, + ) + checkout = service.create_checkout(order.id, portal_token="portal-token") + payload, headers = service.provider.build_webhook_request( + payment_id=checkout.payment_id, + amount_cents=order.amount_cents, + provider_trade_no=f"MOCK-{order.id}", + ) + service.handle_webhook(payload, headers["X-Mock-Signature"]) + + +def _attach_ready_delivery(settings, tmp_path: Path, order_id: str) -> None: + report_path = tmp_path / f"{order_id}-report.html" + pdf_path = tmp_path / f"{order_id}-report.pdf" + report_path.write_text("

ready

", encoding="utf-8") + pdf_path.write_bytes(b"%PDF-1.4\nready\n") + with OrdersDAO.connect(settings.orders_db_path) as dao: + dao.update( + order_id, + {"audit_report": str(report_path), "pdf_path": str(pdf_path)}, + actor="test", + reason="attach_report", + ) + dao.transition_status(order_id, "serving", actor="test", reason="processing") + dao.transition_status( + order_id, "delivered", actor="test", reason="report_ready" + ) + + +def test_dispatch_ready_station_event_marks_sent(settings, tmp_path): + order = _seed_order(settings.orders_db_path) + _mark_paid(settings, order) + IntakeStore.for_db(settings.orders_db_path).save( + order_id=order.id, payload={"candidate_score": 578}, submit=True + ) + _attach_ready_delivery(settings, tmp_path, order.id) + + from data.notifications.dispatcher import DeliveryDispatcher + + dispatcher = DeliveryDispatcher.for_db(settings.orders_db_path) + try: + result = dispatcher.dispatch_ready_events(channel="station") + finally: + dispatcher.close() + + assert result.processed == 1 + assert result.sent == 1 + assert result.failed == 0 + + notification_service = DeliveryNotificationService.for_db(settings.orders_db_path) + try: + event = notification_service.list_events(order.id)[0] + finally: + notification_service.close() + assert event.status == "sent" + assert event.attempt_count == 1 + + +def test_dispatch_ready_station_event_marks_failed_when_pdf_missing(settings, tmp_path): + order = _seed_order(settings.orders_db_path, order_id="GKO-20260614-DISPATCH-FAIL") + _mark_paid(settings, order) + IntakeStore.for_db(settings.orders_db_path).save( + order_id=order.id, payload={"candidate_score": 578}, submit=True + ) + _attach_ready_delivery(settings, tmp_path, order.id) + + with OrdersDAO.connect(settings.orders_db_path) as dao: + current = dao.get(order.id) + assert current.pdf_path is not None + Path(current.pdf_path).unlink() + + from data.notifications.dispatcher import DeliveryDispatcher + + dispatcher = DeliveryDispatcher.for_db(settings.orders_db_path) + try: + result = dispatcher.dispatch_ready_events(channel="station") + finally: + dispatcher.close() + + assert result.processed == 1 + assert result.sent == 0 + assert result.failed == 1 + + notification_service = DeliveryNotificationService.for_db(settings.orders_db_path) + try: + failed_event = notification_service.list_events(order.id)[0] + finally: + notification_service.close() + assert failed_event.status == "failed" + assert failed_event.attempt_count == 2 + assert failed_event.failure_reason == "delivery artifact missing" + + +def test_delivery_dispatch_script_prints_summary(settings, tmp_path): + order = _seed_order(settings.orders_db_path, order_id="GKO-20260614-DISPATCH-CLI") + _mark_paid(settings, order) + IntakeStore.for_db(settings.orders_db_path).save( + order_id=order.id, payload={"candidate_score": 578}, submit=True + ) + _attach_ready_delivery(settings, tmp_path, order.id) + + env = { + **__import__("os").environ, + "GAOKAO_ORDERS_DB_PATH": settings.orders_db_path, + } + proc = subprocess.run( + ["python3", "scripts/gaokao-delivery-dispatch.py", "--channel", "station"], + cwd=PROJECT_ROOT, + text=True, + capture_output=True, + env=env, + check=False, + ) + + assert proc.returncode == 0, proc.stderr + assert '"processed": 1' in proc.stdout + assert '"sent": 1' in proc.stdout