747dd55225
- 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>
83 lines
3.0 KiB
Python
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")
|