feat: 添加飞书用户模块和订阅功能支持 - 新增feishu_users模块用于处理飞书用户身份验证和权限管理 - 新增subscriptions模块用于处理订阅相关功能 - 新增personalization模块用于个性化服务 - 在alembic迁移配置中注册新的模型模块 - 在API路由器中添加feishu_users和subscriptions路由 - 实现事件调度服务的改进,包括错误处理和状态更新优化 - 添加飞书命令处理的权限检查机制 - 实现飞书应用票据事件处理 - 改进审计日志记录功能 ```
138 lines
5.0 KiB
Python
138 lines
5.0 KiB
Python
from app.modules.events.constants import (
|
|
EventAggregateType,
|
|
EventPayloadKey,
|
|
EventType,
|
|
)
|
|
from app.modules.events.models import DomainEvent
|
|
|
|
|
|
class UnsupportedEventTypeError(ValueError):
|
|
"""Raised when no application handler is registered for an event type."""
|
|
|
|
|
|
class EventHandlerMixin:
|
|
def _handle_event(self, record: DomainEvent) -> None:
|
|
if record.event_type == EventType.RISK_ACTION_RECORDED:
|
|
self._handle_risk_action(record)
|
|
return
|
|
if record.event_type in {
|
|
EventType.REPORT_PUSH_SUCCEEDED,
|
|
EventType.REPORT_PUSH_FAILED,
|
|
EventType.REPORT_GENERATED,
|
|
}:
|
|
self._handle_report_event(record)
|
|
return
|
|
if record.event_type in {
|
|
EventType.LEGACY_SYNC_COMPLETED,
|
|
EventType.LEGACY_SYNC_FAILED,
|
|
}:
|
|
self._handle_legacy_sync_event(record)
|
|
return
|
|
if record.event_type == EventType.AI_MEMORY_WRITTEN:
|
|
self._handle_ai_memory_event(record)
|
|
return
|
|
if record.event_type == EventType.ENTERPRISE_ANALYTICS_GENERATED:
|
|
self._handle_enterprise_analytics_event(record)
|
|
return
|
|
raise UnsupportedEventTypeError(
|
|
f"Unsupported domain event type: {record.event_type}"
|
|
)
|
|
|
|
def _handle_risk_action(self, record: DomainEvent) -> None:
|
|
from app.modules.risk.constants import RiskEventActionValue
|
|
from app.modules.workflows.service import WorkflowService
|
|
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
|
|
|
|
payload = record.payload or {}
|
|
action = str(payload.get(EventPayloadKey.ACTION) or "")
|
|
if action == RiskEventActionValue.CLOSE:
|
|
workflow_status = WorkflowStatus.COMPLETED
|
|
elif action == RiskEventActionValue.RESOLVE:
|
|
workflow_status = WorkflowStatus.WAITING_REVIEW
|
|
else:
|
|
workflow_status = WorkflowStatus.RUNNING
|
|
WorkflowService(self.db).start_or_update(
|
|
workflow_type=WorkflowType.RISK_EVENT_REVIEW,
|
|
aggregate_type=EventAggregateType.RISK_EVENT,
|
|
aggregate_id=record.aggregate_id,
|
|
status_value=workflow_status,
|
|
action=action or record.event_type,
|
|
actor=record.actor,
|
|
payload=payload,
|
|
commit=False,
|
|
source_event_id=record.event_id,
|
|
)
|
|
|
|
def _handle_report_event(self, record: DomainEvent) -> None:
|
|
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
|
|
|
|
workflow_status = (
|
|
WorkflowStatus.FAILED
|
|
if record.event_type == EventType.REPORT_PUSH_FAILED
|
|
else WorkflowStatus.COMPLETED
|
|
)
|
|
self._track_operational_workflow(
|
|
record,
|
|
workflow_type=WorkflowType.REPORT_DELIVERY,
|
|
workflow_status=workflow_status,
|
|
)
|
|
|
|
def _handle_legacy_sync_event(self, record: DomainEvent) -> None:
|
|
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
|
|
|
|
workflow_status = (
|
|
WorkflowStatus.FAILED
|
|
if record.event_type == EventType.LEGACY_SYNC_FAILED
|
|
else WorkflowStatus.COMPLETED
|
|
)
|
|
self._track_operational_workflow(
|
|
record,
|
|
workflow_type=WorkflowType.LEGACY_SYNC_MONITOR,
|
|
workflow_status=workflow_status,
|
|
)
|
|
|
|
def _handle_ai_memory_event(self, record: DomainEvent) -> None:
|
|
from app.modules.ai_memory.constants import AIMemoryPayloadKey, AIMemoryStatus
|
|
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
|
|
|
|
payload = record.payload or {}
|
|
workflow_status = (
|
|
WorkflowStatus.BLOCKED
|
|
if payload.get(AIMemoryPayloadKey.STATUS) == AIMemoryStatus.REJECTED
|
|
else WorkflowStatus.COMPLETED
|
|
)
|
|
self._track_operational_workflow(
|
|
record,
|
|
workflow_type=WorkflowType.AI_MEMORY_CAPTURE,
|
|
workflow_status=workflow_status,
|
|
)
|
|
|
|
def _handle_enterprise_analytics_event(self, record: DomainEvent) -> None:
|
|
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
|
|
|
|
self._track_operational_workflow(
|
|
record,
|
|
workflow_type=WorkflowType.ENTERPRISE_ANALYTICS,
|
|
workflow_status=WorkflowStatus.COMPLETED,
|
|
)
|
|
|
|
def _track_operational_workflow(
|
|
self,
|
|
record: DomainEvent,
|
|
workflow_type: str,
|
|
workflow_status: str,
|
|
) -> None:
|
|
from app.modules.workflows.service import WorkflowService
|
|
|
|
WorkflowService(self.db).start_or_update(
|
|
workflow_type=workflow_type,
|
|
aggregate_type=record.aggregate_type,
|
|
aggregate_id=record.aggregate_id,
|
|
status_value=workflow_status,
|
|
action=record.event_type,
|
|
actor=record.actor,
|
|
payload=record.payload or {},
|
|
commit=False,
|
|
source_event_id=record.event_id,
|
|
)
|