diff --git a/app/integrations/vendor_push.py b/app/integrations/vendor_push.py index 54c0e0e..e7167b1 100644 --- a/app/integrations/vendor_push.py +++ b/app/integrations/vendor_push.py @@ -22,7 +22,7 @@ import uuid from collections.abc import Callable from dataclasses import dataclass from typing import Any -from urllib.parse import quote +from urllib.parse import quote, urlsplit import httpx @@ -69,6 +69,18 @@ class _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: if not push_vendor: @@ -115,8 +127,16 @@ def send_notification( if mock: logger.info( - "[mock push] vendor=%s token=%s... title=%s body=%s extras=%s", - vendor, token[:12], title, body, extras, + "vendor push mock vendor=%s request=%s", + vendor, + _log_summary( + { + "push_token": token, + "title": title, + "body": body, + "extras": extras, + } + ), ) return { "mock": True, @@ -133,7 +153,28 @@ def send_notification( "xiaomi": _send_xiaomi, "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( @@ -159,13 +200,120 @@ def _require(value: str, name: str) -> str: 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"" + + +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[""] = 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"") + 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]}…" + 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]}…" + 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]}…" + 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 _request_json( method: str, url: str, *, + vendor: str, + operation: str, expected_status: tuple[int, ...] = (200,), **kwargs: Any, ) -> dict[str, Any]: + started = time.perf_counter() + request_summary = _request_summary(kwargs) try: resp = httpx.request( method, @@ -174,41 +322,82 @@ def _request_json( **kwargs, ) 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, + "", + _elapsed_ms(started), + type(e).__name__, + ) 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]) + 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}") try: - return resp.json() + data = resp.json() 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 + logger.info( + "vendor push http completed vendor=%s operation=%s method=%s endpoint=%s " + "outcome=success 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, data), + _elapsed_ms(started), + ) + return data def _request_form( method: str, url: str, *, + vendor: str, + operation: str, expected_status: tuple[int, ...] = (200,), **kwargs: Any, ) -> dict[str, Any]: - try: - resp = httpx.request( - method, - url, - timeout=settings.PUSH_REQUEST_TIMEOUT_SEC, - **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 + return _request_json( + method, + url, + vendor=vendor, + operation=operation, + expected_status=expected_status, + **kwargs, + ) def _cache_get(key: str) -> str | None: @@ -234,6 +423,8 @@ def _honor_access_token() -> str: data = _request_form( "POST", settings.HONOR_PUSH_TOKEN_ENDPOINT, + vendor="honor", + operation="authenticate", data={ "grant_type": "client_credentials", "client_id": client_id, @@ -243,7 +434,7 @@ def _honor_access_token() -> str: ) token = data.get("access_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")) @@ -270,6 +461,8 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di data = _request_json( "POST", settings.HONOR_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id), + vendor="honor", + operation="send_notification", json=payload, headers={ "Content-Type": "application/json; charset=UTF-8", @@ -279,7 +472,7 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di ) code = data.get("code") 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 @@ -294,6 +487,8 @@ def _huawei_access_token() -> str: data = _request_form( "POST", settings.HUAWEI_PUSH_TOKEN_ENDPOINT, + vendor="huawei", + operation="authenticate", data={ "grant_type": "client_credentials", "client_id": app_id, @@ -303,7 +498,7 @@ def _huawei_access_token() -> str: ) token = data.get("access_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")) @@ -333,6 +528,8 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d data = _request_json( "POST", settings.HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id), + vendor="huawei", + operation="send_notification", json=payload, headers={ "Content-Type": "application/json; charset=UTF-8", @@ -340,7 +537,7 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d }, ) 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 @@ -357,6 +554,8 @@ def _vivo_auth_token() -> str: data = _request_json( "POST", settings.VIVO_PUSH_AUTH_ENDPOINT, + vendor="vivo", + operation="authenticate", json={ "appId": app_id, "appKey": app_key, @@ -366,10 +565,10 @@ def _vivo_auth_token() -> str: headers={"Content-Type": "application/json"}, ) 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") 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) @@ -401,6 +600,8 @@ def _send_vivo(token: str, title: str, body: str, extras: dict[str, str]) -> dic data = _request_json( "POST", settings.VIVO_PUSH_SEND_ENDPOINT, + vendor="vivo", + operation="send_notification", json=payload, headers={ "Content-Type": "application/json", @@ -408,7 +609,7 @@ def _send_vivo(token: str, title: str, body: str, extras: dict[str, str]) -> dic }, ) 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 @@ -446,14 +647,16 @@ def _send_xiaomi(token: str, title: str, body: str, extras: dict[str, str]) -> d data = _request_form( "POST", settings.XIAOMI_PUSH_SEND_ENDPOINT, + vendor="xiaomi", + operation="send_notification", data=form, headers={"Authorization": f"key={app_secret}"}, ) code = data.get("code") 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"): - raise VendorPushError(f"xiaomi push failed: {data}") + raise VendorPushError(f"xiaomi push failed: {_raw_log_summary(data)}") return data @@ -522,6 +725,8 @@ def _oppo_auth_token() -> str: data = _request_form( "POST", settings.OPPO_PUSH_AUTH_ENDPOINT, + vendor="oppo", + operation="authenticate", data={ "app_key": app_key, "timestamp": timestamp, @@ -530,10 +735,10 @@ def _oppo_auth_token() -> str: headers={"Content-Type": "application/x-www-form-urlencoded"}, ) 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") 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) @@ -573,6 +778,8 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic data = _request_form( "POST", settings.OPPO_PUSH_SEND_ENDPOINT, + vendor="oppo", + operation="send_notification", data={ "auth_token": auth_token, "message": json.dumps(message, ensure_ascii=False), @@ -580,5 +787,5 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic headers={"Content-Type": "application/x-www-form-urlencoded"}, ) 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 diff --git a/tests/test_push_center.py b/tests/test_push_center.py index 9b44ac1..585f8dc 100644 --- a/tests/test_push_center.py +++ b/tests/test_push_center.py @@ -7,7 +7,9 @@ mock/真发/变量覆盖/站内联动/设备反查/参数校验。厂商 HTTP from __future__ import annotations import json +import logging +import httpx import pytest 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") +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 " 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=" 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: assert vendor_push.normalize_vendor("华为") == "huawei" assert vendor_push.normalize_vendor("HMS") == "huawei"