feat: add account events and private message management
douyin-release-gate / verify (push) Failing after 18m54s
douyin-release-gate / verify (push) Failing after 18m54s
Add Douyin notification polling, event details, and manual multi-account private messaging. Refine environment memory settings, account operations, login collection recovery, and message UI; update tests and documentation.
This commit is contained in:
@@ -0,0 +1,83 @@
|
||||
import importlib.util
|
||||
import unittest
|
||||
from unittest.mock import Mock
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
MODULE_PATH = Path(__file__).with_name("validate_douyin_event_listener.py")
|
||||
spec = importlib.util.spec_from_file_location("listener_probe", MODULE_PATH)
|
||||
probe = importlib.util.module_from_spec(spec)
|
||||
assert spec and spec.loader
|
||||
spec.loader.exec_module(probe)
|
||||
|
||||
|
||||
class ListenerProbeTests(unittest.TestCase):
|
||||
def test_accepts_only_explicit_local_cdp_endpoints(self):
|
||||
probe.validate_arguments("http://127.0.0.1:19000", "99491952055")
|
||||
probe.validate_arguments("http://localhost:9222/", "123")
|
||||
for endpoint in (
|
||||
"https://127.0.0.1:19000",
|
||||
"http://0.0.0.0:19000",
|
||||
"http://192.168.1.5:19000",
|
||||
"http://user:pass@localhost:19000",
|
||||
"http://127.0.0.1:19000/json",
|
||||
"http://127.0.0.1:19000?token=secret",
|
||||
):
|
||||
with self.subTest(endpoint=endpoint), self.assertRaises(probe.ProbeError):
|
||||
probe.validate_arguments(endpoint, "123")
|
||||
|
||||
def test_validates_expected_uid_as_decimal_string(self):
|
||||
for uid in ("", "0", "123", "abc", "1" * 21, 123):
|
||||
with self.subTest(uid=uid), self.assertRaises(probe.ProbeError):
|
||||
probe.validate_arguments("http://127.0.0.1:19000", uid)
|
||||
|
||||
def test_classifies_readable_api_without_claiming_event_delivery(self):
|
||||
self.assertEqual(probe.classify_pages([]), "no_matching_logged_in_page")
|
||||
self.assertEqual(probe.classify_pages([{"matches_expected_uid": True}]), "notification_api_unavailable")
|
||||
groups = [{"group":group,"http_status":200,"status_code":0,"ids_valid":True,"account_identity_matches":True} for group in (700,960,961)]
|
||||
pages = [{"matches_expected_uid": True, "notice_list_read_probe": groups}]
|
||||
self.assertEqual(probe.classify_pages(pages), "notification_api_readable_delivery_unverified")
|
||||
pages[0]['polling_read_probe'] = {'outcome':'failed'}
|
||||
self.assertEqual(probe.classify_pages(pages), 'notification_polling_failed')
|
||||
pages[0]['polling_read_probe'] = {'outcome':'read_and_checkpoint_restore_verified'}
|
||||
self.assertEqual(probe.classify_pages(pages), 'notification_polling_read_verified_delivery_unverified')
|
||||
pages[0]['polling_read_probe']['reason'] = 'missing checkpoint'
|
||||
self.assertEqual(probe.classify_pages(pages), 'notification_polling_gap')
|
||||
groups[0]['account_identity_matches'] = False
|
||||
self.assertEqual(probe.classify_pages(pages), "notification_api_unavailable")
|
||||
self.assertEqual(probe.classify_pages([
|
||||
{"matches_expected_uid": True}, {"matches_expected_uid": True}
|
||||
]), "multiple_matching_pages")
|
||||
|
||||
def test_readonly_probe_acknowledges_all_chunks_before_restoring_checkpoint(self):
|
||||
session = Mock()
|
||||
session.poll.side_effect = [
|
||||
[{'delivery_id':str(i), 'kind':'notice'} for i in range(100)],
|
||||
[{'delivery_id':str(i), 'kind':'notice'} for i in range(100,151)] +
|
||||
[{'delivery_id':'checkpoint', 'kind':'checkpoint', 'checkpoints':{'700':'1'}, 'ignored_types':{}, 'reason':''}],
|
||||
]
|
||||
summary = probe.drain_readonly_session(session)
|
||||
self.assertEqual(summary['kinds'], {'notice':151, 'checkpoint':1})
|
||||
self.assertEqual(session.ack.call_count, 2)
|
||||
self.assertEqual(summary['checkpoint']['checkpoints'], {'700':'1'})
|
||||
session.poll.side_effect = None
|
||||
session.poll.return_value = []
|
||||
with self.assertRaises(probe.ProbeError):
|
||||
probe.drain_readonly_session(session)
|
||||
|
||||
def test_probe_does_not_enable_listener_or_send_messages(self):
|
||||
program = probe._harness_program("123")
|
||||
self.assertIn("did_not_install_listener", program)
|
||||
self.assertNotIn("install_expression", program)
|
||||
self.assertNotIn("/douyin/events", program)
|
||||
self.assertNotIn("sendMessage", program)
|
||||
self.assertIn("is_mark_read", program)
|
||||
self.assertNotIn("method:", program)
|
||||
self.assertIn('"notice_list_read_probe"', program)
|
||||
self.assertIn('response.text()', program)
|
||||
self.assertNotIn('response.json()', program)
|
||||
compile(program, '<harness probe>', 'exec')
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,189 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Read-only notification API probe through browser-harness; not delivery acceptance."""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
from datetime import datetime, timezone
|
||||
from collections import Counter
|
||||
from pathlib import Path
|
||||
from urllib.parse import urlsplit
|
||||
|
||||
LOCAL_HOSTS = {'127.0.0.1', 'localhost', '::1'}
|
||||
UID_RE = re.compile(r'^[1-9][0-9]{0,19}$')
|
||||
|
||||
|
||||
class ProbeError(ValueError):
|
||||
pass
|
||||
|
||||
|
||||
def validate_arguments(cdp_url: str, expected_uid: str) -> None:
|
||||
try:
|
||||
parsed = urlsplit(cdp_url)
|
||||
port = parsed.port
|
||||
except ValueError as exc:
|
||||
raise ProbeError('CDP 端口无效') from exc
|
||||
if (parsed.scheme != 'http' or parsed.hostname not in LOCAL_HOSTS or not port
|
||||
or parsed.username or parsed.password or parsed.path not in ('', '/')
|
||||
or parsed.query or parsed.fragment):
|
||||
raise ProbeError('只允许明确指定本机 http://127.0.0.1:端口 CDP 地址')
|
||||
if not isinstance(expected_uid, str) or not UID_RE.fullmatch(expected_uid):
|
||||
raise ProbeError('expected UID 必须为十进制字符串')
|
||||
|
||||
|
||||
def classify_pages(pages: list[dict[str, object]]) -> str:
|
||||
matched = [p for p in pages if p.get('matches_expected_uid')]
|
||||
if len(matched) > 1:
|
||||
return 'multiple_matching_pages'
|
||||
if not matched:
|
||||
return 'no_matching_logged_in_page'
|
||||
groups = matched[0].get('notice_list_read_probe', [])
|
||||
if (len(groups) == 3 and {g.get('group') for g in groups} == {700, 960, 961}
|
||||
and all(g.get('http_status') == 200 and g.get('status_code') == 0
|
||||
and g.get('ids_valid') and g.get('account_identity_matches') for g in groups)):
|
||||
polling = matched[0].get('polling_read_probe')
|
||||
if polling:
|
||||
if polling.get('outcome') != 'read_and_checkpoint_restore_verified':
|
||||
return 'notification_polling_failed'
|
||||
if polling.get('reason'):
|
||||
return 'notification_polling_gap'
|
||||
return 'notification_polling_read_verified_delivery_unverified'
|
||||
return 'notification_api_readable_delivery_unverified'
|
||||
return 'notification_api_unavailable'
|
||||
|
||||
|
||||
def drain_readonly_session(session):
|
||||
"""Drain the local queue, including history spanning more than one batch."""
|
||||
kinds = Counter()
|
||||
while True:
|
||||
batch = session.poll(100, 0)
|
||||
if not batch:
|
||||
raise ProbeError('通知读取没有返回确认边界,不能声称历史已读完')
|
||||
kinds.update(d['kind'] for d in batch)
|
||||
checkpoint = next((d for d in batch if d['kind'] == 'checkpoint'), None)
|
||||
session.ack([d['delivery_id'] for d in batch])
|
||||
if checkpoint is not None:
|
||||
return {'kinds':dict(kinds), 'checkpoint':checkpoint}
|
||||
|
||||
|
||||
def _harness_program(expected_uid: str) -> str:
|
||||
# Raw JSON is decoded in Python; JavaScript would round 64-bit notice IDs.
|
||||
fetch = r'''(async()=>{
|
||||
const response=await fetch(URL_VALUE,{credentials:"include",redirect:"error",cache:"no-store",signal:AbortSignal.timeout(10000)});
|
||||
return {status:response.status,body:await response.text()};
|
||||
})()'''
|
||||
return f'''import json,re,sys
|
||||
from collections import Counter
|
||||
from urllib.parse import urlsplit,urlencode
|
||||
sys.path.insert(0,{json.dumps(str(Path(__file__).resolve().parent.parent))})
|
||||
from browser_gateway.platform.douyin import BrowserResponse,DouyinBrowser
|
||||
from browser_gateway.platform.notice_polling import NoticePollingSession
|
||||
from scripts.validate_douyin_event_listener import drain_readonly_session
|
||||
expected_uid={json.dumps(expected_uid)}
|
||||
fetch_expression={json.dumps(fetch)}
|
||||
reports=[]
|
||||
for target in cdp("Target.getTargets")["targetInfos"]:
|
||||
url=target.get("url", "")
|
||||
if target.get("type")!="page" or urlsplit(url).netloc!="www.douyin.com": continue
|
||||
target_id=target["targetId"]
|
||||
report={{"page_path":urlsplit(url).path}}
|
||||
def read(path):
|
||||
response=js(fetch_expression.replace("URL_VALUE",json.dumps(path)),target_id=target_id)
|
||||
return response["status"],json.loads(response["body"])
|
||||
try:
|
||||
status,identity=read("/aweme/v1/web/user/profile/self/?device_platform=webapp&aid=6383&channel=channel_pc_web")
|
||||
uid=str((identity.get("user") or {{}}).get("uid") or "")
|
||||
report["identity_status"]="read" if status==200 and identity.get("status_code")==0 and uid else "login_required"
|
||||
report["matches_expected_uid"]=report["identity_status"]=="read" and uid==expected_uid
|
||||
if report["matches_expected_uid"]:
|
||||
groups=[]
|
||||
for group in (700,960,961):
|
||||
query=urlencode(dict(device_platform="webapp",aid=6383,channel="channel_pc_web",is_new_notice=1,is_mark_read=0,count=1,min_time=0,max_time=0,notice_group=group))
|
||||
try:
|
||||
status,body=read("/aweme/v1/web/notice/?"+query)
|
||||
items=body.get("notice_list_v2")
|
||||
valid=isinstance(items,list)
|
||||
groups.append({{"group":group,"http_status":status,"status_code":body.get("status_code"),
|
||||
"has_more":body.get("has_more"),"sample_count":len(items) if valid else None,
|
||||
"notice_types":sorted({{str(n.get("type","unknown")) for n in items}}) if valid else [],
|
||||
"ids_valid":valid and all(re.fullmatch(r"[1-9][0-9]{{0,29}}",str(n.get("nid_str") or n.get("nid") or "")) for n in items),
|
||||
"account_identity_matches":valid and all(str(n.get("user_id"))==expected_uid for n in items)}})
|
||||
except Exception as exc: groups.append({{"group":group,"error":str(exc)[:240]}})
|
||||
report["notice_list_read_probe"]=groups
|
||||
class ReadOnlyBrowser:
|
||||
identity=DouyinBrowser.identity
|
||||
def get(self,alias,url):
|
||||
raw=js(fetch_expression.replace("URL_VALUE",json.dumps(url)),target_id=target_id)
|
||||
return BrowserResponse(raw["status"],raw["body"],False)
|
||||
sessions=[]
|
||||
try:
|
||||
browser=ReadOnlyBrowser()
|
||||
session=NoticePollingSession(browser,"readonly",expected_uid)
|
||||
sessions.append(session)
|
||||
first=drain_readonly_session(session)
|
||||
restored=NoticePollingSession(browser,"readonly",expected_uid,session.boundary_at,session.checkpoints)
|
||||
sessions.append(restored)
|
||||
second=drain_readonly_session(restored)
|
||||
report["polling_read_probe"]={{"outcome":"read_and_checkpoint_restore_verified",
|
||||
"first_kinds":first["kinds"], "restored_kinds":second["kinds"],
|
||||
"checkpoint_groups":list(second["checkpoint"]["checkpoints"]),
|
||||
"ignored_types":first["checkpoint"]["ignored_types"],
|
||||
"reason":second["checkpoint"]["reason"]}}
|
||||
except Exception as exc:
|
||||
report["polling_read_probe"]={{"outcome":"failed","error":str(exc)[:240]}}
|
||||
finally:
|
||||
for session in sessions: session.stop()
|
||||
except Exception as exc:
|
||||
report.update(identity_status="probe_error",error=str(exc)[:240])
|
||||
reports.append(report)
|
||||
print("CREATORHUB_LISTENER_PROBE="+json.dumps({{"expected_uid":expected_uid,"pages":reports,
|
||||
"read_only":True,"did_not_install_listener":True}},ensure_ascii=False))
|
||||
'''
|
||||
|
||||
|
||||
def run_probe(cdp_url: str, expected_uid: str, harness_name: str) -> dict[str, object]:
|
||||
validate_arguments(cdp_url, expected_uid)
|
||||
executable = shutil.which('browser-harness')
|
||||
if not executable:
|
||||
raise ProbeError('找不到 browser-harness;请按 docs/research/douyin-event-listener.md 安装')
|
||||
env = os.environ.copy()
|
||||
env.update(BU_CDP_URL=cdp_url, BU_NAME=harness_name, BH_RECORD='0')
|
||||
try:
|
||||
process = subprocess.run([executable], input=_harness_program(expected_uid), text=True,
|
||||
capture_output=True, check=False, timeout=90, env=env)
|
||||
except subprocess.TimeoutExpired as exc:
|
||||
raise ProbeError('browser-harness 超时;未切换页面或启动监听') from exc
|
||||
if process.returncode:
|
||||
detail = process.stderr.strip().splitlines()
|
||||
raise ProbeError(detail[-1][:400] if detail else 'browser-harness 连接失败')
|
||||
for line in reversed(process.stdout.splitlines()):
|
||||
if line.startswith('CREATORHUB_LISTENER_PROBE='):
|
||||
report = json.loads(line.split('=', 1)[1])
|
||||
report['outcome'] = classify_pages(report['pages'])
|
||||
report['checked_at'] = datetime.now(timezone.utc).isoformat()
|
||||
return report
|
||||
raise ProbeError('browser-harness 未返回诊断结果')
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument('--cdp-url', required=True, help='该账号浏览器的本机 CDP 地址')
|
||||
parser.add_argument('--expected-uid', required=True, help='该浏览器绑定的抖音 UID')
|
||||
parser.add_argument('--harness-name', default='creatorhub-event-listener-readonly')
|
||||
args = parser.parse_args(argv)
|
||||
try:
|
||||
report = run_probe(args.cdp_url, args.expected_uid, args.harness_name)
|
||||
except (ProbeError, json.JSONDecodeError) as exc:
|
||||
print(json.dumps({'outcome':'probe_failed','error':str(exc)},ensure_ascii=False),file=sys.stderr)
|
||||
return 2
|
||||
print(json.dumps(report,ensure_ascii=False,indent=2))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
raise SystemExit(main())
|
||||
Reference in New Issue
Block a user