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

1631 lines
68 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

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

#!/usr/bin/env python3
import argparse
import csv
import glob
import json
import os
import sqlite3
import sys
import time
from typing import Any, List, Optional
import redis as redis_lib
PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
if PROJECT_ROOT not in sys.path:
sys.path.insert(0, PROJECT_ROOT)
from analytics import (
AnalyticsRepository,
AnalyticsService,
DEFAULT_INCREMENTAL_TOP_N,
LIGHT_RESTRICTED_MIN_NUM_NODES,
_humanize_error_type,
_incremental_batch_tags_to_json,
_last_updated_ts,
_normalize_app_magic_label,
_normalize_country_codes,
_normalize_incremental_batch_tags,
_normalize_last_updated,
_normalize_task_device_type,
_resolve_downloads_for_package,
build_task_payload_from_row,
)
from config import ANALYTICS_TRAFFIC_ROOT, MONITORING_DB_PATH, REDIS_HOST, REDIS_PORT, REDIS_DB, TASK_CSV_PATH, channel_name
def _build_service(db_path: str) -> AnalyticsService:
return AnalyticsService(db_path=db_path, start_worker=False)
def _load_active_catalog_packages(service: AnalyticsService) -> List[str]:
with service.repo._connect() as connection:
rows = connection.execute(
"""
SELECT package_name
FROM app_catalog
WHERE COALESCE(is_active, 1) = 1
ORDER BY package_name ASC
"""
).fetchall()
return [str(row["package_name"]).strip() for row in rows if row["package_name"]]
def _load_catalog_snapshot(service: AnalyticsService) -> dict:
with service.repo._connect() as connection:
summaries = service.repo._list_catalog_summaries(connection)
snapshot = {}
for row in summaries:
package_name = str(row["package_name"] or "").strip()
if not package_name:
continue
snapshot[package_name] = {
"catalog_active": 1 if row["catalog_active"] else 0,
"app_magic_label": str(row["app_magic_label"] or "").strip(),
"last_updated": _normalize_last_updated(row["last_updated"] or ""),
"collection_status": str(row["collection_status"] or "").strip(),
"latest_status": str(row["latest_status"] or "").strip(),
"restriction_status": str(row["restriction_status"] or "").strip(),
"collection_task_type": str(row["collection_task_type"] or "new_app").strip() or "new_app",
}
return snapshot
def _read_csv_rows(csv_path: str) -> List[dict]:
with open(csv_path, "r", newline="", encoding="utf-8-sig") as handle:
reader = csv.DictReader(handle)
rows = []
for row in reader:
normalized = {}
for key, value in (row or {}).items():
normalized[str(key or "").strip()] = str(value or "").strip()
rows.append(normalized)
return rows
def _load_csv_package_names(csv_path: str) -> List[str]:
if not csv_path or not os.path.exists(csv_path):
return []
packages: List[str] = []
seen = set()
for row in _read_csv_rows(csv_path):
package_name = str(
row.get("package_name")
or row.get("包名")
or row.get("package")
or ""
).strip()
if not package_name or package_name in seen:
continue
seen.add(package_name)
packages.append(package_name)
return packages
def _csv_value_matches_stored_original(key: str, raw_value: Any, stored_value: Any) -> bool:
normalized_key = str(key or "").strip()
raw_text = str(raw_value or "").strip()
stored_text = str(stored_value or "").strip()
if normalized_key in {"last_updated", "最后更新", "更新时间"}:
raw_date = _normalize_last_updated(raw_text)
stored_date = _normalize_last_updated(stored_text)
return raw_date == stored_date if raw_date or stored_date else raw_text == stored_text
if normalized_key in {"device_type", "设备类型"}:
return _normalize_task_device_type(raw_text) == _normalize_task_device_type(stored_text)
return raw_text == stored_text
def _validate_catalog_payload_storage(service: AnalyticsService, incoming_payloads: List[tuple]) -> List[str]:
issues: List[str] = []
with service.repo._connect() as connection:
for _, row, payload in incoming_payloads:
package_name = str(payload.get("package_name") or "").strip()
if not package_name:
continue
catalog_row = connection.execute(
"""
SELECT last_updated, country_code, device_type, task_payload_json
FROM app_catalog
WHERE package_name = ?
""",
(package_name,),
).fetchone()
if not catalog_row or not str(catalog_row["task_payload_json"] or "").strip():
issues.append(package_name)
continue
incoming_last_updated = _normalize_last_updated(payload.get("last_updated", ""))
if incoming_last_updated and _normalize_last_updated(catalog_row["last_updated"] or "") != incoming_last_updated:
issues.append(package_name)
continue
if str(catalog_row["country_code"] or "").strip() != str(payload.get("country_code") or "").strip():
issues.append(package_name)
continue
if _normalize_task_device_type(catalog_row["device_type"] or "") != _normalize_task_device_type(payload.get("device_type") or ""):
issues.append(package_name)
continue
stored_payload = {}
try:
stored_payload = json.loads(catalog_row["task_payload_json"] or "{}")
except (TypeError, ValueError):
issues.append(package_name)
continue
original_row = stored_payload.get("original_row")
if not isinstance(original_row, dict):
issues.append(package_name)
continue
for key, value in row.items():
normalized_key = str(key or "").strip()
if not normalized_key:
continue
if normalized_key not in original_row:
issues.append(package_name)
break
if not _csv_value_matches_stored_original(normalized_key, value, original_row.get(normalized_key, "")):
issues.append(package_name)
break
return sorted(set(issues))
def _count_pending_collection_task_rows(service: AnalyticsService, package_names: List[str]) -> int:
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 service.repo._connect() as connection:
row = connection.execute(
f"""
SELECT COUNT(*) AS total
FROM collection_task
WHERE package_name IN ({placeholders})
AND task_status = 'pending'
""",
normalized,
).fetchone()
return int((row["total"] or 0) if row else 0)
def _normalize_manual_test_status(raw_status: str) -> str:
normalized = str(raw_status or "").strip()
if normalized in {"测试正常"}:
return "success"
if normalized in {"轻度受限"}:
return "light_restricted"
if normalized in {"严重受限"}:
return "severe_restricted"
raise ValueError("test_status must be 测试正常, 轻度受限, or 严重受限")
def _find_traffic_files_under_root(traffic_root: str, package_name: str) -> List[str]:
normalized_root = str(traffic_root or "").strip()
normalized_package_name = str(package_name or "").strip()
if not normalized_root or not normalized_package_name or not os.path.isdir(normalized_root):
return []
pattern = os.path.join(normalized_root, "**", normalized_package_name, "traffic_count*.txt")
return sorted(path for path in glob.glob(pattern, recursive=True) if os.path.isfile(path))
def add_app(args) -> int:
service = _build_service(args.db_path)
try:
created = service.add_catalog_app(
app_name=args.app_name,
package_name=args.package_name,
country_code=args.country_code,
device_type=args.device_type,
app_magic_label=getattr(args, "app_magic_label", ""),
last_updated=getattr(args, "last_updated", ""),
)
print(f"added={created} package={args.package_name}")
return 0
finally:
service.close()
def sync_from_csv(args) -> int:
csv_path = str(args.csv_path or "").strip() or TASK_CSV_PATH
batch_tags = str(
getattr(args, "batch_tags", "")
or getattr(args, "incremental_batch_tag", "")
or ""
).strip()
dry_run = bool(getattr(args, "dry_run", False))
service = _build_service(args.db_path)
try:
rows = _read_csv_rows(csv_path)
catalog_snapshot = _load_catalog_snapshot(service)
active_before = {
package_name
for package_name, item in catalog_snapshot.items()
if int(item.get("catalog_active") or 0) == 1
}
history_packages = set(catalog_snapshot.keys())
incoming_payloads = []
seen_packages = set()
for index, row in enumerate(rows):
payload = build_task_payload_from_row(row)
package_name = str(payload.get("package_name") or "").strip()
if not package_name or package_name in seen_packages:
continue
seen_packages.add(package_name)
incoming_payloads.append((index, row, payload))
incoming_packages = {str(payload["package_name"]).strip() for _, _, payload in incoming_payloads}
history_new_packages = sorted(incoming_packages - history_packages)
added_packages = history_new_packages
added_count = len(added_packages)
reduced_count = len(active_before - incoming_packages)
reactivated_packages = sorted((incoming_packages & history_packages) - active_before)
app_magic_label_updates = 0
last_updated_fills = 0
version_updates = 0
older_last_updated_rows = 0
unchanged_existing = 0
pending_task_packages = set(added_packages)
for _, _, payload in incoming_payloads:
package_name = str(payload.get("package_name") or "").strip()
existing = catalog_snapshot.get(package_name)
if not existing:
continue
changed = False
incoming_magic_label = _normalize_app_magic_label(payload.get("app_magic_label", ""))
existing_magic_label = _normalize_app_magic_label(existing.get("app_magic_label", ""))
if incoming_magic_label and incoming_magic_label != existing_magic_label:
app_magic_label_updates += 1
changed = True
incoming_last_updated = _normalize_last_updated(payload.get("last_updated", ""))
existing_last_updated = _normalize_last_updated(existing.get("last_updated", ""))
if incoming_last_updated and not existing_last_updated:
last_updated_fills += 1
changed = True
pending_task_packages.add(package_name)
elif incoming_last_updated and existing_last_updated:
incoming_ts = _last_updated_ts(incoming_last_updated)
existing_ts = _last_updated_ts(existing_last_updated)
if incoming_ts > existing_ts:
version_updates += 1
changed = True
pending_task_packages.add(package_name)
elif incoming_ts < existing_ts:
older_last_updated_rows += 1
if existing.get("collection_status") in {"pending", "failed_terminal"}:
pending_task_packages.add(package_name)
if not changed:
unchanged_existing += 1
if dry_run:
print(
"dry_run=1 csv={csv} listed={listed} upserted=0 added={added} reduced={reduced} "
"deactivated={deactivated} history_new={history_new} reactivated={reactivated} "
"app_magic_label_updates={app_magic_label_updates} last_updated_fills={last_updated_fills} "
"version_updates={version_updates} older_last_updated_rows={older_last_updated_rows} "
"unchanged_existing={unchanged_existing} would_create_pending_tasks={pending_tasks}{batch_suffix}".format(
csv=csv_path,
listed=len(incoming_packages),
added=added_count,
reduced=reduced_count,
deactivated=reduced_count,
history_new=len(history_new_packages),
reactivated=len(reactivated_packages),
app_magic_label_updates=app_magic_label_updates,
last_updated_fills=last_updated_fills,
version_updates=version_updates,
older_last_updated_rows=older_last_updated_rows,
unchanged_existing=unchanged_existing,
pending_tasks=len(pending_task_packages),
batch_suffix=(
f" batch_tags={batch_tags} incremental_marked={len(incoming_packages)}"
if batch_tags
else ""
),
)
)
return 0
pending_tasks_before = _count_pending_collection_task_rows(service, sorted(pending_task_packages))
result = service.replace_catalog_from_rows(rows)
incremental_marked = 0
if batch_tags:
service.clear_incremental_batch_tags(batch_tags)
incremental_marked = service.set_incremental_batch_tag_for_packages(sorted(incoming_packages), batch_tags)
pending_tasks_after = _count_pending_collection_task_rows(service, sorted(pending_task_packages))
pending_tasks_created = max(0, pending_tasks_after - pending_tasks_before)
payload_storage_issues = _validate_catalog_payload_storage(service, incoming_payloads)
print(
"csv={csv} listed={listed} upserted={upserted} added={added} reduced={reduced} "
"deactivated={deactivated} history_new={history_new} reactivated={reactivated} "
"pending_tasks_created={pending_tasks_created} payload_storage_issues={payload_issues}{batch_suffix}".format(
csv=csv_path,
listed=int(result.get("listed") or 0),
upserted=int(result.get("upserted") or 0),
added=added_count,
reduced=reduced_count,
deactivated=int(result.get("deactivated") or 0),
history_new=len(history_new_packages),
reactivated=len(reactivated_packages),
pending_tasks_created=pending_tasks_created,
payload_issues=len(payload_storage_issues),
batch_suffix=(
f" batch_tags={batch_tags} incremental_marked={incremental_marked}"
if batch_tags
else ""
),
)
)
try:
r = redis_lib.Redis(host=REDIS_HOST, port=REDIS_PORT, db=REDIS_DB, decode_responses=True)
r.publish(channel_name("catalog:updated"), json.dumps({"added": added_count}))
except Exception:
pass
return 0
finally:
service.close()
def enqueue_high_priority(args) -> int:
csv_path = str(args.csv_path or "").strip()
batch_tags = str(getattr(args, "batch_tags", "") or "").strip()
force_pending = bool(getattr(args, "force_pending", False))
if not csv_path:
raise ValueError("csv_path is required")
service = _build_service(args.db_path)
try:
rows = _read_csv_rows(csv_path)
catalog_snapshot = _load_catalog_snapshot(service) if not force_pending else {}
updated = 0
skipped = 0
seen_packages = set()
for row in rows:
payload = build_task_payload_from_row(row)
package_name = str(payload.get("package_name") or "").strip()
if not package_name or package_name in seen_packages:
continue
seen_packages.add(package_name)
if not force_pending:
existing = catalog_snapshot.get(package_name)
if existing and existing.get("collection_status", "") not in ("", "pending", "failed_terminal"):
skipped += 1
continue
app_name = str(payload.get("app_name") or "").strip()
service.repo.upsert_high_priority_app(package_name, app_name, payload)
updated += 1
incremental_marked = 0
if batch_tags:
service.clear_incremental_batch_tags(batch_tags)
incremental_marked = service.set_incremental_batch_tag_for_packages(sorted(seen_packages), batch_tags)
batch_suffix = f" batch_tags={batch_tags} incremental_marked={incremental_marked}" if batch_tags else ""
force_suffix = " force_pending=1" if force_pending else ""
skip_suffix = f" skipped={skipped}" if skipped else ""
print(f"enqueue-high-priority csv={csv_path} updated={updated}{skip_suffix}{force_suffix}{batch_suffix}")
try:
r = redis_lib.Redis(host=REDIS_HOST, port=REDIS_PORT, db=REDIS_DB, decode_responses=True)
r.publish(channel_name("catalog:updated"), json.dumps({"added": updated, "priority": "high"}))
except Exception:
pass
return 0
finally:
service.close()
def mark_added_from_csv_diff(args) -> int:
old_csv_path = str(getattr(args, "old_csv_path", "") or "").strip()
new_csv_path = str(args.new_csv_path or "").strip()
compare_with_db = bool(getattr(args, "compare_with_db", False))
batch_tags = str(getattr(args, "batch_tags", "") or "").strip()
if not new_csv_path:
raise ValueError("new_csv_path is required")
if not batch_tags:
raise ValueError("batch_tags is required")
if not compare_with_db and not old_csv_path:
raise ValueError("old_csv_path is required when --compare-with-db is not set")
new_packages = set(_load_csv_package_names(new_csv_path))
service = _build_service(args.db_path)
try:
if compare_with_db:
# 与数据库历史采集全量应用进行比较,找出新增包名
db_packages = set(_load_active_catalog_packages(service))
# 同时加载非活跃的包名,确保覆盖全量历史数据
with service.repo._connect() as connection:
all_rows = connection.execute(
"SELECT package_name FROM app_catalog"
).fetchall()
all_db_packages = {str(row["package_name"]).strip() for row in all_rows if row["package_name"]}
old_packages = all_db_packages
source_label = "db(all_history)"
else:
old_packages = set(_load_csv_package_names(old_csv_path))
source_label = f"csv({old_csv_path})"
added_packages = sorted(new_packages - old_packages)
service.clear_incremental_batch_tags(batch_tags)
marked = service.set_incremental_batch_tag_for_packages(added_packages, batch_tags)
finally:
service.close()
print(
"baseline={baseline} new_csv={new_csv} new_total={new_total} added={csv_added} batch_tags={batch_tag} marked={marked}".format(
baseline=source_label,
new_csv=new_csv_path,
new_total=len(new_packages),
csv_added=len(added_packages),
batch_tag=batch_tags,
marked=marked,
)
)
return 0
def show_overview(args) -> int:
batch_tags = str(getattr(args, "batch_tags", "") or "").strip()
top_n = int(getattr(args, "top_n", DEFAULT_INCREMENTAL_TOP_N) or DEFAULT_INCREMENTAL_TOP_N)
service = _build_service(args.db_path)
try:
overview = service.get_overview(
batch_tags=batch_tags,
top_n=top_n,
)
finally:
service.close()
print(
"scope={scope} batch_tag={batch_tag} top_n={top_n} batch_total={batch_total} missing_rank={missing_rank} total={total} qualified={qualified} abnormal={abnormal} pending={pending} analyzed={analyzed}".format(
scope=overview.get("scope") or ("incremental_topn" if batch_tags else "topn_catalog"),
batch_tag=batch_tags or "-",
top_n=int(overview.get("top_n") or top_n),
batch_total=int(overview.get("base_scope_app_count") or 0),
missing_rank=int(overview.get("base_scope_missing_source_order_count") or 0),
total=int(overview.get("app_count") or 0),
qualified=int(overview.get("qualified_count") or 0),
abnormal=int(overview.get("failed_terminal_count") or 0),
pending=int(overview.get("pending_count") or 0),
analyzed=int(overview.get("analyzed_app_count") or 0),
)
)
return 0
def set_batch_tag_from_csv(args) -> int:
csv_path = str(args.csv_path).strip()
tag = str(args.tag).strip()
if not csv_path or not os.path.exists(csv_path):
print(f"错误: CSV 文件不存在: {csv_path}")
return 1
if not tag:
print("错误: --tag 不能为空")
return 1
service = _build_service(args.db_path)
try:
catalog_snapshot = _load_catalog_snapshot(service)
rows = _read_csv_rows(csv_path)
to_tag: List[str] = []
skipped = 0
new_apps = 0
updated_apps = 0
reclassified = 0
seen = set()
for row in rows:
package_name = str(row.get("package_name") or row.get("包名") or "").strip()
if not package_name or package_name in seen:
continue
seen.add(package_name)
existing = catalog_snapshot.get(package_name)
if not existing:
to_tag.append(package_name)
new_apps += 1
continue
incoming_last_updated = _normalize_last_updated(row.get("last_updated", ""))
existing_last_updated = _normalize_last_updated(existing.get("last_updated", ""))
incoming_ts = _last_updated_ts(incoming_last_updated) if incoming_last_updated else 0.0
existing_ts = _last_updated_ts(existing_last_updated) if existing_last_updated else 0.0
is_version_update = incoming_ts > existing_ts > 0.0
is_fill_update = bool(incoming_last_updated) and not existing_last_updated
is_version_or_fill = is_version_update or is_fill_update
incoming_magic_label = _normalize_app_magic_label(row.get("app_magic_label", ""))
existing_magic_label = _normalize_app_magic_label(existing.get("app_magic_label", ""))
is_magic_label_changed = bool(incoming_magic_label) and incoming_magic_label != existing_magic_label
if is_version_or_fill:
to_tag.append(package_name)
updated_apps += 1
elif is_magic_label_changed:
to_tag.append(package_name)
reclassified += 1
else:
skipped += 1
if not to_tag:
print(f"tag={tag} csv_total={len(seen)} new={new_apps} updated={updated_apps} reclassified={reclassified} skipped={skipped} → 无需打tag的应用")
return 0
marked = service.set_incremental_batch_tag_for_packages(to_tag, tag)
print(
f"tag={tag} csv_total={len(seen)} new={new_apps} updated={updated_apps} "
f"reclassified={reclassified} skipped={skipped} marked={marked}"
)
finally:
service.close()
return 0
def clear_batch_tag(args) -> int:
tag = str(args.tag).strip()
if not tag:
print("错误: --tag 不能为空")
return 1
service = _build_service(args.db_path)
try:
with service.repo._write_lock, service.repo._connect() as conn:
conn.row_factory = sqlite3.Row
rows = conn.execute(
"""
SELECT package_name, batch_tags
FROM app_catalog
WHERE COALESCE(batch_tags, '') LIKE '%"'
|| ?
|| '"%'
"""
,
(tag,),
).fetchall()
cleared = 0
skipped = 0
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
for row in rows:
summary = service.repo.get_collection_row(row["package_name"]) or {}
tags = _normalize_incremental_batch_tags(row["batch_tags"])
task_type = str(summary.get("collection_task_type") or "").strip()
is_new_app = task_type == "new_app"
has_multiple_tags = len(tags) > 1
if is_new_app and has_multiple_tags:
# 仅对 old batch 就已标记为 new_app 且被错误追加了第二个 tag 的包做清除
new_tags = [t for t in tags if t != tag]
conn.execute(
"""
UPDATE app_catalog
SET batch_tags = ?,
updated_at = ?
WHERE package_name = ?
""",
(_incremental_batch_tags_to_json(new_tags), now_iso, row["package_name"]),
)
cleared += 1
else:
skipped += 1
print(f"tag={tag} cleared={cleared} skipped={skipped}")
finally:
service.close()
return 0
def deactivate_stale_tagged(args) -> int:
"""将指定 tag 下 new_app 且多 tag被错误追加到旧批次的应用 deactivate"""
tag = str(args.tag).strip()
if not tag:
print("错误: --tag 不能为空")
return 1
service = _build_service(args.db_path)
try:
with service.repo._write_lock, service.repo._connect() as conn:
conn.row_factory = sqlite3.Row
rows = conn.execute(
"""
SELECT package_name, batch_tags
FROM app_catalog
WHERE COALESCE(is_active, 1) = 1
AND COALESCE(batch_tags, '') LIKE '%"'
|| ?
|| '"%'
"""
,
(tag,),
).fetchall()
to_deactivate = []
for r in rows:
summary = service.repo.get_collection_row(r["package_name"]) or {}
tags = _normalize_incremental_batch_tags(r["batch_tags"])
task_type = str(summary.get("collection_task_type") or "").strip()
if task_type == "new_app" and len(tags) > 1:
to_deactivate.append(r["package_name"])
if not to_deactivate:
print(f"tag={tag} 无需 deactivate 的应用")
return 0
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
placeholders = ",".join("?" for _ in to_deactivate)
cursor = conn.execute(
f"""
UPDATE app_catalog
SET is_active = 0, updated_at = ?
WHERE package_name IN ({placeholders})
""",
[now_iso, *to_deactivate],
)
deactivated = int(cursor.rowcount or 0)
print(f"tag={tag} deactivated={deactivated}")
finally:
service.close()
return 0
def set_pending(args) -> int:
packages: List[str] = [str(item or "").strip() for item in (args.packages or []) if str(item or "").strip()]
tag_filter = str(getattr(args, "tag", "") or "").strip()
tag_exact = bool(getattr(args, "tag_exact", False))
error_type_prefix = str(getattr(args, "error_type_prefix", "") or "").strip()
catalog_active = getattr(args, "catalog_active", None)
collection_task_type = str(getattr(args, "collection_task_type", "") or "").strip()
collection_status_filter = str(getattr(args, "status", "") or "").strip()
reason = str(args.reason or "manual_pending").strip() or "manual_pending"
dry_run = bool(getattr(args, "dry_run", False))
has_filters = any([tag_filter, error_type_prefix, catalog_active is not None, collection_task_type, collection_status_filter])
service = _build_service(args.db_path)
try:
summaries = _list_catalog_summaries(service)
selected = summaries
if packages:
package_set = set(packages)
selected = [item for item in selected if item["package_name"] in package_set]
if tag_filter:
if tag_exact:
selected = [
item for item in selected
if _normalize_incremental_batch_tags(item.get("incremental_batch_tags")) == [tag_filter]
]
else:
selected = [
item for item in selected
if tag_filter in (item.get("incremental_batch_tags") or [])
]
if error_type_prefix:
selected = [
item for item in selected
if str(item.get("latest_failure_type") or "").startswith(error_type_prefix)
]
if catalog_active is not None:
selected = [item for item in selected if int(bool(item.get("catalog_active"))) == int(catalog_active)]
if collection_task_type:
if collection_task_type in ("new_app", "new_app_or_null"):
selected = [item for item in selected if (item.get("collection_task_type") or "new_app") == "new_app"]
else:
selected = [item for item in selected if item.get("collection_task_type") == collection_task_type]
if collection_status_filter:
selected = [item for item in selected if item.get("collection_status") == collection_status_filter]
if packages and not has_filters:
if dry_run:
print(f"would_update={len(selected)} status=pending [DRY RUN]")
return 0
updated = service.set_packages_pending([item["package_name"] for item in selected], reason=reason)
print(f"updated={updated} status=pending")
return 0
if not packages and not has_filters:
print("No filter criteria or packages provided.")
return 1
if dry_run:
print(f"would_update={len(selected)} status=pending [DRY RUN]")
return 0
updated = service.set_packages_pending([item["package_name"] for item in selected], reason=reason)
print(f"updated={updated} status=pending")
return 0
finally:
service.close()
def _extract_country_codes_for_export(task_payload_json: str, fallback_country_code: str) -> List[str]:
payload = {}
if task_payload_json:
try:
payload = json.loads(task_payload_json)
except (TypeError, ValueError):
payload = {}
if not isinstance(payload, dict):
payload = {}
raw_country_codes = payload.get("country_codes") or payload.get("country_code") or fallback_country_code
return _normalize_country_codes(raw_country_codes)
def _build_google_play_urls(package_name: str, country_codes: List[str]) -> str:
normalized_package_name = str(package_name or "").strip()
if not normalized_package_name:
return ""
normalized_country_codes = _normalize_country_codes(country_codes)
if not normalized_country_codes:
return f"https://play.google.com/store/apps/details?id={normalized_package_name}"
return ",".join(
f"https://play.google.com/store/apps/details?id={normalized_package_name}&gl={country_code}"
for country_code in normalized_country_codes
)
def _normalize_optional_top_n(value: Any) -> Optional[int]:
if value is None:
return None
normalized_text = str(value).strip()
if not normalized_text:
return None
try:
normalized_top_n = int(normalized_text)
except (TypeError, ValueError) as exc:
raise ValueError("top_n must be a positive integer") from exc
if normalized_top_n <= 0:
raise ValueError("top_n must be a positive integer")
return normalized_top_n
def _list_catalog_summaries(service: AnalyticsService) -> List[dict]:
with service.repo._connect() as connection:
return service.repo._list_catalog_summaries(connection)
def _filter_summaries_by_scope(summaries: List[dict], *, tag: str = "", top_n: Any = None) -> List[dict]:
normalized_tag = str(tag or "").strip()
normalized_top_n = _normalize_optional_top_n(top_n)
filtered = list(summaries)
if normalized_tag:
filtered = [item for item in filtered if normalized_tag in (item.get("incremental_batch_tags") or [])]
if normalized_top_n is not None:
filtered = [
item for item in filtered
if item.get("source_order") is not None and int(item.get("source_order") or 0) < normalized_top_n
]
return filtered
def _scope_parts(tag: str = "", top_n: Any = None) -> List[str]:
parts: List[str] = []
if str(tag or "").strip():
parts.append(f"tag={str(tag).strip()}")
normalized_top_n = _normalize_optional_top_n(top_n)
if normalized_top_n is not None:
parts.append(f"top_n={normalized_top_n}")
return parts
def _set_device_type_for_packages(service: AnalyticsService, summaries: List[dict], device_type: str) -> int:
normalized_device_type = _normalize_task_device_type(device_type)
if not normalized_device_type:
return 0
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
updated = 0
with service.repo._write_lock, service.repo._connect() as connection:
for summary in summaries:
package_name = str(summary.get("package_name") or "").strip()
if not package_name:
continue
payload = dict(summary.get("task_payload") or {})
payload["device_type"] = normalized_device_type
original_row = payload.get("original_row")
if isinstance(original_row, dict):
original_row_copy = dict(original_row)
original_row_copy["device_type"] = normalized_device_type
payload["original_row"] = original_row_copy
connection.execute(
"""
UPDATE app_catalog
SET device_type = ?,
task_payload_json = ?,
updated_at = ?
WHERE package_name = ?
""",
(normalized_device_type, json.dumps(payload, ensure_ascii=False), now_iso, package_name),
)
updated += 1
return updated
def set_pending_by_error_type(args) -> int:
error_types: List[str] = [str(item or "").strip() for item in args.error_types if str(item or "").strip()]
error_type_prefix = str(getattr(args, "error_type_prefix", "") or "").strip()
all_non_retryable = bool(getattr(args, "all_non_retryable", False))
non_retryable_only = bool(getattr(args, "non_retryable_only", False))
pending_only = bool(getattr(args, "pending_only", False))
tag_filter = str(getattr(args, "tag", "") or "").strip()
if not error_types and not error_type_prefix and not all_non_retryable:
raise ValueError("at least one exact error type, error type prefix, or --all-non-retryable is required")
if non_retryable_only and pending_only:
raise ValueError("--non-retryable-only and --pending-only cannot be used together")
if pending_only and all_non_retryable:
raise ValueError("--pending-only cannot be used with --all-non-retryable")
reason = str(args.reason or "manual_pending_by_error_type").strip() or "manual_pending_by_error_type"
requested_device_type = str(getattr(args, "device_type", "") or "").strip()
normalized_device_type = _normalize_task_device_type(requested_device_type) if requested_device_type else ""
if requested_device_type and not normalized_device_type:
raise ValueError("device_type must be emulator, physical, or any")
scope_name = "failed_terminal_non_retryable"
if pending_only:
scope_name = "pending_only"
scope_parts: List[str] = []
if error_types:
scope_parts.append(f"error_types={','.join(error_types)}")
if error_type_prefix:
scope_parts.append(f"error_type_prefix={error_type_prefix}")
if all_non_retryable:
scope_parts.append("all_non_retryable=1")
if non_retryable_only:
scope_parts.append("non_retryable_only=1")
if pending_only:
scope_parts.append("pending_only=1")
if normalized_device_type:
scope_parts.append(f"device_type={normalized_device_type}")
scope_parts.extend(_scope_parts(tag_filter, getattr(args, "top_n", None)))
service = _build_service(args.db_path)
try:
summaries = _filter_summaries_by_scope(
_list_catalog_summaries(service),
tag=tag_filter,
top_n=getattr(args, "top_n", None),
)
if pending_only:
summaries = [item for item in summaries if item.get("collection_status") == "pending"]
else:
summaries = [
item for item in summaries
if item.get("collection_status") == "failed_terminal"
and item.get("restriction_status") == "severe_restricted"
and item.get("retryability") == "non_retryable"
]
if error_types:
summaries = [item for item in summaries if item.get("latest_failure_type") in error_types]
if error_type_prefix:
summaries = [
item for item in summaries
if str(item.get("latest_failure_type") or "").startswith(error_type_prefix)
]
if normalized_device_type:
_set_device_type_for_packages(service, summaries, normalized_device_type)
updated = service.set_packages_pending([item["package_name"] for item in summaries], reason=reason)
finally:
service.close()
scope_text = " ".join(scope_parts) if scope_parts else "all_non_retryable=0"
print(
f"updated={updated} status=pending "
f"scope={scope_name} "
f"{scope_text}"
)
return 0
def set_pending_success_zero_self_ratio(args) -> int:
reason = (
str(args.reason or "manual_pending_success_zero_self_ratio").strip()
or "manual_pending_success_zero_self_ratio"
)
tag_filter = str(getattr(args, "tag", "") or "").strip()
scope_parts = _scope_parts(tag_filter, getattr(args, "top_n", None))
service = _build_service(args.db_path)
try:
summaries = _filter_summaries_by_scope(
_list_catalog_summaries(service),
tag=tag_filter,
top_n=getattr(args, "top_n", None),
)
selected = [
item for item in summaries
if item.get("collection_status") == "qualified"
and item.get("latest_status") == "success"
and item.get("restriction_status") in {"success", "light_restricted"}
and float(item.get("self_ratio") or 0.0) <= float(args.max_self_ratio)
and int(item.get("num_nodes") or 0) < LIGHT_RESTRICTED_MIN_NUM_NODES
]
updated = service.set_packages_pending([item["package_name"] for item in selected], reason=reason)
finally:
service.close()
scope_text = f" {' '.join(scope_parts)}" if scope_parts else ""
print(
"updated="
f"{updated} status=pending latest_status=success max_self_ratio={float(args.max_self_ratio)} "
f"max_num_nodes={LIGHT_RESTRICTED_MIN_NUM_NODES - 1}{scope_text}"
)
return 0
def set_pending_retryable(args) -> int:
reason = str(args.reason or "manual_pending_retryable").strip() or "manual_pending_retryable"
tag_filter = str(getattr(args, "tag", "") or "").strip()
scope_parts = _scope_parts(tag_filter, getattr(args, "top_n", None))
current_status = str(args.current_status or "").strip()
if current_status:
scope_parts.append(f"current_status={current_status}")
service = _build_service(args.db_path)
try:
summaries = _filter_summaries_by_scope(
_list_catalog_summaries(service),
tag=tag_filter,
top_n=getattr(args, "top_n", None),
)
selected = [
item for item in summaries
if item.get("retryability") == "retryable"
and item.get("latest_status") != "success"
]
if current_status:
selected = [item for item in selected if item.get("collection_status") == current_status]
updated = service.set_packages_pending([item["package_name"] for item in selected], reason=reason)
finally:
service.close()
scope_text = f" {' '.join(scope_parts)}" if scope_parts else ""
print(f"updated={updated} status=pending retryability=retryable{scope_text}")
return 0
def set_existing_non_retryable_light_restricted(args) -> int:
reason = (
str(args.reason or "manual_existing_non_retryable_light_restricted").strip()
or "manual_existing_non_retryable_light_restricted"
)
min_num_nodes = int(args.min_num_nodes)
tag_filter = str(getattr(args, "tag", "") or "").strip()
scope_parts = _scope_parts(tag_filter, getattr(args, "top_n", None))
service = _build_service(args.db_path)
try:
summaries = _filter_summaries_by_scope(
_list_catalog_summaries(service),
tag=tag_filter,
top_n=getattr(args, "top_n", None),
)
selected = [
item for item in summaries
if item.get("retryability") == "non_retryable"
and item.get("restriction_status") == "severe_restricted"
and item.get("collection_status") == "failed_terminal"
and int(item.get("num_nodes") or 0) >= min_num_nodes
]
# 新 schema 不落库保存派生状态;满足该条件的包重新入 pending 以便重采。
updated = service.set_packages_pending([item["package_name"] for item in selected], reason=reason)
finally:
service.close()
scope_text = f" {' '.join(scope_parts)}" if scope_parts else ""
print(
f"updated={updated} status=pending "
f"reason={reason} min_num_nodes={min_num_nodes}{scope_text}"
)
return 0
def export_apps_by_error_type(args) -> int:
error_types: List[str] = [str(item or "").strip() for item in args.error_types if str(item or "").strip()]
error_type_prefix = str(getattr(args, "error_type_prefix", "") or "").strip()
output_path = str(getattr(args, "output", "") or "").strip()
non_retryable_only = bool(getattr(args, "non_retryable_only", False))
tag_filter = str(getattr(args, "tag", "") or "").strip()
if not error_types and not error_type_prefix and not non_retryable_only:
raise ValueError("at least one exact error type, error type prefix, or --non-retryable-only is required")
if not output_path:
raise ValueError("output path is required")
output_dir = os.path.dirname(output_path)
if output_dir:
os.makedirs(output_dir, exist_ok=True)
scope_parts = _scope_parts(tag_filter, getattr(args, "top_n", None))
service = _build_service(args.db_path)
try:
rows = _filter_summaries_by_scope(
_list_catalog_summaries(service),
tag=tag_filter,
top_n=getattr(args, "top_n", None),
)
finally:
service.close()
exported_rows = []
for row in rows:
if non_retryable_only and not (
row.get("restriction_status") == "severe_restricted"
and row.get("retryability") == "non_retryable"
):
continue
package_name = str(row["package_name"] or "").strip()
matched_error_type = str(row["latest_failure_type"] or "").strip()
if error_types and matched_error_type not in error_types:
continue
if error_type_prefix and not matched_error_type.startswith(error_type_prefix):
continue
downloads = _resolve_downloads_for_package(
package_name,
row["task_payload_json"],
downloads=row["downloads"],
)
country_codes = _extract_country_codes_for_export(
row["task_payload_json"],
row["country_code"] or "",
)
exported_rows.append(
{
"app_name": row["app_name"] or "",
"package_name": package_name,
"error_type": matched_error_type,
"error_type_description": _humanize_error_type(matched_error_type),
"reason_message": row["latest_task_detail"] or "",
"downloads": "" if downloads is None else int(downloads),
"google_play_urls": _build_google_play_urls(package_name, country_codes),
}
)
with open(output_path, "w", newline="", encoding="utf-8-sig") as handle:
writer = csv.DictWriter(
handle,
fieldnames=[
"app_name",
"package_name",
"error_type",
"error_type_description",
"reason_message",
"downloads",
"google_play_urls",
],
)
writer.writeheader()
writer.writerows(exported_rows)
exact_text = ",".join(error_types) if error_types else "-"
prefix_text = error_type_prefix or "-"
print(
f"exported=1 rows={len(exported_rows)} output={output_path} "
f"error_types={exact_text} error_type_prefix={prefix_text} "
f"non_retryable_only={int(non_retryable_only)}"
f"{(' ' + ' '.join(scope_parts)) if scope_parts else ''}"
)
return 0
def rebuild_from_csv(args) -> int:
csv_path = str(args.csv_path or "").strip() or TASK_CSV_PATH
packages = _load_csv_package_names(csv_path)
if not packages:
print(f"rebuilt=0 partial=0 failed=0 csv={csv_path}")
return 0
service = _build_service(args.db_path)
rebuilt = 0
partial = 0
failed = 0
failed_packages: List[str] = []
try:
for package_name in packages:
try:
summary = service.rebuild_package_now(package_name)
rebuilt += 1
if summary.get("artifact_status") != "complete":
partial += 1
except Exception:
failed += 1
failed_packages.append(package_name)
suffix = f" failed_packages={','.join(failed_packages[:10])}" if failed_packages else ""
print(f"rebuilt={rebuilt} partial={partial} failed={failed} csv={csv_path}{suffix}")
return 0 if failed == 0 else 1
finally:
service.close()
def rebuild_status_from_csv(args) -> int:
csv_path = str(args.csv_path or "").strip()
traffic_root = str(args.traffic_root or "").strip() or ANALYTICS_TRAFFIC_ROOT
rows = _read_csv_rows(csv_path)
normalized_rows = [row for row in rows if str(row.get("package_name") or "").strip()]
if not normalized_rows:
print(f"updated=0 mode=manual-status csv={csv_path} traffic_root={traffic_root}")
return 0
service = AnalyticsService(
db_path=args.db_path,
traffic_root=traffic_root,
start_worker=False,
)
repo = AnalyticsRepository(db_path=args.db_path)
original_find_traffic_files = service._find_traffic_files
service._find_traffic_files = lambda package_name: _find_traffic_files_under_root(traffic_root, package_name)
updated = 0
failed = 0
failed_packages: List[str] = []
try:
for row in normalized_rows:
package_name = str(row.get("package_name") or "").strip()
manual_status = _normalize_manual_test_status(row.get("test_status", ""))
exception_reason = str(row.get("exception_reason") or "").strip()
if manual_status == "severe_restricted" and not exception_reason:
failed += 1
failed_packages.append(package_name)
continue
existing_row = repo.get_collection_row(package_name) or {}
app_name = str(existing_row.get("app_name") or package_name).strip() or package_name
app_magic_label = str(row.get("app_magic_label") or existing_row.get("app_magic_label") or "").strip()
last_updated = str(row.get("last_updated") or existing_row.get("last_updated") or "").strip()
latest_task_override = {
"task_key": f"{app_name}_{package_name}",
"app_name": app_name,
"package_name": package_name,
"worker_id": "manual_csv",
"latest_time": time.time(),
}
if manual_status == "success":
latest_task_override.update(
{
"status": "success",
"error_type": "",
"result_detail": "",
"download_errors": {},
}
)
elif manual_status == "light_restricted":
latest_task_override.update(
{
"status": "failed",
"error_type": "",
"result_detail": "",
"download_errors": {},
}
)
else:
latest_task_override.update(
{
"status": "failed",
"error_type": exception_reason,
"result_detail": exception_reason,
"download_errors": {},
}
)
try:
service.rebuild_package_now(package_name, latest_task_override=latest_task_override)
now_iso = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())
with repo._write_lock, repo._connect() as connection:
connection.execute(
"""
UPDATE app_catalog
SET app_magic_label = COALESCE(NULLIF(?, ''), app_magic_label),
last_updated = COALESCE(NULLIF(?, ''), last_updated),
updated_at = ?
WHERE package_name = ?
""",
(app_magic_label, last_updated, now_iso, package_name),
)
updated += 1
except Exception:
failed += 1
failed_packages.append(package_name)
finally:
service._find_traffic_files = original_find_traffic_files
service.close()
suffix = f" failed_packages={','.join(failed_packages[:10])}" if failed_packages else ""
print(
f"updated={updated} failed={failed} mode=manual-status "
f"csv={csv_path} traffic_root={traffic_root}{suffix}"
)
return 0 if failed == 0 else 1
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="Manage app_catalog and collection_task entries.")
parser.add_argument("--db-path", default=MONITORING_DB_PATH, help="Target monitoring sqlite path.")
subparsers = parser.add_subparsers(dest="command", required=True)
add_parser = subparsers.add_parser("add-app", help="Add one app to app_catalog and create a pending collection_task.")
add_parser.add_argument("--app-name", required=True, help="App name.")
add_parser.add_argument("--package-name", required=True, help="Package name.")
add_parser.add_argument("--app-magic-label", default="", help="Related app magic label.")
add_parser.add_argument("--last-updated", default="", help="App last updated date, e.g. 2026/04/07.")
add_parser.add_argument("--country-code", default="", help="Country code string, e.g. US,JP.")
add_parser.add_argument("--device-type", default="", help="Device type: emulator, physical, or any.")
add_parser.set_defaults(func=add_app)
sync_parser = subparsers.add_parser(
"sync-from-csv",
help="Sync current catalog from csv, keep missing apps in DB but hide them from queue and dashboard.",
)
sync_parser.add_argument(
"--csv-path",
default=TASK_CSV_PATH,
help="CSV path containing app_name/package_name rows.",
)
sync_parser.add_argument(
"--incremental-batch-tag",
default="",
help="Optional batch tag to write into app_catalog.batch_tags for packages in this CSV sync.",
)
sync_parser.add_argument(
"--dry-run",
action="store_true",
help="Preview sync impact without changing app_catalog or batch tags.",
)
sync_parser.set_defaults(func=sync_from_csv)
mark_added_parser = subparsers.add_parser(
"mark-added-from-csv-diff",
help="Compare old/new CSV or new CSV vs DB history, and mark added packages with an incremental batch tag.",
)
mark_added_parser.add_argument("--old-csv-path", default="", help="Baseline CSV path (not required when --compare-with-db is set).")
mark_added_parser.add_argument("--new-csv-path", required=True, help="Latest CSV path.")
mark_added_parser.add_argument(
"--compare-with-db",
action="store_true",
help="Use all historical packages in DB as baseline instead of old CSV.",
)
mark_added_parser.add_argument(
"--incremental-batch-tag",
required=True,
help="Batch tag to write onto packages that exist in new CSV but not in baseline.",
)
mark_added_parser.set_defaults(func=mark_added_from_csv_diff)
show_overview_parser = subparsers.add_parser(
"show-overview",
help="Show analytics overview for the selected incremental batch within TopN of the full catalog.",
)
show_overview_parser.add_argument(
"--incremental-batch-tag",
default="",
help="Optional incremental batch tag. Empty means use the current active catalog.",
)
show_overview_parser.add_argument(
"--top-n",
type=int,
default=DEFAULT_INCREMENTAL_TOP_N,
help=f"Only count apps whose full-catalog rank is within TopN. Default: {DEFAULT_INCREMENTAL_TOP_N}.",
)
show_overview_parser.set_defaults(func=show_overview)
pending_parser = subparsers.add_parser("set-pending", help="Mark existing apps as pending for recollection.")
pending_parser.add_argument("packages", nargs="*", help="Package names to mark pending (optional, use with filter flags).")
pending_parser.add_argument("--tag", default="", help="Only affect apps with this batch_tags.")
pending_parser.add_argument("--error-type-prefix", default="", help="Prefix match for latest_failure_type, e.g. DOWNLOAD_ERROR.")
pending_parser.add_argument("--catalog-active", type=int, choices=[0, 1], default=None, help="Filter by catalog_active (0 or 1).")
pending_parser.add_argument("--collection-task-type", default="", help="Filter by collection_task_type, e.g. new_app or app_update.")
pending_parser.add_argument("--reason", default="manual_pending", help="collection_status_reason to write.")
pending_parser.add_argument("--dry-run", action="store_true", help="Preview what would be updated without making changes.")
pending_parser.add_argument("--tag-exact", action="store_true", help="When set, use exact JSON array match ([\"tag\"]) instead of LIKE for --tag.")
pending_parser.add_argument(
"--status",
default="",
choices=["qualified", "failed_terminal", "pending"],
help="Filter by current collection_status.",
)
pending_parser.set_defaults(func=set_pending)
rebuild_from_csv_parser = subparsers.add_parser(
"rebuild-from-csv",
help="Rebuild analytics by package list, or apply manual status updates from package_name/test_status/exception_reason csv.",
)
rebuild_from_csv_parser.add_argument(
"--csv-path",
default=TASK_CSV_PATH,
help="CSV path containing package_name rows.",
)
rebuild_from_csv_parser.set_defaults(func=rebuild_from_csv)
rebuild_status_from_csv_parser = subparsers.add_parser(
"rebuild-status-from-csv",
help="Read package_name/test_status/exception_reason from csv, rebuild traffic metrics from a folder, and apply manual statuses.",
)
rebuild_status_from_csv_parser.add_argument(
"--csv-path",
required=True,
help="CSV path containing package_name,test_status,exception_reason rows.",
)
rebuild_status_from_csv_parser.add_argument(
"--traffic-root",
default=ANALYTICS_TRAFFIC_ROOT,
help="Folder to recursively search for <package>/traffic_count*.txt files.",
)
rebuild_status_from_csv_parser.set_defaults(func=rebuild_status_from_csv)
pending_by_error_parser = subparsers.add_parser(
"set-pending-by-error-type",
help="Mark failed terminal non_retryable apps as pending by latest_failure_type.",
)
pending_by_error_parser.add_argument("error_types", nargs="*", help="Exact latest_failure_type values to match.")
pending_by_error_parser.add_argument(
"--error-type-prefix",
default="",
help="Prefix match for latest_failure_type, e.g. DOWNLOAD_ERROR/.",
)
pending_by_error_parser.add_argument(
"--all-non-retryable",
action="store_true",
help="Apply to all current failed terminal non_retryable apps.",
)
pending_by_error_parser.add_argument(
"--non-retryable-only",
action="store_true",
help="Only modify current failed_terminal non_retryable apps.",
)
pending_by_error_parser.add_argument(
"--pending-only",
action="store_true",
help="Only modify current pending apps.",
)
pending_by_error_parser.add_argument(
"--reason",
default="manual_pending_by_error_type",
help="collection_status_reason to write.",
)
pending_by_error_parser.add_argument(
"--device-type",
default="",
help="Optional retry device type override: emulator, physical, or any.",
)
pending_by_error_parser.add_argument(
"--tag",
default="",
help="Only affect apps with this batch_tags.",
)
pending_by_error_parser.add_argument(
"--top-n",
type=int,
default=None,
help="Only affect rows whose source_order is within TopN, using source_order < N.",
)
pending_by_error_parser.set_defaults(func=set_pending_by_error_type)
pending_success_zero_self_ratio_parser = subparsers.add_parser(
"set-pending-success-zero-self-ratio",
help=(
"Mark qualified success apps with self_ratio <= threshold and "
f"num_nodes < {LIGHT_RESTRICTED_MIN_NUM_NODES} as pending."
),
)
pending_success_zero_self_ratio_parser.add_argument(
"--max-self-ratio",
type=float,
default=0.0,
help="Maximum self_ratio to match. Default: 0.0.",
)
pending_success_zero_self_ratio_parser.add_argument(
"--reason",
default="manual_pending_success_zero_self_ratio",
help="collection_status_reason to write.",
)
pending_success_zero_self_ratio_parser.add_argument(
"--tag",
default="",
help="Only affect apps with this batch_tags.",
)
pending_success_zero_self_ratio_parser.add_argument(
"--top-n",
type=int,
default=None,
help="Only affect rows whose source_order is within TopN, using source_order < N.",
)
pending_success_zero_self_ratio_parser.set_defaults(func=set_pending_success_zero_self_ratio)
pending_retryable_parser = subparsers.add_parser(
"set-pending-retryable",
help="Mark retryable non-success apps back to pending, useful after removing retry caps.",
)
pending_retryable_parser.add_argument(
"--current-status",
default="",
help="Only reset rows currently in this collection_status, e.g. failed_terminal.",
)
pending_retryable_parser.add_argument(
"--reason",
default="manual_pending_retryable",
help="collection_status_reason to write.",
)
pending_retryable_parser.add_argument(
"--tag",
default="",
help="Only affect apps with this batch_tags.",
)
pending_retryable_parser.add_argument(
"--top-n",
type=int,
default=None,
help="Only affect rows whose source_order is within TopN, using source_order < N.",
)
pending_retryable_parser.set_defaults(func=set_pending_retryable)
existing_non_retryable_parser = subparsers.add_parser(
"set-existing-non-retryable-light-restricted",
help="Requeue matching failed non_retryable apps by creating pending collection_task rows.",
)
existing_non_retryable_parser.add_argument(
"--min-num-nodes",
type=int,
default=LIGHT_RESTRICTED_MIN_NUM_NODES,
help=f"Minimum num_nodes to match. Default: {LIGHT_RESTRICTED_MIN_NUM_NODES}.",
)
existing_non_retryable_parser.add_argument(
"--reason",
default="manual_existing_non_retryable_light_restricted",
help="collection_status_reason to write.",
)
existing_non_retryable_parser.add_argument(
"--tag",
default="",
help="Only affect apps with this batch_tags.",
)
existing_non_retryable_parser.add_argument(
"--top-n",
type=int,
default=None,
help="Only affect rows whose source_order is within TopN, using source_order < N.",
)
existing_non_retryable_parser.set_defaults(func=set_existing_non_retryable_light_restricted)
export_by_error_parser = subparsers.add_parser(
"export-by-error-type",
help="Export matching apps by latest_failure_type to csv.",
)
export_by_error_parser.add_argument("error_types", nargs="*", help="Exact latest_failure_type values to match.")
export_by_error_parser.add_argument(
"--error-type-prefix",
default="",
help="Prefix match for latest_failure_type, e.g. DOWNLOAD_ERROR/.",
)
export_by_error_parser.add_argument(
"--non-retryable-only",
action="store_true",
help="Only export rows currently judged as severe_restricted non_retryable; can be used alone to export all non-retryable apps.",
)
export_by_error_parser.add_argument(
"--tag",
default="",
help="Only export apps with this batch_tags.",
)
export_by_error_parser.add_argument(
"--top-n",
type=int,
default=None,
help="Only export rows whose source_order is within TopN, using source_order < N.",
)
export_by_error_parser.add_argument("--output", required=True, help="Target csv output path.")
export_by_error_parser.set_defaults(func=export_apps_by_error_type)
set_batch_tag_parser = subparsers.add_parser(
"set-batch-tag",
help="Add tag only to new/updated apps in a CSV — skips unchanged existing to avoid cross-batch pollution.",
)
set_batch_tag_parser.add_argument("--csv-path", required=True, help="CSV path containing package_name column.")
set_batch_tag_parser.add_argument("--tag", required=True, help="batch_tags to set on matching apps.")
set_batch_tag_parser.set_defaults(func=set_batch_tag_from_csv)
clear_batch_tag_parser = subparsers.add_parser(
"clear-batch-tag",
help="Remove tag only from new_app apps with multiple tags (cross-batch pollution).",
)
clear_batch_tag_parser.add_argument("--tag", required=True, help="batch_tags to remove.")
clear_batch_tag_parser.set_defaults(func=clear_batch_tag)
deactivate_stale_parser = subparsers.add_parser(
"deactivate-stale-tagged",
help="Deactivate new_app apps that were cross-batch tagged (multi-tag pollution).",
)
deactivate_stale_parser.add_argument("--tag", required=True, help="batch_tags to inspect.")
deactivate_stale_parser.set_defaults(func=deactivate_stale_tagged)
enqueue_high_parser = subparsers.add_parser(
"enqueue-high-priority",
help="Mark apps from CSV as high priority and enqueue for immediate collection.",
)
enqueue_high_parser.add_argument("--csv-path", required=True, help="CSV path containing package_name column.")
enqueue_high_parser.add_argument(
"--incremental-batch-tag",
default="",
help="Optional batch tag to mark packages enqueued by this command.",
)
enqueue_high_parser.add_argument(
"--force-pending",
action="store_true",
help="Force all apps in CSV to pending, including ones already successfully collected.",
)
enqueue_high_parser.set_defaults(func=enqueue_high_priority)
load_block_parser = subparsers.add_parser(
"load-block-tasks",
help="Load qualified apps as block test tasks (from CSV or directly from DB).",
)
load_block_parser.add_argument(
"--csv-path",
default="",
help="Optional: CSV file with package_name column. If not provided, loads all qualified apps from DB.",
)
load_block_parser.add_argument(
"--all",
dest="load_all",
action="store_true",
help="加载数据库中所有应用(忽略 collection_status与 --csv-path 配合时跳过 qualified 校验。",
)
load_block_parser.set_defaults(func=load_block_tasks)
force_reapk_parser = subparsers.add_parser(
"force-reapk",
help="清除指定应用的本地APK缓存并触发US重新下载。",
)
force_reapk_parser.add_argument("package_name", help="要重新下载的应用包名")
force_reapk_parser.set_defaults(func=force_reapk)
return parser
def force_reapk(args) -> int:
pkg = str(args.package_name or "").strip()
if not pkg:
print("package_name is required")
return 1
try:
from redis_task_distribute import RedisTaskDispatcher
dispatcher = RedisTaskDispatcher()
n = dispatcher.force_reapk(pkg)
print(f"[force_reapk] 已重置 {pkg} ({n} 个任务),清除本地缓存并等待中心 APK 缓存刷新")
return 0
except Exception as e:
print(f"[force_reapk] 失败: {e}")
return 1
def load_block_tasks(args) -> int:
csv_path = str(args.csv_path or "").strip()
load_all = bool(getattr(args, "load_all", False))
service = _build_service(args.db_path)
try:
if csv_path:
rows = _read_csv_rows(csv_path)
if load_all:
db_apps = service.repo.list_all_catalog_apps()
else:
db_apps = service.repo.list_qualified_apps()
db_lookup = {item["package_name"]: item for item in db_apps}
packages = []
skipped = 0
seen = set()
for row in rows:
payload = build_task_payload_from_row(row)
package_name = str(payload.get("package_name") or "").strip()
if not package_name or package_name in seen:
continue
seen.add(package_name)
if package_name not in db_lookup:
skipped += 1
print(f"[跳过] {package_name}: 未在数据库中")
continue
db_entry = db_lookup[package_name]
app_name = str(db_entry.get("app_name") or payload.get("app_name") or "").strip()
packages.append({
"package_name": package_name,
"app_name": app_name,
"task_payload": db_entry.get("task_payload") or payload,
})
if skipped:
print(f"[CSV校验] 共跳过 {skipped} 个未在数据库中的包名")
else:
if load_all:
packages = service.repo.list_all_catalog_apps()
else:
packages = service.repo.list_qualified_apps()
if not packages:
label = "数据库" if load_all else "qualified"
print(f"[Block测试] 没有符合条件的{label}应用")
return 0
print(f"[Block测试] 准备加载 {len(packages)} 个任务")
try:
from redis_task_distribute import RedisTaskDispatcher
dispatcher = RedisTaskDispatcher()
count = dispatcher.load_block_tasks(packages)
if count > 0:
r = redis_lib.Redis(host=REDIS_HOST, port=REDIS_PORT, db=REDIS_DB, decode_responses=True)
r.publish(channel_name("catalog:updated"), json.dumps({"added": count, "priority": "high_block"}))
print(f"[Block测试] 成功加载 {count} 个 block 任务到高优队列")
except Exception as e:
print(f"[Block测试] 加载失败: {e}")
return 1
return 0
finally:
service.close()
def main() -> int:
parser = build_parser()
args = parser.parse_args()
return int(args.func(args))
if __name__ == "__main__":
raise SystemExit(main())