diff --git a/backend/alembic/versions/k1l2m3n4o5p6_add_user_daily_seen.py b/backend/alembic/versions/k1l2m3n4o5p6_add_user_daily_seen.py new file mode 100644 index 0000000..8f4498b --- /dev/null +++ b/backend/alembic/versions/k1l2m3n4o5p6_add_user_daily_seen.py @@ -0,0 +1,53 @@ +"""add_user_daily_seen + +Revision ID: k1l2m3n4o5p6 +Revises: j1k2l3m4n5o6 +Create Date: 2026-09-21 + +每日用户活动表(user_daily_seen) +-------------------------------- +一天一人一行,记录当天首次 / 末次活动时刻,供日活报表计算 +「上线时间 / 下线时间」。 + +为什么不复用 audit_logs: + · 上线/下线时间不能取登录时间 —— Refresh Token 有效期 7 天,用户不必每天 + 重新登录,「登录次数 0 却操作 35 次」的报表没有意义。 + · 也不能只取写操作时间 —— 审计中间件只记写操作,普通 GET 不入账, + 当天只翻看的人会被漏掉。 + · 更不能把活动写进审计表 —— 「末次活动」是需要不断 UPDATE 的状态, + 而审计流水必须只增不改;能改的审计记录等于没有审计价值。 + +存量数据无需回填:本表从上线时刻开始记录;日活接口对更早的日期会自动 +回退到审计表的写操作时间去推算。 +""" +from typing import Sequence, Union + +import sqlalchemy as sa +from alembic import op + +revision: str = "k1l2m3n4o5p6" +down_revision: Union[str, None] = "j1k2l3m4n5o6" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "user_daily_seen", + sa.Column("user_id", sa.String(64), primary_key=True, + comment="操作人账号(逻辑外键→MOM)"), + sa.Column("day", sa.Date(), primary_key=True, + comment="北京时间自然日"), + sa.Column("first_seen_at", sa.DateTime(timezone=True), nullable=False, + comment="当天首次活动时刻"), + sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, + comment="当天末次活动时刻"), + ) + # 日活查询按日期区间扫,给 day 单独建索引。 + # (主键是 (user_id, day),前缀是 user_id,按 day 过滤用不上,故需补一条) + op.create_index("ix_user_daily_seen_day", "user_daily_seen", ["day"]) + + +def downgrade() -> None: + op.drop_index("ix_user_daily_seen_day", table_name="user_daily_seen") + op.drop_table("user_daily_seen") diff --git a/backend/app/core/middleware.py b/backend/app/core/middleware.py index d8c8996..36b74b0 100644 --- a/backend/app/core/middleware.py +++ b/backend/app/core/middleware.py @@ -9,6 +9,7 @@ from starlette.middleware.base import BaseHTTPMiddleware from starlette.requests import Request from app.core.logging import request_id_var, user_var +from app.services.audit_service import touch_daily_seen access_log = logging.getLogger("track.access") @@ -42,14 +43,29 @@ class RequestContextMiddleware(BaseHTTPMiddleware): response.headers["X-Request-ID"] = request_id self._log_access(request, status_code, started) logged = True + await self._touch_activity(request) return response finally: # 异常路径也要留下访问记录,否则接口 500 时日志里反而没有痕迹 if not logged: self._log_access(request, status_code, started) + await self._touch_activity(request) request_id_var.reset(rid_token) user_var.reset(user_token) + async def _touch_activity(self, request: Request) -> None: + """记录「该用户今天活动过」,供日活报表算上线/下线时间。 + + 为什么挂在这一层:本中间件是最外层,能覆盖**所有**请求 —— + 包括不被审计的普通 GET。而审计中间件只记写操作,当天只翻看、 + 没做写操作的人会被日活完全漏掉。 + + user 同样只能从 request.state 取:本中间件在独立 task 中执行, + 路由内写的 contextvar 不会回流(详见 _log_access 的说明)。 + 未认证请求取不到 user,自然跳过。 + """ + await touch_daily_seen(getattr(request.state, "audit_user", None)) + def _log_access(self, request: Request, status_code: int, started: float) -> None: path = request.url.path duration_ms = round((time.perf_counter() - started) * 1000, 1) diff --git a/backend/app/models/__init__.py b/backend/app/models/__init__.py index ac5bba5..951d631 100644 --- a/backend/app/models/__init__.py +++ b/backend/app/models/__init__.py @@ -9,6 +9,7 @@ from app.models.app_version import AppVersion from app.models.message import ProductMessage from app.models.holiday import Holiday from app.models.audit_log import AuditLog +from app.models.user_daily_seen import UserDailySeen __all__ = [ "Base", "ProductionOrder", @@ -21,4 +22,5 @@ __all__ = [ "ProductMessage", "Holiday", "AuditLog", + "UserDailySeen", ] diff --git a/backend/app/models/user_daily_seen.py b/backend/app/models/user_daily_seen.py new file mode 100644 index 0000000..38cf73d --- /dev/null +++ b/backend/app/models/user_daily_seen.py @@ -0,0 +1,47 @@ +"""每日用户活动表 —— 一天一人一行,只记录"今天来过"这件事。 + +为什么需要它(而不是复用 audit_logs) +-------------------------------------- +日活报表要的「上线时间 / 下线时间」,两个都不能从审计表直接得出: + +1. **上线/下线时间不能取登录时间**:Refresh Token 有效期 7 天,用户不必每天 + 重新登录。按登录算会出现「登录次数 0、上线时间空,但操作次数 35」的 + 自相矛盾报表。 + +2. **也不能只取写操作时间**:审计中间件只记录写操作(及导出/打印这类敏感读), + 普通 GET 不入账。当天只翻看、没做写操作的人会被整条漏掉。 + +3. **更不能把活动记录写进 audit_logs**: + · 「上线时间」是**事件**(INSERT 一次即可),但「下线时间」是**状态** + (每次活动都要刷新同一个值)。往审计流水里做 UPDATE,等于承认审计记录 + 可以被改写 —— 那审计本身就失去可信度了。 + · 若改为每个请求 INSERT 一条,表会随访问量线性膨胀。 + +于是单开一张"可变的小状态表":一人一天一行,首见 INSERT、其后只 +UPDATE last_seen_at。50 人 × 365 天 ≈ 1.8 万行/年,可忽略。 +""" +from __future__ import annotations + +from datetime import date, datetime + +from sqlalchemy import Date, DateTime, String +from sqlalchemy.orm import Mapped, mapped_column + +from app.models.base import Base + + +class UserDailySeen(Base): + """用户在某个北京时间自然日的首末活动时刻""" + + __tablename__ = "user_daily_seen" + + # 联合主键即 UPSERT 的冲突目标,也是"一天一人一行"的保证 + user_id: Mapped[str] = mapped_column(String(64), primary_key=True) + day: Mapped[date] = mapped_column(Date, primary_key=True, comment="北京时间自然日") + + first_seen_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, comment="当天首次活动时刻", + ) + last_seen_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, comment="当天末次活动时刻", + ) diff --git a/backend/app/services/audit_service.py b/backend/app/services/audit_service.py index e032902..7021c94 100644 --- a/backend/app/services/audit_service.py +++ b/backend/app/services/audit_service.py @@ -15,6 +15,7 @@ Track 改为:响应生成后,用**独立 session** 写入审计。 from __future__ import annotations import logging +import time import uuid from datetime import datetime @@ -22,7 +23,9 @@ from sqlalchemy import and_, func, select from sqlalchemy.ext.asyncio import AsyncSession from app.core.database import AsyncSessionLocal +from app.core.time_utils import get_beijing_time from app.models.audit_log import AuditLog +from app.models.user_daily_seen import UserDailySeen logger = logging.getLogger("track.audit") @@ -252,6 +255,68 @@ async def export_audit_logs( return list(rows[:limit]), truncated +# ============================================================ +# 每日活动打点(日活报表的「上线时间 / 下线时间」来源) +# ============================================================ + +# 同一用户两次落盘之间的最小间隔(秒)。 +# +# 打点挂在「每个请求」上,但不希望每个请求都写一次数据库 —— 那会把 +# user_daily_seen 变成热点。这里用进程内缓存做节流:同一用户 2 分钟内 +# 只落盘一次。代价是「末次活动时间」最多落后真实值 2 分钟, +# 对"日活统计"这个精度要求完全够用。 +# +# 多 worker 部署时每个进程各持一份缓存,实际写库频率最多放大到 worker 数倍 +# (4 worker × 每人每 2 分钟 1 次),依然可忽略。 +_TOUCH_INTERVAL_S = 120.0 +_touch_cache: dict[str, float] = {} + +# 缓存只增不减会缓慢泄漏(键是 user_id,量级 = 用户数,实际很小)。 +# 超过阈值就整体清空 —— 代价只是多写几次库,换来内存有界。 +_TOUCH_CACHE_MAX = 5000 + + +async def touch_daily_seen(user_id: str | None) -> None: + """记录「该用户此刻活动过」。首次 INSERT、其后只刷新 last_seen_at。 + + 唯一的消费方是日活报表的上线/下线时间(见 get_daily_usage)。 + 刻意不写进 audit_logs:那是只增不改的审计流水,而本表是需要不断 + UPDATE 的状态(详见 UserDailySeen 模型注释)。 + + 任何异常都吞掉 —— 活动打点失败绝不能影响业务请求本身。 + """ + if not user_id: + return + + now_mono = time.monotonic() + last = _touch_cache.get(user_id) + if last is not None and now_mono - last < _TOUCH_INTERVAL_S: + return # 节流窗口内,跳过 + if len(_touch_cache) > _TOUCH_CACHE_MAX: + _touch_cache.clear() + # 先占位再写库:同一用户的并发请求不会同时打进来 + _touch_cache[user_id] = now_mono + + try: + from sqlalchemy.dialects.postgresql import insert as pg_insert + + now = get_beijing_time() + day = now.date() # 北京时间自然日(与报表分日口径一致) + async with AsyncSessionLocal() as db: + await db.execute( + pg_insert(UserDailySeen) + .values(user_id=user_id, day=day, first_seen_at=now, last_seen_at=now) + # 冲突时只刷新 last_seen_at,first_seen_at 保持当天首次值不变 + .on_conflict_do_update( + index_elements=["user_id", "day"], + set_={"last_seen_at": now}, + ) + ) + await db.commit() + except Exception: # noqa: BLE001 —— 打点失败不影响业务 + logger.exception("记录每日活动失败(已忽略)") + + # ============================================================ # 日活 / 使用统计 # ============================================================ @@ -270,15 +335,18 @@ async def get_daily_usage( 全部指标由**一个 GROUP BY 查询**算出,不用窗口函数: · 登录/登出次数 = 成功登录 / 成功登出数(最终凭证是 login_count,不是"上线次数") · 操作频次 = 当天该用户的全部审计记录数(代表系统使用深度) - · 上线/下线时间 = 当天**首次 / 末次活动**时间(任意审计记录) + · 登录/登出次数 = 成功登录 / 成功登出数 + · 上线/下线时间 = 当天**首次 / 末次活动**(优先取 user_daily_seen) ⚠️ 上线/下线时间【不能】取登录/登出时间。 Access/Refresh Token 有效期内(refresh 7 天)用户无需重新登录, 于是"周一登录、周二到周日继续用"会导致周二~周日: 登录次数=0、登录时间=空,但操作次数却是几十 —— 报表自相矛盾。 - 改用活动口径后,工人当天的第一次操作就是真实上线时间。 - (审计中间件是全站的,移动端的接收/转交/完工同样入账,故自动覆盖移动端, - 无需前端上报心跳。) + + ⚠️ 也不能只取审计表的写操作时间:审计中间件只记写操作,普通 GET 不入账, + 当天只翻看、没做写操作的人会被整条漏掉。 + 故上线/下线时间优先取 user_daily_seen(挂在每个请求上打点), + 仅对本表上线前的历史数据回退到审计表的写操作时间。 为什么用 `count(*) FILTER (WHERE ...)`:分组内一次扫描同时算出多个条件计数, 比多次子查询或 UNION 简单得多,且语义一眼可读。Postgres 原生支持。 @@ -323,17 +391,62 @@ async def get_daily_usage( ) rows = (await db.execute(stmt)).all() - return [ - { - "day": r.day.strftime("%Y-%m-%d") if hasattr(r.day, "strftime") else str(r.day), - "user_id": r.user_id, - "display_name": r.display_name, - "role": r.role, - "login_count": r.login_count or 0, - "logout_count": r.logout_count or 0, - "op_count": r.op_count or 0, - "first_active_at": r.first_active_at, - "last_active_at": r.last_active_at, - } + + # ── 活动表:当天首次/末次活动(覆盖"只翻看不操作"的人)── + # day 列是北京时间 DATE,与上面的 day_col 口径一致,可直接按 (user_id, day) 对齐。 + # 取 start.date() ~ end.date()(end 是次日 00:00 的半开上界,故用 <)。 + seen_rows = ( + await db.execute( + select( + UserDailySeen.user_id, UserDailySeen.day, + UserDailySeen.first_seen_at, UserDailySeen.last_seen_at, + ).where( + UserDailySeen.day >= start.date(), + UserDailySeen.day < end.date(), + ) + ) + ).all() + seen = { + (s.user_id, s.day.strftime("%Y-%m-%d")): (s.first_seen_at, s.last_seen_at) + for s in seen_rows + } + + audit = { + (r.user_id, r.day.strftime("%Y-%m-%d") if hasattr(r.day, "strftime") else str(r.day)): r for r in rows - ] + } + + # ── 合并 ── + # 并集:只有审计记录的人(本表上线前的历史数据)和只有活动记录的人 + # (当天只翻看、没做写操作)都要出现,各自缺的部分留空/计 0。 + items: list[dict] = [] + for key in set(audit) | set(seen): + user_id, day = key + a = audit.get(key) + first_seen, last_seen = seen.get(key, (None, None)) + + # 取「两者的最早/最晚」,而不是简单地"活动表优先": + # 活动表靠请求触发且有 2 分钟节流,极端情况(跨零点被节流、 + # 打点写库失败被吞掉)可能晚于当天第一次写操作。 + # 取 min/max 后,结果永远不会比任一来源更差,也不需要为兜底写分支逻辑。 + audit_first = a.first_active_at if a else None + audit_last = a.last_active_at if a else None + first_candidates = [t for t in (first_seen, audit_first) if t is not None] + last_candidates = [t for t in (last_seen, audit_last) if t is not None] + + items.append({ + "day": day, + "user_id": user_id, + # 姓名字段只有审计记录里有(活动表为了轻量刻意不冗余存) + "display_name": a.display_name if a else None, + "role": a.role if a else None, + "login_count": (a.login_count or 0) if a else 0, + "logout_count": (a.logout_count or 0) if a else 0, + "op_count": (a.op_count or 0) if a else 0, + "first_active_at": min(first_candidates) if first_candidates else None, + "last_active_at": max(last_candidates) if last_candidates else None, + }) + + # 与 SQL 里的排序保持一致:日期倒序 → 操作次数倒序 + items.sort(key=lambda x: (x["day"], x["op_count"]), reverse=True) + return items