feat: add analysis phase API plumbing (#1573)

This commit is contained in:
Alfred
2026-06-03 19:51:07 +08:00
committed by GitHub
parent d8a067ba47
commit 7e54c6c3ba
19 changed files with 589 additions and 18 deletions
+13 -1
View File
@@ -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),
+24 -4
View File
@@ -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"]
}
})
+2
View File
@@ -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,
@@ -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'],
});
});
});
+2
View File
@@ -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()) {
+7
View File
@@ -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 */
+1
View File
@@ -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 设置页仅展示已注册且带说明的配置项。
+55
View File
@@ -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": [
+11
View File
@@ -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`(共享装配路径)。该接口当前仅提供单进程/单机级防重复能力,若为多实例部署需通过外部任务队列或分布式锁补齐全局幂等。
+11
View File
@@ -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.
+3 -1
View File
@@ -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)
+44 -7
View File
@@ -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,
)
+6
View File
@@ -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,
+10
View File
@@ -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
+255 -4
View File
@@ -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)
+26 -1
View File
@@ -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"
+25
View File
@@ -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")
@@ -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()
+65
View File
@@ -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,