This commit is contained in:
Your Name
2026-08-07 15:35:02 +08:00
parent 3fc94c4a89
commit 6119fdd767
25 changed files with 2126 additions and 355 deletions
+43 -21
View File
@@ -30,17 +30,26 @@ from models.db_migrate import (
migrate_account_videos_table as _migrate_account_videos_table,
migrate_message_logs_table as _migrate_message_logs_table,
migrate_payment_orders_table as _migrate_payment_orders_table,
migrate_roles_table as _migrate_roles_table,
migrate_rules_table as _migrate_rules_table,
migrate_users_table as _migrate_users_table,
)
from models.db_config import database_config_to_response
from models.models import Account, AccountProfileDetail, AccountVideo, AutoReplyRule, MessageLog, ReceivedMessageLog, SystemLog, User
from auth.router import router as auth_router, users_router
from auth.router import router as auth_router, users_router, roles_router
from auth.settings_router import router as settings_router
from auth.role_service import seed_builtin_roles
from payments.router import router as payments_router
from desktop_router import router as desktop_router
from link_cards_router import router as link_cards_router, UPLOAD_DIR as LINK_CARD_UPLOAD_DIR
from auth.dependencies import get_current_user, require_admin, require_write
from auth.dependencies import (
get_current_user,
require_accounts_write,
require_admin,
require_messages_write,
require_rules_write,
require_write,
)
from auth.account_limits import ensure_can_add_account
from auth.scopes import (
accounts_for_user,
@@ -52,7 +61,8 @@ from auth.scopes import (
received_logs_for_user,
system_logs_for_user,
)
from auth.roles import is_admin
from auth.roles import has_permission, is_admin
from auth.permissions import LOGS_READ, RECEIVED_MESSAGES_READ, SYSTEM_LOGS_READ
from auth.passwords import hash_password
from rpa_engine.batch_start import BatchStartQueue
from rpa_engine.playwright_worker import DouyinWorker
@@ -201,6 +211,7 @@ app.mount("/api/media/link-cards", StaticFiles(directory=LINK_CARD_UPLOAD_DIR),
app.include_router(auth_router)
app.include_router(users_router)
app.include_router(roles_router)
app.include_router(settings_router)
app.include_router(payments_router)
app.include_router(desktop_router)
@@ -610,6 +621,7 @@ async def startup():
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
await conn.run_sync(_migrate_roles_table)
await conn.run_sync(_migrate_accounts_table)
await conn.run_sync(_migrate_rules_table)
await conn.run_sync(_migrate_message_logs_table)
@@ -618,6 +630,8 @@ async def startup():
await conn.run_sync(_migrate_payment_orders_table)
await conn.run_sync(_migrate_accounts_quota_disabled)
await _seed_app_config()
async with AsyncSessionLocal() as db:
await seed_builtin_roles(db)
await _seed_admin_user()
# 进程启动时没有任何内存 Worker;复位异常退出遗留的运行状态。
# 同时,账号数量/并发限制已移除,清理历史“额度停用”标记。
@@ -1618,7 +1632,7 @@ async def update_account(
account_id: int,
body: AccountUpdate,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
account = await get_owned_account(db, user, account_id, write=True)
follow_config_changed = bool(
@@ -1776,7 +1790,7 @@ async def send_account_queued_reply_now(
account_id: int,
job_id: str,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_messages_write),
):
"""将指定任务原子移入紧急队列,并把它后面的普通任务前移一槽。"""
await get_owned_account(db, user, account_id, write=True)
@@ -1931,7 +1945,7 @@ async def update_account_cookie(
account_id: int,
body: AccountCookieUpdate,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
# Validate first: malformed input must not take a healthy hosted account
# offline. Filesystem and database mutations happen only after the worker
@@ -1980,7 +1994,7 @@ async def update_account_cookie(
async def delete_account_cookie(
account_id: int,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
# Deleting credentials uses the same preparation lock as starting a
# worker, preventing a new worker from appearing after stop_worker but
@@ -2009,7 +2023,7 @@ async def delete_account_cookie(
async def create_account(
account_in: AccountCreate,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
cookie_data = (account_in.cookie_data or "").strip()
standard_json_str = None
@@ -2047,7 +2061,7 @@ async def create_account(
async def delete_account(
account_id: int,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
await get_owned_account(db, user, account_id, write=True)
# 停止运行中的任务
@@ -2082,7 +2096,7 @@ async def validate_account_credential(
async def reset_account_credentials(
account_id: int,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
await get_owned_account(db, user, account_id, write=True)
await _release_db_connection(db)
@@ -2237,7 +2251,7 @@ async def start_account_rpa(
account_id: int,
body: StartAccountRequest = StartAccountRequest(),
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
account = await get_owned_account(db, user, account_id, write=True)
# Cancelling waits for an in-flight queued start, and the preparation lock
@@ -2253,7 +2267,7 @@ async def start_account_rpa(
async def submit_account_start_batch(
body: BatchStartRequest,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
requested_ids = list(
dict.fromkeys(int(value) for value in body.account_ids if int(value) > 0)
@@ -2317,7 +2331,7 @@ async def get_account_start_batch(
async def stop_account_rpa(
account_id: int,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_accounts_write),
):
account = await get_owned_account(db, user, account_id, write=True)
@@ -2395,7 +2409,7 @@ async def get_rules(
async def create_rule(
rule_in: RuleCreate,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_rules_write),
):
if rule_in.match_type == "default":
rule_in.keyword = ""
@@ -2429,7 +2443,7 @@ async def update_rule(
rule_in: RuleCreate,
is_active: Optional[bool] = None,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_rules_write),
):
rule = await get_accessible_rule(db, user, rule_id, write=True)
@@ -2457,7 +2471,7 @@ async def update_rule(
async def toggle_rule(
rule_id: int,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_rules_write),
):
rule = await get_accessible_rule(db, user, rule_id, write=True)
@@ -2470,7 +2484,7 @@ async def toggle_rule(
async def delete_rule(
rule_id: int,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_rules_write),
):
await get_accessible_rule(db, user, rule_id, write=True)
await db.execute(delete(AutoReplyRule).where(AutoReplyRule.id == rule_id))
@@ -2487,7 +2501,7 @@ async def move_rule(
rule_id: int,
body: RuleMove,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_rules_write),
):
"""在同账号规则内上移/下移一位(服务端交换排序,适配前端分页)。"""
rule = await get_accessible_rule(db, user, rule_id, write=True)
@@ -2512,7 +2526,7 @@ async def move_rule(
async def reorder_rules(
body: RuleReorder,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_rules_write),
):
if not body.rule_ids:
return {"message": "No rules to reorder."}
@@ -2531,6 +2545,8 @@ async def get_logs_stats(
user: User = Depends(get_current_user),
):
"""消息日志全量统计(数据库计数,不受列表 limit 限制)。"""
if not has_permission(user.role, LOGS_READ):
raise HTTPException(status_code=403, detail="缺少权限:logs.read")
if account_id is not None:
await get_owned_account(db, user, account_id)
# Select only the indexed status column and calculate both counters in one
@@ -2564,6 +2580,8 @@ async def get_logs(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
if not has_permission(user.role, LOGS_READ):
raise HTTPException(status_code=403, detail="缺少权限:logs.read")
if account_id is not None:
await get_owned_account(db, user, account_id)
limit = max(1, min(int(limit or 50), 500))
@@ -2585,6 +2603,8 @@ async def get_received_messages(
user: User = Depends(get_current_user),
):
"""接收消息原始日志:仅包含收到的消息,内容为接口/通道原样记录。"""
if not has_permission(user.role, RECEIVED_MESSAGES_READ):
raise HTTPException(status_code=403, detail="缺少权限:received_messages.read")
if account_id is not None:
await get_owned_account(db, user, account_id)
limit = max(1, min(int(limit or 100), 500))
@@ -2607,6 +2627,8 @@ async def get_system_logs(
user: User = Depends(get_current_user),
):
"""系统诊断日志:私信收发 / 实时连接 / 鉴权 等链路事件,用于排查失败原因。"""
if not has_permission(user.role, SYSTEM_LOGS_READ):
raise HTTPException(status_code=403, detail="缺少权限:system_logs.read")
if account_id is not None:
await get_owned_account(db, user, account_id)
entries = system_logger.get_logs(
@@ -2688,7 +2710,7 @@ async def upload_message_image(
account_id: int,
file: UploadFile = File(...),
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_messages_write),
):
account = await get_owned_account(db, user, account_id, write=True)
if not file.content_type or not file.content_type.startswith("image/"):
@@ -2779,7 +2801,7 @@ async def send_account_message(
account_id: int,
body: SendMessageRequest,
db: AsyncSession = Depends(get_db),
user: User = Depends(require_write),
user: User = Depends(require_messages_write),
):
account = await get_owned_account(db, user, account_id, write=True)