1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143
| import asyncio import json import base64
class CDPServiceWorkerManager: """CDP Service Worker 管理器"""
def __init__(self, ws): self.ws = ws self._cmd_id = 0 self.versions = {} self.registrations = {} self.error_log = []
async def _cmd(self, method, params=None): """发送 CDP 命令(浏览器级别)""" self._cmd_id += 1 msg = {"id": self._cmd_id, "method": method, "params": params or {}} await self.ws.send(json.dumps(msg)) async for resp in self.ws: data = json.loads(resp) if data.get("id") == self._cmd_id: return data.get("result", {})
async def _cmd_session(self, method, params, session_id): """发送 CDP 命令(页面级别,带 sessionId)""" self._cmd_id += 1 msg = { "id": self._cmd_id, "method": method, "params": params or {}, "sessionId": session_id } await self.ws.send(json.dumps(msg)) async for resp in self.ws: data = json.loads(resp) if data.get("id") == self._cmd_id: return data.get("result", {})
async def enable(self): """启用 Service Worker 监听""" return await self._cmd("ServiceWorker.enable")
async def disable(self): """禁用 Service Worker 监听""" return await self._cmd("ServiceWorker.disable")
async def skip_waiting(self, scope_url=""): """跳过等待,激活新版本""" return await self._cmd("ServiceWorker.skipWaiting", { "scopeURL": scope_url })
async def unregister(self, scope_url): """注销 Service Worker""" return await self._cmd("ServiceWorker.unregister", { "scopeURL": scope_url })
async def dispatch_sync(self, tag, last_chance=False): """触发后台同步事件""" return await self._cmd("ServiceWorker.dispatchSyncEvent", { "tag": tag, "lastChance": last_chance })
async def dispatch_periodic_sync(self, tag): """触发定期后台同步事件""" return await self._cmd("ServiceWorker.dispatchPeriodicSyncEvent", { "tag": tag })
async def inspect_worker(self, version_id): """获取 Worker 调试 URL""" return await self._cmd("ServiceWorker.inspectWorker", { "versionId": version_id })
async def list_caches(self, session_id, origin=""): """列出缓存(需要 sessionId)""" return await self._cmd_session( "CacheStorage.requestCacheNames", {"securityOrigin": origin}, session_id )
async def read_cache_entries(self, session_id, cache_id, skip=0, page_size=50): """读取缓存条目(需要 sessionId)""" return await self._cmd_session( "CacheStorage.requestEntries", {"cacheId": cache_id, "skipCount": skip, "pageSize": page_size}, session_id )
async def delete_cache(self, session_id, cache_name): """删除缓存(需要 sessionId)""" return await self._cmd_session( "CacheStorage.deleteCache", {"cacheName": cache_name}, session_id )
async def delete_cache_entry(self, session_id, cache_id, request_url): """删除缓存条目(需要 sessionId)""" return await self._cmd_session( "CacheStorage.deleteEntry", {"cacheId": cache_id, "request": request_url}, session_id )
async def collect_events(self, duration=30): """收集一段时间内的 Service Worker 事件""" await self.enable() events = [] start = asyncio.get_event_loop().time()
while asyncio.get_event_loop().time() - start < duration: try: resp = await asyncio.wait_for( self.ws.recv(), timeout=1.0 ) data = json.loads(resp) method = data.get("method", "")
if method.startswith("ServiceWorker."): events.append({ "method": method, "params": data.get("params", {}), "time": asyncio.get_event_loop().time() })
if method == "ServiceWorker.onWorkerVersionUpdated": for v in data["params"]["versions"]: self.versions[v["versionId"]] = v
elif method == "ServiceWorker.onWorkerRegistrationUpdated": for r in data["params"]["registrations"]: self.registrations[r["registrationId"]] = r
elif method == "ServiceWorker.onWorkerErrorReported": self.error_log.append( data["params"]["errorMessage"] ) except asyncio.TimeoutError: pass
return events
|