feat: improve audit exports and admin analytics

This commit is contained in:
2026-08-25 18:49:56 +08:00
parent 910e67bd89
commit 5740bd26e7
55 changed files with 1654 additions and 269 deletions

View File

@@ -10,7 +10,9 @@ from sqlalchemy.orm import Session
from app.core.database import get_db
from app.core.dependencies import get_current_admin
from app.core.responses import api_success
from app.core.time_utils import as_utc_naive, business_day_boundary
from app.models.admin import Admin
from app.models.chat import ChatMessage
from app.models.knowledge import (
HumanAttentionHistory,
HumanAttentionRecord,
@@ -95,7 +97,9 @@ def _cleanup_before(payload: dict) -> datetime:
if not raw:
raise HTTPException(status_code=400, detail="必须指定清理日期")
try:
return datetime.fromisoformat(raw.replace("Z", "+00:00")).replace(tzinfo=None)
if len(raw) == 10:
return business_day_boundary(datetime.fromisoformat(raw).date())
return as_utc_naive(datetime.fromisoformat(raw.replace("Z", "+00:00")))
except ValueError as exc:
raise HTTPException(status_code=400, detail="清理日期格式错误") from exc
@@ -151,21 +155,46 @@ def _knowledge_type_label(value: str) -> str:
@router.get("/attention/list")
def attention_list(
priority: str = Query(default=""),
statusValue: str = Query(default="", alias="status"),
priority: str = Query(default="", pattern="^(|urgent|important|normal)$"),
statusValue: str = Query(default="", alias="status", pattern="^(|pending|processing|resolved|ignored)$"),
dateFrom: datetime | None = Query(default=None),
dateTo: datetime | None = Query(default=None),
page: int = Query(default=1, ge=1),
pageSize: int = Query(default=20, ge=10, le=100),
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> dict:
dateFrom = as_utc_naive(dateFrom) if dateFrom is not None else None
dateTo = as_utc_naive(dateTo) if dateTo is not None else None
if dateFrom is not None and dateTo is not None and dateFrom > dateTo:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="开始时间不能晚于结束时间")
query = select(HumanAttentionRecord)
if priority:
query = query.where(HumanAttentionRecord.priority == priority)
if statusValue:
query = query.where(HumanAttentionRecord.status == statusValue)
if dateFrom is not None:
query = query.where(HumanAttentionRecord.created_at >= dateFrom)
if dateTo is not None:
query = query.where(HumanAttentionRecord.created_at <= dateTo)
total = db.scalar(select(func.count()).select_from(query.subquery())) or 0
rows = db.scalars(query.order_by(HumanAttentionRecord.id.desc()).offset((page - 1) * pageSize).limit(pageSize)).all()
return api_success(page_result([_attention_dict(item) for item in rows], total=total, page=page, page_size=pageSize))
message_keys = {
(message_id, session_id)
for message_id, session_id in db.execute(
select(ChatMessage.id, ChatMessage.session_id).where(
ChatMessage.id.in_([item.message_id for item in rows])
)
).all()
} if rows else set()
items = [
_attention_dict(
item,
target_available=(item.message_id, item.session_id) in message_keys,
)
for item in rows
]
return api_success(page_result(items, total=total, page=page, page_size=pageSize))
@router.put("/attention/{attention_id}")
@@ -251,8 +280,8 @@ def _retrieval_dict(item: KnowledgeRetrievalLog, *, detail: bool) -> dict:
return data
def _attention_dict(item: HumanAttentionRecord) -> dict:
return {
def _attention_dict(item: HumanAttentionRecord, *, target_available: bool | None = None) -> dict:
data = {
"id": item.id,
"sessionId": item.session_id,
"messageId": item.message_id,
@@ -267,6 +296,9 @@ def _attention_dict(item: HumanAttentionRecord) -> dict:
"createdAt": item.created_at,
"updatedAt": item.updated_at,
}
if target_available is not None:
data["targetAvailable"] = target_available
return data
def _json(raw: str | None):

View File

@@ -18,6 +18,7 @@ from app.services.admin_service import AdminDashboardService
from app.models.logs import StorageSnapshot
from app.services.redis_client import get_sync_redis_client
from app.services.request_traffic_service import RequestTrafficService
from app.core.time_utils import business_day_boundary
router = APIRouter()
@@ -29,12 +30,8 @@ def dashboard(
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> dict:
start_dt = datetime.strptime(start, "%Y-%m-%d") if start else None
if end:
end_dt = datetime.strptime(end, "%Y-%m-%d")
end_dt = end_dt.replace(hour=23, minute=59, second=59)
else:
end_dt = None
start_dt = business_day_boundary(datetime.strptime(start, "%Y-%m-%d").date()) if start else None
end_dt = business_day_boundary(datetime.strptime(end, "%Y-%m-%d").date(), end_exclusive=True) if end else None
stats = AdminDashboardService.stats(db, start_dt, end_dt)
return api_success(DashboardStats.model_validate(stats).model_dump())

View File

@@ -14,6 +14,7 @@ from sqlalchemy.orm import Session, load_only
from app.core.database import get_db
from app.core.dependencies import get_current_admin
from app.core.responses import api_success
from app.core.time_utils import as_utc_naive, to_business_naive
from app.models.admin import Admin
from app.models.chat import ChatMessage, ChatSession, TopicSession
from app.models.growth import ShareDraft, TeacherHelpCard, TopicSummary
@@ -27,6 +28,7 @@ from app.services.question_insight_export_service import QuestionInsightExportSe
from app.services.growth_profile_service import topic_dict, topic_summary_dict
from app.services.help_card_service import help_card_dict
from app.services.share_draft_service import share_draft_dict
from app.services.chat_export_service import ChatExportService
router = APIRouter()
@@ -45,6 +47,8 @@ def chat_list(
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> dict:
dateFrom = as_utc_naive(dateFrom) if dateFrom is not None else None
dateTo = as_utc_naive(dateTo) if dateTo is not None else None
query = _chat_query(
keyword=keyword,
user_id=userId,
@@ -72,6 +76,8 @@ def export_chats(
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> Response:
dateFrom = as_utc_naive(dateFrom) if dateFrom is not None else None
dateTo = as_utc_naive(dateTo) if dateTo is not None else None
rows = db.execute(
_chat_query(
keyword=keyword,
@@ -98,8 +104,8 @@ def export_chats(
source_client.name if source_client else "千问千答直接访问",
session.title,
session.message_count,
session.last_message_at,
session.updated_at,
to_business_naive(session.last_message_at),
to_business_naive(session.updated_at),
]
)
content = "\ufeff" + output.getvalue()
@@ -132,11 +138,53 @@ def chat_messages(
return api_success(page_result([_message_dict(item) for item in messages], total=total, page=page, page_size=pageSize))
@router.get("/chat/{session_id}/export")
def export_chat_detail(
session_id: int,
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> StreamingResponse:
session_row = db.execute(
select(ChatSession, User, SsoClient)
.join(User, User.id == ChatSession.user_id, isouter=True)
.join(SsoClient, SsoClient.id == ChatSession.source_client_id, isouter=True)
.where(ChatSession.id == session_id)
).first()
if session_row is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="会话不存在")
session, user, source_client = session_row
messages = db.scalars(
select(ChatMessage)
.where(ChatMessage.session_id == session_id)
.order_by(ChatMessage.created_at.asc(), ChatMessage.id.asc())
).all()
stream = ChatExportService.build(
session=session,
user=user,
source_name=source_client.name if source_client else "千问千答直接访问",
messages=messages,
)
OperationLogService.write(
db,
admin_id=current_admin.id,
module="records",
action="export_chat_detail",
target_id=session.id,
)
db.commit()
return StreamingResponse(
stream,
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
headers={"Content-Disposition": f'attachment; filename="chat_session_{session.id}.xlsx"'},
)
@router.get("/chat/{session_id}")
def chat_detail(
session_id: int,
messagePage: int = Query(default=1, ge=1),
messagePageSize: int = Query(default=20, ge=10, le=100),
focusMessageId: int | None = Query(default=None, gt=0),
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> dict:
@@ -154,6 +202,28 @@ def chat_detail(
session, user, source_client = session_row
message_query = select(ChatMessage).where(ChatMessage.session_id == session_id)
message_total = db.scalar(select(func.count()).select_from(message_query.subquery())) or 0
if focusMessageId is not None:
focused_message = db.scalar(
select(ChatMessage).where(
ChatMessage.id == focusMessageId,
ChatMessage.session_id == session_id,
)
)
if focused_message is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="触发消息不属于该会话或已不存在")
focused_position = db.scalar(
select(func.count()).select_from(ChatMessage).where(
ChatMessage.session_id == session_id,
or_(
ChatMessage.created_at < focused_message.created_at,
(
(ChatMessage.created_at == focused_message.created_at)
& (ChatMessage.id <= focused_message.id)
),
),
)
) or 1
messagePage = (focused_position - 1) // messagePageSize + 1
messages = db.scalars(
message_query
.order_by(ChatMessage.created_at.asc(), ChatMessage.id.asc())
@@ -198,6 +268,7 @@ def chat_detail(
"topics": topics,
"helpCards": [help_card_dict(item) for item in help_cards],
"shareDrafts": [share_draft_dict(item) for item in share_drafts],
"focusedMessageId": focusMessageId,
}
)
@@ -288,6 +359,8 @@ def question_insights(
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> dict:
dateFrom = as_utc_naive(dateFrom) if dateFrom is not None else None
dateTo = as_utc_naive(dateTo) if dateTo is not None else None
return api_success(
QuestionInsightService.summarize(
db,
@@ -309,6 +382,8 @@ def refresh_question_insights(
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> dict:
dateFrom = as_utc_naive(dateFrom) if dateFrom is not None else None
dateTo = as_utc_naive(dateTo) if dateTo is not None else None
result = QuestionInsightService.refresh(
db,
date_from=dateFrom,
@@ -334,6 +409,8 @@ def export_question_insights(
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> StreamingResponse:
dateFrom = as_utc_naive(dateFrom) if dateFrom is not None else None
dateTo = as_utc_naive(dateTo) if dateTo is not None else None
if dateFrom is not None and dateTo is not None and dateFrom > dateTo:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="开始时间不能晚于结束时间")
result = QuestionInsightService.summarize(
@@ -387,9 +464,9 @@ def _chat_query(
if source_client_id is not None:
query = query.where(ChatSession.source_type == "sso", ChatSession.source_client_id == source_client_id)
if date_from is not None:
query = query.where(ChatSession.updated_at >= date_from.replace(tzinfo=None))
query = query.where(ChatSession.updated_at >= date_from)
if date_to is not None:
query = query.where(ChatSession.updated_at <= date_to.replace(tzinfo=None))
query = query.where(ChatSession.updated_at <= date_to)
if status:
query = query.where(
exists()

View File

@@ -3,12 +3,15 @@ from __future__ import annotations
from datetime import date
from fastapi import APIRouter, Depends, HTTPException, Query, status
from fastapi.responses import Response
from sqlalchemy.orm import Session
from app.core.database import get_db
from app.core.dependencies import get_current_admin
from app.core.responses import api_success
from app.models.admin import Admin
from app.services.admin_service import OperationLogService
from app.services.user_behavior_export_service import UserBehaviorExportService
from app.services.user_behavior_service import UserBehaviorService
@@ -30,6 +33,8 @@ def behavior_users(
start: date | None = Query(default=None),
end: date | None = Query(default=None),
keyword: str = Query(default="", max_length=50),
sortBy: str = Query(default="lastEventAt", pattern="^(userName|phone|eventCount|lastEventAt)$"),
sortOrder: str = Query(default="desc", pattern="^(asc|desc)$"),
page: int = Query(default=1, ge=1),
pageSize: int = Query(default=20, ge=5, le=100),
db: Session = Depends(get_db),
@@ -43,10 +48,42 @@ def behavior_users(
keyword=keyword,
page=page,
page_size=pageSize,
sort_by=sortBy,
sort_order=sortOrder,
)
)
@router.get("/user-behavior/export")
def export_behavior_report(
start: date | None = Query(default=None),
end: date | None = Query(default=None),
keyword: str = Query(default="", max_length=50),
sortBy: str = Query(default="lastEventAt", pattern="^(userName|phone|eventCount|lastEventAt)$"),
sortOrder: str = Query(default="desc", pattern="^(asc|desc)$"),
db: Session = Depends(get_db),
current_admin: Admin = Depends(get_current_admin),
) -> Response:
stream = UserBehaviorExportService.build_workbook(
db,
start=start,
end=end,
keyword=keyword,
sort_by=sortBy,
sort_order=sortOrder,
)
OperationLogService.write(db, admin_id=current_admin.id, module="behavior", action="export")
db.commit()
start_label = start.isoformat() if start else "default"
end_label = end.isoformat() if end else "today"
filename = f"user_behavior_report_{start_label}_{end_label}.xlsx"
return Response(
content=stream.getvalue(),
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
headers={"Content-Disposition": f'attachment; filename="{filename}"'},
)
@router.get("/user-behavior/user/{user_id}/timeline")
def behavior_timeline(
user_id: int,

View File

@@ -17,6 +17,7 @@ from sqlalchemy.orm import Session
from app.core.database import get_db
from app.core.dependencies import get_current_admin
from app.core.responses import api_success
from app.core.time_utils import as_utc_naive, business_datetime_to_utc_naive, business_day_boundary
from app.models.admin import Admin
from app.models.ai_config import SystemConfig
from app.models.entitlement import EntitlementPlan, UserEntitlement
@@ -475,7 +476,7 @@ def update_user(
if payload.dailyChatLimit is not None:
user.daily_chat_limit = payload.dailyChatLimit
if payload.expiredAt is not None:
user.expired_at = payload.expiredAt.replace(tzinfo=None)
user.expired_at = as_utc_naive(payload.expiredAt)
db.add(user)
if payload.entitlementPlanId is not None:
EntitlementService.assign_user_plan(
@@ -698,7 +699,7 @@ def _apply_user_payload(
user.nickname = payload.nickname.strip() if payload.nickname else None
user.status = payload.status
user.daily_chat_limit = payload.dailyChatLimit if payload.dailyChatLimit is not None else default_daily_limit
user.expired_at = payload.expiredAt.replace(tzinfo=None) if payload.expiredAt is not None else None
user.expired_at = as_utc_naive(payload.expiredAt) if payload.expiredAt is not None else None
def _default_daily_chat_limit(db: Session) -> int:
@@ -763,14 +764,17 @@ def _parse_optional_datetime(value: object | None) -> datetime | None:
if value is None or _cell_text(value) == "":
return None
if isinstance(value, datetime):
return value.replace(tzinfo=None)
return business_datetime_to_utc_naive(value)
if isinstance(value, date):
return datetime.combine(value, datetime.min.time())
return business_day_boundary(value, end_exclusive=True) - timedelta(microseconds=1)
text = _cell_text(value)
for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%d", "%Y/%m/%d %H:%M:%S", "%Y/%m/%d"):
try:
return datetime.strptime(text, fmt)
parsed = datetime.strptime(text, fmt)
if "%H" not in fmt:
return business_day_boundary(parsed.date(), end_exclusive=True) - timedelta(microseconds=1)
return business_datetime_to_utc_naive(parsed)
except ValueError:
continue
raise ValueError("有效期格式应为 YYYY-MM-DD 或 YYYY-MM-DD HH:mm:ss")

View File

@@ -1,6 +1,6 @@
from __future__ import annotations
from datetime import UTC, date, datetime, time, timedelta
from datetime import UTC, date, datetime, timedelta
from io import BytesIO
from fastapi import APIRouter, Depends, HTTPException, Query, status
@@ -18,6 +18,7 @@ from app.api.pagination import page_result
from app.core.database import get_db
from app.core.dependencies import get_current_admin, get_current_user
from app.core.responses import api_success
from app.core.time_utils import business_day_boundary, to_business_naive
from app.models.admin import Admin
from app.models.chat import ChatMessage, ChatSession
from app.models.feedback import MessageFeedback
@@ -136,9 +137,9 @@ def _feedback_query(read_status: str, start_date: date | None, end_date: date |
if read_status != "all":
query = query.where(MessageFeedback.is_read == (1 if read_status == "read" else 0))
if start_date is not None:
query = query.where(MessageFeedback.created_at >= datetime.combine(start_date, time.min))
query = query.where(MessageFeedback.created_at >= business_day_boundary(start_date))
if end_date is not None:
query = query.where(MessageFeedback.created_at < datetime.combine(end_date + timedelta(days=1), time.min))
query = query.where(MessageFeedback.created_at < business_day_boundary(end_date, end_exclusive=True))
return query
@@ -159,8 +160,8 @@ def _feedback_workbook(rows: list[tuple]) -> Workbook:
_excel_safe_text(session.title),
message.id,
_excel_safe_text(ensure_ai_generated_notice(message.content)),
feedback.created_at,
feedback.read_at,
to_business_naive(feedback.created_at),
to_business_naive(feedback.read_at),
])
header_fill = PatternFill("solid", fgColor="1F6F5F")
for cell in sheet[1]:

View File

@@ -2,7 +2,7 @@ from __future__ import annotations
from collections.abc import Generator
from sqlalchemy import create_engine
from sqlalchemy import create_engine, event
from sqlalchemy.orm import Session, sessionmaker
from app.core.config import get_settings
@@ -20,6 +20,16 @@ if not settings.database_url.startswith("sqlite"):
pool_recycle=settings.db_pool_recycle,
)
engine = create_engine(settings.database_url, **engine_options)
if engine.dialect.name == "mysql":
@event.listens_for(engine, "connect")
def _set_mysql_session_timezone(dbapi_connection, _connection_record) -> None:
"""Keep CURRENT_TIMESTAMP and application-written values on one UTC contract."""
cursor = dbapi_connection.cursor()
try:
cursor.execute("SET time_zone = '+00:00'")
finally:
cursor.close()
SessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False, future=True)
register_knowledge_cache_invalidation_events()

View File

@@ -126,7 +126,9 @@ def enforce_admin_access(
else:
permission = "attention.view" if method == "GET" else "attention.edit"
elif path.startswith("user-behavior"):
permission = "behavior.view"
permission = "behavior.export" if path.startswith("user-behavior/export") else "behavior.view"
elif path == "chat/export" or (path.startswith("chat/") and path.endswith("/export")) or path.startswith("question-insights/export"):
permission = "records.export"
else:
permission = "records.view"
require_permission(admin, permission)

View File

@@ -2,10 +2,12 @@ from __future__ import annotations
from typing import Any
from app.core.time_utils import normalize_api_datetimes
def api_success(data: Any = None, message: str = "success") -> dict[str, Any]:
return {"code": 0, "message": message, "data": data}
return {"code": 0, "message": message, "data": normalize_api_datetimes(data)}
def api_error(code: int, message: str, data: Any = None) -> dict[str, Any]:
return {"code": code, "message": message, "data": data}
return {"code": code, "message": message, "data": normalize_api_datetimes(data)}

View File

@@ -0,0 +1,64 @@
from __future__ import annotations
from datetime import UTC, date, datetime, time, timedelta
from typing import Any
from zoneinfo import ZoneInfo
BUSINESS_TIMEZONE = ZoneInfo("Asia/Shanghai")
def utc_now_naive() -> datetime:
"""Return the database timestamp convention: UTC without tzinfo."""
return datetime.now(UTC).replace(tzinfo=None)
def as_utc_naive(value: datetime) -> datetime:
"""Normalize an API datetime to the database UTC-naive convention."""
if value.tzinfo is None:
return value
return value.astimezone(UTC).replace(tzinfo=None)
def business_day_boundary(value: date, *, end_exclusive: bool = False) -> datetime:
"""Convert a Beijing calendar-day boundary to UTC-naive database time."""
local_value = datetime.combine(value, time.min, tzinfo=BUSINESS_TIMEZONE)
if end_exclusive:
local_value += timedelta(days=1)
return local_value.astimezone(UTC).replace(tzinfo=None)
def business_datetime_to_utc_naive(value: datetime) -> datetime:
"""Interpret a timezone-free operator input as Beijing local time."""
if value.tzinfo is not None:
return as_utc_naive(value)
return value.replace(tzinfo=BUSINESS_TIMEZONE).astimezone(UTC).replace(tzinfo=None)
def to_business_naive(value: datetime | None) -> datetime | None:
"""Convert UTC database time to a timezone-free Beijing value for Excel."""
if value is None:
return None
aware = value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
return aware.astimezone(BUSINESS_TIMEZONE).replace(tzinfo=None)
def api_datetime(value: datetime) -> str:
"""Serialize datetimes with an explicit offset; naive values are database UTC."""
aware = value.replace(tzinfo=UTC) if value.tzinfo is None else value
return aware.astimezone(UTC).isoformat(timespec="milliseconds").replace("+00:00", "Z")
def normalize_api_datetimes(value: Any) -> Any:
"""Recursively enforce the API time contract for ordinary JSON responses."""
if isinstance(value, datetime):
return api_datetime(value)
if isinstance(value, date):
return value.isoformat()
if isinstance(value, dict):
return {key: normalize_api_datetimes(item) for key, item in value.items()}
if isinstance(value, list):
return [normalize_api_datetimes(item) for item in value]
if isinstance(value, tuple):
return [normalize_api_datetimes(item) for item in value]
return value

View File

@@ -248,6 +248,9 @@ class KnowledgeRetrievalCandidate(Base):
class HumanAttentionRecord(Base):
__tablename__ = "sys_human_attention_record"
__table_args__ = (
Index("ix_human_attention_created", "created_at", "id"),
)
id: Mapped[int] = mapped_column(PRIMARY_KEY_TYPE, primary_key=True, autoincrement=True)
session_id: Mapped[int] = mapped_column(BigInteger, index=True, nullable=False)

View File

@@ -17,11 +17,11 @@ PERMISSION_TREE = [
{"code": "content-generation", "name": "内容生成", "children": [{"code": "content-generation.view", "name": "查看配置"}, {"code": "content-generation.edit", "name": "编辑/测试配置"}]},
{"code": "configs", "name": "系统配置", "children": [{"code": "configs.view", "name": "查看配置"}, {"code": "configs.edit", "name": "修改配置"}]},
{"code": "sso", "name": "应用接入", "children": [{"code": "sso.view", "name": "查看应用"}, {"code": "sso.edit", "name": "管理应用"}]},
{"code": "records", "name": "记录审计", "children": [{"code": "records.view", "name": "查看/导出记录"}]},
{"code": "records", "name": "记录审计", "children": [{"code": "records.view", "name": "查看记录"}, {"code": "records.export", "name": "导出记录/洞察"}]},
{"code": "retrievals", "name": "检索日志", "children": [{"code": "retrievals.view", "name": "查看检索日志"}]},
{"code": "attention", "name": "人工关注", "children": [{"code": "attention.view", "name": "查看关注项"}, {"code": "attention.edit", "name": "处理/删除关注项"}, {"code": "attention.config", "name": "查看/修改筛选规则"}, {"code": "attention.preview", "name": "使用历史消息预览筛选效果"}]},
{"code": "feedback", "name": "反馈管理", "children": [{"code": "feedback.view", "name": "查看反馈列表/筛选分页"}, {"code": "feedback.detail", "name": "查看详情/标记已读"}, {"code": "feedback.export", "name": "导出反馈"}, {"code": "feedback.delete", "name": "删除反馈"}]},
{"code": "behavior", "name": "用户行为分析", "children": [{"code": "behavior.view", "name": "查看行为总览和用户轨迹"}]},
{"code": "behavior", "name": "用户行为分析", "children": [{"code": "behavior.view", "name": "查看行为总览和用户轨迹"}, {"code": "behavior.export", "name": "导出用户行为报表"}]},
{"code": "admins", "name": "管理员与权限", "superOnly": True, "children": [{"code": "admins.view", "name": "查看管理员"}, {"code": "admins.edit", "name": "新增/编辑管理员"}, {"code": "admins.delete", "name": "删除管理员"}]},
]

View File

@@ -141,10 +141,10 @@ class AdminDashboardService:
msg_filter.append(ChatMessage.created_at >= start)
ai_filter.append(AiRequestLog.created_at >= start)
if end:
user_filter.append(User.created_at <= end)
session_filter.append(ChatSession.created_at <= end)
msg_filter.append(ChatMessage.created_at <= end)
ai_filter.append(AiRequestLog.created_at <= end)
user_filter.append(User.created_at < end)
session_filter.append(ChatSession.created_at < end)
msg_filter.append(ChatMessage.created_at < end)
ai_filter.append(AiRequestLog.created_at < end)
user_where = and_(true(), *user_filter)
session_where = and_(true(), *session_filter)

View File

@@ -0,0 +1,131 @@
from __future__ import annotations
from collections.abc import Iterable
from io import BytesIO
from openpyxl import Workbook
from openpyxl.styles import Alignment, Font, PatternFill
from openpyxl.worksheet.table import Table, TableStyleInfo
from app.core.ai_content_label import ensure_ai_generated_notice
from app.core.time_utils import to_business_naive
from app.models.chat import ChatMessage, ChatSession
from app.models.user import User
class ChatExportService:
"""Build a complete, pagination-independent workbook for one chat session."""
@staticmethod
def build(
*,
session: ChatSession,
user: User | None,
source_name: str,
messages: Iterable[ChatMessage],
) -> BytesIO:
message_rows = list(messages)
workbook = Workbook()
info_sheet = workbook.active
info_sheet.title = "会话信息"
_write_session_info(
info_sheet,
session=session,
user=user,
source_name=source_name,
message_count=len(message_rows),
)
message_sheet = workbook.create_sheet("完整对话")
_write_messages(message_sheet, message_rows)
stream = BytesIO()
workbook.save(stream)
stream.seek(0)
return stream
def _write_session_info(
sheet,
*,
session: ChatSession,
user: User | None,
source_name: str,
message_count: int,
) -> None:
sheet.sheet_view.showGridLines = False
sheet.append(["项目", "内容"])
rows = [
("会话ID", session.id),
("会话标题", _excel_safe(session.title)),
("用户ID", session.user_id),
("用户姓名", _excel_safe(user.name if user else "")),
("手机号", _excel_safe(user.phone if user else "")),
("会话来源", _excel_safe(source_name)),
("消息总数", message_count),
("最后消息时间", to_business_naive(session.last_message_at)),
("会话更新时间", to_business_naive(session.updated_at)),
]
for label, value in rows:
sheet.append([label, value])
_style_header(sheet[1])
sheet.freeze_panes = "A2"
sheet.column_dimensions["A"].width = 20
sheet.column_dimensions["B"].width = 72
for row in sheet.iter_rows(min_row=2):
row[0].font = Font(bold=True, color="31584F")
row[1].alignment = Alignment(vertical="top", wrap_text=True)
for row_number in (9, 10):
sheet.cell(row_number, 2).number_format = "yyyy-mm-dd hh:mm:ss"
def _write_messages(sheet, messages: Iterable[ChatMessage]) -> None:
headers = ["序号", "消息ID", "发送方", "发送时间", "消息内容", "状态", "输入Token", "输出Token", "耗时(ms)"]
sheet.append(headers)
count = 0
for count, message in enumerate(messages, start=1):
content = ensure_ai_generated_notice(message.content) if message.role == "assistant" else message.content
sheet.append(
[
count,
message.id,
"用户" if message.role == "user" else "大本营答疑助手",
to_business_naive(message.created_at),
_excel_safe(content),
message.message_status,
message.token_input,
message.token_output,
message.response_time_ms,
]
)
_style_header(sheet[1])
sheet.freeze_panes = "A2"
sheet.auto_filter.ref = f"A1:I{max(1, sheet.max_row)}"
widths = (8, 12, 20, 22, 88, 14, 14, 14, 14)
for index, width in enumerate(widths, start=1):
sheet.column_dimensions[chr(64 + index)].width = width
for row in sheet.iter_rows(min_row=2):
row[3].number_format = "yyyy-mm-dd hh:mm:ss"
row[4].alignment = Alignment(vertical="top", wrap_text=True)
if count:
table = Table(displayName="ChatMessages", ref=f"A1:I{sheet.max_row}")
table.tableStyleInfo = TableStyleInfo(
name="TableStyleMedium4",
showFirstColumn=False,
showLastColumn=False,
showRowStripes=True,
)
sheet.add_table(table)
def _style_header(cells) -> None:
fill = PatternFill("solid", fgColor="176B57")
for cell in cells:
cell.fill = fill
cell.font = Font(color="FFFFFF", bold=True)
cell.alignment = Alignment(horizontal="center", vertical="center")
def _excel_safe(value: object | None) -> str:
text = str(value or "")
return f"'{text}" if text.lstrip().startswith(("=", "+", "-", "@")) else text

View File

@@ -11,6 +11,7 @@ from sqlalchemy.orm import Session
from app.models.entitlement import EntitlementPlan, UserEntitlement, UserEntitlementLog
from app.models.user import User
from app.core.time_utils import as_utc_naive
DEFAULT_PLAN_TYPE = "basic"
@@ -260,12 +261,12 @@ class EntitlementService:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="权益版本不存在或已停用")
now = _now()
if effective_at is not None:
effective_at = effective_at.replace(tzinfo=None)
effective_at = as_utc_naive(effective_at)
if expired_at is None and plan.validity_days:
start = effective_at or now
expired_at = start + timedelta(days=plan.validity_days)
elif expired_at is not None:
expired_at = expired_at.replace(tzinfo=None)
expired_at = as_utc_naive(expired_at)
if expired_at is not None and effective_at is not None and expired_at < effective_at:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="权益到期时间不能早于生效时间")

View File

@@ -13,6 +13,7 @@ from app.services.redis_client import get_sync_redis_client
from app.services.entitlement_service import EntitlementService
from app.services.user_behavior_service import UserBehaviorService
from app.core.config import get_settings
from app.core.time_utils import utc_now_naive
class MaintenanceService:
@@ -43,11 +44,12 @@ class MaintenanceService:
policy = db.scalar(select(LogRetentionPolicy).order_by(LogRetentionPolicy.id).limit(1))
if not policy or not policy.enabled or not policy.retention_days:
return
if policy.last_run_at and policy.last_run_at.date() >= datetime.now().date():
now = utc_now_naive()
if policy.last_run_at and policy.last_run_at.date() >= now.date():
return
before = datetime.now() - timedelta(days=policy.retention_days)
before = now - timedelta(days=policy.retention_days)
deleted = cls.delete_retrieval_logs(db, before)
policy.last_run_at = datetime.now()
policy.last_run_at = now
policy.last_result = json.dumps({"status": "success", "deleted": deleted, "before": before.isoformat()}, ensure_ascii=False)
db.add(policy)
db.commit()
@@ -55,7 +57,7 @@ class MaintenanceService:
with SessionLocal() as db:
policy = db.scalar(select(LogRetentionPolicy).order_by(LogRetentionPolicy.id).limit(1))
if policy:
policy.last_run_at = datetime.now()
policy.last_run_at = utc_now_naive()
policy.last_result = json.dumps({"status": "failed", "error": str(exc)[:1000]}, ensure_ascii=False)
db.add(policy)
db.commit()

View File

@@ -629,4 +629,4 @@ def _local_datetime(value: datetime) -> datetime:
except ZoneInfoNotFoundError:
timezone = ZoneInfo("Asia/Shanghai")
aware_utc = value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
return aware_utc.astimezone(timezone).replace(tzinfo=None)
return aware_utc.astimezone(timezone)

View File

@@ -8,6 +8,7 @@ from openpyxl.styles import Alignment, Font, PatternFill
from openpyxl.worksheet.table import Table, TableStyleInfo
from app.core.ai_content_label import AI_GENERATED_NOTICE
from app.core.time_utils import to_business_naive, utc_now_naive
class QuestionInsightExportService:
@@ -116,7 +117,7 @@ def _append_summary(sheet, result: dict) -> None:
("导出问题组", summary.get("visibleClusterCount", 0)),
("清洗规则版本", summary.get("cleanerVersion", "")),
("内容标识", AI_GENERATED_NOTICE),
("导出时间", datetime.now()),
("导出时间", to_business_naive(utc_now_naive())),
("说明", "导出结果按所选日期范围和最低频次生成,包含全部符合条件的问题组,不受页面分页影响。"),
]
for label, value in rows:
@@ -165,6 +166,6 @@ def _excel_safe(value: object | None) -> str:
def _excel_datetime(value: object) -> object:
if isinstance(value, datetime) and value.tzinfo is not None:
return value.replace(tzinfo=None)
if isinstance(value, datetime):
return to_business_naive(value)
return value

View File

@@ -17,6 +17,7 @@ from app.models.chat import ChatMessage, ChatSession
from app.models.insight import QuestionInsightCleanedQuestion
from app.models.logs import AiRequestLog
from app.models.user import User
from app.core.time_utils import as_utc_naive
CLEANER_VERSION = "v1"
@@ -249,9 +250,9 @@ def _load_unprocessed_user_messages(
)
)
if date_from is not None:
query = query.where(ChatMessage.created_at >= date_from.replace(tzinfo=None))
query = query.where(ChatMessage.created_at >= as_utc_naive(date_from))
if date_to is not None:
query = query.where(ChatMessage.created_at <= date_to.replace(tzinfo=None))
query = query.where(ChatMessage.created_at <= as_utc_naive(date_to))
return list(
db.execute(
query.order_by(ChatMessage.created_at.desc(), ChatMessage.id.desc()).limit(limit)
@@ -268,9 +269,9 @@ def _load_persisted_cleaned_questions(
) -> list[QuestionInsightCleanedQuestion]:
range_filters = [QuestionInsightCleanedQuestion.cleaner_version == CLEANER_VERSION]
if date_from is not None:
range_filters.append(QuestionInsightCleanedQuestion.source_created_at >= date_from.replace(tzinfo=None))
range_filters.append(QuestionInsightCleanedQuestion.source_created_at >= as_utc_naive(date_from))
if date_to is not None:
range_filters.append(QuestionInsightCleanedQuestion.source_created_at <= date_to.replace(tzinfo=None))
range_filters.append(QuestionInsightCleanedQuestion.source_created_at <= as_utc_naive(date_to))
latest_messages = (
select(
@@ -312,9 +313,9 @@ def _load_ai_logs(
) -> list[AiRequestLog]:
query = select(AiRequestLog)
if date_from is not None:
query = query.where(AiRequestLog.created_at >= date_from.replace(tzinfo=None))
query = query.where(AiRequestLog.created_at >= as_utc_naive(date_from))
if date_to is not None:
query = query.where(AiRequestLog.created_at <= date_to.replace(tzinfo=None))
query = query.where(AiRequestLog.created_at <= as_utc_naive(date_to))
return list(db.scalars(query.order_by(AiRequestLog.created_at.desc(), AiRequestLog.id.desc()).limit(limit)).all())

View File

@@ -0,0 +1,203 @@
from __future__ import annotations
from datetime import date, datetime
from io import BytesIO
from copy import copy
from openpyxl import Workbook
from openpyxl.styles import Alignment, Font, PatternFill
from openpyxl.utils import get_column_letter
from sqlalchemy.orm import Session
from app.services.user_behavior_service import EVENT_CATALOG, UserBehaviorService
from app.core.time_utils import to_business_naive, utc_now_naive
TITLE_FILL = PatternFill("solid", fgColor="176B55")
HEADER_FILL = PatternFill("solid", fgColor="DDEFE9")
TITLE_FONT = Font(color="FFFFFF", bold=True, size=16)
HEADER_FONT = Font(color="214239", bold=True)
EVENT_TYPE_LABELS = {"page": "页面", "dialog": "弹窗", "button": "按钮"}
SORT_LABELS = {
"userName": "用户名称",
"phone": "手机号",
"eventCount": "交互次数",
"lastEventAt": "最后交互时间",
}
FONT_NAME = "Hiragino Sans GB"
class UserBehaviorExportService:
@staticmethod
def build_workbook(
db: Session,
*,
start: date | None,
end: date | None,
keyword: str,
sort_by: str,
sort_order: str,
) -> BytesIO:
overview = UserBehaviorService.overview(
db,
start=start,
end=end,
ranking_limit=len(EVENT_CATALOG),
)
users = UserBehaviorService.users_for_export(
db,
start=start,
end=end,
keyword=keyword,
sort_by=sort_by,
sort_order=sort_order,
)
workbook = Workbook()
summary = workbook.active
summary.title = "报表说明"
_write_summary(summary, overview, keyword=keyword, sort_by=sort_by, sort_order=sort_order)
_write_daily(workbook.create_sheet("每日趋势"), overview["daily"])
_write_ranking(workbook.create_sheet("功能使用频率"), overview["eventRanking"])
_write_users(workbook.create_sheet("用户汇总"), users)
stream = BytesIO()
workbook.save(stream)
stream.seek(0)
return stream
def _write_summary(sheet, overview: dict, *, keyword: str, sort_by: str, sort_order: str) -> None:
sheet.sheet_view.showGridLines = False
sheet.merge_cells("A1:D1")
title = sheet["A1"]
title.value = "用户行为分析报表"
title.fill = TITLE_FILL
title.font = TITLE_FONT
title.alignment = Alignment(horizontal="center", vertical="center")
sheet.row_dimensions[1].height = 30
metadata = [
("统计周期", f'{overview["startDate"]}{overview["endDate"]}'),
("用户筛选", _excel_safe(keyword.strip()) if keyword.strip() else "全部用户"),
("用户汇总排序", f'{SORT_LABELS[sort_by]}·{"升序" if sort_order == "asc" else "降序"}'),
("数据保留", f'{overview["retentionDays"]}'),
("导出时间", to_business_naive(utc_now_naive())),
]
for row_index, (label, value) in enumerate(metadata, start=3):
sheet.cell(row_index, 1, label).font = HEADER_FONT
sheet.cell(row_index, 2, value)
metrics = [
("交互事件", overview["totalEvents"]),
("活跃用户", overview["activeUsers"]),
("页面/弹窗打开", overview["pageDialogOpens"]),
("关键按钮点击", overview["buttonClicks"]),
]
sheet.cell(10, 1, "核心指标").fill = HEADER_FILL
sheet.cell(10, 1).font = HEADER_FONT
sheet.cell(10, 2, "数值").fill = HEADER_FILL
sheet.cell(10, 2).font = HEADER_FONT
for row_index, (label, value) in enumerate(metrics, start=11):
sheet.cell(row_index, 1, label)
sheet.cell(row_index, 2, value).number_format = "#,##0"
sheet.column_dimensions["A"].width = 22
sheet.column_dimensions["B"].width = 34
sheet.column_dimensions["C"].width = 16
sheet.column_dimensions["D"].width = 16
sheet["B7"].number_format = "yyyy-mm-dd hh:mm:ss"
_finish_sheet(sheet, landscape=True)
def _write_daily(sheet, rows: list[dict]) -> None:
_write_tabular_sheet(sheet, ["日期", "交互事件", "活跃用户"])
for item in rows:
sheet.append([date.fromisoformat(item["date"]), item["eventCount"], item["activeUsers"]])
for cell in sheet["A"][1:]:
cell.number_format = "yyyy-mm-dd"
for column in ("B", "C"):
for cell in sheet[column][1:]:
cell.number_format = "#,##0"
_finish_tabular_sheet(sheet, widths=(15, 15, 15))
def _write_ranking(sheet, rows: list[dict]) -> None:
_write_tabular_sheet(sheet, ["排名", "功能名称", "事件类型", "使用次数", "用户数", "事件代码"])
for index, item in enumerate(rows, start=1):
sheet.append(
[
index,
_excel_safe(item["eventName"]),
EVENT_TYPE_LABELS.get(item["eventType"], item["eventType"]),
item["count"],
item["userCount"],
item["eventCode"],
]
)
for column in ("A", "D", "E"):
for cell in sheet[column][1:]:
cell.number_format = "#,##0"
_finish_tabular_sheet(sheet, widths=(10, 24, 14, 14, 14, 30))
def _write_users(sheet, rows: list[dict]) -> None:
_write_tabular_sheet(sheet, ["用户ID", "用户", "手机号", "交互次数", "最后交互时间"])
for item in rows:
sheet.append(
[
item["userId"],
_excel_safe(item["userName"]),
_excel_safe(item["phone"]),
item["eventCount"],
to_business_naive(item["lastEventAt"]),
]
)
for cell in sheet["A"][1:]:
cell.number_format = "0"
for cell in sheet["C"][1:]:
cell.number_format = "@"
for cell in sheet["D"][1:]:
cell.number_format = "#,##0"
for cell in sheet["E"][1:]:
cell.number_format = "yyyy-mm-dd hh:mm:ss"
_finish_tabular_sheet(sheet, widths=(12, 22, 18, 14, 22))
def _write_tabular_sheet(sheet, headers: list[str]) -> None:
sheet.sheet_view.showGridLines = False
sheet.append(headers)
for cell in sheet[1]:
cell.fill = HEADER_FILL
cell.font = HEADER_FONT
cell.alignment = Alignment(horizontal="center", vertical="center")
sheet.row_dimensions[1].height = 24
sheet.freeze_panes = "A2"
def _finish_tabular_sheet(sheet, *, widths: tuple[int, ...]) -> None:
sheet.auto_filter.ref = sheet.dimensions
for index, width in enumerate(widths, start=1):
sheet.column_dimensions[get_column_letter(index)].width = width
for row in sheet.iter_rows(min_row=2):
for cell in row:
cell.alignment = Alignment(vertical="top", wrap_text=True)
_finish_sheet(sheet, landscape=True)
def _finish_sheet(sheet, *, landscape: bool) -> None:
for row in sheet.iter_rows():
for cell in row:
font = copy(cell.font)
font.name = FONT_NAME
cell.font = font
sheet.sheet_properties.pageSetUpPr.fitToPage = True
sheet.page_setup.fitToWidth = 1
sheet.page_setup.fitToHeight = 0
sheet.page_setup.orientation = "landscape" if landscape else "portrait"
sheet.page_setup.paperSize = sheet.PAPERSIZE_A4
sheet.page_margins.left = 0.3
sheet.page_margins.right = 0.3
sheet.page_margins.top = 0.5
sheet.page_margins.bottom = 0.5
sheet.print_area = sheet.dimensions
def _excel_safe(value: object) -> str:
text = "" if value is None else str(value)
return f"'{text}" if text.startswith(("=", "+", "-", "@")) else text

View File

@@ -1,7 +1,6 @@
from __future__ import annotations
from datetime import UTC, date, datetime, time, timedelta
from zoneinfo import ZoneInfo
from datetime import UTC, date, datetime, timedelta
from sqlalchemy import case, delete, func, or_, select
from sqlalchemy.exc import IntegrityError
@@ -11,6 +10,7 @@ from app.models.behavior import UserBehaviorEvent
from app.models.user import User
from app.schemas.behavior import UserBehaviorEventCreate
from app.core.config import get_settings
from app.core.time_utils import BUSINESS_TIMEZONE, business_day_boundary, to_business_naive
EVENT_CATALOG: dict[str, tuple[str, str]] = {
@@ -48,6 +48,8 @@ EVENT_CATALOG: dict[str, tuple[str, str]] = {
TARGET_TYPES = {"session", "message", "help_card", "share_draft", "report"}
DEFAULT_RANGE_DAYS = 7
MAX_RANGE_DAYS = 31
USER_SORT_FIELDS = {"userName", "phone", "eventCount", "lastEventAt"}
USER_SORT_ORDERS = {"asc", "desc"}
class UserBehaviorService:
@@ -97,7 +99,7 @@ class UserBehaviorService:
return accepted
@staticmethod
def overview(db: Session, *, start: date | None, end: date | None) -> dict:
def overview(db: Session, *, start: date | None, end: date | None, ranking_limit: int = 20) -> dict:
start_dt, end_dt = _date_range(start, end)
base = (UserBehaviorEvent.occurred_at >= start_dt, UserBehaviorEvent.occurred_at < end_dt)
total, active_users, page_opens, button_clicks = db.execute(
@@ -119,7 +121,7 @@ class UserBehaviorService:
.where(*base)
.group_by(UserBehaviorEvent.event_code, UserBehaviorEvent.event_name, UserBehaviorEvent.event_type)
.order_by(func.count(UserBehaviorEvent.id).desc(), UserBehaviorEvent.event_code)
.limit(20)
.limit(max(1, min(ranking_limit, len(EVENT_CATALOG))))
).all()
if db.bind and db.bind.dialect.name == "mysql":
day_expression = func.date(func.convert_tz(UserBehaviorEvent.occurred_at, "+00:00", "+08:00"))
@@ -135,8 +137,8 @@ class UserBehaviorService:
.group_by(day_expression)
.order_by(day_expression)
).all()
local_start = (start_dt + timedelta(hours=8)).date()
local_end = (end_dt + timedelta(hours=8) - timedelta(days=1)).date()
local_start = to_business_naive(start_dt).date()
local_end = (to_business_naive(end_dt) - timedelta(days=1)).date()
daily_map = {str(day): (event_count, users) for day, event_count, users in daily_rows}
daily = []
cursor = local_start
@@ -160,47 +162,57 @@ class UserBehaviorService:
}
@staticmethod
def users(db: Session, *, start: date | None, end: date | None, keyword: str, page: int, page_size: int) -> dict:
start_dt, end_dt = _date_range(start, end)
filters = [UserBehaviorEvent.occurred_at >= start_dt, UserBehaviorEvent.occurred_at < end_dt]
if keyword.strip():
pattern = f"%{keyword.strip()}%"
filters.append(or_(User.name.like(pattern), User.nickname.like(pattern), User.phone.like(pattern)))
grouped = (
select(
User.id.label("user_id"),
User.name,
User.nickname,
User.phone,
func.count(UserBehaviorEvent.id).label("event_count"),
func.max(UserBehaviorEvent.occurred_at).label("last_event_at"),
)
.join(UserBehaviorEvent, UserBehaviorEvent.user_id == User.id)
.where(*filters)
.group_by(User.id, User.name, User.nickname, User.phone)
def users(
db: Session,
*,
start: date | None,
end: date | None,
keyword: str,
page: int,
page_size: int,
sort_by: str = "lastEventAt",
sort_order: str = "desc",
) -> dict:
grouped, order_expression = _user_summary_query(
start=start,
end=end,
keyword=keyword,
sort_by=sort_by,
sort_order=sort_order,
)
total = db.scalar(select(func.count()).select_from(grouped.subquery())) or 0
rows = db.execute(
grouped.order_by(func.max(UserBehaviorEvent.occurred_at).desc(), User.id.desc())
grouped.order_by(order_expression, User.id.asc())
.offset((page - 1) * page_size)
.limit(page_size)
).all()
return {
"items": [
{
"userId": row.user_id,
"userName": row.nickname or row.name,
"phone": row.phone,
"eventCount": row.event_count,
"lastEventAt": row.last_event_at,
}
for row in rows
],
"items": [_user_summary_dict(row) for row in rows],
"total": total,
"page": page,
"pageSize": page_size,
}
@staticmethod
def users_for_export(
db: Session,
*,
start: date | None,
end: date | None,
keyword: str,
sort_by: str,
sort_order: str,
) -> list[dict]:
grouped, order_expression = _user_summary_query(
start=start,
end=end,
keyword=keyword,
sort_by=sort_by,
sort_order=sort_order,
)
rows = db.execute(grouped.order_by(order_expression, User.id.asc())).all()
return [_user_summary_dict(row) for row in rows]
@staticmethod
def timeline(db: Session, *, user_id: int, start: date | None, end: date | None, page: int, page_size: int) -> dict:
user = db.get(User, user_id)
@@ -249,18 +261,67 @@ class UserBehaviorService:
def _date_range(start: date | None, end: date | None) -> tuple[datetime, datetime]:
today = datetime.now(ZoneInfo("Asia/Shanghai")).date()
today = datetime.now(BUSINESS_TIMEZONE).date()
end_date = end or today
start_date = start or (end_date - timedelta(days=DEFAULT_RANGE_DAYS - 1))
if end_date < start_date:
start_date, end_date = end_date, start_date
if (end_date - start_date).days >= MAX_RANGE_DAYS:
start_date = end_date - timedelta(days=MAX_RANGE_DAYS - 1)
# Admin date filters are Beijing calendar days; persisted timestamps are UTC-naive.
return (
datetime.combine(start_date, time.min) - timedelta(hours=8),
datetime.combine(end_date + timedelta(days=1), time.min) - timedelta(hours=8),
return business_day_boundary(start_date), business_day_boundary(end_date, end_exclusive=True)
def _user_summary_query(
*,
start: date | None,
end: date | None,
keyword: str,
sort_by: str,
sort_order: str,
):
if sort_by not in USER_SORT_FIELDS:
raise ValueError(f"unsupported user behavior sort field: {sort_by}")
if sort_order not in USER_SORT_ORDERS:
raise ValueError(f"unsupported user behavior sort order: {sort_order}")
start_dt, end_dt = _date_range(start, end)
filters = [UserBehaviorEvent.occurred_at >= start_dt, UserBehaviorEvent.occurred_at < end_dt]
if keyword.strip():
pattern = f"%{keyword.strip()}%"
filters.append(or_(User.name.like(pattern), User.nickname.like(pattern), User.phone.like(pattern)))
user_name = func.coalesce(func.nullif(User.nickname, ""), User.name)
event_count = func.count(UserBehaviorEvent.id)
last_event_at = func.max(UserBehaviorEvent.occurred_at)
grouped = (
select(
User.id.label("user_id"),
User.name,
User.nickname,
User.phone,
event_count.label("event_count"),
last_event_at.label("last_event_at"),
)
.join(UserBehaviorEvent, UserBehaviorEvent.user_id == User.id)
.where(*filters)
.group_by(User.id, User.name, User.nickname, User.phone)
)
sort_expressions = {
"userName": user_name,
"phone": User.phone,
"eventCount": event_count,
"lastEventAt": last_event_at,
}
expression = sort_expressions[sort_by]
return grouped, expression.asc() if sort_order == "asc" else expression.desc()
def _user_summary_dict(row) -> dict:
return {
"userId": row.user_id,
"userName": row.nickname or row.name,
"phone": row.phone,
"eventCount": row.event_count,
"lastEventAt": row.last_event_at,
}
def _event_dict(item: UserBehaviorEvent) -> dict: