refactor(Dockerfile): 使用requirements.txt替代硬编码依赖

将Dockerfile中的硬编码pip包列表替换为通过requirements.txt文件安装,
提高依赖管理的灵活性和可维护性。

feat(scheduling): 移除内置APScheduler,采用独立调度系统

移除app/core/background/scheduler.py中原来的APScheduler实现,
改为使用新的应用级调度系统app.application.scheduling。

refactor(task_queue): 调整任务队列模块结构和导入路径

将任务队列相关常量从app.core.background.task_queue.constants迁移至
app.tasks.constants,并更新所有相关导入路径和引用。

refactor(events): 将事件服务重构为独立的应用层组件

将事件分发逻辑从核心层迁移到应用层,使用app.application.events.EventDispatchService
替代原有的app.modules.events.services.EventService。

feat(ai_memory): 增强AI记忆自动写入的安全策略

新增ai_memory_blocked_content_terms配置项用于阻止敏感内容,
添加TTL过期机制控制自动写入条目的生命周期。

fix(security): 强化生产环境安全验证机制

增加model_validator确保生产环境中数据库连接、API密钥、CORS设置等
关键安全配置符合要求。

feat(risks): 优化风险事件操作动作的外键约束

为RiskEventAction模型的风险事件ID字段添加外键约束,
防止孤立记录并增强数据完整性。

refactor(audit): 优化审计服务方法命名和事务处理

将AuditService的log方法重命名为record以反映其阶段行为,
并调整事务提交时机以提高性能。

feat(events): 增强领域事件并发处理和响应模型

