feat: harden delivery and deletion workflows
This commit is contained in:
13
README.md
13
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` 并替换密钥/密码。
|
||||
|
||||
@@ -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__
|
||||
|
||||
145
admin/tests/test_order_deletion.py
Normal file
145
admin/tests/test_order_deletion.py
Normal file
@@ -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("<h1>report</h1>", 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()
|
||||
69
data/notifications/dispatcher.py
Normal file
69
data/notifications/dispatcher.py
Normal file
@@ -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
|
||||
@@ -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]
|
||||
|
||||
148
data/orders/deletion_service.py
Normal file
148
data/orders/deletion_service.py
Normal file
@@ -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
|
||||
@@ -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 本地验证环境仍未完全固化为单命令体验
|
||||
|
||||
|
||||
@@ -80,4 +80,6 @@
|
||||
当前状态:
|
||||
|
||||
- **文档基线已补齐**
|
||||
- **代码主链分级展示仍待继续实现**
|
||||
- **`risk_report` 已输出 `quality_level / quality_label`**
|
||||
- **`gaokao-data-trace --human` 已输出质量等级**
|
||||
- **低置信省份已通过 confidence 分级与 warning 机制降级**
|
||||
|
||||
@@ -40,13 +40,15 @@
|
||||
|
||||
- `OrdersDAO.delete()` 删除订单能力
|
||||
- ON DELETE CASCADE 清理部分关联表
|
||||
- 后台 `DELETE /api/orders/{id}?mode=delete|anonymize&reason=...` 最小执行入口
|
||||
- 删除时自动清理 HTML/PDF 报告文件
|
||||
- `order_deletion_audits` 最小审计表
|
||||
|
||||
尚缺:
|
||||
|
||||
- 前台/后台删除工单流程
|
||||
- 报告文件自动清理脚本
|
||||
- 删除操作审计记录
|
||||
- 前台/客服删除工单流程
|
||||
- 数据保留期自动清理任务
|
||||
- 更细粒度的匿名化策略(如案例长期保留脱敏版)
|
||||
|
||||
## 5. MVP 最低执行要求
|
||||
|
||||
|
||||
@@ -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. 再决定是否扩到多通道
|
||||
|
||||
37
scripts/gaokao-delivery-dispatch.py
Normal file
37
scripts/gaokao-delivery-dispatch.py
Normal file
@@ -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())
|
||||
152
tests/test_delivery_dispatcher.py
Normal file
152
tests/test_delivery_dispatcher.py
Normal file
@@ -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("<h1>ready</h1>", 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
|
||||
Reference in New Issue
Block a user