refactor(wx-finance): 会话存档轮询拆独立进程(Go .so 嵌 app-server 会 core-dump)
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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)
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
+3
-7
@@ -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)
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user