diff --git a/api/v1/endpoints/analysis.py b/api/v1/endpoints/analysis.py index ffe797aa6..9a53c5c32 100644 --- a/api/v1/endpoints/analysis.py +++ b/api/v1/endpoints/analysis.py @@ -66,7 +66,7 @@ from src.analysis_context_pack_overview import ( extract_analysis_context_pack_overview, sanitize_context_snapshot_for_api, ) -from src.market_phase_summary import extract_market_phase_summary +from src.market_phase_summary import extract_market_phase_summary, render_market_phase_summary from src.report_language import get_localized_stock_name, normalize_report_language from src.services.name_to_code_resolver import resolve_name_to_code from src.services.stock_code_utils import is_code_like @@ -342,6 +342,7 @@ def _handle_async_analysis_batch( selection_source = request.selection_source if (is_single or preserve_batch_metadata) else None notify = getattr(request, "notify", True) skills = getattr(request, "skills", None) + analysis_phase = request.analysis_phase submit_kwargs = dict( stock_codes=stock_codes, @@ -349,6 +350,7 @@ def _handle_async_analysis_batch( original_query=original_query, selection_source=selection_source, report_type=request.report_type, + analysis_phase=analysis_phase, force_refresh=request.force_refresh, notify=notify, ) @@ -364,6 +366,7 @@ def _handle_async_analysis_batch( stock_code=task.stock_code, status="pending", message=f"分析任务已加入队列: {task.stock_code}", + analysis_phase=task.analysis_phase, ) for task in accepted_tasks ] @@ -397,6 +400,7 @@ def _handle_async_analysis_batch( trace_id=accepted[0].trace_id, status="pending", message=accepted[0].message, + analysis_phase=accepted[0].analysis_phase, ) return JSONResponse( status_code=202, @@ -438,6 +442,7 @@ def _handle_sync_analysis( query_id=query_id, send_notification=getattr(request, "notify", True), skills=getattr(request, "skills", None), + analysis_phase=request.analysis_phase, ) if result is None: @@ -618,6 +623,8 @@ def get_task_list( error=t.error, original_query=t.original_query, selection_source=t.selection_source, + analysis_phase=t.analysis_phase, + skills=getattr(t, "skills", None), ) for t in all_tasks ] @@ -868,6 +875,7 @@ def get_analysis_status(task_id: str) -> TaskStatus: stock_name=task.stock_name, original_query=task.original_query, selection_source=task.selection_source, + analysis_phase=task.analysis_phase, skills=getattr(task, "skills", None), ) @@ -1107,6 +1115,10 @@ def _build_analysis_report( if change_pct is None: change_pct = realtime_fields.get("change_pct") market_phase_summary = extract_market_phase_summary(context_snapshot) + if market_phase_summary is None: + meta_phase_summary = meta_data.get("market_phase_summary") + if meta_phase_summary is not None: + market_phase_summary = render_market_phase_summary(meta_phase_summary) meta = ReportMeta( query_id=meta_data.get("query_id", query_id), diff --git a/api/v1/schemas/analysis.py b/api/v1/schemas/analysis.py index b4cd4f77b..ecc3f063c 100644 --- a/api/v1/schemas/analysis.py +++ b/api/v1/schemas/analysis.py @@ -10,7 +10,7 @@ 3. 定义异步任务队列相关模型 """ -from typing import Optional, List, Any +from typing import Optional, List, Any, Literal from enum import Enum from pydantic import AliasChoices, BaseModel, ConfigDict, Field @@ -25,6 +25,9 @@ class TaskStatusEnum(str, Enum): FAILED = "failed" +AnalysisPhase = Literal["auto", "premarket", "intraday", "postmarket"] + + class AnalyzeRequest(BaseModel): """Analysis request parameters""" @@ -51,6 +54,10 @@ class AnalyzeRequest(BaseModel): False, description="是否使用异步模式" ) + analysis_phase: AnalysisPhase = Field( + "auto", + description="分析阶段覆盖:auto(自动推断) / premarket(盘前) / intraday(盘中) / postmarket(盘后)", + ) stock_name: Optional[str] = Field( None, description="用户选中的股票名称(自动补全时提供)", @@ -84,6 +91,7 @@ class AnalyzeRequest(BaseModel): "report_type": "detailed", "force_refresh": False, "async_mode": False, + "analysis_phase": "auto", "stock_name": "贵州茅台", "original_query": "茅台", "selection_source": "autocomplete", @@ -156,12 +164,14 @@ class TaskAccepted(BaseModel): pattern="^(pending|processing)$" ) message: Optional[str] = Field(None, description="提示信息") + analysis_phase: AnalysisPhase = Field("auto", description="请求的分析阶段") model_config = ConfigDict(json_schema_extra={ "example": { "task_id": "task_abc123", "status": "pending", - "message": "Analysis task accepted" + "message": "Analysis task accepted", + "analysis_phase": "auto" } }) @@ -178,13 +188,15 @@ class BatchTaskAcceptedItem(BaseModel): pattern="^(pending|processing)$" ) message: Optional[str] = Field(None, description="提示信息") + analysis_phase: AnalysisPhase = Field("auto", description="请求的分析阶段") model_config = ConfigDict(json_schema_extra={ "example": { "task_id": "task_abc123", "stock_code": "600519", "status": "pending", - "message": "分析任务已加入队列: 600519" + "message": "分析任务已加入队列: 600519", + "analysis_phase": "auto" } }) @@ -219,7 +231,8 @@ class BatchTaskAcceptedResponse(BaseModel): "task_id": "task_abc123", "stock_code": "600519", "status": "pending", - "message": "分析任务已加入队列: 600519" + "message": "分析任务已加入队列: 600519", + "analysis_phase": "auto" } ], "duplicates": [ @@ -269,6 +282,10 @@ class TaskStatus(BaseModel): description="选择来源", pattern=SELECTION_SOURCE_PATTERN, ) + analysis_phase: Optional[AnalysisPhase] = Field( + None, + description="请求的分析阶段;无持久化字段的历史 DB fallback 可能为空", + ) skills: Optional[List[str]] = Field(None, description="本次任务使用的策略 skill ID 列表") model_config = ConfigDict(json_schema_extra={ @@ -282,6 +299,7 @@ class TaskStatus(BaseModel): "stock_name": "贵州茅台", "original_query": "茅台", "selection_source": "autocomplete", + "analysis_phase": "auto", "skills": ["bull_trend"] } }) @@ -312,6 +330,7 @@ class TaskInfo(BaseModel): description="选择来源", pattern=SELECTION_SOURCE_PATTERN, ) + analysis_phase: AnalysisPhase = Field("auto", description="请求的分析阶段") skills: Optional[List[str]] = Field(None, description="本次任务使用的策略 skill ID 列表") model_config = ConfigDict(json_schema_extra={ @@ -329,6 +348,7 @@ class TaskInfo(BaseModel): "error": None, "original_query": "茅台", "selection_source": "autocomplete", + "analysis_phase": "auto", "skills": ["bull_trend"] } }) diff --git a/apps/dsa-web/src/api/analysis.ts b/apps/dsa-web/src/api/analysis.ts index f3b069901..40eadce42 100644 --- a/apps/dsa-web/src/api/analysis.ts +++ b/apps/dsa-web/src/api/analysis.ts @@ -27,6 +27,7 @@ export const analysisApi = { report_type: data.reportType || 'detailed', force_refresh: data.forceRefresh || false, async_mode: data.asyncMode || false, + analysis_phase: data.analysisPhase || 'auto', stock_name: data.stockName, original_query: data.originalQuery, selection_source: data.selectionSource, @@ -61,6 +62,7 @@ export const analysisApi = { report_type: data.reportType || 'detailed', force_refresh: data.forceRefresh || false, async_mode: true, + analysis_phase: data.analysisPhase || 'auto', stock_name: data.stockName, original_query: data.originalQuery, selection_source: data.selectionSource, diff --git a/apps/dsa-web/src/hooks/__tests__/useTaskStream.test.tsx b/apps/dsa-web/src/hooks/__tests__/useTaskStream.test.tsx index 127558ec4..25fc778fe 100644 --- a/apps/dsa-web/src/hooks/__tests__/useTaskStream.test.tsx +++ b/apps/dsa-web/src/hooks/__tests__/useTaskStream.test.tsx @@ -77,7 +77,9 @@ describe('useTaskStream', () => { progress: 72, message: 'LLM 正在生成分析结果', report_type: 'detailed', + analysis_phase: 'intraday', created_at: '2026-03-29T08:00:00Z', + skills: ['growth_quality'], }), }), ); @@ -97,6 +99,8 @@ describe('useTaskStream', () => { error: undefined, originalQuery: undefined, selectionSource: undefined, + analysisPhase: 'intraday', + skills: ['growth_quality'], }); }); }); diff --git a/apps/dsa-web/src/hooks/useTaskStream.ts b/apps/dsa-web/src/hooks/useTaskStream.ts index 7d52e0c15..3b0debaa0 100644 --- a/apps/dsa-web/src/hooks/useTaskStream.ts +++ b/apps/dsa-web/src/hooks/useTaskStream.ts @@ -123,6 +123,8 @@ export function useTaskStream(options: UseTaskStreamOptions = {}): UseTaskStream error: data.error as string | undefined, originalQuery: data.original_query as string | undefined, selectionSource: data.selection_source as string | undefined, + analysisPhase: data.analysis_phase as TaskInfo['analysisPhase'], + skills: Array.isArray(data.skills) ? data.skills.map(String) : undefined, }; if (typeof data.trace_id === 'string' && data.trace_id.trim()) { diff --git a/apps/dsa-web/src/types/analysis.ts b/apps/dsa-web/src/types/analysis.ts index 922e975a0..936b53e08 100644 --- a/apps/dsa-web/src/types/analysis.ts +++ b/apps/dsa-web/src/types/analysis.ts @@ -7,6 +7,7 @@ export type StockReportType = 'simple' | 'detailed' | 'full' | 'brief'; export type ReportType = StockReportType | 'market_review'; +export type AnalysisPhase = 'auto' | 'premarket' | 'intraday' | 'postmarket'; export interface AnalysisRequest { stockCode?: string; @@ -14,6 +15,7 @@ export interface AnalysisRequest { reportType?: StockReportType; forceRefresh?: boolean; asyncMode?: boolean; + analysisPhase?: AnalysisPhase; stockName?: string; originalQuery?: string; selectionSource?: 'manual' | 'autocomplete' | 'import' | 'image'; @@ -254,6 +256,7 @@ export interface TaskAccepted { traceId?: string; status: 'pending' | 'processing'; message?: string; + analysisPhase?: AnalysisPhase; } export interface BatchTaskAcceptedItem { @@ -262,6 +265,7 @@ export interface BatchTaskAcceptedItem { stockCode: string; status: 'pending' | 'processing'; message?: string; + analysisPhase?: AnalysisPhase; } export interface BatchDuplicateTaskItem { @@ -292,6 +296,7 @@ export interface TaskStatus { stockName?: string; originalQuery?: string; selectionSource?: string; + analysisPhase?: AnalysisPhase | null; skills?: string[]; } @@ -311,6 +316,8 @@ export interface TaskInfo { error?: string; originalQuery?: string; selectionSource?: string; + analysisPhase?: AnalysisPhase; + skills?: string[]; } /** Task list response */ diff --git a/docs/CHANGELOG.md b/docs/CHANGELOG.md index 049decdfd..4cd7f074c 100644 --- a/docs/CHANGELOG.md +++ b/docs/CHANGELOG.md @@ -42,6 +42,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/). - [改进] AnalysisContextPack P5 增加数据质量评分、`fetch_failed` 状态、Prompt 数据限制区块和 Web 低敏质量展示。 - [改进] #1386 P2-full 在 AnalysisContextPack Prompt 数据限制中追加市场阶段与降级数据的交叉约束,并修正中文分析 Prompt 的阶段化行情标签。 - [新功能] #1386 P5 为个股分析报告新增 `dashboard.phase_decision` 盘中决策护栏,并在保存历史前按市场阶段与数据质量限制高置信盘中买卖结论。 +- [新功能] #1386 P4a 新增 `analysis_phase=auto|premarket|intraday|postmarket` API 参数,并在异步任务 accepted、内存 status、list、SSE 与分析 pipeline 中透传请求阶段。 - [文档] 明确同股历史趋势新增模型字段为历史快照展示元数据,不影响运行时 LLM Provider/Model/Base URL 路由与配置迁移清理;回退方式为按常规发布回滚本变更。 - [修复] 收口 Web 中文界面残留英文文案与设置页 help 缺口,回测页改为中文展示,并让 Web 设置页仅展示已注册且带说明的配置项。 diff --git a/docs/architecture/api_spec.json b/docs/architecture/api_spec.json index 100ef4294..b8b4391f8 100644 --- a/docs/architecture/api_spec.json +++ b/docs/architecture/api_spec.json @@ -793,6 +793,17 @@ "type": "boolean", "default": false, "description": "是否使用异步模式" + }, + "analysis_phase": { + "type": "string", + "enum": [ + "auto", + "premarket", + "intraday", + "postmarket" + ], + "default": "auto", + "description": "分析阶段覆盖:auto(自动推断) / premarket(盘前) / intraday(盘中) / postmarket(盘后)" } } }, @@ -842,6 +853,17 @@ "message": { "type": "string", "example": "Analysis task accepted" + }, + "analysis_phase": { + "type": "string", + "enum": [ + "auto", + "premarket", + "intraday", + "postmarket" + ], + "default": "auto", + "description": "请求的分析阶段" } }, "required": [ @@ -869,6 +891,17 @@ }, "message": { "type": "string" + }, + "analysis_phase": { + "type": "string", + "enum": [ + "auto", + "premarket", + "intraday", + "postmarket" + ], + "default": "auto", + "description": "请求的分析阶段" } }, "required": [ @@ -950,6 +983,17 @@ "error": { "type": "string", "description": "错误信息(仅在 failed 时存在)" + }, + "analysis_phase": { + "type": "string", + "enum": [ + "auto", + "premarket", + "intraday", + "postmarket" + ], + "nullable": true, + "description": "请求的分析阶段;无持久化字段的历史 DB fallback 可能为空" } }, "required": [ @@ -1019,6 +1063,17 @@ "error": { "type": "string", "description": "错误信息" + }, + "analysis_phase": { + "type": "string", + "enum": [ + "auto", + "premarket", + "intraday", + "postmarket" + ], + "default": "auto", + "description": "请求的分析阶段" } }, "required": [ diff --git a/docs/full-guide.md b/docs/full-guide.md index 88cf8aede..cbf7a9738 100644 --- a/docs/full-guide.md +++ b/docs/full-guide.md @@ -770,6 +770,16 @@ P3 补齐普通分析主路径使用的实时行情质量元数据,但仍不 整源 fallback 的语义固定为:`source` 保留实际成功的数据源 token,`fallback_from` 记录本轮失败的最高优先级整源 token;首选源成功后只从后续源补字段时不写 `fallback_from`。`AnalysisContextBuilder` 只映射这些上游 artifact,不重新取数、不做质量评分;quote block 状态按 `STALE > FALLBACK > AVAILABLE` 归并。盘中实时价覆盖 `today` 时会标记 `is_partial_bar`、`is_estimated`、`estimated_fields`、`realtime_source` 和 quote 元数据;`daily_bars` block 仍表示 storage 中完整日线窗口,partial/estimated 只进入 technical block。freshness scoring、盘中 cache TTL 分级、Agent 工具级复用和 API/Web 展示留给后续阶段。 +#### 分析阶段入口与任务队列透传(Issue #1386 P4a) + +P4a 新增 `analysis_phase=auto|premarket|intraday|postmarket` 请求参数,默认 `auto`,用于让 API 调用方显式覆盖本次分析阶段。该参数目前接入 `POST /api/v1/analysis/analyze`、异步任务队列、`AnalysisService`、普通分析 pipeline 和市场阶段上下文;Web 前端类型和 API mapper 已承接该字段,但不新增页面 selector,Bot、schedule、GitHub Actions 和 DB migration 也不在本阶段范围内。 + +`analysis_phase` 是请求覆盖值;最终报告阶段仍以 `report.meta.market_phase_summary.phase` 为准。异步 accepted response、内存任务 status、任务列表和 SSE payload 会回显请求阶段;历史 DB fallback 不新增持久化字段,旧记录仍可能为空。同股不同 phase 仍按同一个股票任务去重,避免并发重复分析。 + +内部阶段上下文构造仍兼容旧参数 `analysis_intent`:仅当 `analysis_phase` 保持 `auto` 时,非 `auto` 的 `analysis_intent` 会被归一为本次请求阶段;外部调用方应优先使用 `analysis_phase`。 + +`auto` 保持既有交易日历推断;非 `auto` 只覆盖 phase 并重算 `is_trading_day`、`is_market_open_now`、`is_partial_bar`、`minutes_to_open` 和 `minutes_to_close`。覆盖不会改写真实 `market_local_time` 或 `effective_daily_bar_date`;如果当前日期不是交易日或日历不支持对应 session,分钟字段可以为空。 + #### AnalysisContextPack Prompt 摘要(Issue #1389 P3) P3 在普通分析和 Agent 初始上下文中接入 `AnalysisContextPack` 低敏摘要。Pipeline 会用已获取的行情、日线、趋势、筹码、基本面、新闻和市场阶段 artifacts 组装 pack,再把 `analysis_context_pack_summary` 插入 Prompt;在这个新增的 pack 摘要区块中,LLM 只看到 subject、版本、各数据块的状态/来源/warning/missing reason 和新闻结果数,不会通过该区块看到完整 `news.content`、`trend_result`、筹码或基本面原始 payload。既有 `news_context`、Agent pre-fetched JSON 和 `enhanced_context` 原始数据通道保持 P3 前行为,不由本摘要替代或脱敏。 @@ -1295,6 +1305,7 @@ FastAPI 提供 RESTful API 服务,支持配置管理和触发分析。 > 说明:`POST /api/v1/analysis/analyze` 在 `async_mode=false` 时仅支持单只股票;批量 `stock_codes` 需使用 `async_mode=true`。异步 `202` 响应对单股返回 `task_id`,对批量返回 `accepted` / `duplicates` 汇总结构。 > 说明:`POST /api/v1/analysis/analyze` 支持使用 `skills` 传入策略 skill ID 列表;若未传则按服务端默认策略执行。为兼容历史调用,`strategies` 字段仍作为兼容别名保留。 +> 说明:`POST /api/v1/analysis/analyze` 支持 `analysis_phase=auto|premarket|intraday|postmarket`,默认 `auto`。非 `auto` 只覆盖本次分析阶段与派生阶段标记,不改写真实交易日历时间;accepted response、内存 task status、任务列表和 SSE 会回显请求阶段,最终报告阶段以 `report.meta.market_phase_summary.phase` 为准。 > 说明:Web 侧首页策略下拉为显式可选策略入口。用户未手动选择时不会携带 `skills`,与历史客户端行为一致;选择策略后将透传到该接口并在任务状态与历史快照中保留。 > 说明:`POST /api/v1/analysis/market-review` 采用后端与 CLI/Bot 共用的配置路径(`GeminiAnalyzer(config=...)` 与同样的搜索/提示词构造入口)。Provider 兼容路由会优先识别并使用 `litellm_model`、`llm_model_list`,若未配置则回退 legacy `GEMINI_*`、`OPENAI_*`、`ANTHROPIC_*`、`DEEPSEEK_*` 键;不会新增/调整 provider、Base URL 或 LiteLLM 路由语义。 > 审计依据:优先级与回退语义以 `src/config.py` 的 `Config._load_from_env()` 为准(`LITELLM_CONFIG` > `LLM_CHANNELS` > legacy)。配套回归见 `tests/test_llm_channel_config.py`(配置源解析)与 `tests/test_market_review_runtime.py`(共享装配路径)。该接口当前仅提供单进程/单机级防重复能力,若为多实例部署需通过外部任务队列或分布式锁补齐全局幂等。 diff --git a/docs/full-guide_EN.md b/docs/full-guide_EN.md index 2c0a2e5d5..1c7a853b1 100644 --- a/docs/full-guide_EN.md +++ b/docs/full-guide_EN.md @@ -647,6 +647,16 @@ P3 adds realtime quote quality metadata for the regular analysis path, but still Whole-source fallback semantics are fixed: `source` keeps the actual successful provider token, while `fallback_from` records the highest-priority whole source that failed in the current attempt; if the primary source succeeds and later providers only supplement missing fields, `fallback_from` is not set. `AnalysisContextBuilder` only maps these upstream artifacts, performs no extra fetches, and does no quality scoring; quote block status collapses as `STALE > FALLBACK > AVAILABLE`. When realtime price overlays `today`, the pipeline marks `is_partial_bar`, `is_estimated`, `estimated_fields`, `realtime_source`, and quote metadata. The `daily_bars` block still represents the complete daily-bar window in storage; partial/estimated markers only enter the technical block. Freshness scoring, intraday cache TTL tiers, Agent tool-level reuse, and API/Web display remain follow-ups. +### Analysis Phase Entrypoint and Task Queue Pass-Through (Issue #1386 P4a) + +P4a adds an `analysis_phase=auto|premarket|intraday|postmarket` request parameter, defaulting to `auto`, so API callers can explicitly override the phase for the current analysis. The parameter is wired through `POST /api/v1/analysis/analyze`, the async task queue, `AnalysisService`, the regular analysis pipeline, and market-phase context construction. Web frontend types and API mapping accept the field, but this phase does not add a page selector; Bot, schedule, GitHub Actions, and DB migrations remain out of scope. + +`analysis_phase` is the requested override value; the final report phase remains `report.meta.market_phase_summary.phase`. Async accepted responses, in-memory task status, task list responses, and SSE payloads echo the requested phase. DB history fallback does not add a persisted phase field, so older records may still return it empty. Duplicate detection remains stock-only, so the same stock submitted with different phases is still treated as a duplicate in-flight task. + +Market-phase context construction still supports the legacy internal `analysis_intent` argument: only when `analysis_phase` remains `auto`, a non-`auto` `analysis_intent` is normalized as the requested phase for this run. External callers should prefer `analysis_phase`. + +`auto` preserves existing calendar inference. Non-`auto` values only override the phase and recompute `is_trading_day`, `is_market_open_now`, `is_partial_bar`, `minutes_to_open`, and `minutes_to_close`. The override does not rewrite the real `market_local_time` or `effective_daily_bar_date`; if the current date is not a trading session or the calendar cannot support the session, minute fields may be empty. + ### AnalysisContextPack Prompt Summary (Issue #1389 P3) P3 injects a low-sensitivity `AnalysisContextPack` summary into regular analysis and Agent initial prompts. The pipeline builds the pack from already-fetched quote, daily-bar, trend, chip, fundamentals, news, and market-phase artifacts, then passes `analysis_context_pack_summary` downstream; in this new pack-summary section, the LLM only sees subject, version, data-block status/source/warnings/missing reason, and news result count, not full `news.content`, `trend_result`, chip, or fundamentals raw payloads through that section. Existing `news_context`, Agent pre-fetched JSON, and `enhanced_context` raw-payload channels keep their pre-P3 behavior and are not replaced or sanitized by this summary. @@ -1127,6 +1137,7 @@ FastAPI provides RESTful API service for configuration management and triggering > Note: `POST /api/v1/analysis/analyze` supports only one stock when `async_mode=false`; batch `stock_codes` requires `async_mode=true`. The async `202` response returns a single `task_id` for one stock, or an `accepted` / `duplicates` summary for batch requests. > Note: `POST /api/v1/analysis/analyze` accepts `skills` as an array of strategy IDs; if omitted, server defaults are used. The legacy field `strategies` is still accepted for backward compatibility. +> Note: `POST /api/v1/analysis/analyze` accepts `analysis_phase=auto|premarket|intraday|postmarket`, defaulting to `auto`. Non-`auto` only overrides the phase and derived phase flags for this run; it does not rewrite real trading-calendar timestamps. Accepted responses, in-memory task status, task lists, and SSE echo the requested phase, while the final report phase remains `report.meta.market_phase_summary.phase`. > Note: The Web Home page exposes an explicit strategy selector. When users do not pick one, `skills` is not sent and legacy behavior is preserved; when selected, it is passed through to this endpoint and persisted in task status/history snapshots. > Note: `POST /api/v1/analysis/market-review` follows the same runtime configuration path as CLI/Bot market review (`GeminiAnalyzer(config=...)`, search setup, and prompt/rendering pipeline). The provider compatibility path prioritizes `litellm_model` and `llm_model_list`, then falls back to existing legacy keys (`GEMINI_*`, `OPENAI_*`, `ANTHROPIC_*`, `DEEPSEEK_*`) when those are not set; provider names, Base URL, and LiteLLM routing semantics are otherwise unchanged. > Audit note: priority and fallback are defined by `Config._load_from_env()` in `src/config.py` (`LITELLM_CONFIG` > `LLM_CHANNELS` > legacy). Regression coverage is in `tests/test_llm_channel_config.py` (configuration source parsing) and `tests/test_market_review_runtime.py` (shared runtime assembly). The endpoint lock is process/host-level only; multi-instance deployments still need external distributed idempotency controls. diff --git a/src/core/pipeline.py b/src/core/pipeline.py index dd01d1d0c..f76d05d8a 100644 --- a/src/core/pipeline.py +++ b/src/core/pipeline.py @@ -103,6 +103,7 @@ class StockAnalysisPipeline: save_context_snapshot: Optional[bool] = None, progress_callback: Optional[Callable[[int, str], None]] = None, analysis_skills: Optional[List[str]] = None, + analysis_phase: str = "auto", ): """ 初始化调度器 @@ -122,6 +123,7 @@ class StockAnalysisPipeline: ) self.progress_callback = progress_callback self.analysis_skills = list(analysis_skills) if analysis_skills is not None else None + self.analysis_phase = analysis_phase or "auto" # 初始化各模块 self.db = get_db() @@ -296,7 +298,7 @@ class StockAnalysisPipeline: market=market, current_time=current_time, trigger_source=self.query_source, - analysis_intent="auto", + analysis_phase=getattr(self, "analysis_phase", "auto"), ) market_phase_context_dict = market_phase_context.to_dict() market_phase_summary = render_market_phase_summary(market_phase_context_dict) diff --git a/src/core/trading_calendar.py b/src/core/trading_calendar.py index e258ec6be..b08d33b16 100644 --- a/src/core/trading_calendar.py +++ b/src/core/trading_calendar.py @@ -49,6 +49,12 @@ MARKET_TIMEZONE = { # regular-session inference layer; it does not change existing fail-open # trading-day filtering or effective-date behavior. _CLOSING_AUCTION_WINDOW_MINUTES = {"cn": 3, "hk": 10, "us": 5} +_SUPPORTED_ANALYSIS_PHASES = { + "auto", + "premarket", + "intraday", + "postmarket", +} class MarketPhase(str, Enum): @@ -412,20 +418,45 @@ def _phase_minutes( return None, None, False +def _normalize_analysis_phase( + analysis_phase: Optional[str], + analysis_intent: Optional[str], +) -> str: + def _coerce(value: Optional[str]) -> str: + if isinstance(value, MarketPhase): + return value.value + return str(value or "").strip().lower() + + requested = _coerce(analysis_phase) or "auto" + legacy_intent = _coerce(analysis_intent) + if requested == "auto" and legacy_intent and legacy_intent != "auto": + requested = legacy_intent + if requested not in _SUPPORTED_ANALYSIS_PHASES: + raise ValueError( + f"invalid analysis_phase: {requested}. " + f"Must be one of {sorted(_SUPPORTED_ANALYSIS_PHASES)}" + ) + return requested + + def build_market_phase_context( *, market: Optional[str], current_time: Optional[datetime] = None, trigger_source: str = "system", analysis_intent: str = "auto", + analysis_phase: str = "auto", ) -> MarketPhaseContext: """ Build a JSON-safe runtime market-phase context for analysis plumbing. - This helper does not change prompt wording, API schema, history metadata, - or task status behavior. Calendar failures degrade to ``unknown`` with - stable warning codes so later P1b/P2 work can consume the same contract. + ``analysis_phase="auto"`` keeps calendar inference. Explicit supported + phases override only the phase and derived flags/minute fields; they do + not rewrite market-local time or the effective daily-bar date. The legacy + ``analysis_intent`` argument remains a compatibility alias when + ``analysis_phase`` is left as ``auto``. """ + requested_phase = _normalize_analysis_phase(analysis_phase, analysis_intent) market_now = get_market_now(market, current_time=current_time) warnings: List[str] = [] @@ -435,9 +466,15 @@ def build_market_phase_context( else: if not _XCALS_AVAILABLE: _add_warning_code(warnings, "calendar_unavailable") - phase = infer_market_phase(market, current_time=current_time) - if phase == MarketPhase.UNKNOWN and _XCALS_AVAILABLE: - _add_warning_code(warnings, "calendar_error") + if requested_phase == "auto": + phase = infer_market_phase(market, current_time=current_time) + if phase == MarketPhase.UNKNOWN and _XCALS_AVAILABLE: + _add_warning_code(warnings, "calendar_error") + else: + phase = MarketPhase(requested_phase) + + if requested_phase != "auto" and phase == MarketPhase.UNKNOWN: + phase = MarketPhase(requested_phase) effective_daily_bar_date = get_effective_trading_date( market, @@ -464,7 +501,7 @@ def build_market_phase_context( minutes_to_open=minutes_to_open, minutes_to_close=minutes_to_close, trigger_source=trigger_source or "system", - analysis_intent=analysis_intent or "auto", + analysis_intent=requested_phase, warnings=warnings, ) diff --git a/src/services/analysis_service.py b/src/services/analysis_service.py index 7fdd1e649..e8f3ded9d 100644 --- a/src/services/analysis_service.py +++ b/src/services/analysis_service.py @@ -22,6 +22,7 @@ from src.report_language import ( localize_trend_prediction, normalize_report_language, ) +from src.market_phase_summary import extract_market_phase_summary from src.services.run_diagnostics import ( activate_run_diagnostic_context, build_run_diagnostic_summary, @@ -54,6 +55,7 @@ class AnalysisService: send_notification: bool = True, progress_callback: Optional[Callable[[int, str], None]] = None, skills: Optional[List[str]] = None, + analysis_phase: str = "auto", ) -> Optional[Dict[str, Any]]: """ 执行股票分析 @@ -64,6 +66,7 @@ class AnalysisService: force_refresh: 是否强制刷新 query_id: 查询 ID(可选) send_notification: 是否发送通知(API 触发默认发送) + analysis_phase: 请求的分析阶段覆盖(auto/premarket/intraday/postmarket) Returns: 分析结果字典,包含: @@ -102,6 +105,7 @@ class AnalysisService: query_source="api", progress_callback=progress_callback, analysis_skills=skills, + analysis_phase=analysis_phase, ) # 确定报告类型 (API: simple/detailed/full/brief -> ReportType) @@ -165,6 +169,7 @@ class AnalysisService: trace_id = diagnostic_context.trace_id if diagnostic_context is not None else query_id diagnostic_snapshot = diagnostic_context.snapshot() if diagnostic_context is not None else None diagnostic_context_snapshot = getattr(result, "diagnostic_context_snapshot", None) + market_phase_summary = extract_market_phase_summary(diagnostic_context_snapshot) if isinstance(diagnostic_context_snapshot, dict): context_snapshot = dict(diagnostic_context_snapshot) if diagnostic_snapshot is not None: @@ -193,6 +198,7 @@ class AnalysisService: "current_price": result.current_price, "change_pct": result.change_pct, "model_used": getattr(result, "model_used", None), + "market_phase_summary": market_phase_summary, }, "summary": { "analysis_summary": result.analysis_summary, diff --git a/src/services/task_queue.py b/src/services/task_queue.py index ebc38a323..41b8c1b22 100644 --- a/src/services/task_queue.py +++ b/src/services/task_queue.py @@ -71,6 +71,7 @@ class TaskInfo: result: Optional[Dict[str, Any]] = None error: Optional[str] = None report_type: str = "detailed" + analysis_phase: str = "auto" created_at: datetime = field(default_factory=datetime.now) started_at: Optional[datetime] = None completed_at: Optional[datetime] = None @@ -90,6 +91,7 @@ class TaskInfo: "progress": self.progress, "message": self.message, "report_type": self.report_type, + "analysis_phase": self.analysis_phase, "created_at": self.created_at.isoformat(), "started_at": self.started_at.isoformat() if self.started_at else None, "completed_at": self.completed_at.isoformat() if self.completed_at else None, @@ -111,6 +113,7 @@ class TaskInfo: result=self.result, error=self.error, report_type=self.report_type, + analysis_phase=self.analysis_phase, created_at=self.created_at, started_at=self.started_at, completed_at=self.completed_at, @@ -310,6 +313,7 @@ class AnalysisTaskQueue: original_query: Optional[str] = None, selection_source: Optional[str] = None, report_type: str = "detailed", + analysis_phase: str = "auto", force_refresh: bool = False, skills: Optional[List[str]] = None, ) -> TaskInfo: @@ -322,6 +326,7 @@ class AnalysisTaskQueue: original_query: Optional raw user input selection_source: Optional source label report_type: Report type + analysis_phase: Requested analysis phase override force_refresh: Whether to bypass cache Returns: @@ -340,6 +345,7 @@ class AnalysisTaskQueue: original_query=original_query, selection_source=selection_source, report_type=report_type, + analysis_phase=analysis_phase, force_refresh=force_refresh, skills=skills, ) @@ -354,6 +360,7 @@ class AnalysisTaskQueue: original_query: Optional[str] = None, selection_source: Optional[str] = None, report_type: str = "detailed", + analysis_phase: str = "auto", force_refresh: bool = False, notify: bool = True, skills: Optional[List[str]] = None, @@ -393,6 +400,7 @@ class AnalysisTaskQueue: status=TaskStatus.PENDING, message="任务已加入队列", report_type=report_type, + analysis_phase=analysis_phase or "auto", original_query=original_query, selection_source=selection_source, skills=task_skills, @@ -613,6 +621,7 @@ class AnalysisTaskQueue: if not task: return None trace_id = task.trace_id or task_id + analysis_phase = task.analysis_phase task.status = TaskStatus.PROCESSING task.started_at = datetime.now() task.message = "正在分析中..." @@ -648,6 +657,7 @@ class AnalysisTaskQueue: send_notification=notify, progress_callback=_on_progress, skills=skills, + analysis_phase=analysis_phase, ) reset_run_diagnostic_context(diag_token) diag_token = None diff --git a/tests/test_analysis_api_contract.py b/tests/test_analysis_api_contract.py index 9f19af745..712e5d9ec 100644 --- a/tests/test_analysis_api_contract.py +++ b/tests/test_analysis_api_contract.py @@ -25,6 +25,7 @@ try: _build_analysis_report, _load_sync_fundamental_sources, get_analysis_status, + get_task_list, ) except Exception: # pragma: no cover - optional dependency environments create_app = None @@ -35,6 +36,7 @@ except Exception: # pragma: no cover - optional dependency environments _build_analysis_report = None _load_sync_fundamental_sources = None get_analysis_status = None + get_task_list = None from src.enums import ReportType from src.services.analysis_service import AnalysisService @@ -346,6 +348,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): error=None, original_query=None, selection_source=None, + analysis_phase="auto", ) with patch("api.v1.endpoints.analysis.get_task_queue", return_value=queue): @@ -378,6 +381,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): error=None, original_query=None, selection_source=None, + analysis_phase="auto", created_at=created_at, completed_at=datetime(2026, 5, 21, 17, 45, 0), ) @@ -419,6 +423,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): error=None, original_query=None, selection_source=None, + analysis_phase="auto", created_at=created_at, completed_at=datetime(2026, 5, 21, 17, 45, 0), ) @@ -773,6 +778,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): report_type="detailed", force_refresh=False, notify=True, + analysis_phase="auto", ), ) @@ -823,9 +829,14 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, notify=True, skills=None, + analysis_phase="intraday", ), ) + self.assertEqual( + service_instance.analyze_stock.call_args.kwargs["analysis_phase"], + "intraday", + ) details = result.report["details"] self.assertEqual(result.report["meta"]["market_phase_summary"]["phase"], "intraday") self.assertEqual( @@ -892,6 +903,72 @@ class AnalysisApiContractTestCase(unittest.TestCase): news_component = result["diagnostic_summary"]["components"]["news"] self.assertEqual(news_component["status"], "unknown") + def test_build_analysis_response_includes_market_phase_summary_from_result_snapshot(self) -> None: + service = AnalysisService() + phase_summary = _market_phase_summary() + + result = service._build_analysis_response( + SimpleNamespace( + code="600519", + name="贵州茅台", + current_price=1234.56, + change_pct=1.23, + model_used="test-model", + analysis_summary="summary", + operation_advice="hold", + trend_prediction="up", + sentiment_score=80, + news_summary="news", + technical_analysis="tech", + fundamental_analysis="fundamental", + risk_warning="risk", + diagnostic_context_snapshot={"market_phase_summary": phase_summary}, + get_sniper_points=lambda: {}, + ), + "q1", + report_type="full", + ) + + self.assertEqual( + result["report"]["meta"]["market_phase_summary"]["phase"], + "intraday", + ) + + def test_analysis_service_passes_analysis_phase_to_pipeline(self) -> None: + service = AnalysisService() + pipeline_instance = MagicMock() + pipeline_instance.process_single_stock.return_value = SimpleNamespace( + success=True, + code="600519", + name="贵州茅台", + current_price=1234.56, + change_pct=1.23, + model_used="test-model", + analysis_summary="summary", + operation_advice="hold", + trend_prediction="up", + sentiment_score=80, + news_summary="news", + technical_analysis="tech", + fundamental_analysis="fundamental", + risk_warning="risk", + get_sniper_points=lambda: {}, + ) + + with patch("src.config.get_config", return_value=SimpleNamespace()), patch( + "src.core.pipeline.StockAnalysisPipeline", + return_value=pipeline_instance, + ) as pipeline_cls: + result = service.analyze_stock( + "600519", + report_type="detailed", + send_notification=False, + analysis_phase="postmarket", + ) + + self.assertIsNotNone(result) + self.assertEqual(pipeline_cls.call_args.kwargs["analysis_phase"], "postmarket") + def test_build_analysis_report_extracts_fundamental_fields_from_snapshot(self) -> None: if _build_analysis_report is None: self.skipTest("analysis endpoint helpers unavailable in this environment") @@ -1036,6 +1113,61 @@ class AnalysisApiContractTestCase(unittest.TestCase): report.details.context_snapshot, ) + def test_build_analysis_report_falls_back_to_sanitized_report_meta_phase_summary(self) -> None: + if _build_analysis_report is None: + self.skipTest("analysis endpoint helpers unavailable in this environment") + + phase_summary = { + **_market_phase_summary(), + "warnings": ["api_key=secret"], + "market_phase_context": {"raw": True}, + } + + report = _build_analysis_report( + report_data={ + "meta": {"market_phase_summary": phase_summary}, + "summary": {}, + "strategy": {}, + "details": {}, + }, + query_id="q-meta-phase", + stock_code="600519", + stock_name="贵州茅台", + context_snapshot=None, + fallback_fundamental_payload=None, + ) + + self.assertIsNotNone(report.meta.market_phase_summary) + self.assertEqual(report.meta.market_phase_summary.phase, "intraday") + self.assertEqual(report.meta.market_phase_summary.warnings, ["[REDACTED]"]) + + def test_build_analysis_report_prefers_snapshot_phase_summary_over_report_meta(self) -> None: + if _build_analysis_report is None: + self.skipTest("analysis endpoint helpers unavailable in this environment") + + snapshot_summary = _market_phase_summary() + report = _build_analysis_report( + report_data={ + "meta": { + "market_phase_summary": { + **snapshot_summary, + "phase": "postmarket", + }, + }, + "summary": {}, + "strategy": {}, + "details": {}, + }, + query_id="q-snapshot-phase", + stock_code="600519", + stock_name="贵州茅台", + context_snapshot={"market_phase_summary": snapshot_summary}, + fallback_fundamental_payload=None, + ) + + self.assertIsNotNone(report.meta.market_phase_summary) + self.assertEqual(report.meta.market_phase_summary.phase, "intraday") + def test_build_analysis_report_merges_partial_top_level_context_with_fallback(self) -> None: if _build_analysis_report is None: self.skipTest("analysis endpoint helpers unavailable in this environment") @@ -1365,6 +1497,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): self.assertEqual(status.status, "completed") self.assertEqual(status.result.report["meta"]["current_price"], 1888.0) self.assertEqual(status.result.report["meta"]["change_pct"], 1.56) + self.assertIsNone(status.analysis_phase) self.assertEqual( status.result.report["meta"]["market_phase_summary"]["phase"], "intraday", @@ -1443,6 +1576,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): error=None, original_query=None, selection_source=None, + analysis_phase="auto", skills=None, created_at=datetime(2026, 4, 10, 12, 0, 0), completed_at=datetime(2026, 4, 10, 12, 1, 0), @@ -1494,16 +1628,18 @@ class AnalysisApiContractTestCase(unittest.TestCase): limit=1, ) - def test_get_analysis_status_in_memory_task_without_db_snapshot_omits_phase_summary(self) -> None: + def test_get_analysis_status_in_memory_task_without_db_snapshot_preserves_service_phase_summary(self) -> None: if get_analysis_status is None: self.skipTest("analysis endpoint helpers unavailable in this environment") + phase_summary = _market_phase_summary() task = SimpleNamespace( task_id="task_no_snapshot_in_memory_1", stock_code="600519", stock_name="贵州茅台", status=TaskStatus.COMPLETED, progress=100, + analysis_phase="intraday", result={ "stock_code": "600519", "stock_name": "贵州茅台", @@ -1512,6 +1648,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): "query_id": "task_no_snapshot_in_memory_1", "stock_code": "600519", "stock_name": "贵州茅台", + "market_phase_summary": phase_summary, }, "summary": {"analysis_summary": "summary"}, }, @@ -1533,9 +1670,11 @@ class AnalysisApiContractTestCase(unittest.TestCase): status = get_analysis_status("task_no_snapshot_in_memory_1") self.assertEqual(status.status, "completed") + self.assertEqual(status.analysis_phase, "intraday") self.assertIsNotNone(status.result) - self.assertIsNone( - status.result.report["meta"].get("market_phase_summary"), + self.assertEqual( + status.result.report["meta"]["market_phase_summary"]["phase"], + "intraday", ) load_sources.assert_called_once_with( query_id="task_no_snapshot_in_memory_1", @@ -1600,6 +1739,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): report_type="detailed", force_refresh=False, async_mode=False, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1623,6 +1763,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): report_type="detailed", force_refresh=False, async_mode=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1645,6 +1786,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): report_type="detailed", force_refresh=False, async_mode=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1673,6 +1815,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1685,6 +1828,50 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="AAPL.US", selection_source="manual", report_type="detailed", + analysis_phase="auto", + force_refresh=False, + notify=True, + ) + + def test_trigger_analysis_async_passes_and_returns_analysis_phase(self) -> None: + if trigger_analysis is None: + self.skipTest("fastapi is not installed in this test environment") + + task = SimpleNamespace( + task_id="task-phase-1", + trace_id="trace-phase-1", + stock_code="600519", + analysis_phase="intraday", + ) + queue = MagicMock() + queue.submit_tasks_batch.return_value = ([task], []) + + with patch("api.v1.endpoints.analysis.get_task_queue", return_value=queue): + response = trigger_analysis( + request=SimpleNamespace( + stock_code="600519", + stock_codes=None, + stock_name=None, + original_query=None, + selection_source=None, + report_type="detailed", + force_refresh=False, + async_mode=True, + notify=True, + analysis_phase="intraday", + ), + config=SimpleNamespace(), + ) + + self.assertEqual(response.status_code, 202) + self.assertEqual(json.loads(response.body)["analysis_phase"], "intraday") + queue.submit_tasks_batch.assert_called_once_with( + stock_codes=["600519"], + stock_name=None, + original_query=None, + selection_source=None, + report_type="detailed", + analysis_phase="intraday", force_refresh=False, notify=True, ) @@ -1708,6 +1895,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): report_type="detailed", force_refresh=False, async_mode=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1720,6 +1908,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="00700", selection_source="autocomplete", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -1744,6 +1933,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1756,6 +1946,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="920493", selection_source="autocomplete", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -1782,6 +1973,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1810,6 +2002,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): report_type="detailed", force_refresh=False, async_mode=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1822,6 +2015,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="HK00700", selection_source="manual", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -1846,6 +2040,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1857,6 +2052,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="西安奕材-U", selection_source="manual", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -1881,6 +2077,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1892,6 +2089,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="贵州茅台", selection_source="manual", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -1915,6 +2113,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1926,6 +2125,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="uploaded.csv", selection_source="import", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -1952,6 +2152,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -1966,6 +2167,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -2005,6 +2207,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): force_refresh=False, async_mode=True, notify=True, + analysis_phase="auto", ), config=SimpleNamespace(), ) @@ -2016,6 +2219,7 @@ class AnalysisApiContractTestCase(unittest.TestCase): original_query="茅台,平安银行", selection_source="import", report_type="detailed", + analysis_phase="auto", force_refresh=False, notify=True, ) @@ -2110,6 +2314,44 @@ class AnalysisApiContractTestCase(unittest.TestCase): mock_task_queue.unsubscribe.assert_called_once_with(never_queue) + def test_get_task_list_includes_analysis_phase_and_skills(self) -> None: + if get_task_list is None: + self.skipTest("analysis endpoint helpers unavailable in this environment") + + task = SimpleNamespace( + task_id="task-list-phase", + trace_id="trace-list-phase", + stock_code="600519", + stock_name="贵州茅台", + status=TaskStatus.PROCESSING, + progress=42, + message="running", + report_type="detailed", + created_at=datetime(2026, 4, 10, 12, 0, 0), + started_at=datetime(2026, 4, 10, 12, 0, 1), + completed_at=None, + error=None, + original_query="茅台", + selection_source="manual", + analysis_phase="postmarket", + skills=["growth_quality"], + ) + queue = MagicMock() + queue.list_all_tasks.return_value = [task] + queue.get_task_stats.return_value = { + "total": 1, + "pending": 0, + "processing": 1, + "completed": 0, + "failed": 0, + } + + with patch("api.v1.endpoints.analysis.get_task_queue", return_value=queue): + response = get_task_list(status=None, limit=20) + + self.assertEqual(response.tasks[0].analysis_phase, "postmarket") + self.assertEqual(response.tasks[0].skills, ["growth_quality"]) + class BatchTaskQueueContractTestCase(unittest.TestCase): def setUp(self) -> None: @@ -2172,11 +2414,15 @@ class BatchTaskQueueContractTestCase(unittest.TestCase): accepted, duplicates = queue.submit_tasks_batch( ["600519"], report_type="detailed", + analysis_phase="intraday", skills=request_skills, ) request_skills.append("mutated_after_submit") self.assertEqual(duplicates, []) + self.assertEqual(accepted[0].analysis_phase, "intraday") + self.assertEqual(accepted[0].to_dict()["analysis_phase"], "intraday") + self.assertEqual(accepted[0].copy().analysis_phase, "intraday") self.assertEqual(accepted[0].skills, ["growth_quality"]) self.assertIs(executor.calls[0][1][-1], accepted[0].skills) @@ -2190,6 +2436,7 @@ class BatchTaskQueueContractTestCase(unittest.TestCase): accepted[0].skills, ) self.assertEqual(service_instance.analyze_stock.call_args.kwargs["skills"], ["growth_quality"]) + self.assertEqual(service_instance.analyze_stock.call_args.kwargs["analysis_phase"], "intraday") def test_batch_submit_deduplicates_equivalent_stock_code_shapes(self) -> None: queue = AnalysisTaskQueue(max_workers=1) @@ -2202,7 +2449,11 @@ class BatchTaskQueueContractTestCase(unittest.TestCase): self.assertTrue(queue.is_analyzing("600519.SH")) self.assertEqual(queue.get_analyzing_task_id("600519.SH"), accepted[0].task_id) - accepted_again, duplicates_again = queue.submit_tasks_batch(["600519.SH"], report_type="detailed") + accepted_again, duplicates_again = queue.submit_tasks_batch( + ["600519.SH"], + report_type="detailed", + analysis_phase="intraday", + ) self.assertEqual(accepted_again, []) self.assertEqual(len(duplicates_again), 1) diff --git a/tests/test_analysis_integration.py b/tests/test_analysis_integration.py index 55321af8a..d29a19b2f 100644 --- a/tests/test_analysis_integration.py +++ b/tests/test_analysis_integration.py @@ -48,7 +48,7 @@ class TestAnalysisIntegration: """Test flow: User enters stock name -> resolved to code -> task submitted.""" # Setup mock behavior mock_task_queue.submit_tasks_batch.return_value = ( - [MagicMock(task_id="test_task_123", stock_code="600519")], + [MagicMock(task_id="test_task_123", stock_code="600519", analysis_phase="auto")], [] ) @@ -78,6 +78,7 @@ class TestAnalysisIntegration: assert kwargs["original_query"] == "贵州茅台" assert kwargs["selection_source"] == "manual" assert kwargs["report_type"] == "detailed" + assert kwargs["analysis_phase"] == "auto" assert kwargs["force_refresh"] is False assert kwargs["notify"] is True @@ -98,6 +99,7 @@ class TestAnalysisIntegration: args, kwargs = mock_task_queue.submit_tasks_batch.call_args assert len(kwargs["stock_codes"]) == 1 assert kwargs["stock_codes"] == ["600519"] + assert kwargs["analysis_phase"] == "auto" def test_trigger_analysis_dos_protection(self, client): """Test that excessive stock codes are rejected.""" @@ -133,3 +135,26 @@ class TestAnalysisIntegration: assert kwargs["stock_name"] is None assert kwargs["original_query"] is None assert kwargs["selection_source"] is None + assert kwargs["analysis_phase"] == "auto" + + def test_trigger_analysis_explicit_analysis_phase(self, client, mock_task_queue): + """Explicit analysis_phase is passed through to the task queue.""" + mock_task_queue.submit_tasks_batch.return_value = ( + [MagicMock(task_id="test_task_phase", stock_code="600519", analysis_phase="intraday")], + [] + ) + + response = client.post( + "/api/v1/analysis/analyze", + json={ + "stock_code": "600519", + "async_mode": True, + "analysis_phase": "intraday", + }, + ) + + assert response.status_code == 202 + assert response.json()["analysis_phase"] == "intraday" + mock_task_queue.submit_tasks_batch.assert_called_once() + _, kwargs = mock_task_queue.submit_tasks_batch.call_args + assert kwargs["analysis_phase"] == "intraday" diff --git a/tests/test_api_schema_pydantic.py b/tests/test_api_schema_pydantic.py index 282a9159f..3c74f032a 100644 --- a/tests/test_api_schema_pydantic.py +++ b/tests/test_api_schema_pydantic.py @@ -16,6 +16,13 @@ def test_schema_examples_remain_in_openapi_schema() -> None: assert root_schema["example"]["version"] == "1.0.0" assert analyze_schema["properties"]["stock_code"]["example"] == "600519" assert analyze_schema["properties"]["skills"]["example"] == ["bull_trend", "growth_quality"] + assert analyze_schema["properties"]["analysis_phase"]["default"] == "auto" + assert analyze_schema["properties"]["analysis_phase"]["enum"] == [ + "auto", + "premarket", + "intraday", + "postmarket", + ] assert history_schema["example"]["stock_code"] == "600519" assert quote_schema["example"]["stock_name"] == "贵州茅台" @@ -27,3 +34,21 @@ def test_analyze_request_supports_legacy_strategies_dict_input() -> None: }) assert request.skills == ["bull_trend", "growth_quality"] + + +def test_analyze_request_analysis_phase_defaults_to_auto() -> None: + request = AnalyzeRequest(stock_code="600519") + + assert request.analysis_phase == "auto" + + +def test_analyze_request_rejects_invalid_analysis_phase() -> None: + try: + AnalyzeRequest.model_validate({ + "stock_code": "600519", + "analysis_phase": "lunch_break", + }) + except Exception as exc: + assert "analysis_phase" in str(exc) + else: + raise AssertionError("invalid analysis_phase should be rejected") diff --git a/tests/test_pipeline_market_phase_context.py b/tests/test_pipeline_market_phase_context.py index e76e2770d..29ed02b0e 100644 --- a/tests/test_pipeline_market_phase_context.py +++ b/tests/test_pipeline_market_phase_context.py @@ -68,6 +68,7 @@ def _make_pipeline(*, agent_mode: bool = False, save_context_snapshot: bool = Tr pipeline.save_context_snapshot = save_context_snapshot pipeline.progress_callback = None pipeline.analysis_skills = None + pipeline.analysis_phase = "auto" pipeline.social_sentiment_service = None pipeline.fetcher_manager = MagicMock() @@ -282,6 +283,30 @@ class PipelineMarketPhaseContextTestCase(unittest.TestCase): self.assertIsInstance(result.dashboard["phase_decision"]["watch_conditions"], list) self.assertIn("daily_bars: missing", result.dashboard["phase_decision"]["data_limitations"]) + def test_pipeline_passes_configured_analysis_phase_to_market_context(self): + pipeline = _make_pipeline(agent_mode=False, save_context_snapshot=True) + pipeline.analysis_phase = "postmarket" + phase_payload = { + **_phase_payload(), + "phase": "postmarket", + "analysis_intent": "postmarket", + "is_market_open_now": False, + "is_partial_bar": False, + "minutes_to_close": None, + } + phase_context = SimpleNamespace(to_dict=MagicMock(return_value=phase_payload)) + + with patch("src.core.pipeline.build_market_phase_context", return_value=phase_context) as mock_build: + result = pipeline.analyze_stock( + "600519", + ReportType.SIMPLE, + "q-runtime-phase", + current_time=datetime(2026, 3, 27, 16, 0), + ) + + self.assertIsNotNone(result) + self.assertEqual(mock_build.call_args.kwargs["analysis_phase"], "postmarket") + def test_legacy_pipeline_fail_open_when_pack_summary_generation_fails(self): pipeline = _make_pipeline(agent_mode=False, save_context_snapshot=True) phase_payload = _phase_payload() diff --git a/tests/test_trading_calendar.py b/tests/test_trading_calendar.py index e022f0cf4..c56dc5653 100644 --- a/tests/test_trading_calendar.py +++ b/tests/test_trading_calendar.py @@ -610,6 +610,71 @@ class MarketPhaseContextTestCase(unittest.TestCase): effective_date.isoformat(), ) + def test_manual_analysis_phase_overrides_non_trading_day_without_rewriting_calendar_fields(self): + fake_calendar = _FakeCalendar( + sessions=[date(2026, 3, 26), date(2026, 3, 27)], + close_hour=15, + tz_name="Asia/Shanghai", + open_time=time(9, 30), + break_start=time(11, 30), + break_end=time(13, 0), + ) + + with patch.object(trading_calendar, "_XCALS_AVAILABLE", True), patch.object( + trading_calendar, + "xcals", + _calendar_namespace(fake_calendar), + create=True, + ): + ctx = trading_calendar.build_market_phase_context( + market="cn", + current_time=datetime(2026, 3, 28, 10, 0, tzinfo=ZoneInfo("Asia/Shanghai")), + trigger_source="api", + analysis_phase="intraday", + ) + + payload = ctx.to_dict() + self.assertEqual(payload["phase"], "intraday") + self.assertEqual(payload["analysis_intent"], "intraday") + self.assertEqual(payload["market_local_time"], "2026-03-28T10:00:00+08:00") + self.assertEqual(payload["effective_daily_bar_date"], "2026-03-27") + self.assertTrue(payload["is_trading_day"]) + self.assertTrue(payload["is_market_open_now"]) + self.assertTrue(payload["is_partial_bar"]) + self.assertIsNone(payload["minutes_to_open"]) + self.assertIsNone(payload["minutes_to_close"]) + + def test_legacy_analysis_intent_alias_can_override_phase(self): + fake_calendar = _FakeCalendar( + sessions=[date(2026, 3, 26), date(2026, 3, 27)], + close_hour=15, + tz_name="Asia/Shanghai", + open_time=time(9, 30), + ) + + with patch.object(trading_calendar, "_XCALS_AVAILABLE", True), patch.object( + trading_calendar, + "xcals", + _calendar_namespace(fake_calendar), + create=True, + ): + ctx = trading_calendar.build_market_phase_context( + market="cn", + current_time=datetime(2026, 3, 27, 10, 0, tzinfo=ZoneInfo("Asia/Shanghai")), + analysis_intent="postmarket", + ) + + self.assertEqual(ctx.phase, trading_calendar.MarketPhase.POSTMARKET) + self.assertEqual(ctx.analysis_intent, "postmarket") + + def test_invalid_manual_analysis_phase_raises_value_error(self): + with self.assertRaisesRegex(ValueError, "invalid analysis_phase"): + trading_calendar.build_market_phase_context( + market="cn", + current_time=datetime(2026, 3, 27, 10, 0, tzinfo=ZoneInfo("Asia/Shanghai")), + analysis_phase="lunch_break", + ) + def test_unknown_market_uses_null_tristate_flags_and_warning_code(self): ctx = trading_calendar.build_market_phase_context( market=None,