From ba0706dc4744aed834c8f669841cb7458fb225b6 Mon Sep 17 00:00:00 2001 From: chengma Date: Mon, 10 Aug 2026 17:40:44 +0800 Subject: [PATCH] =?UTF-8?q?perf:=20=E9=87=8F=E5=8C=96=E5=95=86=E5=93=81?= =?UTF-8?q?=E9=A1=B5=E6=89=93=E5=BC=80=E8=80=97=E6=97=B6=E5=B9=B6=E5=8E=BB?= =?UTF-8?q?=E9=87=8D=E6=A3=80=E6=9F=A5=20(#109)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/buyer_main.py | 5 +- client/src/pdd_collect_service.py | 99 +++++++++---- client/src/pdd_device_service.py | 18 ++- client/src/pdd_u2_purchase_adapter.py | 85 +++++++---- client/src/performance_timing.py | 149 +++++++++++++++++++ client/src/task_dispatcher.py | 38 ++++- client/test/test_pdd_collect_service.py | 4 + client/test/test_pdd_device_service.py | 24 +++ client/test/test_pdd_u2_purchase_adapter.py | 32 ++++ client/test/test_performance_timing.py | 107 ++++++++++++++ client/test/test_task_dispatcher.py | 24 +++ client/tools/measure_goods_open.py | 153 ++++++++++++++++++++ docs/client/06-quality-security.md | 27 ++++ 13 files changed, 693 insertions(+), 72 deletions(-) create mode 100644 client/src/performance_timing.py create mode 100644 client/test/test_performance_timing.py create mode 100644 client/tools/measure_goods_open.py diff --git a/client/buyer_main.py b/client/buyer_main.py index fe84d78..21d9ecf 100644 --- a/client/buyer_main.py +++ b/client/buyer_main.py @@ -1,3 +1,6 @@ import sys +from src.performance_timing import configure_performance_logging from src.ui_main import ui_main -sys.exit(ui_main()) \ No newline at end of file + +configure_performance_logging() +sys.exit(ui_main()) diff --git a/client/src/pdd_collect_service.py b/client/src/pdd_collect_service.py index eb2e474..a048caa 100644 --- a/client/src/pdd_collect_service.py +++ b/client/src/pdd_collect_service.py @@ -11,6 +11,7 @@ import math import re import time import xml.etree.ElementTree as ET +from contextlib import nullcontext from dataclasses import dataclass from datetime import datetime, timezone from decimal import Decimal, InvalidOperation, ROUND_HALF_UP @@ -19,6 +20,7 @@ from typing import Any, Callable, Iterable, Mapping, Optional, Sequence from urllib.parse import parse_qs, urlparse from .pdd_device_service import PDD_PACKAGE_NAME, PddDeviceError, PddDeviceService +from .performance_timing import current_performance_trace from .util.get_size_panle_coord import get_size_panel_coord @@ -691,8 +693,11 @@ class PddCollectService: self._check_cancelled() try: - with self._device_service.connect(self._device_address) as device: - self._open_goods(device, goods_url) + session = self._device_service.connect(self._device_address) + with session as device: + self._open_goods( + device, goods_url, session.initial_app_state + ) goods = self._collect_goods_details(device) if not goods.title: raise PddCollectError("PDD_DATA_TITLE_MISSING", "商品页没有可识别的标题") @@ -779,14 +784,32 @@ class PddCollectService: ): raise PddCollectError("PDD_PAGE_OVERALL_TIMEOUT", "PDD 采集超过 10 分钟") - def _open_goods(self, device: Any, goods_url: str) -> None: + def _open_goods( + self, + device: Any, + goods_url: str, + initial_app_state: Optional[Mapping[str, Any]] = None, + ) -> None: + trace = current_performance_trace() try: - current = device.app_current() + current = ( + dict(initial_app_state) + if initial_app_state is not None + else device.app_current() + ) if current.get("package") != PDD_PACKAGE_NAME: - device.app_start(PDD_PACKAGE_NAME) - if not device.app_wait(PDD_PACKAGE_NAME, timeout=10): - raise PddCollectError("DEVICE_APP_START_FAILED", "PDD 应用启动失败") - device.open_url(goods_url) + stage = trace.stage("pdd_start_or_wait") if trace else nullcontext() + with stage: + device.app_start(PDD_PACKAGE_NAME) + if not device.app_wait(PDD_PACKAGE_NAME, timeout=10): + raise PddCollectError( + "DEVICE_APP_START_FAILED", "PDD 应用启动失败" + ) + elif trace is not None: + trace.record("pdd_start_or_wait", 0, "already_foreground") + stage = trace.stage("open_url") if trace else nullcontext() + with stage: + device.open_url(goods_url) except PddCollectError: raise except Exception as exc: @@ -797,30 +820,44 @@ class PddCollectService: ) from exc raise PddCollectError("DEVICE_APP_START_FAILED", f"无法打开 PDD 商品链接:{exc}") from exc - deadline = self._monotonic() + self._page_timeout - while self._monotonic() < deadline: - self._check_cancelled() - current = device.app_current() - xml_data = self._dump_hierarchy(device) - root = _parse_xml(xml_data) - last_labels = _all_labels(root) - pdd_node_count = sum( - 1 - for node in root.iter("node") - if node.get("package") == PDD_PACKAGE_NAME + ready_stage = trace.stage("goods_page_ready") if trace else nullcontext() + with ready_stage: + deadline = self._monotonic() + self._page_timeout + first_dump = True + while self._monotonic() < deadline: + self._check_cancelled() + current = device.app_current() + if first_dump and trace is not None: + with trace.stage("first_dump_hierarchy"): + xml_data = self._dump_hierarchy(device) + first_dump = False + else: + xml_data = self._dump_hierarchy(device) + root = _parse_xml(xml_data) + last_labels = _all_labels(root) + pdd_node_count = sum( + 1 + for node in root.iter("node") + if node.get("package") == PDD_PACKAGE_NAME + ) + # 部分 OPPO/ColorOS 设备会一直把无线调试设置页报告为焦点, + # 即使屏幕和无障碍树已经是 PDD。此时以树中真实包名为准。 + is_pdd_hierarchy = pdd_node_count >= 3 + is_pdd_focused = current.get("package") == PDD_PACKAGE_NAME + if is_pdd_focused or is_pdd_hierarchy: + _raise_special_page(last_labels) + combined = " ".join(last_labels) + loading = "加载中" in combined or "正在加载" in combined + if not loading and any( + marker in combined for marker in _READY_MARKERS + ): + if trace is not None: + trace.checkpoint("end_to_end_total") + return + self._sleep(0.5) + raise PddCollectError( + "PDD_PAGE_TIMEOUT", "等待 PDD 商品详情页加载超时" ) - # 部分 OPPO/ColorOS 设备会一直把无线调试设置页报告为焦点, - # 即使屏幕和无障碍树已经是 PDD。此时以树中真实包名为准。 - is_pdd_hierarchy = pdd_node_count >= 3 - is_pdd_focused = current.get("package") == PDD_PACKAGE_NAME - if is_pdd_focused or is_pdd_hierarchy: - _raise_special_page(last_labels) - combined = " ".join(last_labels) - loading = "加载中" in combined or "正在加载" in combined - if not loading and any(marker in combined for marker in _READY_MARKERS): - return - self._sleep(0.5) - raise PddCollectError("PDD_PAGE_TIMEOUT", "等待 PDD 商品详情页加载超时") def _collect_goods_details(self, device: Any) -> GoodsSnapshot: """保留首屏结果,并向下浏览补采评价数量和店铺名。""" diff --git a/client/src/pdd_device_service.py b/client/src/pdd_device_service.py index 008313a..66c71ae 100644 --- a/client/src/pdd_device_service.py +++ b/client/src/pdd_device_service.py @@ -10,6 +10,8 @@ import threading from contextlib import AbstractContextManager from typing import Any, Callable, Optional +from .performance_timing import current_performance_trace + PDD_PACKAGE_NAME = "com.xunmeng.pinduoduo" @@ -30,10 +32,12 @@ class PddDeviceSession(AbstractContextManager[Any]): self, device: Any, serial: str, + initial_app_state: dict[str, Any], release: Callable[[], None], ) -> None: self._device = device self.serial = serial + self.initial_app_state = dict(initial_app_state) self._release = release self._owner_thread_id = threading.get_ident() self._closed = False @@ -107,8 +111,17 @@ class PddDeviceService: self._active_serial = checked_serial try: - device = self._connector(checked_serial) - current = device.app_current() + trace = current_performance_trace() + if trace is None: + device = self._connector(checked_serial) + else: + with trace.stage("uiautomator2_connect"): + device = self._connector(checked_serial) + if trace is None: + current = device.app_current() + else: + with trace.stage("first_app_current"): + current = device.app_current() if not isinstance(current, dict): raise RuntimeError("uiautomator2 未返回有效的设备状态") except PddDeviceError: @@ -130,6 +143,7 @@ class PddDeviceService: return PddDeviceSession( device, checked_serial, + current, lambda: self._release(checked_serial), ) diff --git a/client/src/pdd_u2_purchase_adapter.py b/client/src/pdd_u2_purchase_adapter.py index a8b8f11..ebc8b5f 100644 --- a/client/src/pdd_u2_purchase_adapter.py +++ b/client/src/pdd_u2_purchase_adapter.py @@ -9,6 +9,7 @@ from __future__ import annotations import re import time import xml.etree.ElementTree as ET +from contextlib import nullcontext from decimal import Decimal, InvalidOperation, ROUND_HALF_UP from typing import Any, Callable, Mapping, Optional from urllib.parse import parse_qs, urlparse @@ -18,6 +19,7 @@ from .pdd_device_service import ( PddDeviceError, PddDeviceService, ) +from .performance_timing import current_performance_trace from .pdd_purchase_adapter import ( PddLivePurchaseAdapter, PddPurchaseAdapter, @@ -266,17 +268,24 @@ class U2PddPurchaseAdapter(PddPurchaseAdapter): try: self._session = self._device_service.connect(self._device_address) self._device = self._session.__enter__() - current = self._device.app_current() + current = self._session.initial_app_state + trace = current_performance_trace() if current.get("package") != PDD_PACKAGE_NAME: - self._device.app_start(PDD_PACKAGE_NAME) - if not self._device.app_wait(PDD_PACKAGE_NAME, timeout=10): - raise PddPurchaseError( - "DEVICE_APP_START_FAILED", - "PDD 应用启动失败", - step="purchase_open_goods", - retryable=True, - ) - self._device.open_url(goods_url) + stage = trace.stage("pdd_start_or_wait") if trace else nullcontext() + with stage: + self._device.app_start(PDD_PACKAGE_NAME) + if not self._device.app_wait(PDD_PACKAGE_NAME, timeout=10): + raise PddPurchaseError( + "DEVICE_APP_START_FAILED", + "PDD 应用启动失败", + step="purchase_open_goods", + retryable=True, + ) + elif trace is not None: + trace.record("pdd_start_or_wait", 0, "already_foreground") + stage = trace.stage("open_url") if trace else nullcontext() + with stage: + self._device.open_url(goods_url) self._wait_for_page("goods", self._page_timeout) except PddPurchaseError: raise @@ -436,26 +445,42 @@ class U2PddPurchaseAdapter(PddPurchaseAdapter): session.__exit__(None, None, None) def _wait_for_page(self, expected: str, timeout: float) -> str: - deadline = self._monotonic() + timeout - last_kind = "unknown" - while self._monotonic() < deadline: - self._check_cancelled("purchase_open_goods") - device = self._require_device() - current = device.app_current() - xml_data = self._dump_hierarchy() - root = _parse_xml(xml_data) - last_kind = _page_kind(root, str(current.get("package") or "")) - if last_kind == expected: - return xml_data - if last_kind in {"captcha", "login_required", "risk_control", "payment"}: - self._raise_special_page(last_kind) - self._sleep(0.25) - raise PddPurchaseError( - "PDD_PAGE_TIMEOUT", - f"等待 PDD {expected} 页面超时,最后页面为 {last_kind}", - step="purchase_open_goods", - retryable=True, - ) + trace = current_performance_trace() + ready_stage = trace.stage("goods_page_ready") if trace else nullcontext() + with ready_stage: + deadline = self._monotonic() + timeout + last_kind = "unknown" + first_dump = True + while self._monotonic() < deadline: + self._check_cancelled("purchase_open_goods") + device = self._require_device() + current = device.app_current() + if first_dump and trace is not None: + with trace.stage("first_dump_hierarchy"): + xml_data = self._dump_hierarchy() + first_dump = False + else: + xml_data = self._dump_hierarchy() + root = _parse_xml(xml_data) + last_kind = _page_kind(root, str(current.get("package") or "")) + if last_kind == expected: + if trace is not None: + trace.checkpoint("end_to_end_total") + return xml_data + if last_kind in { + "captcha", + "login_required", + "risk_control", + "payment", + }: + self._raise_special_page(last_kind) + self._sleep(0.25) + raise PddPurchaseError( + "PDD_PAGE_TIMEOUT", + f"等待 PDD {expected} 页面超时,最后页面为 {last_kind}", + step="purchase_open_goods", + retryable=True, + ) def _wait_for_confirmation_panel(self) -> str: deadline = self._monotonic() + self._panel_timeout diff --git a/client/src/performance_timing.py b/client/src/performance_timing.py new file mode 100644 index 0000000..c4c7466 --- /dev/null +++ b/client/src/performance_timing.py @@ -0,0 +1,149 @@ +"""任务领取到商品页的脱敏分阶段计时。 + +日志字段使用白名单,只包含任务编号、稳定操作名、毫秒耗时和结果分类。 +不要把 URL、控件树、异常正文或业务数据传入本模块。 +""" + +from __future__ import annotations + +import json +import logging +from contextlib import contextmanager +from contextvars import ContextVar +from logging.handlers import RotatingFileHandler +from pathlib import Path +import time +from typing import Callable, Iterator, Optional + +from .db import data_dir + + +PERFORMANCE_LOG_NAME = "task_performance.jsonl" +_LOGGER_NAME = "cmautobuy.performance" +_ACTIVE_TRACE: ContextVar[Optional["TaskPerformanceTrace"]] = ContextVar( + "cmautobuy_task_performance_trace", default=None +) + + +def _default_sink(record: dict[str, object]) -> None: + logging.getLogger(_LOGGER_NAME).info( + json.dumps(record, ensure_ascii=False, separators=(",", ":")) + ) + + +def configure_performance_logging( + directory: Optional[Path] = None, +) -> Path: + """把性能事件写入可轮转 JSONL 文件,重复调用不会增加处理器。""" + + log_directory = directory or data_dir() / "logs" + log_directory.mkdir(parents=True, exist_ok=True) + log_path = log_directory / PERFORMANCE_LOG_NAME + logger = logging.getLogger(_LOGGER_NAME) + if not any( + getattr(handler, "baseFilename", None) == str(log_path.resolve()) + for handler in logger.handlers + ): + handler = RotatingFileHandler( + log_path, + maxBytes=2 * 1024 * 1024, + backupCount=3, + encoding="utf-8", + ) + handler.setFormatter(logging.Formatter("%(message)s")) + logger.addHandler(handler) + logger.setLevel(logging.INFO) + logger.propagate = False + return log_path + + +class TaskPerformanceTrace: + """收集一条任务的阶段耗时,并在绑定任务编号后输出记录。""" + + def __init__( + self, + *, + monotonic: Callable[[], float] = time.monotonic, + sink: Callable[[dict[str, object]], None] = _default_sink, + ) -> None: + self._monotonic = monotonic + self._sink = sink + self._started_at = monotonic() + self._task_id = "" + self._pending: list[tuple[str, int, str]] = [] + + @property + def task_id(self) -> str: + return self._task_id + + def bind_task(self, task_id: str) -> None: + """绑定稳定任务编号,并输出绑定前缓存的领取阶段记录。""" + + value = str(task_id or "").strip() + if not value: + return + if self._task_id and self._task_id != value: + raise RuntimeError("性能计时不能跨任务复用") + self._task_id = value + pending, self._pending = self._pending, [] + for operation, duration_ms, result in pending: + self._emit(operation, duration_ms, result) + + @contextmanager + def stage(self, operation: str) -> Iterator[None]: + """记录一个阶段;异常时只写 failed,不记录异常正文。""" + + started_at = self._monotonic() + try: + yield + except BaseException: + self.record(operation, self._monotonic() - started_at, "failed") + raise + else: + self.record(operation, self._monotonic() - started_at, "succeeded") + + def record(self, operation: str, seconds: float, result: str) -> None: + """记录白名单阶段数据,负耗时会安全归零。""" + + checked_operation = str(operation or "").strip() + checked_result = str(result or "").strip() + if not checked_operation or not checked_result: + raise ValueError("性能操作名和结果分类不能为空") + duration_ms = max(0, round(float(seconds) * 1000)) + if self._task_id: + self._emit(checked_operation, duration_ms, checked_result) + else: + self._pending.append( + (checked_operation, duration_ms, checked_result) + ) + + def checkpoint(self, operation: str, result: str = "succeeded") -> None: + """记录从本轮调度开始到当前时刻的端到端耗时。""" + + self.record(operation, self._monotonic() - self._started_at, result) + + @contextmanager + def activate(self) -> Iterator["TaskPerformanceTrace"]: + """让同一工作线程内的设备和页面适配层取得本计时对象。""" + + token = _ACTIVE_TRACE.set(self) + try: + yield self + finally: + _ACTIVE_TRACE.reset(token) + + def _emit(self, operation: str, duration_ms: int, result: str) -> None: + self._sink( + { + "task_id": self._task_id, + "operation": operation, + "duration_ms": duration_ms, + "result": result, + } + ) + + +def current_performance_trace() -> Optional[TaskPerformanceTrace]: + """返回当前工作线程的任务计时对象;普通调用返回 None。""" + + return _ACTIVE_TRACE.get() diff --git a/client/src/task_dispatcher.py b/client/src/task_dispatcher.py index ecf9f56..8af708f 100644 --- a/client/src/task_dispatcher.py +++ b/client/src/task_dispatcher.py @@ -27,6 +27,7 @@ from .purchase_reconcile_service import ( PurchaseReconcileFactory, PurchaseReconcileService, ) +from .performance_timing import TaskPerformanceTrace from .task_models import ( NewClaimedTask, OutboxEventRecord, @@ -139,6 +140,9 @@ class TaskDispatcher: cancelled: Callable[[], bool] = lambda: False, device_connection_checker: Optional[Callable[[str], None]] = None, task_saved: Callable[[str], None] = lambda _task_id: None, + performance_trace_factory: Callable[[], TaskPerformanceTrace] = ( + TaskPerformanceTrace + ), ) -> None: self._gateway = gateway self._repository = repository @@ -153,6 +157,7 @@ class TaskDispatcher: self._reconcile_factory = purchase_reconcile_factory self._cancelled = cancelled self._task_saved = task_saved + self._performance_trace_factory = performance_trace_factory self._device_connection_checker = ( device_connection_checker or AndroidDeviceService().require_connected @@ -193,6 +198,15 @@ class TaskDispatcher: def execute_one(self) -> TaskDispatchOutcome: """不可逆采购优先核单,其余按 Outbox、本地任务、Admin 顺序处理。""" + trace = self._performance_trace_factory() + with trace.activate(): + return self._execute_one_traced(trace) + + def _execute_one_traced( + self, trace: TaskPerformanceTrace + ) -> TaskDispatchOutcome: + """在同一任务计时上下文中串行处理一轮任务。""" + unresolved_reader = getattr( self._repository, "unresolved_irreversible_purchase", None ) @@ -211,7 +225,9 @@ class TaskDispatcher: raise AndroidDeviceSearchError( "存在只允许核对的采购任务;请先连接并保存原 Android 设备" ) - self._device_connection_checker(self._device_address) + trace.bind_task(reconcile_task.remote_task_id) + with trace.stage("adb_device_check"): + self._device_connection_checker(self._device_address) if self._reconcile_factory is None: raise RuntimeError( f"任务 {reconcile_task.remote_task_id} 只允许核对订单," @@ -234,7 +250,8 @@ class TaskDispatcher: raise AndroidDeviceSearchError( "请先在设置页选择并保存 Android 设备" ) - self._device_connection_checker(self._device_address) + with trace.stage("adb_device_check"): + self._device_connection_checker(self._device_address) if not self.purchase_ready: pending_purchase = self._repository.next_purchase_task() @@ -248,16 +265,19 @@ class TaskDispatcher: include_purchase=self.purchase_ready ) if task is None: - remote = self._gateway.claim_next( - self._client, self.claim_capabilities() - ) + with trace.stage("admin_claim_request"): + remote = self._gateway.claim_next( + self._client, self.claim_capabilities() + ) if remote is None: return TaskDispatchOutcome("no_task", "暂无可领取的任务") + trace.bind_task(remote.task_id) saved_now = False try: - self._repository.add_claimed_task( - admin_task_to_new_claimed_task(remote) - ) + with trace.stage("sqlite_local_save"): + self._repository.add_claimed_task( + admin_task_to_new_claimed_task(remote) + ) saved_now = True except DuplicateTaskError: pass @@ -273,6 +293,8 @@ class TaskDispatcher: if saved_now: self._task_saved(remote.task_id) + trace.bind_task(task.remote_task_id) + if task.task_type is TaskType.COLLECT: service = CollectTaskService( self._gateway, diff --git a/client/test/test_pdd_collect_service.py b/client/test/test_pdd_collect_service.py index cc551d8..d22511b 100644 --- a/client/test/test_pdd_collect_service.py +++ b/client/test/test_pdd_collect_service.py @@ -34,6 +34,7 @@ class FakeCollectDevice: self.clicks = [] self.swipes = [] self.swipe_panel_states = [] + self.app_wait_calls = 0 def app_current(self): return {"package": "com.xunmeng.pinduoduo"} @@ -42,6 +43,7 @@ class FakeCollectDevice: raise AssertionError("PDD 已经在前台,不应重复启动") def app_wait(self, _package, timeout=10): + self.app_wait_calls += 1 return 123 def open_url(self, url): @@ -424,6 +426,7 @@ class PddCollectParserTest(unittest.TestCase): ) self.assertEqual(data["source"]["device_address"], "USB-001") self.assertIsNone(data["purchase"]) + self.assertEqual(device.app_wait_calls, 0) def test_collect_scrolls_before_spec_panel_and_accumulates_metadata(self): initial_page = self.home_xml.replace( @@ -510,6 +513,7 @@ class PddCollectParserTest(unittest.TestCase): self.assertEqual(result.goods_id, "123") self.assertTrue(device.app_started) + self.assertEqual(device.app_wait_calls, 1) def test_price_is_sampled_once_per_available_color(self): device = FakeCollectDevice(self.home_xml, self.spec_xml) diff --git a/client/test/test_pdd_device_service.py b/client/test/test_pdd_device_service.py index fb147aa..fe29340 100644 --- a/client/test/test_pdd_device_service.py +++ b/client/test/test_pdd_device_service.py @@ -4,10 +4,15 @@ import threading import unittest from src.pdd_device_service import PddDeviceError, PddDeviceService +from src.performance_timing import TaskPerformanceTrace class FakeDevice: + def __init__(self): + self.current_calls = 0 + def app_current(self): + self.current_calls += 1 return {"package": "com.xunmeng.pinduoduo"} @@ -25,6 +30,25 @@ class PddDeviceServiceTest(unittest.TestCase): pass self.assertEqual(connected, ["USB-001", "USB-002"]) + def test_connect_exposes_initial_state_and_records_two_separate_stages(self): + device = FakeDevice() + records = [] + trace = TaskPerformanceTrace(sink=records.append) + trace.bind_task("COL-001") + + with trace.activate(): + session = PddDeviceService(lambda _serial: device).connect("USB-001") + self.assertEqual( + session.initial_app_state["package"], + "com.xunmeng.pinduoduo", + ) + self.assertEqual(device.current_calls, 1) + self.assertEqual( + [record["operation"] for record in records], + ["uiautomator2_connect", "first_app_current"], + ) + session.__exit__(None, None, None) + def test_invalid_serial_is_rejected_before_connect(self): with self.assertRaisesRegex(PddDeviceError, "设备号无效") as raised: PddDeviceService(lambda _serial: FakeDevice()).connect("bad serial") diff --git a/client/test/test_pdd_u2_purchase_adapter.py b/client/test/test_pdd_u2_purchase_adapter.py index 63c7ff8..31e3e62 100644 --- a/client/test/test_pdd_u2_purchase_adapter.py +++ b/client/test/test_pdd_u2_purchase_adapter.py @@ -59,6 +59,7 @@ class FakeDevice: self.quantity = 1 self.clicks = [] self.opened_urls = [] + self.app_wait_calls = 0 def app_current(self): return {"package": "com.xunmeng.pinduoduo"} @@ -67,6 +68,7 @@ class FakeDevice: pass def app_wait(self, _package, timeout=10): + self.app_wait_calls += 1 return True def open_url(self, url): @@ -87,6 +89,25 @@ class FakeDevice: raise AssertionError(f"未预期的 selector: {selector}") +class ColdStartDevice(FakeDevice): + def __init__(self) -> None: + super().__init__() + self.current_calls = 0 + self.app_started = False + + def app_current(self): + self.current_calls += 1 + package = ( + "com.android.launcher" + if self.current_calls == 1 + else "com.xunmeng.pinduoduo" + ) + return {"package": package} + + def app_start(self, _package): + self.app_started = True + + class U2PddPurchaseAdapterTest(unittest.TestCase): def _adapter(self, device, calls): def select_color_fn(_device, _xml, target, **_kwargs): @@ -150,6 +171,7 @@ class U2PddPurchaseAdapterTest(unittest.TestCase): self.assertEqual(calls, [("color", "黑色"), ("size", "3XL【140-165斤】")]) # 只点击一次商品页采购入口;最终“提交订单”从不点击。 self.assertEqual(len(device.clicks), 1) + self.assertEqual(device.app_wait_calls, 0) def test_unsupported_dynamic_option_stops_before_click(self): device = FakeDevice() @@ -163,6 +185,16 @@ class U2PddPurchaseAdapterTest(unittest.TestCase): self.assertEqual(device.clicks, []) adapter.close() + def test_cold_start_waits_for_pdd_once(self): + device = ColdStartDevice() + adapter = self._adapter(device, []) + + adapter.open_goods(GOODS_URL) + + self.assertTrue(device.app_started) + self.assertEqual(device.app_wait_calls, 1) + adapter.close() + def test_captcha_stops_without_waiting_or_clicking(self): device = FakeDevice( ' None: + self.value = 10.0 + + def __call__(self) -> float: + return self.value + + +class TaskPerformanceTraceTest(unittest.TestCase): + def test_stages_are_monotonic_scoped_and_use_only_safe_fields(self): + clock = MutableClock() + records = [] + trace = TaskPerformanceTrace(monotonic=clock, sink=records.append) + + with trace.activate(): + self.assertIs(current_performance_trace(), trace) + with trace.stage("admin_claim_request"): + clock.value = 10.125 + trace.bind_task("COL-001") + clock.value = 10.250 + trace.checkpoint("end_to_end_total") + + self.assertIsNone(current_performance_trace()) + self.assertEqual( + records, + [ + { + "task_id": "COL-001", + "operation": "admin_claim_request", + "duration_ms": 125, + "result": "succeeded", + }, + { + "task_id": "COL-001", + "operation": "end_to_end_total", + "duration_ms": 250, + "result": "succeeded", + }, + ], + ) + self.assertEqual( + set(records[0]), + {"task_id", "operation", "duration_ms", "result"}, + ) + + def test_negative_duration_is_clamped_and_cross_task_reuse_is_rejected(self): + records = [] + trace = TaskPerformanceTrace(sink=records.append) + trace.bind_task("COL-001") + trace.record("sqlite_local_save", -1, "succeeded") + + self.assertEqual(records[0]["duration_ms"], 0) + with self.assertRaisesRegex(RuntimeError, "不能跨任务"): + trace.bind_task("COL-002") + + def test_failed_stage_does_not_log_exception_details(self): + records = [] + trace = TaskPerformanceTrace(sink=records.append) + trace.bind_task("COL-SECRET") + + with self.assertRaisesRegex(RuntimeError, "token-secret"): + with trace.stage("open_url"): + raise RuntimeError("token-secret") + + self.assertEqual(records[0]["result"], "failed") + self.assertNotIn("token-secret", str(records[0])) + + def test_configured_jsonl_log_can_be_viewed_and_has_safe_schema(self): + with tempfile.TemporaryDirectory() as temporary: + path = configure_performance_logging(Path(temporary)) + trace = TaskPerformanceTrace() + trace.bind_task("COL-001") + trace.record("first_app_current", 0.125, "succeeded") + + logger = logging.getLogger("cmautobuy.performance") + for handler in logger.handlers: + handler.flush() + record = json.loads(path.read_text(encoding="utf-8").strip()) + self.assertEqual(record["duration_ms"], 125) + self.assertEqual( + set(record), + {"task_id", "operation", "duration_ms", "result"}, + ) + + for handler in list(logger.handlers): + if getattr(handler, "baseFilename", None) == str(path.resolve()): + logger.removeHandler(handler) + handler.close() + + +if __name__ == "__main__": + unittest.main() diff --git a/client/test/test_task_dispatcher.py b/client/test/test_task_dispatcher.py index 5b543b7..a267049 100644 --- a/client/test/test_task_dispatcher.py +++ b/client/test/test_task_dispatcher.py @@ -12,6 +12,7 @@ from src.pdd_purchase_adapter import ( PddPurchaseAdapter, PurchasePageState, ) +from src.performance_timing import TaskPerformanceTrace from src.task_dispatcher import TaskDispatcher, admin_task_to_new_claimed_task from src.task_models import TaskStatus, TaskType from src.task_repository import TaskRepository @@ -134,6 +135,7 @@ class TaskDispatcherTest(unittest.TestCase): purchase_mode="dry_run", device_checker=lambda _serial: None, task_saved=lambda _task_id: None, + performance_trace_factory=TaskPerformanceTrace, ): purchase_factory = None if purchase_ready: @@ -161,8 +163,30 @@ class TaskDispatcherTest(unittest.TestCase): purchase_mode=purchase_mode, device_connection_checker=device_checker, task_saved=task_saved, + performance_trace_factory=performance_trace_factory, ) + def test_claim_save_and_device_check_have_scoped_timing_records(self): + task = purchase_task() + self.gateway.enqueue_task(task, self.client.client_id) + records = [] + + outcome = self._dispatcher( + purchase_ready=True, + performance_trace_factory=lambda: TaskPerformanceTrace( + sink=records.append + ), + ).execute_one() + + self.assertEqual(outcome.kind, "succeeded") + operations = [record["operation"] for record in records] + self.assertEqual( + operations[:3], + ["adb_device_check", "admin_claim_request", "sqlite_local_save"], + ) + self.assertTrue(all(record["task_id"] == task.task_id for record in records)) + self.assertTrue(all(record["duration_ms"] >= 0 for record in records)) + def test_capability_only_includes_purchase_when_adapter_is_ready(self): collect_only = self._dispatcher(purchase_ready=False) ready = self._dispatcher(purchase_ready=True) diff --git a/client/tools/measure_goods_open.py b/client/tools/measure_goods_open.py new file mode 100644 index 0000000..49c9dd4 --- /dev/null +++ b/client/tools/measure_goods_open.py @@ -0,0 +1,153 @@ +"""测量同一 Android 设备冷/热启动打开 PDD 商品页的耗时。 + +本工具只启动或停止 PDD、打开商品链接并读取控件树;不点击规格、下单或付款。 +输出只包含运行模式、稳定阶段名和毫秒值,不输出设备号、URL 或控件树。 +""" + +from __future__ import annotations + +import argparse +import json +import math +import statistics +import subprocess +import time +import xml.etree.ElementTree as ET +from collections.abc import Callable +from typing import Any + + +PDD_PACKAGE_NAME = "com.xunmeng.pinduoduo" +READY_MARKERS = ("发起拼单", "立即购买", "单独购买", "免拼购买", "快要抢光") +STOP_MARKERS = ("手机号登录", "请完成验证", "安全验证", "拖动滑块") + + +def elapsed_ms(started_at: float) -> int: + return max(0, round((time.monotonic() - started_at) * 1000)) + + +def measure_call(operation: str, call: Callable[[], Any]) -> tuple[Any, int]: + started_at = time.monotonic() + value = call() + return value, elapsed_ms(started_at) + + +def labels(xml_data: str) -> list[str]: + root = ET.fromstring(xml_data) + values = [] + for node in root.iter("node"): + value = str(node.get("text") or node.get("content-desc") or "").strip() + if value: + values.append(value) + return values + + +def wait_goods_ready(device: Any, first_xml: str, timeout: float) -> None: + deadline = time.monotonic() + timeout + xml_data = first_xml + while time.monotonic() < deadline: + current_labels = labels(xml_data) + combined = " ".join(current_labels) + if any(marker in combined for marker in STOP_MARKERS): + raise RuntimeError("PDD 出现登录或安全验证,测量已停止") + loading = "加载中" in combined or "正在加载" in combined + if not loading and any(marker in combined for marker in READY_MARKERS): + return + time.sleep(0.25) + xml_data = str(device.dump_hierarchy()) + raise TimeoutError("等待 PDD 商品页就绪超时") + + +def measure_once(serial: str, goods_url: str, cold: bool) -> dict[str, int]: + import uiautomator2 as u2 + + if cold: + subprocess.run( + ["adb", "-s", serial, "shell", "am", "force-stop", PDD_PACKAGE_NAME], + check=True, + capture_output=True, + text=True, + ) + time.sleep(1) + else: + warm_device = u2.connect(serial) + warm_device.app_start(PDD_PACKAGE_NAME) + if not warm_device.app_wait(PDD_PACKAGE_NAME, timeout=10): + raise RuntimeError("PDD 热启动准备失败") + + total_started_at = time.monotonic() + result: dict[str, int] = {} + _, result["adb_device_check"] = measure_call( + "adb_device_check", + lambda: subprocess.run( + ["adb", "-s", serial, "get-state"], + check=True, + capture_output=True, + text=True, + ), + ) + device, result["uiautomator2_connect"] = measure_call( + "uiautomator2_connect", lambda: u2.connect(serial) + ) + current, result["first_app_current"] = measure_call( + "first_app_current", device.app_current + ) + if current.get("package") != PDD_PACKAGE_NAME: + started_at = time.monotonic() + device.app_start(PDD_PACKAGE_NAME) + if not device.app_wait(PDD_PACKAGE_NAME, timeout=10): + raise RuntimeError("PDD 冷启动失败") + result["pdd_start_or_wait"] = elapsed_ms(started_at) + else: + result["pdd_start_or_wait"] = 0 + + _, result["open_url"] = measure_call( + "open_url", lambda: device.open_url(goods_url) + ) + first_xml, result["first_dump_hierarchy"] = measure_call( + "first_dump_hierarchy", lambda: str(device.dump_hierarchy()) + ) + ready_started_at = time.monotonic() + wait_goods_ready(device, first_xml, timeout=30) + result["goods_page_ready"] = elapsed_ms(ready_started_at) + result["end_to_end_total"] = elapsed_ms(total_started_at) + return result + + +def percentile_95(values: list[int]) -> int: + ordered = sorted(values) + return ordered[max(0, math.ceil(len(ordered) * 0.95) - 1)] + + +def summarize(rows: list[dict[str, int]]) -> dict[str, dict[str, int]]: + return { + operation: { + "median_ms": round(statistics.median(row[operation] for row in rows)), + "p95_ms": percentile_95([row[operation] for row in rows]), + } + for operation in rows[0] + } + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("serial") + parser.add_argument("goods_url") + parser.add_argument("--runs", type=int, default=5) + args = parser.parse_args() + if args.runs < 1: + parser.error("--runs 必须大于 0") + + output = {} + for mode, cold in (("cold", True), ("hot", False)): + rows = [ + measure_once(args.serial, args.goods_url, cold) + for _ in range(args.runs) + ] + output[mode] = {"runs": args.runs, "stages": summarize(rows)} + print(json.dumps(output, ensure_ascii=False, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/docs/client/06-quality-security.md b/docs/client/06-quality-security.md index d8dbd32..4cb86f7 100644 --- a/docs/client/06-quality-security.md +++ b/docs/client/06-quality-security.md @@ -158,6 +158,33 @@ PDD 解析测试优先使用脱敏的 XML 固件,不要求每次连接真实 - `purchase_irreversible_step_entered` - `order_reconcile_completed` +### 6.1 领取到商品页性能日志 + +正式 Client 把任务领取到商品页就绪的白名单耗时写入 +`data/logs/task_performance.jsonl`。每行是一个 JSON 对象,只允许下面四个字段: + +| 字段 | 含义 | +|---|---| +| `task_id` | Admin 分配的稳定任务编号 | +| `operation` | 稳定阶段名,例如 `admin_claim_request`、`uiautomator2_connect` | +| `duration_ms` | 单调时钟计算的非负毫秒值 | +| `result` | `succeeded`、`failed` 或 `already_foreground` 等结果分类 | + +阶段至少包含 Admin 领取、本地保存、ADB 检查、uiautomator2 连接、首次应用状态、 +PDD 启动/等待、打开链接、首次控件树、商品页就绪和 `end_to_end_total`。 +`goods_page_ready` 内部包含首次及后续控件树读取,`end_to_end_total` 又包含所有阶段, +所以阶段耗时不能直接相加后与总耗时比较;它们是有意重叠的嵌套区间。 + +性能日志不得增加 URL、设备号、控件树、异常正文、凭据或个人信息。日志按 2 MiB +轮转并保留 3 个旧文件。需要比较真机冷/热启动时,从 `client/` 执行: + +```powershell +C:/Python310/python.exe tools/measure_goods_open.py <设备号> <商品链接> --runs 5 +``` + +测量工具只启动或停止 PDD、打开商品页和读取控件树,不点击规格、下单或付款; +输出只包含模式、阶段名、中位数和 P95。遇到登录或安全验证会立即停止。 + 失败 Artifact 目录建议: ```text