添加事件锁定机制防止重复处理,更新API响应模型以提供
更准确的数据类型定义。
```
This commit is contained in:
2026-07-15 16:36:42 +08:00
parent 267b01b9f4
commit db751f03b4
73 changed files with 1615 additions and 933 deletions

View File

@@ -0,0 +1 @@
"""Application-level orchestration across domain modules and adapters."""

View File

@@ -0,0 +1,3 @@
from app.application.delivery.reports import ReportDeliveryService
__all__ = ["ReportDeliveryService"]

View File

@@ -0,0 +1,119 @@
from sqlalchemy.orm import Session
from app.modules.audit.constants import AuditAction, AuditSource
from app.modules.audit.schemas import AuditLogCreate
from app.modules.audit.service import AuditService
from app.modules.events.constants import (
EventAggregateType,
EventPayloadKey,
EventSource,
EventType,
)
from app.modules.events.services import EventService
from app.modules.feishu.service import FeishuService
from app.modules.reports.constants import (
ReportPushStatus,
ReportResponseKey,
)
from app.modules.reports.chart import lifecycle_chart_alt, render_lifecycle_chart
from app.modules.reports.services import ReportService
class ReportDeliveryService:
"""Deliver generated reports without coupling report generation to Feishu."""
def __init__(self, db: Session):
self.db = db
self.reports = ReportService(db)
def push_report(
self,
report: dict,
receive_id: str | None,
receive_id_type: str,
actor: str,
push_run_code: str | None = None,
) -> dict:
report_type = str(report.get(ReportResponseKey.REPORT_TYPE) or report.get("type") or "report")
title = report.get(ReportResponseKey.TITLE)
push_run = (
self.reports._get_push_run(push_run_code)
if push_run_code
else self.reports.create_push_run(
report_type=report_type,
title=title,
receive_id=receive_id,
receive_id_type=receive_id_type,
actor=actor,
)
)
try:
feishu = FeishuService(self.db)
image_key = None
image_alt = None
chart_data = report.get("chart_data")
if chart_data:
image_result = feishu.upload_image(render_lifecycle_chart(chart_data), actor)
image_key = (image_result.get("data") or {}).get("image_key")
if not image_key:
raise ValueError("Feishu image upload did not return image_key")
image_alt = lifecycle_chart_alt(chart_data)
card = FeishuService.build_basic_card(
report[ReportResponseKey.TITLE],
report[ReportResponseKey.LINES],
image_key=image_key,
image_alt=image_alt,
)
result = feishu.send_card(card, receive_id, receive_id_type, actor)
except Exception as exc:
failed_run = self.reports.update_push_run(
push_run.code,
ReportPushStatus.FAILED,
error_message=str(exc),
commit=False,
)
EventService(self.db).enqueue(
event_type=EventType.REPORT_PUSH_FAILED,
source=EventSource.REPORTS,
aggregate_type=EventAggregateType.REPORT_PUSH_RUN,
aggregate_id=failed_run.code,
actor=actor,
payload={
EventPayloadKey.CODE: failed_run.code,
EventPayloadKey.STATUS: failed_run.status,
EventPayloadKey.ERROR_MESSAGE: failed_run.error_message,
},
idempotency_key=f"report-push:{failed_run.code}:{failed_run.status}",
)
self.db.commit()
raise
success_run = self.reports.update_push_run(
push_run.code,
ReportPushStatus.SUCCESS,
provider_response=result,
sent=True,
commit=False,
)
EventService(self.db).enqueue(
event_type=EventType.REPORT_PUSH_SUCCEEDED,
source=EventSource.REPORTS,
aggregate_type=EventAggregateType.REPORT_PUSH_RUN,
aggregate_id=success_run.code,
actor=actor,
payload={
EventPayloadKey.CODE: success_run.code,
EventPayloadKey.STATUS: success_run.status,
},
idempotency_key=f"report-push:{success_run.code}:{success_run.status}",
)
AuditService(self.db).record(
AuditLogCreate(
actor=actor,
source=AuditSource.REPORTS,
action=AuditAction.REPORT_PUSH,
target_id=push_run.code,
response_payload={"status": ReportPushStatus.SUCCESS},
)
)
self.db.commit()
return result

View File

@@ -0,0 +1,3 @@
from app.application.events.dispatch import EventDispatchService
__all__ = ["EventDispatchService"]

View File

@@ -0,0 +1,212 @@
from datetime import timedelta
from typing import Any
from uuid import uuid4
from fastapi import HTTPException, status
from sqlalchemy import or_, select
from app.application.events.handlers import EventHandlerMixin
from app.core.config import get_settings
from app.core.constants import ActorValue
from app.core.http.pagination import bounded_limit
from app.core.utils.time import utc_now
from app.modules.audit.constants import (
AuditAction,
AuditRiskLevel,
AuditSource,
AuditTargetType,
)
from app.modules.audit.schemas import AuditLogCreate
from app.modules.audit.service import AuditService
from app.modules.events.constants import (
EventErrorDetail,
EventPayloadKey,
EventStatus,
)
from app.modules.events.models import DomainEvent
from app.modules.events.services import EventService
from app.modules.events.services.serialization import _serialize_event
class EventDispatchService(EventHandlerMixin):
"""Claim and dispatch outbox events at the application boundary."""
def __init__(self, db: Any):
self.db = db
self.events = EventService(db)
def get_event(self, event_id: str) -> DomainEvent:
return self.events.get_event(event_id)
def dispatch_event(
self,
event_id: str,
worker_id: str | None = None,
preclaimed: bool = False,
) -> DomainEvent:
lock_owner = worker_id or f"api:{uuid4().hex}"
record = (
self.get_event(event_id)
if preclaimed
else self._claim_event(event_id, lock_owner)
)
if record.status == EventStatus.PROCESSED:
return record
if preclaimed and record.locked_by != lock_owner:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail=EventErrorDetail.EVENT_LOCKED,
)
if not preclaimed and record.locked_by != lock_owner:
return record
settings = get_settings()
try:
self._handle_event(record)
except Exception as exc:
retryable = record.attempts < self._max_attempts(record)
record.status = EventStatus.PENDING if retryable else EventStatus.FAILED
record.last_error = str(exc)
record.next_attempt_at = (
utc_now() + timedelta(seconds=settings.event_dispatch_retry_delay_seconds)
if retryable
else None
)
record.locked_by = None
record.locked_until = None
self.db.commit()
self.db.refresh(record)
self._audit_dispatch(record)
return record
record.status = EventStatus.PROCESSED
record.last_error = None
record.processed_at = utc_now()
record.next_attempt_at = None
record.locked_by = None
record.locked_until = None
self.db.commit()
self.db.refresh(record)
self._audit_dispatch(record)
return record
def dispatch_pending(
self,
limit: int = 100,
worker_id: str | None = None,
) -> list[dict[str, Any]]:
now = utc_now()
stmt = (
select(DomainEvent)
.where(
DomainEvent.status == EventStatus.PENDING,
or_(
DomainEvent.next_attempt_at.is_(None),
DomainEvent.next_attempt_at <= now,
),
or_(
DomainEvent.locked_until.is_(None),
DomainEvent.locked_until <= now,
),
or_(
DomainEvent.max_attempts.is_(None),
DomainEvent.attempts < DomainEvent.max_attempts,
),
)
.order_by(DomainEvent.id.asc())
.limit(bounded_limit(limit))
.with_for_update(skip_locked=True)
)
records = list(self.db.execute(stmt).scalars())
lock_owner = worker_id or f"worker:{uuid4().hex}"
locked_until = now + timedelta(
seconds=get_settings().event_dispatch_lock_seconds
)
for record in records:
record.locked_by = lock_owner
record.locked_until = locked_until
record.attempts += 1
self.db.commit()
return [
_serialize_event(
self.dispatch_event(
record.event_id,
worker_id=lock_owner,
preclaimed=True,
)
)
for record in records
]
def retry_event(self, event_id: str, actor: str = ActorValue.API) -> DomainEvent:
record = self.get_event(event_id)
if record.status == EventStatus.PROCESSED:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail=EventErrorDetail.EVENT_NOT_RETRYABLE,
)
record.status = EventStatus.PENDING
record.actor = actor
record.attempts = 0
record.max_attempts = record.max_attempts or get_settings().event_dispatch_max_attempts
record.last_error = None
record.locked_by = None
record.locked_until = None
record.next_attempt_at = utc_now()
self.db.commit()
self.db.refresh(record)
return record
def _audit_dispatch(self, record: DomainEvent) -> None:
AuditService(self.db).log(
AuditLogCreate(
actor=record.actor,
source=AuditSource.EVENTS,
action=AuditAction.EVENT_DISPATCH,
target_type=AuditTargetType.DOMAIN_EVENT,
target_id=record.event_id,
risk_level=AuditRiskLevel.LOW,
response_payload={
EventPayloadKey.STATUS: record.status,
EventPayloadKey.ATTEMPTS: record.attempts,
EventPayloadKey.ERROR_MESSAGE: record.last_error,
},
)
)
def _can_attempt(self, record: DomainEvent) -> bool:
return record.attempts < self._max_attempts(record)
def _claim_event(self, event_id: str, lock_owner: str) -> DomainEvent:
now = utc_now()
record = self.db.execute(
select(DomainEvent)
.where(DomainEvent.event_id == event_id)
.with_for_update(skip_locked=True)
).scalar_one_or_none()
if record is None:
self.db.rollback()
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail=EventErrorDetail.EVENT_LOCKED,
)
if record.status == EventStatus.PROCESSED or not self._can_attempt(record):
self.db.commit()
return record
if record.locked_until is not None and record.locked_until > now:
self.db.commit()
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail=EventErrorDetail.EVENT_LOCKED,
)
record.locked_by = lock_owner
record.locked_until = now + timedelta(
seconds=get_settings().event_dispatch_lock_seconds
)
record.status = EventStatus.PENDING
record.attempts += 1
self.db.commit()
self.db.refresh(record)
return record
@staticmethod
def _max_attempts(record: DomainEvent) -> int:
return record.max_attempts or get_settings().event_dispatch_max_attempts

View File

@@ -0,0 +1,128 @@
from app.modules.events.constants import (
EventAggregateType,
EventPayloadKey,
EventType,
)
from app.modules.events.models import DomainEvent
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
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,
)
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,
)

View File

@@ -0,0 +1,4 @@
from app.application.feishu.commands import FeishuCommandService
from app.application.feishu.events import FeishuEventService
__all__ = ["FeishuCommandService", "FeishuEventService"]

View File

@@ -0,0 +1,201 @@
import json
import re
from typing import Any
from sqlalchemy.orm import Session
from app.application.feishu.delivery import send_card_if_configured, send_text_if_configured
from app.application.feishu.handlers import (
handle_finance_command,
handle_market_command,
handle_rule_command,
)
from app.application.feishu.results import command_result
from app.core.constants import ActorValue
from app.modules.ai_agent.constants import AIResponseKey
from app.modules.ai_agent.service import AIService
from app.modules.audit.constants import AuditSource
from app.modules.feishu.constants import (
FEISHU_AI_REPLY_TITLE,
FEISHU_MENTION_PATTERN,
FEISHU_ZERO_WIDTH_SPACE,
FeishuCommandKey,
FeishuCommandName,
FeishuPayloadKey,
FeishuReplyType,
)
from app.modules.feishu.service import FeishuService
from app.modules.reports.constants import ReportResponseKey
from app.modules.reports.services import ReportService
DAILY_REPORT_KEYWORDS = ("日报", "晨报", "经营日报", "经营晨报")
PROJECT_WEEKLY_KEYWORDS = ("周报", "项目周报")
ATTENDANCE_KEYWORDS = ("打卡", "考勤", "attendance")
RISK_KEYWORDS = ("风险", "预警", "risk")
AI_COMMAND_PREFIXES = ("", "ai ", "AI ", "/ask ")
DEFAULT_AI_PROMPT = "请说明你能做什么。"
def _parse_content_text(content: Any) -> str:
"""Extract plain command text from a Feishu message content payload."""
if isinstance(content, dict):
return str(
content.get(FeishuPayloadKey.TEXT) or content.get(FeishuPayloadKey.CONTENT) or ""
)
if not isinstance(content, str):
return ""
try:
data = json.loads(content)
except json.JSONDecodeError:
return content
if isinstance(data, dict):
return str(data.get(FeishuPayloadKey.TEXT) or data.get(FeishuPayloadKey.CONTENT) or "")
return content
def _clean_command_text(text: str) -> str:
"""Remove mentions and invisible characters from Feishu command text."""
text = re.sub(FEISHU_MENTION_PATTERN, "", text or "")
text = text.replace(FEISHU_ZERO_WIDTH_SPACE, "")
return text.strip()
class FeishuCommandService:
"""Route Feishu text commands to focused application handlers."""
def __init__(self, db: Session):
self.db = db
self.feishu = FeishuService(db)
def extract_event_command(self, payload: dict[str, Any]) -> dict[str, Any] | None:
event = payload.get(FeishuPayloadKey.EVENT) or {}
message = event.get(FeishuPayloadKey.MESSAGE) or {}
if not message:
return None
text = _clean_command_text(_parse_content_text(message.get(FeishuPayloadKey.CONTENT)))
if not text:
return None
sender = event.get(FeishuPayloadKey.SENDER) or {}
sender_id = sender.get(FeishuPayloadKey.SENDER_ID) or {}
actor = (
sender_id.get(FeishuPayloadKey.OPEN_ID)
or sender_id.get(FeishuPayloadKey.USER_ID)
or ActorValue.FEISHU
)
return {
FeishuCommandKey.TEXT: text,
FeishuCommandKey.CHAT_ID: message.get(FeishuCommandKey.CHAT_ID),
FeishuCommandKey.ACTOR: actor,
}
def handle_text(
self,
text: str,
chat_id: str | None = None,
actor: str = ActorValue.FEISHU,
auto_reply: bool = True,
) -> dict[str, Any]:
command_text = _clean_command_text(text)
lowered = command_text.lower()
for handler in (
handle_rule_command,
handle_finance_command,
handle_market_command,
):
result = handler(
self.db,
self.feishu,
command_text,
chat_id,
actor,
auto_reply,
)
if result is not None:
return result
report_result = self._handle_report_command(command_text, chat_id, actor, auto_reply)
if report_result is not None:
return report_result
return self._handle_ai_command(command_text, lowered, chat_id, actor, auto_reply)
def _handle_report_command(
self,
command_text: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any] | None:
service = ReportService(self.db)
if any(keyword in command_text for keyword in DAILY_REPORT_KEYWORDS):
command = FeishuCommandName.DAILY_BRIEF
report = service.daily_brief()
elif any(keyword in command_text for keyword in PROJECT_WEEKLY_KEYWORDS):
command = FeishuCommandName.PROJECT_WEEKLY
report = service.project_weekly()
elif any(keyword in command_text for keyword in ATTENDANCE_KEYWORDS):
command = FeishuCommandName.ATTENDANCE_SUMMARY
report = service.attendance_summary()
elif any(keyword in command_text for keyword in RISK_KEYWORDS):
command = FeishuCommandName.RISK_SUMMARY
report = service.risk_progress()
else:
return None
response = (
send_card_if_configured(
self.feishu,
chat_id,
report[ReportResponseKey.TITLE],
report[ReportResponseKey.LINES],
actor,
)
if auto_reply
else None
)
return command_result(
command,
FeishuReplyType.CARD,
report[ReportResponseKey.TITLE],
report[ReportResponseKey.CONTENT],
response,
report[ReportResponseKey.LINES],
)
def _handle_ai_command(
self,
command_text: str,
lowered: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any]:
prompt = command_text
for prefix in AI_COMMAND_PREFIXES:
if command_text.startswith(prefix):
prompt = command_text[len(prefix) :].strip()
break
if not prompt:
prompt = DEFAULT_AI_PROMPT
ai_result = AIService(self.db).ask(
prompt,
context={},
actor=actor,
source=AuditSource.FEISHU,
)
content = ai_result[AIResponseKey.ANSWER]
is_explicit_ai = any(
command_text.startswith(prefix) or lowered.startswith(prefix)
for prefix in AI_COMMAND_PREFIXES
)
response = (
send_text_if_configured(self.feishu, chat_id, content, actor) if auto_reply else None
)
return command_result(
FeishuCommandName.AI_ASK if is_explicit_ai else FeishuCommandName.FALLBACK_AI,
FeishuReplyType.TEXT,
FEISHU_AI_REPLY_TITLE,
content,
response,
)

View File

@@ -0,0 +1,34 @@
from typing import Any
from app.core.config import get_settings
from app.modules.feishu.service import FeishuService
def send_text_if_configured(
feishu: FeishuService,
chat_id: str | None,
text: str,
actor: str,
) -> dict[str, Any] | None:
"""Send a text reply only when Feishu credentials are configured."""
settings = get_settings()
if not (settings.feishu_app_id and settings.feishu_app_secret):
return None
return feishu.send_text(text, receive_id=chat_id, actor=actor)
def send_card_if_configured(
feishu: FeishuService,
chat_id: str | None,
title: str,
lines: list[str],
actor: str,
) -> dict[str, Any] | None:
"""Send a basic card only when Feishu credentials are configured."""
settings = get_settings()
if not (settings.feishu_app_id and settings.feishu_app_secret):
return None
card = FeishuService.build_basic_card(title, lines)
return feishu.send_card(card, receive_id=chat_id, actor=actor)

View File

@@ -0,0 +1,144 @@
from typing import Any
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.application.feishu.commands import FeishuCommandService
from app.core.constants import ActorValue
from app.modules.audit.constants import AuditAction, AuditSource
from app.modules.audit.schemas import AuditLogCreate
from app.modules.feishu.constants import (
FeishuCommandKey,
FeishuEventReceiptKey,
FeishuEventSource,
FeishuPayloadKey,
FeishuResponseKey,
)
from app.modules.feishu.models import FeishuEventReceipt
from app.modules.feishu.service import FeishuService
FEISHU_EVENT_ACTIONS = {
FeishuEventSource.WEBHOOK: AuditAction.FEISHU_WEBHOOK_EVENT,
FeishuEventSource.LONG_CONNECTION: AuditAction.FEISHU_LONG_CONNECTION_EVENT,
}
class FeishuEventService:
"""Handle Feishu message events from webhook or long connection."""
def __init__(self, db: Session):
self.db = db
self.feishu = FeishuService(db)
self.commands = FeishuCommandService(db)
def handle_event(
self,
payload: dict[str, Any],
source: str | FeishuEventSource,
auto_reply: bool = True,
) -> dict[str, Any]:
self.feishu.verify_event(payload)
challenge = payload.get(FeishuPayloadKey.CHALLENGE)
if challenge:
return {FeishuResponseKey.CHALLENGE: challenge}
source_value = _normalize_source(source)
event_identity = _event_identity(payload, source)
if event_identity and not self._register_event(event_identity):
return {
FeishuResponseKey.OK: True,
FeishuResponseKey.HANDLED: False,
FeishuResponseKey.DUPLICATE: True,
}
self.feishu.audit.log(
AuditLogCreate(
actor=ActorValue.FEISHU,
source=AuditSource.FEISHU,
action=FEISHU_EVENT_ACTIONS[source_value],
target_type=source_value,
target_id=(
event_identity.get(FeishuEventReceiptKey.EVENT_KEY)
if event_identity
else None
),
request_payload=_audit_event_metadata(payload),
response_payload={FeishuResponseKey.ACCEPTED: True},
)
)
command = self.commands.extract_event_command(payload)
if not command:
return {FeishuResponseKey.OK: True, FeishuResponseKey.HANDLED: False}
result = self.commands.handle_text(
command[FeishuCommandKey.TEXT],
chat_id=command[FeishuCommandKey.CHAT_ID],
actor=command[FeishuCommandKey.ACTOR],
auto_reply=auto_reply,
)
return {
FeishuResponseKey.OK: True,
FeishuResponseKey.HANDLED: True,
FeishuResponseKey.RESULT: result,
}
def _register_event(self, event_identity: dict[str, str | None]) -> bool:
receipt = FeishuEventReceipt(
event_key=str(event_identity[FeishuEventReceiptKey.EVENT_KEY]),
source=str(event_identity[FeishuEventReceiptKey.SOURCE]),
event_id=event_identity.get(FeishuEventReceiptKey.EVENT_ID),
message_id=event_identity.get(FeishuEventReceiptKey.MESSAGE_ID),
)
self.db.add(receipt)
try:
self.db.flush()
except IntegrityError:
self.db.rollback()
return False
return True
def _audit_event_metadata(payload: dict[str, Any]) -> dict[str, Any]:
"""Keep webhook audit evidence without storing message content or tokens."""
header = payload.get(FeishuPayloadKey.HEADER) or {}
event = payload.get(FeishuPayloadKey.EVENT) or {}
message = event.get(FeishuPayloadKey.MESSAGE) or {}
sender = event.get(FeishuPayloadKey.SENDER) or {}
sender_id = sender.get(FeishuPayloadKey.SENDER_ID) or {}
return {
"schema": payload.get("schema"),
FeishuPayloadKey.EVENT_ID: header.get(FeishuPayloadKey.EVENT_ID),
FeishuPayloadKey.EVENT_TYPE: header.get(FeishuPayloadKey.EVENT_TYPE),
FeishuPayloadKey.MESSAGE_ID: message.get(FeishuPayloadKey.MESSAGE_ID),
FeishuCommandKey.CHAT_ID: message.get(FeishuCommandKey.CHAT_ID),
FeishuPayloadKey.MESSAGE_TYPE: message.get(FeishuPayloadKey.MESSAGE_TYPE),
FeishuPayloadKey.OPEN_ID: sender_id.get(FeishuPayloadKey.OPEN_ID),
}
def _normalize_source(source: str | FeishuEventSource) -> FeishuEventSource:
return FeishuEventSource(source)
def _event_identity(
payload: dict[str, Any],
source: str | FeishuEventSource,
) -> dict[str, str | None] | None:
source_value = _normalize_source(source)
header = payload.get(FeishuPayloadKey.HEADER) or {}
event = payload.get(FeishuPayloadKey.EVENT) or {}
message = event.get(FeishuPayloadKey.MESSAGE) or {}
event_id = header.get(FeishuPayloadKey.EVENT_ID)
message_id = message.get(FeishuPayloadKey.MESSAGE_ID)
stable_id = event_id or message_id
if not stable_id:
return None
event_type = header.get(FeishuPayloadKey.EVENT_TYPE)
event_key = ":".join(
str(part)
for part in (source_value, event_type or FeishuPayloadKey.EVENT, stable_id)
)
return {
FeishuEventReceiptKey.EVENT_KEY: event_key,
FeishuEventReceiptKey.SOURCE: source_value,
FeishuEventReceiptKey.EVENT_ID: str(event_id) if event_id else None,
FeishuEventReceiptKey.MESSAGE_ID: str(message_id) if message_id else None,
}

View File

@@ -0,0 +1,9 @@
from app.application.feishu.handlers.finance import handle_finance_command
from app.application.feishu.handlers.market import handle_market_command
from app.application.feishu.handlers.rules import handle_rule_command
__all__ = [
"handle_finance_command",
"handle_market_command",
"handle_rule_command",
]

View File

@@ -0,0 +1,141 @@
import re
from typing import Any
from sqlalchemy.orm import Session
from app.application.feishu.delivery import send_text_if_configured
from app.application.feishu.results import command_result
from app.core.config import get_settings
from app.modules.ai_agent.constants import AIResponseKey
from app.modules.feishu.constants import FeishuCommandName, FeishuReplyType
from app.modules.feishu.service import FeishuService
from app.modules.reports.chart import lifecycle_chart_alt, render_lifecycle_chart
from app.modules.reports.constants import ReportResponseKey
from app.modules.reports.services import ReportService
PROJECT_FINANCE_PATTERN = re.compile(r"^项目资金\s+(.+)$")
FINANCE_COMMANDS = {"资金需求", "未来30天资金需求"}
def handle_finance_command(
db: Session,
feishu: FeishuService,
command_text: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any] | None:
"""Handle project cash-needs commands."""
project_match = PROJECT_FINANCE_PATTERN.fullmatch(command_text)
if command_text not in FINANCE_COMMANDS and project_match is None:
return None
command = (
FeishuCommandName.PROJECT_FINANCE if project_match else FeishuCommandName.FINANCE_NEEDS
)
if not get_settings().finance_needs_enabled:
return _text_result(
feishu,
command,
"项目资金需求分析",
"项目资金需求分析尚未启用,请先配置并启用财务只读同步。",
chat_id,
actor,
auto_reply,
)
project_code = project_match.group(1).strip() if project_match else None
service = ReportService(db)
preview = service.project_finance_needs_report(
project_code=project_code,
include_ai=False,
actor=actor,
)
if project_code and not preview["items"]:
return _text_result(
feishu,
command,
"项目资金需求分析",
f"未找到项目“{project_code}”,请使用稳定项目编号或展示编号。",
chat_id,
actor,
auto_reply,
)
if not preview["summary"]["data_available"]:
return _text_result(
feishu,
command,
"项目资金需求分析",
"项目财务数据未接入或无有效记录,暂不生成资金分析报告。",
chat_id,
actor,
auto_reply,
)
report = service.project_finance_needs_report(
project_code=project_code,
include_ai=True,
actor=actor,
)
ai_analysis = report.get("ai_analysis") or {}
if not ai_analysis.get(AIResponseKey.OK):
return _text_result(
feishu,
command,
"AI 暂不可用",
"AI 当前不可用,本次项目资金分析报告未发送。请检查模型服务。",
chat_id,
actor,
auto_reply,
)
provider_response = None
if auto_reply:
provider_response = _send_finance_card(feishu, chat_id, report, actor)
return command_result(
command,
FeishuReplyType.CARD,
report[ReportResponseKey.TITLE],
report[ReportResponseKey.CONTENT],
provider_response,
report[ReportResponseKey.LINES],
)
def _text_result(
feishu: FeishuService,
command: FeishuCommandName,
title: str,
content: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any]:
response = send_text_if_configured(feishu, chat_id, content, actor) if auto_reply else None
return command_result(command, FeishuReplyType.TEXT, title, content, response)
def _send_finance_card(
feishu: FeishuService,
chat_id: str | None,
report: dict[str, Any],
actor: str,
) -> dict[str, Any] | None:
settings = get_settings()
if not (settings.feishu_app_id and settings.feishu_app_secret):
return None
chart_data = {
"period": report.get("as_of"),
"finance": report.get("finance_chart_data"),
}
image_result = feishu.upload_image(render_lifecycle_chart(chart_data), actor)
image_key = (image_result.get("data") or {}).get("image_key")
if not image_key:
raise ValueError("Feishu image upload did not return image_key")
card = FeishuService.build_basic_card(
report[ReportResponseKey.TITLE],
report[ReportResponseKey.LINES],
image_key=image_key,
image_alt=lifecycle_chart_alt(chart_data),
)
return feishu.send_card(card, receive_id=chat_id, actor=actor)

View File

@@ -0,0 +1,243 @@
import re
from typing import Any
from fastapi import HTTPException
from sqlalchemy.orm import Session
from app.application.feishu.delivery import send_text_if_configured
from app.application.feishu.results import command_result
from app.core.config import get_settings
from app.modules.feishu.constants import FeishuCommandName, FeishuReplyType
from app.modules.feishu.service import FeishuService
from app.modules.market.chart import render_market_chart
from app.modules.market.service import MarketService
STOCK_ANALYSIS_PATTERN = re.compile(
r"^(?:股票分析|估值分析|财报分析)\s+([0-9]{6}(?:\.(?:SH|SZ|BJ))?)$", re.IGNORECASE
)
WATCHLIST_ADD_PATTERN = re.compile(r"^加入自选\s+([0-9]{6}(?:\.(?:SH|SZ|BJ))?)$", re.IGNORECASE)
MARKET_COMMANDS = {
"市场分析",
"今日收盘分析",
"本周市场分析",
"宏观金融分析",
"最新公告",
}
INDUSTRY_ANALYSIS_PATTERN = re.compile(r"^行业分析\s+(.+)$")
STOCK_COMPARE_PATTERN = re.compile(
r"^股票对比\s+([0-9]{6}(?:\.(?:SH|SZ|BJ))?)\s+"
r"([0-9]{6}(?:\.(?:SH|SZ|BJ))?)$",
re.IGNORECASE,
)
def handle_market_command(
db: Session,
feishu: FeishuService,
text: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any] | None:
"""Handle market analysis and watchlist commands."""
stock = STOCK_ANALYSIS_PATTERN.fullmatch(text)
add = WATCHLIST_ADD_PATTERN.fullmatch(text)
industry = INDUSTRY_ANALYSIS_PATTERN.fullmatch(text)
comparison = STOCK_COMPARE_PATTERN.fullmatch(text)
if (
text not in MARKET_COMMANDS
and text != "查看自选"
and not stock
and not add
and not industry
and not comparison
):
return None
if not get_settings().market_analysis_enabled:
return _text_result(
feishu,
FeishuCommandName.MARKET_OVERVIEW,
"市场分析",
"市场分析尚未启用,请配置市场数据源后启用。",
chat_id,
actor,
auto_reply,
)
service = MarketService(db)
if add:
if get_settings().read_only_mode:
return _text_result(
feishu,
FeishuCommandName.WATCHLIST_ADD,
"自选股",
"当前为只读模式,不能修改自选股。请由管理员启用操作后重试。",
chat_id,
actor,
auto_reply,
)
item = service.add_watchlist(actor, add.group(1))
return _text_result(
feishu,
FeishuCommandName.WATCHLIST_ADD,
"自选股",
f"已加入自选:{item['symbol']}",
chat_id,
actor,
auto_reply,
)
if text == "查看自选":
items = service.watchlist(actor)
content = "自选股:" + ("".join(item["symbol"] for item in items) or "暂无")
return _text_result(
feishu,
FeishuCommandName.WATCHLIST_LIST,
"自选股",
content,
chat_id,
actor,
auto_reply,
)
if text == "最新公告":
items = service.announcements(limit=10)["items"]
content = (
"最新公告:\n"
+ "\n".join(
f"- {item['announcement_date']} {item['symbol'] or '市场'}{item['title']}"
for item in items
)
if items
else "公告元数据尚未接入。"
)
return _text_result(
feishu,
FeishuCommandName.MARKET_ANNOUNCEMENTS,
"最新公告",
content,
chat_id,
actor,
auto_reply,
)
if industry:
try:
data = service.industry_analysis(industry.group(1).strip())
content = (
f"{data['industry']} 平均涨跌 {data['average_pct_change']}%\n"
+ "\n".join(
f"- {item['name']}{item['symbol']}{item['pct_change']}%"
for item in data["items"][:10]
)
)
except HTTPException:
content = "未找到该行业的最新市场数据。"
return _text_result(
feishu,
FeishuCommandName.MARKET_OVERVIEW,
"行业分析",
content,
chat_id,
actor,
auto_reply,
)
if comparison:
try:
content = service.compare_stocks([comparison.group(1), comparison.group(2)])["content"]
except HTTPException:
content = "至少一只股票缺少可用行情,暂时无法比较。"
return _text_result(
feishu,
FeishuCommandName.STOCK_ANALYSIS,
"股票对比",
content,
chat_id,
actor,
auto_reply,
)
return _analysis_result(service, feishu, text, stock.group(1) if stock else None, chat_id, actor, auto_reply)
def _analysis_result(
service: MarketService,
feishu: FeishuService,
text: str,
symbol: str | None,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any]:
command = FeishuCommandName.STOCK_ANALYSIS if symbol else FeishuCommandName.MARKET_OVERVIEW
if text == "宏观金融分析":
command = FeishuCommandName.MARKET_MACRO
try:
if symbol:
report = service.stock_analysis(symbol, True, actor)
elif text == "本周市场分析":
report = service.weekly_overview(include_ai=True, actor=actor)
elif text == "宏观金融分析":
report = service.macro_analysis(include_ai=True, actor=actor)
else:
report = service.market_overview(include_ai=True, actor=actor)
except HTTPException:
return _text_result(
feishu,
command,
"股票分析",
"未找到该股票的可用行情,请确认代码或先执行行情同步。",
chat_id,
actor,
auto_reply,
)
ai = report.get("ai_analysis") or {}
if not report.get("data_available"):
content = "市场数据未接入,暂不生成分析报告。"
elif not ai.get("ok"):
content = "AI 当前不可用,本次市场分析报告未发送。"
else:
content = report["content"]
response = _send_analysis(feishu, text, report, content, chat_id, actor) if auto_reply else None
return command_result(
command,
FeishuReplyType.CARD if ai.get("ok") and text != "宏观金融分析" else FeishuReplyType.TEXT,
report["title"],
content,
response,
report["lines"] if ai.get("ok") else None,
)
def _send_analysis(
feishu: FeishuService,
text: str,
report: dict[str, Any],
content: str,
chat_id: str | None,
actor: str,
) -> dict[str, Any] | None:
ai = report.get("ai_analysis") or {}
if not ai.get("ok") or text == "宏观金融分析":
return send_text_if_configured(feishu, chat_id, content, actor)
image = feishu.upload_image(render_market_chart(report), actor)
image_key = (image.get("data") or {}).get("image_key")
if not image_key:
raise ValueError("Feishu image upload did not return image_key")
card = FeishuService.build_basic_card(
report["title"],
report["lines"],
image_key=image_key,
image_alt=report["title"],
)
return feishu.send_card(card, receive_id=chat_id, actor=actor)
def _text_result(
feishu: FeishuService,
command: FeishuCommandName,
title: str,
content: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any]:
response = send_text_if_configured(feishu, chat_id, content, actor) if auto_reply else None
return command_result(command, FeishuReplyType.TEXT, title, content, response)

View File

@@ -0,0 +1,192 @@
import re
from typing import Any
from fastapi import HTTPException
from sqlalchemy.orm import Session
from app.application.feishu.delivery import send_text_if_configured
from app.application.feishu.results import command_result
from app.core.config import get_settings
from app.modules.ai_memory.constants import AIMemoryStatus
from app.modules.ai_memory.service import AIMemoryService
from app.modules.feishu.constants import FeishuCommandName, FeishuReplyType
from app.modules.feishu.service import FeishuService
RULE_TITLE = "AI 学习规则"
RULE_CREATE_PATTERN = re.compile(r"^学习规则(?:\s+(\d{1,3}))?\s*[:]\s*(.*)$")
MARKET_RULE_CREATE_PATTERN = re.compile(
r"^学习市场规则(?:\s+(\d{1,3}))?\s*[:]\s*(.*)$"
)
RULE_DISABLE_PATTERN = re.compile(r"^停用规则\s+(MEM-[A-Za-z0-9-]+)$", re.IGNORECASE)
RULE_ENABLE_PATTERN = re.compile(r"^启用规则\s+(MEM-[A-Za-z0-9-]+)$", re.IGNORECASE)
RULE_LIST_COMMANDS = {"查看规则", "规则列表", "查看市场规则"}
RULE_COMMAND_PREFIXES = (
"学习市场规则",
"学习规则",
"查看市场规则",
"查看规则",
"规则列表",
"停用规则",
"启用规则",
)
RULE_COMMAND_HELP = (
"规则指令格式:\n"
"学习规则:<规则内容>\n"
"学习规则 80<规则内容>\n"
"学习市场规则 80<仅用于市场分析的规则内容>\n"
"查看规则\n"
"停用规则 <规则编号>\n"
"启用规则 <规则编号>"
)
def handle_rule_command(
db: Session,
feishu: FeishuService,
command_text: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any] | None:
"""Handle persistent AI rule commands."""
if not command_text.startswith(RULE_COMMAND_PREFIXES):
return None
command = _command_name(command_text)
if command in {
FeishuCommandName.RULE_CREATE,
FeishuCommandName.RULE_DISABLE,
FeishuCommandName.RULE_ENABLE,
} and get_settings().read_only_mode:
return _result(
feishu,
command,
"当前为只读模式,不能新增或修改学习规则。请由管理员启用操作后重试。",
chat_id,
actor,
auto_reply,
)
content = RULE_COMMAND_HELP
try:
command, content = _execute(db, command_text, command, actor)
except HTTPException as exc:
detail = str(exc.detail)
if "secret-like" in detail:
content = "规则疑似包含密码、令牌或其他密钥信息,已拒绝学习。"
elif exc.status_code == 404:
content = "没有找到该规则,请先发送“查看规则”确认规则编号。"
elif "priority" in detail:
content = "规则优先级必须在 1 到 100 之间。"
else:
content = "规则未保存,请检查指令内容后重试。"
return _result(feishu, command, content, chat_id, actor, auto_reply)
def _execute(
db: Session,
command_text: str,
command: FeishuCommandName,
actor: str,
) -> tuple[FeishuCommandName, str]:
market_create_match = MARKET_RULE_CREATE_PATTERN.fullmatch(command_text)
create_match = market_create_match or RULE_CREATE_PATTERN.fullmatch(command_text)
disable_match = RULE_DISABLE_PATTERN.fullmatch(command_text)
enable_match = RULE_ENABLE_PATTERN.fullmatch(command_text)
memory = AIMemoryService(db)
if create_match:
return command, _create_rule(memory, create_match, market_create_match is not None, actor)
if command_text in RULE_LIST_COMMANDS:
return FeishuCommandName.RULE_LIST, _list_rules(memory, command_text)
if disable_match or enable_match:
enabled = enable_match is not None
match = enable_match or disable_match
rule = memory.update_rule(
code=match.group(1),
content=None,
priority=None,
tags=None,
enabled=enabled,
actor=actor,
)
state = "已启用" if enabled else "已停用"
return (
FeishuCommandName.RULE_ENABLE if enabled else FeishuCommandName.RULE_DISABLE,
f"规则{state}\n"
f"编号:{rule['code']}\n"
f"优先级:{rule['importance']}\n"
f"范围:{rule['scope']} / {rule['subject']}\n"
f"状态:{state}",
)
return command, RULE_COMMAND_HELP
def _create_rule(
memory: AIMemoryService,
match: re.Match[str],
market_rule: bool,
actor: str,
) -> str:
priority = int(match.group(1) or 50)
content = match.group(2).strip()
if not content:
return f"规则内容不能为空。\n\n{RULE_COMMAND_HELP}"
if not 1 <= priority <= 100:
return "规则优先级必须在 1 到 100 之间。"
rule = memory.create_rule(
content=content,
scope="market" if market_rule else "global",
subject="market" if market_rule else "company",
priority=priority,
tags=["feishu", *(["market"] if market_rule else [])],
actor=actor,
)
return (
"规则已学习。\n"
f"编号:{rule['code']}\n"
f"优先级:{rule['importance']}\n"
f"范围:{rule['scope']} / {rule['subject']}\n"
"状态:已启用"
)
def _list_rules(memory: AIMemoryService, command_text: str) -> str:
rules = memory.list_rules(
scope="market" if command_text == "查看市场规则" else None,
status_filter=AIMemoryStatus.ACTIVE,
limit=20,
)
if not rules:
return "当前没有已启用的学习规则。"
lines = ["当前已启用的学习规则:"]
for rule in rules:
rule_text = str(rule["content"])
if len(rule_text) > 80:
rule_text = f"{rule_text[:80]}"
lines.append(
f"{rule['code']}|优先级 {rule['importance']}"
f"{rule['scope']}/{rule['subject']}\n{rule_text}"
)
return "\n\n".join(lines)
def _command_name(command_text: str) -> FeishuCommandName:
if command_text.startswith("停用规则"):
return FeishuCommandName.RULE_DISABLE
if command_text.startswith("启用规则"):
return FeishuCommandName.RULE_ENABLE
if command_text.startswith(("查看市场规则", "查看规则", "规则列表")):
return FeishuCommandName.RULE_LIST
return FeishuCommandName.RULE_CREATE
def _result(
feishu: FeishuService,
command: FeishuCommandName,
content: str,
chat_id: str | None,
actor: str,
auto_reply: bool,
) -> dict[str, Any]:
response = send_text_if_configured(feishu, chat_id, content, actor) if auto_reply else None
return command_result(command, FeishuReplyType.TEXT, RULE_TITLE, content, response)

View File

@@ -0,0 +1,29 @@
from typing import Any
from app.modules.feishu.constants import (
FeishuCommandName,
FeishuCommandResultKey,
FeishuReplyType,
)
def command_result(
command: FeishuCommandName,
reply_type: FeishuReplyType,
title: str,
content: str,
provider_response: dict[str, Any] | None = None,
lines: list[str] | None = None,
) -> dict[str, Any]:
"""Build the stable command response contract."""
result: dict[str, Any] = {
FeishuCommandResultKey.COMMAND: command,
FeishuCommandResultKey.REPLY_TYPE: reply_type,
FeishuCommandResultKey.TITLE: title,
FeishuCommandResultKey.CONTENT: content,
FeishuCommandResultKey.PROVIDER_RESPONSE: provider_response,
}
if lines is not None:
result[FeishuCommandResultKey.LINES] = lines
return result

View File

@@ -0,0 +1,4 @@
from app.application.pipelines.lifecycle import LifecyclePipelineService
from app.application.pipelines.market import MarketPipelineService
__all__ = ["LifecyclePipelineService", "MarketPipelineService"]

View File

@@ -0,0 +1,237 @@
from datetime import date
from typing import Any
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.core.config import get_settings
from app.core.security import business_mutations_enabled
from app.core.constants import ActorValue
from app.core.utils.time import utc_now
from app.application.delivery import ReportDeliveryService
from app.modules.feishu.service import FeishuService
from app.modules.legacy_mysql.intasect import IntasectSyncService
from app.modules.reports.constants import ReportPushStatus, ReportType
from app.modules.reports.services import ReportService
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
from app.modules.workflows.models import WorkflowInstance
from app.modules.workflows.service import WorkflowService
class LifecyclePipelineService:
def __init__(self, db: Session):
self.db = db
self.workflows = WorkflowService(db)
def period_key(self, report_type: str, reference_date: date | None = None) -> str:
start, end = ReportService(self.db)._management_period(report_type, reference_date)
period_value = end.isoformat() if report_type == ReportType.DAILY else start.isoformat()
return f"{report_type}:{period_value}"
def find(self, period_key: str) -> WorkflowInstance | None:
return self.db.execute(
select(WorkflowInstance).where(
WorkflowInstance.workflow_type == WorkflowType.LIFECYCLE_REPORT,
WorkflowInstance.aggregate_type == "report_period",
WorkflowInstance.aggregate_id == period_key,
)
).scalar_one_or_none()
def prepare(
self,
report_type: str,
actor: str,
force: bool = False,
) -> tuple[WorkflowInstance, str, bool]:
period_key = self.period_key(report_type)
existing = self.find(period_key)
if existing is not None and existing.status == WorkflowStatus.COMPLETED and not force:
return existing, period_key, True
workflow = self.workflows.start_or_update(
workflow_type=WorkflowType.LIFECYCLE_REPORT,
aggregate_type="report_period",
aggregate_id=period_key,
status_value=WorkflowStatus.RUNNING,
action="queued",
actor=actor,
payload={"report_type": report_type, "period_key": period_key, "force": force},
)
return workflow, period_key, False
def run(
self,
report_type: str,
receive_id: str | None = None,
receive_id_type: str = "chat_id",
force: bool = False,
actor: str = ActorValue.SCHEDULER,
) -> dict[str, Any]:
if not business_mutations_enabled():
return {
"period_key": self.period_key(report_type),
"deduplicated": False,
"status": "operations_disabled",
}
workflow, period_key, deduplicated = self.prepare(report_type, actor, force)
if deduplicated:
return {
"workflow_code": workflow.code,
"period_key": period_key,
"deduplicated": True,
"status": workflow.status,
}
try:
workflow = self.workflows.start_or_update(
workflow_type=WorkflowType.LIFECYCLE_REPORT,
aggregate_type="report_period",
aggregate_id=period_key,
status_value=WorkflowStatus.RUNNING,
action="source_sync",
actor=actor,
payload={"report_type": report_type},
)
sync_result = IntasectSyncService(self.db).sync_all(
run_code=workflow.code,
force_full=False,
)
self.workflows.start_or_update(
workflow_type=WorkflowType.LIFECYCLE_REPORT,
aggregate_type="report_period",
aggregate_id=period_key,
status_value=WorkflowStatus.RUNNING,
action="analysis",
actor=actor,
payload={
"datasets": {name: result["processed"] for name, result in sync_result.items()}
},
)
report_service = ReportService(self.db)
report = report_service.management_lifecycle_report(
report_type=report_type,
actor=actor,
include_ai=True,
)
settings = get_settings()
target_receive_id = receive_id or settings.feishu_default_chat_id
ai_analysis = report.get("ai_analysis") or {}
if not ai_analysis.get("ok"):
notified = self._notify_ai_unavailable(
target_receive_id,
receive_id_type,
period_key,
actor,
)
failed = self.workflows.start_or_update(
workflow_type=WorkflowType.LIFECYCLE_REPORT,
aggregate_type="report_period",
aggregate_id=period_key,
status_value=WorkflowStatus.FAILED,
action="ai_unavailable",
actor=actor,
payload={
"ai_unavailable": True,
"notified": notified,
"error_type": ai_analysis.get("type") or "AIUnavailable",
},
)
return {
"workflow_code": failed.code,
"period_key": period_key,
"deduplicated": False,
"status": failed.status,
"ai_unavailable": True,
"notified": notified,
}
idempotency_key = period_key
if force:
idempotency_key = f"{period_key}:force:{utc_now():%Y%m%d%H%M%S%f}"
push_run = report_service.create_push_run(
report_type=report_type,
title=str(report["title"]),
receive_id=target_receive_id,
receive_id_type=receive_id_type,
actor=actor,
status=ReportPushStatus.PENDING,
idempotency_key=idempotency_key,
)
if push_run.status != ReportPushStatus.SUCCESS:
ReportDeliveryService(self.db).push_report(
report,
target_receive_id,
receive_id_type,
actor,
push_run_code=push_run.code,
)
completed = self.workflows.start_or_update(
workflow_type=WorkflowType.LIFECYCLE_REPORT,
aggregate_type="report_period",
aggregate_id=period_key,
status_value=WorkflowStatus.COMPLETED,
action="pushed",
actor=actor,
payload={"push_run_code": push_run.code, "idempotency_key": idempotency_key},
)
return {
"workflow_code": completed.code,
"period_key": period_key,
"push_run_code": push_run.code,
"deduplicated": False,
"status": completed.status,
}
except Exception as exc:
self.db.rollback()
self.workflows.start_or_update(
workflow_type=WorkflowType.LIFECYCLE_REPORT,
aggregate_type="report_period",
aggregate_id=period_key,
status_value=WorkflowStatus.FAILED,
action="failed",
actor=actor,
payload={"error": str(exc)[:2000]},
)
self._notify_failure(receive_id, receive_id_type, period_key, exc, actor)
raise
def _notify_ai_unavailable(
self,
receive_id: str | None,
receive_id_type: str,
period_key: str,
actor: str,
) -> bool:
settings = get_settings()
if (
not receive_id
or not settings.feishu_app_id
or not settings.feishu_app_secret
):
return False
FeishuService(self.db).send_text(
f"生命周期报告 {period_key}AI 当前不可用,本次分析报告未发送。请检查模型服务。",
receive_id,
receive_id_type,
actor,
)
return True
def _notify_failure(
self,
receive_id: str | None,
receive_id_type: str,
period_key: str,
error: Exception,
actor: str,
) -> None:
settings = get_settings()
target = receive_id or settings.feishu_default_chat_id
if not target or not settings.feishu_app_id or not settings.feishu_app_secret:
return
try:
FeishuService(self.db).send_text(
f"生命周期报告 {period_key} 执行失败:{type(error).__name__}",
target,
receive_id_type,
actor,
)
except Exception:
self.db.rollback()

View File

@@ -0,0 +1,264 @@
from datetime import date, timedelta
from typing import Any
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.core.config import get_settings
from app.core.constants import ActorValue
from app.core.security import business_mutations_enabled
from app.core.utils.time import utc_now
from app.modules.feishu.service import FeishuService
from app.modules.market.chart import render_market_chart
from app.modules.market.service import MarketService
from app.modules.reports.constants import ReportPushStatus
from app.modules.reports.services import ReportService
from app.modules.workflows.constants import WorkflowStatus, WorkflowType
from app.modules.workflows.models import WorkflowInstance
from app.modules.workflows.service import WorkflowService
MARKET_REPORT_TYPES = {"premarket", "close", "weekly"}
class MarketPipelineService:
def __init__(self, db: Session, market: MarketService | None = None):
self.db = db
self.market = market or MarketService(db)
self.workflows = WorkflowService(db)
def period_key(self, report_type: str, reference_date: date) -> str:
self._validate_type(report_type)
period = reference_date
if report_type == "weekly":
period = reference_date - timedelta(days=reference_date.weekday())
return f"market:{report_type}:{period.isoformat()}"
def find(self, period_key: str) -> WorkflowInstance | None:
return self.db.execute(
select(WorkflowInstance).where(
WorkflowInstance.workflow_type == WorkflowType.MARKET_ANALYSIS,
WorkflowInstance.aggregate_type == "market_period",
WorkflowInstance.aggregate_id == period_key,
)
).scalar_one_or_none()
def run(
self,
report_type: str,
reference_date: date | None = None,
force: bool = False,
actor: str = ActorValue.SCHEDULER,
) -> dict[str, Any]:
target = reference_date or date.today()
period_key = self.period_key(report_type, target)
if not business_mutations_enabled():
return {
"period_key": period_key,
"status": "operations_disabled",
"deduplicated": False,
}
existing = self.find(period_key)
if existing is not None and existing.status == WorkflowStatus.COMPLETED and not force:
return {
"workflow_code": existing.code,
"period_key": period_key,
"status": existing.status,
"deduplicated": True,
}
self._step(period_key, report_type, "source_sync", WorkflowStatus.RUNNING, actor)
try:
sync_result = self._sync(report_type, target)
if sync_result.get("market_closed"):
workflow = self._step(
period_key,
report_type,
"market_closed",
WorkflowStatus.COMPLETED,
actor,
sync_result,
)
return {
"workflow_code": workflow.code,
"period_key": period_key,
"status": "market_closed",
"deduplicated": False,
}
self._step(
period_key,
report_type,
"ai_analysis",
WorkflowStatus.RUNNING,
actor,
sync_result,
)
report = (
self.market.weekly_overview(target, include_ai=True, actor=actor)
if report_type == "weekly"
else self.market.market_overview(
target if report_type == "close" else None,
include_ai=True,
actor=actor,
)
)
ai = report.get("ai_analysis") or {}
if not report.get("data_available"):
return self._fail(period_key, report_type, "market_data_unavailable", actor)
if not ai.get("ok"):
self._notify("AI 当前不可用,本次市场分析报告未发送。", actor)
return self._fail(period_key, report_type, "ai_unavailable", actor)
return self._deliver(report_type, period_key, report, force, actor, sync_result)
except Exception as exc:
self.db.rollback()
self._step(
period_key,
report_type,
"failed",
WorkflowStatus.FAILED,
actor,
{"error_type": type(exc).__name__, "error": str(exc)[:1000]},
)
self._notify(f"市场分析 {period_key} 执行失败:{type(exc).__name__}", actor)
raise
def _sync(self, report_type: str, target: date) -> dict[str, Any]:
result: dict[str, Any] = {}
if report_type == "premarket" and not self.market.is_trading_day(target):
return {"market_closed": True}
if report_type == "close":
result["daily"] = self.market.sync_daily(target)
if result["daily"].get("market_closed"):
result["market_closed"] = True
return result
result["macro"] = self.market.sync_macro(target)
start = target - timedelta(days=6 if report_type == "weekly" else 1)
result["announcements"] = self.market.sync_announcements(start, target)
result["financials"] = self.market.sync_watchlist_financials()
return result
def _deliver(
self,
report_type: str,
period_key: str,
report: dict[str, Any],
force: bool,
actor: str,
sync_result: dict[str, Any],
) -> dict[str, Any]:
settings = get_settings()
if not (
settings.feishu_default_chat_id
and settings.feishu_app_id
and settings.feishu_app_secret
):
return self._fail(period_key, report_type, "delivery_not_configured", actor)
idempotency_key = (
period_key if not force else f"{period_key}:force:{utc_now():%Y%m%d%H%M%S%f}"
)
reports = ReportService(self.db)
push_run = reports.create_push_run(
report_type=f"market_{report_type}",
title=report["title"],
receive_id=settings.feishu_default_chat_id,
receive_id_type="chat_id",
actor=actor,
idempotency_key=idempotency_key,
)
if push_run.status != ReportPushStatus.SUCCESS:
try:
feishu = FeishuService(self.db)
image = feishu.upload_image(render_market_chart(report), actor)
image_key = (image.get("data") or {}).get("image_key")
if not image_key:
raise ValueError("Feishu image upload did not return image_key")
card = FeishuService.build_basic_card(
report["title"],
report["lines"],
image_key=image_key,
image_alt=report["title"],
)
response = feishu.send_card(
card,
settings.feishu_default_chat_id,
receive_id_type="chat_id",
actor=actor,
)
reports.update_push_run(
push_run.code,
ReportPushStatus.SUCCESS,
provider_response=response,
sent=True,
)
except Exception as exc:
reports.update_push_run(
push_run.code, ReportPushStatus.FAILED, error_message=str(exc)[:2000]
)
raise
workflow = self._step(
period_key,
report_type,
"pushed",
WorkflowStatus.COMPLETED,
actor,
{"push_run_code": push_run.code, "sync": sync_result},
)
return {
"workflow_code": workflow.code,
"period_key": period_key,
"push_run_code": push_run.code,
"status": workflow.status,
"deduplicated": False,
}
def _fail(self, period_key: str, report_type: str, action: str, actor: str) -> dict[str, Any]:
workflow = self._step(
period_key, report_type, action, WorkflowStatus.FAILED, actor
)
return {
"workflow_code": workflow.code,
"period_key": period_key,
"status": workflow.status,
"reason": action,
"deduplicated": False,
}
def _step(
self,
period_key: str,
report_type: str,
action: str,
status_value: str,
actor: str,
payload: dict[str, Any] | None = None,
) -> WorkflowInstance:
return self.workflows.start_or_update(
workflow_type=WorkflowType.MARKET_ANALYSIS,
aggregate_type="market_period",
aggregate_id=period_key,
status_value=status_value,
action=action,
actor=actor,
payload={"report_type": report_type, **(payload or {})},
)
@staticmethod
def _validate_type(report_type: str) -> None:
if report_type not in MARKET_REPORT_TYPES:
raise ValueError("Market report type must be premarket, close or weekly")
def _notify(self, content: str, actor: str) -> None:
settings = get_settings()
if not (
settings.feishu_default_chat_id
and settings.feishu_app_id
and settings.feishu_app_secret
):
return
try:
FeishuService(self.db).send_text(
content,
settings.feishu_default_chat_id,
receive_id_type="chat_id",
actor=actor,
)
except Exception:
self.db.rollback()

View File

@@ -0,0 +1,3 @@
from app.application.scheduling.scheduler import attach_scheduler, create_scheduler
__all__ = ["attach_scheduler", "create_scheduler"]

View File

@@ -0,0 +1,345 @@
from collections.abc import Callable
from datetime import date
from socket import gethostname
from typing import Any
from fastapi import FastAPI
from app.core.config import get_settings
from app.core.constants import ActorValue
from app.modules.feishu.constants import FeishuReceiveIdType
from app.application.delivery import ReportDeliveryService
from app.modules.observability.constants import HeartbeatComponent
def attach_scheduler(app: FastAPI) -> None:
"""Attach optional APScheduler jobs to the FastAPI application."""
settings = get_settings()
if not settings.scheduler_enabled:
return
scheduler = create_scheduler(app)
@app.on_event("startup")
def start_scheduler() -> None:
scheduler.start()
@app.on_event("shutdown")
def stop_scheduler() -> None:
scheduler.shutdown(wait=False)
def create_scheduler(app: FastAPI | None = None) -> Any:
"""Create the V2/V3 scheduler without requiring a FastAPI process."""
from apscheduler.schedulers.background import BackgroundScheduler
from app.core.database import SessionLocal
from app.core.background.task_queue import (
enqueue_attendance_summary_push,
enqueue_daily_brief_push,
enqueue_event_dispatch,
enqueue_legacy_project_sync,
enqueue_legacy_task_sync,
enqueue_lifecycle_report,
enqueue_market_report,
enqueue_project_weekly_push,
enqueue_risk_progress_push,
enqueue_work_daily_push,
enqueue_work_weekly_push,
)
from app.modules.observability.service import ObservabilityService
from app.modules.reports.services import ReportService
settings = get_settings()
scheduler = BackgroundScheduler(timezone="Asia/Shanghai")
def deliver_report(
state_key: str,
build_report: Callable[[Any], dict[str, Any]],
enqueue_push: Callable[..., dict[str, Any]],
) -> None:
db = SessionLocal()
try:
service = ReportService(db)
report = build_report(service)
_set_state(app, state_key, report)
if not _feishu_delivery_configured(settings):
return
if settings.task_queue_enabled:
dispatch = enqueue_push(
receive_id=settings.feishu_default_chat_id,
receive_id_type=FeishuReceiveIdType.CHAT_ID,
actor=ActorValue.SCHEDULER,
)
_set_state(app, f"{state_key}_dispatch", dispatch)
return
ReportDeliveryService(db).push_report(
report,
receive_id=settings.feishu_default_chat_id,
receive_id_type=FeishuReceiveIdType.CHAT_ID,
actor=ActorValue.SCHEDULER,
)
finally:
db.close()
def run_daily_brief() -> None:
deliver_report(
"last_daily_brief",
lambda service: service.daily_brief(),
enqueue_daily_brief_push,
)
def run_project_weekly() -> None:
deliver_report(
"last_project_weekly",
lambda service: service.project_weekly(),
enqueue_project_weekly_push,
)
def run_attendance_summary() -> None:
deliver_report(
"last_attendance_summary",
lambda service: service.attendance_summary(),
enqueue_attendance_summary_push,
)
def run_risk_progress() -> None:
deliver_report(
"last_risk_progress",
lambda service: service.risk_progress(),
enqueue_risk_progress_push,
)
def run_work_daily() -> None:
deliver_report(
"last_work_daily",
lambda service: service.work_daily_report(
reporter=ActorValue.SCHEDULER,
actor=ActorValue.SCHEDULER,
),
enqueue_work_daily_push,
)
def run_work_weekly() -> None:
deliver_report(
"last_work_weekly",
lambda service: service.work_weekly_report(
reporter=ActorValue.SCHEDULER,
actor=ActorValue.SCHEDULER,
),
enqueue_work_weekly_push,
)
def run_legacy_project_sync() -> None:
dispatch = enqueue_legacy_project_sync(actor=ActorValue.SCHEDULER)
_set_state(app, "last_legacy_project_sync_dispatch", dispatch)
def run_legacy_task_sync() -> None:
dispatch = enqueue_legacy_task_sync(actor=ActorValue.SCHEDULER)
_set_state(app, "last_legacy_task_sync_dispatch", dispatch)
def run_daily_lifecycle() -> None:
dispatch = enqueue_lifecycle_report(
report_type="daily",
receive_id=settings.feishu_default_chat_id,
receive_id_type=FeishuReceiveIdType.CHAT_ID,
actor=ActorValue.SCHEDULER,
)
_set_state(app, "last_daily_lifecycle_dispatch", dispatch)
def run_weekly_lifecycle() -> None:
dispatch = enqueue_lifecycle_report(
report_type="weekly",
receive_id=settings.feishu_default_chat_id,
receive_id_type=FeishuReceiveIdType.CHAT_ID,
actor=ActorValue.SCHEDULER,
)
_set_state(app, "last_weekly_lifecycle_dispatch", dispatch)
def run_market_premarket() -> None:
dispatch = enqueue_market_report("premarket")
_set_state(app, "last_market_premarket_dispatch", dispatch)
def run_market_close() -> None:
dispatch = enqueue_market_report("close")
_set_state(app, "last_market_close_dispatch", dispatch)
def run_market_weekly() -> None:
dispatch = enqueue_market_report("weekly", date.today())
_set_state(app, "last_market_weekly_dispatch", dispatch)
def run_event_dispatch() -> None:
dispatch = enqueue_event_dispatch(
limit=settings.event_dispatch_batch_size,
actor=ActorValue.SCHEDULER,
)
_set_state(app, "last_event_dispatch", dispatch)
def record_scheduler_heartbeat() -> None:
db = SessionLocal()
try:
heartbeat = ObservabilityService(db).record_heartbeat(
component=HeartbeatComponent.SCHEDULER,
instance_id=gethostname(),
actor=ActorValue.SCHEDULER,
)
_set_state(app, "last_scheduler_heartbeat", heartbeat)
finally:
db.close()
if settings.lifecycle_pipeline_enabled and not settings.read_only_mode:
scheduler.add_job(
run_daily_lifecycle,
trigger="cron",
hour=settings.daily_brief_cron_hour,
minute=settings.daily_brief_cron_minute,
id="daily_lifecycle_pipeline",
replace_existing=True,
)
scheduler.add_job(
run_weekly_lifecycle,
trigger="cron",
day_of_week=settings.weekly_project_report_day_of_week,
hour=settings.weekly_project_report_cron_hour,
minute=settings.weekly_project_report_cron_minute,
id="weekly_lifecycle_pipeline",
replace_existing=True,
)
else:
scheduler.add_job(
run_daily_brief,
trigger="cron",
hour=settings.daily_brief_cron_hour,
minute=settings.daily_brief_cron_minute,
id="daily_brief_push",
replace_existing=True,
)
scheduler.add_job(
run_project_weekly,
trigger="cron",
day_of_week=settings.weekly_project_report_day_of_week,
hour=settings.weekly_project_report_cron_hour,
minute=settings.weekly_project_report_cron_minute,
id="project_weekly_push",
replace_existing=True,
)
scheduler.add_job(
run_attendance_summary,
trigger="cron",
hour=settings.attendance_summary_cron_hour,
minute=settings.attendance_summary_cron_minute,
id="attendance_summary_push",
replace_existing=True,
)
scheduler.add_job(
run_risk_progress,
trigger="cron",
hour=settings.risk_progress_cron_hour,
minute=settings.risk_progress_cron_minute,
id="risk_progress_push",
replace_existing=True,
)
scheduler.add_job(
run_work_daily,
trigger="cron",
hour=settings.work_daily_cron_hour,
minute=settings.work_daily_cron_minute,
id="work_daily_push",
replace_existing=True,
)
scheduler.add_job(
run_work_weekly,
trigger="cron",
day_of_week=settings.work_weekly_report_day_of_week,
hour=settings.work_weekly_cron_hour,
minute=settings.work_weekly_cron_minute,
id="work_weekly_push",
replace_existing=True,
)
if settings.event_dispatch_enabled:
scheduler.add_job(
run_event_dispatch,
trigger="cron",
minute=settings.event_dispatch_cron_minute,
id="event_dispatch",
replace_existing=True,
)
if settings.market_analysis_enabled and not settings.read_only_mode:
scheduler.add_job(
run_market_premarket,
trigger="cron",
day_of_week="mon-fri",
hour=settings.market_premarket_cron_hour,
minute=settings.market_premarket_cron_minute,
id="market_premarket_analysis",
replace_existing=True,
)
scheduler.add_job(
run_market_close,
trigger="cron",
day_of_week="mon-fri",
hour=settings.market_close_cron_hour,
minute=settings.market_close_cron_minute,
id="market_close_analysis",
replace_existing=True,
)
scheduler.add_job(
run_market_weekly,
trigger="cron",
day_of_week=settings.market_weekly_day_of_week,
hour=settings.market_weekly_cron_hour,
minute=settings.market_weekly_cron_minute,
id="market_weekly_analysis",
replace_existing=True,
)
scheduler.add_job(
record_scheduler_heartbeat,
trigger="interval",
seconds=settings.heartbeat_interval_seconds,
id="scheduler_heartbeat",
replace_existing=True,
)
if (
not settings.lifecycle_pipeline_enabled
and not settings.read_only_mode
and settings.legacy_sync_enabled
and settings.legacy_project_query
):
scheduler.add_job(
run_legacy_project_sync,
trigger="cron",
hour=settings.legacy_project_sync_cron_hour,
minute=settings.legacy_project_sync_cron_minute,
id="legacy_project_sync",
replace_existing=True,
)
if (
not settings.lifecycle_pipeline_enabled
and not settings.read_only_mode
and settings.legacy_sync_enabled
and settings.legacy_task_query
):
scheduler.add_job(
run_legacy_task_sync,
trigger="cron",
hour=settings.legacy_task_sync_cron_hour,
minute=settings.legacy_task_sync_cron_minute,
id="legacy_task_sync",
replace_existing=True,
)
return scheduler
def _set_state(app: FastAPI | None, key: str, value: Any) -> None:
if app is not None:
setattr(app.state, key, value)
def _feishu_delivery_configured(settings: Any) -> bool:
return bool(
settings.feishu_app_id
and settings.feishu_app_secret
and settings.feishu_default_chat_id
)