autool-dispatcher/analytics.py
2026-06-17 19:50:39 +08:00

3656 lines
164 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import csv
import glob
import json
import os
import sqlite3
import threading
import time
from contextlib import contextmanager
from datetime import datetime
from typing import Any, Dict, Iterable, List, Optional, Tuple
from zoneinfo import ZoneInfo
from config import (
ANALYTICS_ARTIFACT_WAIT_SECONDS,
ANALYTICS_ARTIFACT_FILE_WAIT_SECONDS,
ANALYTICS_JOB_POLL_SECONDS,
ANALYTICS_TPDPI_APP_LIST,
ANALYTICS_TPDPI_URL_LIB,
ANALYTICS_TRAFFIC_ROOT,
ANALYTICS_TRAVERSAL_ROOT,
MODEL_TRAFFIC_THRESHOLD,
MONITORING_DB_PATH,
MONITORING_TIMEZONE,
normalize_worker_id,
)
from log_manager import logger
from result_codes import (
APP_ERROR_DESC,
BUSINESS_ERROR_DESC,
DOWNLOAD_ERROR_DESC,
INFRA_ERROR_DESC,
AppError,
BusinessError,
DownloadError,
ErrorCategory,
ErrorInfo,
InfraError,
)
SHANGHAI_TZ = ZoneInfo(MONITORING_TIMEZONE)
NON_RETRYABLE_DOWNLOAD_BUCKETS = (
("lt_1k", "<1K", 0, 1_000),
("1k_to_10k", "1K-10K", 1_000, 10_000),
("10k_to_100k", "10K-100K", 10_000, 100_000),
("100k_to_1m", "100K-1M", 100_000, 1_000_000),
("1m_to_10m", "1M-10M", 1_000_000, 10_000_000),
("10m_to_100m", "10M-100M", 10_000_000, 100_000_000),
("100m_plus", "100M+", 100_000_000, None),
)
NON_RETRYABLE_DOWNLOAD_SPLIT_THRESHOLD = 10_000
LIGHT_RESTRICTED_MIN_NUM_NODES = 5
AUTO_PENDING_PHYSICAL_REASON = "auto_pending_physical_reroute"
ACTIVE_CATALOG_WHERE = "COALESCE(is_active, 1) = 1"
DEFAULT_INCREMENTAL_TOP_N = 3000
AUTO_PENDING_PHYSICAL_FAILURE_TYPES = {
f"APP_ERROR/{int(AppError.CRASH)}",
f"APP_ERROR/{int(AppError.ROOT_MODE_UNSUPPORTED)}",
}
def _safe_json_dumps(payload: Any) -> str:
return json.dumps(payload, ensure_ascii=False, sort_keys=True)
def _safe_json_loads(payload: Optional[str], fallback: Optional[Any] = None) -> Any:
if not payload:
return fallback
try:
return json.loads(payload)
except (TypeError, ValueError):
return fallback
def _format_ts(value: Optional[float]) -> str:
if not value:
return ""
return datetime.fromtimestamp(float(value), SHANGHAI_TZ).strftime("%Y-%m-%d %H:%M:%S")
def _normalize_incremental_batch_tag(value: Any) -> str:
"""返回单个标签字符串(向后兼容)。"""
return str(value or "").strip()
def _normalize_incremental_batch_tags(value: Any) -> List[str]:
"""将存储的标签值解析为去重后的有序列表。支持 JSON 数组或旧版单标签字符串。"""
if value is None:
return []
if isinstance(value, (list, tuple)):
result: List[str] = []
seen = set()
for item in value:
tag = str(item or "").strip()
if tag and tag not in seen:
seen.add(tag)
result.append(tag)
return result
text = str(value).strip()
if not text:
return []
if text.startswith("["):
try:
parsed = json.loads(text)
if isinstance(parsed, list):
return _normalize_incremental_batch_tags(parsed)
except (TypeError, ValueError):
pass
return [text]
def _incremental_batch_tags_to_json(tags: List[str]) -> str:
"""将标签列表序列化为 JSON 数组字符串。"""
normalized = _normalize_incremental_batch_tags(tags)
return json.dumps(normalized, ensure_ascii=False)
def _incremental_batch_tag_in_tags_sql(tag: str) -> str:
"""返回在 SQL 中检查 incremental_batch_tag JSON 数组是否包含指定 tag 的表达式。"""
escaped = str(tag or "").replace("'", "''")
return f"incremental_batch_tag LIKE '%\"{escaped}\"%'"
def _add_incremental_batch_tag(existing_tags: Any, new_tag: str) -> List[str]:
"""向现有标签列表中添加一个新标签(去重)。"""
tags = _normalize_incremental_batch_tags(existing_tags)
normalized_new = str(new_tag or "").strip()
if normalized_new and normalized_new not in tags:
tags.append(normalized_new)
return tags
def _remove_incremental_batch_tag(existing_tags: Any, tag_to_remove: str) -> List[str]:
"""从现有标签列表中移除指定标签。"""
tags = _normalize_incremental_batch_tags(existing_tags)
normalized_remove = str(tag_to_remove or "").strip()
if not normalized_remove:
return tags
return [t for t in tags if t != normalized_remove]
def _normalize_app_magic_label(value: Any) -> str:
return str(value or "").strip()
def _normalize_last_updated(value: Any) -> str:
text = str(value or "").strip()
if not text:
return ""
normalized = text.replace("-", "/").replace(".", "/")
for fmt in ("%Y/%m/%d", "%Y/%-m/%-d"):
try:
parsed = datetime.strptime(normalized, fmt)
return parsed.strftime("%Y/%m/%d")
except ValueError:
continue
parts = normalized.split("/")
if len(parts) == 3:
try:
year, month, day = [int(part) for part in parts]
return datetime(year, month, day).strftime("%Y/%m/%d")
except (TypeError, ValueError):
return ""
return ""
def _last_updated_ts(value: Any) -> float:
normalized = _normalize_last_updated(value)
if not normalized:
return 0.0
try:
return datetime.strptime(normalized, "%Y/%m/%d").timestamp()
except ValueError:
return 0.0
def _last_updated_interval_days(new_value: Any, old_value: Any) -> int:
new_ts = _last_updated_ts(new_value)
old_ts = _last_updated_ts(old_value)
if not new_ts or not old_ts or new_ts <= old_ts:
return 0
return max(1, int((new_ts - old_ts) // 86400))
def _normalize_top_n(value: Any, default: int = DEFAULT_INCREMENTAL_TOP_N) -> int:
try:
normalized = int(value)
except (TypeError, ValueError):
normalized = default
return normalized if normalized > 0 else default
def _top_n_label(top_n: int) -> str:
return f"榜单 Top{int(top_n)}"
def _decode_mixed_csv_line(raw_line: bytes) -> str:
for encoding in ("utf-8-sig", "utf-8", "gb18030"):
try:
return raw_line.decode(encoding)
except UnicodeDecodeError:
continue
return raw_line.decode("gb18030", errors="replace")
def _read_csv_lines(csv_path: str) -> List[str]:
with open(csv_path, "rb") as handle:
return [_decode_mixed_csv_line(line) for line in handle]
def _read_dict_rows(csv_path: str) -> Iterable[Dict[str, str]]:
reader = csv.DictReader(_read_csv_lines(csv_path))
for row in reader:
normalized = {}
for key, value in (row or {}).items():
normalized[str(key or "").strip()] = str(value or "").strip()
yield normalized
def _parse_int_value(value: Any) -> int:
try:
return int(float(value or 0))
except (TypeError, ValueError):
return 0
def _parse_float_value(value: Any) -> float:
try:
return float(value or 0.0)
except (TypeError, ValueError):
return 0.0
def _parse_downloads_value(value: Any) -> Optional[int]:
text = str(value or "").strip().upper()
if not text:
return None
normalized = text.replace(",", "").replace(" ", "")
if normalized.endswith("+"):
normalized = normalized[:-1]
multiplier = 1
if normalized.endswith("K"):
multiplier = 1_000
normalized = normalized[:-1]
elif normalized.endswith("M"):
multiplier = 1_000_000
normalized = normalized[:-1]
elif normalized.endswith("B"):
multiplier = 1_000_000_000
normalized = normalized[:-1]
elif normalized.endswith("W") or normalized.endswith(""):
multiplier = 10_000
normalized = normalized[:-1]
elif normalized.endswith("亿"):
multiplier = 100_000_000
normalized = normalized[:-1]
try:
return int(float(normalized) * multiplier)
except (TypeError, ValueError):
return None
def _extract_downloads_from_task_payload(task_payload: Any) -> Optional[int]:
if isinstance(task_payload, dict):
payload = task_payload
else:
payload = _safe_json_loads(task_payload, {}) or {}
if not isinstance(payload, dict):
return None
original_row = payload.get("original_row")
if isinstance(original_row, dict):
for key in ("downloads", "download", "下载量"):
if key in original_row:
downloads = _parse_downloads_value(original_row.get(key))
if downloads is not None:
return downloads
for key in ("downloads", "download", "下载量"):
if key in payload:
downloads = _parse_downloads_value(payload.get(key))
if downloads is not None:
return downloads
return None
def _extract_downloads_from_row(row: Dict[str, Any]) -> Optional[int]:
if not isinstance(row, dict):
return None
return _parse_downloads_value(
row.get("downloads")
or row.get("download")
or row.get("下载量")
)
def _normalize_download_errors(download_errors: Any) -> Dict[str, Dict[str, Any]]:
return download_errors if isinstance(download_errors, dict) else {}
def _resolve_downloads_for_package(
package_name: str,
task_payload: Any = None,
*,
downloads: Any = None,
) -> Optional[int]:
del package_name
normalized_downloads = _parse_downloads_value(downloads)
if normalized_downloads is not None:
return normalized_downloads
return _extract_downloads_from_task_payload(task_payload)
def _find_download_bucket_spec(downloads: int) -> Tuple[str, str, int, Optional[int]]:
normalized_downloads = max(0, int(downloads or 0))
for bucket_key, label, lower_bound, upper_bound in NON_RETRYABLE_DOWNLOAD_BUCKETS:
if upper_bound is None and normalized_downloads >= lower_bound:
return bucket_key, label, lower_bound, upper_bound
if upper_bound is not None and lower_bound <= normalized_downloads < upper_bound:
return bucket_key, label, lower_bound, upper_bound
return NON_RETRYABLE_DOWNLOAD_BUCKETS[0]
def _build_non_retryable_download_distribution(
rows: Iterable[sqlite3.Row],
) -> Dict[str, Any]:
rows = list(rows)
if not rows:
return {
"distribution": [],
"matched_downloads_apps": 0,
"missing_downloads_apps": 0,
}
distribution = [
{
"bucket_key": bucket_key,
"label": label,
"min_downloads": lower_bound,
"max_downloads": upper_bound,
"app_count": 0,
}
for bucket_key, label, lower_bound, upper_bound in NON_RETRYABLE_DOWNLOAD_BUCKETS
]
bucket_index = {item["bucket_key"]: item for item in distribution}
unknown_bucket = {
"bucket_key": "unknown",
"label": "未知",
"min_downloads": None,
"max_downloads": None,
"app_count": 0,
}
matched_downloads_apps = 0
for row in rows:
package_name = str(row["package_name"] or "").strip()
downloads = _resolve_downloads_for_package(
package_name,
row["task_payload_json"],
downloads=row["downloads"],
)
if downloads is None:
unknown_bucket["app_count"] += 1
continue
matched_downloads_apps += 1
bucket_key, _, _, _ = _find_download_bucket_spec(downloads)
bucket_index[bucket_key]["app_count"] += 1
total_count = len(rows)
distribution.append(unknown_bucket)
return {
"distribution": [
{
**item,
"app_count": int(item["app_count"] or 0),
"share_percent": round((int(item["app_count"] or 0) / total_count) * 100, 2)
if total_count > 0
else 0.0,
}
for item in distribution
],
"matched_downloads_apps": matched_downloads_apps,
"missing_downloads_apps": int(unknown_bucket["app_count"] or 0),
}
def _merge_detail_text(*parts: Any) -> str:
normalized_parts = []
seen = set()
for part in parts:
text = str(part or "").strip()
if not text or text in seen:
continue
seen.add(text)
normalized_parts.append(text)
return " | ".join(normalized_parts)
def _parse_traffic_size(size_str: str) -> int:
size_str = str(size_str or "").strip()
if not size_str:
return 0
number = []
unit = []
for char in size_str:
if char.isdigit() or char == ".":
number.append(char)
elif not char.isspace():
unit.append(char)
if not number:
return 0
try:
value = float("".join(number))
except ValueError:
return 0
normalized_unit = "".join(unit).upper()
multipliers = {
"": 1,
"B": 1,
"K": 1024,
"KB": 1024,
"M": 1024 ** 2,
"MB": 1024 ** 2,
"G": 1024 ** 3,
"GB": 1024 ** 3,
"T": 1024 ** 4,
"TB": 1024 ** 4,
}
return int(value * multipliers.get(normalized_unit, 1))
def _is_model_traffic_line(fields: List[str]) -> bool:
if len(fields) < 7 or str(fields[2] or "").strip() != "model_data:":
return False
protocol = str(fields[5] or "").strip().upper()
if protocol not in ("TCP", "UDP"):
return False
flow_info = str(fields[3] or "").strip()
if not flow_info:
return False
parts = [part.strip() for part in flow_info.split("-") if part.strip()]
if len(parts) < 2:
return False
dst_part = parts[1]
dst_ip = dst_part.rsplit(":", 1)[0].strip() if ":" in dst_part else ""
if not dst_ip:
return False
for prefix in ("10.", "192.", "198.18", "198.19"):
if dst_ip.startswith(prefix):
return False
return True
def _determine_domain_or_model(fields: List[str]) -> str:
domain = fields[2].strip() if len(fields) > 2 else ""
flow_info = fields[3].strip() if len(fields) > 3 else ""
if domain != "model_data:":
return domain
return f"model_data:{flow_info}"
def _extract_second_level_domain(domain: str) -> str:
cleaned = str(domain or "").strip().lower()
if not cleaned:
return ""
cleaned = cleaned.replace("http://", "").replace("https://", "").strip().strip("\"'")
parts = [part for part in cleaned.split(".") if part]
if len(parts) <= 2:
return cleaned
return "*." + ".".join(parts[-2:])
def _pick_match_state(states: Iterable[str]) -> str:
normalized_states = {str(item or "").strip() for item in states if str(item or "").strip()}
if "self" in normalized_states:
return "self"
if "server" in normalized_states:
return "server"
if "unmatched" in normalized_states:
return "unmatched"
if "ip" in normalized_states:
return "ip"
return ""
def _resolve_data_root(path: str) -> str:
raw_path = str(path or "").strip()
if not raw_path or os.path.isdir(raw_path):
return raw_path
normalized = raw_path.replace("\\", "/").lstrip("/")
parts = [part for part in normalized.split("/") if part]
if len(parts) >= 2 and "." in parts[0]:
# macOS: /Volumes/<share>/...
candidate = os.path.join("/Volumes", parts[1], *parts[2:])
if os.path.isdir(candidate):
return candidate
# Linux: try common SMB mount root
candidate = os.path.join("/srv/samba", parts[1], *parts[2:])
if os.path.isdir(candidate):
return candidate
if len(parts) >= 3 and parts[0].endswith(":") and parts[1].lower() == "share":
candidate = os.path.join("/Volumes", "dpi-sync", *parts[2:])
if os.path.isdir(candidate):
return candidate
return raw_path
def _aggregate_second_level_domain_rows(domain_rows: Iterable[sqlite3.Row]) -> List[Dict[str, Any]]:
aggregated: Dict[str, Dict[str, Any]] = {}
for row in domain_rows:
domain = str(row["domain"] or "")
if str(row["domain_type"] or "") == "model_data" or domain.startswith("model_data:"):
continue
second_level_domain = _extract_second_level_domain(domain)
if not second_level_domain:
continue
bucket = aggregated.setdefault(
second_level_domain,
{
"domain": second_level_domain,
"traffic_bytes": 0,
"flow_count": 0,
"organization": "",
"matched_pattern": "",
"states": set(),
},
)
bucket["traffic_bytes"] += int(row["traffic_bytes"] or 0)
bucket["flow_count"] += int(row["flow_count"] or 0)
if row["organization"] and not bucket["organization"]:
bucket["organization"] = row["organization"]
if row["matched_pattern"] and not bucket["matched_pattern"]:
bucket["matched_pattern"] = row["matched_pattern"]
if row["match_state"]:
bucket["states"].add(row["match_state"])
aggregated_rows = []
for bucket in aggregated.values():
aggregated_rows.append(
{
"domain": bucket["domain"],
"traffic_bytes": int(bucket["traffic_bytes"]),
"flow_count": int(bucket["flow_count"]),
"organization": bucket["organization"],
"matched_pattern": bucket["matched_pattern"],
"match_state": _pick_match_state(bucket["states"]),
}
)
return sorted(aggregated_rows, key=lambda item: (-int(item["traffic_bytes"]), item["domain"]))
def _classify_restriction_status(
latest_status: str,
self_ratio: float,
latest_failure_type: str,
*,
num_nodes: int = 0,
download_errors: Any = None,
) -> Tuple[str, str]:
normalized_status = str(latest_status or "").strip().lower()
if normalized_status == "success":
return "success", "not_applicable"
try:
normalized_self_ratio = float(self_ratio or 0.0)
except (TypeError, ValueError):
normalized_self_ratio = 0.0
normalized_num_nodes = _parse_int_value(num_nodes)
if normalized_self_ratio > 0 or normalized_num_nodes >= LIGHT_RESTRICTED_MIN_NUM_NODES:
return "light_restricted", "not_applicable"
retryability = "retryable"
error_type = str(latest_failure_type or "").strip()
if error_type:
category_name, _, code_text = error_type.partition("/")
try:
error_info = ErrorInfo(ErrorCategory[category_name], int(code_text), error_type)
retryability = "non_retryable" if error_info.is_no_retry(download_errors=download_errors) else "retryable"
except (KeyError, TypeError, ValueError):
retryability = "retryable"
return "severe_restricted", retryability
def _make_task_key(app_name: str, package_name: str) -> str:
return f"{str(app_name or '').strip()}_{str(package_name or '').strip()}"
def _normalize_country_codes(raw_country_code: Any) -> List[str]:
if isinstance(raw_country_code, (list, tuple)):
raw_items = raw_country_code
else:
raw_items = str(raw_country_code or "").split(",")
normalized: List[str] = []
seen = set()
for item in raw_items:
candidate = str(item or "").strip().upper()
if not candidate or candidate in seen:
continue
seen.add(candidate)
normalized.append(candidate)
return normalized
def _normalize_device_type(raw_device_type: Any) -> str:
candidate = str(raw_device_type or "").strip().lower()
if candidate in {"emulator", "physical"}:
return candidate
return ""
def _normalize_task_device_type(raw_device_type: Any) -> str:
candidate = str(raw_device_type or "").strip().lower()
if candidate in {"emulator", "physical", "any"}:
return candidate
return "emulator"
def _build_task_payload_with_device_type(
task_payload: Any,
*,
app_name: str,
package_name: str,
country_code: str,
device_type: str,
) -> Dict[str, Any]:
normalized_device_type = _normalize_task_device_type(device_type)
if isinstance(task_payload, dict) and task_payload:
payload = dict(task_payload)
else:
payload = build_task_payload_from_row(
{
"app_name": app_name,
"package_name": package_name,
"country_code": country_code,
"device_type": normalized_device_type,
}
)
payload["task_key"] = str(payload.get("task_key") or _make_task_key(app_name, package_name)).strip()
payload["app_name"] = app_name
payload["package_name"] = package_name
payload["country_code"] = country_code
payload["country_codes"] = _normalize_country_codes(payload.get("country_codes") or country_code)
payload["device_type"] = normalized_device_type
payload["available_sources"] = payload.get("available_sources") or ["google_play", "local"]
original_row = payload.get("original_row")
normalized_original_row = dict(original_row) if isinstance(original_row, dict) else {}
normalized_original_row["app_name"] = app_name
normalized_original_row["package_name"] = package_name
normalized_original_row["country_code"] = country_code
normalized_original_row["device_type"] = normalized_device_type
payload["original_row"] = normalized_original_row
return payload
def build_task_payload_from_row(row: Dict[str, Any]) -> Dict[str, Any]:
app_name = str(row.get("app_name") or row.get("应用名称") or "").strip()
package_name = str(row.get("package_name") or row.get("包名") or "").strip()
country_code = str(row.get("country_code") or row.get("国家") or "").strip()
app_magic_label = _normalize_app_magic_label(row.get("app_magic_label") or row.get("应用魔法标签") or "")
last_updated = _normalize_last_updated(row.get("last_updated") or row.get("最后更新") or row.get("更新时间") or "")
device_type = _normalize_task_device_type(row.get("device_type") or row.get("设备类型") or "")
downloads = _extract_downloads_from_row(row)
available_sources = row.get("available_sources")
if isinstance(available_sources, str):
normalized_sources = [item.strip() for item in available_sources.split(",") if item.strip()]
elif isinstance(available_sources, (list, tuple)):
normalized_sources = [str(item).strip() for item in available_sources if str(item).strip()]
else:
normalized_sources = []
if not normalized_sources:
normalized_sources = ["google_play", "local"]
normalized_row = {str(key or "").strip(): str(value or "").strip() for key, value in (row or {}).items()}
normalized_row["device_type"] = device_type
payload = {
"task_key": _make_task_key(app_name, package_name),
"app_name": app_name,
"package_name": package_name,
"country_code": country_code,
"country_codes": _normalize_country_codes(country_code),
"app_magic_label": app_magic_label,
"last_updated": last_updated,
"device_type": device_type,
"available_sources": normalized_sources,
"original_row": normalized_row,
}
if downloads is not None:
payload["downloads"] = downloads
return payload
def _derive_collection_status(
restriction_status: str,
retryability: str,
latest_failure_type: str,
*,
allow_legacy_download_error_1_pending: bool = False,
) -> Tuple[str, str]:
normalized_restriction = str(restriction_status or "").strip()
normalized_retryability = str(retryability or "").strip()
normalized_failure_type = str(latest_failure_type or "").strip().upper()
if allow_legacy_download_error_1_pending and normalized_failure_type == "DOWNLOAD_ERROR/1":
return "pending", "legacy_download_error_1"
if normalized_restriction in {"success", "light_restricted"}:
return "qualified", normalized_restriction or "qualified"
if normalized_restriction == "severe_restricted" and normalized_retryability == "non_retryable":
return "failed_terminal", normalized_failure_type or "severe_non_retryable"
return "pending", normalized_restriction or "pending"
def _should_preserve_non_retryable_reason(
existing_row: Dict[str, Any],
new_restriction_status: str,
new_retryability: str,
) -> bool:
previous_restriction_status = str(existing_row.get("restriction_status") or "").strip()
previous_retryability = str(existing_row.get("retryability") or "").strip()
previous_failure_type = str(existing_row.get("latest_failure_type") or "").strip()
normalized_new_restriction = str(new_restriction_status or "").strip()
normalized_new_retryability = str(new_retryability or "").strip()
return (
previous_restriction_status == "severe_restricted"
and previous_retryability == "non_retryable"
and normalized_new_restriction == "severe_restricted"
and normalized_new_retryability == "retryable"
and bool(previous_failure_type)
)
def _should_auto_reroute_to_physical(
current_device_type: str,
latest_failure_type: str,
restriction_status: str,
retryability: str,
) -> bool:
normalized_device_type = _normalize_task_device_type(current_device_type)
normalized_failure_type = str(latest_failure_type or "").strip().upper()
normalized_restriction = str(restriction_status or "").strip()
normalized_retryability = str(retryability or "").strip()
return (
normalized_device_type != "physical"
and normalized_failure_type in AUTO_PENDING_PHYSICAL_FAILURE_TYPES
and normalized_restriction == "severe_restricted"
and normalized_retryability == "non_retryable"
)
def _humanize_error_type(error_type: str) -> str:
normalized = str(error_type or "").strip()
if not normalized:
return "未标记错误"
category_name, _, code_text = normalized.partition("/")
try:
code = int(code_text)
except (TypeError, ValueError):
return normalized
try:
if category_name == "INFRA_ERROR":
return INFRA_ERROR_DESC.get(InfraError(code), normalized)
if category_name == "APP_ERROR":
return APP_ERROR_DESC.get(AppError(code), normalized)
if category_name == "BUSINESS_ERROR":
return BUSINESS_ERROR_DESC.get(BusinessError(code), normalized)
if category_name == "DOWNLOAD_ERROR":
return DOWNLOAD_ERROR_DESC.get(DownloadError(code), normalized)
except ValueError:
return normalized
return normalized
class DpiRuleMatcher:
def __init__(self, app_list_path: str, url_lib_path: str):
self.app_list_path = app_list_path
self.url_lib_path = url_lib_path
self._loaded = False
self._lock = threading.Lock()
self.tp_packages: Dict[str, set] = {}
self.tp_pkg_string: Dict[str, str] = {}
self.pkg_to_tp_mark: Dict[str, str] = {}
self._trie: Dict[str, Any] = {}
def ensure_loaded(self):
if self._loaded:
return
with self._lock:
if self._loaded:
return
self._load_app_list()
self._load_url_rules()
self._loaded = True
def _load_app_list(self):
if not os.path.exists(self.app_list_path):
logger.warning("Analytics app list is missing: %s", self.app_list_path)
return
for row in _read_dict_rows(self.app_list_path):
tp_mark = row.get("app_name", "").strip()
package_string = row.get("package_name", "").strip()
if not tp_mark:
continue
self.tp_pkg_string[tp_mark] = package_string
packages = {item.strip() for item in package_string.split(",") if item.strip()}
self.tp_packages.setdefault(tp_mark, set()).update(packages)
for package_name in packages:
self.pkg_to_tp_mark[package_name] = tp_mark
def _load_url_rules(self):
if not os.path.exists(self.url_lib_path):
logger.warning("Analytics url lib is missing: %s", self.url_lib_path)
return
rule_count = 0
for row in _read_dict_rows(self.url_lib_path):
suffix = row.get("url", "").strip().lower()
tp_mark = row.get("app", "").strip()
if not suffix or not tp_mark:
continue
self._insert_suffix(suffix, tp_mark)
rule_count += 1
logger.info("[Analytics] loaded %s dpi url rules", rule_count)
def _insert_suffix(self, suffix: str, tp_mark: str):
node = self._trie
for part in reversed([item for item in suffix.split(".") if item]):
node = node.setdefault(part, {})
node["$"] = (suffix, tp_mark)
def match_domain(self, domain: str) -> Tuple[Optional[str], Optional[str]]:
self.ensure_loaded()
cleaned = str(domain or "").strip().lower()
if not cleaned or cleaned.startswith("model_data:"):
return None, None
cleaned = cleaned.replace("http://", "").replace("https://", "").strip().strip("\"'")
node = self._trie
best = None
for part in reversed([item for item in cleaned.split(".") if item]):
node = node.get(part)
if node is None:
break
if "$" in node:
best = node["$"]
if not best:
return None, None
return best[1], best[0]
class AnalyticsRepository:
def __init__(self, db_path: str = MONITORING_DB_PATH):
self.db_path = db_path
self._write_lock = threading.Lock()
self._initialize_schema()
@contextmanager
def _connect(self):
connection = sqlite3.connect(self.db_path, check_same_thread=False)
connection.row_factory = sqlite3.Row
try:
yield connection
connection.commit()
except Exception:
connection.rollback()
raise
finally:
connection.close()
def _table_exists(self, connection: sqlite3.Connection, table_name: str) -> bool:
row = connection.execute(
"""
SELECT 1
FROM sqlite_master
WHERE type = 'table' AND name = ?
""",
(table_name,),
).fetchone()
return bool(row)
def _table_columns(self, connection: sqlite3.Connection, table_name: str) -> set:
if not self._table_exists(connection, table_name):
return set()
return {row["name"] for row in connection.execute(f"PRAGMA table_info({table_name})").fetchall()}
def _initialize_schema(self):
"""
初始化数据库 schema。
analytics 直接使用新表 app_catalog / collection_task不创建兼容视图。
"""
os.makedirs(os.path.dirname(self.db_path) or ".", exist_ok=True)
with self._connect() as connection:
connection.execute("PRAGMA journal_mode=WAL")
connection.executescript(
"""
CREATE TABLE IF NOT EXISTS collection_task (
package_name TEXT NOT NULL,
batch_tag TEXT NOT NULL,
run_kind TEXT NOT NULL DEFAULT 'ranking'
CHECK(run_kind IN ('ranking', 'block', 'model', 'manual')),
attempt INTEGER NOT NULL DEFAULT 1,
task_key TEXT,
app_name TEXT,
app_magic_label TEXT,
is_new_app INTEGER DEFAULT 0,
task_status TEXT DEFAULT 'pending',
execution_status TEXT,
worker_id TEXT,
created_at TEXT DEFAULT (datetime('now', 'localtime')),
started_at TEXT,
completed_at TEXT,
duration_seconds REAL,
download_duration_seconds REAL,
execution_duration_seconds REAL,
analysis_duration_seconds REAL,
error_category TEXT,
error_code INTEGER,
error_reason TEXT,
error_details TEXT,
crashed_source TEXT,
droidbot_steps INTEGER DEFAULT 0,
gui_agent_steps INTEGER DEFAULT 0,
total_steps INTEGER DEFAULT 0,
num_nodes INTEGER DEFAULT 0,
num_reached_activities INTEGER DEFAULT 0,
app_num_total_activities INTEGER DEFAULT 0,
total_traffic_bytes INTEGER DEFAULT 0,
self_traffic_bytes INTEGER DEFAULT 0,
server_traffic_bytes INTEGER DEFAULT 0,
unrecognized_traffic_bytes INTEGER DEFAULT 0,
model_flow_count INTEGER DEFAULT 0,
model_traffic_bytes INTEGER DEFAULT 0,
login_count INTEGER DEFAULT 0,
register_count INTEGER DEFAULT 0,
stuck_reason_code INTEGER,
guiagent_message TEXT,
scenario_triggered INTEGER DEFAULT 0,
download_source TEXT,
is_retry INTEGER DEFAULT 0,
exit_code INTEGER,
trace_json TEXT,
traffic_file_paths TEXT,
PRIMARY KEY (package_name, batch_tag, run_kind, attempt)
);
CREATE INDEX IF NOT EXISTS idx_collection_task_status
ON collection_task(task_status, run_kind, batch_tag);
CREATE INDEX IF NOT EXISTS idx_collection_task_execution
ON collection_task(execution_status, run_kind);
CREATE INDEX IF NOT EXISTS idx_collection_task_pkg
ON collection_task(package_name, run_kind);
CREATE INDEX IF NOT EXISTS idx_collection_task_time
ON collection_task(created_at);
CREATE INDEX IF NOT EXISTS idx_collection_task_worker
ON collection_task(worker_id, completed_at);
CREATE INDEX IF NOT EXISTS idx_collection_task_magic
ON collection_task(app_magic_label);
CREATE INDEX IF NOT EXISTS idx_collection_task_batch
ON collection_task(batch_tag, run_kind);
CREATE TABLE IF NOT EXISTS app_catalog (
package_name TEXT PRIMARY KEY,
app_name TEXT,
batch_tags TEXT,
app_magic_label TEXT,
last_updated TEXT,
country_code TEXT,
device_type TEXT,
task_payload_json TEXT,
last_update_interval_days INTEGER DEFAULT 0,
task_queue TEXT DEFAULT 'default',
task_priority INTEGER DEFAULT 50,
source_order INTEGER,
is_active INTEGER DEFAULT 1,
is_blocked INTEGER DEFAULT 0,
category TEXT,
downloads INTEGER,
created_at TEXT DEFAULT (datetime('now', 'localtime')),
updated_at TEXT DEFAULT (datetime('now', 'localtime'))
);
CREATE INDEX IF NOT EXISTS idx_app_catalog_priority
ON app_catalog(task_priority DESC);
CREATE INDEX IF NOT EXISTS idx_app_catalog_magic
ON app_catalog(app_magic_label);
CREATE INDEX IF NOT EXISTS idx_app_catalog_active
ON app_catalog(is_active, source_order);
CREATE TABLE IF NOT EXISTS apk_registry (
package_name TEXT PRIMARY KEY,
download_date TEXT NOT NULL,
download_time REAL NOT NULL,
local_dir TEXT,
smb_dir TEXT,
version_name TEXT,
source TEXT,
apk_files_json TEXT,
updated_at REAL NOT NULL
);
CREATE TABLE IF NOT EXISTS worker_activity_summary (
worker_id TEXT NOT NULL,
stat_date TEXT NOT NULL,
stat_hour INTEGER NOT NULL,
idle_duration_seconds INTEGER DEFAULT 0,
busy_duration_seconds INTEGER DEFAULT 0,
offline_duration_seconds INTEGER DEFAULT 0,
task_count INTEGER DEFAULT 0,
success_count INTEGER DEFAULT 0,
failed_count INTEGER DEFAULT 0,
PRIMARY KEY (worker_id, stat_date, stat_hour)
);
CREATE INDEX IF NOT EXISTS idx_worker_activity_date
ON worker_activity_summary(stat_date DESC);
CREATE TABLE IF NOT EXISTS analytics_job (
id INTEGER PRIMARY KEY AUTOINCREMENT,
job_type TEXT NOT NULL,
status TEXT NOT NULL,
package_name TEXT,
task_key TEXT,
worker_id TEXT,
trigger_source TEXT,
report_time REAL,
attempt_count INTEGER DEFAULT 0,
error_message TEXT,
payload_json TEXT,
scheduled_after REAL DEFAULT 0,
created_at REAL NOT NULL,
updated_at REAL NOT NULL,
finished_at REAL
);
CREATE INDEX IF NOT EXISTS idx_analytics_job_status_created_at
ON analytics_job(status, created_at);
CREATE INDEX IF NOT EXISTS idx_analytics_job_package_name
ON analytics_job(package_name, created_at);
CREATE TABLE IF NOT EXISTS analytics_source_file (
file_path TEXT PRIMARY KEY,
file_type TEXT NOT NULL,
package_name TEXT,
worker_id TEXT,
size_bytes INTEGER,
mtime REAL,
checksum TEXT,
last_seen_at REAL,
last_processed_at REAL
);
CREATE INDEX IF NOT EXISTS idx_analytics_source_file_package_name
ON analytics_source_file(package_name, file_type);
CREATE TABLE IF NOT EXISTS app_domain_traffic (
package_name TEXT NOT NULL,
domain TEXT NOT NULL,
domain_type TEXT,
traffic_bytes INTEGER DEFAULT 0,
flow_count INTEGER DEFAULT 0,
traffic_ratio REAL DEFAULT 0,
domain_traffic_ratio REAL DEFAULT 0,
organization TEXT,
matched_tp_mark TEXT,
matched_pattern TEXT,
match_state TEXT,
updated_at REAL NOT NULL,
PRIMARY KEY(package_name, domain)
);
CREATE INDEX IF NOT EXISTS idx_app_domain_traffic_package_name
ON app_domain_traffic(package_name, traffic_bytes DESC);
CREATE TABLE IF NOT EXISTS app_traffic_component (
package_name TEXT NOT NULL,
component_name TEXT NOT NULL,
component_package_names TEXT,
is_self INTEGER DEFAULT 0,
traffic_bytes INTEGER DEFAULT 0,
share_percent REAL DEFAULT 0,
updated_at REAL NOT NULL,
PRIMARY KEY(package_name, component_name)
);
CREATE INDEX IF NOT EXISTS idx_app_traffic_component_package_name
ON app_traffic_component(package_name, traffic_bytes DESC);
"""
)
self._ensure_columns(
connection,
"app_catalog",
{
"last_updated": "TEXT",
"country_code": "TEXT",
"device_type": "TEXT",
"task_payload_json": "TEXT",
"last_update_interval_days": "INTEGER DEFAULT 0",
},
)
def _ensure_columns(self, connection: sqlite3.Connection, table_name: str, columns: Dict[str, str]) -> None:
existing_columns = self._table_columns(connection, table_name)
for column_name, column_type in columns.items():
if column_name in existing_columns:
continue
connection.execute(f"ALTER TABLE {table_name} ADD COLUMN {column_name} {column_type}")
def create_job(
self,
*,
job_type: str,
status: str,
package_name: Optional[str] = None,
task_key: Optional[str] = None,
worker_id: Optional[str] = None,
trigger_source: Optional[str] = None,
report_time: Optional[float] = None,
payload: Optional[Dict[str, Any]] = None,
scheduled_after: Optional[float] = None,
) -> Dict[str, Any]:
now = time.time()
payload_json = _safe_json_dumps(payload or {})
with self._write_lock, self._connect() as connection:
cursor = connection.execute(
"""
INSERT INTO analytics_job (
job_type, status, package_name, task_key, worker_id, trigger_source,
report_time, attempt_count, error_message, payload_json, scheduled_after,
created_at, updated_at, finished_at
) VALUES (?, ?, ?, ?, ?, ?, ?, 0, '', ?, ?, ?, ?, NULL)
""",
(
job_type,
status,
package_name,
task_key,
worker_id,
trigger_source or "",
report_time,
payload_json,
float(scheduled_after or 0),
now,
now,
),
)
job_id = int(cursor.lastrowid)
return self.get_job(job_id)
def get_job(self, job_id: int) -> Dict[str, Any]:
with self._connect() as connection:
row = connection.execute(
"SELECT * FROM analytics_job WHERE id = ?",
(job_id,),
).fetchone()
return self._serialize_job(row) if row else {}
def _serialize_job(self, row: sqlite3.Row) -> Dict[str, Any]:
return {
"id": int(row["id"]),
"job_type": row["job_type"],
"status": row["status"],
"package_name": row["package_name"] or "",
"task_key": row["task_key"] or "",
"worker_id": row["worker_id"] or "",
"trigger_source": row["trigger_source"] or "",
"report_time": _format_ts(row["report_time"]),
"attempt_count": int(row["attempt_count"] or 0),
"error_message": row["error_message"] or "",
"payload": _safe_json_loads(row["payload_json"], {}) or {},
"created_at": _format_ts(row["created_at"]),
"updated_at": _format_ts(row["updated_at"]),
"finished_at": _format_ts(row["finished_at"]),
}
def list_jobs(
self,
*,
job_type: Optional[str] = None,
status: Optional[str] = None,
package_name: Optional[str] = None,
limit: int = 50,
) -> List[Dict[str, Any]]:
clauses = []
params: List[Any] = []
if job_type:
clauses.append("job_type = ?")
params.append(job_type)
if status:
clauses.append("status = ?")
params.append(status)
if package_name:
clauses.append("package_name = ?")
params.append(package_name)
where_clause = f"WHERE {' AND '.join(clauses)}" if clauses else ""
with self._connect() as connection:
rows = connection.execute(
f"""
SELECT *
FROM analytics_job
{where_clause}
ORDER BY created_at DESC, id DESC
LIMIT ?
""",
[*params, max(1, min(limit, 500))],
).fetchall()
return [self._serialize_job(row) for row in rows]
def claim_next_job(self, artifact_wait_seconds: int) -> Optional[Dict[str, Any]]:
now = time.time()
with self._write_lock, self._connect() as connection:
row = connection.execute(
"""
SELECT *
FROM analytics_job
WHERE (status = 'queued'
AND COALESCE(scheduled_after, 0) <= ?)
OR (
status = 'waiting_artifacts'
AND COALESCE(report_time, created_at) <= ?
)
ORDER BY CASE status WHEN 'queued' THEN 0 ELSE 1 END, created_at ASC, id ASC
LIMIT 1
""",
(now, now - max(0, artifact_wait_seconds)),
).fetchone()
if not row:
return None
connection.execute(
"""
UPDATE analytics_job
SET status = 'running',
attempt_count = COALESCE(attempt_count, 0) + 1,
updated_at = ?
WHERE id = ?
""",
(now, row["id"]),
)
return self.get_job(int(row["id"]))
def mark_job_waiting(self, job_id: int, error_message: str = ""):
now = time.time()
with self._write_lock, self._connect() as connection:
connection.execute(
"""
UPDATE analytics_job
SET status = 'waiting_artifacts',
error_message = ?,
updated_at = ?,
finished_at = NULL
WHERE id = ?
""",
(error_message, now, job_id),
)
def mark_job_ready(self, *, task_key: Optional[str] = None, package_name: Optional[str] = None, worker_id: Optional[str] = None) -> int:
clauses = ["status = 'waiting_artifacts'"]
params: List[Any] = []
if task_key:
clauses.append("task_key = ?")
params.append(task_key)
if package_name:
clauses.append("package_name = ?")
params.append(package_name)
if worker_id:
clauses.append("worker_id = ?")
params.append(worker_id)
if len(clauses) == 1:
return 0
now = time.time()
with self._write_lock, self._connect() as connection:
cursor = connection.execute(
f"""
UPDATE analytics_job
SET status = 'queued',
scheduled_after = ?,
updated_at = ?,
error_message = ''
WHERE {' AND '.join(clauses)}
""",
[now + ANALYTICS_ARTIFACT_FILE_WAIT_SECONDS, now, *params],
)
return int(cursor.rowcount or 0)
def finish_job(self, job_id: int, *, status: str, error_message: str = ""):
now = time.time()
with self._write_lock, self._connect() as connection:
connection.execute(
"""
UPDATE analytics_job
SET status = ?,
error_message = ?,
updated_at = ?,
finished_at = ?
WHERE id = ?
""",
(status, error_message, now, now, job_id),
)
def get_latest_task_execution(self, package_name: str) -> Optional[Dict[str, Any]]:
with self._connect() as connection:
if not self._table_exists(connection, "task_execution"):
return None
row = connection.execute(
"""
SELECT
execution_id,
task_key,
app_name,
package_name,
last_updated,
collection_task_type,
worker_id,
status,
error_type,
error_message,
result_detail,
download_duration_seconds,
collect_duration_seconds,
total_duration_seconds,
failed_total_duration_seconds,
droidbot_steps,
guiagent_steps,
total_steps,
num_nodes,
num_reached_activities,
app_num_total_activities,
trace_json,
COALESCE(task_ended_at, task_started_at, last_updated_at) AS latest_time
FROM task_execution
WHERE package_name = ?
ORDER BY COALESCE(task_ended_at, task_started_at, last_updated_at) DESC, execution_id DESC
LIMIT 1
""",
(package_name,),
).fetchone()
if not row:
return None
trace = _safe_json_loads(row["trace_json"], {})
download_errors = trace.get("download_errors") if isinstance(trace, dict) else {}
return {
"execution_id": row["execution_id"],
"task_key": row["task_key"] or "",
"app_name": row["app_name"] or "",
"package_name": row["package_name"] or package_name,
"last_updated": row["last_updated"] or "",
"collection_task_type": row["collection_task_type"] or "",
"worker_id": row["worker_id"] or "",
"status": row["status"] or "",
"error_type": row["error_type"] or "",
"error_message": row["error_message"] or "",
"result_detail": row["result_detail"] or row["error_message"] or "",
"download_duration_seconds": float(row["download_duration_seconds"] or 0.0),
"collect_duration_seconds": float(row["collect_duration_seconds"] or 0.0),
"total_duration_seconds": float(row["total_duration_seconds"] or row["failed_total_duration_seconds"] or 0.0),
"droidbot_steps": int(row["droidbot_steps"] or 0),
"guiagent_steps": int(row["guiagent_steps"] or 0),
"total_steps": int(row["total_steps"] or 0),
"num_nodes": int(row["num_nodes"] or 0),
"num_reached_activities": int(row["num_reached_activities"] or 0),
"app_num_total_activities": int(row["app_num_total_activities"] or 0),
"latest_time": float(row["latest_time"] or 0.0),
"download_errors": download_errors if isinstance(download_errors, dict) else {},
"trace": trace if isinstance(trace, dict) else {},
}
def _load_latest_download_errors_by_packages(
self,
connection: sqlite3.Connection,
package_names: List[str],
) -> Dict[str, Dict[str, Any]]:
normalized = [str(item or "").strip() for item in package_names if str(item or "").strip()]
if not normalized or not self._table_exists(connection, "task_execution"):
return {}
placeholders = ",".join("?" for _ in normalized)
rows = connection.execute(
f"""
SELECT package_name, trace_json
FROM (
SELECT
package_name,
trace_json,
ROW_NUMBER() OVER (
PARTITION BY package_name
ORDER BY COALESCE(task_ended_at, task_started_at, last_updated_at) DESC, execution_id DESC
) AS row_num
FROM task_execution
WHERE package_name IN ({placeholders})
)
WHERE row_num = 1
""",
normalized,
).fetchall()
latest_download_errors: Dict[str, Dict[str, Any]] = {}
for row in rows:
package_name = str(row["package_name"] or "").strip()
if not package_name:
continue
trace = _safe_json_loads(row["trace_json"], {}) or {}
if not isinstance(trace, dict):
continue
latest_download_errors[package_name] = _normalize_download_errors(trace.get("download_errors"))
return latest_download_errors
def list_all_task_packages(self) -> List[str]:
with self._connect() as connection:
if not self._table_exists(connection, "task_execution"):
return []
rows = connection.execute(
"""
SELECT DISTINCT package_name
FROM task_execution
WHERE package_name IS NOT NULL AND package_name != ''
"""
).fetchall()
return sorted({row["package_name"] for row in rows if row["package_name"]})
def _parse_db_datetime(self, value: Any) -> float:
text = str(value or "").strip()
if not text:
return 0.0
for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S"):
try:
return datetime.strptime(text[:19], fmt).replace(tzinfo=SHANGHAI_TZ).timestamp()
except ValueError:
continue
try:
return float(text)
except (TypeError, ValueError):
return 0.0
def _latest_task_row(self, connection: sqlite3.Connection, package_name: str) -> Optional[sqlite3.Row]:
return connection.execute(
"""
SELECT *
FROM collection_task
WHERE package_name = ?
ORDER BY
datetime(COALESCE(NULLIF(completed_at, ''), NULLIF(started_at, ''), created_at)) DESC,
attempt DESC,
batch_tag DESC,
run_kind DESC
LIMIT 1
""",
(package_name,),
).fetchone()
def _failure_type_from_task(self, task_row: Optional[sqlite3.Row]) -> str:
if not task_row:
return ""
category = str(task_row["error_category"] or "").strip()
code = task_row["error_code"]
if not category:
return ""
if code is None or str(code).strip() == "":
return category
return f"{category}/{int(code)}"
def _download_errors_from_task(self, task_row: Optional[sqlite3.Row]) -> Dict[str, Any]:
if not task_row:
return {}
trace = _safe_json_loads(task_row["trace_json"], {}) or {}
if not isinstance(trace, dict):
return {}
return _normalize_download_errors(trace.get("download_errors"))
def _ratio_from_task(self, task_row: Optional[sqlite3.Row], numerator_key: str) -> float:
if not task_row:
return 0.0
total = int(task_row["total_traffic_bytes"] or 0)
if total <= 0:
return 0.0
numerator = int(task_row[numerator_key] or 0)
return round((numerator / total) * 100, 2)
def _recognition_ratio_from_task(self, task_row: Optional[sqlite3.Row]) -> float:
if not task_row:
return 0.0
total = int(task_row["total_traffic_bytes"] or 0)
if total <= 0:
return 0.0
recognized = int(task_row["self_traffic_bytes"] or 0) + int(task_row["server_traffic_bytes"] or 0)
return round((recognized / total) * 100, 2)
def _state_from_task(self, task_row: Optional[sqlite3.Row]) -> Dict[str, str]:
if not task_row:
return {
"collection_status": "pending",
"collection_status_reason": "no_collection_task",
"restriction_status": "",
"retryability": "not_applicable",
}
task_status = str(task_row["task_status"] or "").strip().lower()
execution_status = str(task_row["execution_status"] or "").strip().lower()
if task_status in {"pending", "running"} or not execution_status:
return {
"collection_status": "pending",
"collection_status_reason": str(task_row["error_reason"] or "pending"),
"restriction_status": "",
"retryability": "not_applicable",
}
failure_type = self._failure_type_from_task(task_row)
restriction_status, retryability = _classify_restriction_status(
execution_status,
self._ratio_from_task(task_row, "self_traffic_bytes"),
failure_type,
num_nodes=int(task_row["num_nodes"] or 0),
download_errors=self._download_errors_from_task(task_row),
)
if (
execution_status == "success"
and self._ratio_from_task(task_row, "self_traffic_bytes") <= 0
and int(task_row["num_nodes"] or 0) < LIGHT_RESTRICTED_MIN_NUM_NODES
):
restriction_status = "light_restricted"
retryability = "not_applicable"
collection_status, derived_reason = _derive_collection_status(
restriction_status,
retryability,
failure_type,
)
return {
"collection_status": collection_status,
"collection_status_reason": str(task_row["error_reason"] or derived_reason or ""),
"restriction_status": restriction_status,
"retryability": retryability,
}
def _catalog_payload(self, catalog_row: sqlite3.Row, collection_task_type: str = "new_app") -> Dict[str, Any]:
payload = _safe_json_loads(catalog_row["task_payload_json"], {}) or {}
if not isinstance(payload, dict) or not payload:
payload = build_task_payload_from_row(
{
"app_name": catalog_row["app_name"] or catalog_row["package_name"],
"package_name": catalog_row["package_name"],
"app_magic_label": catalog_row["app_magic_label"] or "",
"last_updated": catalog_row["last_updated"] or "",
"country_code": catalog_row["country_code"] or "",
"device_type": catalog_row["device_type"] or "",
"downloads": catalog_row["downloads"],
}
)
else:
payload = dict(payload)
payload["app_name"] = payload.get("app_name") or catalog_row["app_name"] or catalog_row["package_name"]
payload["package_name"] = payload.get("package_name") or catalog_row["package_name"]
payload["app_magic_label"] = payload.get("app_magic_label") or catalog_row["app_magic_label"] or ""
payload["last_updated"] = payload.get("last_updated") or catalog_row["last_updated"] or ""
payload["country_code"] = payload.get("country_code") or catalog_row["country_code"] or ""
payload["device_type"] = _normalize_task_device_type(payload.get("device_type") or catalog_row["device_type"] or "")
original_row = payload.get("original_row")
if isinstance(original_row, dict):
original_row_copy = dict(original_row)
original_row_copy["app_magic_label"] = original_row_copy.get("app_magic_label") or payload["app_magic_label"]
original_row_copy["last_updated"] = original_row_copy.get("last_updated") or payload["last_updated"]
original_row_copy["country_code"] = original_row_copy.get("country_code") or payload["country_code"]
original_row_copy["device_type"] = original_row_copy.get("device_type") or payload["device_type"]
payload["original_row"] = original_row_copy
payload["collection_task_type"] = collection_task_type or "new_app"
payload["last_update_interval_days"] = int(catalog_row["last_update_interval_days"] or 0)
return payload
def _task_recency_epoch(self, task_row: Optional[sqlite3.Row]) -> float:
if not task_row:
return 0.0
return max(
self._parse_db_datetime(task_row["completed_at"]),
self._parse_db_datetime(task_row["started_at"]),
self._parse_db_datetime(task_row["created_at"]),
)
def _source_file_count(self, connection: sqlite3.Connection, package_name: str) -> int:
row = connection.execute(
"SELECT COUNT(*) AS total FROM analytics_source_file WHERE package_name = ?",
(package_name,),
).fetchone()
return int((row["total"] or 0) if row else 0)
def _artifact_status(self, connection: sqlite3.Connection, package_name: str, task_row: Optional[sqlite3.Row]) -> str:
has_task = bool(task_row and str(task_row["task_status"] or "").strip())
traffic_file_paths = _safe_json_loads(task_row["traffic_file_paths"], []) if task_row else []
has_files = bool(traffic_file_paths) or self._source_file_count(connection, package_name) > 0
if has_task and has_files:
return "complete"
if has_task or has_files:
return "partial"
return "missing"
def _catalog_summary_from_rows(
self,
connection: sqlite3.Connection,
catalog_row: sqlite3.Row,
task_row: Optional[sqlite3.Row],
) -> Dict[str, Any]:
package_name = str(catalog_row["package_name"] or "").strip()
state = self._state_from_task(task_row)
collection_task_type = "new_app"
if task_row:
collection_task_type = "new_app" if int(task_row["is_new_app"] or 0) else "app_update"
payload = self._catalog_payload(catalog_row, collection_task_type)
latest_test_time = self._parse_db_datetime(task_row["completed_at"]) if task_row else 0.0
latest_status = str(task_row["execution_status"] or "").strip() if task_row else ""
latest_failure_type = self._failure_type_from_task(task_row)
domain_rows = connection.execute(
"""
SELECT domain, domain_type
FROM app_domain_traffic
WHERE package_name = ?
""",
(package_name,),
).fetchall()
unique_domains = sorted(
{
str(row["domain"] or "").strip()
for row in domain_rows
if str(row["domain"] or "").strip()
and str(row["domain_type"] or "") != "model_data"
and not str(row["domain"] or "").startswith("model_data:")
}
)
unique_second_level_domains = sorted(
{
_extract_second_level_domain(domain)
for domain in unique_domains
if _extract_second_level_domain(domain)
}
)
return {
"package_name": package_name,
"app_name": catalog_row["app_name"] or "",
"app_magic_label": _normalize_app_magic_label(catalog_row["app_magic_label"] or ""),
"last_updated": _normalize_last_updated(catalog_row["last_updated"] or ""),
"country_code": catalog_row["country_code"] or "",
"device_type": _normalize_task_device_type(catalog_row["device_type"] or ""),
"source_order": catalog_row["source_order"],
"task_payload": payload,
"task_payload_json": _safe_json_dumps(payload),
"downloads": _parse_downloads_value(catalog_row["downloads"]),
"collection_task_type": collection_task_type,
"last_update_interval_days": int(catalog_row["last_update_interval_days"] or 0),
"incremental_batch_tag": catalog_row["batch_tags"] or "[]",
"incremental_batch_tags": _normalize_incremental_batch_tags(catalog_row["batch_tags"]),
"incremental_batch_marked_at": 0.0,
"catalog_active": bool(catalog_row["is_active"]),
"is_active": int(catalog_row["is_active"] if catalog_row["is_active"] is not None else 1),
"task_queue": str(catalog_row["task_queue"] or "default").strip() or "default",
"collection_status": state["collection_status"],
"collection_status_reason": state["collection_status_reason"],
"restriction_status": state["restriction_status"],
"retryability": state["retryability"],
"latest_task_key": task_row["task_key"] or "" if task_row else "",
"latest_status": latest_status,
"latest_test_time": latest_test_time,
"latest_worker_id": task_row["worker_id"] or "" if task_row else "",
"latest_task_detail": task_row["error_details"] or "" if task_row else "",
"latest_failure_type": latest_failure_type,
"unique_domain_count": len(unique_domains),
"unique_domain_names": unique_domains,
"unique_second_level_domain_count": len(unique_second_level_domains),
"unique_second_level_domains": unique_second_level_domains,
"droidbot_steps": int(task_row["droidbot_steps"] or 0) if task_row else 0,
"gui_agent_steps": int(task_row["gui_agent_steps"] or 0) if task_row else 0,
"duration_seconds": float(task_row["duration_seconds"] or 0.0) if task_row else 0.0,
"num_nodes": int(task_row["num_nodes"] or 0) if task_row else 0,
"num_reached_activities": int(task_row["num_reached_activities"] or 0) if task_row else 0,
"app_num_total_activities": int(task_row["app_num_total_activities"] or 0) if task_row else 0,
"total_traffic_bytes": int(task_row["total_traffic_bytes"] or 0) if task_row else 0,
"self_traffic_bytes": int(task_row["self_traffic_bytes"] or 0) if task_row else 0,
"server_traffic_bytes": int(task_row["server_traffic_bytes"] or 0) if task_row else 0,
"unrecognized_traffic_bytes": int(task_row["unrecognized_traffic_bytes"] or 0) if task_row else 0,
"self_ratio": self._ratio_from_task(task_row, "self_traffic_bytes"),
"recognition_ratio": self._recognition_ratio_from_task(task_row),
"model_flow_count": int(task_row["model_flow_count"] or 0) if task_row else 0,
"model_traffic_bytes": int(task_row["model_traffic_bytes"] or 0) if task_row else 0,
"model_eligible": 1
if task_row and 0 < int(task_row["model_flow_count"] or 0) <= MODEL_TRAFFIC_THRESHOLD
else 0,
"artifact_status": self._artifact_status(connection, package_name, task_row),
"updated_at": self._parse_db_datetime(catalog_row["updated_at"]),
"latest_task_row": task_row,
}
def _list_catalog_summaries(self, connection: sqlite3.Connection) -> List[Dict[str, Any]]:
rows = connection.execute(
"""
SELECT *
FROM app_catalog
ORDER BY
CASE WHEN source_order IS NULL THEN 1 ELSE 0 END ASC,
source_order ASC,
package_name ASC
"""
).fetchall()
return [
self._catalog_summary_from_rows(connection, row, self._latest_task_row(connection, row["package_name"]))
for row in rows
]
def _latest_catalog_status(self, connection: sqlite3.Connection, package_name: str) -> Dict[str, Any]:
catalog_row = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if not catalog_row:
return {}
return self._catalog_summary_from_rows(connection, catalog_row, self._latest_task_row(connection, package_name))
def _next_attempt(self, connection: sqlite3.Connection, package_name: str, batch_tag: str, run_kind: str) -> int:
row = connection.execute(
"""
SELECT COALESCE(MAX(attempt), 0) AS max_attempt
FROM collection_task
WHERE package_name = ?
AND batch_tag = ?
AND run_kind = ?
""",
(package_name, batch_tag, run_kind),
).fetchone()
return int((row["max_attempt"] or 0) if row else 0) + 1
def _catalog_batch_tag(self, catalog_row: sqlite3.Row, fallback: str = "manual") -> str:
tags = _normalize_incremental_batch_tags(catalog_row["batch_tags"])
if tags:
return tags[-1]
return fallback
def _insert_pending_collection_task(
self,
connection: sqlite3.Connection,
catalog_row: sqlite3.Row,
*,
reason: str,
run_kind: str = "ranking",
batch_tag: str = "",
collection_task_type: str = "",
) -> None:
package_name = str(catalog_row["package_name"] or "").strip()
if not package_name:
return
batch_tag = batch_tag or self._catalog_batch_tag(catalog_row, "manual")
existing = self._latest_catalog_status(connection, package_name)
existing_task_type = collection_task_type or str(existing.get("collection_task_type") or "new_app")
is_new_app = 1 if existing_task_type == "new_app" else 0
attempt = self._next_attempt(connection, package_name, batch_tag, run_kind)
app_name = str(catalog_row["app_name"] or package_name).strip() or package_name
task_key = _make_task_key(app_name, package_name)
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
connection.execute(
"""
INSERT INTO collection_task (
package_name, batch_tag, run_kind, attempt,
task_key, app_name, app_magic_label, is_new_app,
task_status, execution_status, error_reason, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'pending', NULL, ?, ?)
""",
(
package_name,
batch_tag,
run_kind,
attempt,
task_key,
app_name,
_normalize_app_magic_label(catalog_row["app_magic_label"] or ""),
is_new_app,
reason,
now_iso,
),
)
def get_collection_row(self, package_name: str) -> Optional[Dict[str, Any]]:
"""查询应用完整信息(新架构:直接读取 app_catalog + collection_task"""
with self._connect() as connection:
summary = self._latest_catalog_status(connection, package_name)
if not summary:
return None
summary.pop("latest_task_row", None)
return summary
def upsert_catalog_entries(self, entries: List[Dict[str, Any]]) -> int:
if not entries:
return 0
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
count = 0
with self._write_lock, self._connect() as connection:
for entry in entries:
package_name = str(entry.get("package_name") or "").strip()
if not package_name:
continue
payload = entry.get("task_payload") or {}
downloads = entry.get("downloads")
if downloads is None:
downloads = _extract_downloads_from_task_payload(payload)
last_updated = _normalize_last_updated(entry.get("last_updated") or payload.get("last_updated"))
country_code = str(entry.get("country_code") or payload.get("country_code") or "").strip()
device_type = _normalize_task_device_type(entry.get("device_type") or payload.get("device_type") or "")
last_update_interval_days = int(entry.get("last_update_interval_days") or payload.get("last_update_interval_days") or 0)
connection.execute(
"""
INSERT INTO app_catalog (
package_name, app_name, app_magic_label, batch_tags,
last_updated, country_code, device_type, task_payload_json,
last_update_interval_days, task_queue, task_priority, source_order, is_active,
downloads, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(package_name) DO UPDATE SET
app_name = COALESCE(NULLIF(excluded.app_name, ''), app_catalog.app_name),
app_magic_label = COALESCE(NULLIF(excluded.app_magic_label, ''), app_catalog.app_magic_label),
batch_tags = COALESCE(NULLIF(excluded.batch_tags, '[]'), app_catalog.batch_tags),
last_updated = COALESCE(NULLIF(excluded.last_updated, ''), app_catalog.last_updated),
country_code = excluded.country_code,
device_type = excluded.device_type,
task_payload_json = excluded.task_payload_json,
last_update_interval_days = COALESCE(excluded.last_update_interval_days, app_catalog.last_update_interval_days),
task_queue = COALESCE(NULLIF(excluded.task_queue, ''), app_catalog.task_queue),
task_priority = COALESCE(excluded.task_priority, app_catalog.task_priority),
source_order = excluded.source_order,
is_active = COALESCE(excluded.is_active, app_catalog.is_active),
downloads = COALESCE(excluded.downloads, app_catalog.downloads),
updated_at = excluded.updated_at
""",
(
package_name,
entry.get("app_name", "") or package_name,
_normalize_app_magic_label(entry.get("app_magic_label", "")),
_incremental_batch_tags_to_json(entry.get("batch_tags") or entry.get("incremental_batch_tags") or []),
last_updated,
country_code,
device_type,
_safe_json_dumps(payload),
last_update_interval_days,
str(entry.get("task_queue") or "default").strip() or "default",
int(entry.get("task_priority") or 50),
entry.get("source_order"),
entry.get("catalog_active", entry.get("is_active", 1)),
downloads,
now_iso,
now_iso,
),
)
catalog_row = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if catalog_row and str(entry.get("collection_status") or "pending").strip() == "pending":
self._insert_pending_collection_task(
connection,
catalog_row,
reason=str(entry.get("collection_status_reason") or "catalog_sync"),
collection_task_type=str(entry.get("collection_task_type") or "new_app"),
)
count += 1
return count
def replace_catalog_entries(self, entries: List[Dict[str, Any]]) -> Dict[str, int]:
normalized_entries: List[Dict[str, Any]] = []
seen_packages = set()
for entry in entries:
package_name = str(entry.get("package_name") or "").strip()
if not package_name or package_name in seen_packages:
continue
seen_packages.add(package_name)
normalized_entries.append(dict(entry))
if not normalized_entries:
return {"listed": 0, "upserted": 0, "deactivated": 0}
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
upserted = 0
listed_packages = [str(entry["package_name"]).strip() for entry in normalized_entries]
placeholders = ",".join("?" for _ in listed_packages)
with self._write_lock, self._connect() as connection:
cursor = connection.execute(
f"""
UPDATE app_catalog
SET is_active = 0,
updated_at = ?
WHERE {ACTIVE_CATALOG_WHERE}
AND package_name NOT IN ({placeholders})
""",
[now_iso, *listed_packages],
)
deactivated = int(cursor.rowcount or 0)
for entry in normalized_entries:
package_name = str(entry.get("package_name") or "").strip()
payload = entry.get("task_payload") or {}
downloads = entry.get("downloads")
if downloads is None:
downloads = _extract_downloads_from_task_payload(payload)
existing_summary = self._latest_catalog_status(connection, package_name)
existing_status = str(existing_summary.get("collection_status") or "").strip()
collection_task_type = str(entry.get("collection_task_type") or existing_summary.get("collection_task_type") or "new_app")
should_mark_pending = not existing_summary or existing_status in {"pending", "failed_terminal"}
pending_reason = str(entry.get("collection_status_reason") or "catalog_sync")
incoming_last_updated = _normalize_last_updated(entry.get("last_updated") or (payload or {}).get("last_updated"))
previous_last_updated = _normalize_last_updated(existing_summary.get("last_updated", "")) if existing_summary else ""
effective_last_updated = incoming_last_updated or previous_last_updated
last_update_interval_days = int(entry.get("last_update_interval_days") or 0)
if incoming_last_updated and previous_last_updated:
incoming_ts = _last_updated_ts(incoming_last_updated)
previous_ts = _last_updated_ts(previous_last_updated)
if incoming_ts > previous_ts:
last_update_interval_days = _last_updated_interval_days(incoming_last_updated, previous_last_updated)
previous_success = (
existing_summary.get("collection_status") == "qualified"
or existing_summary.get("latest_status") == "success"
or existing_summary.get("restriction_status") in {"success", "light_restricted"}
)
collection_task_type = "app_update" if previous_success else "new_app"
should_mark_pending = True
pending_reason = f"catalog_version_update:{previous_last_updated}->{incoming_last_updated}"
elif incoming_last_updated and not previous_last_updated and existing_summary:
should_mark_pending = True
pending_reason = f"catalog_last_updated_fill:{incoming_last_updated}"
if isinstance(payload, dict):
payload = dict(payload)
payload["last_updated"] = effective_last_updated
payload["collection_task_type"] = collection_task_type
payload["last_update_interval_days"] = last_update_interval_days
original_row = payload.get("original_row")
if isinstance(original_row, dict):
original_row_copy = dict(original_row)
original_row_copy["last_updated"] = effective_last_updated
payload["original_row"] = original_row_copy
country_code = str(entry.get("country_code") or (payload or {}).get("country_code") or "").strip()
device_type = _normalize_task_device_type(entry.get("device_type") or (payload or {}).get("device_type") or "")
connection.execute(
"""
INSERT INTO app_catalog (
package_name, app_name, app_magic_label, batch_tags,
last_updated, country_code, device_type, task_payload_json,
last_update_interval_days, task_queue, task_priority, source_order, is_active,
downloads, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?)
ON CONFLICT(package_name) DO UPDATE SET
app_name = COALESCE(NULLIF(excluded.app_name, ''), app_catalog.app_name),
app_magic_label = COALESCE(NULLIF(excluded.app_magic_label, ''), app_catalog.app_magic_label),
batch_tags = COALESCE(NULLIF(excluded.batch_tags, '[]'), app_catalog.batch_tags),
last_updated = COALESCE(NULLIF(excluded.last_updated, ''), app_catalog.last_updated),
country_code = excluded.country_code,
device_type = excluded.device_type,
task_payload_json = excluded.task_payload_json,
last_update_interval_days = COALESCE(excluded.last_update_interval_days, app_catalog.last_update_interval_days),
task_queue = COALESCE(NULLIF(excluded.task_queue, ''), app_catalog.task_queue),
task_priority = COALESCE(excluded.task_priority, app_catalog.task_priority),
source_order = excluded.source_order,
is_active = 1,
downloads = COALESCE(excluded.downloads, app_catalog.downloads),
updated_at = excluded.updated_at
""",
(
package_name,
entry.get("app_name", "") or package_name,
_normalize_app_magic_label(entry.get("app_magic_label", "")),
_incremental_batch_tags_to_json(entry.get("batch_tags") or entry.get("incremental_batch_tags") or []),
effective_last_updated,
country_code,
device_type,
_safe_json_dumps(payload),
last_update_interval_days,
str(entry.get("task_queue") or "default").strip() or "default",
int(entry.get("task_priority") or 50),
entry.get("source_order"),
downloads,
now_iso,
now_iso,
),
)
catalog_row = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if catalog_row and should_mark_pending:
self._insert_pending_collection_task(
connection,
catalog_row,
reason=pending_reason,
collection_task_type=collection_task_type,
)
upserted += 1
return {
"listed": len(normalized_entries),
"upserted": upserted,
"deactivated": deactivated,
}
def list_active_catalog_packages(self) -> List[str]:
with self._connect() as connection:
rows = connection.execute(
"""
SELECT package_name
FROM app_catalog
WHERE COALESCE(is_active, 1) = 1
ORDER BY
CASE WHEN source_order IS NULL THEN 1 ELSE 0 END ASC,
source_order ASC,
package_name ASC
"""
).fetchall()
return [str(row["package_name"]).strip() for row in rows if row["package_name"]]
def clear_incremental_batch_tags(self, batch_tag: str = "") -> int:
normalized_tag = _normalize_incremental_batch_tag(batch_tag)
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
with self._write_lock, self._connect() as connection:
if not normalized_tag:
cursor = connection.execute(
"""
UPDATE app_catalog
SET batch_tags = '[]',
updated_at = ?
WHERE COALESCE(batch_tags, '') NOT IN ('', '[]')
""",
(now_iso,),
)
return int(cursor.rowcount or 0)
rows = connection.execute(
"""
SELECT package_name, batch_tags
FROM app_catalog
WHERE COALESCE(batch_tags, '') != ''
"""
).fetchall()
count = 0
for row in rows:
next_tags = _remove_incremental_batch_tag(row["batch_tags"], normalized_tag)
if next_tags == _normalize_incremental_batch_tags(row["batch_tags"]):
continue
connection.execute(
"""
UPDATE app_catalog
SET batch_tags = ?,
updated_at = ?
WHERE package_name = ?
""",
(_incremental_batch_tags_to_json(next_tags), now_iso, row["package_name"]),
)
count += 1
return count
def set_incremental_batch_tag_for_packages(self, package_names: List[str], batch_tag: str) -> int:
normalized_tag = _normalize_incremental_batch_tag(batch_tag)
normalized_packages = [str(item or "").strip() for item in package_names if str(item or "").strip()]
if not normalized_tag or not normalized_packages:
return 0
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
count = 0
with self._write_lock, self._connect() as connection:
for package_name in normalized_packages:
row = connection.execute(
"SELECT batch_tags FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if not row:
continue
next_tags = _add_incremental_batch_tag(row["batch_tags"], normalized_tag)
connection.execute(
"""
UPDATE app_catalog
SET batch_tags = ?,
updated_at = ?
WHERE package_name = ?
""",
(_incremental_batch_tags_to_json(next_tags), now_iso, package_name),
)
count += 1
return count
def update_collection_statuses(self, updates: List[Dict[str, Any]]) -> int:
if not updates:
return 0
count = 0
with self._write_lock, self._connect() as connection:
for item in updates:
package_name = str(item.get("package_name") or "").strip()
if not package_name:
continue
if str(item.get("collection_status") or "").strip() != "pending":
continue
catalog_row = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if not catalog_row:
continue
self._insert_pending_collection_task(
connection,
catalog_row,
reason=str(item.get("collection_status_reason") or "manual_pending"),
)
count += 1
return count
def upsert_high_priority_app(self, package_name: str, app_name: str, payload: Dict[str, Any]) -> None:
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
app_magic_label = _normalize_app_magic_label(payload.get("app_magic_label", ""))
downloads = payload.get("downloads")
if downloads is None:
downloads = _extract_downloads_from_task_payload(payload)
last_updated = _normalize_last_updated(payload.get("last_updated", ""))
country_code = str(payload.get("country_code") or "").strip()
device_type = _normalize_task_device_type(payload.get("device_type") or "")
collection_task_type = "new_app"
if isinstance(payload, dict):
payload = dict(payload)
payload["collection_task_type"] = collection_task_type
payload["last_update_interval_days"] = int(payload.get("last_update_interval_days") or 0)
with self._write_lock, self._connect() as connection:
connection.execute(
"""
INSERT INTO app_catalog (
package_name, app_name, app_magic_label, batch_tags,
last_updated, country_code, device_type, task_payload_json,
last_update_interval_days,
task_queue, task_priority, source_order, is_active,
downloads, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'high', 100, 0, 1, ?, ?, ?)
ON CONFLICT(package_name) DO UPDATE SET
app_name = COALESCE(NULLIF(excluded.app_name, ''), app_catalog.app_name),
app_magic_label = COALESCE(NULLIF(excluded.app_magic_label, ''), app_catalog.app_magic_label),
batch_tags = COALESCE(app_catalog.batch_tags, excluded.batch_tags),
last_updated = COALESCE(NULLIF(excluded.last_updated, ''), app_catalog.last_updated),
country_code = excluded.country_code,
device_type = excluded.device_type,
task_payload_json = excluded.task_payload_json,
last_update_interval_days = COALESCE(excluded.last_update_interval_days, app_catalog.last_update_interval_days),
source_order = COALESCE(excluded.source_order, app_catalog.source_order),
downloads = COALESCE(excluded.downloads, app_catalog.downloads),
task_queue = 'high',
task_priority = 100,
is_active = 1,
updated_at = excluded.updated_at
""",
(
package_name,
app_name or package_name,
app_magic_label,
"[]",
last_updated,
country_code,
device_type,
_safe_json_dumps(payload),
int(payload.get("last_update_interval_days") or 0),
downloads,
now_iso,
now_iso,
),
)
catalog_row = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if catalog_row:
self._insert_pending_collection_task(
connection,
catalog_row,
reason="high_priority_enqueue",
collection_task_type=collection_task_type,
)
def set_collection_status_for_packages(self, package_names: List[str], collection_status: str, reason: str = "") -> int:
normalized = [str(item or "").strip() for item in package_names if str(item or "").strip()]
if not normalized:
return 0
if str(collection_status or "").strip() != "pending":
return 0
count = 0
with self._write_lock, self._connect() as connection:
for package_name in normalized:
catalog_row = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
if not catalog_row:
continue
self._insert_pending_collection_task(
connection,
catalog_row,
reason=reason or "manual_pending",
)
count += 1
return count
def has_successful_related_magic_label(self, app_magic_label: str, package_name: str) -> bool:
normalized_label = _normalize_app_magic_label(app_magic_label)
normalized_package = str(package_name or "").strip()
if not normalized_label or not normalized_package:
return False
with self._connect() as connection:
row = connection.execute(
"""
SELECT 1
FROM app_catalog ac
JOIN collection_task ct
ON ct.package_name = ac.package_name
WHERE COALESCE(ac.app_magic_label, '') = ?
AND ac.package_name != ?
AND ct.execution_status = 'success'
LIMIT 1
""",
(normalized_label, normalized_package),
).fetchone()
return bool(row)
def list_pending_collection_tasks(self) -> List[Dict[str, Any]]:
with self._connect() as connection:
summaries = self._list_catalog_summaries(connection)
pending = [
item for item in summaries
if item["catalog_active"] and item["collection_status"] == "pending"
]
pending.sort(
key=lambda item: (
0 if item["collection_task_type"] == "new_app" else 1,
-int(item["last_update_interval_days"] or 0),
-int(item["downloads"] or 0),
1 if item["source_order"] is None else 0,
int(item["source_order"] or 0),
item["package_name"],
)
)
return [
{
"package_name": item["package_name"],
"app_name": item["app_name"],
"source_order": item["source_order"],
"collection_task_type": item["collection_task_type"],
"last_update_interval_days": int(item["last_update_interval_days"] or 0),
"task_payload": item["task_payload"],
"task_queue": item["task_queue"],
}
for item in pending
]
def list_qualified_apps(self) -> List[Dict[str, Any]]:
with self._connect() as connection:
summaries = self._list_catalog_summaries(connection)
return [
{
"package_name": item["package_name"],
"app_name": item["app_name"],
"task_payload": item["task_payload"],
}
for item in sorted(summaries, key=lambda value: value["package_name"])
if item["collection_status"] == "qualified"
]
def list_all_catalog_apps(self) -> List[Dict[str, Any]]:
with self._connect() as connection:
summaries = self._list_catalog_summaries(connection)
return [
{
"package_name": item["package_name"],
"app_name": item["app_name"],
"task_payload": item["task_payload"],
}
for item in sorted(summaries, key=lambda value: value["package_name"])
]
def list_model_eligible_apps(self) -> List[Dict[str, Any]]:
with self._connect() as connection:
summaries = self._list_catalog_summaries(connection)
eligible = [item for item in summaries if item["model_eligible"]]
eligible.sort(key=lambda item: (-int(item["downloads"] or 0), int(item["source_order"] or 0), item["package_name"]))
return [
{
"package_name": item["package_name"],
"app_name": item["app_name"],
"downloads": int(item["downloads"] or 0),
"task_payload": item["task_payload"],
}
for item in eligible
]
def replace_package_snapshot(
self,
package_name: str,
summary: Dict[str, Any],
domain_rows: List[Dict[str, Any]],
component_rows: List[Dict[str, Any]],
source_rows: List[Dict[str, Any]],
):
"""
替换应用快照(新架构)。
写入 collection_task执行记录 + app_catalog元数据
"""
now = time.time()
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(now))
latest_test_time = _parse_float_value(summary.get("latest_test_time"))
completed_at = (
time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(latest_test_time))
if latest_test_time > 0
else now_iso
)
# 解析主键字段从 latest_task_key
task_key = summary.get("latest_task_key", "")
batch_tag, run_kind, attempt = self._parse_task_key(
task_key,
batch_tag=summary.get("batch_tag"),
run_kind=summary.get("run_kind"),
attempt=summary.get("attempt"),
)
# 解析错误信息
error_category = None
error_code = None
failure_type = summary.get("latest_failure_type", "")
if "/" in failure_type:
parts = failure_type.split("/", 1)
error_category = parts[0].strip() or None
try:
error_code = int(parts[1].strip())
except (ValueError, IndexError):
pass
elif str(failure_type or "").strip():
error_category = str(failure_type).strip()
traffic_file_paths = summary.get("traffic_file_paths")
if not traffic_file_paths:
traffic_file_paths = [
str(row.get("file_path") or "").strip()
for row in source_rows
if str(row.get("file_path") or "").strip()
]
traffic_file_paths_json = _safe_json_dumps(traffic_file_paths or [])
with self._write_lock, self._connect() as connection:
# 1. 删除旧的详细数据
connection.execute("DELETE FROM app_domain_traffic WHERE package_name = ?", (package_name,))
connection.execute("DELETE FROM app_traffic_component WHERE package_name = ?", (package_name,))
connection.execute("DELETE FROM analytics_source_file WHERE package_name = ?", (package_name,))
# 2. 插入域名流量数据
for row in domain_rows:
connection.execute("""
INSERT INTO app_domain_traffic (
package_name, domain, domain_type, traffic_bytes, flow_count,
traffic_ratio, domain_traffic_ratio, organization,
matched_tp_mark, matched_pattern, match_state, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", (
package_name, row["domain"], row.get("domain_type", ""),
int(row.get("traffic_bytes", 0)), int(row.get("flow_count", 0)),
float(row.get("traffic_ratio", 0)), float(row.get("domain_traffic_ratio", 0)),
row.get("organization", ""), row.get("matched_tp_mark", ""),
row.get("matched_pattern", ""), row.get("match_state", ""),
now,
))
# 3. 插入组件流量数据
for row in component_rows:
connection.execute("""
INSERT INTO app_traffic_component (
package_name, component_name, component_package_names,
is_self, traffic_bytes, share_percent, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?)
""", (
package_name, row["component_name"],
row.get("component_package_names", ""),
1 if row.get("is_self") else 0,
int(row.get("traffic_bytes", 0)), float(row.get("share_percent", 0)),
now,
))
# 4. 插入 source 文件索引
for row in source_rows:
file_path = str(row.get("file_path") or "").strip()
if not file_path:
continue
connection.execute(
"""
INSERT INTO analytics_source_file (
file_path, file_type, package_name, worker_id,
size_bytes, mtime, checksum, last_seen_at, last_processed_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(file_path) DO UPDATE SET
file_type = excluded.file_type,
package_name = excluded.package_name,
worker_id = excluded.worker_id,
size_bytes = excluded.size_bytes,
mtime = excluded.mtime,
checksum = excluded.checksum,
last_seen_at = excluded.last_seen_at,
last_processed_at = excluded.last_processed_at
""",
(
file_path,
row.get("file_type", "traffic"),
package_name,
normalize_worker_id(row.get("worker_id", "")),
int(row.get("size_bytes", 0) or 0),
float(row.get("mtime", 0.0) or 0.0),
row.get("checksum", ""),
now,
now,
),
)
# 5. 插入执行记录到 collection_task
connection.execute("""
INSERT INTO collection_task (
package_name, batch_tag, run_kind, attempt,
task_key, app_name, app_magic_label, is_new_app,
task_status, execution_status, worker_id, created_at, completed_at,
error_category, error_code, error_reason, error_details,
duration_seconds, download_duration_seconds, execution_duration_seconds,
analysis_duration_seconds, droidbot_steps, gui_agent_steps, total_steps,
num_nodes, num_reached_activities, app_num_total_activities,
total_traffic_bytes, self_traffic_bytes, server_traffic_bytes,
unrecognized_traffic_bytes, model_flow_count, model_traffic_bytes,
login_count, register_count, stuck_reason_code, guiagent_message,
scenario_triggered, download_source, is_retry, exit_code, trace_json,
traffic_file_paths
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
ON CONFLICT(package_name, batch_tag, run_kind, attempt) DO UPDATE SET
task_key=excluded.task_key,
app_name=excluded.app_name,
app_magic_label=excluded.app_magic_label,
is_new_app=excluded.is_new_app,
task_status=excluded.task_status,
execution_status=excluded.execution_status,
worker_id=excluded.worker_id,
completed_at=excluded.completed_at,
error_category=excluded.error_category,
error_code=excluded.error_code,
error_reason=excluded.error_reason,
error_details=excluded.error_details,
duration_seconds=excluded.duration_seconds,
download_duration_seconds=excluded.download_duration_seconds,
execution_duration_seconds=excluded.execution_duration_seconds,
analysis_duration_seconds=excluded.analysis_duration_seconds,
droidbot_steps=excluded.droidbot_steps,
gui_agent_steps=excluded.gui_agent_steps,
total_steps=excluded.total_steps,
total_traffic_bytes=excluded.total_traffic_bytes,
self_traffic_bytes=excluded.self_traffic_bytes,
server_traffic_bytes=excluded.server_traffic_bytes,
unrecognized_traffic_bytes=excluded.unrecognized_traffic_bytes,
model_flow_count=excluded.model_flow_count,
model_traffic_bytes=excluded.model_traffic_bytes,
num_nodes=excluded.num_nodes,
num_reached_activities=excluded.num_reached_activities,
app_num_total_activities=excluded.app_num_total_activities,
traffic_file_paths=excluded.traffic_file_paths
""", (
package_name, batch_tag, run_kind, attempt,
task_key, summary.get("app_name"), summary.get("app_magic_label"),
1 if summary.get("collection_task_type") == "new_app" else 0,
"completed", summary.get("latest_status"), normalize_worker_id(summary.get("latest_worker_id")), now_iso, completed_at,
error_category, error_code, summary.get("collection_status_reason"), summary.get("latest_task_detail"),
summary.get("duration_seconds", 0),
summary.get("download_duration_seconds", 0),
summary.get("execution_duration_seconds", summary.get("collect_duration_seconds", 0)),
summary.get("analysis_duration_seconds", 0),
summary.get("droidbot_steps", 0), summary.get("gui_agent_steps", 0),
int(summary.get("droidbot_steps", 0) or 0) + int(summary.get("gui_agent_steps", 0) or 0),
summary.get("num_nodes", 0), summary.get("num_reached_activities", 0), summary.get("app_num_total_activities", 0),
summary.get("total_traffic_bytes", 0), summary.get("self_traffic_bytes", 0), summary.get("server_traffic_bytes", 0),
summary.get("unrecognized_traffic_bytes", 0), summary.get("model_flow_count", 0), summary.get("model_traffic_bytes", 0),
summary.get("login_count", 0), summary.get("register_count", 0), summary.get("stuck_reason_code"),
summary.get("guiagent_message", ""), 1 if summary.get("scenario_triggered") else 0,
summary.get("download_source", ""), 1 if summary.get("is_retry") else 0,
summary.get("exit_code"), _safe_json_dumps(summary.get("trace") or {}),
traffic_file_paths_json,
))
# 6. 更新 app_catalog元数据
existing_catalog = connection.execute(
"SELECT * FROM app_catalog WHERE package_name = ?",
(package_name,),
).fetchone()
batch_tags = _normalize_incremental_batch_tags(existing_catalog["batch_tags"] if existing_catalog else [])
if batch_tag and batch_tag not in batch_tags:
batch_tags.append(batch_tag)
task_payload = summary.get("task_payload")
if not isinstance(task_payload, dict):
task_payload = _safe_json_loads(existing_catalog["task_payload_json"], {}) if existing_catalog else {}
last_updated = _normalize_last_updated(
summary.get("last_updated")
or (task_payload or {}).get("last_updated")
or (existing_catalog["last_updated"] if existing_catalog else "")
)
country_code = str(
summary.get("country_code")
or (task_payload or {}).get("country_code")
or (existing_catalog["country_code"] if existing_catalog else "")
or ""
).strip()
device_type = _normalize_task_device_type(
summary.get("device_type")
or (task_payload or {}).get("device_type")
or (existing_catalog["device_type"] if existing_catalog else "")
or ""
)
last_update_interval_days = int(
summary.get("last_update_interval_days")
or ((task_payload or {}).get("last_update_interval_days") if isinstance(task_payload, dict) else 0)
or (existing_catalog["last_update_interval_days"] if existing_catalog else 0)
or 0
)
if isinstance(task_payload, dict) and task_payload:
task_payload = dict(task_payload)
task_payload["last_updated"] = task_payload.get("last_updated") or last_updated
task_payload["country_code"] = task_payload.get("country_code") or country_code
task_payload["device_type"] = _normalize_task_device_type(task_payload.get("device_type") or device_type)
task_payload["collection_task_type"] = summary.get("collection_task_type") or task_payload.get("collection_task_type") or "new_app"
task_payload["last_update_interval_days"] = last_update_interval_days
connection.execute("""
INSERT INTO app_catalog (
package_name, app_name, app_magic_label, batch_tags,
last_updated, country_code, device_type, task_payload_json,
last_update_interval_days, downloads, source_order, is_active,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(package_name) DO UPDATE SET
app_name=COALESCE(excluded.app_name, app_catalog.app_name),
app_magic_label=COALESCE(NULLIF(excluded.app_magic_label, ''), app_catalog.app_magic_label),
batch_tags=excluded.batch_tags,
last_updated=COALESCE(NULLIF(excluded.last_updated, ''), app_catalog.last_updated),
country_code=COALESCE(NULLIF(excluded.country_code, ''), app_catalog.country_code),
device_type=COALESCE(NULLIF(excluded.device_type, ''), app_catalog.device_type),
task_payload_json=COALESCE(NULLIF(excluded.task_payload_json, '{}'), app_catalog.task_payload_json),
last_update_interval_days=COALESCE(excluded.last_update_interval_days, app_catalog.last_update_interval_days),
downloads=COALESCE(excluded.downloads, app_catalog.downloads),
source_order=COALESCE(excluded.source_order, app_catalog.source_order),
is_active=1,
updated_at=excluded.updated_at
""", (
package_name, summary.get('app_name'), summary.get('app_magic_label'),
_incremental_batch_tags_to_json(batch_tags),
last_updated, country_code, device_type, _safe_json_dumps(task_payload or {}),
last_update_interval_days, summary.get('downloads'), summary.get('source_order'), 1, now_iso, now_iso,
))
def _parse_task_key(
self,
task_key: str,
*,
batch_tag: Any = None,
run_kind: Any = None,
attempt: Any = None,
):
"""从 task_key 解析 batch_tag, run_kind, attempt"""
explicit_batch_tag = str(batch_tag or "").strip()
explicit_run_kind = str(run_kind or "").strip()
explicit_attempt = _parse_int_value(attempt)
if explicit_batch_tag and explicit_run_kind in {"ranking", "block", "model", "manual"} and explicit_attempt > 0:
return explicit_batch_tag, explicit_run_kind, explicit_attempt
parts = str(task_key or "").split("_")
if len(parts) >= 3:
try:
attempt = int(parts[-1])
except (ValueError, IndexError):
attempt = 1
run_kind_candidate = parts[-2]
if run_kind_candidate in ('ranking', 'block', 'model', 'manual'):
run_kind = run_kind_candidate
batch_tag = "_".join(parts[:-2]) or "unknown"
return batch_tag, run_kind, attempt
else:
run_kind = 'ranking'
return 'unknown', run_kind, attempt
if str(task_key or "").endswith("_block"):
return 'unknown', 'block', 1
if str(task_key or "").endswith("_model"):
return 'unknown', 'model', 1
return 'unknown', 'ranking', 1
def delete_snapshots_not_in(self, package_names: List[str]) -> int:
"""删除不在列表中的应用快照(新架构:从 app_catalog 删除)"""
normalized = [str(item or "").strip() for item in package_names if str(item or "").strip()]
if not normalized:
return 0
placeholders = ",".join("?" for _ in normalized)
with self._write_lock, self._connect() as connection:
stale_rows = connection.execute(
f"""
SELECT package_name
FROM app_catalog
WHERE package_name NOT IN ({placeholders})
""",
normalized,
).fetchall()
stale_packages = [row["package_name"] for row in stale_rows if row["package_name"]]
if not stale_packages:
return 0
stale_placeholders = ",".join("?" for _ in stale_packages)
for table_name in (
"app_domain_traffic",
"app_traffic_component",
"analytics_source_file",
"app_catalog",
):
connection.execute(
f"DELETE FROM {table_name} WHERE package_name IN ({stale_placeholders})",
stale_packages,
)
return len(stale_packages)
def clear_all_snapshots(self) -> int:
"""清空所有快照数据(新架构:只清理辅助表,激活 app_catalog"""
with self._write_lock, self._connect() as connection:
row = connection.execute("SELECT COUNT(*) AS total FROM app_catalog").fetchone()
cleared = int((row["total"] or 0) if row else 0)
# 删除快照辅助表数据
for table_name in (
"app_domain_traffic",
"app_traffic_component",
"analytics_source_file",
):
connection.execute(f"DELETE FROM {table_name}")
# 激活所有应用以便重新采集
# 注意:不修改 batch_tags保持原有批次标签
connection.execute(
"""
UPDATE app_catalog
SET is_active = 1,
updated_at = datetime('now','localtime')
""",
)
return cleared
def get_overview(
self,
*,
incremental_batch_tag: str = "",
top_n: Any = DEFAULT_INCREMENTAL_TOP_N,
) -> Dict[str, Any]:
normalized_batch_tag = _normalize_incremental_batch_tag(incremental_batch_tag)
normalized_top_n = _normalize_top_n(top_n)
top_n_label = _top_n_label(normalized_top_n)
pending_preview_limit = 12
with self._connect() as connection:
summaries = self._list_catalog_summaries(connection)
if normalized_batch_tag:
base_scope = [item for item in summaries if normalized_batch_tag in item["incremental_batch_tags"]]
else:
base_scope = [item for item in summaries if item["catalog_active"]]
filtered = [
item for item in base_scope
if item["source_order"] is not None and int(item["source_order"]) < normalized_top_n
]
def _count(predicate) -> int:
return sum(1 for item in filtered if predicate(item))
severe_non_retryable_count = _count(
lambda item: item["restriction_status"] == "severe_restricted"
and item["retryability"] == "non_retryable"
)
breakdown_source_rows = [
{
"package_name": item["package_name"],
"error_type": item["latest_failure_type"],
"downloads": item["downloads"],
"task_payload_json": item["task_payload_json"],
}
for item in filtered
if item["restriction_status"] == "severe_restricted"
and item["retryability"] == "non_retryable"
]
non_retryable_download_rows = sorted(
breakdown_source_rows,
key=lambda item: (item.get("package_name") or ""),
)
pending_preview_rows = sorted(
[item for item in filtered if item["collection_status"] == "pending"],
key=lambda item: (
1 if item["source_order"] is None else 0,
int(item["source_order"] or 0),
item["package_name"],
),
)[:pending_preview_limit]
batch_buckets: Dict[str, Dict[str, Any]] = {}
for item in summaries:
for tag in item["incremental_batch_tags"]:
bucket = batch_buckets.setdefault(tag, {"tag": tag, "app_count": 0, "marked_at": 0.0})
bucket["app_count"] += 1
bucket["marked_at"] = max(float(bucket["marked_at"] or 0.0), float(item["updated_at"] or 0.0))
available_batch_rows = sorted(
batch_buckets.values(),
key=lambda item: (-float(item["marked_at"] or 0.0), item["tag"]),
)
non_retryable_breakdown_buckets: Dict[str, Dict[str, Any]] = {}
for breakdown_row in breakdown_source_rows:
normalized_error_type = str(breakdown_row["error_type"] or "").strip() or "UNMARKED"
bucket = non_retryable_breakdown_buckets.setdefault(
normalized_error_type,
{
"error_type": normalized_error_type,
"label": _humanize_error_type(normalized_error_type),
"app_count": 0,
"downloads_below_10k_count": 0,
"downloads_10k_or_above_count": 0,
"downloads_unknown_count": 0,
},
)
bucket["app_count"] += 1
downloads = _resolve_downloads_for_package(
breakdown_row["package_name"],
breakdown_row["task_payload_json"],
downloads=breakdown_row["downloads"],
)
if downloads is None:
bucket["downloads_unknown_count"] += 1
elif downloads < NON_RETRYABLE_DOWNLOAD_SPLIT_THRESHOLD:
bucket["downloads_below_10k_count"] += 1
else:
bucket["downloads_10k_or_above_count"] += 1
non_retryable_breakdown = sorted(
(
{
"error_type": item["error_type"],
"label": item["label"],
"app_count": int(item["app_count"] or 0),
"downloads_below_10k_count": int(item["downloads_below_10k_count"] or 0),
"downloads_10k_or_above_count": int(item["downloads_10k_or_above_count"] or 0),
"downloads_unknown_count": int(item["downloads_unknown_count"] or 0),
"share_percent": round((int(item["app_count"] or 0) / severe_non_retryable_count) * 100, 2)
if severe_non_retryable_count > 0
else 0.0,
}
for item in non_retryable_breakdown_buckets.values()
),
key=lambda item: (-int(item["app_count"]), str(item["error_type"])),
)
download_distribution = _build_non_retryable_download_distribution(non_retryable_download_rows)
self_ratios = [float(item["self_ratio"] or 0.0) for item in filtered]
last_updated_at = max([float(item["updated_at"] or 0.0) for item in filtered] or [0.0])
return {
"scope": "incremental_topn" if normalized_batch_tag else "topn_catalog",
"incremental_batch_tag": normalized_batch_tag,
"top_n": normalized_top_n,
"top_n_label": top_n_label,
"base_scope_app_count": len(base_scope),
"base_scope_missing_source_order_count": sum(1 for item in base_scope if item["source_order"] is None),
"app_count": len(filtered),
"pending_count": _count(lambda item: item["collection_status"] == "pending"),
"qualified_count": _count(lambda item: item["collection_status"] == "qualified"),
"failed_terminal_count": _count(lambda item: item["collection_status"] == "failed_terminal"),
"analyzed_app_count": _count(lambda item: bool(item["restriction_status"])),
"success_app_count": _count(lambda item: item["latest_status"] == "success"),
"non_success_app_count": _count(lambda item: item["latest_status"] not in {"", "success"}),
"restriction_success_count": _count(lambda item: item["restriction_status"] == "success"),
"light_restricted_count": _count(lambda item: item["restriction_status"] == "light_restricted"),
"severe_restricted_count": _count(lambda item: item["restriction_status"] == "severe_restricted"),
"severe_retryable_count": _count(
lambda item: item["restriction_status"] == "severe_restricted"
and item["retryability"] == "retryable"
),
"severe_non_retryable_count": severe_non_retryable_count,
"pending_non_retryable_count": _count(
lambda item: item["collection_status"] == "pending"
and item["restriction_status"] == "severe_restricted"
and item["retryability"] == "non_retryable"
),
"partial_app_count": _count(lambda item: item["artifact_status"] == "partial"),
"domain_count": sum(int(item["unique_domain_count"] or 0) for item in filtered),
"total_traffic_bytes": sum(int(item["total_traffic_bytes"] or 0) for item in filtered),
"avg_self_ratio": round((sum(self_ratios) / len(self_ratios)) if self_ratios else 0.0, 2),
"last_updated_at": _format_ts(last_updated_at),
"pending_preview_limit": pending_preview_limit,
"pending_preview": [
{
"package_name": pending_row["package_name"],
"app_name": pending_row["app_name"] or "",
"source_order": int(pending_row["source_order"])
if pending_row["source_order"] is not None
else None,
}
for pending_row in pending_preview_rows
],
"non_retryable_breakdown": non_retryable_breakdown,
"non_retryable_download_split_threshold": NON_RETRYABLE_DOWNLOAD_SPLIT_THRESHOLD,
"non_retryable_download_distribution": download_distribution["distribution"],
"non_retryable_downloads_matched_count": download_distribution["matched_downloads_apps"],
"non_retryable_downloads_missing_count": download_distribution["missing_downloads_apps"],
"available_incremental_batches": [
{
"tag": batch_row["tag"] or "",
"app_count": int(batch_row["app_count"] or 0),
"marked_at": _format_ts(batch_row["marked_at"]),
}
for batch_row in available_batch_rows
if str(batch_row["tag"] or "").strip()
],
}
def list_apps(
self,
*,
q: str = "",
latest_status: str = "",
artifact_status: str = "",
collection_status: str = "",
restriction_status: str = "",
retryability: str = "",
incremental_batch_tag: str = "",
top_n: Any = DEFAULT_INCREMENTAL_TOP_N,
sort: str = "updated_at",
order: str = "desc",
page: int = 1,
page_size: int = 50,
) -> Dict[str, Any]:
sort_map = {
"updated_at": "updated_at",
"app_name": "app_name",
"package_name": "package_name",
"collection_status": "collection_status",
"latest_status": "latest_status",
"restriction_status": "restriction_status",
"retryability": "retryability",
"latest_test_time": "latest_test_time",
"self_ratio": "self_ratio",
"recognition_ratio": "recognition_ratio",
"total_traffic_bytes": "total_traffic_bytes",
"unique_domain_count": "unique_domain_count",
"source_order": "source_order",
}
sort_key = sort_map.get(sort, "updated_at")
reverse_order = str(order).lower() != "asc"
normalized_batch_tag = _normalize_incremental_batch_tag(incremental_batch_tag)
normalized_top_n = _normalize_top_n(top_n)
top_n_label = _top_n_label(normalized_top_n)
page = max(1, int(page or 1))
page_size = max(1, min(int(page_size or 50), 200))
offset = (page - 1) * page_size
with self._connect() as connection:
summaries = self._list_catalog_summaries(connection)
if normalized_batch_tag:
filtered = [item for item in summaries if normalized_batch_tag in item["incremental_batch_tags"]]
else:
filtered = [item for item in summaries if item["catalog_active"]]
filtered = [
item for item in filtered
if item["source_order"] is not None and int(item["source_order"]) < normalized_top_n
]
if q:
needle = str(q).lower()
filtered = [
item for item in filtered
if needle in str(item["app_name"] or "").lower()
or needle in str(item["package_name"] or "").lower()
]
if latest_status:
filtered = [item for item in filtered if item["latest_status"] == latest_status]
if artifact_status:
filtered = [item for item in filtered if item["artifact_status"] == artifact_status]
if collection_status:
filtered = [item for item in filtered if item["collection_status"] == collection_status]
if restriction_status:
filtered = [item for item in filtered if item["restriction_status"] == restriction_status]
if retryability:
filtered = [item for item in filtered if item["retryability"] == retryability]
def _sortable_value(item: Dict[str, Any]) -> Any:
value = item.get(sort_key)
if sort_key in {
"updated_at",
"latest_test_time",
"self_ratio",
"recognition_ratio",
"total_traffic_bytes",
"unique_domain_count",
"source_order",
}:
return float(value or 0)
return str(value or "")
if sort_key == "source_order":
filtered.sort(
key=lambda item: (
1 if item["source_order"] is None else 0,
float(item["source_order"] or 0),
item["package_name"],
),
reverse=False,
)
if reverse_order:
filtered = list(reversed(filtered))
else:
filtered.sort(key=lambda item: (_sortable_value(item), item["package_name"]), reverse=reverse_order)
total = len(filtered)
rows = filtered[offset:offset + page_size]
items = []
for row in rows:
items.append(
{
"app_name": row["app_name"] or "",
"package_name": row["package_name"],
"app_magic_label": row["app_magic_label"] or "",
"last_updated": row["last_updated"] or "",
"collection_task_type": row["collection_task_type"] or "new_app",
"last_update_interval_days": int(row["last_update_interval_days"] or 0),
"device_type": _normalize_task_device_type(row["device_type"] or ""),
"source_order": int(row["source_order"]) if row["source_order"] is not None else None,
"collection_status": row["collection_status"] or "pending",
"latest_status": row["latest_status"] or "",
"restriction_status": row["restriction_status"] or "",
"retryability": row["retryability"] or "not_applicable",
"latest_test_time": _format_ts(row["latest_test_time"]),
"latest_worker_id": row["latest_worker_id"] or "",
"downloads": _parse_downloads_value(row["downloads"]),
"incremental_batch_tag": ", ".join(_normalize_incremental_batch_tags(row["incremental_batch_tag"])) or "",
"incremental_batch_tags": _normalize_incremental_batch_tags(row["incremental_batch_tag"]),
"incremental_batch_marked_at": _format_ts(row["incremental_batch_marked_at"]),
"self_traffic_bytes": int(row["self_traffic_bytes"] or 0),
"server_traffic_bytes": int(row["server_traffic_bytes"] or 0),
"self_ratio": round(float(row["self_ratio"] or 0.0), 2),
"recognition_ratio": round(float(row["recognition_ratio"] or 0.0), 2),
"total_traffic_bytes": int(row["total_traffic_bytes"] or 0),
"unique_domain_count": int(row["unique_domain_count"] or 0),
"unique_second_level_domain_count": int(row["unique_second_level_domain_count"] or 0),
"duration_seconds": round(float(row["duration_seconds"] or 0.0), 2),
"num_nodes": int(row["num_nodes"] or 0),
"num_reached_activities": int(row["num_reached_activities"] or 0),
"app_num_total_activities": int(row["app_num_total_activities"] or 0),
"artifact_status": row["artifact_status"] or "missing",
"updated_at": _format_ts(row["updated_at"]),
}
)
return {
"items": items,
"page": page,
"page_size": page_size,
"total": total,
"incremental_batch_tag": normalized_batch_tag,
"top_n": normalized_top_n,
"top_n_label": top_n_label,
}
def get_app_detail(self, package_name: str) -> Optional[Dict[str, Any]]:
with self._connect() as connection:
summary_row = self._latest_catalog_status(connection, package_name)
if not summary_row:
return None
domain_rows = connection.execute(
"""
SELECT domain, domain_type, traffic_bytes, flow_count,
organization, matched_tp_mark, matched_pattern, match_state
FROM app_domain_traffic
WHERE package_name = ?
ORDER BY traffic_bytes DESC, domain ASC
""",
(package_name,),
).fetchall()
component_rows = connection.execute(
"""
SELECT *
FROM app_traffic_component
WHERE package_name = ?
ORDER BY traffic_bytes DESC, component_name ASC
""",
(package_name,),
).fetchall()
raw_domains = [
{
"domain": row["domain"],
"domain_type": row["domain_type"] or "",
"traffic_bytes": int(row["traffic_bytes"] or 0),
"flow_count": int(row["flow_count"] or 0),
"organization": row["organization"] or "",
"matched_pattern": row["matched_pattern"] or "",
"match_state": row["match_state"] or "",
}
for row in domain_rows
]
top_raw_domains = raw_domains[:12]
top_second_level_domains = _aggregate_second_level_domain_rows(domain_rows)[:12]
components = [
{
"component_name": row["component_name"],
"component_package_names": row["component_package_names"] or "",
"is_self": bool(row["is_self"]),
"traffic_bytes": int(row["traffic_bytes"] or 0),
"share_percent": round(float(row["share_percent"] or 0.0), 2),
"updated_at": _format_ts(row["updated_at"]),
}
for row in component_rows
]
summary = {
"app_name": summary_row["app_name"] or "",
"package_name": summary_row["package_name"],
"app_magic_label": summary_row["app_magic_label"] or "",
"last_updated": summary_row["last_updated"] or "",
"collection_task_type": summary_row["collection_task_type"] or "new_app",
"last_update_interval_days": int(summary_row["last_update_interval_days"] or 0),
"country_code": summary_row["country_code"] or "",
"device_type": _normalize_task_device_type(summary_row["device_type"] or ""),
"downloads": _parse_downloads_value(summary_row["downloads"]),
"source_order": int(summary_row["source_order"]) if summary_row["source_order"] is not None else None,
"incremental_batch_tag": ", ".join(_normalize_incremental_batch_tags(summary_row["incremental_batch_tag"])) or "",
"incremental_batch_tags": _normalize_incremental_batch_tags(summary_row["incremental_batch_tag"]),
"incremental_batch_marked_at": _format_ts(summary_row["incremental_batch_marked_at"]),
"collection_status": summary_row["collection_status"] or "pending",
"collection_status_reason": summary_row["collection_status_reason"] or "",
"latest_task_key": summary_row["latest_task_key"] or "",
"latest_status": summary_row["latest_status"] or "",
"restriction_status": summary_row["restriction_status"] or "",
"retryability": summary_row["retryability"] or "not_applicable",
"latest_test_time": _format_ts(summary_row["latest_test_time"]),
"latest_worker_id": summary_row["latest_worker_id"] or "",
"latest_task_detail": summary_row["latest_task_detail"] or "",
"latest_failure_type": summary_row["latest_failure_type"] or "",
"unique_domain_count": int(summary_row["unique_domain_count"] or 0),
"unique_domain_names": summary_row["unique_domain_names"] or [],
"unique_second_level_domain_count": int(summary_row["unique_second_level_domain_count"] or 0),
"unique_second_level_domains": summary_row["unique_second_level_domains"] or [],
"droidbot_steps": int(summary_row["droidbot_steps"] or 0),
"gui_agent_steps": int(summary_row["gui_agent_steps"] or 0),
"duration_seconds": round(float(summary_row["duration_seconds"] or 0.0), 2),
"num_nodes": int(summary_row["num_nodes"] or 0),
"num_reached_activities": int(summary_row["num_reached_activities"] or 0),
"app_num_total_activities": int(summary_row["app_num_total_activities"] or 0),
"total_traffic_bytes": int(summary_row["total_traffic_bytes"] or 0),
"self_traffic_bytes": int(summary_row["self_traffic_bytes"] or 0),
"server_traffic_bytes": int(summary_row["server_traffic_bytes"] or 0),
"unrecognized_traffic_bytes": int(summary_row["unrecognized_traffic_bytes"] or 0),
"self_ratio": round(float(summary_row["self_ratio"] or 0.0), 2),
"recognition_ratio": round(float(summary_row["recognition_ratio"] or 0.0), 2),
"artifact_status": summary_row["artifact_status"] or "missing",
"updated_at": _format_ts(summary_row["updated_at"]),
}
return {
"summary": summary,
"top_raw_domains": top_raw_domains,
"top_second_level_domains": top_second_level_domains,
"components": components,
"top_unmatched_domains": [item for item in top_second_level_domains if item["match_state"] == "unmatched"][:20],
}
class AnalyticsService:
def __init__(
self,
*,
db_path: str = MONITORING_DB_PATH,
traffic_root: str = ANALYTICS_TRAFFIC_ROOT,
traffic_root_block: str = "",
traversal_root: str = ANALYTICS_TRAVERSAL_ROOT,
app_list_path: str = ANALYTICS_TPDPI_APP_LIST,
url_lib_path: str = ANALYTICS_TPDPI_URL_LIB,
artifact_wait_seconds: int = ANALYTICS_ARTIFACT_WAIT_SECONDS,
job_poll_seconds: int = ANALYTICS_JOB_POLL_SECONDS,
on_snapshot_callback=None,
start_worker: bool = True,
):
self.db_path = db_path
self.traffic_root = _resolve_data_root(traffic_root)
self.traffic_root_block = _resolve_data_root(traffic_root_block) if traffic_root_block else ""
self.traversal_root = _resolve_data_root(traversal_root)
self.app_list_path = app_list_path
self.url_lib_path = url_lib_path
self.artifact_wait_seconds = max(0, int(artifact_wait_seconds or 0))
self.job_poll_seconds = max(0.2, float(job_poll_seconds or 1))
self.repo = AnalyticsRepository(db_path=db_path)
self.on_snapshot_callback = on_snapshot_callback
self._rules: Optional[DpiRuleMatcher] = None
self._rules_lock = threading.Lock()
self._stop_event = threading.Event()
self._wake_event = threading.Event()
self._thread: Optional[threading.Thread] = None
if start_worker:
self.start()
def start(self):
if self._thread and self._thread.is_alive():
return
self._stop_event.clear()
self._thread = threading.Thread(target=self._run_loop, name="analytics-worker", daemon=True)
self._thread.start()
def close(self):
self._stop_event.set()
self._wake_event.set()
if self._thread:
self._thread.join(timeout=3)
def _ensure_rules(self) -> DpiRuleMatcher:
if self._rules is not None:
return self._rules
with self._rules_lock:
if self._rules is None:
self._rules = DpiRuleMatcher(self.app_list_path, self.url_lib_path)
return self._rules
def enqueue_incremental(
self,
*,
package_name: str,
task_key: Optional[str] = None,
worker_id: Optional[str] = None,
report_time: Optional[float] = None,
trigger_source: str = "worker_report",
waiting_for_artifacts: bool = True,
scheduled_after: Optional[float] = None,
) -> Dict[str, Any]:
if not package_name:
return {}
job = self.repo.create_job(
job_type="incremental",
status="waiting_artifacts" if waiting_for_artifacts else "queued",
package_name=package_name,
task_key=task_key,
worker_id=worker_id,
trigger_source=trigger_source,
report_time=report_time,
scheduled_after=scheduled_after,
)
self._wake_event.set()
return job
def enqueue_backfill(self, packages: Optional[List[str]] = None, trigger_source: str = "api") -> Dict[str, Any]:
payload = {"scope": "packages" if packages else "all", "packages": packages or []}
job = self.repo.create_job(
job_type="backfill",
status="queued",
trigger_source=trigger_source,
payload=payload,
)
self._wake_event.set()
return job
def enqueue_manual_rebuild(self, package_name: str, trigger_source: str = "api") -> Dict[str, Any]:
job = self.repo.create_job(
job_type="manual_rebuild",
status="queued",
package_name=package_name,
trigger_source=trigger_source,
)
self._wake_event.set()
return job
def handle_worker_event(self, payload: Dict[str, Any]) -> int:
if str(payload.get("event_type") or "").strip() != "artifacts_synced":
return 0
package_name = str(payload.get("package_name") or "").strip()
task_key = str(payload.get("task_key") or "").strip()
worker_id = normalize_worker_id(payload.get("worker_id"))
promoted = self.repo.mark_job_ready(task_key=task_key or None, package_name=package_name or None, worker_id=worker_id or None)
if promoted == 0 and package_name:
self.enqueue_incremental(
package_name=package_name,
task_key=task_key or None,
worker_id=worker_id or None,
report_time=float(payload.get("event_time") or time.time()),
trigger_source="artifacts_synced",
waiting_for_artifacts=False,
scheduled_after=time.time() + ANALYTICS_ARTIFACT_FILE_WAIT_SECONDS,
)
self._wake_event.set()
return promoted
def get_overview(self, **kwargs) -> Dict[str, Any]:
return self.repo.get_overview(**kwargs)
def list_apps(self, **kwargs) -> Dict[str, Any]:
return self.repo.list_apps(**kwargs)
def get_app_detail(self, package_name: str) -> Optional[Dict[str, Any]]:
return self.repo.get_app_detail(package_name)
def clear_incremental_batch_tags(self, batch_tag: str = "") -> int:
return self.repo.clear_incremental_batch_tags(batch_tag)
def set_incremental_batch_tag_for_packages(self, package_names: List[str], batch_tag: str) -> int:
return self.repo.set_incremental_batch_tag_for_packages(package_names, batch_tag)
def list_jobs(self, **kwargs) -> List[Dict[str, Any]]:
return self.repo.list_jobs(**kwargs)
def list_pending_collection_tasks(self) -> List[Dict[str, Any]]:
return self.repo.list_pending_collection_tasks()
def list_model_eligible_apps(self) -> List[Dict[str, Any]]:
return self.repo.list_model_eligible_apps()
def run_backfill_now(self, packages: Optional[List[str]] = None) -> Tuple[str, str]:
return self._process_backfill_job({"payload": {"scope": "packages" if packages else "all", "packages": packages or []}})
def sync_catalog_from_rows(self, rows: List[Dict[str, Any]]) -> int:
result = self.replace_catalog_from_rows(rows)
return int(result.get("upserted") or 0)
def replace_catalog_from_rows(self, rows: List[Dict[str, Any]]) -> Dict[str, int]:
entries = []
seen = set()
for index, row in enumerate(rows):
payload = build_task_payload_from_row(row)
package_name = payload["package_name"]
if not package_name or package_name in seen:
continue
seen.add(package_name)
entries.append(
{
"package_name": package_name,
"app_name": payload["app_name"] or package_name,
"app_magic_label": payload.get("app_magic_label", ""),
"last_updated": payload.get("last_updated", ""),
"country_code": payload["country_code"],
"device_type": payload["device_type"],
"source_order": index,
"task_payload": payload,
"downloads": payload.get("downloads"),
"collection_status": "pending",
"collection_status_reason": "catalog_sync",
}
)
return self.repo.replace_catalog_entries(entries)
def refresh_collection_statuses(
self,
*,
package_names: Optional[List[str]] = None,
allow_legacy_download_error_1_pending: bool = False,
) -> int:
del package_names, allow_legacy_download_error_1_pending
# 新 schema 不落库保存派生状态;读取时直接从 collection_task 计算。
return 0
def add_catalog_app(
self,
*,
app_name: str,
package_name: str,
country_code: str = "",
device_type: str = "",
app_magic_label: str = "",
last_updated: str = "",
) -> int:
normalized_device_type = _normalize_task_device_type(device_type)
payload = build_task_payload_from_row(
{
"app_name": app_name,
"package_name": package_name,
"app_magic_label": app_magic_label,
"last_updated": last_updated,
"country_code": country_code,
"device_type": normalized_device_type,
}
)
with self.repo._connect() as connection:
row = connection.execute("SELECT COALESCE(MAX(source_order), -1) AS max_order FROM app_catalog").fetchone()
max_order = int((row["max_order"] or -1) if row else -1)
return self.repo.upsert_catalog_entries(
[
{
"package_name": package_name,
"app_name": app_name or package_name,
"app_magic_label": _normalize_app_magic_label(app_magic_label),
"last_updated": _normalize_last_updated(last_updated),
"country_code": country_code,
"device_type": normalized_device_type,
"source_order": max_order + 1,
"task_payload": payload,
"downloads": payload.get("downloads"),
"collection_status": "pending",
"collection_status_reason": "manual_add",
}
]
)
def set_packages_pending(self, package_names: List[str], reason: str = "manual_pending") -> int:
return self.repo.set_collection_status_for_packages(package_names, "pending", reason)
def _run_loop(self):
while not self._stop_event.is_set():
try:
job = self.repo.claim_next_job(self.artifact_wait_seconds)
if not job:
self._wake_event.wait(timeout=self.job_poll_seconds)
self._wake_event.clear()
continue
self._process_job(job)
except Exception:
logger.exception("Analytics background loop failed")
self._wake_event.wait(timeout=self.job_poll_seconds)
self._wake_event.clear()
def _process_job(self, job: Dict[str, Any]):
job_id = int(job["id"])
try:
if job["job_type"] == "backfill":
status, message = self._process_backfill_job(job)
else:
package_name = str(job.get("package_name") or "").strip()
if not package_name:
raise ValueError("package_name is required for package rebuild jobs")
summary = self.rebuild_package_now(package_name)
if callable(self.on_snapshot_callback):
try:
self.on_snapshot_callback(package_name, summary)
except Exception:
logger.exception("Analytics snapshot callback failed for %s", package_name)
status = "succeeded" if summary.get("artifact_status") == "complete" else "partial"
message = f"rebuilt {package_name} ({summary.get('artifact_status', 'missing')})"
self.repo.finish_job(job_id, status=status, error_message=message)
except Exception as exc:
logger.exception("Analytics job %s failed", job_id)
self.repo.finish_job(job_id, status="failed", error_message=str(exc))
def _process_backfill_job(self, job: Dict[str, Any]) -> Tuple[str, str]:
payload = job.get("payload") or {}
packages = payload.get("packages") or []
cleared = 0
if not packages:
packages = self.repo.list_active_catalog_packages()
if not packages:
return "failed", "no packages found in active catalog"
cleared = self.repo.clear_all_snapshots()
rebuilt = 0
partial = 0
failed_packages: List[str] = []
for package_name in packages:
try:
summary = self.rebuild_package_now(package_name)
if callable(self.on_snapshot_callback):
try:
self.on_snapshot_callback(package_name, summary)
except Exception:
logger.exception("Analytics snapshot callback failed for %s", package_name)
rebuilt += 1
if summary.get("artifact_status") != "complete":
partial += 1
except Exception:
logger.exception("Analytics backfill failed for package %s", package_name)
failed_packages.append(package_name)
if failed_packages and rebuilt == 0:
return "failed", f"failed packages: {', '.join(failed_packages[:10])}"
if failed_packages or partial:
message_parts = [f"rebuilt={rebuilt}", f"partial={partial}"]
if cleared:
message_parts.append(f"cleared={cleared}")
if failed_packages:
message_parts.append(f"failed={len(failed_packages)}")
return "partial", ", ".join(message_parts)
message_parts = [f"rebuilt={rebuilt}"]
if cleared:
message_parts.append(f"cleared={cleared}")
return "succeeded", ", ".join(message_parts)
def _discover_packages(self) -> set:
packages = set(self.repo.list_all_task_packages())
for root in (self.traffic_root, self.traffic_root_block):
if not root or not os.path.isdir(root):
continue
for worker_dir in os.scandir(root):
if not worker_dir.is_dir():
continue
try:
for package_dir in os.scandir(worker_dir.path):
if package_dir.is_dir():
packages.add(package_dir.name)
except OSError:
continue
return {item for item in packages if item}
def _find_traffic_files(self, package_name: str) -> List[str]:
files: List[str] = []
for root in (self.traffic_root, self.traffic_root_block):
if not root or not os.path.isdir(root):
continue
pattern = os.path.join(root, "*", package_name, "traffic_count*.txt")
files += [path for path in glob.glob(pattern) if os.path.isfile(path)]
return sorted(files)
def _build_source_rows(self, package_name: str, traffic_files: List[str]) -> List[Dict[str, Any]]:
rows = []
for file_path in traffic_files:
if not file_path or not os.path.exists(file_path):
continue
stat = os.stat(file_path)
worker_id = ""
try:
if self.traffic_root_block and os.path.abspath(file_path).startswith(os.path.abspath(self.traffic_root_block)):
relative_path = os.path.relpath(file_path, self.traffic_root_block)
else:
relative_path = os.path.relpath(file_path, self.traffic_root)
worker_id = normalize_worker_id(relative_path.split(os.sep, 1)[0])
except ValueError:
worker_id = ""
rows.append(
{
"file_path": os.path.abspath(file_path),
"file_type": "traffic",
"package_name": package_name,
"worker_id": worker_id,
"size_bytes": int(stat.st_size),
"mtime": float(stat.st_mtime),
"checksum": f"{int(stat.st_size)}:{int(stat.st_mtime)}",
}
)
return rows
def _parse_package_traffic(self, package_name: str, traffic_files: List[str]) -> Dict[str, Any]:
aggregated: Dict[str, Dict[str, Any]] = {}
unique_domains = set()
unique_second_level_domains = set()
app_name = ""
model_flow_count = 0
model_traffic_bytes = 0
for file_path in traffic_files:
try:
with open(file_path, "r", encoding="utf-8", errors="ignore") as handle:
for raw_line in handle:
line = raw_line.strip()
if not line:
continue
fields = [item.strip() for item in line.split(",")]
if len(fields) < 7:
continue
app_id = fields[0].strip()
if package_name and app_id and app_id != package_name:
continue
if fields[1].strip() and not app_name:
app_name = fields[1].strip()
domain = _determine_domain_or_model(fields)
traffic_bytes = _parse_traffic_size(fields[6])
if _is_model_traffic_line(fields):
model_flow_count += 1
model_traffic_bytes += traffic_bytes
organization = fields[7].strip() if len(fields) > 7 else ""
bucket = aggregated.setdefault(
domain,
{"traffic_bytes": 0, "flow_count": 0, "organization": ""},
)
bucket["traffic_bytes"] += traffic_bytes
bucket["flow_count"] += 1
if organization and not bucket["organization"]:
bucket["organization"] = organization
if not domain.startswith("model_data:"):
unique_domains.add(domain)
second_level = _extract_second_level_domain(domain)
if second_level:
unique_second_level_domains.add(second_level)
except OSError as exc:
logger.warning("Failed to read traffic file %s: %s", file_path, exc)
total_traffic_bytes = sum(item["traffic_bytes"] for item in aggregated.values())
total_domain_bytes = sum(item["traffic_bytes"] for key, item in aggregated.items() if not key.startswith("model_data:"))
return {
"app_name": app_name,
"aggregated": aggregated,
"unique_domains": sorted(unique_domains),
"unique_second_level_domains": sorted(unique_second_level_domains),
"total_traffic_bytes": total_traffic_bytes,
"total_domain_bytes": total_domain_bytes,
"model_flow_count": model_flow_count,
"model_traffic_bytes": model_traffic_bytes,
}
def rebuild_package_now(
self,
package_name: str,
*,
latest_task_override: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
rules = self._ensure_rules()
rules.ensure_loaded()
existing_row = self.repo.get_collection_row(package_name) or {}
latest_task = self.repo.get_latest_task_execution(package_name) or {}
if latest_task_override:
latest_task = {**latest_task, **latest_task_override}
traffic_files = self._find_traffic_files(package_name)
traffic = self._parse_package_traffic(package_name, traffic_files)
domain_rows: List[Dict[str, Any]] = []
component_buckets: Dict[str, Dict[str, Any]] = {}
total_traffic_bytes = int(traffic["total_traffic_bytes"] or 0)
total_domain_bytes = int(traffic["total_domain_bytes"] or 0)
self_traffic_bytes = 0
server_traffic_bytes = 0
unrecognized_traffic_bytes = 0
self_component_name = (
rules.pkg_to_tp_mark.get(package_name)
or latest_task.get("app_name")
or traffic.get("app_name")
or package_name
)
self_component_packages = rules.tp_pkg_string.get(self_component_name, package_name)
sorted_domains = sorted(
traffic["aggregated"].items(),
key=lambda item: (-int(item[1]["traffic_bytes"]), item[0]),
)
for domain, bucket in sorted_domains:
domain_type = "model_data" if domain.startswith("model_data:") else "domain"
traffic_bytes = int(bucket["traffic_bytes"] or 0)
flow_count = int(bucket["flow_count"] or 0)
matched_tp_mark = ""
matched_pattern = ""
match_state = ""
if domain.startswith("model_data:"):
match_state = "ip"
unrecognized_traffic_bytes += traffic_bytes
else:
matched_tp_mark, matched_pattern = rules.match_domain(domain)
matched_tp_mark = matched_tp_mark or ""
matched_pattern = matched_pattern or ""
if matched_tp_mark:
owner_packages = rules.tp_packages.get(matched_tp_mark, set())
is_self = package_name in owner_packages
match_state = "self" if is_self else "server"
if is_self:
self_traffic_bytes += traffic_bytes
else:
server_traffic_bytes += traffic_bytes
component = component_buckets.setdefault(
matched_tp_mark,
{
"component_name": matched_tp_mark,
"component_package_names": rules.tp_pkg_string.get(matched_tp_mark, ""),
"is_self": is_self,
"traffic_bytes": 0,
},
)
component["traffic_bytes"] += traffic_bytes
component["is_self"] = bool(component["is_self"]) or is_self
else:
# 未匹配三方规则的域名默认按应用自有域名处理。
match_state = "self"
self_traffic_bytes += traffic_bytes
component = component_buckets.setdefault(
self_component_name,
{
"component_name": self_component_name,
"component_package_names": self_component_packages,
"is_self": True,
"traffic_bytes": 0,
},
)
component["traffic_bytes"] += traffic_bytes
component["is_self"] = True
domain_rows.append(
{
"domain": domain,
"domain_type": domain_type,
"traffic_bytes": traffic_bytes,
"flow_count": flow_count,
"traffic_ratio": round((traffic_bytes / total_traffic_bytes) * 100, 2) if total_traffic_bytes > 0 else 0.0,
"domain_traffic_ratio": round((traffic_bytes / total_domain_bytes) * 100, 2)
if (not domain.startswith("model_data:") and total_domain_bytes > 0)
else 0.0,
"organization": bucket.get("organization", ""),
"matched_tp_mark": matched_tp_mark,
"matched_pattern": matched_pattern,
"match_state": match_state,
}
)
component_rows = []
for component in sorted(component_buckets.values(), key=lambda item: (-int(item["traffic_bytes"]), item["component_name"])):
component_rows.append(
{
"component_name": component["component_name"],
"component_package_names": component["component_package_names"],
"is_self": bool(component["is_self"]),
"traffic_bytes": int(component["traffic_bytes"]),
"share_percent": round((int(component["traffic_bytes"]) / total_traffic_bytes) * 100, 2)
if total_traffic_bytes > 0
else 0.0,
}
)
artifact_status = "missing"
if traffic_files and latest_task:
artifact_status = "complete"
elif traffic_files or latest_task:
artifact_status = "partial"
summary = {
"app_name": latest_task.get("app_name") or traffic.get("app_name") or existing_row.get("app_name") or package_name,
"app_magic_label": _normalize_app_magic_label(existing_row.get("app_magic_label", "")),
"last_updated": _normalize_last_updated(existing_row.get("last_updated", "")),
"country_code": existing_row.get("country_code"),
"device_type": _normalize_task_device_type(existing_row.get("device_type")),
"source_order": existing_row.get("source_order"),
"task_payload": existing_row.get("task_payload"),
"downloads": existing_row.get("downloads"),
"collection_task_type": existing_row.get("collection_task_type", "new_app"),
"last_update_interval_days": int(existing_row.get("last_update_interval_days") or 0),
"latest_task_key": latest_task.get("task_key", ""),
"latest_status": latest_task.get("status", ""),
"latest_test_time": latest_task.get("latest_time", 0.0),
"latest_worker_id": normalize_worker_id(latest_task.get("worker_id", "")),
"latest_task_detail": latest_task.get("result_detail", ""),
"latest_failure_type": latest_task.get("error_type", ""),
"unique_domain_count": len(traffic["unique_domains"]),
"unique_domain_names": traffic["unique_domains"],
"unique_second_level_domain_count": len(traffic["unique_second_level_domains"]),
"unique_second_level_domains": traffic["unique_second_level_domains"],
"droidbot_steps": int(latest_task.get("droidbot_steps") or 0),
"gui_agent_steps": int(latest_task.get("guiagent_steps") or 0),
"duration_seconds": round(float(latest_task.get("total_duration_seconds") or 0.0), 2),
"num_nodes": int(latest_task.get("num_nodes") or 0),
"num_reached_activities": int(latest_task.get("num_reached_activities") or 0),
"app_num_total_activities": int(latest_task.get("app_num_total_activities") or 0),
"total_traffic_bytes": total_traffic_bytes,
"self_traffic_bytes": self_traffic_bytes,
"server_traffic_bytes": server_traffic_bytes,
"unrecognized_traffic_bytes": unrecognized_traffic_bytes,
"self_ratio": round((self_traffic_bytes / total_traffic_bytes) * 100, 2) if total_traffic_bytes > 0 else 0.0,
"recognition_ratio": round(((self_traffic_bytes + server_traffic_bytes) / total_traffic_bytes) * 100, 2)
if total_traffic_bytes > 0
else 0.0,
"artifact_status": artifact_status,
"model_flow_count": int(traffic.get("model_flow_count", 0)),
"model_traffic_bytes": int(traffic.get("model_traffic_bytes", 0)),
}
summary["model_eligible"] = 1 if 0 < summary["model_flow_count"] <= MODEL_TRAFFIC_THRESHOLD else 0
restriction_status, retryability = _classify_restriction_status(
summary["latest_status"],
summary["self_ratio"],
summary["latest_failure_type"],
num_nodes=summary["num_nodes"],
download_errors=latest_task.get("download_errors"),
)
if _should_preserve_non_retryable_reason(existing_row, restriction_status, retryability):
summary["latest_failure_type"] = str(existing_row.get("latest_failure_type") or "").strip()
if str(existing_row.get("latest_task_detail") or "").strip():
summary["latest_task_detail"] = str(existing_row.get("latest_task_detail") or "").strip()
restriction_status = "severe_restricted"
retryability = "non_retryable"
related_magic_label_collected = (
str(summary["latest_status"] or "").strip().lower() != "success"
and bool(summary["app_magic_label"])
and self.repo.has_successful_related_magic_label(summary["app_magic_label"], package_name)
)
if related_magic_label_collected:
restriction_status = "light_restricted"
retryability = "not_applicable"
summary["restriction_status"] = restriction_status
summary["retryability"] = retryability
if not related_magic_label_collected and _should_auto_reroute_to_physical(
summary["device_type"],
summary["latest_failure_type"],
restriction_status,
retryability,
):
summary["device_type"] = "physical"
summary["task_payload"] = _build_task_payload_with_device_type(
summary.get("task_payload"),
app_name=summary["app_name"],
package_name=package_name,
country_code=summary.get("country_code") or "",
device_type="physical",
)
collection_status = "pending"
collection_reason = f"{AUTO_PENDING_PHYSICAL_REASON}:{summary['latest_failure_type']}"
else:
collection_status, collection_reason = _derive_collection_status(
restriction_status,
retryability,
summary["latest_failure_type"],
)
if related_magic_label_collected:
collection_reason = "关联应用已采集"
summary["collection_status"] = collection_status
summary["collection_status_reason"] = collection_reason
source_rows = self._build_source_rows(package_name, traffic_files)
self.repo.replace_package_snapshot(package_name, summary, domain_rows, component_rows, source_rows)
return summary