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.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 get_settings().read_only_mode: 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()