#!/usr/bin/env python3 """ Fake Worker - 模拟真实 Worker 上报采集结果 功能: 1. 模拟成功采集 2. 模拟失败采集(各种错误类型) 3. 模拟重试场景 4. 测试 analytics.replace_package_snapshot 的完整流程 """ import time import json import sys import os # 添加项目路径 sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from analytics import AnalyticsService class FakeWorker: """模拟 Worker 上报采集结果""" def __init__(self, worker_id="FAKE-WORKER-001"): self.worker_id = worker_id self.analytics = AnalyticsService(db_path="runtime/main/monitoring.sqlite3") def report_success(self, package_name: str, batch_tag: str = "test-batch"): """模拟成功采集上报""" print(f"\n[{self.worker_id}] 上报成功采集: {package_name}") summary = { "package_name": package_name, "app_name": f"Test App {package_name}", "app_magic_label": f"magic_{package_name}", "latest_task_key": f"{batch_tag}_{package_name}_ranking_1", "latest_status": "success", "restriction_status": "success", "retryability": "retryable", "latest_test_time": time.time(), "latest_worker_id": self.worker_id, "latest_task_detail": "采集成功", "latest_failure_type": "", "collection_status": "qualified", "collection_status_reason": "success", # 流量数据 "total_traffic_bytes": 5000000, "self_traffic_bytes": 3000000, "server_traffic_bytes": 1500000, "unrecognized_traffic_bytes": 500000, "self_ratio": 60.0, "recognition_ratio": 90.0, # 执行数据 "droidbot_steps": 50, "gui_agent_steps": 20, "duration_seconds": 180.5, "num_nodes": 25, "num_reached_activities": 8, "app_num_total_activities": 15, # 模型数据 "model_flow_count": 10, "model_traffic_bytes": 2000000, "model_eligible": 1, # 元数据 "downloads": 1000000, "source_order": 1, "collection_task_type": "new_app", } domain_rows = [ { "domain": "api.example.com", "domain_type": "api", "traffic_bytes": 2000000, "flow_count": 50, "traffic_ratio": 40.0, "domain_traffic_ratio": 66.7, "organization": "Example Inc", "matched_pattern": "*.example.com", "match_state": "matched", }, { "domain": "cdn.test.com", "domain_type": "cdn", "traffic_bytes": 1000000, "flow_count": 30, "traffic_ratio": 20.0, "domain_traffic_ratio": 33.3, }, ] component_rows = [ { "component_name": package_name, "is_self": True, "traffic_bytes": 3000000, "share_percent": 60.0, }, { "component_name": "com.google.android.gms", "is_self": False, "traffic_bytes": 1500000, "share_percent": 30.0, }, ] source_rows = [] self.analytics.repo.replace_package_snapshot( package_name, summary, domain_rows, component_rows, source_rows ) print(f" ✅ 上报完成") return summary def report_failure(self, package_name: str, error_type: str = "DOWNLOAD_ERROR/404"): """模拟失败采集上报""" print(f"\n[{self.worker_id}] 上报失败采集: {package_name} ({error_type})") summary = { "package_name": package_name, "app_name": f"Test App {package_name}", "latest_task_key": f"test_{package_name}_ranking_1", "latest_status": "failed", "restriction_status": "severe_restricted", "retryability": "non_retryable" if "404" in error_type else "retryable", "latest_test_time": time.time(), "latest_worker_id": self.worker_id, "latest_task_detail": f"下载失败: {error_type}", "latest_failure_type": error_type, "collection_status": "restricted", "collection_status_reason": error_type, # 失败任务无流量数据 "total_traffic_bytes": 0, "self_traffic_bytes": 0, "server_traffic_bytes": 0, "num_nodes": 0, "downloads": 500000, "source_order": 10, } self.analytics.repo.replace_package_snapshot( package_name, summary, [], [], [] ) print(f" ✅ 上报完成") return summary def report_light_restricted(self, package_name: str): """模拟轻度受限采集(成功但质量不足)""" print(f"\n[{self.worker_id}] 上报轻度受限采集: {package_name}") summary = { "package_name": package_name, "app_name": f"Test App {package_name}", "latest_task_key": f"test_{package_name}_ranking_1", "latest_status": "success", "restriction_status": "light_restricted", "retryability": "retryable", "latest_test_time": time.time(), "latest_worker_id": self.worker_id, "latest_task_detail": "采集成功但质量不足", "collection_status": "restricted", "collection_status_reason": "low_quality", # 质量不足:nodes<5 或 self_ratio<30 "total_traffic_bytes": 1000000, "self_traffic_bytes": 200000, # self_ratio=20% < 30% "server_traffic_bytes": 500000, "unrecognized_traffic_bytes": 300000, "self_ratio": 20.0, "recognition_ratio": 70.0, "num_nodes": 3, # < 5 "downloads": 800000, "source_order": 5, } self.analytics.repo.replace_package_snapshot( package_name, summary, [], [], [] ) print(f" ✅ 上报完成") return summary def test_basic_flow(): """测试基本流程""" print("=" * 60) print("测试1:基本采集流程") print("=" * 60) worker = FakeWorker() # 1. 成功采集 worker.report_success("com.test.app1", "batch-2026-06-15") # 2. 失败采集 worker.report_failure("com.test.app2", "DOWNLOAD_ERROR/404") # 3. 轻度受限 worker.report_light_restricted("com.test.app3") print("\n✅ 基本流程测试完成") def test_retry_scenario(): """测试重试场景""" print("\n" + "=" * 60) print("测试2:重试场景(同一应用多次上报)") print("=" * 60) worker = FakeWorker() package = "com.test.retry.app" # 第1次:失败 print("\n--- 第1次尝试(失败)---") worker.report_failure(package, "APP_ERROR/1") time.sleep(1) # 第2次:成功 print("\n--- 第2次尝试(成功)---") worker.report_success(package, "batch-retry") print("\n✅ 重试场景测试完成") def test_query_after_report(): """测试上报后查询""" print("\n" + "=" * 60) print("测试3:上报后查询验证") print("=" * 60) worker = FakeWorker() package = "com.test.query.app" # 上报 worker.report_success(package, "batch-query-test") # 查询验证 print(f"\n查询应用详情: {package}") analytics = AnalyticsService(db_path="runtime/main/monitoring.sqlite3") result = analytics.repo.get_collection_row(package) if result: print(f" ✅ 查询成功:") print(f" app_name: {result.get('app_name')}") print(f" latest_status: {result.get('latest_status')}") print(f" restriction_status: {result.get('restriction_status')}") print(f" collection_status: {result.get('collection_status')}") print(f" total_traffic_bytes: {result.get('total_traffic_bytes')}") print(f" num_nodes: {result.get('num_nodes')}") else: print(f" ❌ 查询失败") print("\n✅ 查询验证测试完成") def test_batch_report(): """测试批量上报""" print("\n" + "=" * 60) print("测试4:批量上报(模拟多个Worker)") print("=" * 60) workers = [ FakeWorker(f"WORKER-{i:03d}") for i in range(1, 4) ] for i, worker in enumerate(workers, 1): print(f"\n--- Worker {i} 上报 ---") worker.report_success(f"com.test.batch.app{i}", "batch-multi") print("\n✅ 批量上报测试完成") if __name__ == "__main__": print("\n" + "🚀" * 30) print("Fake Worker 测试开始") print("🚀" * 30) try: # 运行所有测试 test_basic_flow() test_retry_scenario() test_query_after_report() test_batch_report() print("\n" + "🎉" * 30) print("所有测试通过!") print("🎉" * 30) except Exception as e: print(f"\n❌ 测试失败: {e}") import traceback traceback.print_exc() sys.exit(1)