From ac8935c641e2b3566509942f078312b4653b7e41 Mon Sep 17 00:00:00 2001 From: zzhyyyyy <2685922758@qq.com> Date: Fri, 17 Jul 2026 22:32:07 +0800 Subject: [PATCH] =?UTF-8?q?refactor(wx-finance):=20=E4=BC=9A=E8=AF=9D?= =?UTF-8?q?=E5=AD=98=E6=A1=A3=E8=BD=AE=E8=AF=A2=E6=8B=86=E7=8B=AC=E7=AB=8B?= =?UTF-8?q?=E8=BF=9B=E7=A8=8B(Go=20.so=20=E5=B5=8C=20app-server=20?= =?UTF-8?q?=E4=BC=9A=20core-dump)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WeWorkFinanceSdk 是 Go c-shared 库、自带 Go runtime, 嵌进 app-server(uvicorn)后台线程跑会 segfault 把主进程带崩(2026-07-17 实测:开 worker 后 app-server core-dump 每 5s 循环)。改成 轮询搬到【独立进程】(Go runtime 独占进程+主线程, 崩了 systemd 只重启本进程、不动 app-server), 拉到卡片/图后 POST app-server 内部端点打信号(set_pending 在 app-server 进程内存, 心跳才 pop 得到)。 - 删 app/core/wx_finance_worker.py(进程内 worker), main lifespan 不再起它 - 新增内部端点 POST /internal/wx-finance/pending(X-Internal-Secret)→ set_pending(复用现有内部端点范式) - 新增 scripts/wx_finance_poller.py(独立 poller:GetChatData 轮询→解密→weapp/image→POST 内部端点) - 新增 deploy/wx-finance-poller.service(独立 systemd 单元, 与 app-server 隔离) - config 加 WX_FINANCE_INTERNAL_URL(poller→app-server 内部地址, 默认 127.0.0.1:8770) - SDK 封装 wx_finance_sdk.py 原样复用(探针已验证 Init/GetChatData 通) Co-Authored-By: Claude Opus 4.8 --- app/api/internal/wx_finance.py | 58 ++++++++++++++ app/core/config.py | 1 + app/core/wx_finance_worker.py | 116 ---------------------------- app/main.py | 10 +-- app/schemas/wx_finance.py | 13 ++++ deploy/wx-finance-poller.service | 18 +++++ scripts/wx_finance_poller.py | 128 +++++++++++++++++++++++++++++++ 7 files changed, 221 insertions(+), 123 deletions(-) create mode 100644 app/api/internal/wx_finance.py delete mode 100644 app/core/wx_finance_worker.py create mode 100644 app/schemas/wx_finance.py create mode 100644 deploy/wx-finance-poller.service create mode 100644 scripts/wx_finance_poller.py diff --git a/app/api/internal/wx_finance.py b/app/api/internal/wx_finance.py new file mode 100644 index 0000000..baa0da5 --- /dev/null +++ b/app/api/internal/wx_finance.py @@ -0,0 +1,58 @@ +"""会话存档独立 poller → app-server 的内部信号端点(server→server, 非客户端接口)。 + +会话存档轮询在【独立进程】跑(WeWorkFinanceSdk 是 Go c-shared 库, 嵌进 app-server 会 segfault +把主进程带崩, 见 scripts/wx_finance_poller.py)。poller 拉到美团卡片/截图后 POST 到这里, 由 +app-server 打比价信号 —— set_pending 在 app-server 进程内存, 心跳才 pop 得到, 故必须回到本进程。 +靠共享密钥头 X-Internal-Secret(== settings.INTERNAL_API_SECRET)校验;未配置 → 503。 +""" +from __future__ import annotations + +import hmac +import logging +from typing import Annotated + +from fastapi import APIRouter, Header, HTTPException, status + +from app.core import wx_poc_signal +from app.core.config import settings +from app.schemas.wx_finance import WxFinancePendingIn, WxFinancePendingOut + +logger = logging.getLogger("shagua.internal.wx_finance") + +router = APIRouter(prefix="/internal", tags=["internal"]) + + +def _check_secret(x_internal_secret: str | None) -> None: + """共享密钥校验。未配置 → 503(挡裸奔);不匹配 → 401(常量时间比较)。""" + configured = settings.INTERNAL_API_SECRET + if not configured: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail="internal api not configured", + ) + if not x_internal_secret or not hmac.compare_digest(x_internal_secret, configured): + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, detail="invalid internal secret" + ) + + +@router.post( + "/wx-finance/pending", + response_model=WxFinancePendingOut, + summary="会话存档 poller 触发比价信号(独立进程→app-server, 打 set_pending)", +) +def wx_finance_pending( + payload: WxFinancePendingIn, + x_internal_secret: Annotated[str | None, Header()] = None, +) -> WxFinancePendingOut: + _check_secret(x_internal_secret) + device_id = settings.WX_POC_TEST_DEVICE_ID + if not device_id: + logger.warning("wx_finance pending: WX_POC_TEST_DEVICE_ID 未配置, 忽略信号") + return WxFinancePendingOut(ok=False) + wx_poc_signal.set_pending(device_id, payload.source) + logger.info( + "wx_finance PoC: 已给测试设备 %s 打比价信号 source=%s (触发=%s)", + device_id, payload.source, payload.kind, + ) + return WxFinancePendingOut(ok=True) diff --git a/app/core/config.py b/app/core/config.py index 3c583ff..3bd3446 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -168,6 +168,7 @@ class Settings(BaseSettings): WX_FINANCE_POLL_INTERVAL_SEC: int = 3 # GetChatData 轮询间隔(无回调, 自己定频) WX_FINANCE_SEQ_FILE: str = "./data/wx_finance_seq.txt" # seq 游标持久化(防丢/重复) WX_FINANCE_RECEIVER_USERID: str = "" # 接收成员 userid(万朗杰);只处理【非他发】的消息, 空=不过滤 + WX_FINANCE_INTERNAL_URL: str = "http://127.0.0.1:8770" # poller 通知 app-server 打信号的内部地址(同机) @property def wx_mp_configured(self) -> bool: diff --git a/app/core/wx_finance_worker.py b/app/core/wx_finance_worker.py deleted file mode 100644 index 7909bfa..0000000 --- a/app/core/wx_finance_worker.py +++ /dev/null @@ -1,116 +0,0 @@ -"""微信会话内容存档轮询 worker:定时 GetChatData → 解密 → 识别【外部用户发给成员】的 weapp/image -→ 打比价信号(复用 wx_poc_signal)。下游(心跳 → 弹窗 → 比价)完全复用, 与 wx_kf / 服务号版一致。 - -在【独立线程】跑:GetChatData 是阻塞式 C 调用, 不放进 asyncio 事件循环。gate 在 -settings.wx_finance_configured(总开关 + corpid + secret)。seq 游标持久化到文件(防重启重拉 / 丢消息)。 -""" -from __future__ import annotations - -import logging -import threading -from pathlib import Path - -from app.core import wx_poc_signal -from app.core.config import settings -from app.integrations.wx_finance_sdk import WxFinanceError, WxFinanceSdk - -logger = logging.getLogger("shagua.wx_finance") - -_stop = threading.Event() -_thread: threading.Thread | None = None - - -def _load_seq() -> int: - try: - return int(Path(settings.WX_FINANCE_SEQ_FILE).read_text().strip() or "0") - except Exception: - return 0 - - -def _save_seq(seq: int) -> None: - try: - p = Path(settings.WX_FINANCE_SEQ_FILE) - p.parent.mkdir(parents=True, exist_ok=True) - p.write_text(str(seq)) - except Exception: - logger.exception("wx_finance seq 持久化失败") - - -def _handle_msg(msg: dict) -> None: - """明文消息:外部用户(from≠接收成员)发来的 weapp(美团卡片)/ image(截图)→ 打比价信号。 - - 只处理"发给成员"的:from==接收成员 = 成员自己发的, 跳过(避免成员回消息也触发)。 - 源平台/订单识别留后续(可据 weapp.username=appid 判平台), PoC 写死 meituan。 - """ - msgtype = msg.get("msgtype") or "" - frm = msg.get("from") or "" - receiver = settings.WX_FINANCE_RECEIVER_USERID - if receiver and frm == receiver: - return # 成员自己发的, 不触发 - if msgtype not in ("weapp", "image"): - return # 只关心美团卡片 / 截图 - logger.info("wx_finance 命中触发 from=%s type=%s", frm, msgtype) - poc_dev = settings.WX_POC_TEST_DEVICE_ID - if poc_dev: - wx_poc_signal.set_pending(poc_dev, "meituan") - logger.info( - "wx_finance PoC: 已给测试设备 %s 打比价信号 source=meituan (触发=%s)", poc_dev, msgtype - ) - - -def _run() -> None: - try: - sdk = WxFinanceSdk( - settings.WX_FINANCE_SDK_PATH, - settings.WX_KF_CORP_ID, - settings.WX_FINANCE_SECRET, - Path(settings.WX_FINANCE_PRIVATE_KEY_PATH).read_bytes(), - ) - except Exception: - logger.exception( - "wx_finance SDK 初始化失败, worker 退出(检查 .so / corpid / 存档Secret / 私钥 / 可信IP)" - ) - return - seq = _load_seq() - logger.info("wx_finance worker 启动, 从 seq=%d 开始轮询, 间隔 %ds", seq, settings.WX_FINANCE_POLL_INTERVAL_SEC) - while not _stop.is_set(): - try: - msgs = sdk.get_chat_data(seq, limit=1000) - if msgs: - logger.info("wx_finance 拉到 %d 条(seq>%d)", len(msgs), seq) - for item in msgs: - try: - plain = sdk.decrypt(item["encrypt_random_key"], item["encrypt_chat_msg"]) - _handle_msg(plain) - except Exception: - logger.exception("wx_finance 解密/处理单条失败 seq=%s", item.get("seq")) - seq = max(seq, int(item.get("seq", seq))) - if msgs: - _save_seq(seq) - except WxFinanceError: - logger.exception("wx_finance GetChatData 出错") - except Exception: - logger.exception("wx_finance 轮询异常") - _stop.wait(settings.WX_FINANCE_POLL_INTERVAL_SEC) - try: - sdk.close() - except Exception: - pass - logger.info("wx_finance worker 已停止") - - -def start_wx_finance_worker() -> None: - """gate 在 wx_finance_configured;未开启则直接返回(零影响)。""" - global _thread - if not settings.wx_finance_configured: - return - _stop.clear() - _thread = threading.Thread(target=_run, name="wx-finance-worker", daemon=True) - _thread.start() - logger.info("wx_finance worker 线程已启动") - - -def stop_wx_finance_worker() -> None: - _stop.set() - if _thread is not None: - _thread.join(timeout=5) diff --git a/app/main.py b/app/main.py index c8c3502..7589ac7 100644 --- a/app/main.py +++ b/app/main.py @@ -19,6 +19,7 @@ from app.api.internal.app_version import router as internal_app_version_router from app.api.internal.launch_confirm import router as internal_launch_confirm_router from app.api.internal.price import router as internal_price_router from app.api.internal.store import router as internal_store_router +from app.api.internal.wx_finance import router as internal_wx_finance_router from app.api.v1.ad import router as ad_router from app.api.v1.analytics import router as analytics_router from app.api.v1.auth import router as auth_router @@ -51,10 +52,6 @@ from app.core.heartbeat_monitor_worker import ( start_heartbeat_monitor, stop_heartbeat_monitor, ) -from app.core.wx_finance_worker import ( - start_wx_finance_worker, - stop_wx_finance_worker, -) from app.core.logging import setup_logging from app.core.pricebot_client import aclose_pricebot_client, get_pricebot_client from app.core.withdraw_reconcile_worker import ( @@ -86,12 +83,9 @@ async def lifespan(_: FastAPI) -> AsyncIterator[None]: reconcile_task = start_withdraw_reconcile_worker() heartbeat_task = start_heartbeat_monitor() daily_exchange_task = start_daily_exchange_worker() - # 会话存档轮询 worker(独立线程;gate 在 wx_finance_configured, 未开启则 no-op、不加载 .so) - start_wx_finance_worker() try: yield finally: - stop_wx_finance_worker() await stop_heartbeat_monitor(heartbeat_task) await stop_withdraw_reconcile_worker(reconcile_task) await stop_daily_exchange_worker(daily_exchange_task) @@ -146,6 +140,8 @@ app.include_router(internal_price_router) app.include_router(internal_store_router) app.include_router(internal_app_version_router) app.include_router(internal_launch_confirm_router) +# 会话存档独立 poller → 打比价信号(server→server, X-Internal-Secret;轮询进程在 app 外, 见 scripts/wx_finance_poller.py) +app.include_router(internal_wx_finance_router) app.include_router(platform_router) # CPS 群发短链跳转 /c/{code}(公网无鉴权:记点击 → 302 跳美团) app.include_router(cps_redirect_router) diff --git a/app/schemas/wx_finance.py b/app/schemas/wx_finance.py new file mode 100644 index 0000000..244cf79 --- /dev/null +++ b/app/schemas/wx_finance.py @@ -0,0 +1,13 @@ +"""会话存档独立 poller → app-server 内部信号端点的 schema。""" +from __future__ import annotations + +from pydantic import BaseModel + + +class WxFinancePendingIn(BaseModel): + source: str = "meituan" # 源平台代号(pop 后前端 launch 用;PoC 写死 meituan) + kind: str = "" # 触发物类型 weapp / image(仅日志用) + + +class WxFinancePendingOut(BaseModel): + ok: bool = True diff --git a/deploy/wx-finance-poller.service b/deploy/wx-finance-poller.service new file mode 100644 index 0000000..fbf8160 --- /dev/null +++ b/deploy/wx-finance-poller.service @@ -0,0 +1,18 @@ +[Unit] +# 会话存档轮询 poller —— 独立进程跑 WeWorkFinanceSdk(Go c-shared .so),与 app-server 隔离: +# 嵌进 app-server 会 segfault 把主进程带崩,拆出来后崩了 systemd 只重启本服务、不影响 app-server。 +Description=Shaguabijia 会话存档轮询 poller (WeWorkFinanceSdk, isolated from app-server) +After=network.target shaguabijia-app-server.service + +[Service] +Type=simple +User=root +WorkingDirectory=/opt/shaguabijia-app-server +Environment="PATH=/opt/shaguabijia-app-server/.venv/bin:/usr/bin:/bin" +EnvironmentFile=/opt/shaguabijia-app-server/.env +ExecStart=/opt/shaguabijia-app-server/.venv/bin/python scripts/wx_finance_poller.py +Restart=on-failure +RestartSec=5 + +[Install] +WantedBy=multi-user.target diff --git a/scripts/wx_finance_poller.py b/scripts/wx_finance_poller.py new file mode 100644 index 0000000..83198eb --- /dev/null +++ b/scripts/wx_finance_poller.py @@ -0,0 +1,128 @@ +"""会话存档【独立进程】轮询 poller —— 单独 systemd 服务跑, 与 app-server 隔离。 + +为什么独立进程:WeWorkFinanceSdk 是 Go 编译的 c-shared 库, 自带 Go runtime。嵌进 app-server(uvicorn) +后台线程里跑会 segfault、把主进程带崩(2026-07-17 实测 core-dump 每 5s 循环)。独立进程让 Go runtime +独占进程 + 主线程, 崩了 systemd 只重启本 poller、不动 app-server。拉到美团卡片(weapp)/截图(image) +→ POST app-server 内部端点 /internal/wx-finance/pending 打比价信号(下游心跳→弹窗→比价 完全复用)。 + + 手动跑: cd /opt/shaguabijia-app-server && .venv/bin/python scripts/wx_finance_poller.py + 生产: deploy/wx-finance-poller.service 常驻(systemctl enable --now wx-finance-poller) +""" +from __future__ import annotations + +import logging +import signal +import time +from pathlib import Path + +import httpx + +from app.core.config import settings +from app.core.logging import setup_logging +from app.integrations.wx_finance_sdk import WxFinanceError, WxFinanceSdk + +logger = logging.getLogger("shagua.wx_finance_poller") + +_stop = False + + +def _on_signal(_signum: int, _frame: object) -> None: + global _stop + _stop = True + + +def _load_seq() -> int: + try: + return int(Path(settings.WX_FINANCE_SEQ_FILE).read_text().strip() or "0") + except Exception: + return 0 + + +def _save_seq(seq: int) -> None: + try: + p = Path(settings.WX_FINANCE_SEQ_FILE) + p.parent.mkdir(parents=True, exist_ok=True) + p.write_text(str(seq)) + except Exception: + logger.exception("seq 持久化失败") + + +def _notify(client: httpx.Client, source: str, kind: str) -> None: + """POST 内部端点让 app-server 打信号(set_pending 在 app-server 进程内存)。""" + url = f"{settings.WX_FINANCE_INTERNAL_URL.rstrip('/')}/internal/wx-finance/pending" + try: + resp = client.post( + url, + json={"source": source, "kind": kind}, + headers={"X-Internal-Secret": settings.INTERNAL_API_SECRET}, + timeout=10, + ) + logger.info("notify app-server (%s) → %s", kind, resp.status_code) + except Exception: + logger.exception("notify app-server 失败") + + +def _handle(client: httpx.Client, msg: dict) -> None: + """明文消息:外部用户(from≠接收成员)发来的 weapp/image → 通知 app-server 打信号。""" + msgtype = msg.get("msgtype") or "" + frm = msg.get("from") or "" + receiver = settings.WX_FINANCE_RECEIVER_USERID + if receiver and frm == receiver: + return # 成员自己发的, 跳过 + if msgtype not in ("weapp", "image"): + return + logger.info("命中触发 from=%s type=%s", frm, msgtype) + _notify(client, "meituan", msgtype) + + +def main() -> None: + setup_logging(debug=settings.APP_DEBUG) + if not settings.wx_finance_configured: + logger.warning("WX_FINANCE 未配齐(ENABLED/corpid/存档Secret), poller 退出") + return + if not settings.INTERNAL_API_SECRET: + logger.error("INTERNAL_API_SECRET 未配置, poller 无法通知 app-server, 退出") + return + signal.signal(signal.SIGTERM, _on_signal) + signal.signal(signal.SIGINT, _on_signal) + + sdk = WxFinanceSdk( + settings.WX_FINANCE_SDK_PATH, + settings.WX_KF_CORP_ID, + settings.WX_FINANCE_SECRET, + Path(settings.WX_FINANCE_PRIVATE_KEY_PATH).read_bytes(), + ) + seq = _load_seq() + logger.info( + "wx_finance poller 启动, 从 seq=%d 轮询, 间隔 %ds", seq, settings.WX_FINANCE_POLL_INTERVAL_SEC + ) + with httpx.Client() as client: + while not _stop: + try: + msgs = sdk.get_chat_data(seq, limit=1000) + if msgs: + logger.info("拉到 %d 条(seq>%d)", len(msgs), seq) + for item in msgs: + try: + plain = sdk.decrypt(item["encrypt_random_key"], item["encrypt_chat_msg"]) + _handle(client, plain) + except Exception: + logger.exception("解密/处理单条失败 seq=%s", item.get("seq")) + seq = max(seq, int(item.get("seq", seq))) + if msgs: + _save_seq(seq) + except WxFinanceError: + logger.exception("GetChatData 出错") + except Exception: + logger.exception("轮询异常") + # 可打断的 sleep(收到 SIGTERM 尽快退) + for _ in range(max(1, settings.WX_FINANCE_POLL_INTERVAL_SEC) * 2): + if _stop: + break + time.sleep(0.5) + sdk.close() + logger.info("wx_finance poller 已停止") + + +if __name__ == "__main__": + main()