3656 lines
164 KiB
Python
3656 lines
164 KiB
Python
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
|