|
20 | 20 | import json |
21 | 21 | import logging |
22 | 22 | import re |
| 23 | +import uuid |
23 | 24 | from datetime import datetime |
24 | 25 | from pathlib import Path |
25 | 26 | from typing import Optional, Union, Dict, Any |
@@ -116,20 +117,24 @@ def _run_market_review_background( |
116 | 117 | override_region: Optional[str] = None, |
117 | 118 | lock_token: Optional[_MarketReviewExecutionLock] = None, |
118 | 119 | config: Optional[Config] = None, |
| 120 | + query_id: Optional[str] = None, |
119 | 121 | ) -> None: |
120 | 122 | """Run market review after the API response has been accepted.""" |
121 | 123 | from src.core.market_review import run_market_review |
122 | 124 |
|
123 | 125 | runtime_config = config or get_config_dep() |
124 | 126 | try: |
125 | 127 | notifier, analyzer, search_service = _build_market_review_runtime(runtime_config) |
126 | | - report = run_market_review( |
127 | | - notifier=notifier, |
128 | | - analyzer=analyzer, |
129 | | - search_service=search_service, |
130 | | - send_notification=send_notification, |
131 | | - override_region=override_region, |
132 | | - ) |
| 128 | + review_kwargs = { |
| 129 | + "notifier": notifier, |
| 130 | + "analyzer": analyzer, |
| 131 | + "search_service": search_service, |
| 132 | + "send_notification": send_notification, |
| 133 | + "override_region": override_region, |
| 134 | + } |
| 135 | + if query_id: |
| 136 | + review_kwargs["query_id"] = query_id |
| 137 | + report = run_market_review(**review_kwargs) |
133 | 138 | if not report: |
134 | 139 | raise RuntimeError("大盘复盘未返回可持久化报告") |
135 | 140 | return {"result": report} |
@@ -500,16 +505,19 @@ def trigger_market_review( |
500 | 505 | ) |
501 | 506 |
|
502 | 507 | try: |
| 508 | + task_id = uuid.uuid4().hex |
503 | 509 | task = get_task_queue().submit_background_task( |
504 | 510 | lambda: _run_market_review_background( |
505 | 511 | request.send_notification, |
506 | 512 | override_region=override_region, |
507 | 513 | lock_token=lock_token, |
508 | 514 | config=config, |
| 515 | + query_id=task_id, |
509 | 516 | ), |
510 | 517 | stock_code="market_review", |
511 | 518 | stock_name="大盘复盘", |
512 | 519 | message="大盘复盘任务已提交", |
| 520 | + task_id=task_id, |
513 | 521 | ) |
514 | 522 | except Exception: |
515 | 523 | _release_market_review_lock(lock_token) |
@@ -751,6 +759,25 @@ def get_analysis_status(task_id: str) -> TaskStatus: |
751 | 759 | if records: |
752 | 760 | record = records[0] |
753 | 761 | raw_result = parse_json_field(record.raw_result) |
| 762 | + if getattr(record, "report_type", None) == "market_review": |
| 763 | + market_review_report = None |
| 764 | + if isinstance(raw_result, dict): |
| 765 | + report_text = raw_result.get("raw_response") or raw_result.get("market_review_report") |
| 766 | + if isinstance(report_text, str) and report_text.strip(): |
| 767 | + market_review_report = report_text |
| 768 | + if not market_review_report and record.news_content: |
| 769 | + market_review_report = record.news_content |
| 770 | + |
| 771 | + return TaskStatus( |
| 772 | + task_id=task_id, |
| 773 | + status="completed", |
| 774 | + progress=100, |
| 775 | + result=None, |
| 776 | + market_review_report=market_review_report, |
| 777 | + error=None, |
| 778 | + stock_name=record.name, |
| 779 | + ) |
| 780 | + |
754 | 781 | model_used = normalize_model_used( |
755 | 782 | (raw_result or {}).get("model_used") if isinstance(raw_result, dict) else None |
756 | 783 | ) |
|
0 commit comments