Compare commits

..

10 Commits

Author SHA1 Message Date
exinglang 55cdbd07f1 fix(push): 补齐小米通道与审核推送日志 2026-07-30 12:17:12 +08:00
exinglang 6ca8ca4ed2 Merge remote-tracking branch 'origin/main' into fix-pushlogxiaomi
# Conflicts:
#	app/integrations/vendor_push.py
2026-07-30 12:14:45 +08:00
guke b6ddb275f4 fix(comparison): 比价 best 无 is_best 时兜底回落「含源最低价」 (#201)
背景
_derive_from_platforms 是比价记录的单一真相源派生函数(pricebot done 帧带 platforms 时走它),正常路径直接取 platforms 里 is_best=true 的那行当最低价。

线上出现一类帧(如记录 id 3304):全平台 has_dish_diff(菜品「相似替换 / 价格仅供参考」),pricebot 认为无法认定权威最低价,于是一个 is_best 都不标。此时旧逻辑 best=None,导致 best_platform_id / best_price_cents / saved_amount_cents / is_source_best 整条落 NULL,连锁反应:

首页比价价显示 0.00
记录页无最低价红框
省额丢失、「累计发现可省」漏计
方案
无 is_best 时兜底:在有价行里取最低价当参考 best。

关键设计点 —— 候选池含源(而非仅目标):源平台常年全菜、价可信。若源本身最便宜(其余都是更贵的相似替换),则 best 回落到源、saved=0、is_source_best=True。这与仓库老派生函数 _derive 的既有语义(「全目标缺菜 → 回落源、不虚报省」)完全一致,三个派生函数行为对齐。

⚠️ 若像分支首个提交那样排除源、强选最低目标,当源最便宜时会选中更贵目标 → saved 变负,倒扣 get_stats 的「累计发现可省」(该聚合按 status='success' 求和、不带 >0 过滤)。第二个提交据此修正为含源。

影响面 / 兼容性
只影响「带 platforms 且无任何 is_best」这一条兜底分支;正常有 is_best 的路径不变。
老客户端不带 platforms → 走 _derive,不受影响。
纯派生逻辑,无 schema / 无迁移,回滚成本低。

---------

Co-authored-by: guke <guke@autohome.com.cn>
Reviewed-on: #201
2026-07-30 11:00:35 +08:00
zuochenyong ebacf01742 feat(新手引导): 支持十圈进度与逐圈发奖 (#200)
Co-authored-by: exinglang <exinglang@qq.com>
Reviewed-on: #200
Co-authored-by: zuochenyong <zuochenyong@wonderable.ai>
Co-committed-by: zuochenyong <zuochenyong@wonderable.ai>
2026-07-30 10:49:24 +08:00
zuochenyong 89f266419b fix(push): 完善厂商推送参数与消息通知 (#199)
Co-authored-by: exinglang <exinglang@qq.com>
Reviewed-on: #199
Co-authored-by: zuochenyong <zuochenyong@wonderable.ai>
Co-committed-by: zuochenyong <zuochenyong@wonderable.ai>
2026-07-30 10:49:08 +08:00
zuochenyong 675c7ecf81 fix(push): 完善厂商推送排障日志 (#198)
Co-authored-by: exinglang <exinglang@qq.com>
Reviewed-on: #198
Co-authored-by: zuochenyong <zuochenyong@wonderable.ai>
Co-committed-by: zuochenyong <zuochenyong@wonderable.ai>
2026-07-30 10:41:43 +08:00
exinglang 7474c2bebc fix(push): 补齐厂商请求日志参数 2026-07-30 10:26:26 +08:00
exinglang 7bfe703760 chore(tests): 取消推送测试文件改动 2026-07-30 10:25:05 +08:00
exinglang 69eeb43fe2 fix(push): 完善厂商推送参数与消息通知 2026-07-30 10:25:05 +08:00
exinglang a0196b8d64 fix(push): 完善厂商推送排障日志 2026-07-29 19:52:06 +08:00
8 changed files with 831 additions and 63 deletions
+6
View File
@@ -175,3 +175,9 @@ PANGLE_REPORT_SITE_ID_TEST=5832303
# APPLOG_MAX_BATCH=500 # 单批最大条数(超 → 422;导入期常量,改需重启) # APPLOG_MAX_BATCH=500 # 单批最大条数(超 → 422;导入期常量,改需重启)
# APPLOG_MAX_BODY_BYTES=2097152 # 请求体上限 2MB(超 → 413;运行期可调) # APPLOG_MAX_BODY_BYTES=2097152 # 请求体上限 2MB(超 → 413;运行期可调)
# APPLOG_MAX_MSG_BYTES=8192 # 单条 msg 超此字节数截断 # APPLOG_MAX_MSG_BYTES=8192 # 单条 msg 超此字节数截断
# ===== 华为 Push =====
HUAWEI_PUSH_APP_ID=
HUAWEI_PUSH_APP_SECRET=
# 0=正式消息(默认);1=测试消息(仅开发联调,生产必须保持 0)
HUAWEI_PUSH_TARGET_USER_TYPE=0
+109 -6
View File
@@ -1,6 +1,7 @@
"""admin 反馈工单:列表筛选 + 审核采纳/拒绝(带金币发放与审计)。""" """admin 反馈工单:列表筛选 + 审核采纳/拒绝(带金币发放与审计)。"""
from __future__ import annotations from __future__ import annotations
import logging
from datetime import datetime from datetime import datetime
from typing import Annotated from typing import Annotated
@@ -25,6 +26,8 @@ from app.models.feedback import Feedback
from app.repositories import wallet as wallet_repo from app.repositories import wallet as wallet_repo
from app.services import notification_events from app.services import notification_events
logger = logging.getLogger("shagua.admin.feedback")
router = APIRouter( router = APIRouter(
prefix="/admin/api/feedbacks", prefix="/admin/api/feedbacks",
tags=["admin-feedback"], tags=["admin-feedback"],
@@ -46,10 +49,24 @@ def _approve_feedback(
*, *,
bulk: bool = False, bulk: bool = False,
) -> FeedbackOut: ) -> FeedbackOut:
logger.info(
"feedback approve started feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
fb = db.get(Feedback, feedback_id, with_for_update=True) fb = db.get(Feedback, feedback_id, with_for_update=True)
if fb is None: if fb is None:
logger.warning(
"feedback approve rejected not found feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
raise HTTPException(status_code=404, detail="反馈不存在") raise HTTPException(status_code=404, detail="反馈不存在")
_ensure_pending(fb) if fb.status not in {"pending", "new"}:
logger.warning(
"feedback approve rejected invalid status feedback_id=%s user_id=%s "
"admin_id=%s status=%s bulk=%s",
feedback_id, fb.user_id, admin.id, fb.status, bulk,
)
_ensure_pending(fb)
before = fb.status before = fb.status
mutations.review_feedback( mutations.review_feedback(
@@ -91,8 +108,18 @@ def _approve_feedback(
) )
db.commit() db.commit()
db.refresh(fb) db.refresh(fb)
logger.info(
"feedback approve committed feedback_id=%s user_id=%s admin_id=%s "
"before=%s after=%s reward_coins=%s bulk=%s",
feedback_id, fb.user_id, admin.id, before, fb.status, payload.reward_coins, bulk,
)
out = FeedbackOut.model_validate(fb) out = FeedbackOut.model_validate(fb)
notification_events.notify_feedback_reward(db, fb) notification_events.notify_feedback_reward(db, fb)
logger.info(
"feedback approve notification dispatch returned feedback_id=%s user_id=%s "
"admin_id=%s notification_type=feedback_reward bulk=%s",
feedback_id, fb.user_id, admin.id, bulk,
)
return out return out
@@ -105,10 +132,24 @@ def _reject_feedback(
*, *,
bulk: bool = False, bulk: bool = False,
) -> FeedbackOut: ) -> FeedbackOut:
logger.info(
"feedback reject started feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
fb = db.get(Feedback, feedback_id, with_for_update=True) fb = db.get(Feedback, feedback_id, with_for_update=True)
if fb is None: if fb is None:
logger.warning(
"feedback reject rejected not found feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
raise HTTPException(status_code=404, detail="反馈不存在") raise HTTPException(status_code=404, detail="反馈不存在")
_ensure_pending(fb) if fb.status not in {"pending", "new"}:
logger.warning(
"feedback reject rejected invalid status feedback_id=%s user_id=%s "
"admin_id=%s status=%s bulk=%s",
feedback_id, fb.user_id, admin.id, fb.status, bulk,
)
_ensure_pending(fb)
before = fb.status before = fb.status
mutations.review_feedback( mutations.review_feedback(
@@ -142,8 +183,18 @@ def _reject_feedback(
) )
db.commit() db.commit()
db.refresh(fb) db.refresh(fb)
logger.info(
"feedback reject committed feedback_id=%s user_id=%s admin_id=%s "
"before=%s after=%s bulk=%s",
feedback_id, fb.user_id, admin.id, before, fb.status, bulk,
)
out = FeedbackOut.model_validate(fb) out = FeedbackOut.model_validate(fb)
notification_events.notify_feedback_reply(db, fb) notification_events.notify_feedback_reply(db, fb)
logger.info(
"feedback reject notification dispatch returned feedback_id=%s user_id=%s "
"admin_id=%s notification_type=feedback_reply bulk=%s",
feedback_id, fb.user_id, admin.id, bulk,
)
return out return out
@@ -201,6 +252,10 @@ def bulk_approve_feedbacks(
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
) -> FeedbackBulkResult: ) -> FeedbackBulkResult:
logger.info(
"feedback bulk approve started admin_id=%s item_count=%s",
admin.id, len(body.ids),
)
results: list[FeedbackBulkItemResult] = [] results: list[FeedbackBulkItemResult] = []
ip = get_client_ip(request) ip = get_client_ip(request)
for feedback_id in body.ids: for feedback_id in body.ids:
@@ -209,11 +264,24 @@ def bulk_approve_feedbacks(
results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status)) results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status))
except HTTPException as exc: except HTTPException as exc:
db.rollback() db.rollback()
logger.warning(
"feedback bulk approve item failed feedback_id=%s admin_id=%s error=%s",
feedback_id, admin.id, exc.detail,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail))) results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail)))
except Exception: # noqa: BLE001 - 单笔失败不打断整批 except Exception: # noqa: BLE001 - 单笔失败不打断整批
db.rollback() db.rollback()
logger.exception(
"feedback bulk approve item failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常")) results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常"))
return _bulk_result(results) result = _bulk_result(results)
logger.info(
"feedback bulk approve completed admin_id=%s total=%s success=%s failed=%s",
admin.id, result.total, result.success, result.failed,
)
return result
@router.post("/bulk/reject", response_model=FeedbackBulkResult, summary="批量拒绝采纳反馈") @router.post("/bulk/reject", response_model=FeedbackBulkResult, summary="批量拒绝采纳反馈")
@@ -223,6 +291,10 @@ def bulk_reject_feedbacks(
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
) -> FeedbackBulkResult: ) -> FeedbackBulkResult:
logger.info(
"feedback bulk reject started admin_id=%s item_count=%s",
admin.id, len(body.ids),
)
results: list[FeedbackBulkItemResult] = [] results: list[FeedbackBulkItemResult] = []
ip = get_client_ip(request) ip = get_client_ip(request)
for feedback_id in body.ids: for feedback_id in body.ids:
@@ -231,11 +303,24 @@ def bulk_reject_feedbacks(
results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status)) results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status))
except HTTPException as exc: except HTTPException as exc:
db.rollback() db.rollback()
logger.warning(
"feedback bulk reject item failed feedback_id=%s admin_id=%s error=%s",
feedback_id, admin.id, exc.detail,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail))) results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail)))
except Exception: # noqa: BLE001 - 单笔失败不打断整批 except Exception: # noqa: BLE001 - 单笔失败不打断整批
db.rollback() db.rollback()
logger.exception(
"feedback bulk reject item failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常")) results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常"))
return _bulk_result(results) result = _bulk_result(results)
logger.info(
"feedback bulk reject completed admin_id=%s total=%s success=%s failed=%s",
admin.id, result.total, result.success, result.failed,
)
return result
@router.post("/{feedback_id}/handle", response_model=OkResponse, summary="标记反馈已处理") @router.post("/{feedback_id}/handle", response_model=OkResponse, summary="标记反馈已处理")
@@ -258,7 +343,16 @@ def approve_feedback(
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
) -> FeedbackOut: ) -> FeedbackOut:
return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request)) try:
return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request))
except HTTPException:
raise
except Exception:
logger.exception(
"feedback approve failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
raise
@router.post("/{feedback_id}/reject", response_model=FeedbackOut, summary="拒绝采纳反馈") @router.post("/{feedback_id}/reject", response_model=FeedbackOut, summary="拒绝采纳反馈")
@@ -269,4 +363,13 @@ def reject_feedback(
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
) -> FeedbackOut: ) -> FeedbackOut:
return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request)) try:
return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request))
except HTTPException:
raise
except Exception:
logger.exception(
"feedback reject failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
raise
+2
View File
@@ -88,6 +88,8 @@ class Settings(BaseSettings):
# (OAuth 换 token 时 client_id 即 AppId)。发送走 v1 messages:send,成功码 80000000。 # (OAuth 换 token 时 client_id 即 AppId)。发送走 v1 messages:send,成功码 80000000。
HUAWEI_PUSH_APP_ID: str = "" HUAWEI_PUSH_APP_ID: str = ""
HUAWEI_PUSH_APP_SECRET: str = "" HUAWEI_PUSH_APP_SECRET: str = ""
# 0=正式消息(默认,受正式消息频控);1=测试消息(仅开发联调,勿用于生产)。
HUAWEI_PUSH_TARGET_USER_TYPE: int = Field(default=0, ge=0, le=1)
HUAWEI_PUSH_TOKEN_ENDPOINT: str = "https://oauth-login.cloud.huawei.com/oauth2/v3/token" HUAWEI_PUSH_TOKEN_ENDPOINT: str = "https://oauth-login.cloud.huawei.com/oauth2/v3/token"
HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE: str = ( HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE: str = (
"https://push-api.cloud.huawei.com/v1/{app_id}/messages:send" "https://push-api.cloud.huawei.com/v1/{app_id}/messages:send"
+455 -46
View File
@@ -22,7 +22,7 @@ import uuid
from collections.abc import Callable from collections.abc import Callable
from dataclasses import dataclass from dataclasses import dataclass
from typing import Any from typing import Any
from urllib.parse import quote from urllib.parse import quote, urlsplit
import httpx import httpx
@@ -31,6 +31,7 @@ from app.core.config import settings
logger = logging.getLogger("shagua.vendor_push") logger = logging.getLogger("shagua.vendor_push")
TYPE_ACCESSIBILITY_DISABLED = "accessibility_disabled" TYPE_ACCESSIBILITY_DISABLED = "accessibility_disabled"
DATA_EVENT_NOTIFICATION_CREATED = "notification_created"
SUPPORTED_VENDORS = frozenset({"honor", "huawei", "vivo", "xiaomi", "oppo"}) SUPPORTED_VENDORS = frozenset({"honor", "huawei", "vivo", "xiaomi", "oppo"})
# vendor key → 中文名(测试/配置状态接口展示用) # vendor key → 中文名(测试/配置状态接口展示用)
@@ -69,6 +70,18 @@ class _CachedToken:
_token_cache: dict[str, _CachedToken] = {} _token_cache: dict[str, _CachedToken] = {}
_LOG_SUMMARY_MAX_CHARS = 1500
_LOG_STRING_MAX_CHARS = 200
_SENSITIVE_LOG_KEYS = {
"accesstoken",
"appkey",
"authtoken",
"authorization",
"clientsecret",
"mastersecret",
"sign",
}
def normalize_vendor(push_vendor: str | None) -> str | None: def normalize_vendor(push_vendor: str | None) -> str | None:
if not push_vendor: if not push_vendor:
@@ -115,8 +128,16 @@ def send_notification(
if mock: if mock:
logger.info( logger.info(
"[mock push] vendor=%s token=%s... title=%s body=%s extras=%s", "vendor push mock vendor=%s request=%s",
vendor, token[:12], title, body, extras, vendor,
_log_summary(
{
"push_token": token,
"title": title,
"body": body,
"extras": extras,
}
),
) )
return { return {
"mock": True, "mock": True,
@@ -133,7 +154,28 @@ def send_notification(
"xiaomi": _send_xiaomi, "xiaomi": _send_xiaomi,
"oppo": _send_oppo, "oppo": _send_oppo,
} }
return dispatch[vendor](token, title, body, extras) started = time.perf_counter()
try:
result = dispatch[vendor](token, title, body, extras)
except Exception:
logger.exception(
"vendor push dispatch failed vendor=%s token=%s notification_type=%s "
"elapsed_ms=%.1f",
vendor,
token,
extras.get("type", ""),
_elapsed_ms(started),
)
raise
logger.info(
"vendor push dispatch succeeded vendor=%s token=%s notification_type=%s "
"elapsed_ms=%.1f",
vendor,
token,
extras.get("type", ""),
_elapsed_ms(started),
)
return result
def send_accessibility_disabled( def send_accessibility_disabled(
@@ -153,19 +195,188 @@ def send_accessibility_disabled(
) )
def send_data_event(
push_vendor: str,
push_token: str,
*,
event: str,
notification_id: str,
mock: bool = False,
) -> dict[str, Any]:
"""发送无界面的轻量事件;当前只允许站内消息创建事件。"""
vendor = normalize_vendor(push_vendor)
token = push_token.strip() if push_token else ""
if not vendor or vendor not in SUPPORTED_VENDORS:
raise VendorPushError(f"unsupported push vendor: {push_vendor}")
if not token:
raise VendorPushError("push token is empty")
if event != DATA_EVENT_NOTIFICATION_CREATED:
raise VendorPushError(f"unsupported data event: {event}")
if not notification_id:
raise VendorPushError("notification_id is empty")
payload = {"event": event, "notificationId": str(notification_id)}
if mock:
return {"mock": True, "vendor": vendor, "payload": payload}
# vivo 的通知已设置 foregroundShow=falseApp 在前台时必走
# onForegroundMessageArrived,正好就是本事件需要的刷新信号;不重复占用一次推送配额。
if vendor == "vivo":
return {"skipped": True, "reason": "foreground notification callback"}
# OPush 当前只支持通知栏消息,没有服务端透传单推接口。
# 客户端在 OPPO 首页可见时轮询未读数;这里必须安全跳过,不能请求不存在的
# /message/transparent/unicast(该地址会稳定返回 HTTP 404)。
if vendor == "oppo":
return {"skipped": True, "reason": "oppo does not support data messages"}
dispatch: dict[str, Callable[[str, dict[str, str]], dict[str, Any]]] = {
"honor": _send_honor_data,
"huawei": _send_huawei_data,
"xiaomi": _send_xiaomi_data,
}
return dispatch[vendor](token, payload)
def _require(value: str, name: str) -> str: def _require(value: str, name: str) -> str:
if not value: if not value:
raise VendorPushError(f"{name} not configured") raise VendorPushError(f"{name} not configured")
return value return value
def _elapsed_ms(started: float) -> float:
return (time.perf_counter() - started) * 1000
def _redacted_value(value: Any) -> str:
raw = json.dumps(value, ensure_ascii=False, sort_keys=True, default=str)
digest = hashlib.sha256(raw.encode()).hexdigest()[:10]
return f"<redacted len={len(raw)} sha256={digest}>"
def _sanitize_for_log(value: Any, *, key: str = "", depth: int = 0) -> Any:
"""生成有排障价值、但不暴露服务端鉴权凭据的紧凑日志摘要。"""
normalized_key = "".join(char for char in key.lower() if char.isalnum())
if normalized_key in _SENSITIVE_LOG_KEYS or normalized_key.endswith("secret"):
return _redacted_value(value)
if depth >= 10:
return f"<{type(value).__name__}>"
if isinstance(value, dict):
items = list(value.items())
sanitized = {
str(item_key): _sanitize_for_log(item_value, key=str(item_key), depth=depth + 1)
for item_key, item_value in items[:30]
}
if len(items) > 30:
sanitized["<omitted_keys>"] = len(items) - 30
return sanitized
if isinstance(value, (list, tuple)):
items = list(value)
sanitized = [_sanitize_for_log(item, key=key, depth=depth + 1) for item in items[:20]]
if len(items) > 20:
sanitized.append(f"<omitted_items={len(items) - 20}>")
return sanitized
if isinstance(value, str):
if value.lstrip().startswith(("{", "[")):
try:
decoded = json.loads(value)
except ValueError:
pass
else:
return _sanitize_for_log(decoded, key=key, depth=depth + 1)
if len(value) > _LOG_STRING_MAX_CHARS:
return f"{value[:_LOG_STRING_MAX_CHARS]}…<len={len(value)}>"
return value
def _log_summary(value: Any) -> str:
rendered = json.dumps(
_sanitize_for_log(value),
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
default=str,
)
if len(rendered) > _LOG_SUMMARY_MAX_CHARS:
return f"{rendered[:_LOG_SUMMARY_MAX_CHARS]}…<len={len(rendered)}>"
return rendered
def _raw_log_summary(value: Any) -> str:
"""厂商响应摘要:保留原始字段和值,仅限制单条日志长度。"""
rendered = json.dumps(
value,
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
default=str,
)
if len(rendered) > _LOG_SUMMARY_MAX_CHARS:
return f"{rendered[:_LOG_SUMMARY_MAX_CHARS]}…<len={len(rendered)}>"
return rendered
def _endpoint_for_log(url: str) -> str:
parsed = urlsplit(url)
return f"{parsed.scheme}://{parsed.netloc}{parsed.path}"
def _request_summary(kwargs: dict[str, Any]) -> str:
summary: dict[str, Any] = {}
if "json" in kwargs:
summary["body_type"] = "json"
summary["request_payload"] = kwargs["json"]
elif "data" in kwargs:
summary["body_type"] = "form"
summary["request_payload"] = kwargs["data"]
if kwargs.get("params"):
summary["params"] = kwargs["params"]
if kwargs.get("headers"):
summary["headers"] = kwargs["headers"]
return _log_summary(summary)
def _response_summary(resp: Any, parsed: Any | None = None) -> str:
if parsed is None:
try:
parsed = resp.json()
except ValueError:
pass
if parsed is not None:
return _raw_log_summary(parsed)
return _raw_log_summary(getattr(resp, "text", ""))
def _vendor_response_failed(vendor: str, data: dict[str, Any]) -> bool:
"""识别 HTTP 200 中明确的厂商业务失败,避免将失败请求记录成 success。"""
if data.get("error") or data.get("success") is False:
return True
if vendor == "xiaomi":
code = data.get("code")
result = str(data.get("result", "ok")).lower()
return code not in (0, "0", None) or result not in ("ok", "success")
if vendor == "oppo" and data.get("code") is not None:
return int(data["code"]) != 0
if vendor == "vivo" and data.get("result") is not None:
return int(data["result"]) != 0
if vendor == "honor" and data.get("code") is not None:
return int(data["code"]) != 200
if vendor == "huawei" and data.get("code") is not None:
return str(data["code"]) != "80000000"
return False
def _request_json( def _request_json(
method: str, method: str,
url: str, url: str,
*, *,
vendor: str,
operation: str,
expected_status: tuple[int, ...] = (200,), expected_status: tuple[int, ...] = (200,),
**kwargs: Any, **kwargs: Any,
) -> dict[str, Any]: ) -> dict[str, Any]:
started = time.perf_counter()
request_summary = _request_summary(kwargs)
try: try:
resp = httpx.request( resp = httpx.request(
method, method,
@@ -174,41 +385,85 @@ def _request_json(
**kwargs, **kwargs,
) )
except httpx.HTTPError as e: except httpx.HTTPError as e:
logger.error(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=network_error request=%s response=%s elapsed_ms=%.1f error=%s",
vendor,
operation,
method,
_endpoint_for_log(url),
request_summary,
"<no_response>",
_elapsed_ms(started),
type(e).__name__,
)
raise VendorPushError(f"push http error: {e}") from e raise VendorPushError(f"push http error: {e}") from e
if resp.status_code not in expected_status: if resp.status_code not in expected_status:
logger.error("vendor push http failed url=%s http=%s body=%s", url, resp.status_code, resp.text[:500]) logger.error(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=http_error http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
resp.status_code,
request_summary,
_response_summary(resp),
_elapsed_ms(started),
)
raise VendorPushError(f"push http {resp.status_code}") raise VendorPushError(f"push http {resp.status_code}")
try: try:
return resp.json() data = resp.json()
except ValueError as e: except ValueError as e:
logger.error(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=invalid_json http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
resp.status_code,
request_summary,
_response_summary(resp),
_elapsed_ms(started),
)
raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e
vendor_failed = _vendor_response_failed(vendor, data)
log = logger.warning if vendor_failed else logger.info
log(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=%s http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
"vendor_error" if vendor_failed else "success",
resp.status_code,
request_summary,
_response_summary(resp, data),
_elapsed_ms(started),
)
return data
def _request_form( def _request_form(
method: str, method: str,
url: str, url: str,
*, *,
vendor: str,
operation: str,
expected_status: tuple[int, ...] = (200,), expected_status: tuple[int, ...] = (200,),
**kwargs: Any, **kwargs: Any,
) -> dict[str, Any]: ) -> dict[str, Any]:
try: return _request_json(
resp = httpx.request( method,
method, url,
url, vendor=vendor,
timeout=settings.PUSH_REQUEST_TIMEOUT_SEC, operation=operation,
**kwargs, expected_status=expected_status,
) **kwargs,
except httpx.HTTPError as e: )
raise VendorPushError(f"push http error: {e}") from e
if resp.status_code not in expected_status:
logger.error("vendor push http failed url=%s http=%s body=%s", url, resp.status_code, resp.text[:500])
raise VendorPushError(f"push http {resp.status_code}")
try:
return resp.json()
except ValueError as e:
raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e
def _cache_get(key: str) -> str | None: def _cache_get(key: str) -> str | None:
@@ -234,6 +489,8 @@ def _honor_access_token() -> str:
data = _request_form( data = _request_form(
"POST", "POST",
settings.HONOR_PUSH_TOKEN_ENDPOINT, settings.HONOR_PUSH_TOKEN_ENDPOINT,
vendor="honor",
operation="authenticate",
data={ data={
"grant_type": "client_credentials", "grant_type": "client_credentials",
"client_id": client_id, "client_id": client_id,
@@ -243,16 +500,22 @@ def _honor_access_token() -> str:
) )
token = data.get("access_token") token = data.get("access_token")
if not token: if not token:
raise VendorPushError(f"honor auth failed: {data}") raise VendorPushError(f"honor auth failed: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), data.get("expires_in")) return _cache_put(cache_key, str(token), data.get("expires_in"))
def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> dict[str, Any]: def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> dict[str, Any]:
app_id = _require(settings.HONOR_PUSH_APP_ID, "HONOR_PUSH_APP_ID") app_id = _require(settings.HONOR_PUSH_APP_ID, "HONOR_PUSH_APP_ID")
access_token = _honor_access_token() access_token = _honor_access_token()
click_action: dict[str, Any]
if extras.get("notificationId"):
# type=3 只负责打开首页,不保证 data 会变成目标 Activity 的 extras。消息中心通知必须
# 用自定义页面(type=1)+ intent URI,把 notificationId/type/feedbackId 等直接带给
# MainActivity;否则华为/荣耀系统能展示通知,但用户点击后客户端收不到任何导航参数。
click_action = {"type": 1, "intent": _click_intent_uri(extras)}
else:
click_action = {"type": 3}
payload = { payload = {
# clickAction type=3(打开应用首页)时,荣耀点击会把 data JSON 的键值对注入启动 intent 的
# extras(与 HMS 同机制)→ MainActivity.consumeNavTarget 读 notif_id/notif_type 直达落地。
"data": json.dumps(_click_extras(extras), ensure_ascii=False), "data": json.dumps(_click_extras(extras), ensure_ascii=False),
"notification": {"title": title, "body": body}, "notification": {"title": title, "body": body},
"android": { "android": {
@@ -261,7 +524,7 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di
"notification": { "notification": {
"title": title, "title": title,
"body": body, "body": body,
"clickAction": {"type": 3}, "clickAction": click_action,
"importance": "NORMAL", "importance": "NORMAL",
}, },
}, },
@@ -270,6 +533,8 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di
data = _request_json( data = _request_json(
"POST", "POST",
settings.HONOR_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id), settings.HONOR_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id),
vendor="honor",
operation="send_notification",
json=payload, json=payload,
headers={ headers={
"Content-Type": "application/json; charset=UTF-8", "Content-Type": "application/json; charset=UTF-8",
@@ -279,7 +544,28 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di
) )
code = data.get("code") code = data.get("code")
if code is not None and int(code) != 200: if code is not None and int(code) != 200:
raise VendorPushError(f"honor push failed: {data}") raise VendorPushError(f"honor push failed: {_raw_log_summary(data)}")
return data
def _send_honor_data(token: str, payload: dict[str, str]) -> dict[str, Any]:
app_id = _require(settings.HONOR_PUSH_APP_ID, "HONOR_PUSH_APP_ID")
access_token = _honor_access_token()
data = _request_json(
"POST",
settings.HONOR_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id),
vendor="honor",
operation="send_data_event",
json={"data": json.dumps(payload, ensure_ascii=False), "token": [token]},
headers={
"Content-Type": "application/json; charset=UTF-8",
"Authorization": f"Bearer {access_token}",
"timestamp": str(int(time.time() * 1000)),
},
)
code = data.get("code")
if code is not None and int(code) != 200:
raise VendorPushError(f"honor data push failed: {_raw_log_summary(data)}")
return data return data
@@ -294,6 +580,8 @@ def _huawei_access_token() -> str:
data = _request_form( data = _request_form(
"POST", "POST",
settings.HUAWEI_PUSH_TOKEN_ENDPOINT, settings.HUAWEI_PUSH_TOKEN_ENDPOINT,
vendor="huawei",
operation="authenticate",
data={ data={
"grant_type": "client_credentials", "grant_type": "client_credentials",
"client_id": app_id, "client_id": app_id,
@@ -303,7 +591,7 @@ def _huawei_access_token() -> str:
) )
token = data.get("access_token") token = data.get("access_token")
if not token: if not token:
raise VendorPushError(f"huawei auth failed: {data}") raise VendorPushError(f"huawei auth failed: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), data.get("expires_in")) return _cache_put(cache_key, str(token), data.get("expires_in"))
@@ -312,18 +600,24 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d
'80100000' 为部分成功(单 token 场景仍视为失败,错误里带原始响应便于排障)。""" '80100000' 为部分成功(单 token 场景仍视为失败,错误里带原始响应便于排障)。"""
app_id = _require(settings.HUAWEI_PUSH_APP_ID, "HUAWEI_PUSH_APP_ID") app_id = _require(settings.HUAWEI_PUSH_APP_ID, "HUAWEI_PUSH_APP_ID")
access_token = _huawei_access_token() access_token = _huawei_access_token()
click_action: dict[str, Any]
if extras.get("notificationId"):
# type=3 只打开首页,Mate 20/EMUI 不会把 message.data 自动拆成启动 Intent extras。
# type=1 的自定义 intent 才能稳定携带反馈记录 id,且冷启动/onNewIntent 都走同一路由。
click_action = {"type": 1, "intent": _click_intent_uri(extras)}
else:
click_action = {"type": 3}
payload = { payload = {
"validate_only": False, "validate_only": False,
"message": { "message": {
# click_action type=3(打开应用首页)时,HMS 点击会把 data JSON 的键值对注入启动 intent
# 的 extras → MainActivity.consumeNavTarget 读 notif_id/notif_type 直达落地。
"data": json.dumps(_click_extras(extras), ensure_ascii=False), "data": json.dumps(_click_extras(extras), ensure_ascii=False),
"android": { "android": {
"target_user_type": settings.HUAWEI_PUSH_TARGET_USER_TYPE,
"ttl": f"{settings.PUSH_TIME_TO_LIVE_SEC}s", "ttl": f"{settings.PUSH_TIME_TO_LIVE_SEC}s",
"notification": { "notification": {
"title": title, "title": title,
"body": body, "body": body,
"click_action": {"type": 3}, "click_action": click_action,
"importance": "NORMAL", "importance": "NORMAL",
}, },
}, },
@@ -333,6 +627,8 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d
data = _request_json( data = _request_json(
"POST", "POST",
settings.HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id), settings.HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id),
vendor="huawei",
operation="send_notification",
json=payload, json=payload,
headers={ headers={
"Content-Type": "application/json; charset=UTF-8", "Content-Type": "application/json; charset=UTF-8",
@@ -340,7 +636,32 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d
}, },
) )
if str(data.get("code", "")) != "80000000": if str(data.get("code", "")) != "80000000":
raise VendorPushError(f"huawei push failed: {data}") raise VendorPushError(f"huawei push failed: {_raw_log_summary(data)}")
return data
def _send_huawei_data(token: str, payload: dict[str, str]) -> dict[str, Any]:
app_id = _require(settings.HUAWEI_PUSH_APP_ID, "HUAWEI_PUSH_APP_ID")
access_token = _huawei_access_token()
data = _request_json(
"POST",
settings.HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id),
vendor="huawei",
operation="send_data_event",
json={
"validate_only": False,
"message": {
"data": json.dumps(payload, ensure_ascii=False),
"token": [token],
},
},
headers={
"Content-Type": "application/json; charset=UTF-8",
"Authorization": f"Bearer {access_token}",
},
)
if str(data.get("code", "")) != "80000000":
raise VendorPushError(f"huawei data push failed: {_raw_log_summary(data)}")
return data return data
@@ -357,6 +678,8 @@ def _vivo_auth_token() -> str:
data = _request_json( data = _request_json(
"POST", "POST",
settings.VIVO_PUSH_AUTH_ENDPOINT, settings.VIVO_PUSH_AUTH_ENDPOINT,
vendor="vivo",
operation="authenticate",
json={ json={
"appId": app_id, "appId": app_id,
"appKey": app_key, "appKey": app_key,
@@ -366,10 +689,10 @@ def _vivo_auth_token() -> str:
headers={"Content-Type": "application/json"}, headers={"Content-Type": "application/json"},
) )
if int(data.get("result", -1)) != 0: if int(data.get("result", -1)) != 0:
raise VendorPushError(f"vivo auth failed: {data}") raise VendorPushError(f"vivo auth failed: {_raw_log_summary(data)}")
token = data.get("authToken") token = data.get("authToken")
if not token: if not token:
raise VendorPushError(f"vivo auth missing authToken: {data}") raise VendorPushError(f"vivo auth missing authToken: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), 24 * 3600) return _cache_put(cache_key, str(token), 24 * 3600)
@@ -385,15 +708,22 @@ def _send_vivo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
"timeToLive": settings.PUSH_TIME_TO_LIVE_SEC, "timeToLive": settings.PUSH_TIME_TO_LIVE_SEC,
"requestId": uuid.uuid4().hex, "requestId": uuid.uuid4().hex,
"pushMode": settings.VIVO_PUSH_MODE, "pushMode": settings.VIVO_PUSH_MODE,
# vivo SDK 4.1.5 只有在前台展示关闭时才调用
# OpenClientPushMessageReceiver.onForegroundMessageArrived。
# App 在该回调中刷新服务端未读数并补发一条本地通知;后台/锁屏时仍由
# vivo 系统展示,因而任何场景都只会出现一条通知。
"foregroundShow": False,
# vivo 官方 VPush 角标字段:离线收到通知时先由桌面自动 +1;
# App 启动/同步后再由统一未读数通过系统 API 精确校准。
"addBadge": True,
"clientCustomMap": extras, "clientCustomMap": extras,
} }
# 点击落地:消息中心推送(带 notificationId)→ skipType=4 + skipContent=intent uri,由 vivo # vivo Push SDK 480+ 已不再回调 skipType=3;必须使用 skipType=4 的完整 Intent URI。
# 系统直启 MainActivity 并携带 S. extras(与小米 notify_effect=2 同机制)。不依赖客户端 # URI 同时包含 data/scheme、显式 component 和 S. extrasOriginOS 才会把业务参数
# VivoPushReceiver.onNotificationMessageClicked 里的后台 startActivity——Android 10+ BAL # 原样交给 MainActivity(仅靠隐式 deeplink 会打开 App,但可能剥掉 extras)。
# 会静默拦掉,receiver 路径仅作兜底。无 notificationId 的召回类保持 skipType=1 仅打开首页。
if extras.get("notificationId"): if extras.get("notificationId"):
payload["skipType"] = 4 payload["skipType"] = 4
payload["skipContent"] = _click_intent_uri(extras) payload["skipContent"] = _vivo_click_intent_uri(extras)
else: else:
payload["skipType"] = 1 payload["skipType"] = 1
if settings.VIVO_PUSH_CATEGORY: if settings.VIVO_PUSH_CATEGORY:
@@ -401,14 +731,37 @@ def _send_vivo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
data = _request_json( data = _request_json(
"POST", "POST",
settings.VIVO_PUSH_SEND_ENDPOINT, settings.VIVO_PUSH_SEND_ENDPOINT,
vendor="vivo",
operation="send_notification",
json=payload, json=payload,
headers={ headers={
"Content-Type": "application/json", "Content-Type": "application/json",
"authToken": auth_token, "authToken": auth_token,
}, },
) )
# 10089 = 当前 vivo 应用尚未开通 VPush「离线自动角标」能力。
# 角标只是通知中心未读数的镜像,不能因此让整条通知发送失败:去掉 addBadge
# 原样重试,通知到达/应用同步后仍由客户端系统 API 按统一未读数精确设置。
if int(data.get("result", -1)) == 10089:
logger.warning("vivo VPush badge not enabled (10089); retrying without addBadge")
fallback_payload = dict(payload)
fallback_payload.pop("addBadge", None)
fallback_payload["requestId"] = uuid.uuid4().hex
data = _request_json(
"POST",
settings.VIVO_PUSH_SEND_ENDPOINT,
vendor="vivo",
operation="send_notification",
json=fallback_payload,
headers={
"Content-Type": "application/json",
"authToken": auth_token,
},
)
if int(data.get("result", -1)) == 0:
data = {**data, "badgeFallback": True}
if int(data.get("result", -1)) != 0: if int(data.get("result", -1)) != 0:
raise VendorPushError(f"vivo push failed: {data}") raise VendorPushError(f"vivo push failed: {_raw_log_summary(data)}")
return data return data
@@ -446,14 +799,43 @@ def _send_xiaomi(token: str, title: str, body: str, extras: dict[str, str]) -> d
data = _request_form( data = _request_form(
"POST", "POST",
settings.XIAOMI_PUSH_SEND_ENDPOINT, settings.XIAOMI_PUSH_SEND_ENDPOINT,
vendor="xiaomi",
operation="send_notification",
data=form, data=form,
headers={"Authorization": f"key={app_secret}"}, headers={"Authorization": f"key={app_secret}"},
) )
code = data.get("code") code = data.get("code")
if code not in (0, "0", None): if code not in (0, "0", None):
raise VendorPushError(f"xiaomi push failed: {data}") raise VendorPushError(f"xiaomi push failed: {_raw_log_summary(data)}")
if str(data.get("result", "ok")).lower() not in ("ok", "success"): if str(data.get("result", "ok")).lower() not in ("ok", "success"):
raise VendorPushError(f"xiaomi push failed: {data}") raise VendorPushError(f"xiaomi push failed: {_raw_log_summary(data)}")
return data
def _send_xiaomi_data(token: str, payload: dict[str, str]) -> dict[str, Any]:
app_secret = _require(settings.XIAOMI_PUSH_APP_SECRET, "XIAOMI_PUSH_APP_SECRET")
form = {
"registration_id": token,
"restricted_package_name": settings.ANDROID_PACKAGE_NAME,
"payload": json.dumps(payload, ensure_ascii=False),
"pass_through": "1",
"time_to_live": str(settings.PUSH_TIME_TO_LIVE_SEC * 1000),
}
if settings.XIAOMI_PUSH_CHANNEL_ID.strip():
form["extra.channel_id"] = settings.XIAOMI_PUSH_CHANNEL_ID.strip()
data = _request_form(
"POST",
settings.XIAOMI_PUSH_SEND_ENDPOINT,
vendor="xiaomi",
operation="send_data_event",
data=form,
headers={"Authorization": f"key={app_secret}"},
)
code = data.get("code")
if code not in (0, "0", None):
raise VendorPushError(f"xiaomi data push failed: {_raw_log_summary(data)}")
if str(data.get("result", "ok")).lower() not in ("ok", "success"):
raise VendorPushError(f"xiaomi data push failed: {_raw_log_summary(data)}")
return data return data
@@ -490,6 +872,24 @@ def _click_intent_uri(extras: dict[str, str]) -> str:
return ";".join(parts) return ";".join(parts)
def _vivo_click_intent_uri(extras: dict[str, str]) -> str:
"""vivo skipType=4 使用官方要求的完整 intent deeplink 形态。
`?#Intent` 中的问号不能省;OriginOS 对缺少它的 URI 会退化为“仅打开应用”,
不把 S. 参数交给目标 Activity。
"""
pkg = settings.ANDROID_PACKAGE_NAME
parts = [
"intent://push/detail?#Intent",
"scheme=shaguabijia",
f"component={pkg}/{pkg}.MainActivity",
"launchFlags=0x14000000",
]
parts += [f"S.{key}={quote(str(value), safe='')}" for key, value in _click_extras(extras).items()]
parts.append("end")
return ";".join(parts)
def _xiaomi_template_param(title: str, alert: str) -> str: def _xiaomi_template_param(title: str, alert: str) -> str:
rendered = ( rendered = (
settings.XIAOMI_PUSH_TEMPLATE_PARAM_JSON settings.XIAOMI_PUSH_TEMPLATE_PARAM_JSON
@@ -522,6 +922,8 @@ def _oppo_auth_token() -> str:
data = _request_form( data = _request_form(
"POST", "POST",
settings.OPPO_PUSH_AUTH_ENDPOINT, settings.OPPO_PUSH_AUTH_ENDPOINT,
vendor="oppo",
operation="authenticate",
data={ data={
"app_key": app_key, "app_key": app_key,
"timestamp": timestamp, "timestamp": timestamp,
@@ -530,10 +932,10 @@ def _oppo_auth_token() -> str:
headers={"Content-Type": "application/x-www-form-urlencoded"}, headers={"Content-Type": "application/x-www-form-urlencoded"},
) )
if int(data.get("code", -1)) != 0: if int(data.get("code", -1)) != 0:
raise VendorPushError(f"oppo auth failed: {data}") raise VendorPushError(f"oppo auth failed: {_raw_log_summary(data)}")
token = (data.get("data") or {}).get("auth_token") or data.get("auth_token") token = (data.get("data") or {}).get("auth_token") or data.get("auth_token")
if not token: if not token:
raise VendorPushError(f"oppo auth missing auth_token: {data}") raise VendorPushError(f"oppo auth missing auth_token: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), 24 * 3600) return _cache_put(cache_key, str(token), 24 * 3600)
@@ -548,6 +950,11 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
"off_line_ttl": ttl_hours, "off_line_ttl": ttl_hours,
"action_parameters": json.dumps(_click_extras(extras), ensure_ascii=False), "action_parameters": json.dumps(_click_extras(extras), ensure_ascii=False),
} }
if extras.get("notificationId"):
# 只有已落库的站内消息才计入角标;保活提醒等临时推送不应留下无法消除的未读数。
# 客户端打开/已读后会按服务端真实未读总数覆盖,避免长期漂移。
notification["badge_operation_type"] = 1
notification["badge_message_count"] = 1
# 点击落地:OPPO SDK 没有点击回调,参数只能靠服务端点击动作配置送达——action_parameters 的 # 点击落地:OPPO SDK 没有点击回调,参数只能靠服务端点击动作配置送达——action_parameters 的
# 键值对仅在 click_action_type=1/4 时才会注入目标 Activity 的 intent extras(type=0「启动应用」 # 键值对仅在 click_action_type=1/4 时才会注入目标 Activity 的 intent extras(type=0「启动应用」
# 会忽略它,extras 全丢 → 点了没反应,与小米 notify_effect=1 同款坑)。 # 会忽略它,extras 全丢 → 点了没反应,与小米 notify_effect=1 同款坑)。
@@ -573,6 +980,8 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
data = _request_form( data = _request_form(
"POST", "POST",
settings.OPPO_PUSH_SEND_ENDPOINT, settings.OPPO_PUSH_SEND_ENDPOINT,
vendor="oppo",
operation="send_notification",
data={ data={
"auth_token": auth_token, "auth_token": auth_token,
"message": json.dumps(message, ensure_ascii=False), "message": json.dumps(message, ensure_ascii=False),
@@ -580,5 +989,5 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
headers={"Content-Type": "application/x-www-form-urlencoded"}, headers={"Content-Type": "application/x-www-form-urlencoded"},
) )
if int(data.get("code", -1)) != 0: if int(data.get("code", -1)) != 0:
raise VendorPushError(f"oppo push failed: {data}") raise VendorPushError(f"oppo push failed: {_raw_log_summary(data)}")
return data return data
+10
View File
@@ -375,6 +375,16 @@ def _derive_from_platforms(
rows = [p for p in (platforms or []) if isinstance(p, dict)] rows = [p for p in (platforms or []) if isinstance(p, dict)]
src = next((p for p in rows if p.get("role") == "source"), None) src = next((p for p in rows if p.get("role") == "source"), None)
best = next((p for p in rows if p.get("is_best")), None) best = next((p for p in rows if p.get("is_best")), None)
if best is None:
# pricebot 没标 is_best(如全平台 has_dish_diff「相似替换/仅供参考」→ 不认定权威最低价)
# 但仍有有价行 → 兜底取有价行里最低价当参考 best,避免记录级 best_*/saved 整条落 NULL
# (否则首页价 0.00 / 记录页无最低红框 / 省额丢失 /「累计发现可省」漏计)。**候选含源**:
# 源常年全菜、价可信,源本身最便宜时 best 回落源(saved=0、is_source_best=True),与 _derive
# 「全目标缺菜回落源、不虚报省」同一语义;若排除源强选更贵目标,saved 会变负、倒扣「累计
# 发现可省」(get_stats 对 saved_amount_cents 求和不带 >0 过滤)。
priced = [p for p in rows if p.get("price") is not None]
if priced:
best = min(priced, key=lambda p: (p["price"], p.get("display_order") or 0))
source_price_cents = _yuan_to_cents(src.get("price")) if src else None source_price_cents = _yuan_to_cents(src.get("price")) if src else None
best_price_cents = _yuan_to_cents(best.get("price")) if best else None best_price_cents = _yuan_to_cents(best.get("price")) if best else None
saved_amount_cents = None saved_amount_cents = None
+85 -11
View File
@@ -85,6 +85,10 @@ def _dispatch(
push_vars: dict[str, str] | None = None, push_vars: dict[str, str] | None = None,
) -> Notification | None: ) -> Notification | None:
"""落一条站内消息并向该用户设备直推。返回落库行;去重命中/失败返回 None。""" """落一条站内消息并向该用户设备直推。返回落库行;去重命中/失败返回 None。"""
logger.info(
"notification dispatch started user_id=%s type=%s dedup_key=%s",
user_id, type_key, dedup_key,
)
try: try:
row = notif_repo.create_notification( row = notif_repo.create_notification(
db, db,
@@ -111,6 +115,10 @@ def _dispatch(
logger.exception("rollback after notification failure also failed") logger.exception("rollback after notification failure also failed")
return None return None
logger.info(
"notification created user_id=%s type=%s notification_id=%s dedup_key=%s",
user_id, type_key, row.id, dedup_key,
)
_push_to_user_devices(db, row, push_vars) _push_to_user_devices(db, row, push_vars)
return row return row
@@ -124,29 +132,95 @@ def _push_to_user_devices(db: Session, row: Notification, push_vars: dict[str, s
extras.update({str(k): str(v) for k, v in (row.extra or {}).items()}) extras.update({str(k): str(v) for k, v in (row.extra or {}).items()})
extras["notificationId"] = str(row.id) extras["notificationId"] = str(row.id)
for dev in device_repo.list_push_targets(db, user_id=row.user_id): targets = device_repo.list_push_targets(db, user_id=row.user_id)
logger.info(
"push targets resolved user_id=%s type=%s notification_id=%s target_count=%s",
row.user_id, row.type, row.id, len(targets),
)
if not targets:
logger.warning(
"push skipped no targets user_id=%s type=%s notification_id=%s",
row.user_id, row.type, row.id,
)
return
sent = failed = skipped = data_sent = data_failed = 0
for dev in targets:
vendor = vendor_push.normalize_vendor(dev.push_vendor) vendor = vendor_push.normalize_vendor(dev.push_vendor)
if not vendor or vendor not in vendor_push.SUPPORTED_VENDORS: if not vendor or vendor not in vendor_push.SUPPORTED_VENDORS:
continue skipped += 1
if vendor_push.missing_settings(vendor): logger.warning(
# 本地/测试环境凭据不齐 → 只落站内消息,不发真推送(与 push/vendors 的报缺口径一致) "push target skipped unsupported vendor user_id=%s type=%s "
logger.info( "notification_id=%s device_id=%s raw_vendor=%s normalized_vendor=%s",
"skip push (vendor %s not configured) user_id=%s type=%s", row.user_id, row.type, row.id, dev.device_id, dev.push_vendor, vendor,
vendor, row.user_id, row.type,
) )
continue continue
missing = vendor_push.missing_settings(vendor)
if missing:
skipped += 1
# 本地/测试环境凭据不齐 → 只落站内消息,不发真推送(与 push/vendors 的报缺口径一致)
logger.warning(
"push target skipped vendor not configured user_id=%s type=%s "
"notification_id=%s device_id=%s vendor=%s missing_settings=%s",
row.user_id, row.type, row.id, dev.device_id, vendor, missing,
)
continue
logger.info(
"push send started user_id=%s type=%s notification_id=%s "
"device_id=%s vendor=%s",
row.user_id, row.type, row.id, dev.device_id, vendor,
)
try: try:
vendor_push.send_notification( response = vendor_push.send_notification(
vendor, dev.push_token, title=title, body=body, extras=extras vendor, dev.push_token, title=title, body=body, extras=extras
) )
if vendor == "huawei" and row.type in {"feedback_reply", "feedback_reward"}:
# 华为反馈推送联调日志:保留厂商返回码/requestId 等排障信息,
# 请求本身的 token、Authorization 和应用密钥不会进入 response。
logger.info(
"huawei feedback push response user_id=%s type=%s "
"notification_id=%s response=%s",
row.user_id, row.type, row.id, response,
)
logger.info( logger.info(
"push sent user_id=%s type=%s vendor=%s notification_id=%s", "push sent user_id=%s type=%s vendor=%s notification_id=%s device_id=%s",
row.user_id, row.type, vendor, row.id, row.user_id, row.type, vendor, row.id, dev.device_id,
)
sent += 1
except vendor_push.VendorPushError as e:
failed += 1
logger.warning(
"push failed user_id=%s type=%s vendor=%s notification_id=%s "
"device_id=%s error=%s",
row.user_id, row.type, vendor, row.id, dev.device_id, e,
)
try:
data_response = vendor_push.send_data_event(
vendor,
dev.push_token,
event=vendor_push.DATA_EVENT_NOTIFICATION_CREATED,
notification_id=str(row.id),
)
data_sent += 1
logger.info(
"push data event completed user_id=%s type=%s vendor=%s "
"notification_id=%s device_id=%s response=%s",
row.user_id, row.type, vendor, row.id, dev.device_id, data_response,
) )
except vendor_push.VendorPushError as e: except vendor_push.VendorPushError as e:
data_failed += 1
# 透传只负责前台铃铛实时刷新,失败不能影响通知栏消息或站内消息。
logger.warning( logger.warning(
"push failed user_id=%s type=%s vendor=%s: %s", row.user_id, row.type, vendor, e "push data event failed user_id=%s type=%s vendor=%s notification_id=%s "
"device_id=%s error=%s",
row.user_id, row.type, vendor, row.id, dev.device_id, e,
) )
logger.info(
"push dispatch completed user_id=%s type=%s notification_id=%s targets=%s "
"sent=%s failed=%s skipped=%s data_sent=%s data_failed=%s",
row.user_id, row.type, row.id, len(targets),
sent, failed, skipped, data_sent, data_failed,
)
except Exception: # noqa: BLE001 — 渲染/查设备等意外失败同样不外抛 except Exception: # noqa: BLE001 — 渲染/查设备等意外失败同样不外抛
logger.exception("push notification failed user_id=%s type=%s", row.user_id, row.type) logger.exception("push notification failed user_id=%s type=%s", row.user_id, row.type)
+69
View File
@@ -140,6 +140,75 @@ def test_harvest_done_failed_derives_fail_reason(client) -> None:
assert rec.information == "比价过程出错,请稍后重试" # 原文案仍留存 assert rec.information == "比价过程出错,请稍后重试" # 原文案仍留存
def _done_params_platforms_no_isbest() -> dict:
"""id 3304 型:done 帧 platforms 全平台「相似替换/仅供参考」(has_dish_diff),pricebot
一个 is_best 都没标,但有有价目标(美团 57.8 < 源淘宝闪购 60.8)"""
return {
"platforms": [
{"role": "source", "platform_id": "eleme", "platform_name": "淘宝闪购",
"package": "me.ele", "price": 60.8, "is_best": False, "has_dish_diff": False,
"store_name": "窑鸡王", "items": [{"name": "招牌窑鸡 整只-香辣", "qty": 1}]},
{"role": "target", "platform_id": "meituan_waimai", "platform_name": "美团外卖",
"package": "com.sankuai.meituan.takeoutnew", "price": 57.8, "is_best": False,
"has_dish_diff": True},
{"role": "target", "platform_id": "jd_waimai_standalone", "platform_name": "京东外卖",
"package": "com.jd.waimai", "price": 74.9, "is_best": False, "has_dish_diff": True},
],
"information": "美团更便宜(含相似商品替换)",
}
def test_harvest_done_platforms_no_isbest_falls_back_to_cheapest_target(client) -> None:
"""回归(id 3304):platforms 全无 is_best(全平台相似替换)但有有价目标 → best 应兜底取
有价目标里最低那家,不能让 best_*/saved 整条落 NULL(否则首页价 0.00 / 记录页无最低红框 /
省额丢失)"""
tid = _tid()
with SessionLocal() as db:
crud.harvest_running(db, trace_id=tid, user_id=None)
rec, _ = crud.harvest_done(db, trace_id=tid, user_id=None,
done_params=_done_params_platforms_no_isbest())
assert rec.status == "success"
assert rec.best_platform_id == "meituan_waimai" # 有价目标里最低
assert rec.best_price_cents == 5780 # 57.8 元
assert rec.source_price_cents == 6080 # 源淘宝闪购 60.8
assert rec.saved_amount_cents == 300 # 60.8 - 57.8
assert rec.is_source_best is False # 兜底选的是目标,非源
def _done_params_no_isbest_source_cheapest() -> dict:
"""无 is_best 且源本身最便宜:源淘宝闪购 50.0 < 全部 dish-diff 目标(美团 57.8 / 京东 74.9)。
此时不能强选更贵的目标当 best(否则 saved 变负污染累计发现可省),应回落源"""
return {
"platforms": [
{"role": "source", "platform_id": "eleme", "platform_name": "淘宝闪购",
"package": "me.ele", "price": 50.0, "is_best": False, "has_dish_diff": False,
"store_name": "窑鸡王", "items": [{"name": "招牌窑鸡 整只-香辣", "qty": 1}]},
{"role": "target", "platform_id": "meituan_waimai", "platform_name": "美团外卖",
"package": "com.sankuai.meituan.takeoutnew", "price": 57.8, "is_best": False,
"has_dish_diff": True},
{"role": "target", "platform_id": "jd_waimai_standalone", "platform_name": "京东外卖",
"package": "com.jd.waimai", "price": 74.9, "is_best": False, "has_dish_diff": True},
],
"information": "源平台已是最低(其余为相似替换)",
}
def test_harvest_done_no_isbest_source_cheapest_falls_back_to_source(client) -> None:
"""回归:无 is_best 且源最便宜 → best 回落源(与 _derive「全目标缺菜回落源、不虚报省」同语义),
saved=0is_source_best=True,绝不因强选更贵目标而让 saved 变负倒扣累计发现可省"""
tid = _tid()
with SessionLocal() as db:
crud.harvest_running(db, trace_id=tid, user_id=None)
rec, _ = crud.harvest_done(db, trace_id=tid, user_id=None,
done_params=_done_params_no_isbest_source_cheapest())
assert rec.status == "success"
assert rec.best_platform_id == "eleme" # 回落到源(源最便宜)
assert rec.best_price_cents == 5000 # 源 50.0
assert rec.source_price_cents == 5000
assert rec.saved_amount_cents == 0 # 没省到,绝不为负
assert rec.is_source_best is True # 源就是最便宜
def test_harvest_abort_cancels_running(client) -> None: def test_harvest_abort_cancels_running(client) -> None:
tid = _tid() tid = _tid()
with SessionLocal() as db: with SessionLocal() as db:
+95
View File
@@ -7,7 +7,9 @@ mock/真发/变量覆盖/站内联动/设备反查/参数校验。厂商 HTTP
from __future__ import annotations from __future__ import annotations
import json import json
import logging
import httpx
import pytest import pytest
from fastapi.testclient import TestClient from fastapi.testclient import TestClient
@@ -110,6 +112,99 @@ def test_huawei_non_success_code_raises(monkeypatch) -> None:
vendor_push.send_notification("huawei", "bad-token", title="t", body="b") vendor_push.send_notification("huawei", "bad-token", title="t", body="b")
def test_vendor_http_logs_device_token_raw_response_and_elapsed(monkeypatch, caplog) -> None:
vendor_push._token_cache.clear()
def _fake_request(method, url, **kwargs): # noqa: ANN001
if url == vendor_push.settings.HUAWEI_PUSH_TOKEN_ENDPOINT:
return _Resp({"access_token": "raw-access-token", "expires_in": 3600})
return _Resp({"code": "80000000", "msg": "Success", "requestId": "vendor-request-1"})
monkeypatch.setattr(vendor_push.settings, "HUAWEI_PUSH_APP_ID", "10086001")
monkeypatch.setattr(vendor_push.settings, "HUAWEI_PUSH_APP_SECRET", "raw-app-secret")
monkeypatch.setattr(vendor_push.httpx, "request", _fake_request)
with caplog.at_level(logging.INFO, logger=vendor_push.logger.name):
vendor_push.send_notification(
"huawei",
"raw-device-token",
title="需要进入日志的标题",
body="需要进入日志的正文",
extras={"type": "withdraw_success", "notificationId": "90001"},
)
messages = "\n".join(record.getMessage() for record in caplog.records)
assert "operation=authenticate" in messages
assert "operation=send_notification" in messages
assert "request=" in messages
assert "response=" in messages
assert "elapsed_ms=" in messages
assert "vendor-request-1" in messages
assert "withdraw_success" in messages
assert "raw-access-token" in messages
assert "raw-device-token" in messages
assert "需要进入日志的标题" in messages
assert "需要进入日志的正文" in messages
assert "<redacted" in messages
assert "raw-app-secret" not in messages
send_http_log = next(
record.getMessage()
for record in caplog.records
if "operation=send_notification" in record.getMessage()
and "vendor push http completed" in record.getMessage()
)
assert "raw-device-token" in send_http_log
assert "raw-access-token" not in send_http_log # 请求头中的厂商鉴权 token 仍脱敏
def test_vendor_http_network_error_logs_request_context(monkeypatch, caplog) -> None:
def _fake_request(*args, **kwargs): # noqa: ANN002, ANN003
raise httpx.ConnectError("connection refused")
monkeypatch.setattr(vendor_push.settings, "XIAOMI_PUSH_APP_SECRET", "raw-xiaomi-secret")
monkeypatch.setattr(vendor_push.httpx, "request", _fake_request)
with (
caplog.at_level(logging.ERROR, logger=vendor_push.logger.name),
pytest.raises(vendor_push.VendorPushError, match="push http error"),
):
vendor_push.send_notification(
"xiaomi",
"raw-xiaomi-token",
title="title",
body="body",
extras={"type": "push_test"},
)
messages = "\n".join(record.getMessage() for record in caplog.records)
assert "outcome=network_error" in messages
assert "response=<no_response>" in messages
assert "operation=send_notification" in messages
assert "elapsed_ms=" in messages
assert "raw-xiaomi-secret" not in messages
assert "raw-xiaomi-token" in messages
@pytest.mark.parametrize(
("request_kwargs", "device_token"),
[
({"json": {"token": ["honor-device-token"]}}, "honor-device-token"),
({"json": {"message": {"token": ["huawei-device-token"]}}}, "huawei-device-token"),
({"json": {"regId": "vivo-device-token"}}, "vivo-device-token"),
({"data": {"registration_id": "xiaomi-device-token"}}, "xiaomi-device-token"),
(
{"data": {"message": '{"target_value":"oppo-device-token"}'}},
"oppo-device-token",
),
],
)
def test_all_vendor_device_token_fields_remain_visible_in_request_summary(
request_kwargs, device_token
) -> None:
summary = vendor_push._request_summary(request_kwargs)
assert device_token in summary
def test_vendor_aliases_normalize() -> None: def test_vendor_aliases_normalize() -> None:
assert vendor_push.normalize_vendor("华为") == "huawei" assert vendor_push.normalize_vendor("华为") == "huawei"
assert vendor_push.normalize_vendor("HMS") == "huawei" assert vendor_push.normalize_vendor("HMS") == "huawei"