feat: add push subscription mode for interaction notifications

This commit is contained in:
2026-09-06 01:31:15 +08:00
parent ce3e3869e9
commit 9ccb5d659c
5 changed files with 614 additions and 3 deletions
+18 -3
View File
@@ -1,5 +1,5 @@
#!/usr/bin/env python3
"""自动翻页获取互动通知,遇到 last-id 停止;不指定则获取全部,不标记已读。"""
"""拉取互动通知,或用 --sub 订阅新通知;不标记已读,不含私信。"""
import argparse
import json
@@ -132,10 +132,16 @@ def collect_notifications(
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__, allow_abbrev=False)
parser.add_argument(
"--last-id", help="从最新开始,遇到此通知 ID 停止且不包含该条;未指定则获取全部"
"--sub", action="store_true", help="订阅互动通知推送;不指定则保持原拉取行为"
)
parser.add_argument(
"--output", type=Path, default=Path(__file__).with_name("notifications.json")
"--last-id",
help="历史拉取截至此 ID(不含);--sub 时用于一次启动补拉,不指定则只接收新推送",
)
parser.add_argument(
"--output",
type=Path,
help="拉取默认保存 src/notifications.json;订阅时可选追加 JSONL 文件(默认仅 stdout)",
)
args = parser.parse_args()
if args.last_id is not None and (
@@ -143,6 +149,12 @@ def main() -> int:
):
parser.error("--last-id 必须为非空的数字通知 ID")
try:
if args.sub:
from subscribe_notifications import subscribe
subscribe(output=args.output, last_id=args.last_id)
return 0
args.output = args.output or Path(__file__).with_name("notifications.json")
# 复用已有登录态校验;失败就退出,不尝试登录。
user = get_user_from_browser()
notices, raw_pages, last_id_found = collect_notifications(
@@ -168,6 +180,9 @@ def main() -> int:
json.dumps(result, ensure_ascii=False, indent=2) + "\n", encoding="utf-8"
)
temporary.replace(args.output)
except KeyboardInterrupt:
print("已停止。", file=sys.stderr)
return 130
except (
OSError,
ValueError,
+255
View File
@@ -0,0 +1,255 @@
"""互动通知推送桥接;入口:python3 src/get_notifications.py --sub。"""
import json
import subprocess
import sys
import time
import uuid
from collections import OrderedDict
from contextlib import suppress
from datetime import datetime, timezone
from urllib.parse import urlencode
from get_current_user import get_user_from_browser
# 不改写站点 onmessage、不新建连接;退出时仅卸载自己的监听器。
INSTALL = r"""(() => {
if (location.origin !== 'https://www.douyin.com') throw Error('请选择已登录的抖音标签页');
const key = KEY, uid = UID, chunks = window.webpackChunkdouyin_web;
if (!chunks) throw Error('未找到抖音运行时');
let require;
chunks.push([['notice-sub-' + key + '-' + Date.now()], {}, r => { require = r; }]);
chunks.pop();
const entries = Object.entries(require.m);
const entry = entries.find(([, f]) => String(f).includes('NOTICE_PUSH_EVENT_NAMES:function'));
const codec = entries.find(([, f]) => {
const s = String(f); return s.includes('.decodedFrame=') && s.includes('.encodeFrame=');
});
if (!entry || !codec) throw Error('通知 SDK 已变化,无法订阅');
const C = require(entry[0]).NoticeFrontier;
const decode = require(codec[0]).decodedFrame;
const f = C.frontierInstance;
if (!f || String(f._options.deviceID) !== uid) throw Error('通知连接未就绪或登录账号已变化');
const state = {queue: [], error: null, wake: null, C, f, uid};
const emit = event => {
if (state.queue.length >= 1000) state.error = '订阅队列已满,可能丢失事件';
else state.queue.push(event);
if (state.wake) state.wake();
};
const message = event => {
try {
const frame = decode(new Uint8Array(event.data));
if (frame.service === 20313 || frame.service === 20003)
emit({kind: 'push', service: frame.service, payload: new TextDecoder().decode(frame.payload)});
} catch (_) { state.error = '通知帧解码失败'; if (state.wake) state.wake(); }
};
const open = () => emit({kind: 'open'});
const close = () => emit({kind: 'close'});
f.addEventListener('message', message);
f.addEventListener('open', open);
f.addEventListener('close', close);
state.dispose = () => {
f.removeEventListener('message', message);
f.removeEventListener('open', open);
f.removeEventListener('close', close);
if (state.wake) state.wake();
};
window[key] = state;
return JSON.stringify({connected: f.readyState === f.OPEN});
})()"""
WAIT = r"""(async () => {
const s = window[KEY];
if (!s || s.C.frontierInstance !== s.f || String(s.f._options.deviceID) !== s.uid)
throw Error('页面刷新、登录失效或账号变化,请确认登录后重新运行');
if (!s.queue.length && !s.error) await new Promise(resolve => {
// 本地桥接心跳,不请求通知接口;小于 harness 默认 IPC 超时。
const timer = setTimeout(done, 2000);
function done() { clearTimeout(timer); s.wake = null; resolve(); }
s.wake = done;
});
if (s.error) throw Error(s.error);
if (s.C.frontierInstance !== s.f) throw Error('通知连接已销毁,请检查登录态');
return JSON.stringify(s.queue.splice(0));
})()"""
def browser(expression):
expression = (
"(async()=>{try{return await ("
+ expression
+ ");}catch(e){return JSON.stringify({bridge_error:String(e.message)});}})()"
)
result = subprocess.run(
["browser-harness"],
input=f"print(js({expression!r}))\n",
text=True,
capture_output=True,
timeout=30,
check=False,
)
if result.returncode:
# 不回显包含账号等信息的完整表达式。
raise RuntimeError(
"浏览器订阅操作失败:页面刷新、登录失效、SDK 变化或连接超时;请确认后重启"
)
try:
value = json.loads(result.stdout.strip().splitlines()[-1])
except (ValueError, IndexError) as exc:
raise RuntimeError("浏览器返回内容无法解析") from exc
if isinstance(value, dict) and "bridge_error" in value:
raise RuntimeError(value["bridge_error"])
return value
def notice_ids(event):
"""只处理互动通知服务,保留字符串 ID,避免 64 位精度损失。"""
try:
payload = json.loads(event["payload"])
except (ValueError, TypeError) as exc:
raise RuntimeError("通知推送内容无法解析") from exc
if not isinstance(payload, dict):
raise TypeError("通知推送格式异常")
if event["service"] == 20313:
notices = payload.get("notices", [])
if not isinstance(notices, list) or any(
not isinstance(n, dict) or not isinstance(n.get("effect_groups", []), list)
for n in notices
):
raise RuntimeError("通知推送列表格式异常")
items = [
n
for n in notices
if {str(g) for g in n.get("effect_groups", [])} & {"960", "961"}
]
elif event["service"] == 20003 and payload.get("notice_type") in (
45,
31,
9009,
9002,
514,
9067,
):
items = [payload]
else:
return []
ids = [n.get("notice_id_str") for n in items]
if any(not isinstance(n, str) or not n.isascii() or not n.isdecimal() for n in ids):
raise RuntimeError("推送通知 ID 格式异常")
return list(dict.fromkeys(ids))
def details(ids, uid):
url = "/aweme/v1/web/notice/detail/?" + urlencode(
{
"device_platform": "webapp",
"aid": 6383,
"is_mark_read": 0,
"id_list": json.dumps([{"notice_id_str": n, "type": 0} for n in ids]),
}
)
expression = """(async () => {
if (location.origin !== 'https://www.douyin.com') throw Error('抖音页面已关闭');
const r = await fetch(URL, {credentials: 'include', signal: AbortSignal.timeout(20000)});
return JSON.stringify({status: r.status, body: await r.text()});
})()""".replace("URL", json.dumps(url))
# 推送可能先于详情入库;仅对这次事件重试,不做定时通知扫描。
for attempt in range(3):
response = browser(expression)
try:
payload = json.loads(response["body"])
except (ValueError, TypeError) as exc:
raise RuntimeError("通知详情无法解析,请检查登录态") from exc
if not isinstance(payload, dict):
raise TypeError("通知详情格式异常")
if response["status"] != 200 or payload.get("status_code") != 0:
raise RuntimeError("通知详情请求失败,请检查登录态;不会自动登录")
notices = payload.get("notice_list_v2")
if not isinstance(notices, list) or any(
not isinstance(n, dict) for n in notices
):
raise RuntimeError("通知详情格式异常")
if any(str(n.get("user_id")) != str(uid) for n in notices):
raise RuntimeError("账号变化,停止订阅")
found = {n.get("nid_str") or str(n.get("nid")) for n in notices}
if found == set(ids):
return notices
if attempt < 2:
time.sleep(0.5 * (attempt + 1))
raise RuntimeError("推送详情未完整返回,未处理通知 ID:" + ",".join(ids))
def subscribe(output=None, last_id=None):
from get_notifications import collect_notifications, summarize
user = get_user_from_browser()
key = "__douyin_notice_sub_" + uuid.uuid4().hex
seen = OrderedDict()
stream = None
def emit(notice):
item = summarize(notice)
# 同一 ID 的聚合内容更新也输出;两路推送的相同详情只处理一次。
signature = json.dumps(notice, sort_keys=True, ensure_ascii=False)
if seen.get(item["id"]) == signature:
return
record = {
"received_at": datetime.now(timezone.utc).isoformat(),
"notification": item,
"raw": notice,
}
line = json.dumps(record, ensure_ascii=False)
if stream:
stream.write(line + "\n")
stream.flush()
print(line, flush=True)
seen[item["id"]] = signature
seen.move_to_end(item["id"])
# ponytail: 仅内存保留最近 4096 个 ID;跨重启幂等由消费方持久化处理。
if len(seen) > 4096:
seen.popitem(last=False)
try:
if output:
stream = output.open("a", encoding="utf-8")
state = browser(
INSTALL.replace("KEY", json.dumps(key)).replace(
"UID", json.dumps(str(user["uid"]))
)
)
print(
"订阅已安装,连接"
+ ("已建立" if state["connected"] else "建立中")
+ ";Ctrl+C 退出。stdout 为通知 JSONL,不自动标记已读。",
file=sys.stderr,
)
if last_id:
notices, _, found = collect_notifications(user["uid"], last_id)
for notice in reversed(notices):
emit(notice)
if not found:
print("未找到 last-id,已补拉接口可获取的全部通知。", file=sys.stderr)
while True:
for event in browser(WAIT.replace("KEY", json.dumps(key))):
if event["kind"] != "push":
print(
"通知连接事件:"
+ event["kind"]
+ ";断线期间可能漏消息,可用 --last-id 补拉。",
file=sys.stderr,
)
continue
ids = notice_ids(event)
for offset in range(0, len(ids), 3):
for notice in details(ids[offset : offset + 3], user["uid"]):
emit(notice)
finally:
# 页面已关闭时监听器也已销毁,不关闭站点自己的连接。
with suppress(OSError, ValueError, RuntimeError, subprocess.SubprocessError):
browser(
"(() => {const k="
+ json.dumps(key)
+ '; if(window[k]) {window[k].dispose(); delete window[k];} return "null";})()'
)
if stream:
stream.close()
+211
View File
@@ -0,0 +1,211 @@
"""离线检查(Python 标准库 + Node.js):python3 src/test_subscribe_notifications.py。"""
import io
import json
import subprocess
from contextlib import redirect_stdout
from pathlib import Path
from tempfile import TemporaryDirectory
from unittest.mock import patch
import subscribe_notifications as sub
from get_notifications import main
def push(service=20313, **payload):
return {"kind": "push", "service": service, "payload": json.dumps(payload)}
def check_raises(call):
try:
call()
except (RuntimeError, TypeError):
return
raise AssertionError("必须拒绝异常数据")
def test_ids():
nid = "7680568855210951729"
item = {"notice_id_str": nid, "effect_groups": [960]}
assert sub.notice_ids(push(notices=[item, item])) == [nid]
assert sub.notice_ids(push(notices=[{**item, "effect_groups": ["961"]}])) == [nid]
assert sub.notice_ids(push(notices=[{**item, "effect_groups": [100]}])) == []
assert sub.notice_ids(push(20003, notice_type=45, notice_id_str=nid)) == [nid]
assert sub.notice_ids(push(20003, notice_type=999, notice_id_str=nid)) == []
assert sub.notice_ids(push(999)) == []
for invalid in [int(nid), None, "12", "abc"]:
check_raises(
lambda invalid=invalid: sub.notice_ids(
push(notices=[{**item, "notice_id_str": invalid}])
)
)
check_raises(lambda: sub.notice_ids({"payload": "invalid"}))
check_raises(lambda: sub.notice_ids({"payload": "[]"}))
check_raises(lambda: sub.notice_ids(push(notices=[None])))
check_raises(lambda: sub.notice_ids(push(notices=[{"effect_groups": None}])))
def test_details():
notice = {"nid_str": "42", "user_id": "u"}
good = {
"status": 200,
"body": json.dumps({"status_code": 0, "notice_list_v2": [notice]}),
}
empty = {
"status": 200,
"body": json.dumps({"status_code": 0, "notice_list_v2": []}),
}
with (
patch.object(sub, "browser", side_effect=[empty, good]) as call,
patch.object(sub.time, "sleep"),
):
assert sub.details(["42"], "u") == [notice]
assert call.call_count == 2
expression = call.call_args.args[0]
assert "/notice/detail/" in expression and "is_mark_read=0" in expression
assert "id_list=" in expression and "/notice/?" not in expression
with patch.object(sub, "browser", return_value=good):
check_raises(lambda: sub.details(["42"], "other"))
with (
patch.object(sub, "browser", return_value=empty) as call,
patch.object(sub.time, "sleep"),
):
check_raises(lambda: sub.details(["42"], "u"))
assert call.call_count == 3
with patch.object(sub, "browser", return_value={"status": 403, "body": "{}"}):
check_raises(lambda: sub.details(["42"], "u"))
with patch.object(sub, "browser", return_value={"status": 200, "body": "[]"}):
check_raises(lambda: sub.details(["42"], "u"))
def test_subscription_and_cli():
event = push(notices=[{"notice_id_str": "42", "effect_groups": [960]}])
notice = {"nid_str": "42", "user_id": "u", "create_time": 1788271791}
calls = []
waits = iter([[], [event, event], [event], KeyboardInterrupt()])
def browser(expression):
calls.append(expression)
if "const entries" in expression:
return {"connected": True}
if "s.queue.splice" in expression:
value = next(waits)
if isinstance(value, BaseException):
raise value
return value
return None
with TemporaryDirectory() as folder:
output = Path(folder) / "events.jsonl"
output.write_text('{"existing":true}\n', encoding="utf-8")
with (
patch.object(sub, "browser", side_effect=browser),
patch.object(sub, "get_user_from_browser", return_value={"uid": "u"}),
patch.object(
sub,
"details",
side_effect=[[notice], [notice], [{**notice, "type": 2}]],
),
patch("sys.argv", ["get_notifications", "--sub", "--output", str(output)]),
patch("get_notifications.collect_notifications") as collect,
redirect_stdout(io.StringIO()) as stdout,
):
assert main() == 130
collect.assert_not_called()
lines = [json.loads(line) for line in stdout.getvalue().splitlines()]
assert len(lines) == 2 # 同内容重复丢弃,同 ID 内容变化保留。
assert lines[0]["notification"]["id"] == "42"
assert len(output.read_text(encoding="utf-8").splitlines()) == 3
assert ".dispose()" in calls[-1]
# 不带 --sub 仍走原分页拉取和 JSON 快照保存路径。
with TemporaryDirectory() as folder:
output = Path(folder) / "snapshot.json"
with (
patch("sys.argv", ["get_notifications", "--output", str(output)]),
patch("get_notifications.get_user_from_browser", return_value={"uid": "u"}),
patch(
"get_notifications.collect_notifications",
return_value=([notice], [{"has_more": 0}], False),
),
patch.object(sub, "subscribe") as subscribe,
redirect_stdout(io.StringIO()),
):
assert main() == 0
subscribe.assert_not_called()
assert json.loads(output.read_text(encoding="utf-8"))["count"] == 1
with (
patch.object(sub, "subscribe") as subscribe,
patch("sys.argv", ["get_notifications", "--sub"]),
):
assert main() == 0
subscribe.assert_called_once_with(output=None, last_id=None)
with (
patch.object(sub, "subscribe") as subscribe,
patch("sys.argv", ["get_notifications", "--sub", "--last-id", "42"]),
):
assert main() == 0
subscribe.assert_called_once_with(output=None, last_id="42")
def test_js_bridge():
# 假 Webpack/FWS;只测自有监听器,不向真实账号发送通知。
script = r"""
const assert = require('node:assert/strict');
const input = JSON.parse(require('node:fs').readFileSync(0, 'utf8'));
const listeners = {};
const f = {_options:{deviceID:'u'}, readyState:1, OPEN:1,
addEventListener(type, callback) { (listeners[type] ||= new Set()).add(callback); },
removeEventListener(type, callback) { listeners[type].delete(callback); }
};
const C = {frontierInstance:f};
const decode = data => {const frame=JSON.parse(new TextDecoder().decode(data)); frame.payload=new Uint8Array(frame.payload); return frame;};
const r = id => id === '1' ? {NoticeFrontier:C} : {decodedFrame:decode};
r.m = {'1':function(){/* NOTICE_PUSH_EVENT_NAMES:function */},
'2':function(){/* exports.decodedFrame= exports.encodeFrame= */}};
global.window = {webpackChunkdouyin_web:[]};
window.webpackChunkdouyin_web.push = function(chunk) { chunk[2](r); return Array.prototype.push.call(this, chunk); };
global.location = {origin:'https://www.douyin.com'};
(async () => {
assert.equal(JSON.parse(eval(input.install)).connected, true);
const send = frame => {for(const callback of listeners.message) callback({data:new TextEncoder().encode(JSON.stringify(frame))});};
send({service:999,payload:[]});
assert.equal(window.probe.queue.length,0);
const text = JSON.stringify({notices:[{notice_id_str:'7680568855210951729',effect_groups:[960]}],text:'测试'});
send({service:20313,payload:[...new TextEncoder().encode(text)]});
const events = JSON.parse(await eval(input.wait));
assert.equal(events.length,1);
assert.equal(events[0].payload,text);
const start=Date.now();
const pending=eval(input.wait);
setTimeout(()=>send({service:20003,payload:[...new TextEncoder().encode('{}')]}),20);
assert.equal(JSON.parse(await pending)[0].service,20003);
assert.ok(Date.now()-start<1000); // 收到事件立即唤醒,不等本地心跳。
window.probe.dispose();
assert.equal(listeners.message.size,0);
assert.equal(listeners.open.size,0);
assert.equal(listeners.close.size,0);
assert.equal(f.readyState,1);
C.frontierInstance=null;
await assert.rejects(eval(input.wait),/登录失效/);
console.log('JS bridge checks passed');
})().catch(e=>{console.error(e);process.exitCode=1;});
"""
data = {
"install": sub.INSTALL.replace("KEY", '"probe"').replace("UID", '"u"'),
"wait": sub.WAIT.replace("KEY", '"probe"'),
}
subprocess.run(
["node", "-e", script],
input=json.dumps(data),
text=True,
check=True,
timeout=15,
)
if __name__ == "__main__":
test_ids()
test_details()
test_subscription_and_cli()
test_js_bridge()
print("subscription checks passed")