242 lines
8.3 KiB
Python
242 lines
8.3 KiB
Python
import json
|
|
|
|
from analytics import AnalyticsService, _make_task_key
|
|
from redis_task_distribute import RedisTaskDispatcher
|
|
|
|
|
|
class FakeRedis:
|
|
def __init__(self):
|
|
self.lists = {}
|
|
self.hashes = {}
|
|
self.sets = {}
|
|
|
|
def delete(self, key):
|
|
self.lists.pop(key, None)
|
|
self.hashes.pop(key, None)
|
|
self.sets.pop(key, None)
|
|
|
|
def lpush(self, key, value):
|
|
self.lists.setdefault(key, []).insert(0, value)
|
|
|
|
def rpush(self, key, value):
|
|
self.lists.setdefault(key, []).append(value)
|
|
|
|
def rpop(self, key):
|
|
values = self.lists.setdefault(key, [])
|
|
return values.pop() if values else None
|
|
|
|
def lrange(self, key, start, end):
|
|
values = self.lists.get(key, [])
|
|
stop = None if end == -1 else end + 1
|
|
return values[start:stop]
|
|
|
|
def llen(self, key):
|
|
return len(self.lists.get(key, []))
|
|
|
|
def lpos(self, key, value):
|
|
try:
|
|
return self.lists.get(key, []).index(value)
|
|
except ValueError:
|
|
return None
|
|
|
|
def lrem(self, key, count, value):
|
|
values = self.lists.get(key, [])
|
|
original_len = len(values)
|
|
self.lists[key] = [item for item in values if item != value]
|
|
return original_len - len(self.lists[key])
|
|
|
|
def hset(self, key, field, value):
|
|
self.hashes.setdefault(key, {})[field] = value
|
|
|
|
def hget(self, key, field):
|
|
return self.hashes.get(key, {}).get(field)
|
|
|
|
def hgetall(self, key):
|
|
return dict(self.hashes.get(key, {}))
|
|
|
|
def hvals(self, key):
|
|
return list(self.hashes.get(key, {}).values())
|
|
|
|
def hlen(self, key):
|
|
return len(self.hashes.get(key, {}))
|
|
|
|
def hdel(self, key, field):
|
|
self.hashes.get(key, {}).pop(field, None)
|
|
|
|
def sadd(self, key, value):
|
|
self.sets.setdefault(key, set()).add(value)
|
|
|
|
def srem(self, key, value):
|
|
self.sets.setdefault(key, set()).discard(value)
|
|
|
|
def scard(self, key):
|
|
return len(self.sets.get(key, set()))
|
|
|
|
def smembers(self, key):
|
|
return set(self.sets.get(key, set()))
|
|
|
|
def scan_iter(self, match=None):
|
|
return iter(())
|
|
|
|
|
|
class FakeMonitor:
|
|
def set_worker_state(self, *args, **kwargs):
|
|
return None
|
|
|
|
def set_failure_bucket(self, *args, **kwargs):
|
|
return None
|
|
|
|
|
|
class FakeWorker:
|
|
def __init__(self, analytics, dispatcher, worker_id="fake-worker"):
|
|
self.analytics = analytics
|
|
self.dispatcher = dispatcher
|
|
self.worker_id = worker_id
|
|
|
|
def report(
|
|
self,
|
|
*,
|
|
app_name,
|
|
package_name,
|
|
status,
|
|
failure_type="",
|
|
num_nodes=12,
|
|
self_traffic_bytes=600,
|
|
server_traffic_bytes=300,
|
|
total_traffic_bytes=1000,
|
|
model_flow_count=0,
|
|
):
|
|
summary = {
|
|
"app_name": app_name,
|
|
"app_magic_label": f"magic:{package_name}",
|
|
"latest_task_key": _make_task_key(app_name, package_name),
|
|
"latest_status": status,
|
|
"latest_worker_id": self.worker_id,
|
|
"latest_failure_type": failure_type,
|
|
"latest_task_detail": failure_type,
|
|
"collection_status_reason": failure_type or "success",
|
|
"collection_task_type": "new_app",
|
|
"duration_seconds": 42,
|
|
"droidbot_steps": 30,
|
|
"gui_agent_steps": 12,
|
|
"num_nodes": num_nodes,
|
|
"total_traffic_bytes": total_traffic_bytes,
|
|
"self_traffic_bytes": self_traffic_bytes,
|
|
"server_traffic_bytes": server_traffic_bytes,
|
|
"unrecognized_traffic_bytes": max(0, total_traffic_bytes - self_traffic_bytes - server_traffic_bytes),
|
|
"model_flow_count": model_flow_count,
|
|
"model_traffic_bytes": model_flow_count * 100,
|
|
}
|
|
self.analytics.repo.replace_package_snapshot(package_name, summary, [], [], [])
|
|
persisted = self.analytics.repo.get_collection_row(package_name)
|
|
self.dispatcher._handle_analytics_snapshot(package_name, persisted)
|
|
return persisted
|
|
|
|
|
|
def make_dispatcher(redis_conn, analytics):
|
|
dispatcher = RedisTaskDispatcher.__new__(RedisTaskDispatcher)
|
|
dispatcher.redis = redis_conn
|
|
dispatcher.analytics = analytics
|
|
dispatcher.worker_inventory = {}
|
|
dispatcher.managed_worker_ids = set()
|
|
dispatcher.only_managed_workers_can_dispatch = False
|
|
dispatcher.worker_online_timeout = 300
|
|
dispatcher.task_routing_rules = {"package_name": {}, "task_key": {}}
|
|
dispatcher.monitor = FakeMonitor()
|
|
dispatcher.notifier = None
|
|
dispatcher._apk_registry = None
|
|
dispatcher._minio_storage = None
|
|
return dispatcher
|
|
|
|
|
|
def task_status(redis_conn, task_key):
|
|
payload = redis_conn.hget("task:status", task_key)
|
|
return json.loads(payload) if payload else {}
|
|
|
|
|
|
def test_repository_initializes_new_tables_without_views(tmp_path):
|
|
analytics = AnalyticsService(db_path=str(tmp_path / "analytics.sqlite3"), start_worker=False)
|
|
|
|
with analytics.repo._connect() as connection:
|
|
names = {
|
|
row["name"]
|
|
for row in connection.execute("SELECT name FROM sqlite_master WHERE type IN ('table', 'view')")
|
|
}
|
|
|
|
assert "app_catalog" in names
|
|
assert "collection_task" in names
|
|
assert "app_collect_summary" not in names
|
|
assert "v_collection_latest" not in names
|
|
|
|
analytics.replace_catalog_from_rows(
|
|
[{"app_name": "Tagged App", "package_name": "com.example.tagged"}]
|
|
)
|
|
assert analytics.set_incremental_batch_tag_for_packages(["com.example.tagged"], "batch-a") == 1
|
|
assert "batch-a" in analytics.repo.get_collection_row("com.example.tagged")["incremental_batch_tags"]
|
|
assert analytics.clear_incremental_batch_tags("batch-a") == 1
|
|
assert "batch-a" not in analytics.repo.get_collection_row("com.example.tagged")["incremental_batch_tags"]
|
|
|
|
|
|
def test_fake_worker_reports_drive_new_schema_and_redis_states(tmp_path):
|
|
analytics = AnalyticsService(db_path=str(tmp_path / "analytics.sqlite3"), start_worker=False)
|
|
analytics.replace_catalog_from_rows(
|
|
[
|
|
{"app_name": "Good App", "package_name": "com.example.good", "downloads": "100K"},
|
|
{"app_name": "Gone App", "package_name": "com.example.gone", "downloads": "1K"},
|
|
{"app_name": "Retry App", "package_name": "com.example.retry", "downloads": "2M"},
|
|
{"app_name": "Model App", "package_name": "com.example.model", "downloads": "3M"},
|
|
]
|
|
)
|
|
redis_conn = FakeRedis()
|
|
dispatcher = make_dispatcher(redis_conn, analytics)
|
|
|
|
assert dispatcher.load_tasks_from_app_summary() == 4
|
|
|
|
worker = FakeWorker(analytics, dispatcher)
|
|
good = worker.report(app_name="Good App", package_name="com.example.good", status="success")
|
|
gone = worker.report(
|
|
app_name="Gone App",
|
|
package_name="com.example.gone",
|
|
status="failed",
|
|
failure_type="APP_ERROR/3",
|
|
num_nodes=0,
|
|
total_traffic_bytes=0,
|
|
self_traffic_bytes=0,
|
|
server_traffic_bytes=0,
|
|
)
|
|
retry = worker.report(
|
|
app_name="Retry App",
|
|
package_name="com.example.retry",
|
|
status="failed",
|
|
failure_type="DOWNLOAD_ERROR/5",
|
|
num_nodes=0,
|
|
total_traffic_bytes=0,
|
|
self_traffic_bytes=0,
|
|
server_traffic_bytes=0,
|
|
)
|
|
model = worker.report(
|
|
app_name="Model App",
|
|
package_name="com.example.model",
|
|
status="success",
|
|
model_flow_count=2,
|
|
)
|
|
|
|
assert good["collection_status"] == "qualified"
|
|
assert gone["collection_status"] == "failed_terminal"
|
|
assert retry["collection_status"] == "pending"
|
|
assert model["model_eligible"] == 1
|
|
|
|
good_key = _make_task_key("Good App", "com.example.good")
|
|
gone_key = _make_task_key("Gone App", "com.example.gone")
|
|
retry_key = _make_task_key("Retry App", "com.example.retry")
|
|
|
|
assert task_status(redis_conn, good_key)["status"] == "qualified"
|
|
assert good_key in redis_conn.smembers("task:completed")
|
|
assert task_status(redis_conn, gone_key)["status"] == "failed"
|
|
assert gone_key in redis_conn.smembers("task:failed")
|
|
assert task_status(redis_conn, retry_key)["status"] == "pending"
|
|
assert retry_key in redis_conn.lrange("task:queue:default", 0, -1)
|
|
|
|
assert dispatcher.refresh_tasks_from_app_summary() == 1
|
|
assert _make_task_key("Model App", "com.example.model") + "_model" in redis_conn.lrange("task:queue:low", 0, -1)
|