feat(backend): RBAC admin roles + runtime logs service

This commit is contained in:
Simon
2026-07-09 13:46:15 +08:00
parent 7ffc6c8e23
commit d480aa4c1d
12 changed files with 726 additions and 25 deletions
+29 -3
View File
@@ -15,6 +15,7 @@
# ============================================================================= # =============================================================================
import logging import logging
from datetime import datetime
from typing import Optional from typing import Optional
from uuid import UUID from uuid import UUID
@@ -895,12 +896,37 @@ async def get_agent_performance(
async def get_system_logs( async def get_system_logs(
page: int = Query(1, ge=1, description="页码"), page: int = Query(1, ge=1, description="页码"),
page_size: int = Query(50, ge=1, le=200, description="每页条数"), page_size: int = Query(50, ge=1, le=200, description="每页条数"),
config_key: Optional[str] = Query(None, description="按配置键精确筛选"),
changed_by: Optional[str] = Query(None, description="按操作人 agent_id 精确筛选"),
from_time: Optional[datetime] = Query(None, alias="from", description="变更时间起始(ISO8601, 闭区间)"),
to_time: Optional[datetime] = Query(None, alias="to", description="变更时间截止(ISO8601, 闭区间)"),
admin: Agent = Depends(require_admin), admin: Agent = Depends(require_admin),
db: AsyncSession = Depends(get_db), db: AsyncSession = Depends(get_db),
): ):
"""获取系统日志(配置变更日志)""" """获取系统日志(配置变更日志),支持配置键/操作人/时间范围筛选。
result = await admin_service.get_system_logs(db, page=page, page_size=page_size)
return success_response(data=result)# ---------- GET /api/admin/integrations/lianruan/terminals/{devname}/detail ---------- Args:
config_key: 按配置键精确匹配(可选)
changed_by: 按操作人 agent_id 精确匹配(可选)
from_time: 变更时间起始(别名 from,可选)
to_time: 变更时间截止(别名 to,可选)
Returns:
Dict: 统一响应格式,包含筛选后的配置变更日志
"""
result = await admin_service.get_system_logs(
db,
page=page,
page_size=page_size,
config_key=config_key,
changed_by=changed_by,
from_time=from_time,
to_time=to_time,
)
return success_response(data=result)
# ---------- GET /api/admin/integrations/lianruan/terminals/{devname}/detail ----------
@router.get("/integrations/lianruan/terminals/{devname}/detail") @router.get("/integrations/lianruan/terminals/{devname}/detail")
async def get_lianruan_terminal_detail( async def get_lianruan_terminal_detail(
devname: str, devname: str,
+79 -8
View File
@@ -13,12 +13,14 @@ import logging
from datetime import datetime from datetime import datetime
from typing import List, Optional from typing import List, Optional
import redis.asyncio as aioredis
from fastapi import APIRouter, Depends, Query from fastapi import APIRouter, Depends, Query
from sqlalchemy import select, func from sqlalchemy import select, func
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from app.dependencies import get_current_user, UserInfo, require_role from app.dependencies import get_current_user, UserInfo, require_role, dep_redis
from app.database import get_db from app.database import get_db
from app.models.employee import Employee
from app.models.role import Role from app.models.role import Role
from app.models.role_mapping_rule import RoleMappingRule from app.models.role_mapping_rule import RoleMappingRule
from app.models.user_role import UserRole from app.models.user_role import UserRole
@@ -30,6 +32,7 @@ from app.schemas.role import (
RoleResponse, RoleResponse,
UserRoleResponse, UserRoleResponse,
) )
from app.services.employee_directory import get_org_directory, resolve_target
from app.utils.response import AppException, success_response from app.utils.response import AppException, success_response
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -146,22 +149,27 @@ async def list_user_role_assignments(
""" """
# 使用 LEFT OUTER JOIN:即使 user_roles.role_id 在 roles 表中已不存在 # 使用 LEFT OUTER JOIN:即使 user_roles.role_id 在 roles 表中已不存在
# (例如角色被重建导致 UUID 变化),也保留该条分配记录,避免已分配用户被静默隐藏。 # (例如角色被重建导致 UUID 变化),也保留该条分配记录,避免已分配用户被静默隐藏。
# 同时 LEFT JOIN employees 表获取员工姓名。
# assigned_at 定义为 NOT NULLnulls_last() 无意义,直接降序即可(SQLite/PG 通用)。 # assigned_at 定义为 NOT NULLnulls_last() 无意义,直接降序即可(SQLite/PG 通用)。
stmt = ( stmt = (
select(UserRole, Role) select(UserRole, Role, Employee)
.outerjoin(Role, UserRole.role_id == Role.id) .outerjoin(Role, UserRole.role_id == Role.id)
.outerjoin(Employee, UserRole.employee_id == Employee.employee_id)
.order_by(UserRole.assigned_at.desc()) .order_by(UserRole.assigned_at.desc())
) )
result = await db.execute(stmt) result = await db.execute(stmt)
rows = result.all() rows = result.all()
assignments = [] assignments = []
for user_role, role in rows: for user_role, role, employee in rows:
# role 可能为 None(孤儿记录),做兜底展示,而不是丢弃该用户 # role 可能为 None(孤儿记录),做兜底展示,而不是丢弃该用户
role_name = role.name if role else "unknown" role_name = role.name if role else "unknown"
role_display = (role.display_name or role.name) if role else "未知角色" role_display = (role.display_name or role.name) if role else "未知角色"
# employee 可能为 None(员工已从组织架构移除)
employee_name = employee.name if employee else ""
assignments.append({ assignments.append({
"employee_id": user_role.employee_id, "employee_id": user_role.employee_id,
"employee_name": employee_name,
"role_name": role_name, "role_name": role_name,
"role_display_name": role_display, "role_display_name": role_display,
"source": user_role.source or "manual", "source": user_role.source or "manual",
@@ -183,22 +191,46 @@ async def assign_role(
body: RoleAssignRequest, body: RoleAssignRequest,
admin: UserInfo = Depends(require_admin), admin: UserInfo = Depends(require_admin),
db: AsyncSession = Depends(get_db), db: AsyncSession = Depends(get_db),
redis: aioredis.Redis = Depends(dep_redis),
): ):
"""手动分配角色。 """手动分配角色。
为指定用户分配角色,记录分配者和分配原因。 为指定用户分配角色,记录分配者和分配原因。
支持输入「员工账号 或 姓名」:
- 先经 employee_directory.resolve_target 解析为企微 UserID
- 再用企微通讯录实时校验该员工确实属于企微组织架构(不在组织内则拒绝)。
安全限制:禁止管理员给自己分配角色。 安全限制:禁止管理员给自己分配角色。
Args: Args:
body: 分配角色请求 body: 分配角色请求target=账号/姓名,或兼容 employee_id
admin: 管理员(权限校验) admin: 管理员(权限校验)
db: 数据库会话 db: 数据库会话
redis: Redis 客户端(组织目录缓存用)
Returns: Returns:
Dict: 统一响应格式 Dict: 统一响应格式
""" """
# 取待解析的目标(优先 target,兼容旧 employee_id
raw = (body.target or body.employee_id or "").strip()
if not raw:
raise AppException(4001, "请输入员工账号或姓名")
# 解析 + 实时校验:是否为企微组织架构内的真实员工
resolved = await resolve_target(raw, db, redis)
if resolved.get("ambiguous"):
cands = resolved["candidates"]
names = "".join(f"{c['name']}({c['employee_id']})" for c in cands[:5])
raise AppException(
4016,
f"匹配到多名员工,请改用员工账号精确分配:{names}",
)
if not resolved.get("found"):
raise AppException(4017, resolved.get("reason", "员工不存在"))
employee_id = resolved["employee_id"]
# 安全限制:禁止管理员给自己分配角色 # 安全限制:禁止管理员给自己分配角色
if body.employee_id == admin.employee_id: if admin.employee_id and employee_id == admin.employee_id:
raise AppException(4014, "不能给自己分配角色") raise AppException(4014, "不能给自己分配角色")
# 查询目标角色 # 查询目标角色
@@ -211,7 +243,7 @@ async def assign_role(
# 检查是否已拥有该角色 # 检查是否已拥有该角色
existing_stmt = select(UserRole).where( existing_stmt = select(UserRole).where(
UserRole.employee_id == body.employee_id, UserRole.employee_id == employee_id,
UserRole.role_id == role.id, UserRole.role_id == role.id,
) )
existing_result = await db.execute(existing_stmt) existing_result = await db.execute(existing_stmt)
@@ -222,7 +254,7 @@ async def assign_role(
# 创建用户角色关联 # 创建用户角色关联
user_role = UserRole( user_role = UserRole(
employee_id=body.employee_id, employee_id=employee_id,
role_id=role.id, role_id=role.id,
source="manual", source="manual",
assigned_by=admin.employee_id, assigned_by=admin.employee_id,
@@ -230,11 +262,50 @@ async def assign_role(
db.add(user_role) db.add(user_role)
await db.commit() await db.commit()
logger.info(f"管理员 {_mask_sensitive_data(admin.employee_id)} 为用户 {_mask_sensitive_data(body.employee_id)} 分配角色 {body.role_name},原因:{body.reason}") logger.info(f"管理员 {_mask_sensitive_data(admin.employee_id)} 为用户 {_mask_sensitive_data(employee_id)} 分配角色 {body.role_name},原因:{body.reason}")
return success_response(message=f"角色 {body.role_name} 分配成功") return success_response(message=f"角色 {body.role_name} 分配成功")
# ---------- GET /api/admin/roles/employees/search ----------
@router.get("/employees/search")
async def search_employees(
q: str = Query("", description="姓名或员工账号关键字"),
admin: UserInfo = Depends(require_admin),
db: AsyncSession = Depends(get_db),
redis: aioredis.Redis = Depends(dep_redis),
):
"""员工目录搜索(分配角色时自动补全用)。
返回匹配关键字的员工候选(账号 + 姓名 + 部门),并在 full_directory 字段
标明当前目录是否覆盖全公司(True=企微全组织;False=仅已登录员工,需开通
通讯录读取权限才能搜索全公司)。
Args:
q: 搜索关键字(姓名或账号,至少 1 个字符)
admin: 管理员(权限校验)
db: 数据库会话
redis: Redis 客户端
Returns:
Dict: { items: [...], full_directory: bool|None }
"""
q = (q or "").strip()
if len(q) < 1:
return success_response(data={"items": [], "full_directory": None})
directory, full = await get_org_directory(db, redis)
ql = q.lower()
hits = [
m for m in directory
if ql in (m["employee_id"] or "").lower()
or ql in (m["name"] or "").lower()
]
# 优先精确账号命中,其次姓名包含
items = hits[:20]
return success_response(data={"items": items, "full_directory": full})
# ---------- POST /api/admin/roles/revoke ---------- # ---------- POST /api/admin/roles/revoke ----------
@router.post("/revoke") @router.post("/revoke")
async def revoke_role( async def revoke_role(
+1
View File
@@ -243,6 +243,7 @@ async def send_message(
try: try:
# 构建消息载荷 # 构建消息载荷
msg_payload = MessageResponse.model_validate(message).model_dump() msg_payload = MessageResponse.model_validate(message).model_dump()
msg_payload["message_id"] = msg_payload["id"]
# 构建 WebSocket 事件 # 构建 WebSocket 事件
ws_event = { ws_event = {
+6
View File
@@ -243,6 +243,12 @@ api_router.include_router(auth_wecom_sso_router, tags=["企微SSO"])
from app.api.audit_logs import router as audit_logs_router from app.api.audit_logs import router as audit_logs_router
api_router.include_router(audit_logs_router, tags=["审计日志"]) api_router.include_router(audit_logs_router, tags=["审计日志"])
# 运行期日志 API (standard SOP 第三阶段 — 管理后台日志体系)
# GET /api/admin/runtime-logs — 分页 + 多条件筛选(级别/时间/关键字)
# GET /api/admin/runtime-logs?download=true — 命中行文本下载(附件)
from app.api.runtime_logs import router as runtime_logs_router
api_router.include_router(runtime_logs_router, tags=["运行期日志"])
# 阶段5 自动化闭环 API # 阶段5 自动化闭环 API
# POST /itportal/automation/sessions — 创建自动化会话 # POST /itportal/automation/sessions — 创建自动化会话
# GET /itportal/automation/sessions — 会话列表 # GET /itportal/automation/sessions — 会话列表
+71
View File
@@ -0,0 +1,71 @@
# =============================================================================
# 企微IT智能服务台 — 运行期日志 API (standard SOP 第三阶段)
# =============================================================================
# 说明:管理后台查看后端运行期日志
# GET /admin/runtime-logs 正常查询(分页 + 级别/时间/关键字筛选)
# GET /admin/runtime-logs?download=true 下载命中行(text/plain 附件)
# 权限:require_admin(非 admin 返回 403,由依赖装饰器强制)
# =============================================================================
import logging
from datetime import datetime
from typing import Optional
from fastapi import APIRouter, Depends, Query
from fastapi.responses import StreamingResponse
from app.dependencies import require_admin, get_current_user, UserInfo
from app.services.runtime_log_service import (
iter_runtime_log_lines,
query_runtime_logs,
)
from app.utils.response import success_response
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/runtime-logs", tags=["运行期日志"])
@router.get("")
@require_admin
async def get_runtime_logs(
level: str = Query("INFO", description="日志级别阈值(DEBUG/INFO/WARNING/ERROR/CRITICAL),返回 >= 该级别"),
from_time: Optional[datetime] = Query(None, alias="from", description="起始时间(ISO8601)"),
to_time: Optional[datetime] = Query(None, alias="to", description="结束时间(ISO8601)"),
keyword: Optional[str] = Query(None, description="按消息关键字子串筛选(大小写不敏感)"),
page: int = Query(1, ge=1, description="页码"),
page_size: int = Query(50, ge=1, le=500, description="每页条数"),
download: bool = Query(False, description="true 时返回命中行文本下载"),
current_user: UserInfo = Depends(get_current_user),
):
"""查询后端运行期日志(分页 + 级别/时间/关键字筛选)。
权限:仅 admin 角色可访问,非 admin 由 require_admin 装饰器返回 403。
"""
# 下载模式:流式返回命中行文本,触发浏览器文件下载
if download:
ts = datetime.now().strftime("%Y%m%d-%H%M%S")
headers = {
"Content-Disposition": f'attachment; filename="runtime-logs-{ts}.log"'
}
return StreamingResponse(
iter_runtime_log_lines(
level=level,
from_time=from_time,
to_time=to_time,
keyword=keyword,
),
media_type="text/plain; charset=utf-8",
headers=headers,
)
# 普通查询:分页返回结构化条目
result = await query_runtime_logs(
level=level,
from_time=from_time,
to_time=to_time,
keyword=keyword,
page=page,
page_size=page_size,
)
return success_response(data=result)
+9
View File
@@ -6,6 +6,7 @@
# 所有配置项集中管理,避免散落在代码各处 # 所有配置项集中管理,避免散落在代码各处
# ============================================================================= # =============================================================================
import os
from typing import List from typing import List
import redis.asyncio as aioredis import redis.asyncio as aioredis
@@ -72,6 +73,14 @@ class Settings(BaseSettings):
# CORS 允许的源地址(逗号分隔的字符串) # CORS 允许的源地址(逗号分隔的字符串)
cors_origins: str = "http://localhost:5173,http://localhost:5174,http://localhost:5175" cors_origins: str = "http://localhost:5173,http://localhost:5174,http://localhost:5175"
# ----------------------------------------------------------------------
# 运行期日志目录(供"运行期日志"管理页面读取)
# ----------------------------------------------------------------------
# 后端运行期日志(JSON 格式,由 utils.logging_config.setup_logging 写入)
# 所在目录。容器内默认 /app/logs,可通过环境变量 RUNTIME_LOG_DIR 覆盖;
# 宿主机需将真实日志目录挂载到该路径(见 docker-compose.yml backend 卷)。
RUNTIME_LOG_DIR: str = os.getenv("RUNTIME_LOG_DIR", "/app/logs")
# ---------------------------------------------------------------------- # ----------------------------------------------------------------------
# AI 服务配置(Dify # AI 服务配置(Dify
# ---------------------------------------------------------------------- # ----------------------------------------------------------------------
+14
View File
@@ -28,6 +28,8 @@ from app.api.router import api_router
from app.dependencies import init_shared_services, cleanup_shared_services from app.dependencies import init_shared_services, cleanup_shared_services
# 导入异常处理器和异常类 # 导入异常处理器和异常类
from app.utils.response import AppException, app_exception_handler, success_response from app.utils.response import AppException, app_exception_handler, success_response
# 导入日志配置(运行期日志文件接线:standard SOP 第三阶段 · 运行期日志 D 页依赖)
from app.utils.logging_config import setup_logging
# 导入定时任务 # 导入定时任务
from app.tasks.reminder_task import check_unreplied_sessions from app.tasks.reminder_task import check_unreplied_sessions
@@ -75,6 +77,18 @@ async def lifespan(app: FastAPI):
在应用关闭时执行清理操作。 在应用关闭时执行清理操作。
""" """
# ===== 启动事件 ===== # ===== 启动事件 =====
# 运行期日志接线(standard SOP 第三阶段 · 运行期日志 D 页依赖):
# 在 stdout handler 之外,将结构化 JSON 日志写入 settings.RUNTIME_LOG_DIR/wecom-it-desk.log
# RotatingFileHandler:单文件上限 20MB,保留 5 个备份,复用 JSONFormatter)。
# 目录不可创建/写入时自动降级为仅 stdout,不阻塞启动。
# 注:setup_logging 会清空 root logger 已有 handler(含模块级 basicConfig),
# 并以 console(stdout) + 文件(JSON) 重新配置,避免重复输出。
setup_logging(
level=os.getenv("LOG_LEVEL", "INFO"),
json_format=True,
log_dir=settings.RUNTIME_LOG_DIR,
)
logger.info("🚀 企微IT智能服务台启动中...") logger.info("🚀 企微IT智能服务台启动中...")
# 校验关键配置项(防止生产环境忘记配置导致静默失败) # 校验关键配置项(防止生产环境忘记配置导致静默失败)
+4 -2
View File
@@ -80,12 +80,14 @@ class RoleAssignRequest(BaseModel):
"""角色分配请求 Schema。 """角色分配请求 Schema。
Attributes: Attributes:
employee_id: 企微 UserID target: 员工账号或姓名(自动解析为企微 UserID 并校验组织内存在性)【推荐】
employee_id: 兼容旧字段,直接传企微 UserID
role_name: 角色标识(user/agent/admin role_name: 角色标识(user/agent/admin
reason: 分配原因(可选) reason: 分配原因(可选)
""" """
employee_id: str = Field(..., min_length=1, max_length=100, description="企微 UserID") target: Optional[str] = Field(None, min_length=1, max_length=100, description="员工账号或姓名(自动解析)")
employee_id: Optional[str] = Field(None, min_length=1, max_length=100, description="兼容字段:企微 UserID")
role_name: str = Field(..., min_length=1, max_length=50, description="角色标识") role_name: str = Field(..., min_length=1, max_length=50, description="角色标识")
reason: Optional[str] = Field(None, max_length=500, description="分配原因") reason: Optional[str] = Field(None, max_length=500, description="分配原因")
+67 -9
View File
@@ -22,6 +22,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
from app.models.agent import Agent from app.models.agent import Agent
from app.models.config_change_log import ConfigChangeLog from app.models.config_change_log import ConfigChangeLog
from app.models.audit_log import AuditLog
from app.models.conversation import Conversation from app.models.conversation import Conversation
from app.models.message import Message from app.models.message import Message
from app.models.quick_reply_template import QuickReplyTemplate from app.models.quick_reply_template import QuickReplyTemplate
@@ -451,7 +452,10 @@ async def update_config(
old_value = config.config_value old_value = config.config_value
# 写入变更日志 # ===== A+B 双写(决策2:配置变更同时进入 config_change_logs 与 audit_logs=====
# 表 Aconfig_change_logs(既有,供"配置变更历史"页)
# 表 Baudit_logs 的 config_change 事件(供"安全审计日志"页,复用 record_audit_log 统一构造)
# 两表在同一 DB session 内 add,由调用方统一 commit(与 record_audit_log 约定一致:本函数不 commit)。
change_log = ConfigChangeLog( change_log = ConfigChangeLog(
config_key=key, config_key=key,
old_value=old_value, old_value=old_value,
@@ -460,6 +464,26 @@ async def update_config(
) )
db.add(change_log) db.add(change_log)
# 表 Baudit_logs —— 严格复用 app.services.audit_log_service.record_audit_log 的构造方式
# (字段命名 / 默认值保持一致)。update_config 为服务层、无 FastAPI Request 上下文,
# 故 ip_address / user_agent 置 Nonerecord_audit_log 在无 request 时同样取 None)。
audit_log = AuditLog(
employee_id=agent_id,
action="config_change",
resource="system_config",
resource_id=key,
details={
"old": old_value,
"new": value,
"operator": agent_id,
},
result="success",
ip_address=None,
user_agent=None,
created_at=datetime.now(),
)
db.add(audit_log)
# 更新配置值 # 更新配置值
config.config_value = value config.config_value = value
config.updated_at = datetime.now() config.updated_at = datetime.now()
@@ -1690,19 +1714,53 @@ async def get_system_logs(
db: AsyncSession, db: AsyncSession,
page: int = 1, page: int = 1,
page_size: int = 50, page_size: int = 50,
config_key: Optional[str] = None,
changed_by: Optional[str] = None,
from_time: Optional[datetime] = None,
to_time: Optional[datetime] = None,
) -> Dict[str, Any]: ) -> Dict[str, Any]:
"""获取系统日志(配置变更日志)""" """获取系统日志(配置变更日志),支持配置键/操作人/时间范围筛选。
count_result = await db.execute(select(func.count(ConfigChangeLog.id)))
total = count_result.scalar() or 0
offset = (page - 1) * page_size Args:
result = await db.execute( db: 数据库会话
select(ConfigChangeLog) page: 页码,从 1 开始
page_size: 每页条数
config_key: 按配置键精确筛选(可选,空则不追加条件)
changed_by: 按操作人 agent_id 精确筛选(可选)
from_time: 变更时间起始(闭区间,可选)
to_time: 变更时间截止(闭区间,可选)
Returns:
Dict: {items, total, page, page_size}
"""
# 动态构建筛选条件(空参数不追加条件,保持旧行为)
conditions = []
if config_key:
conditions.append(ConfigChangeLog.config_key == config_key)
if changed_by:
conditions.append(ConfigChangeLog.changed_by == changed_by)
if from_time:
conditions.append(ConfigChangeLog.changed_at >= from_time)
if to_time:
conditions.append(ConfigChangeLog.changed_at <= to_time)
# 总数(带筛选)
count_stmt = select(func.count(ConfigChangeLog.id))
if conditions:
count_stmt = count_stmt.where(and_(*conditions))
total = (await db.execute(count_stmt)).scalar() or 0
# 分页查询(带筛选,按变更时间倒序)
query = select(ConfigChangeLog)
if conditions:
query = query.where(and_(*conditions))
query = (
query
.order_by(ConfigChangeLog.changed_at.desc()) .order_by(ConfigChangeLog.changed_at.desc())
.offset(offset) .offset((page - 1) * page_size)
.limit(page_size) .limit(page_size)
) )
logs = list(result.scalars().all()) logs = list((await db.execute(query)).scalars().all())
agent_ids = list({log.changed_by for log in logs if log.changed_by}) agent_ids = list({log.changed_by for log in logs if log.changed_by})
agent_names: Dict[str, str] = {} agent_names: Dict[str, str] = {}
+196
View File
@@ -0,0 +1,196 @@
# =============================================================================
# 员工目录解析服务
# =============================================================================
# 说明:角色分配时,将管理员输入的「员工账号 或 姓名」解析为企微 UserID,
# 并校验该员工确实属于企微组织架构(需求:分配角色时按姓名/账号自动转换 + 校验)。
#
# 数据源优先级(自动适配,无需改代码即可在权限开通后升级):
# 1. 企微通讯录(实时):
# - get_user_info(userid) 校验账号是否为组织内真实员工
# - get_department_members(1, 1) 拉取全组织架构,用于「姓名 -> 账号」匹配
# - 需要企微应用具备「通讯录读取」权限;权限不足(errcode 60011)时自动降级
# 2. 本地 employees 表(仅登录过的员工):作为降级目录,保证功能在缺权限时仍可用
#
# 设计目标:无论企微权限是否齐全,分配功能都可用;权限齐全时自动获得全公司
# 姓名搜索能力(full_directory=True),缺权限时仅覆盖已登录员工。
# =============================================================================
import json
import logging
from typing import Any, Dict, List, Optional, Tuple
import redis.asyncio as aioredis
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.employee import Employee
from app.services.wecom_service import WecomService
logger = logging.getLogger(__name__)
# 组织目录 Redis 缓存 key 与 TTL(10 分钟,避免频繁调用企微通讯录 API)
ORG_DIRECTORY_CACHE_KEY = "wecom:org_directory"
ORG_DIRECTORY_CACHE_TTL = 600
async def get_org_directory(
db: AsyncSession,
redis: Optional[aioredis.Redis],
) -> Tuple[List[Dict[str, Any]], bool]:
"""获取「组织目录」(员工账号 + 姓名 列表),用于姓名 -> 账号匹配。
优先返回缓存;缓存未命中时尝试从企微通讯录拉全组织(需通讯录读取权限)。
若企微权限不足或调用失败,降级到本地 employees 表。
Returns:
(directory, full_directory)
- directory: [{"employee_id": str, "name": str, "department": str}, ...]
- full_directory: True=来自企微全组织(覆盖全公司);False=仅本地已登录员工
"""
# 1. 尝试命中缓存(缓存一定来自企微全组织,full=True)
if redis:
try:
raw = await redis.get(ORG_DIRECTORY_CACHE_KEY)
if raw:
logger.debug("命中组织目录缓存")
return json.loads(raw.decode("utf-8")), True
except Exception as e:
logger.warning(f"读取组织目录缓存失败(降级): {e}")
# 2. 尝试从企微通讯录拉全组织
wecom = WecomService(redis_client=redis)
try:
members = await wecom.get_department_members(1, 1)
directory = [
{
"employee_id": m.get("userid", ""),
"name": m.get("name", "") or "",
"department": ",".join(str(d) for d in (m.get("department") or [])),
}
for m in members
if m.get("userid")
]
# 写入缓存(仅全组织结果缓存,降级结果不缓存以免长期误用)
if redis:
try:
await redis.setex(
ORG_DIRECTORY_CACHE_KEY,
ORG_DIRECTORY_CACHE_TTL,
json.dumps(directory, ensure_ascii=False),
)
except Exception as e:
logger.warning(f"写入组织目录缓存失败: {e}")
logger.info(f"组织目录来自企微全组织,共 {len(directory)}")
return directory, True
except Exception as e:
err_text = str(e)
if "60011" in err_text or "privilege" in err_text.lower():
logger.warning("企微通讯录部门读取权限不足,降级到本地 employees 表")
else:
logger.warning(f"企微通讯录获取失败,降级本地: {err_text}")
# 3. 降级:本地 employees 表(仅登录过的员工)
try:
result = await db.execute(
select(Employee.employee_id, Employee.name).where(Employee.employee_id != "")
)
rows = result.all()
directory = [
{"employee_id": r[0], "name": r[1] or "", "department": ""} for r in rows
]
logger.info(f"组织目录降级到本地 employees 表,共 {len(directory)}")
return directory, False
except Exception as e:
logger.error(f"本地 employees 表查询失败: {e}")
return [], False
async def resolve_target(
target: str,
db: AsyncSession,
redis: Optional[aioredis.Redis] = None,
) -> Dict[str, Any]:
"""将输入的「员工账号 或 姓名」解析为企微 UserID,并校验组织内存在性。
Returns(结构化结果,由调用方翻译为响应/异常):
{"found": True, "employee_id": str, "name": str, "source": str}
{"found": False, "reason": str, "suggestion": str}
{"ambiguous": True, "candidates": [{"employee_id","name","department"}, ...]}
"""
target = (target or "").strip()
if not target:
return {
"found": False,
"reason": "请输入员工账号或姓名",
"suggestion": "请填写企微员工账号或姓名后重试",
}
wecom = WecomService(redis_client=redis)
# 1) 先尝试按 userid 实时校验(企微组织内真实员工)
try:
info = await wecom.get_user_info(target)
# 成功 -> target 本身就是有效 userid
return {
"found": True,
"employee_id": info.get("userid") or target,
"name": info.get("name") or "",
"source": "wecom_userid",
}
except Exception as e:
logger.debug(f"get_user_info('{target}') 未命中(将尝试按姓名解析): {e}")
# 2) 按姓名(包含)解析
directory, full = await get_org_directory(db, redis)
ql = target.lower()
# 精确 userid 匹配优先(目录里可能存在)
exact = [m for m in directory if m["employee_id"] and m["employee_id"].lower() == ql]
# 姓名包含匹配
name_hits = [m for m in directory if m["name"] and ql in m["name"].lower()]
matches = exact if exact else name_hits
if len(matches) == 1:
m = matches[0]
# 二次实时校验该 userid 确实在组织内(网络可用时)
try:
info = await wecom.get_user_info(m["employee_id"])
return {
"found": True,
"employee_id": info.get("userid") or m["employee_id"],
"name": info.get("name") or m.get("name", ""),
"source": "wecom_name" if full else "local_name",
}
except Exception:
# 实时校验失败(网络/权限),但目录里有 -> 仍可用
return {
"found": True,
"employee_id": m["employee_id"],
"name": m.get("name", ""),
"source": "local_name",
}
if len(matches) > 1:
return {
"ambiguous": True,
"candidates": [
{
"employee_id": m["employee_id"],
"name": m["name"],
"department": m.get("department", ""),
}
for m in matches[:10]
],
}
# 未找到
if full:
return {
"found": False,
"reason": f"企微组织架构中未找到匹配「{target}」的员工",
"suggestion": "请确认姓名/账号拼写,或改为输入员工账号",
}
return {
"found": False,
"reason": f"未找到匹配「{target}」的员工",
"suggestion": "当前仅能按姓名搜索已登录过本系统的员工;请直接输入员工账号,或为企微应用开通「通讯录读取」权限以搜索全公司",
}
+220
View File
@@ -0,0 +1,220 @@
# =============================================================================
# 企微IT智能服务台 — 运行期日志服务
# =============================================================================
# 说明:读取后端运行期日志目录(config.RUNTIME_LOG_DIR)下的 *.log 文件,
# 逐行 json.loads(容错跳过非 JSON 行),按 级别 >= 阈值 / 时间区间 /
# 关键字 过滤,按时间倒序分页返回结构化条目;
# download 模式返回原始 JSON 行(供流式下载)。
# 目录不存在 / 不可读时返回空结果(不抛 500)。
# =============================================================================
import json
import logging
import os
import re
from datetime import datetime
from typing import Any, AsyncGenerator, Dict, List, Optional
from app.config import settings
logger = logging.getLogger(__name__)
# 日志级别数值映射(用于 >= 阈值 比较)
_LEVEL_ORDER = {
"DEBUG": 10,
"INFO": 20,
"WARNING": 30,
"ERROR": 40,
"CRITICAL": 50,
}
# 轮转备份命名形如 wecom-it-desk.log.1 / .2 ...
_ROTATED_RE = re.compile(r"\.log\.\d+$")
def _level_num(level: str) -> int:
"""将级别名转为数值(用于 >= 阈值比较)。"""
return _LEVEL_ORDER.get((level or "INFO").upper(), 20)
def _parse_ts(ts: Optional[str]) -> datetime:
"""解析 JSONFormatter 写入的时间戳(ISO8601,可能带 Z)。
解析失败时返回 datetime.min,使其排在最后(不影响过滤语义)。
"""
if not ts:
return datetime.min
try:
return datetime.fromisoformat(ts.replace("Z", "+00:00"))
except (ValueError, TypeError):
try:
return datetime.strptime(ts, "%Y-%m-%dT%H:%M:%S.%f")
except (ValueError, TypeError):
return datetime.min
def _list_log_files() -> List[str]:
"""列出运行期日志目录下所有 .log 文件(含轮转备份)。
目录不存在 / 不可读时返回空列表(不抛异常)。
"""
log_dir = getattr(settings, "RUNTIME_LOG_DIR", None) or "/app/logs"
if not log_dir or not os.path.isdir(log_dir):
return []
files: List[str] = []
try:
for name in os.listdir(log_dir):
if name.endswith(".log") or _ROTATED_RE.search(name):
full = os.path.join(log_dir, name)
if os.path.isfile(full):
files.append(full)
except (OSError, PermissionError) as e:
logger.warning(f"读取运行期日志目录失败: {log_dir}: {e}")
return []
# 按文件名排序,保证读取顺序稳定(主日志优先于轮转备份)
files.sort()
return files
def _iter_raw_lines() -> List[str]:
"""逐文件读取所有日志行(原始文本,已去除行尾换行)。"""
lines: List[str] = []
for path in _list_log_files():
try:
with open(path, "r", encoding="utf-8", errors="ignore") as f:
for line in f:
stripped = line.rstrip("\n").rstrip("\r")
if stripped:
lines.append(stripped)
except (OSError, PermissionError) as e:
logger.warning(f"读取运行期日志文件失败: {path}: {e}")
continue
return lines
def _match(
entry: Dict[str, Any],
level_th: int,
from_time: Optional[datetime],
to_time: Optional[datetime],
keyword: Optional[str],
) -> bool:
"""判断单条日志是否命中筛选条件。"""
# 级别阈值:仅保留 >= 阈值的条目
if _level_num(entry.get("level", "INFO")) < level_th:
return False
# 时间区间(闭区间)
ts = _parse_ts(entry.get("timestamp"))
if from_time and ts < from_time:
return False
if to_time and ts > to_time:
return False
# 关键字:消息子串(大小写不敏感);为空则跳过该过滤
if keyword:
msg = str(entry.get("message", ""))
if keyword.lower() not in msg.lower():
return False
return True
async def query_runtime_logs(
level: str = "INFO",
from_time: Optional[datetime] = None,
to_time: Optional[datetime] = None,
keyword: Optional[str] = None,
page: int = 1,
page_size: int = 50,
) -> Dict[str, Any]:
"""查询运行期日志(分页 + 多条件筛选)。
读取 config.RUNTIME_LOG_DIR 下 *.log(含轮转备份),逐行 json.loads
(容错跳过非 JSON 行),按 级别>=阈值 / 时间区间 / 关键字 过滤,
按时间倒序分页返回结构化条目。
目录不存在 / 不可读时返回空列表(不抛 500)。
Args:
level: 级别阈值(DEBUG/INFO/WARNING/ERROR/CRITICAL),返回 >= 该级别
from_time: 起始时间(可选)
to_time: 结束时间(可选)
keyword: 消息关键字子串(可选,大小写不敏感)
page: 页码,从 1 开始
page_size: 每页条数
Returns:
Dict: {items, total, page, page_size}
"""
level_th = _level_num(level)
entries: List[Dict[str, Any]] = []
for raw in _iter_raw_lines():
try:
entry = json.loads(raw)
except (json.JSONDecodeError, ValueError):
# 容错:跳过非 JSON 行
continue
if not isinstance(entry, dict):
continue
if _match(entry, level_th, from_time, to_time, keyword):
entries.append(entry)
# 按时间倒序(新 -> 旧)
entries.sort(key=lambda e: _parse_ts(e.get("timestamp")), reverse=True)
total = len(entries)
page = max(1, page)
page_size = max(1, page_size)
start = (page - 1) * page_size
page_items = entries[start:start + page_size]
items = []
for e in page_items:
items.append({
"timestamp": e.get("timestamp", ""),
"level": e.get("level", ""),
"logger": e.get("logger", ""),
"message": e.get("message", ""),
"module": e.get("module", ""),
"function": e.get("function", ""),
"line": e.get("line", ""),
"request_id": e.get("request_id"),
"user_id": e.get("user_id"),
})
return {
"items": items,
"total": total,
"page": page,
"page_size": page_size,
}
async def iter_runtime_log_lines(
level: str = "INFO",
from_time: Optional[datetime] = None,
to_time: Optional[datetime] = None,
keyword: Optional[str] = None,
) -> AsyncGenerator[str, None]:
"""流式返回命中的原始 JSON 行(供下载)。
逐行读取并过滤,命中即 yield 原始行文本(带换行),
便于 StreamingResponse 直接流式输出为 .log 文件。
Args:
level: 级别阈值
from_time: 起始时间(可选)
to_time: 结束时间(可选)
keyword: 消息关键字子串(可选)
Yields:
str: 命中的原始 JSON 日志行(含行尾换行)
"""
level_th = _level_num(level)
for raw in _iter_raw_lines():
try:
entry = json.loads(raw)
except (json.JSONDecodeError, ValueError):
continue
if not isinstance(entry, dict):
continue
if _match(entry, level_th, from_time, to_time, keyword):
yield raw + "\n"
+30 -3
View File
@@ -6,9 +6,11 @@
import json import json
import logging import logging
import os
import sys import sys
from datetime import datetime from datetime import datetime
from typing import Any from logging.handlers import RotatingFileHandler
from typing import Any, Optional
class JSONFormatter(logging.Formatter): class JSONFormatter(logging.Formatter):
@@ -51,12 +53,16 @@ class PlainFormatter(logging.Formatter):
) )
def setup_logging(level: str = "INFO", json_format: bool = False) -> None: def setup_logging(level: str = "INFO", json_format: bool = False, log_dir: Optional[str] = None) -> None:
"""配置日志系统 """配置日志系统
Args: Args:
level: 日志级别 (DEBUG, INFO, WARNING, ERROR, CRITICAL) level: 日志级别 (DEBUG, INFO, WARNING, ERROR, CRITICAL)
json_format: 是否使用 JSON 格式输出 json_format: 是否使用 JSON 格式输出
log_dir: 运行期日志文件目录。若目录存在或可被创建,则在 stdout handler
之外额外写入 <log_dir>/wecom-it-desk.log
RotatingFileHandler,单文件上限 20MB,保留 5 个备份,复用 JSONFormatter)。
为空 / 不可创建时仅保留 stdout handler,不影响启动。
""" """
log_level = getattr(logging, level.upper(), logging.INFO) log_level = getattr(logging, level.upper(), logging.INFO)
@@ -68,7 +74,7 @@ def setup_logging(level: str = "INFO", json_format: bool = False) -> None:
for handler in root_logger.handlers[:]: for handler in root_logger.handlers[:]:
root_logger.removeHandler(handler) root_logger.removeHandler(handler)
# 创建 console handler # 创建 console handler(保持原有 stdout 行为不变)
console_handler = logging.StreamHandler(sys.stdout) console_handler = logging.StreamHandler(sys.stdout)
console_handler.setLevel(log_level) console_handler.setLevel(log_level)
@@ -81,6 +87,27 @@ def setup_logging(level: str = "INFO", json_format: bool = False) -> None:
console_handler.setFormatter(formatter) console_handler.setFormatter(formatter)
root_logger.addHandler(console_handler) root_logger.addHandler(console_handler)
# 运行期日志文件 handler(独立于 stdout,便于"运行期日志"管理页面读取)
# 始终使用 JSONFormatter,保证日志文件为结构化 JSON 行,便于解析与下载。
if log_dir:
try:
os.makedirs(log_dir, exist_ok=True)
log_file = os.path.join(log_dir, "wecom-it-desk.log")
file_handler = RotatingFileHandler(
log_file,
maxBytes=20 * 1024 * 1024, # 单文件上限 20MB
backupCount=5, # 保留 5 个轮转备份
encoding="utf-8",
)
file_handler.setLevel(log_level)
file_handler.setFormatter(JSONFormatter())
root_logger.addHandler(file_handler)
except (OSError, PermissionError) as e:
# 目录不可创建 / 不可写:不阻塞启动,仅记录告警,降级为仅 stdout 日志
logging.getLogger(__name__).warning(
f"无法为运行期日志创建文件 handlerlog_dir={log_dir}: {e};仅保留 stdout 日志。"
)
# 设置第三方库日志级别 # 设置第三方库日志级别
logging.getLogger("uvicorn").setLevel(logging.WARNING) logging.getLogger("uvicorn").setLevel(logging.WARNING)
logging.getLogger("fastapi").setLevel(logging.WARNING) logging.getLogger("fastapi").setLevel(logging.WARNING)