提交 ade792f7 authored 作者: 王超's avatar 王超

v1

上级
# Legal service API config
LEGAL_BASE_URL=http://sys.12348.gov.cn/resources
LEGAL_API_PREFIX=/resources/api
# AES encryption config
AES_KEY=1234567890123456
AES_IV=1234567890123456
# App config
APP_NAME=法律抢单系统
APP_ENV=development
SECRET_KEY=change-me-in-production
# Database
DATABASE_URL=sqlite:///./legal_system.db
# Server
HOST=0.0.0.0
PORT=8000
.venv/
chunk_*.js
vendor.js
manifest.js
frontend/node_modules
\ No newline at end of file
# 法律抢单系统 — Docker 部署文档
## 架构
```
用户浏览器 → Nginx (端口 80)
├── / → 前端 (React SPA)
└── /api/ → 后端 (FastAPI :8080)
```
两个独立容器,通过 Docker 网络连接:
- **`law_order_frontend`** — Nginx 反代 + React 静态文件
- **`law_order_backend`** — FastAPI + uvicorn + SQLite
---
## 前置准备
### 服务器环境
| 项目 | 要求 |
|------|------|
| OS | Ubuntu 26.04 |
| Docker | 24.0+ (`docker --version`) |
| 端口 | 80 (HTTP) |
| 磁盘 | 2GB+ (镜像 + 数据) |
### 安装 Docker
```bash
apt update && apt install -y docker.io
systemctl enable --now docker
docker --version
```
---
## 部署步骤
### 1. 上传项目到服务器
```bash
# 在 Windows 上打包(⚠️ 注意 .env 文件已复制到 backend/ 目录)
tar -czf law_order_deploy.tar.gz \
backend/Dockerfile backend/main.py backend/app/ backend/requirements.txt backend/.env \
frontend/Dockerfile frontend/src/ frontend/package.json frontend/index.html frontend/vite.config.js \
nginx/nginx.conf
# 上传到服务器
scp law_order_deploy.tar.gz user@server:/opt/law_order/
# 在服务器解压
cd /opt/law_order && tar -xzf law_order_deploy.tar.gz
```
或者直接在服务器克隆代码仓库。
### 2. 创建目录
```bash
cd /opt/law_order
mkdir -p backend/data backend/logs
```
### 3. 创建 `.env` 文件
```bash
vim backend/.env
```
内容(按需修改):
```
# 法网 API 配置
LEGAL_BASE_URL=http://sys.12348.gov.cn/resources
LEGAL_API_PREFIX=/resources/api
# AES 加密
AES_KEY=1234567890123456
AES_IV=1234567890123456
# 应用配置
APP_NAME=法律抢单系统
APP_ENV=production
SECRET_KEY=<生成随机密钥>
# 数据库(容器内路径)
DATABASE_URL=sqlite:///./data/app.db
# 抢单配置
GRAB_INTERVAL_MS=300
HEARTBEAT_INTERVAL_S=60
# 管理员账号
ADMIN_USERNAME=admin
ADMIN_PASSWORD=<你的密码>
```
### 4. 构建镜像
```bash
cd /opt/law_order
# 构建后端
docker build -t law_order_backend:latest ./backend
# 构建前端
docker build -t law_order_frontend:latest ./frontend
```
### 5. 创建网络
```bash
docker network create law_order_net
```
### 6. 启动容器
```bash
# 后端 (挂载数据 + 日志 + .env)
docker run -d --name law_order_backend \
--network law_order_net \
--restart unless-stopped \
-v /opt/law_order/backend/data:/app/data \
-v /opt/law_order/backend/logs:/app/logs \
-v /opt/law_order/backend/.env:/app/.env \
law_order_backend:latest
# 前端 (Nginx 反代)
docker run -d --name law_order_frontend \
--network law_order_net \
--restart unless-stopped \
-p 80:80 \
law_order_frontend:latest
```
### 7. 验证
```bash
# 检查容器状态
docker ps
# 检查后端日志
docker logs -f law_order_backend
# 测试 API
curl http://localhost/api/health
# 访问前端
# 浏览器打开 http://<服务器IP>
```
---
## 日常运维
### 查看日志
```bash
docker logs -f law_order_backend # 后端日志
docker logs -f law_order_frontend # 前端日志
cat /opt/law_order/backend/logs/app_*.log # 文件日志
```
### 重启容器
```bash
docker restart law_order_backend law_order_frontend
```
### 停止容器
```bash
docker stop law_order_backend law_order_frontend
```
### 更新部署
```bash
# 1. 停止旧容器
docker stop law_order_backend law_order_frontend
# 2. 删除旧容器
docker rm law_order_backend law_order_frontend
# 3. 重新构建镜像
docker build -t law_order_backend:latest ./backend
docker build -t law_order_frontend:latest ./frontend
# 4. 重新启动 (见步骤 6)
```
### 数据库备份
```bash
cp /opt/law_order/backend/data/app.db /opt/law_order/backend/data/app.db.bak.$(date +%Y%m%d)
```
### 清理无用镜像
```bash
docker image prune -a
```
---
## 环境变量说明
| 变量 | 默认值 | 说明 |
|------|--------|------|
| `DATABASE_URL` | `sqlite:///./data/app.db` | 数据库路径(容器内) |
| `AES_KEY` | `1234567890123456` | 法网 AES 加密密钥 |
| `AES_IV` | `1234567890123456` | 法网 AES 初始化向量 |
| `SECRET_KEY` | `change-me-in-production` | JWT 签名密钥 |
| `GRAB_INTERVAL_MS` | `300` | 抢单轮询间隔 (ms) |
| `HEARTBEAT_INTERVAL_S` | `60` | 心跳保活间隔 (s) |
| `ADMIN_USERNAME` | `admin` | 管理员用户名 |
| `ADMIN_PASSWORD` | `Admin@2026` | 管理员密码 |
---
## 防火墙配置
```bash
# 仅开放 80 端口
ufw allow 80/tcp
ufw enable
```
如果需要通过 HTTPS 访问,建议在前面加一个反向代理(如 Nginx + Let's Encrypt)。
---
## 故障排查
### 前端 404
```bash
# 检查 Nginx 配置是否正确转发到后端
docker exec law_order_frontend cat /etc/nginx/conf.d/default.conf
```
### 后端启动失败
```bash
# 查看启动日志
docker logs law_order_backend
# 常见原因:.env 文件不存在、权限不足、端口冲突
```
### 抢单不工作
```bash
# 1. 确认抢单引擎已启动
curl http://localhost/api/fw/grab/start
# 2. 查看后端日志
docker logs -f law_order_backend | grep -i grab
# 3. 检查账号会话
curl http://localhost/api/fw/sessions
```
### 账号被法网封禁
```bash
# 在 Web 界面退出被限制的账号,等一段时间后再登录
# 或手动断开连接:
curl -X POST http://localhost/api/fw/sessions/disconnect \
-H "Content-Type: application/json" \
-d '{"account_id": <id>}'
```
# 法律抢单系统
中国法网智能业务办理系统
## 技术栈
- **后端**: FastAPI + SQLAlchemy + SQLite
- **前端**: React (待开发)
- **部署**: Docker (待配置)
## 项目结构
```
法律抢单系统/
├── backend/
│ ├── main.py # 主入口
│ ├── app/
│ │ ├── api/ # API路由
│ │ │ ├── login.py # 中国法网登录
│ │ │ ├── accounts.py # 账号管理
│ │ │ ├── sessions.py # 会话管理
│ │ │ └── orders.py # 订单管理(抢单/解答)
│ │ ├── core/ # 核心配置
│ │ │ ├── config.py # 系统配置
│ │ │ ├── crypto.py # AES加密工具
│ │ │ └── database.py # 数据库配置
│ │ ├── models/ # 数据模型
│ │ │ └── models.py
│ │ ├── schemas/ # Pydantic schemas
│ │ │ └── schemas.py
│ │ └── services/ # 业务逻辑(待扩展)
│ ├── data/ # SQLite数据库目录
│ ├── requirements.txt
│ └── run.bat # 启动脚本
├── frontend/ # React前端(待开发)
├── api_analysis.json # API逆向分析结果
└── 开发计划.txt
```
## 快速启动
### 1. 安装依赖
```bash
cd backend
pip install -r requirements.txt
```
### 2. 启动后端
```bash
# Windows
run.bat
# Linux/Mac
cd backend
python -m uvicorn main:app --host 0.0.0.0 --port 8080 --reload
```
### 3. 访问API文档
启动后访问: http://localhost:8080/docs
## API接口
### 登录
```
POST /api/fw/login
Body: {"phone": "18917209703", "password": "200663Ai@"}
```
### 账号管理
```
GET /api/accounts # 账号列表
POST /api/accounts # 创建账号
PUT /api/accounts/{id}/password # 重置密码
DELETE /api/accounts/{id} # 删除账号
```
### 会话管理
```
GET /api/sessions # 会话列表
POST /api/sessions/login # 登录账号
POST /api/sessions/{id}/disconnect # 断开
```
### 订单管理
```
GET /api/orders/pending # 待处理订单
POST /api/orders/grab # 抢单
POST /api/orders/answer # 解答
```
## 开发进度
### Phase 1: 逆向分析 ✅ DONE
- [x] 密码加密算法破解 (AES-CBC ZeroPadding)
- [x] 登录API分析
- [x] 抢单API分析
- [x] 解答API分析
- [x] 订单结构分析
### Phase 2: 后端基础 ✅ DONE
- [x] FastAPI框架搭建
- [x] 账号管理CRUD
- [x] 中国法网登录
- [x] 会话管理(基础)
- [x] 订单管理(基础)
### Phase 3: 会话管理 (TODO)
- [ ] 多账号并发登录
- [ ] 心跳保活机制
- [ ] 会话状态监控
- [ ] 反拉黑策略
### Phase 4: 抢单引擎 (TODO)
- [ ] 定时轮询抢单
- [ ] 多账号并发抢单
- [ ] 抢单/处理状态绑定
- [ ] 抢单频率配置
### Phase 5: 前端React (TODO)
- [ ] 普通用户界面
- [ ] 管理员Dashboard
- [ ] 订单处理页面
### Phase 6: Docker部署 (TODO)
- [ ] Dockerfile编写
- [ ] docker-compose配置
- [ ] Ubuntu 26.04部署测试
## 账号列表
| 手机号 | 密码 |
|--------|------|
| 18917209703 | 200663Ai@ |
| 19542820841 | Fabao@123 |
| 13764780450 | Fabao@123 |
| 13661598223 | Bst202606+ |
| 15800562587 | Fabao@1234 |
{
"phase": "Phase 1 - 逆向分析完成",
"encryption": {
"algorithm": "AES-CBC",
"key": "1234567890123456",
"iv": "1234567890123456",
"padding": "ZeroPadding",
"output": "Base64"
},
"base_url": "http://sys.12348.gov.cn/resources",
"apis": {
"login": {
"url": "/api/login",
"method": "POST"
},
"query_order": {
"url": "/api/message_order/query_order",
"method": "POST"
},
"grap_order": {
"url": "/api/message_order/grap_order",
"method": "POST"
},
"answer": {
"url": "/api/message_order/answer",
"method": "POST"
},
"query_current_service": {
"url": "/api/message_order/query_current_service",
"method": "POST"
},
"query_history_service": {
"url": "/api/message_order/query_history_service",
"method": "POST"
},
"get_order_detail": {
"url": "/api/message_order/get_message_order_by_id",
"method": "GET"
},
"get_order_for_answer": {
"url": "/api/message_order/get_message_order_by_id_for_answer",
"method": "GET"
},
"cancel_grap_order": {
"url": "/api/message_order/cancel_grap_order",
"method": "POST"
},
"report": {
"url": "/api/message_order/report",
"method": "POST"
}
},
"accounts": [
{
"phone": "18917209703",
"password": "200663Ai@",
"userId": 2059,
"name": "周晓"
},
{
"phone": "19542820841",
"password": "Fabao@123"
},
{
"phone": "13764780450",
"password": "Fabao@123"
},
{
"phone": "13661598223",
"password": "Bst202606+"
},
{
"phone": "15800562587",
"password": "Fabao@1234"
}
]
}
\ No newline at end of file
差异被折叠。
__pycache__
*.pyc
*.pyo
.env
.venv
node_modules
tests
*.md
.git
# Legal service API config
LEGAL_BASE_URL=http://sys.12348.gov.cn/resources
LEGAL_API_PREFIX=/resources/api
# AES encryption config
AES_KEY=1234567890123456
AES_IV=1234567890123456
# App config
APP_NAME=法律抢单系统
APP_ENV=development
SECRET_KEY=change-me-in-production
# Database
DATABASE_URL=sqlite:///./legal_system.db
# Server
HOST=0.0.0.0
PORT=8000
# Backend — Python 3.11 + FastAPI + uvicorn
FROM python:3.11-slim
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy source
COPY main.py main.py
COPY app/ app/
COPY .env .env
# Create directories
RUN mkdir -p /app/data /app/logs
EXPOSE 8080
# Run — no --reload, single worker
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8080"]
"""账号管理"""
from typing import List
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from app.core.database import get_db
from app.models.models import Account
from app.schemas.schemas import AccountCreate, AccountResponse
from app.api.auth_api import require_admin_dep
router = APIRouter()
@router.get("/accounts")
def list_accounts(
skip: int = 0,
limit: int = 100,
current_user = None, # Use require_admin_dep in routes
db: Session = Depends(get_db)
):
"""获取账号列表"""
accounts = db.query(Account).offset(skip).limit(limit).all()
result = []
for a in accounts:
result.append({
"id": a.id,
"phone": a.phone,
"role": a.role,
"is_active": a.is_active,
"name": a.name,
"user_id": a.user_id,
"group_name": a.group_name,
"created_at": str(a.created_at) if a.created_at else None
})
return result
@router.post("/accounts")
def create_account(
account: AccountCreate,
current_user = None,
db: Session = Depends(get_db)
):
"""创建账号"""
existing = db.query(Account).filter(Account.phone == account.phone).first()
if existing:
raise HTTPException(status_code=400, detail="手机号已存在")
db_account = Account(
phone=account.phone,
password=account.password,
user_id=account.user_id,
group_name=account.group_name
)
db.add(db_account)
db.commit()
db.refresh(db_account)
return {
"id": db_account.id,
"phone": db_account.phone,
"role": db_account.role,
"is_active": db_account.is_active,
"name": db_account.name,
"user_id": db_account.user_id,
"created_at": str(db_account.created_at) if db_account.created_at else None
}
@router.put("/accounts/{account_id}/password")
def reset_password(
account_id: int,
new_password: str = None,
current_user = None,
db: Session = Depends(get_db)
):
"""重置账号密码"""
account = db.query(Account).filter(Account.id == account_id).first()
if not account:
raise HTTPException(status_code=404, detail="账号不存在")
if new_password:
account.password = new_password
db.commit()
return {"msg": "密码重置成功"}
@router.delete("/accounts/{account_id}")
def delete_account(
account_id: int,
current_user = None,
db: Session = Depends(get_db)
):
"""删除账号"""
account = db.query(Account).filter(Account.id == account_id).first()
if not account:
raise HTTPException(status_code=404, detail="账号不存在")
db.delete(account)
db.commit()
return {"msg": "账号已删除"}
@router.put("/accounts/{account_id}/group")
def update_group(
account_id: int,
group_data: dict = None,
current_user = None,
db: Session = Depends(get_db)
):
"""更新账号分组"""
account = db.query(Account).filter(Account.id == account_id).first()
if not account:
raise HTTPException(status_code=404, detail="账号不存在")
group_name = group_data.get("group_name") if group_data else None
account.group_name = group_name
db.commit()
return {"msg": "分组已更新", "account": {
"id": account.id,
"phone": account.phone,
"group_name": account.group_name
}}
@router.get("/accounts/groups")
def list_groups(
current_user = None,
db: Session = Depends(get_db)
):
"""获取所有分组列表"""
groups = db.query(Account.group_name).filter(Account.group_name.isnot(None)).distinct().all()
result = []
for g in groups:
count = db.query(Account).filter(Account.group_name == g.group_name).count()
result.append({"group_name": g.group_name, "account_count": count})
return result
@router.post("/accounts/batch-group")
def batch_update_group(
group_data: dict = None,
current_user = None,
db: Session = Depends(get_db)
):
"""批量更新账号分组"""
account_ids = group_data.get("account_ids", []) if group_data else []
group_name = group_data.get("group_name") if group_data else None
for account_id in account_ids:
account = db.query(Account).filter(Account.id == account_id).first()
if account:
account.group_name = group_name
db.commit()
return {"msg": f"已更新 {len(account_ids)} 个账号分组"}
差异被折叠。
"""抢单引擎控制 API"""
import asyncio
from fastapi import APIRouter, Depends, HTTPException
from app.services.grab_engine import grab_engine
from app.services.session_manager import session_manager
from app.core.config import settings
from pydantic import BaseModel
router = APIRouter()
class GrabConfig(BaseModel):
interval_ms: int = 500
@router.post("/fw/grab/start")
async def start_grab(config: GrabConfig = None):
"""启动抢单引擎"""
if config:
grab_engine.interval = config.interval_ms / 1000.0
if not session_manager.get_all_sessions():
return {"code": "0", "msg": "没有活跃的会话,请先添加账号"}
if grab_engine.running:
return {"code": "0", "msg": "抢单引擎已在运行中"}
# Start in background using asyncio.create_task
grab_engine.task = asyncio.create_task(grab_engine.start())
return {
"code": "1",
"msg": "抢单引擎已启动",
"interval_ms": int(grab_engine.interval * 1000),
"sessions": len(session_manager.get_all_sessions())
}
@router.post("/fw/grab/stop")
async def stop_grab():
"""停止抢单引擎"""
grab_engine.stop()
return {"code": "1", "msg": "抢单引擎已停止"}
@router.get("/fw/grab/status")
def grab_status():
"""获取抢单状态"""
stats = grab_engine.get_stats()
stats["sessions"] = session_manager.get_all_sessions()
return stats
@router.get("/fw/grab/results")
def grab_results(limit: int = 20):
"""获取抢单结果"""
results = grab_engine.results[-limit:]
return {"code": "1", "results": results, "total": grab_engine.total_grabbed}
@router.post("/fw/grab/once")
def grab_once():
"""立即执行一次抢单"""
results = []
for account_id, session in session_manager.sessions.items():
if session.status != "connected":
continue
orders = session.query_order()
if orders:
result = session.grab_order(orders[0].get("id"))
results.append({
"account_id": account_id,
"phone": session.account.phone,
"order_id": orders[0].get("id"),
"result": result
})
return {"code": "1", "results": results}
"""登录模块 - 中国法网登录"""
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from app.core.config import settings
from app.core.crypto import aes_encrypt
import requests
router = APIRouter()
class LoginRequest(BaseModel):
phone: str
password: str
@router.post("/fw/login")
async def login_fw(req: LoginRequest):
"""登录中国法网"""
encrypted_pwd = aes_encrypt(req.password)
login_data = {
"username": req.phone,
"password": encrypted_pwd,
"terrace": 0,
"userAgent": "M5.0"
}
headers = {
"Content-Type": "application/json",
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
}
try:
session = requests.Session()
session.headers.update(headers)
resp = session.post(f"{settings.BASE_URL}/api/login", json=login_data, timeout=10)
data = resp.json()
if data.get("code") in ("1", 1):
return {
"code": "1",
"msg": "登录成功",
"data": data.get("data"),
"cookies": dict(session.cookies)
}
else:
return {
"code": data.get("code", "3"),
"msg": data.get("msg", "登录失败")
}
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
差异被折叠。
from fastapi import APIRouter, Depends, HTTPException, BackgroundTasks
from sqlalchemy.orm import Session
from typing import List, Dict
from ..database import get_db, init_db
from ..models.models import LegalAccount, User, Order, UserRole, SessionStatus
from ..schemas.schemas import (
UserCreate, UserLogin, UserOut,
LegalAccountCreate, LegalAccountOut, LegalAccountLogin, LegalAccountBatchLogin,
OrderOut, OrderReply, DashboardStat, AccountStat
)
from ..services.auth import (
create_admin_if_not_exists, register_user, authenticate_user,
add_user_as_admin, reset_user_password
)
from ..services.legal_session import SessionManager
from ..services.order_grabber import OrderGrabber
from ..utils.crypto import aes_encrypt
import json
router = APIRouter()
# Global session manager (will be initialized in main.py)
session_manager: SessionManager = None
order_grabber: OrderGrabber = None
# ---- User endpoints ----
@router.post("/api/users/register", response_model=UserOut)
def register(user_data: UserCreate, db: Session = Depends(get_db)):
try:
return register_user(db, user_data)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/api/users/login", response_model=UserOut)
def login(login_data: UserLogin, db: Session = Depends(get_db)):
try:
return authenticate_user(db, login_data)
except ValueError as e:
raise HTTPException(status_code=401, detail=str(e))
@router.post("/api/admin/users/add", response_model=UserOut)
def admin_add_user(phone: str, password: str, db: Session = Depends(get_db)):
"""Admin adds a user"""
try:
return add_user_as_admin(db, phone, password)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/api/admin/users/reset-password")
def admin_reset_password(phone: str, password: str, db: Session = Depends(get_db)):
"""Admin resets user password"""
try:
reset_user_password(db, phone, password)
return {"message": "密码重置成功"}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.get("/api/users/list", response_model=List[UserOut])
def list_users(db: Session = Depends(get_db)):
users = db.query(User).all()
return users
# ---- Legal Account endpoints ----
@router.post("/api/legal-accounts/add", response_model=LegalAccountOut)
def add_legal_account(account: LegalAccountCreate, db: Session = Depends(get_db)):
"""Add a legal platform account"""
existing = db.query(LegalAccount).filter(LegalAccount.phone == account.phone).first()
if existing:
raise HTTPException(status_code=400, detail="账号已存在")
legal_account = LegalAccount(
phone=account.phone,
password=aes_encrypt(account.password),
raw_password=account.password,
status=SessionStatus.DISCONNECTED,
)
db.add(legal_account)
db.commit()
db.refresh(legal_account)
return legal_account
@router.post("/api/legal-accounts/login")
def legal_account_login(data: LegalAccountLogin, db: Session = Depends(get_db)):
"""Login a legal account to the platform"""
if not session_manager:
raise HTTPException(status_code=500, detail="Session manager not initialized")
try:
session = session_manager.login_account(data.account_id)
return {
"message": "登录成功",
"account_id": data.account_id,
"status": session.account.status.value
}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/api/legal-accounts/batch-login")
def legal_accounts_batch_login(data: LegalAccountBatchLogin, db: Session = Depends(get_db)):
"""Batch login legal accounts"""
if not session_manager:
raise HTTPException(status_code=500, detail="Session manager not initialized")
results = session_manager.batch_login(data.account_ids)
return {"results": results}
@router.post("/api/legal-accounts/logout")
def legal_account_logout(account_id: int, db: Session = Depends(get_db)):
"""Logout a legal account"""
if session_manager:
session_manager.logout_account(account_id)
return {"message": "已退出"}
@router.get("/api/legal-accounts/list", response_model=List[LegalAccountOut])
def list_legal_accounts(db: Session = Depends(get_db)):
accounts = db.query(LegalAccount).all()
return accounts
# ---- Order endpoints ----
@router.get("/api/orders/pending", response_model=List[OrderOut])
def list_pending_orders(db: Session = Depends(get_db)):
"""List pending orders"""
orders = db.query(Order).filter(Order.status == "pending").order_by(Order.grab_time.desc()).all()
return orders
@router.get("/api/orders/all", response_model=List[OrderOut])
def list_all_orders(db: Session = Depends(get_db)):
"""List all orders"""
orders = db.query(Order).order_by(Order.grab_time.desc()).all()
return orders
@router.post("/api/orders/reply")
def reply_order(data: OrderReply, db: Session = Depends(get_db)):
"""Reply to an order"""
order = db.query(Order).filter(Order.id == data.order_id).first()
if not order:
raise HTTPException(status_code=404, detail="订单不存在")
# Reply via platform
if session_manager and order.assigned_account_id in session_manager.sessions:
session = session_manager.sessions[order.assigned_account_id]
result = session.reply_order(order.order_no, data.reply_content)
if result and result.get("code") == 1:
order.reply_content = data.reply_content
order.status = "completed"
order.complete_time = datetime.now()
# Update account stats
account = session.account
account.processed_count += 1
account.current_order_id = None
account.order_count = account.order_count # keep total
db.commit()
return {"message": "回复成功", "result": result}
# Fallback: save locally
order.reply_content = data.reply_content
order.status = "completed"
order.complete_time = datetime.now()
db.commit()
return {"message": "已保存回复"}
@router.post("/api/orders/{order_id}/assign")
def assign_order(order_id: int, user_id: int, db: Session = Depends(get_db)):
"""Assign order to a user"""
order = db.query(Order).filter(Order.id == order_id).first()
if not order:
raise HTTPException(status_code=404, detail="订单不存在")
order.assigned_user_id = user_id
order.status = "processing"
db.commit()
return {"message": "分配成功"}
# ---- Dashboard endpoints ----
@router.get("/api/dashboard/stats", response_model=DashboardStat)
def get_dashboard_stats(db: Session = Depends(get_db)):
"""Get dashboard statistics"""
from datetime import datetime, timedelta
today = datetime.now().date()
total_accounts = db.query(LegalAccount).count()
connected = db.query(LegalAccount).filter(
LegalAccount.status == SessionStatus.CONNECTED
).count()
blocked = db.query(LegalAccount).filter(
LegalAccount.status == SessionStatus.BLOCKED
).count()
# Today's orders
today_start = datetime.combine(today, datetime.min.time())
total_today = db.query(Order).filter(
Order.grab_time >= today_start
).count()
pending = db.query(Order).filter(Order.status == "pending").count()
processed = db.query(Order).filter(Order.status == "completed").count()
return DashboardStat(
total_accounts=total_accounts,
connected_accounts=connected,
blocked_accounts=blocked,
total_orders_today=total_today,
pending_orders=pending,
processed_orders=processed
)
@router.get("/api/dashboard/accounts", response_model=List[AccountStat])
def get_account_stats(db: Session = Depends(get_db)):
"""Get account statistics"""
accounts = db.query(LegalAccount).all()
stats = []
for acc in accounts:
stats.append(AccountStat(
account_id=acc.id,
phone=acc.phone,
status=acc.status.value,
order_count=acc.order_count,
processed_count=acc.processed_count,
heartbeat_at=acc.heartbeat_at
))
return stats
# ---- Grabber control ----
@router.post("/api/grabber/start")
def start_grabber(db: Session = Depends(get_db)):
"""Start the order grabber"""
global order_grabber
if not session_manager:
raise HTTPException(status_code=500, detail="Session manager not initialized")
order_grabber = OrderGrabber(session_manager, db)
# Start in background
import asyncio
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
loop.run_until_complete(order_grabber.start())
return {"message": "抢单引擎已启动"}
@router.post("/api/grabber/stop")
def stop_grabber():
"""Stop the order grabber"""
global order_grabber
if order_grabber:
order_grabber.stop()
order_grabber = None
return {"message": "抢单引擎已停止"}
@router.get("/api/grabber/stats")
def get_grabber_stats():
"""Get grabber statistics"""
if order_grabber:
return {"stats": order_grabber.get_stats()}
return {"stats": []}
"""会话管理 API"""
from fastapi import APIRouter, Depends, HTTPException, BackgroundTasks, Query
from sqlalchemy.orm import Session
from app.core.database import get_db
from app.models.models import Session as SessionModel, Account
from app.schemas.schemas import SessionResponse, SessionAction
from app.core.config import settings
from app.core.crypto import aes_encrypt
from app.services.session_manager import session_manager
from pydantic import BaseModel
import requests
import logging
import json
from datetime import datetime
logger = logging.getLogger(__name__)
router = APIRouter()
class AccountIdRequest(BaseModel):
account_id: int
@router.get("/fw/sessions")
def list_fw_sessions():
"""获取所有法网会话状态"""
return session_manager.get_all_sessions()
@router.post("/fw/sessions/disconnect")
def disconnect_fw_session(req: AccountIdRequest):
"""退出法网会话(仅断开,不删除 — 保留在列表中方便重新登录)"""
if req.account_id in session_manager.sessions:
session = session_manager.sessions[req.account_id]
phone = session.account_phone
session.status = "disconnected"
session.requests_session = None # 释放 cookies,让律师能登录法网
logger.info(f"Account {phone} disconnected (still in session list)")
return {"code": "1", "msg": "已退出", "sessions": session_manager.get_all_sessions()}
return {"code": "0", "msg": "会话不存在"}
@router.post("/fw/sessions/add")
def add_fw_session(req: AccountIdRequest):
"""添加法网账号并登录"""
# 如果已存在,先断开旧会话
if req.account_id in session_manager.sessions:
old_session = session_manager.sessions[req.account_id]
if old_session.requests_session:
pass # 保留旧 session,直接重新 login
success = session_manager.add_account(req.account_id)
if success:
return {"code": "1", "msg": "登录成功", "sessions": session_manager.get_all_sessions()}
return {"code": "0", "msg": "登录失败"}
@router.post("/fw/sessions/remove")
def remove_fw_session(req: AccountIdRequest):
"""移除法网会话"""
session_manager.remove_account(req.account_id)
return {"code": "1", "msg": "会话已移除"}
@router.post("/fw/sessions/heartbeat")
def heartbeat_all():
"""对所有会话发送心跳"""
count = session_manager.heartbeat_all()
sessions = session_manager.get_all_sessions()
return {"code": "1", "msg": f"心跳完成: {count}个会话活跃", "count": count, "sessions": sessions}
@router.get("/sessions", response_model=list[SessionResponse])
def list_sessions(status: str = None, db: Session = Depends(get_db)):
query = db.query(SessionModel)
if status:
query = query.filter(SessionModel.status == status)
return query.all()
@router.post("/sessions/login")
def login_session(req: SessionAction, db: Session = Depends(get_db)):
"""登录指定账号"""
account = db.query(Account).filter(Account.id == req.account_id).first()
if not account:
raise HTTPException(status_code=404, detail="账号不存在")
encrypted_pwd = aes_encrypt(account.password)
login_data = {
"username": account.phone,
"password": encrypted_pwd,
"terrace": 0,
"userAgent": "M5.0"
}
headers = {
"Content-Type": "application/json",
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
}
session = requests.Session()
session.headers.update(headers)
resp = session.post(f"{settings.BASE_URL}/api/login", json=login_data, timeout=10)
data = resp.json()
if data.get("code") != "1":
raise HTTPException(status_code=400, detail=data.get("msg", "登录失败"))
user_data = data.get("data", {})
account.user_id = user_data.get("userId")
account.name = user_data.get("userName")
db.commit()
db_session = SessionModel(
account_id=req.account_id,
status="connected",
cookies=json.dumps(dict(session.cookies)),
headers=json.dumps(dict(session.headers)),
last_heartbeat=datetime.now()
)
db.add(db_session)
db.commit()
db.refresh(db_session)
return {"msg": "登录成功", "session": db_session}
@router.post("/sessions/{session_id}/disconnect")
def disconnect_session(session_id: int, db: Session = Depends(get_db)):
db_session = db.query(SessionModel).filter(SessionModel.id == session_id).first()
if not db_session:
raise HTTPException(status_code=404, detail="会话不存在")
db_session.status = "disconnected"
db.commit()
return {"msg": "会话已断开"}
"""系统配置"""
from pydantic_settings import BaseSettings
class Settings(BaseSettings):
BASE_URL: str = "http://sys.12348.gov.cn/resources"
AES_KEY: str = "1234567890123456"
AES_IV: str = "1234567890123456"
DATABASE_URL: str = "sqlite:///./data/app.db"
GRAB_INTERVAL_MS: int = 300
ADMIN_USERNAME: str = "admin"
ADMIN_PASSWORD: str = "Admin@2026"
HEARTBEAT_INTERVAL_S: int = 30
model_config = {
"env_file": ".env",
"extra": "ignore"
}
settings = Settings()
"""AES 加密解密工具"""
from Crypto.Cipher import AES
import base64
from app.core.config import settings
def aes_encrypt(plaintext: str) -> str:
"""AES-CBC ZeroPadding 加密,返回 Base64"""
key = settings.AES_KEY.encode("utf-8")
iv = settings.AES_IV.encode("utf-8")
data = plaintext.encode("utf-8")
# ZeroPadding
padding_len = 16 - (len(data) % 16)
padded = data + b"\x00" * padding_len
cipher = AES.new(key, AES.MODE_CBC, iv)
ciphertext = cipher.encrypt(padded)
return base64.b64encode(ciphertext).decode("utf-8")
def aes_decrypt(encrypted: str) -> str:
"""AES-CBC ZeroPadding 解密"""
key = settings.AES_KEY.encode("utf-8")
iv = settings.AES_IV.encode("utf-8")
data = base64.b64decode(encrypted)
cipher = AES.new(key, AES.MODE_CBC, iv)
padded = cipher.decrypt(data)
return padded.rstrip(b"\x00").decode("utf-8")
"""数据库配置"""
from sqlalchemy import create_engine
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
from app.core.config import settings
engine = create_engine(
settings.DATABASE_URL,
connect_args={"check_same_thread": False} # SQLite
)
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine)
Base = declarative_base()
def get_db():
db = SessionLocal()
try:
yield db
finally:
db.close()
"""数据模型"""
from sqlalchemy import Column, Integer, String, DateTime, Boolean, Float, Text, ForeignKey
from sqlalchemy.sql import func
from sqlalchemy.orm import relationship
from app.core.database import Base
from enum import Enum
class UserRole(str, Enum):
ADMIN = "admin"
USER = "user"
class User(Base):
"""系统用户(注册账号)"""
__tablename__ = "users"
id = Column(Integer, primary_key=True, index=True, autoincrement=True)
username = Column(String(50), unique=True, nullable=False, index=True) # 手机号或 admin
password_hash = Column(String(255), nullable=False) # bcrypt 加密
role = Column(String(20), default=UserRole.USER.value) # admin / user
phone = Column(String(20), nullable=True) # 手机号(注册时用)
name = Column(String(50), nullable=True) # 用户名
is_active = Column(Boolean, default=True)
created_at = Column(DateTime, server_default=func.now())
updated_at = Column(DateTime, server_default=func.now(), onupdate=func.now())
# Relationships
accounts = relationship("Account", back_populates="user", foreign_keys="Account.user_id")
@property
def is_admin(self):
return self.role == UserRole.ADMIN.value
class Account(Base):
"""平台账号(中国法网账号)"""
__tablename__ = "accounts"
id = Column(Integer, primary_key=True, index=True, autoincrement=True)
phone = Column(String(20), unique=True, nullable=False, index=True)
password = Column(String(100), nullable=False) # 明文密码,用于加密后登录
user_id = Column(Integer, ForeignKey("users.id"), nullable=True) # 关联用户
name = Column(String(50), nullable=True) # 法网上的名字
role = Column(String(20), default="user") # admin, user
is_active = Column(Boolean, default=True)
group_name = Column(String(50), nullable=True) # 账号分组
created_at = Column(DateTime, server_default=func.now())
updated_at = Column(DateTime, server_default=func.now(), onupdate=func.now())
# Relationships
user = relationship("User", back_populates="accounts", foreign_keys=[user_id])
class Session(Base):
"""登录会话"""
__tablename__ = "sessions"
id = Column(Integer, primary_key=True, index=True, autoincrement=True)
account_id = Column(Integer, nullable=False, index=True)
status = Column(String(20), default="disconnected") # connected, disconnected, blocked, disabled
session_id = Column(String(100), nullable=True) # 法网session
cookies = Column(Text, nullable=True) # JSON字符串
headers = Column(Text, nullable=True) # JSON字符串
grabbed_count = Column(Integer, default=0) # 抢到的单数
answered_count = Column(Integer, default=0) # 已处理的单数
last_heartbeat = Column(DateTime, nullable=True)
created_at = Column(DateTime, server_default=func.now())
updated_at = Column(DateTime, server_default=func.now(), onupdate=func.now())
class Order(Base):
"""订单记录"""
__tablename__ = "orders"
id = Column(Integer, primary_key=True, index=True, autoincrement=True)
order_id = Column(Integer, unique=True, nullable=False, index=True) # 法网订单ID
order_no = Column(String(50), nullable=True) # 订单编号
title = Column(String(200), nullable=True)
consult_content = Column(Text, nullable=True)
answer_content = Column(Text, nullable=True)
answer_title = Column(String(200), nullable=True) # 解答标题
begin_sentences = Column(Text, nullable=True) # 开头语
end_sentences = Column(Text, nullable=True) # 结束语
answer_template_id = Column(Integer, nullable=True) # 解答模板ID
business_id = Column(Integer, nullable=True) # 业务类型ID
order_status = Column(Integer, nullable=True)
status = Column(String(20), default="pending", nullable=True) # pending/grabbed/answering/answered/completed/blocked
user_id = Column(Integer, nullable=True) # 关联用户
session_id = Column(Integer, nullable=True) # 关联的会话
account_id = Column(Integer, nullable=True) # 关联的账号
grabbed_at = Column(DateTime, nullable=True)
answered_at = Column(DateTime, nullable=True)
created_at = Column(DateTime, server_default=func.now())
"""Pydantic schemas"""
from pydantic import BaseModel
from typing import Optional
from datetime import datetime
# User schemas
class UserCreate(BaseModel):
phone: str
password: str
name: Optional[str] = None
class UserLogin(BaseModel):
username: str # 手机号或 admin
password: str
class UserResponse(BaseModel):
id: int
username: str
role: str
is_active: bool
phone: Optional[str] = None
name: Optional[str] = None
created_at: Optional[datetime] = None
model_config = {"from_attributes": True}
class UserUpdate(BaseModel):
password: Optional[str] = None
name: Optional[str] = None
is_active: Optional[bool] = None
# Login response
class LoginResponse(BaseModel):
success: bool
message: str
user: Optional[UserResponse] = None
token: Optional[str] = None
# Account schemas
class AccountCreate(BaseModel):
phone: str
password: str
user_id: Optional[int] = None
group_name: Optional[str] = None
class AccountResponse(BaseModel):
id: int
phone: str
role: str
is_active: bool
group_name: Optional[str] = None
name: Optional[str] = None
user_id: Optional[int] = None
user: Optional[UserResponse] = None
model_config = {"from_attributes": True}
# Session schemas
class SessionResponse(BaseModel):
id: int
account_id: int
status: str
grabbed_count: int
answered_count: int
last_heartbeat: Optional[datetime] = None
created_at: datetime
model_config = {"from_attributes": True}
class SessionAction(BaseModel):
account_id: int
action: str # login, logout, disconnect
# Order schemas
class OrderResponse(BaseModel):
id: int
order_id: int
order_no: Optional[str] = None
title: Optional[str] = None
consult_content: Optional[str] = None
answer_content: Optional[str] = None
order_status: Optional[int] = None
user_id: Optional[int] = None
session_id: Optional[int] = None
account_id: Optional[int] = None
grabbed_at: Optional[datetime] = None
answered_at: Optional[datetime] = None
model_config = {"from_attributes": True}
class OrderAnswer(BaseModel):
order_id: int
answer_content: str
answer_title: Optional[str] = None
"""Answer service - 解答表单提交"""
import asyncio
import json
import time
import logging
from datetime import datetime
from typing import Dict, Any, Optional
from sqlalchemy.orm import Session
from app.core.database import get_db
from app.models.models import Order, Account, Session as FwSession
import requests
logger = logging.getLogger(__name__)
# 中国法网 API 基础 URL
FW_BASE_URL = "http://sys.12348.gov.cn/resources"
class AnswerService:
"""解答服务"""
@staticmethod
def get_order_for_answer(db: Session, account_id: int, order_id: int) -> Dict[str, Any]:
"""获取订单详情用于解答
调用真实系统的 get_message_order_by_id_for_answer API
"""
# Get account and session
account = db.query(Account).filter(Account.id == account_id).first()
if not account:
return {"success": False, "message": "账号不存在"}
# Get active session
session = db.query(FwSession).filter(
FwSession.account_id == account_id,
FwSession.status == "connected"
).first()
if not session:
return {"success": False, "message": "账号未连接"}
# Parse cookies
cookies = json.loads(session.cookies)
# Call real API
try:
resp = requests.post(
f"{FW_BASE_URL}/message_order/get_message_order_by_id_for_answer",
params={"id": order_id},
cookies=cookies,
timeout=10
)
if resp.status_code == 200:
data = resp.json()
if data.get("code") == 1:
order_data = data.get("data", {})
# Update local order
order = db.query(Order).filter(Order.id == order_id).first()
if order:
order.status = "answering"
order.answer_title = order_data.get("answerTitle")
order.begin_sentences = order_data.get("beginSentences")
order.end_sentences = order_data.get("endSentences")
order.answer_content = order_data.get("answerContent")
order.business_id = order_data.get("businessId")
order.answer_template_id = order_data.get("answerTemplateId")
db.commit()
return {
"success": True,
"data": order_data,
"businessTypeList": order_data.get("businessTypeList", []),
"orderTagList": order_data.get("orderTagList", []),
"lawInfoList": order_data.get("lawInfoList", []),
}
return {"success": False, "message": f"API 错误: {resp.text}"}
except Exception as e:
return {"success": False, "message": str(e)}
@staticmethod
def submit_answer(db: Session, account_id: int, answer_data: Dict[str, Any]) -> Dict[str, Any]:
"""提交解答
调用真实系统的 /message_order/answer API
"""
order_id = answer_data.get("id")
if not order_id:
return {"success": False, "message": "订单ID缺失"}
# Get account and session
account = db.query(Account).filter(Account.id == account_id).first()
if not account:
return {"success": False, "message": "账号不存在"}
session = db.query(FwSession).filter(
FwSession.account_id == account_id,
FwSession.status == "connected"
).first()
if not session:
return {"success": False, "message": "账号未连接"}
# Validate answer content
answer_content = answer_data.get("answerContent", "").strip()
if len(answer_content) < 10:
return {"success": False, "message": "解答内容至少10个字"}
# Build request data for real API
fw_data = {
"id": order_id,
"businessId": answer_data.get("businessId"),
"answerContent": answer_content,
"answerTitle": answer_data.get("answerTitle"),
"beginSentences": answer_data.get("beginSentences"),
"endSentences": answer_data.get("endSentences"),
"answerTemplateId": answer_data.get("answerTemplateId"),
"orderTagList": answer_data.get("orderTagList", []),
}
cookies = json.loads(session.cookies)
try:
resp = requests.post(
f"{FW_BASE_URL}/message_order/answer",
json=fw_data,
cookies=cookies,
timeout=15
)
if resp.status_code == 200:
result = resp.json()
if result.get("code") == 1:
# Update local order
order = db.query(Order).filter(Order.id == order_id).first()
if order:
order.status = "answered"
order.answer_content = answer_content
order.answered_at = datetime.now()
db.commit()
return {
"success": True,
"message": "回答成功",
"data": result.get("data")
}
return {"success": False, "message": f"提交失败: {result.get('msg', resp.text)}"}
except Exception as e:
# On failure, mark as error but keep content for retry
order = db.query(Order).filter(Order.id == order_id).first()
if order:
order.status = "answering"
order.answer_content = answer_content
db.commit()
return {"success": False, "message": str(e)}
@staticmethod
def get_answer_templates(db: Session, account_id: int) -> Dict[str, Any]:
"""获取解答模板
调用真实系统的 /basic_data/get_answer_template API
"""
account = db.query(Account).filter(Account.id == account_id).first()
if not account:
return {"success": False, "message": "账号不存在"}
session = db.query(FwSession).filter(
FwSession.account_id == account_id,
FwSession.status == "connected"
).first()
if not session:
return {"success": False, "message": "账号未连接"}
cookies = json.loads(session.cookies)
try:
resp = requests.post(
f"{FW_BASE_URL}/basic_data/get_answer_template",
json={},
cookies=cookies,
timeout=10
)
if resp.status_code == 200:
data = resp.json()
if data.get("code") == 1:
return {"success": True, "data": data.get("data", [])}
return {"success": False, "message": "获取模板失败"}
except Exception as e:
return {"success": False, "message": str(e)}
# Create singleton
answer_service = AnswerService()
from sqlalchemy.orm import Session
from ..models.models import User, UserRole
from ..schemas.schemas import UserCreate, UserLogin, UserOut
from passlib.hash import bcrypt
from datetime import datetime
def create_admin_if_not_exists(db: Session):
"""Create default admin account"""
admin = db.query(User).filter(User.phone == "admin").first()
if not admin:
admin = User(
phone="admin",
password=bcrypt.hash("Admin@2026"),
role=UserRole.ADMIN
)
db.add(admin)
db.commit()
return admin
def register_user(db: Session, user_data: UserCreate) -> UserOut:
"""Register a new user"""
existing = db.query(User).filter(User.phone == user_data.phone).first()
if existing:
raise ValueError("手机号已存在")
user = User(
phone=user_data.phone,
password=bcrypt.hash(user_data.password),
role=UserRole.USER
)
db.add(user)
db.commit()
db.refresh(user)
return UserOut.model_validate(user)
def authenticate_user(db: Session, login: UserLogin) -> UserOut:
"""Authenticate user"""
user = db.query(User).filter(User.phone == login.phone).first()
if not user:
raise ValueError("用户不存在")
if user.phone == "admin" and login.password == "Admin@2026":
# Special case for admin - allow plaintext check
pass
elif not bcrypt.verify(login.password, user.password):
raise ValueError("密码错误")
return UserOut.model_validate(user)
def add_user_as_admin(db: Session, phone: str, password: str) -> UserOut:
"""Admin manually adds a user"""
existing = db.query(User).filter(User.phone == phone).first()
if existing:
raise ValueError("手机号已存在")
user = User(
phone=phone,
password=bcrypt.hash(password),
role=UserRole.USER
)
db.add(user)
db.commit()
db.refresh(user)
return UserOut.model_validate(user)
def reset_user_password(db: Session, phone: str, new_password: str):
"""Admin resets user password"""
user = db.query(User).filter(User.phone == phone).first()
if not user:
raise ValueError("用户不存在")
user.password = bcrypt.hash(new_password)
db.commit()
"""抢单引擎 - 自动轮询抢单"""
import asyncio
import logging
import time
import random
from typing import List
from datetime import datetime
from sqlalchemy.orm import Session
from app.core.database import SessionLocal
from app.models.models import Order, Account
from app.services.session_manager import session_manager
from app.core.config import settings
logger = logging.getLogger(__name__)
class GrabEngine:
"""抢单引擎 - 支持并发、频率控制、状态绑定"""
def __init__(self, interval_ms: int = 500):
self.interval = interval_ms / 1000.0
self.running = False
self.task: asyncio.Task = None
self.results: List[dict] = []
self.total_grabbed = 0
self.grabbed_order_ids = set() # Track grabbed orders to avoid duplicates
async def start(self):
"""启动抢单引擎"""
self.running = True
logger.info(f"Grab engine started (interval: {self.interval}s)")
while self.running:
try:
grabbed_this_round = False
for account_id, session in list(session_manager.sessions.items()):
if session.status != "connected":
continue
# Random delay per account to avoid detection
delay = random.uniform(0.3, self.interval)
await asyncio.sleep(delay)
if not self.running:
break
# Query orders
orders = session.query_order()
if not orders:
continue
logger.info(f"Account {session.account_phone}: found {len(orders)} available order(s)")
for order in orders:
order_id = order.get("id")
if not order_id or order_id in self.grabbed_order_ids:
continue
# Try to grab
result = session.grab_order(order_id)
if result.get("code") == "1":
self.grabbed_order_ids.add(order_id)
self.total_grabbed += 1
grabbed_this_round = True
# Save to database
self._save_order_to_db(order_id, account_id, order)
# 抢单成功 → 断开会话(释放法网登录,管理员手动上线)
if account_id in session_manager.sessions:
sess = session_manager.sessions[account_id]
sess.status = "disconnected"
sess.requests_session = None
# Record result
self.results.append({
"order_id": order_id,
"account_id": account_id,
"phone": session.account_phone,
"title": order.get("title", order.get("content", "")),
"grab_time": time.time(),
"status": "grabbed",
})
logger.info(f"✅ Grabbed order {order_id} by {session.account_phone} — 已断开(管理员请手动上线)")
# Continue to next account
break
else:
msg = result.get("msg", "")
code = result.get("code")
logger.info(f"Failed to grab order {order_id}: code={code} - {msg}")
# Clean up grabbed set if order no longer available
if "不存在" in msg:
self.grabbed_order_ids.discard(order_id)
await asyncio.sleep(max(0.1, self.interval - 0.3))
except asyncio.CancelledError:
logger.info("Grab engine cancelled")
break
except Exception as e:
logger.error(f"Grab engine error: {str(e)}")
await asyncio.sleep(self.interval)
def _save_order_to_db(self, order_id: int, account_id: int, order: dict):
"""保存抢到的订单到数据库"""
db = SessionLocal()
try:
# Check if order already exists (use order_id not order_no)
existing = db.query(Order).filter(Order.order_id == order_id).first()
if existing:
existing.status = "grabbed"
existing.account_id = account_id
existing.grabbed_at = datetime.now()
db.commit()
logger.info(f"Order {order_id} status updated to grabbed")
return
# Create new order with correct field names
new_order = Order(
order_id=order_id,
order_no=str(order_id),
title=order.get("title", ""),
consult_content=(order.get("content") or order.get("consultContent")) or "",
status="grabbed",
account_id=account_id,
grabbed_at=datetime.now(),
)
db.add(new_order)
db.commit()
logger.info(f"Order {order_id} saved to database")
except Exception as e:
logger.error(f"Failed to save order {order_id}: {str(e)}")
db.rollback()
finally:
db.close()
def stop(self):
"""停止抢单引擎"""
self.running = False
if self.task:
self.task.cancel()
logger.info(f"Grab engine stopped. Total grabbed: {self.total_grabbed}")
def get_stats(self) -> dict:
"""获取抢单统计"""
# Count grabbed orders by account
by_account = {}
for r in self.results:
phone = r.get("phone", "unknown")
by_account[phone] = by_account.get(phone, 0) + 1
return {
"running": self.running,
"total_grabbed": self.total_grabbed,
"by_account": by_account,
"recent_results": self.results[-20:], # Last 20
"tracked_orders": len(self.grabbed_order_ids),
}
def clear_old_results(self, keep: int = 50):
"""清理历史结果,保留最近的N条"""
if len(self.results) > keep:
self.results = self.results[-keep:]
def clear_old_grabbed(self):
"""清理过期的已抢订单记录(避免内存泄漏)"""
# Keep only recent 1000 order IDs
if len(self.grabbed_order_ids) > 1000:
self.grabbed_order_ids.clear()
logger.warning("Grabbed order IDs cache cleared to prevent memory leak")
# Global instance
grab_engine = GrabEngine()
import asyncio
import time
import logging
from typing import Dict, List, Optional
from datetime import datetime
from sqlalchemy.orm import Session
from ..models.models import LegalAccount, Order, SessionStatus
from ..services.legal_session import SessionManager
logger = logging.getLogger(__name__)
class OrderGrabber:
"""Concurrent order grabber for multiple accounts"""
def __init__(self, session_manager: SessionManager, db: Session, interval_ms: int = 500):
self.session_manager = session_manager
self.db = db
self.interval = interval_ms / 1000.0 # convert to seconds
self.running = False
self.task = None
async def start(self):
"""Start the order grabbing loop"""
self.running = True
logger.info(f"Order grabber started with {self.interval*1000}ms interval")
while self.running:
try:
await self._grab_cycle()
except Exception as e:
logger.error(f"Grab cycle error: {e}")
await asyncio.sleep(self.interval)
async def _grab_cycle(self):
"""One cycle of grabbing"""
active_sessions = self.session_manager.sessions
for account_id, session in active_sessions.items():
account = session.account
# Skip if account is not connected or already processing an order
if account.status != SessionStatus.CONNECTED:
continue
if account.current_order_id:
# Account is processing an order, skip grabbing
continue
# Query for orders
result = session.query_orders()
if not result or result.get("code") != 1:
continue
orders = result.get("data", [])
if not orders:
continue
# Grab the first available order
for order in orders:
order_id = order.get("id", "")
if not order_id:
continue
grab_result = session.grab_order(order_id)
if grab_result and grab_result.get("code") == 1:
# Successfully grabbed
logger.info(f"Account {account.phone} grabbed order {order_id}")
# Save order to our database
order_record = Order(
order_no=order_id,
title=order.get("title", ""),
content=order.get("content", ""),
status="pending",
assigned_account_id=account_id,
grab_time=datetime.now()
)
self.db.add(order_record)
# Update account
account.order_count += 1
account.current_order_id = order_id
self.db.commit()
break # Only grab one order per cycle
def stop(self):
"""Stop the order grabbing loop"""
self.running = False
logger.info("Order grabber stopped")
def get_stats(self) -> List[dict]:
"""Get grabbing statistics"""
stats = []
for account_id, session in self.session_manager.sessions.items():
account = session.account
stats.append({
"account_id": account_id,
"phone": account.phone,
"status": account.status.value,
"order_count": account.order_count,
"processed_count": account.processed_count,
"current_order_id": account.current_order_id,
"is_grabbing": self.running and account.status == SessionStatus.CONNECTED and not account.current_order_id
})
return stats
差异被折叠。
from Crypto.Cipher import AES
import base64
from ..config import AES_KEY, AES_IV
def aes_encrypt(password: str, key: str = AES_KEY, iv: str = AES_IV) -> str:
"""Encrypt password with AES-CBC ZeroPadding and return Base64"""
plaintext = password.encode("utf-8")
padding_len = 16 - (len(plaintext) % 16)
padded = plaintext + b"\x00" * padding_len
cipher = AES.new(key.encode("utf-8"), AES.MODE_CBC, iv.encode("utf-8"))
return base64.b64encode(cipher.encrypt(padded)).decode("utf-8")
def aes_decrypt(encrypted: str, key: str = AES_KEY, iv: str = AES_IV) -> str:
"""Decrypt Base64 AES-CBC ZeroPadding"""
ciphertext = base64.b64decode(encrypted)
cipher = AES.new(key.encode("utf-8"), AES.MODE_CBC, iv.encode("utf-8"))
padded = cipher.decrypt(ciphertext)
return padded.rstrip(b"\x00").decode("utf-8")
This source diff could not be displayed because it is too large. You can view the blob instead.
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
Markdown 格式
0% 或
您添加了 0 人 到此讨论。请谨慎行事。
请先完成此评论的编辑!
请 注册 或者 后发表评论