207 lines
9.2 KiB
Python
207 lines
9.2 KiB
Python
"""
|
||
replace_package_snapshot 新实现 - 适配新架构
|
||
|
||
核心改动:
|
||
1. 拆分写入:collection_task(执行记录) + app_catalog(元数据)
|
||
2. 从 latest_task_key 解析 batch_tag, run_kind, attempt
|
||
3. 事务保护确保原子性
|
||
"""
|
||
|
||
def _parse_task_key(task_key: str):
|
||
"""
|
||
从 task_key 解析 batch_tag, run_kind, attempt
|
||
|
||
格式: {batch_tag}_{package_name}_{run_kind}_{attempt}
|
||
示例: 2026-6-15_com.test.app_ranking_1
|
||
"""
|
||
parts = str(task_key or "").split("_")
|
||
if len(parts) >= 4:
|
||
# 最后两个是 run_kind 和 attempt
|
||
attempt = 1
|
||
try:
|
||
attempt = int(parts[-1])
|
||
except (ValueError, IndexError):
|
||
pass
|
||
|
||
run_kind = parts[-2] if parts[-2] in ('ranking', 'block', 'model', 'manual') else 'ranking'
|
||
|
||
# batch_tag 是除了最后3部分之外的(去掉package_name, run_kind, attempt)
|
||
# 但 package_name 可能包含多个_,所以简化:取第一个部分作为batch
|
||
batch_tag = parts[0] if parts else 'unknown'
|
||
|
||
return batch_tag, run_kind, attempt
|
||
|
||
# 降级处理
|
||
if '_block' in task_key:
|
||
return 'unknown', 'block', 1
|
||
return 'unknown', 'ranking', 1
|
||
|
||
|
||
def replace_package_snapshot_v2(
|
||
self,
|
||
package_name: str,
|
||
summary: dict,
|
||
domain_rows: list,
|
||
component_rows: list,
|
||
source_rows: list,
|
||
):
|
||
"""
|
||
替换应用快照(新架构版本)
|
||
|
||
核心逻辑:
|
||
1. 写入 collection_task(执行记录)
|
||
2. 更新 app_catalog(元数据)
|
||
3. 写入 app_domain_traffic, app_traffic_component(详细数据)
|
||
"""
|
||
import time
|
||
now = time.time()
|
||
now_iso = time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(now))
|
||
|
||
# 解析主键字段
|
||
task_key = summary.get('latest_task_key', '')
|
||
batch_tag, run_kind, attempt = _parse_task_key(task_key)
|
||
|
||
# 准备 collection_task 数据
|
||
task_data = {
|
||
'package_name': package_name,
|
||
'batch_tag': batch_tag,
|
||
'run_kind': run_kind,
|
||
'attempt': attempt,
|
||
'task_key': task_key,
|
||
'app_name': summary.get('app_name'),
|
||
'app_magic_label': summary.get('app_magic_label'),
|
||
'is_new_app': 1 if summary.get('collection_task_type') == 'new_app' else 0,
|
||
'task_status': 'completed',
|
||
'execution_status': summary.get('latest_status'), # 'success' / 'failed'
|
||
'worker_id': summary.get('latest_worker_id'),
|
||
'completed_at': now_iso,
|
||
# 错误信息(从 latest_failure_type 解析)
|
||
'error_category': None,
|
||
'error_code': None,
|
||
'error_reason': summary.get('collection_status_reason'),
|
||
'error_details': summary.get('latest_task_detail'),
|
||
# 执行统计
|
||
'duration_seconds': summary.get('duration_seconds', 0),
|
||
'droidbot_steps': summary.get('droidbot_steps', 0),
|
||
'gui_agent_steps': summary.get('gui_agent_steps', 0),
|
||
'num_nodes': summary.get('num_nodes', 0),
|
||
'num_reached_activities': summary.get('num_reached_activities', 0),
|
||
'app_num_total_activities': summary.get('app_num_total_activities', 0),
|
||
# 流量统计
|
||
'total_traffic_bytes': summary.get('total_traffic_bytes', 0),
|
||
'self_traffic_bytes': summary.get('self_traffic_bytes', 0),
|
||
'server_traffic_bytes': summary.get('server_traffic_bytes', 0),
|
||
'unrecognized_traffic_bytes': summary.get('unrecognized_traffic_bytes', 0),
|
||
'model_flow_count': summary.get('model_flow_count', 0),
|
||
'model_traffic_bytes': summary.get('model_traffic_bytes', 0),
|
||
}
|
||
|
||
# 解析 error_category 和 error_code
|
||
failure_type = summary.get('latest_failure_type', '')
|
||
if '/' in failure_type:
|
||
cat, code_str = failure_type.split('/', 1)
|
||
task_data['error_category'] = cat.strip()
|
||
try:
|
||
task_data['error_code'] = int(code_str.strip())
|
||
except ValueError:
|
||
pass
|
||
|
||
# 准备 app_catalog 数据(元数据)
|
||
catalog_data = {
|
||
'package_name': package_name,
|
||
'app_name': summary.get('app_name'),
|
||
'app_magic_label': summary.get('app_magic_label'),
|
||
'downloads': summary.get('downloads'),
|
||
'source_order': summary.get('source_order'),
|
||
'is_active': 1,
|
||
}
|
||
|
||
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["domain_type"],
|
||
int(row["traffic_bytes"]), int(row["flow_count"]),
|
||
float(row["traffic_ratio"]), float(row["domain_traffic_ratio"]),
|
||
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["traffic_bytes"]), float(row["share_percent"]),
|
||
now,
|
||
))
|
||
|
||
# 4. 插入执行记录到 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, completed_at,
|
||
error_category, error_code, error_reason, error_details,
|
||
duration_seconds, droidbot_steps, gui_agent_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,
|
||
created_at
|
||
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||
ON CONFLICT(package_name, batch_tag, run_kind, attempt) DO UPDATE SET
|
||
execution_status=excluded.execution_status,
|
||
error_category=excluded.error_category,
|
||
error_code=excluded.error_code,
|
||
total_traffic_bytes=excluded.total_traffic_bytes,
|
||
self_traffic_bytes=excluded.self_traffic_bytes,
|
||
num_nodes=excluded.num_nodes,
|
||
updated_at=datetime('now','localtime')
|
||
""", (
|
||
task_data['package_name'], task_data['batch_tag'], task_data['run_kind'], task_data['attempt'],
|
||
task_data['task_key'], task_data['app_name'], task_data['app_magic_label'], task_data['is_new_app'],
|
||
task_data['task_status'], task_data['execution_status'], task_data['worker_id'], task_data['completed_at'],
|
||
task_data['error_category'], task_data['error_code'], task_data['error_reason'], task_data['error_details'],
|
||
task_data['duration_seconds'], task_data['droidbot_steps'], task_data['gui_agent_steps'],
|
||
task_data['num_nodes'], task_data['num_reached_activities'], task_data['app_num_total_activities'],
|
||
task_data['total_traffic_bytes'], task_data['self_traffic_bytes'], task_data['server_traffic_bytes'],
|
||
task_data['unrecognized_traffic_bytes'], task_data['model_flow_count'], task_data['model_traffic_bytes'],
|
||
now_iso,
|
||
))
|
||
|
||
# 5. 更新 app_catalog(元数据)
|
||
connection.execute("""
|
||
INSERT INTO app_catalog (
|
||
package_name, app_name, app_magic_label, downloads, source_order, is_active, updated_at
|
||
) VALUES (?, ?, ?, ?, ?, ?, datetime('now','localtime'))
|
||
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),
|
||
downloads=COALESCE(excluded.downloads, app_catalog.downloads),
|
||
source_order=COALESCE(excluded.source_order, app_catalog.source_order),
|
||
updated_at=datetime('now','localtime')
|
||
""", (
|
||
catalog_data['package_name'], catalog_data['app_name'], catalog_data['app_magic_label'],
|
||
catalog_data['downloads'], catalog_data['source_order'], catalog_data['is_active'],
|
||
))
|
||
|
||
connection.commit()
|