Files
ctms/backend/app/services/permission_metric_aggregator.py
T
Cheng Zhou 747dd55225 feat(perm): 项目 PM 共管接口权限与监控并细化 PM 写权限边界
- deps 新增 list_active_pm_study_ids、is_active_project_pm 与
  require_admin_or_any_project_pm 依赖,用于把 PM 项目范围带进鉴权与
  监控。get_cra_site_scope 内部延后导入 site CRUD,避免循环依赖。
- system_permissions API 改用 PM/ADMIN 双角色入口,permissions/monitoring
  系统级权限新增 PM 配额,并细化访问日志、告警、监控指标的可见范围。
- members API 调整:项目 PM 仅可管理低于 PM 的项目角色,禁止互相
  改写或授予 PM。
- api_permissions API 增加 GET /api-permissions/me,返回当前用户在该
  项目的有效权限矩阵;保存权限矩阵时校验 PM 行为不被篡改。
- core/api_permissions:新增立项配置接口键、PM 默认拥有的监控/权限
  系统级条目,并在权限元信息中标注 PM 共享角色。
- core/project_permissions:role_has_api_permission 命中默认角色矩阵;
  replace_api_endpoint_permissions 改为部分更新且永不持久化 ADMIN/PM。
- studies setup-config 各端点改用接口级权限装饰器,与新的 setup_config
  权限键对齐。permission_monitor 新增 get_metrics 摘要供 PM 视图调用。
- 测试:新增 test_admin_pm_permissions 覆盖 PM 系统级权限、监控范围和
  成员管理边界;conftest 兼容 SA_UUID 列;权限相关用例同步移除已失效
  的 module_permission 链路。

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-25 12:38:07 +08:00

83 lines
3.0 KiB
Python

"""权限指标小时聚合任务
每小时从 permission_access_logs 聚合数据写入 permission_metric_snapshots。
"""
from __future__ import annotations
import asyncio
import logging
import uuid
from datetime import datetime, timedelta, timezone
from sqlalchemy import func, select
from app.db.session import SessionLocal
from app.models.permission_access_log import PermissionAccessLog
from app.models.permission_metric_snapshot import PermissionMetricSnapshot
logger = logging.getLogger("ctms.permission_metric_aggregator")
async def _aggregate_hour(bucket_start: datetime, bucket_end: datetime) -> None:
async with SessionLocal() as session:
result = await session.execute(
select(
func.count().label("total"),
func.count().filter(PermissionAccessLog.allowed.is_(True)).label("allowed"),
func.count().filter(PermissionAccessLog.allowed.is_(False)).label("denied"),
func.coalesce(func.avg(PermissionAccessLog.elapsed_ms), 0).label("avg_ms"),
func.coalesce(func.max(PermissionAccessLog.elapsed_ms), 0).label("max_ms"),
).where(
PermissionAccessLog.created_at >= bucket_start,
PermissionAccessLog.created_at < bucket_end,
)
)
row = result.one()
if row.total == 0:
return
from app.core.permission_monitor import get_permission_monitor
monitor = get_permission_monitor()
cache_metrics = monitor.metrics.cache_metrics
snapshot = PermissionMetricSnapshot(
id=uuid.uuid4(),
bucket_time=bucket_start,
total_checks=row.total,
allowed_checks=row.allowed,
denied_checks=row.denied,
avg_elapsed_ms=float(row.avg_ms),
max_elapsed_ms=float(row.max_ms),
cache_hits=cache_metrics.cache_hits,
cache_misses=cache_metrics.cache_misses,
error_count=0,
)
session.add(snapshot)
await session.commit()
logger.info("Aggregated permission metrics for bucket %s: %d checks", bucket_start, row.total)
async def run_hourly_metric_aggregation(stop_event: asyncio.Event) -> None:
logger.info("Permission metric aggregator started")
while not stop_event.is_set():
now = datetime.now(timezone.utc)
next_hour = now.replace(minute=0, second=0, microsecond=0) + timedelta(hours=1)
wait_seconds = (next_hour - now).total_seconds()
try:
await asyncio.wait_for(stop_event.wait(), timeout=wait_seconds)
break
except asyncio.TimeoutError:
pass
bucket_end = next_hour
bucket_start = bucket_end - timedelta(hours=1)
try:
await _aggregate_hour(bucket_start, bucket_end)
except Exception:
logger.exception("Failed to aggregate permission metrics for %s", bucket_start)
logger.info("Permission metric aggregator stopped")