From 31b3acb226e0998d585cd9e619946d1d3ee993af Mon Sep 17 00:00:00 2001 From: lpm Date: Sun, 13 Sep 2026 22:21:49 +0800 Subject: [PATCH] =?UTF-8?q?feat:=E5=A2=9E=E5=8A=A0=E5=88=86=E9=A1=B5?= =?UTF-8?q?=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/advisor/audit.py | 2 +- api/advisor/customers.py | 4 ++-- api/advisor/drafts.py | 2 +- api/advisor/todos.py | 2 +- api/advisor/visits.py | 4 ++-- api/routers/advisor_agent.py | 2 +- api/routers/nl2sql.py | 4 +++- api/routers/nl2sql_admin.py | 4 ++-- api/routers/risk.py | 7 ++++--- api/routers/work_order.py | 11 ++++++++--- nl2sql/job_history.py | 20 ++++++++++++++++---- repositories/biz_work_order.py | 23 +++++++++++++++++++++-- repositories/fin_risk_alert.py | 22 +++++++++++++++++++--- service/advisor/agent_client.py | 2 +- service/advisor/audit.py | 10 +++++++--- service/advisor/customers.py | 23 ++++++++++++++--------- service/advisor/drafts.py | 2 +- service/advisor/todos.py | 10 +++++++--- service/advisor/visits.py | 17 +++++++++++------ service/advisor_agent/draft.py | 12 +++++++----- service/product.py | 17 ++++++----------- service/risk/handle.py | 25 ++++++++++++++++++++++--- service/work_order.py | 25 ++++++++++++++++++++++--- 23 files changed, 179 insertions(+), 71 deletions(-) diff --git a/api/advisor/audit.py b/api/advisor/audit.py index 611378d..cea0bd0 100644 --- a/api/advisor/audit.py +++ b/api/advisor/audit.py @@ -22,7 +22,7 @@ async def ledger( start: datetime | None = Query(None), end: datetime | None = Query(None), page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): diff --git a/api/advisor/customers.py b/api/advisor/customers.py index d18e9e2..1fbfcdb 100644 --- a/api/advisor/customers.py +++ b/api/advisor/customers.py @@ -17,7 +17,7 @@ async def list_customers( status: str | None = Query(None, max_length=16), keyword: str | None = Query(None, max_length=64), page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): @@ -51,7 +51,7 @@ async def get_holdings( async def get_reports( customer_id: int, page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): diff --git a/api/advisor/drafts.py b/api/advisor/drafts.py index ff9791a..bb97af8 100644 --- a/api/advisor/drafts.py +++ b/api/advisor/drafts.py @@ -19,7 +19,7 @@ async def list_drafts( customer_id: int | None = Query(None), status: str | None = Query(None, max_length=16), page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): diff --git a/api/advisor/todos.py b/api/advisor/todos.py index 1fba54d..23db5d8 100644 --- a/api/advisor/todos.py +++ b/api/advisor/todos.py @@ -17,7 +17,7 @@ async def list_todos( status: str | None = Query(None, max_length=16), todo_type: str | None = Query(None, max_length=32), page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): diff --git a/api/advisor/visits.py b/api/advisor/visits.py index a92e2e4..e1a69c3 100644 --- a/api/advisor/visits.py +++ b/api/advisor/visits.py @@ -16,7 +16,7 @@ router = APIRouter() async def list_visits( customer_id: int | None = Query(None), page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): @@ -64,7 +64,7 @@ async def talk_templates(user: SysUser = Depends(require_advisor)): async def touch_logs( customer_id: int = Query(...), page: int = Query(1, ge=1), - page_size: int = Query(20, ge=1, le=100), + page_size: int = Query(10, ge=1, le=100), user: SysUser = Depends(require_advisor), db: AsyncSession = Depends(get_db), ): diff --git a/api/routers/advisor_agent.py b/api/routers/advisor_agent.py index 0ad6c48..937736c 100644 --- a/api/routers/advisor_agent.py +++ b/api/routers/advisor_agent.py @@ -322,7 +322,7 @@ async def list_drafts( customer_id: int | None = Query(default=None), status: Literal[DRAFT_STATUS_DRAFT, DRAFT_STATUS_DISCARDED] | None = Query(default=None), page: int = Query(default=1, ge=1), - page_size: int = Query(default=20, ge=1, le=100), + page_size: int = Query(default=10, ge=1, le=100), user: SysUser = Depends(audited_advisor), db: AsyncSession = Depends(get_db), ): diff --git a/api/routers/nl2sql.py b/api/routers/nl2sql.py index 889f780..000ad61 100644 --- a/api/routers/nl2sql.py +++ b/api/routers/nl2sql.py @@ -57,6 +57,7 @@ from tool.llm import llm from utils.exceptions import ForbiddenError, NotFoundError, ParamError from utils.request_id import get_request_id, new_request_id from utils.response import success +from utils.pagination import normalize_pagination router = APIRouter() @@ -474,10 +475,11 @@ async def list_query_history( ): """分页读取当前员工自己的查询历史。""" ensure_query_employee(user) + page, page_size, offset = normalize_pagination(page, page_size) rows = await Nl2SqlPermissionRepo(db).list_query_history( user.id, limit=page_size, - offset=(page - 1) * page_size, + offset=offset, status=status, start_time=start_time, end_time=end_time, diff --git a/api/routers/nl2sql_admin.py b/api/routers/nl2sql_admin.py index 3c62dea..963ea75 100644 --- a/api/routers/nl2sql_admin.py +++ b/api/routers/nl2sql_admin.py @@ -5,7 +5,7 @@ import json import time from datetime import datetime, timedelta, timezone -from fastapi import APIRouter, Depends +from fastapi import APIRouter, Depends, Query from fastapi.responses import PlainTextResponse from sqlalchemy.ext.asyncio import AsyncSession @@ -275,7 +275,7 @@ async def run_admin_job( @router.get("/nl2sql/admin/jobs/history") async def admin_job_history( page: int = 1, - page_size: int = 20, + page_size: int = Query(10, ge=1, le=10), status: str | None = None, user: SysUser = Depends(get_current_user), db: AsyncSession = Depends(get_db), diff --git a/api/routers/risk.py b/api/routers/risk.py index 20ac821..0d5748c 100644 --- a/api/routers/risk.py +++ b/api/routers/risk.py @@ -1,5 +1,5 @@ """风控处置路由:预警列表 + 放行/拦截/冻结(仅风控专员)。""" -from fastapi import APIRouter, Depends +from fastapi import APIRouter, Depends, Query from sqlalchemy.ext.asyncio import AsyncSession from api.deps import require_risk_officer @@ -14,11 +14,12 @@ router = APIRouter(prefix="/risk", tags=["风控"]) @router.get("/alert/list", summary="预警列表") async def list_alerts( status: str | None = None, + page: int = Query(1, ge=1), + page_size: int = Query(10, ge=1, le=10), user: SysUser = Depends(require_risk_officer), db: AsyncSession = Depends(get_db), ): - alerts = await risk_handle.list_alerts(db, status) - return success([a.model_dump(mode="json") for a in alerts]) + return success(await risk_handle.list_alerts(db, status, page=page, page_size=page_size)) @router.post("/alert/{alert_id}/release", summary="放行") diff --git a/api/routers/work_order.py b/api/routers/work_order.py index 4ad730d..cec183f 100644 --- a/api/routers/work_order.py +++ b/api/routers/work_order.py @@ -1,5 +1,5 @@ """业务工单路由:列表 / 详情 / 认领 / 提交审核 / 复核(仅风控专员)。""" -from fastapi import APIRouter, Depends +from fastapi import APIRouter, Depends, Query from sqlalchemy.ext.asyncio import AsyncSession from api.deps import require_risk_officer @@ -15,11 +15,16 @@ router = APIRouter(prefix="/work-order", tags=["工单"]) @router.get("/list", summary="工单列表") async def list_work_orders( status: str | None = None, + page: int = Query(1, ge=1), + page_size: int = Query(10, ge=1, le=10), user: SysUser = Depends(require_risk_officer), db: AsyncSession = Depends(get_db), ): - orders = await work_order_service.list_work_orders(db, status) - return success([o.model_dump(mode="json") for o in orders]) + return success( + await work_order_service.list_work_orders( + db, status, page=page, page_size=page_size + ) + ) @router.get("/{work_order_id}", summary="工单详情") diff --git a/nl2sql/job_history.py b/nl2sql/job_history.py index 750d696..7763c4b 100644 --- a/nl2sql/job_history.py +++ b/nl2sql/job_history.py @@ -7,6 +7,7 @@ from datetime import datetime from sqlalchemy import text from nl2sql.audit import _sanitize +from utils.pagination import normalize_pagination, pagination_result async def record_job_history(db, result, *, elapsed_ms: float, parameter_summary: dict | None = None) -> None: @@ -49,9 +50,15 @@ async def record_job_history_safely(db, result, *, elapsed_ms: float, parameter_ return True -async def list_job_history(db, *, page: int = 1, page_size: int = 20, status: str | None = None) -> list[dict]: +async def list_job_history(db, *, page: int = 1, page_size: int = 10, status: str | None = None) -> dict: """分页查询任务历史,只返回执行摘要,不返回 detail 明细。""" + page, page_size, offset = normalize_pagination(page, page_size) conditions = "WHERE (:status IS NULL OR status = :status)" + count_result = await db.execute( + text(f"SELECT COUNT(*) AS total FROM nl2sql_job_history {conditions}"), + {"status": status}, + ) + total = int(count_result.scalar() or 0) statement = text( f""" SELECT id, job_name, status, attempts, error_type, elapsed_ms, create_time @@ -65,8 +72,13 @@ async def list_job_history(db, *, page: int = 1, page_size: int = 20, status: st statement, { "status": status, - "limit": max(1, min(page_size, 100)), - "offset": max(0, (page - 1) * page_size), + "limit": page_size, + "offset": offset, }, ) - return [dict(row) for row in result.mappings().all()] + return pagination_result( + [dict(row) for row in result.mappings().all()], + total, + page=page, + page_size=page_size, + ) diff --git a/repositories/biz_work_order.py b/repositories/biz_work_order.py index aa645e6..be0acad 100644 --- a/repositories/biz_work_order.py +++ b/repositories/biz_work_order.py @@ -1,7 +1,7 @@ """biz_work_order 工单仓储:查工单 + 流转条件更新(防并发,不 commit)。""" from __future__ import annotations -from sqlalchemy import select, update +from sqlalchemy import func, select, update from model.biz_work_order import BizWorkOrder from repositories.base import BaseRepository @@ -21,6 +21,8 @@ class BizWorkOrderRepo(BaseRepository): handler_id: int | None = None, status: str | None = None, customer_id: int | None = None, + limit: int = 10, + offset: int = 0, ) -> list[BizWorkOrder]: stmt = select(BizWorkOrder) if handler_id is not None: @@ -29,9 +31,26 @@ class BizWorkOrderRepo(BaseRepository): stmt = stmt.where(BizWorkOrder.status == status) if customer_id is not None: stmt = stmt.where(BizWorkOrder.customer_id == customer_id) - stmt = stmt.order_by(BizWorkOrder.id.desc()) + stmt = stmt.order_by(BizWorkOrder.id.desc()).limit(limit).offset(offset) return list((await self.db.scalars(stmt)).all()) + async def count_with_filter( + self, + *, + handler_id: int | None = None, + status: str | None = None, + customer_id: int | None = None, + ) -> int: + """统计工单筛选结果总数,供分页响应使用。""" + stmt = select(func.count()).select_from(BizWorkOrder) + if handler_id is not None: + stmt = stmt.where(BizWorkOrder.handler_id == handler_id) + if status is not None: + stmt = stmt.where(BizWorkOrder.status == status) + if customer_id is not None: + stmt = stmt.where(BizWorkOrder.customer_id == customer_id) + return int((await self.db.scalar(stmt)) or 0) + async def conditional_transition( self, work_order_id: int, *, from_status: str, to_status: str, **fields ) -> bool: diff --git a/repositories/fin_risk_alert.py b/repositories/fin_risk_alert.py index 78233c5..dd880c8 100644 --- a/repositories/fin_risk_alert.py +++ b/repositories/fin_risk_alert.py @@ -3,7 +3,7 @@ from __future__ import annotations from datetime import datetime -from sqlalchemy import select, update +from sqlalchemy import func, select, update from model.fin_risk_alert import FinRiskAlert from repositories.base import BaseRepository @@ -13,16 +13,32 @@ class FinRiskAlertRepo(BaseRepository): model = FinRiskAlert async def list_by_status( - self, status: str | None = None, customer_id: int | None = None + self, + status: str | None = None, + customer_id: int | None = None, + *, + limit: int = 10, + offset: int = 0, ) -> list[FinRiskAlert]: stmt = select(FinRiskAlert) if status is not None: stmt = stmt.where(FinRiskAlert.status == status) if customer_id is not None: stmt = stmt.where(FinRiskAlert.customer_id == customer_id) - stmt = stmt.order_by(FinRiskAlert.id.desc()) + stmt = stmt.order_by(FinRiskAlert.id.desc()).limit(limit).offset(offset) return list((await self.db.scalars(stmt)).all()) + async def count_by_status( + self, status: str | None = None, customer_id: int | None = None + ) -> int: + """统计筛选条件下的预警总数,供后端分页响应使用。""" + stmt = select(func.count()).select_from(FinRiskAlert) + if status is not None: + stmt = stmt.where(FinRiskAlert.status == status) + if customer_id is not None: + stmt = stmt.where(FinRiskAlert.customer_id == customer_id) + return int((await self.db.scalar(stmt)) or 0) + async def conditional_handle( self, alert_id: int, diff --git a/service/advisor/agent_client.py b/service/advisor/agent_client.py index d7e20bd..70fd3f3 100644 --- a/service/advisor/agent_client.py +++ b/service/advisor/agent_client.py @@ -138,7 +138,7 @@ class AdvisorAgentClient: customer_id: int | None = None, status: str | None = None, page: int = 1, - page_size: int = 20, + page_size: int = 10, ) -> dict: params = {"advisor_id": advisor_id, "page": page, "page_size": page_size} if customer_id is not None: diff --git a/service/advisor/audit.py b/service/advisor/audit.py index be67d6b..f9519a6 100644 --- a/service/advisor/audit.py +++ b/service/advisor/audit.py @@ -9,6 +9,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from model.sys_user import SysUser from repositories.audit_log import AuditLogRepo +from utils.pagination import normalize_pagination, pagination_result def _audit_item(a) -> dict: @@ -34,17 +35,20 @@ async def list_ledger( start: datetime | None = None, end: datetime | None = None, page: int = 1, - page_size: int = 20, + page_size: int = 10, ) -> dict: + page, page_size, offset = normalize_pagination(page, page_size) repo = AuditLogRepo(db) items = await repo.list_by_advisor( user_id=user.id, action=action, customer_id=customer_id, keyword=keyword, start=start, end=end, - limit=page_size, offset=(page - 1) * page_size, + limit=page_size, offset=offset, ) total = await repo.count_by_advisor( user_id=user.id, action=action, keyword=keyword, start=start, end=end ) - return {"total": total, "page": page, "page_size": page_size, "items": [_audit_item(a) for a in items]} + return pagination_result( + [_audit_item(a) for a in items], total, page=page, page_size=page_size + ) def _to_csv(rows: list[dict]) -> str: diff --git a/service/advisor/customers.py b/service/advisor/customers.py index 18d8a05..cbda198 100644 --- a/service/advisor/customers.py +++ b/service/advisor/customers.py @@ -28,6 +28,7 @@ from service.advisor.audit_writer import write_audit from service.advisor.masking import mask_name, mask_phone from service.advisor.permissions import ensure_customer_owned, require_owned_relation from utils.exceptions import ForbiddenError, NotFoundError +from utils.pagination import normalize_pagination, pagination_result # 持仓中状态(与 service/holdings.py 口径一致) _HOLDING_STATUS = "持有中" @@ -57,12 +58,13 @@ async def list_customers( status: str | None = None, keyword: str | None = None, page: int = 1, - page_size: int = 20, + page_size: int = 10, ) -> dict: + page, page_size, offset = normalize_pagination(page, page_size) repo = CustomerRelationRepo(db) rows = await repo.list_customer_rows( advisor_id=user.id, status=status, keyword=keyword, - limit=page_size, offset=(page - 1) * page_size, + limit=page_size, offset=offset, ) total = await repo.count_customer_rows(advisor_id=user.id, status=status, keyword=keyword) items = [] @@ -80,7 +82,7 @@ async def list_customers( else None, } ) - return {"total": total, "page": page, "page_size": page_size, "items": items} + return pagination_result(items, total, page=page, page_size=page_size) async def get_customer( @@ -133,19 +135,19 @@ async def get_customer_holdings( async def get_customer_reports( - db: AsyncSession, user: SysUser, customer_id: int, *, page: int = 1, page_size: int = 20 + db: AsyncSession, user: SysUser, customer_id: int, *, page: int = 1, page_size: int = 10 ) -> dict: + page, page_size, offset = normalize_pagination(page, page_size) await ensure_customer_owned(db, user.id, customer_id) repo = AdvisorReportRepo(db) items = await repo.list_by_customer( customer_id, advisor_id=user.id, limit=page_size, - offset=(page - 1) * page_size, + offset=offset, ) - return { - "total": await repo.count_by_advisor(advisor_id=user.id, customer_id=customer_id), - "items": [ + return pagination_result( + [ { "report_id": r.report_id, "intent": r.intent, @@ -155,7 +157,10 @@ async def get_customer_reports( } for r in items ], - } + await repo.count_by_advisor(advisor_id=user.id, customer_id=customer_id), + page=page, + page_size=page_size, + ) async def update_relation( diff --git a/service/advisor/drafts.py b/service/advisor/drafts.py index c57e5a7..d12b259 100644 --- a/service/advisor/drafts.py +++ b/service/advisor/drafts.py @@ -144,7 +144,7 @@ async def list_drafts( customer_id: int | None = None, status: str | None = None, page: int = 1, - page_size: int = 20, + page_size: int = 10, ) -> dict: result = await get_agent_client().draft_list( auth_header=auth_header, diff --git a/service/advisor/todos.py b/service/advisor/todos.py index 9753b80..3399f8e 100644 --- a/service/advisor/todos.py +++ b/service/advisor/todos.py @@ -11,6 +11,7 @@ from model.sys_user import SysUser from repositories.advisor_todo import AdvisorTodoRepo from schemas.advisor import TodoHandleReq from utils.exceptions import ForbiddenError, NotFoundError, ParamError +from utils.pagination import normalize_pagination, pagination_result def _todo_item(t: AdvisorTodo) -> dict: @@ -36,15 +37,18 @@ async def list_todos( status: str | None = None, todo_type: str | None = None, page: int = 1, - page_size: int = 20, + page_size: int = 10, ) -> dict: + page, page_size, offset = normalize_pagination(page, page_size) repo = AdvisorTodoRepo(db) items = await repo.list_by_advisor( advisor_id=user.id, status=status, todo_type=todo_type, - limit=page_size, offset=(page - 1) * page_size, + limit=page_size, offset=offset, ) total = await repo.count_by_advisor(advisor_id=user.id, status=status, todo_type=todo_type) - return {"total": total, "page": page, "page_size": page_size, "items": [_todo_item(t) for t in items]} + return pagination_result( + [_todo_item(t) for t in items], total, page=page, page_size=page_size + ) async def handle_todo(db: AsyncSession, user: SysUser, todo_id: int, req: TodoHandleReq) -> dict: diff --git a/service/advisor/visits.py b/service/advisor/visits.py index f1a0429..8437691 100644 --- a/service/advisor/visits.py +++ b/service/advisor/visits.py @@ -15,6 +15,7 @@ from schemas.advisor import VisitCreateReq, VisitUpdateReq from service.advisor.permissions import ensure_customer_owned from service.advisor.audit_writer import write_audit from utils.exceptions import NotFoundError, ParamError +from utils.pagination import normalize_pagination, pagination_result # 合规话术库(内置,标准化投教/市场解读/调仓沟通)。 # 假设:V1.0 内置常量,后续可迁移 sys_config 运营化;内容不含承诺收益等敏感词。 @@ -40,17 +41,20 @@ def _visit_item(v: AdvisorVisitRecord) -> dict: async def list_visits( db: AsyncSession, user: SysUser, *, customer_id: int | None = None, - page: int = 1, page_size: int = 20, + page: int = 1, page_size: int = 10, ) -> dict: + page, page_size, offset = normalize_pagination(page, page_size) if customer_id is not None: await ensure_customer_owned(db, user.id, customer_id) repo = AdvisorVisitRecordRepo(db) items = await repo.list_by_advisor( advisor_id=user.id, customer_id=customer_id, - limit=page_size, offset=(page - 1) * page_size, + limit=page_size, offset=offset, ) total = await repo.count_by_advisor(advisor_id=user.id, customer_id=customer_id) - return {"total": total, "page": page, "page_size": page_size, "items": [_visit_item(v) for v in items]} + return pagination_result( + [_visit_item(v) for v in items], total, page=page, page_size=page_size + ) async def create_visit(db: AsyncSession, user: SysUser, req: VisitCreateReq) -> dict: @@ -114,19 +118,20 @@ def list_talk_templates() -> list[dict]: async def list_touch_logs( - db: AsyncSession, user: SysUser, customer_id: int, *, page: int = 1, page_size: int = 20 + db: AsyncSession, user: SysUser, customer_id: int, *, page: int = 1, page_size: int = 10 ) -> dict: """触达留痕:站内信(发送触达)+ 回访记录(人工沟通)聚合。 假设:V1.0 触达通道仅站内信与回访;短信/企微/电话统一落在回访记录。 """ + page, page_size, offset = normalize_pagination(page, page_size) await ensure_customer_owned(db, user.id, customer_id) messages = await SysMessageRepo(db).list_by_user( - customer_id, limit=page_size, offset=(page - 1) * page_size + customer_id, limit=page_size, offset=offset ) visits = await AdvisorVisitRecordRepo(db).list_by_advisor( advisor_id=user.id, customer_id=customer_id, - limit=page_size, offset=(page - 1) * page_size, + limit=page_size, offset=offset, ) return { "messages": [ diff --git a/service/advisor_agent/draft.py b/service/advisor_agent/draft.py index 0866e0a..79d50a4 100644 --- a/service/advisor_agent/draft.py +++ b/service/advisor_agent/draft.py @@ -18,6 +18,7 @@ from service.advisor_agent.compliance import ensure_safe_content from repositories.sensitive_word import SensitiveWordRepo from model.advisor_draft import AdvisorDraft from utils.exceptions import ApiError +from utils.pagination import normalize_pagination, pagination_result def build_generated_content(content: str) -> str: @@ -133,18 +134,19 @@ async def list_drafts( customer_id: int | None = None, status: str | None = None, page: int = 1, - page_size: int = 20, + page_size: int = 10, ) -> dict: - page = max(1, page) - page_size = min(100, max(1, page_size)) + page, page_size, offset = normalize_pagination(page, page_size) total, items = await repo.list_drafts( advisor_id=advisor_id, customer_id=customer_id, status=status, limit=page_size, - offset=(page - 1) * page_size, + offset=offset, + ) + return pagination_result( + [summarize_draft(item) for item in items], total, page=page, page_size=page_size ) - return {"total": total, "items": [summarize_draft(item) for item in items]} async def save_draft( diff --git a/service/product.py b/service/product.py index 0bb05b4..24c1993 100644 --- a/service/product.py +++ b/service/product.py @@ -9,6 +9,7 @@ from model.fin_product import FinProduct from repositories.fund_nav import FundNavRepo from repositories.product import ProductRepo from utils.exceptions import NotFoundError, ParamError +from utils.pagination import normalize_pagination, pagination_result def _to_item(p: FinProduct) -> dict: @@ -44,10 +45,7 @@ async def list_products( sort_by: str = "create_time", sort_order: str = "desc", ) -> dict: - if page < 1: - raise ParamError("page 必须 >= 1") - if page_size < 1 or page_size > 100: - raise ParamError("page_size 必须在 1-100 之间") + page, page_size, offset = normalize_pagination(page, page_size, max_page_size=100) if sort_order not in ("asc", "desc"): raise ParamError("sort_order 仅支持 asc/desc") @@ -59,14 +57,11 @@ async def list_products( sort_by=sort_by, sort_order=sort_order, limit=page_size, - offset=(page - 1) * page_size, + offset=offset, + ) + return pagination_result( + [_to_item(p) for p in items], total, page=page, page_size=page_size ) - return { - "total": total, - "page": page, - "page_size": page_size, - "items": [_to_item(p) for p in items], - } def _calc_metrics(rows: list) -> dict: diff --git a/service/risk/handle.py b/service/risk/handle.py index 340948b..7bd4ae0 100644 --- a/service/risk/handle.py +++ b/service/risk/handle.py @@ -24,6 +24,7 @@ from repositories.trade_order import TradeOrderRepo from schemas.risk import RiskAlertResp from service.risk.settle import settle from utils.exceptions import NotFoundError, ParamError +from utils.pagination import normalize_pagination, pagination_result from utils.order_no import gen_order_no # 预警级别 → 工单优先级 @@ -253,6 +254,24 @@ async def freeze(db: AsyncSession, handler: SysUser, alert_id: int) -> dict: raise -async def list_alerts(db: AsyncSession, status: str | None = None) -> list[RiskAlertResp]: - alerts = await FinRiskAlertRepo(db).list_by_status(status) - return [_alert_resp(a) for a in alerts] +async def list_alerts( + db: AsyncSession, + status: str | None = None, + *, + page: int = 1, + page_size: int = 10, +) -> dict: + """分页查询风控预警,默认每页 10 条。""" + page, page_size, offset = normalize_pagination(page, page_size) + repo = FinRiskAlertRepo(db) + alerts = await repo.list_by_status( + status, + limit=page_size, + offset=offset, + ) + return pagination_result( + [_alert_resp(a).model_dump(mode="json") for a in alerts], + await repo.count_by_status(status), + page=page, + page_size=page_size, + ) diff --git a/service/work_order.py b/service/work_order.py index 6bf18a8..84b9a50 100644 --- a/service/work_order.py +++ b/service/work_order.py @@ -15,6 +15,7 @@ from model.sys_user import SysUser from repositories.biz_work_order import BizWorkOrderRepo from schemas.work_order import WorkOrderResp from utils.exceptions import NotFoundError, ParamError +from utils.pagination import normalize_pagination, pagination_result def _audit(db: AsyncSession, user: SysUser, module: str, action: str, target, detail: str) -> None: @@ -49,9 +50,27 @@ def _resp(wo: BizWorkOrder) -> WorkOrderResp: ) -async def list_work_orders(db: AsyncSession, status: str | None = None) -> list[WorkOrderResp]: - orders = await BizWorkOrderRepo(db).list_with_filter(status=status) - return [_resp(o) for o in orders] +async def list_work_orders( + db: AsyncSession, + status: str | None = None, + *, + page: int = 1, + page_size: int = 10, +) -> dict: + """分页查询公共工单,默认每页 10 条。""" + page, page_size, offset = normalize_pagination(page, page_size) + repo = BizWorkOrderRepo(db) + orders = await repo.list_with_filter( + status=status, + limit=page_size, + offset=offset, + ) + return pagination_result( + [_resp(o).model_dump(mode="json") for o in orders], + await repo.count_with_filter(status=status), + page=page, + page_size=page_size, + ) async def get_work_order(db: AsyncSession, work_order_id: int) -> WorkOrderResp: