From 3790821ce0e3e8785db97bc71721f2597daafffa Mon Sep 17 00:00:00 2001 From: LittleGuo Date: Sat, 26 Sep 2026 22:25:02 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E9=98=B6=E6=AE=B5=201=20=E6=9C=80?= =?UTF-8?q?=E5=B0=8F=E9=97=AD=E7=8E=AF=20-=20=E5=BB=BA=E8=A1=A8=20SQL?= =?UTF-8?q?=E3=80=81/wechat=20=E4=BA=8B=E4=BB=B6=E3=80=81/auth=20=E6=8E=A5?= =?UTF-8?q?=E5=8F=A3=E3=80=81=E9=A6=96=E6=AC=A1=E5=85=B3=E6=B3=A8=E5=85=8D?= =?UTF-8?q?=E8=B4=B9=E6=8E=88=E6=9D=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - sql/schema.sql: users/authorizations/auth_scenes/usage_logs + sessions 建表 - db.py: aiomysql 连接池,lifespan 内初始化与释放 - wechat_api.py: access_token 缓存 + 临时二维码创建 - auth.py: /auth/create_scene、/auth/status,handle_scan 事务内幂等处理扫码 - wechat.py: 接入 DB 生命周期,处理 subscribe/SCAN 事件;移除多余的 openid query 参数 - 首次关注赠送 7 天免费授权,has_claimed_free 条件更新保证幂等 - config.py/.env.example: 新增 FREE_AUTH_DAYS/SCENE_TTL_SECONDS/SESSION_TTL_HOURS --- .env.example | 23 ++-- CODEBUDDY.md | 49 ++++++++ REQUIREMENTS.md | 286 +++++++++++++++++++++++++++++++++++++++++++++++ auth.py | 255 ++++++++++++++++++++++++++++++++++++++++++ config.py | 5 + db.py | 58 ++++++++++ requirements.txt | 2 + sql/schema.sql | 99 ++++++++++++++++ wechat.py | 93 +++++++++++---- wechat_api.py | 77 +++++++++++++ 10 files changed, 915 insertions(+), 32 deletions(-) create mode 100644 CODEBUDDY.md create mode 100644 REQUIREMENTS.md create mode 100644 auth.py create mode 100644 db.py create mode 100644 sql/schema.sql create mode 100644 wechat_api.py diff --git a/.env.example b/.env.example index 958a135..1cd4407 100644 --- a/.env.example +++ b/.env.example @@ -1,13 +1,11 @@ -""" -环境变量模板 - -使用方法: -1. 复制此文件为 .env -2. 填入你的实际值 -3. .env 不会进 Git(已在 .gitignore 中排除) - -cp .env.example .env -""" +# 环境变量模板 +# +# 使用方法: +# 1. 复制此文件为 .env +# 2. 填入你的实际值 +# 3. .env 不会进 Git(已在 .gitignore 中排除) +# +# cp .env.example .env # 微信测试号 - 从 https://mp.weixin.qq.com/debug/cgi-bin/sandbox?t=sandbox/login 获取 WECHAT_TOKEN=your_custom_token_here @@ -25,3 +23,8 @@ MYSQL_DB=wechat_api REDIS_HOST=127.0.0.1 REDIS_PORT=6379 REDIS_PASSWORD=your_redis_password_here + +# 业务配置 +FREE_AUTH_DAYS=7 +SCENE_TTL_SECONDS=300 +SESSION_TTL_HOURS=24 diff --git a/CODEBUDDY.md b/CODEBUDDY.md new file mode 100644 index 0000000..ef27403 --- /dev/null +++ b/CODEBUDDY.md @@ -0,0 +1,49 @@ +# CODEBUDDY.md + +This file provides guidance to CodeBuddy Code when working with code in this repository. + +## 项目概述 + +微信公众号扫码授权服务,基于 FastAPI。当前实现微信服务器验证(GET /wechat)与消息/事件接收(POST /wechat)。MySQL 与 Redis 已在配置层声明,但业务代码尚未使用。 + +## 常用命令 + +```bash +# 安装依赖(建议先创建/激活虚拟环境 venv) +pip install -r requirements.txt + +# 本地开发(uvicorn reload,监听 127.0.0.1:8000) +python run_local.py +# 等价于: +uvicorn wechat:app --reload --host 127.0.0.1 --port 8000 + +# 服务器部署(不 reload,单 worker,监听 127.0.0.1:8000,由 Nginx 反代) +python run_server.py + +# 初始化本地配置 +cp .env.example .env # 然后填入真实值 +``` + +当前仓库没有测试框架、lint 或构建配置;如需运行单个测试,需先引入 pytest 等工具。 + +## 架构 + +``` +config.py # 从 .env 读取配置(dotenv),模块级常量 +wechat.py # FastAPI app 本体:路由 + 签名校验 + XML 解析/构造 +run_local.py # 开发启动入口(reload=True) +run_server.py # 生产启动入口(reload=False, workers=1) +``` + +- **配置**:所有敏感值经 `config.py` 从 `.env` 读取,`.env` 已被 `.gitignore` 排除。`.env.example` 是字段模板。新增配置项需同时更新这两处。 +- **应用入口**:两个启动脚本均以 `"wechat:app"` 字符串形式加载 `wechat.py` 中的 `app`,因此模块名/对象名不可随意重命名。 +- **微信交互协议**: + - 所有请求先经 `verify_signature()`(token+timestamp+nonce 字典序拼接后 SHA1 比对)校验,失败返回 403。 + - GET 校验通过后原样返回 `echostr`。 + - POST 解析微信推送的 XML(`MsgType`/`FromUserName`/`Event` 等),通过 `_reply_text()` 构造文本回复 XML 返回。新增消息类型处理应在 `wechat_message()` 的事件/消息分支中扩展。 +- **注意**:`_reply_text()` 中 ToUserName/FromUserName 是反置的(回复时收发方互换),这是微信协议要求。 + +## 约定 + +- 代码注释与文档字符串使用中文。 +- 生产环境仅监听 127.0.0.1,对外暴露依赖 Nginx 反向代理。 diff --git a/REQUIREMENTS.md b/REQUIREMENTS.md new file mode 100644 index 0000000..2ec1311 --- /dev/null +++ b/REQUIREMENTS.md @@ -0,0 +1,286 @@ +# 微信扫码授权服务 — 需求文档 + +## 1. 项目概述 + +为一个 Windows MFC 桌面程序提供微信扫码授权服务。用户通过微信扫码关注公众号(当前使用微信测试号),服务端据此判断授权状态,MFC 端轮询后决定是否放行。后续支持用户充值获得时长或积分。 + +## 2. 技术栈 + +- 操作系统:Alibaba Cloud Linux 3 +- Web 框架:Python 3.10+ / FastAPI +- ASGI 服务器:Uvicorn +- 反向代理:Nginx(已配置,Cloudflare Tunnel 作为当前 HTTPS 入口) +- 数据库:MySQL 8.0(Docker 部署,监听 127.0.0.1:3306) +- 微信侧:微信公众平台测试号(后续迁移正式服务号) + +## 3. 微信测试号配置 + +在.env文件内. +URL:`http://ethereal-realm.top/wechat` | + +## 4. 系统架构 + +``` +MFC 客户端 + │ HTTPS + ▼ +Cloudflare Tunnel / Nginx + │ + ▼ +FastAPI 服务 + ├── /wechat 接收微信事件推送与 URL 验证 + ├── /auth/* MFC 请求授权相关接口 + └── /usage/* MFC 上报使用与扣减 + │ + ▼ +MySQL +``` + +服务端为唯一权威。MFC 端不缓存任何可信授权数据,仅保存会话令牌。 + +## 5. 数据库设计 + +### 5.1 users — 用户表 + +| 字段 | 类型 | 说明 | +| --- | --- | --- | +| id | BIGINT PK AUTO_INCREMENT | | +| openid | VARCHAR(64) UNIQUE NOT NULL | 微信 OpenID | +| nickname | VARCHAR(100) | 昵称(可选) | +| created_at | DATETIME | 首次关注时间 | +| last_seen_at | DATETIME | 最近活跃时间 | +| has_claimed_free | TINYINT DEFAULT 0 | 是否已领过免费授权 | + +### 5.2 authorizations — 授权表 + +| 字段 | 类型 | 说明 | +| --- | --- | --- | +| id | BIGINT PK AUTO_INCREMENT | | +| user_id | BIGINT NOT NULL | 外键 users.id | +| type | ENUM('time','points') | 授权类型 | +| start_at | DATETIME NULL | 时间授权开始 | +| end_at | DATETIME NULL | 时间授权结束 | +| remaining_points | INT DEFAULT 0 | 积分余额 | +| total_points | INT DEFAULT 0 | 积分总量 | +| source | ENUM('free','purchase','admin') | 来源 | +| status | ENUM('pending','active','expired','exhausted','cancelled') | 状态 | +| created_at | DATETIME | | +| updated_at | DATETIME | | + +索引:`idx_user_status (user_id, status)` + +### 5.3 auth_scenes — 扫码场景表 + +| 字段 | 类型 | 说明 | +| --- | --- | --- | +| id | BIGINT PK AUTO_INCREMENT | | +| scene_str | VARCHAR(128) UNIQUE NOT NULL | 二维码场景值 | +| device_id | VARCHAR(128) | MFC 设备标识 | +| status | ENUM('pending','scanned','authorized','expired') | | +| user_id | BIGINT NULL | 扫码用户 | +| created_at | DATETIME | | +| expires_at | DATETIME NOT NULL | 默认 300 秒 | +| authorized_at | DATETIME NULL | | + +索引:`idx_scene_status (scene_str, status)` + +### 5.4 usage_logs — 使用日志 + +| 字段 | 类型 | 说明 | +| --- | --- | --- | +| id | BIGINT PK AUTO_INCREMENT | | +| user_id | BIGINT NOT NULL | | +| device_id | VARCHAR(128) | | +| authorization_id | BIGINT | 本次扣减的授权 | +| cost_type | ENUM('time','points') | | +| cost_points | INT DEFAULT 0 | 时间授权为 0 | +| used_at | DATETIME | | + +索引:`idx_user_time (user_id, used_at)` + +### 5.5 orders — 订单表(阶段 3 使用) + +| 字段 | 类型 | 说明 | +| --- | --- | --- | +| id | BIGINT PK AUTO_INCREMENT | | +| order_no | VARCHAR(64) UNIQUE NOT NULL | | +| user_id | BIGINT NOT NULL | | +| amount | DECIMAL(10,2) NOT NULL | | +| product_id | BIGINT | 商品 ID | +| status | ENUM('pending','paid','failed','refunded') | | +| paid_at | DATETIME NULL | | +| created_at | DATETIME | | + +## 6. API 接口设计 + +### 6.1 微信侧 + +**GET /wechat** +微信服务器 URL 验证。校验 signature,返回 echostr。 + +**POST /wechat** +接收微信事件推送,解析 XML: + +- `subscribe` 事件:EventKey 形如 `qrscene_` +- `SCAN` 事件:EventKey 直接为 `` + +处理逻辑见第 7 节。 + +### 6.2 MFC 侧 + +**POST /auth/create_scene** + +请求: + +```json +{ "device_id": "设备唯一标识" } +``` + +响应: + +```json +{ + "scene_str": "pc_xxxxx", + "qr_url": "https://mp.weixin.qq.com/cgi-bin/showqrcode?ticket=...", + "expires_in": 300 +} +``` + +说明:服务端生成唯一 scene_str,调用微信接口生成临时二维码,写入 auth_scenes,返回二维码图片 URL。 + +**GET /auth/status?scene=xxx** + +响应: + +```json +{ + "status": "pending | authorized | expired | need_purchase", + "session_token": "当 status=authorized 时返回", + "authorization": { + "type": "time | points", + "end_at": "2026-10-03T12:00:00", + "remaining_points": 0 + } +} +``` + +**POST /usage/consume** + +请求: + +```json +{ + "session_token": "xxx", + "device_id": "xxx" +} +``` + +响应: + +```json +{ + "ok": true, + "authorization": { + "type": "time", + "end_at": "2026-10-03T12:00:00" + } +} +``` + +失败时: + +```json +{ + "ok": false, + "reason": "expired | exhausted | invalid_token" +} +``` + +## 7. 核心业务流程 + +### 7.1 首次扫码授权 + +1. MFC 启动,检查本地 session_token,无则调用 `/auth/create_scene` +2. MFC 显示二维码,每 2 秒轮询 `/auth/status` +3. 用户微信扫码 +4. 微信推送事件到 `/wechat`,服务端解析 scene_str +5. 服务端查找或创建 user(按 openid) +6. 如果 `has_claimed_free = 0`:创建时间授权,`start_at = now()`,`end_at = now() + 7 天`,`source = free`,`status = active`,`has_claimed_free = 1` +7. 如果 `has_claimed_free = 1`:读取该用户当前有效授权(active 状态) +8. 将 auth_scenes.status 置为 `authorized`,绑定 user_id +9. MFC 轮询到 authorized,拿到 session_token,关闭弹窗 + +### 7.2 再次扫码 + +触发条件:时间授权到期,或积分耗尽,或换设备。 + +流程同上,但服务端在步骤 6 时: + +- 若无有效授权 → 返回 `need_purchase` +- 若有有效授权 → 正常返回 authorized + +### 7.3 每次使用扣减 + +1. MFC 调用 `/usage/consume` +2. 服务端校验 session_token +3. 查找该用户当前 active 授权 +4. 时间授权:检查 `now() < end_at`,未过期则记录 usage_logs,直接返回 ok +5. 积分授权:事务内 `UPDATE authorizations SET remaining_points = remaining_points - 1, status = IF(remaining_points - 1 <= 0, 'exhausted', 'active') WHERE id = ? AND remaining_points > 0 AND status = 'active'`,影响行数 1 则记录 usage_logs 并返回 ok,否则返回 exhausted +6. 时间授权到期或积分耗尽时,MFC 下次使用会收到失败,弹出二维码引导再次扫码 + +## 8. 授权规则 + +**互斥原则**:同一用户同一时刻最多只有一条 `status = active` 的授权。 + +**新授权下发规则**: + +- 若用户当前无 active 授权 → 新授权直接 active +- 若用户当前有 active 授权 → 新授权以 `status = pending` 保存,待当前授权到期或耗尽后,由定时任务或下次请求时激活 + +**免费授权**:仅限首次关注,每个 openid 一次,7 天时间授权。 + +**时间与积分不叠加**:用户同时拥有时间和积分授权时,以当前 active 的为准,另一条保持 pending。active 结束后自动切换。 + +## 9. 边界与异常处理 + +- scene_str 一次性:auth_scenes 被授权后 status 不再回到 pending +- scene 过期:超过 expires_at 后 status 置为 expired,MFC 轮询返回 expired,提示刷新二维码 +- 重复扫码:同一 scene 被多次扫描,只处理第一次,后续忽略 +- 已关注用户扫码:必须处理 SCAN 事件(EventKey 无 qrscene_ 前缀) +- 未关注用户扫码:处理 subscribe 事件(EventKey 有 qrscene_ 前缀) +- 并发扣减:必须使用数据库事务,禁止应用层先读后写 +- session_token:随机生成,存服务端,有效期 24 小时,过期需重新扫码 + +## 10. 开发阶段划分 + +**阶段 1:最小闭环** + +- 建库建表 +- 实现 `/wechat` 的 GET 验证与 POST 事件接收 +- 实现 `/auth/create_scene`、`/auth/status` +- 实现首次关注赠送 7 天免费授权 +- 用 curl 或 Postman 模拟微信事件,验证状态流转 + +**阶段 2:扣减与再扫码** + +- 实现 `/usage/consume` +- 实现时间到期与积分耗尽后的再次扫码流程 +- 实现 pending 授权自动激活 + +**阶段 3:充值** + +- 实现 orders 表与卡密兑换接口 +- 后续接入微信支付 Native 扫码 + +## 11. 验收标准 + +- 首次扫码关注后,MFC 能收到 authorized 并正常使用 +- 7 天后再次使用,MFC 收到 expired,弹出二维码 +- 已关注用户再次扫码,服务端能收到 SCAN 事件并正确识别 +- 积分扣减在并发请求下不会出现负数 +- scene_str 一次性,重复扫描不产生副作用 +- 所有授权判断均在服务端完成,MFC 本地无法伪造 + +## 12. 待确认 + +- 用户充值获得积分后,若当前处于时间授权 active 状态,积分授权是否进入 pending 等待(建议:是) diff --git a/auth.py b/auth.py new file mode 100644 index 0000000..12a7daf --- /dev/null +++ b/auth.py @@ -0,0 +1,255 @@ +""" +授权接口与核心业务逻辑 + +路由(MFC 侧): + POST /auth/create_scene 生成 scene_str、创建微信临时二维码并落库 + GET /auth/status 轮询扫码授权结果 + +业务函数(微信事件侧,由 wechat.py 调用): + handle_scan() 处理扫码事件:建用户、发免费授权、绑定场景、签发会话 +""" + +import logging +import secrets + +import aiomysql +from fastapi import APIRouter, HTTPException, Query +from pydantic import BaseModel + +import db +from config import FREE_AUTH_DAYS, SCENE_TTL_SECONDS, SESSION_TTL_HOURS +from wechat_api import create_temp_qrcode + +logger = logging.getLogger(__name__) + +router = APIRouter(prefix="/auth", tags=["auth"]) + +SCENE_PREFIX = "pc_" + + +class CreateSceneRequest(BaseModel): + device_id: str + + +@router.post("/create_scene") +async def create_scene(payload: CreateSceneRequest): + """生成唯一 scene_str,调用微信接口创建临时二维码,写入 auth_scenes""" + scene_str = SCENE_PREFIX + secrets.token_hex(16) + try: + qr_url = await create_temp_qrcode(scene_str, SCENE_TTL_SECONDS) + except Exception as exc: + # 微信接口不可用(appid/secret 未配置、网络异常、token 失效等) + logger.exception("创建二维码失败 scene_str=%s", scene_str) + raise HTTPException(status_code=502, detail=f"创建微信二维码失败: {exc}") + + async with db.acquire() as conn: + async with conn.cursor() as cur: + await cur.execute( + "INSERT INTO auth_scenes (scene_str, device_id, status, created_at, expires_at) " + "VALUES (%s, %s, 'pending', NOW(), DATE_ADD(NOW(), INTERVAL %s SECOND))", + (scene_str, payload.device_id, SCENE_TTL_SECONDS), + ) + + return {"scene_str": scene_str, "qr_url": qr_url, "expires_in": SCENE_TTL_SECONDS} + + +@router.get("/status") +async def get_status(scene: str = Query(..., description="create_scene 返回的 scene_str")): + """查询扫码授权状态:pending / authorized / expired / need_purchase""" + async with db.acquire() as conn: + async with conn.cursor(aiomysql.DictCursor) as cur: + await cur.execute( + "SELECT id, status, user_id, (expires_at > NOW()) AS not_expired " + "FROM auth_scenes WHERE scene_str = %s", + (scene,), + ) + scene_row = await cur.fetchone() + if scene_row is None: + raise HTTPException(status_code=404, detail="scene 不存在") + + # 尚未扫码:过期则惰性置为 expired + if scene_row["status"] in ("pending", "scanned"): + if not scene_row["not_expired"]: + await cur.execute( + "UPDATE auth_scenes SET status = 'expired' " + "WHERE id = %s AND status IN ('pending', 'scanned')", + (scene_row["id"],), + ) + return {"status": "expired"} + return {"status": "pending"} + + if scene_row["status"] == "expired": + return {"status": "expired"} + + # 已扫码授权:判断用户当前是否有可用授权 + auth_row = await _get_active_authorization(cur, scene_row["user_id"]) + if auth_row is None: + # 免费已领过且无有效授权 → 引导充值(阶段 3) + return {"status": "need_purchase"} + + await cur.execute( + "SELECT token, (expires_at > NOW()) AS not_expired FROM sessions " + "WHERE scene_id = %s ORDER BY id DESC LIMIT 1", + (scene_row["id"],), + ) + session_row = await cur.fetchone() + if session_row is None or not session_row["not_expired"]: + return {"status": "expired"} + + return { + "status": "authorized", + "session_token": session_row["token"], + "authorization": _serialize_authorization(auth_row), + } + + +# --------------------------------------------------------------------------- +# 业务逻辑 +# --------------------------------------------------------------------------- + +async def handle_scan(scene_str: str, openid: str) -> bool: + """ + 处理扫码事件。 + + 返回 True 表示本次扫码完成授权,False 表示场景无效、已过期或已被处理。 + 整个流程在事务内完成,并对 scene 行加排他锁,保证同一 scene 只被处理一次。 + """ + async with db.acquire() as conn: + await conn.begin() + try: + async with conn.cursor(aiomysql.DictCursor) as cur: + await cur.execute( + "SELECT id, status, device_id, (expires_at > NOW()) AS not_expired " + "FROM auth_scenes WHERE scene_str = %s FOR UPDATE", + (scene_str,), + ) + scene_row = await cur.fetchone() + + if scene_row is None: + await conn.rollback() + return False + + if scene_row["status"] == "authorized": + # 重复扫码:只处理第一次,后续忽略 + await conn.rollback() + return False + + if not scene_row["not_expired"]: + await cur.execute( + "UPDATE auth_scenes SET status = 'expired' " + "WHERE id = %s AND status IN ('pending', 'scanned')", + (scene_row["id"],), + ) + await conn.commit() + return False + + user_id = await _find_or_create_user(cur, openid) + await _grant_free_authorization(cur, user_id) + + await cur.execute( + "UPDATE auth_scenes SET status = 'authorized', user_id = %s, authorized_at = NOW() " + "WHERE id = %s", + (user_id, scene_row["id"]), + ) + + token = secrets.token_urlsafe(32) + await cur.execute( + "INSERT INTO sessions (token, user_id, device_id, scene_id, created_at, expires_at) " + "VALUES (%s, %s, %s, %s, NOW(), DATE_ADD(NOW(), INTERVAL %s HOUR))", + (token, user_id, scene_row["device_id"], scene_row["id"], SESSION_TTL_HOURS), + ) + + await conn.commit() + logger.info("扫码授权成功 scene_str=%s openid=%s", scene_str, openid) + return True + except Exception: + await conn.rollback() + raise + + +async def _find_or_create_user(cur, openid: str) -> int: + """按 openid 查找用户,不存在则创建,并刷新 last_seen_at""" + await cur.execute("SELECT id FROM users WHERE openid = %s", (openid,)) + row = await cur.fetchone() + if row is None: + # INSERT IGNORE + 重查:并发扫码时避免唯一键冲突报错 + await cur.execute( + "INSERT IGNORE INTO users (openid, created_at, last_seen_at, has_claimed_free) " + "VALUES (%s, NOW(), NOW(), 0)", + (openid,), + ) + await cur.execute("SELECT id FROM users WHERE openid = %s", (openid,)) + row = await cur.fetchone() + + await cur.execute("UPDATE users SET last_seen_at = NOW() WHERE id = %s", (row["id"],)) + return row["id"] + + +async def _grant_free_authorization(cur, user_id: int) -> None: + """ + 首次关注赠送 7 天时间授权。 + + 以 has_claimed_free 的条件更新作为幂等闸门:只有把 0 改成 1 的那一次才真正发授权。 + 若用户已有 active 授权(互斥原则),新授权以 pending 保存。 + """ + await cur.execute( + "UPDATE users SET has_claimed_free = 1 WHERE id = %s AND has_claimed_free = 0", + (user_id,), + ) + if cur.rowcount != 1: + return + + status = "pending" if await _has_active_authorization(cur, user_id) else "active" + await cur.execute( + "INSERT INTO authorizations " + "(user_id, type, start_at, end_at, remaining_points, total_points, source, status, created_at, updated_at) " + "VALUES (%s, 'time', NOW(), DATE_ADD(NOW(), INTERVAL %s DAY), 0, 0, 'free', %s, NOW(), NOW())", + (user_id, FREE_AUTH_DAYS, status), + ) + logger.info("已发放免费授权 user_id=%s days=%s status=%s", user_id, FREE_AUTH_DAYS, status) + + +async def _has_active_authorization(cur, user_id: int) -> bool: + await cur.execute( + "SELECT 1 FROM authorizations WHERE user_id = %s AND status = 'active' LIMIT 1", + (user_id,), + ) + return await cur.fetchone() is not None + + +async def _get_active_authorization(cur, user_id: int): + """取用户当前 active 授权;已失效的惰性置为 expired / exhausted 并返回 None""" + await cur.execute( + "SELECT id, type, end_at, remaining_points, total_points, " + "(end_at IS NOT NULL AND end_at > NOW()) AS time_valid " + "FROM authorizations WHERE user_id = %s AND status = 'active' ORDER BY id DESC LIMIT 1", + (user_id,), + ) + row = await cur.fetchone() + if row is None: + return None + + if row["type"] == "time": + if not row["time_valid"]: + await cur.execute( + "UPDATE authorizations SET status = 'expired' WHERE id = %s AND status = 'active'", + (row["id"],), + ) + return None + return row + + if row["remaining_points"] <= 0: + await cur.execute( + "UPDATE authorizations SET status = 'exhausted' WHERE id = %s AND status = 'active'", + (row["id"],), + ) + return None + return row + + +def _serialize_authorization(row) -> dict: + return { + "type": row["type"], + "end_at": row["end_at"].isoformat() if row["end_at"] else None, + "remaining_points": row["remaining_points"], + } diff --git a/config.py b/config.py index 2bafe71..6f12911 100644 --- a/config.py +++ b/config.py @@ -25,3 +25,8 @@ MYSQL_DB = os.getenv("MYSQL_DB", "wechat_api") REDIS_HOST = os.getenv("REDIS_HOST", "127.0.0.1") REDIS_PORT = int(os.getenv("REDIS_PORT", "6379")) REDIS_PASSWORD = os.getenv("REDIS_PASSWORD", "") + +# 业务配置 +FREE_AUTH_DAYS = int(os.getenv("FREE_AUTH_DAYS", "7")) # 首次关注赠送天数 +SCENE_TTL_SECONDS = int(os.getenv("SCENE_TTL_SECONDS", "300")) # 二维码/scene 有效期(秒) +SESSION_TTL_HOURS = int(os.getenv("SESSION_TTL_HOURS", "24")) # session_token 有效期(小时) diff --git a/db.py b/db.py new file mode 100644 index 0000000..12bb870 --- /dev/null +++ b/db.py @@ -0,0 +1,58 @@ +""" +MySQL 连接池管理(aiomysql) + +用法: + await init_pool() # 应用启动时调用 + async with acquire() as conn: + ... + await close_pool() # 应用关闭时调用 + +连接池开启 autocommit,单条语句自动提交;需要事务时显式调用 conn.begin() / conn.commit()。 +""" + +from contextlib import asynccontextmanager +from typing import Optional + +import aiomysql + +from config import MYSQL_DB, MYSQL_HOST, MYSQL_PASSWORD, MYSQL_PORT, MYSQL_USER + +_pool: Optional[aiomysql.Pool] = None + + +async def init_pool() -> None: + """创建全局连接池""" + global _pool + _pool = await aiomysql.create_pool( + host=MYSQL_HOST, + port=MYSQL_PORT, + user=MYSQL_USER, + password=MYSQL_PASSWORD, + db=MYSQL_DB, + charset="utf8mb4", + autocommit=True, + minsize=1, + maxsize=10, + ) + + +async def close_pool() -> None: + """关闭连接池""" + global _pool + if _pool is not None: + _pool.close() + await _pool.wait_closed() + _pool = None + + +def get_pool() -> aiomysql.Pool: + if _pool is None: + raise RuntimeError("MySQL 连接池未初始化,请先调用 init_pool()") + return _pool + + +@asynccontextmanager +async def acquire(): + """从连接池借出一个连接""" + async with get_pool().acquire() as conn: + yield conn diff --git a/requirements.txt b/requirements.txt index 237bfa3..0c8c5aa 100644 --- a/requirements.txt +++ b/requirements.txt @@ -2,3 +2,5 @@ fastapi==0.115.0 uvicorn[standard]==0.30.0 httpx==0.27.0 python-dotenv==1.0.1 +aiomysql==0.2.0 +cryptography==43.0.1 diff --git a/sql/schema.sql b/sql/schema.sql new file mode 100644 index 0000000..0270ee7 --- /dev/null +++ b/sql/schema.sql @@ -0,0 +1,99 @@ +-- 微信扫码授权服务 — 数据库结构(阶段 1) +-- 目标环境:MySQL 8.0 +-- 执行方式:mysql -u root -p < sql/schema.sql + +CREATE DATABASE IF NOT EXISTS `wechat_api` + DEFAULT CHARACTER SET utf8mb4 + DEFAULT COLLATE utf8mb4_unicode_ci; + +USE `wechat_api`; + +-- --------------------------------------------------------------------------- +-- 5.1 users — 用户表 +-- --------------------------------------------------------------------------- +CREATE TABLE IF NOT EXISTS `users` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `openid` VARCHAR(64) NOT NULL COMMENT '微信 OpenID', + `nickname` VARCHAR(100) DEFAULT NULL COMMENT '昵称(可选)', + `created_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '首次关注时间', + `last_seen_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '最近活跃时间', + `has_claimed_free` TINYINT NOT NULL DEFAULT 0 COMMENT '是否已领过免费授权', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_openid` (`openid`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='用户表'; + +-- --------------------------------------------------------------------------- +-- 5.2 authorizations — 授权表 +-- --------------------------------------------------------------------------- +CREATE TABLE IF NOT EXISTS `authorizations` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `user_id` BIGINT UNSIGNED NOT NULL COMMENT '外键 users.id', + `type` ENUM('time','points') NOT NULL COMMENT '授权类型', + `start_at` DATETIME DEFAULT NULL COMMENT '时间授权开始', + `end_at` DATETIME DEFAULT NULL COMMENT '时间授权结束', + `remaining_points` INT NOT NULL DEFAULT 0 COMMENT '积分余额', + `total_points` INT NOT NULL DEFAULT 0 COMMENT '积分总量', + `source` ENUM('free','purchase','admin') NOT NULL COMMENT '来源', + `status` ENUM('pending','active','expired','exhausted','cancelled') + NOT NULL DEFAULT 'pending' COMMENT '状态', + `created_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + `updated_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + PRIMARY KEY (`id`), + KEY `idx_user_status` (`user_id`, `status`), + CONSTRAINT `fk_auth_user` FOREIGN KEY (`user_id`) REFERENCES `users` (`id`) ON DELETE CASCADE +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='授权表'; + +-- --------------------------------------------------------------------------- +-- 5.3 auth_scenes — 扫码场景表 +-- --------------------------------------------------------------------------- +CREATE TABLE IF NOT EXISTS `auth_scenes` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `scene_str` VARCHAR(128) NOT NULL COMMENT '二维码场景值', + `device_id` VARCHAR(128) DEFAULT NULL COMMENT 'MFC 设备标识', + `status` ENUM('pending','scanned','authorized','expired') + NOT NULL DEFAULT 'pending' COMMENT '场景状态', + `user_id` BIGINT UNSIGNED DEFAULT NULL COMMENT '扫码用户', + `created_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + `expires_at` DATETIME NOT NULL COMMENT '过期时间,默认 300 秒', + `authorized_at` DATETIME DEFAULT NULL COMMENT '授权完成时间', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_scene_str` (`scene_str`), + KEY `idx_scene_status` (`scene_str`, `status`), + CONSTRAINT `fk_scene_user` FOREIGN KEY (`user_id`) REFERENCES `users` (`id`) ON DELETE SET NULL +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='扫码场景表'; + +-- --------------------------------------------------------------------------- +-- 5.4 usage_logs — 使用日志 +-- --------------------------------------------------------------------------- +CREATE TABLE IF NOT EXISTS `usage_logs` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `user_id` BIGINT UNSIGNED NOT NULL, + `device_id` VARCHAR(128) DEFAULT NULL, + `authorization_id` BIGINT UNSIGNED DEFAULT NULL COMMENT '本次扣减的授权', + `cost_type` ENUM('time','points') NOT NULL, + `cost_points` INT NOT NULL DEFAULT 0 COMMENT '时间授权为 0', + `used_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (`id`), + KEY `idx_user_time` (`user_id`, `used_at`), + CONSTRAINT `fk_log_user` FOREIGN KEY (`user_id`) REFERENCES `users` (`id`) ON DELETE CASCADE, + CONSTRAINT `fk_log_auth` FOREIGN KEY (`authorization_id`) REFERENCES `authorizations` (`id`) ON DELETE SET NULL +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='使用日志'; + +-- --------------------------------------------------------------------------- +-- sessions — 会话令牌表(需求 9:session_token 服务端存储,有效期 24 小时) +-- --------------------------------------------------------------------------- +CREATE TABLE IF NOT EXISTS `sessions` ( + `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, + `token` VARCHAR(64) NOT NULL COMMENT '随机生成的会话令牌', + `user_id` BIGINT UNSIGNED NOT NULL, + `device_id` VARCHAR(128) DEFAULT NULL COMMENT '签发时绑定的设备', + `scene_id` BIGINT UNSIGNED DEFAULT NULL COMMENT '签发来源场景', + `created_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + `expires_at` DATETIME NOT NULL COMMENT '过期时间,默认 24 小时', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_token` (`token`), + KEY `idx_session_user` (`user_id`), + KEY `idx_session_scene` (`scene_id`), + CONSTRAINT `fk_session_user` FOREIGN KEY (`user_id`) REFERENCES `users` (`id`) ON DELETE CASCADE, + CONSTRAINT `fk_session_scene` FOREIGN KEY (`scene_id`) REFERENCES `auth_scenes` (`id`) ON DELETE SET NULL +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='会话令牌表'; diff --git a/wechat.py b/wechat.py index c62c9f3..ce01646 100644 --- a/wechat.py +++ b/wechat.py @@ -2,19 +2,43 @@ 微信公众号 FastAPI 应用 GET /wechat - 微信服务器验证(签名校验 + 返回 echostr) -POST /wechat - 接收微信推送的消息和事件 +POST /wechat - 接收微信推送的消息和事件(subscribe / SCAN 触发扫码授权) + +授权接口在 auth.py 中定义,通过 include_router 挂载。 """ import hashlib +import logging import time import xml.etree.ElementTree as ET +from contextlib import asynccontextmanager -from fastapi import FastAPI, Request, Query, HTTPException +from fastapi import FastAPI, HTTPException, Query, Request from fastapi.responses import PlainTextResponse +import auth +import db from config import WECHAT_TOKEN -app = FastAPI(title="WeChat API", version="0.1.0") +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)s [%(name)s] %(message)s", +) +logger = logging.getLogger(__name__) + + +@asynccontextmanager +async def lifespan(app: FastAPI): + """应用启动时建立 MySQL 连接池,关闭时释放""" + await db.init_pool() + logger.info("MySQL 连接池已初始化") + yield + await db.close_pool() + logger.info("MySQL 连接池已关闭") + + +app = FastAPI(title="WeChat API", version="0.2.0", lifespan=lifespan) +app.include_router(auth.router) def verify_signature(signature: str, timestamp: str, nonce: str) -> bool: @@ -52,17 +76,16 @@ async def wechat_message( signature: str = Query(...), timestamp: str = Query(...), nonce: str = Query(...), - openid: str = Query(...), ): """POST: 接收微信推送的消息和事件""" - # 验证签名 if not verify_signature(signature, timestamp, nonce): raise HTTPException(status_code=403, detail="Invalid signature") - # 解析 XML 消息体 body = await request.body() - root = ET.fromstring(body) + if not body: + return PlainTextResponse(content="success") + root = ET.fromstring(body) msg_type = root.findtext("MsgType", "") from_user = root.findtext("FromUserName", "") # 发送方(用户 openid) to_user = root.findtext("ToUserName", "") # 接收方(公众号) @@ -70,31 +93,57 @@ async def wechat_message( event = root.findtext("Event", "") event_key = root.findtext("EventKey", "") - print(f"[WeChat] type={msg_type} from={from_user} event={event} content={content} key={event_key}") + logger.info( + "收到微信推送 type=%s event=%s from=%s key=%s", msg_type, event, from_user, event_key + ) - # 处理事件推送 if msg_type == "event": if event == "subscribe": - # 用户关注 - return _reply_text(from_user, to_user, "欢迎关注!") - elif event == "unsubscribe": - # 用户取关 - print(f"[WeChat] 用户取关: {from_user}") + # 未关注用户扫码关注:EventKey 形如 qrscene_ + return await _handle_scan_event(from_user, to_user, event_key, is_subscribe=True) + if event == "SCAN": + # 已关注用户扫码:EventKey 直接是 + return await _handle_scan_event(from_user, to_user, event_key, is_subscribe=False) + if event == "unsubscribe": + logger.info("用户取关 openid=%s", from_user) return PlainTextResponse(content="success") - elif event == "SCAN": - # 已关注用户扫码 - return _reply_text(from_user, to_user, f"扫码成功,场景值: {event_key}") - elif event == "CLICK": - # 菜单点击 + if event == "CLICK": return _reply_text(from_user, to_user, f"点击了: {event_key}") - # 处理文本消息 if msg_type == "text": # 原样返回(echo 模式,方便测试) return _reply_text(from_user, to_user, f"你说: {content}") - # 其他类型暂不处理 - return _reply_text(from_user, to_user, "收到") + return PlainTextResponse(content="success") + + +async def _handle_scan_event( + from_user: str, to_user: str, event_key: str, is_subscribe: bool +) -> PlainTextResponse: + """解析场景值并完成扫码授权,回复用户处理结果""" + scene_str = _parse_scene_key(event_key, is_subscribe) + if not scene_str: + # 无场景值的普通关注 + return _reply_text(from_user, to_user, "欢迎关注!") + + authorized = await auth.handle_scan(scene_str, from_user) + if authorized: + return _reply_text(from_user, to_user, "授权成功,请返回电脑端继续操作。") + return _reply_text(from_user, to_user, "二维码已失效或已被使用,请在电脑端刷新后重新扫码。") + + +def _parse_scene_key(event_key: str, is_subscribe: bool) -> str: + """ + 从 EventKey 中提取 scene_str。 + + subscribe 事件带 qrscene_ 前缀,SCAN 事件不带;无场景值时返回空串。 + """ + if not event_key: + return "" + if is_subscribe: + prefix = "qrscene_" + return event_key[len(prefix):] if event_key.startswith(prefix) else "" + return event_key def _reply_text(from_user: str, to_user: str, content: str) -> PlainTextResponse: diff --git a/wechat_api.py b/wechat_api.py new file mode 100644 index 0000000..b440952 --- /dev/null +++ b/wechat_api.py @@ -0,0 +1,77 @@ +""" +微信公众平台开放接口封装 + +- get_access_token(): 获取并缓存 access_token(有效期 7200 秒) +- create_temp_qrcode(): 创建带字符串场景值的临时二维码,返回图片 URL + +仅服务端调用,appid/secret 来自 .env。 +""" + +import logging +import time +from urllib.parse import quote + +import httpx + +from config import WECHAT_APPID, WECHAT_SECRET + +logger = logging.getLogger(__name__) + +API_BASE = "https://api.weixin.qq.com/cgi-bin" +QRCODE_URL = "https://mp.weixin.qq.com/cgi-bin/showqrcode" + +# access_token 内存缓存(单 worker 部署足够;提前 5 分钟过期避免边界失效) +_token_cache = {"value": "", "expires_at": 0.0} + + +async def get_access_token() -> str: + """获取 access_token,命中缓存则直接返回""" + now = time.time() + if _token_cache["value"] and now < _token_cache["expires_at"]: + return _token_cache["value"] + + if not WECHAT_APPID or not WECHAT_SECRET: + raise RuntimeError("WECHAT_APPID / WECHAT_SECRET 未配置") + + params = { + "grant_type": "client_credential", + "appid": WECHAT_APPID, + "secret": WECHAT_SECRET, + } + async with httpx.AsyncClient(timeout=10) as client: + resp = await client.get(f"{API_BASE}/token", params=params) + data = resp.json() + + if "access_token" not in data: + raise RuntimeError(f"获取 access_token 失败: {data}") + + _token_cache["value"] = data["access_token"] + _token_cache["expires_at"] = now + int(data.get("expires_in", 7200)) - 300 + return _token_cache["value"] + + +async def create_temp_qrcode(scene_str: str, expire_seconds: int) -> str: + """创建临时二维码(字符串场景值),返回扫码图片 URL""" + token = await get_access_token() + payload = { + "expire_seconds": expire_seconds, + "action_name": "QR_STR_SCENE", + "action_info": {"scene": {"scene_str": scene_str}}, + } + async with httpx.AsyncClient(timeout=10) as client: + resp = await client.post( + f"{API_BASE}/qrcode/create", + params={"access_token": token}, + json=payload, + ) + data = resp.json() + + ticket = data.get("ticket") + if not ticket: + # 40001 等 token 失效错误:清空缓存,下次重新获取 + if data.get("errcode") in (40001, 42001): + _token_cache["value"] = "" + raise RuntimeError(f"创建二维码失败: {data}") + + logger.info("已创建临时二维码 scene_str=%s", scene_str) + return f"{QRCODE_URL}?ticket={quote(ticket)}"