feat: link decision signals with alerts and portfolio risk (#1715)

Co-authored-by: mumu <42829555+ZhuLinsen@users.noreply.github.com>
This commit is contained in:
Alfred
2026-06-19 09:43:40 +08:00
committed by GitHub
parent 626b63c144
commit a7876fe11c
28 changed files with 2047 additions and 5 deletions

View File

@@ -93,6 +93,8 @@ from api.v1.schemas.portfolio import (
PortfolioImportBrokerItem,
PortfolioImportBrokerListResponse,
PortfolioFxRefreshResponse,
PortfolioDecisionSignalRiskBlock,
PortfolioDecisionSignalRiskItem,
PortfolioRiskResponse,
)
from api.v1.schemas.alerts import (
@@ -204,6 +206,8 @@ __all__ = [
"PortfolioImportBrokerItem",
"PortfolioImportBrokerListResponse",
"PortfolioFxRefreshResponse",
"PortfolioDecisionSignalRiskBlock",
"PortfolioDecisionSignalRiskItem",
"PortfolioRiskResponse",
# alerts
"AlertDeleteResponse",

View File

@@ -117,6 +117,7 @@ class AlertTriggerItem(BaseModel):
"evaluator_snapshot / legacy_text / null"
),
)
decision_signal_summary: Optional[Dict[str, Any]] = None
class AlertTriggerListResponse(BaseModel):

View File

@@ -263,6 +263,20 @@ class PortfolioFxRefreshResponse(BaseModel):
error_count: int
class PortfolioDecisionSignalRiskItem(BaseModel):
account_id: Optional[int] = None
symbol: str
market: str
signal: Dict[str, Any] = Field(default_factory=dict)
class PortfolioDecisionSignalRiskBlock(BaseModel):
available: bool = True
total: int = 0
actions: Dict[str, int] = Field(default_factory=dict)
items: List[PortfolioDecisionSignalRiskItem] = Field(default_factory=list)
class PortfolioRiskResponse(BaseModel):
as_of: str
account_id: Optional[int] = None
@@ -273,3 +287,4 @@ class PortfolioRiskResponse(BaseModel):
sector_concentration: Dict[str, Any] = Field(default_factory=dict)
drawdown: Dict[str, Any] = Field(default_factory=dict)
stop_loss: Dict[str, Any] = Field(default_factory=dict)
decision_signal_risk: PortfolioDecisionSignalRiskBlock = Field(default_factory=PortfolioDecisionSignalRiskBlock)

View File

@@ -60,8 +60,8 @@ export const ALERT_MARKET_REGION_OPTIONS: Record<UiLanguage, Array<Option<Market
export const ALERT_MARKET_LIGHT_STATUS_OPTIONS: Record<UiLanguage, Array<Option<MarketLightStatus>>> = { zh: [{ value: 'red', label: '红灯' }, { value: 'yellow', label: '黄灯' }], en: [{ value: 'red', label: 'Red' }, { value: 'yellow', label: 'Yellow' }] };
export const PORTFOLIO_TEXT = {
zh: { documentTitle: '持仓分析 - DSA', title: '持仓管理', description: '组合快照、手工录入、CSV 导入与风险分析(支持全组合 / 单账户切换)', accountView: '账户视图', allAccounts: '全部账户', costMethod: '成本口径', fifo: '先进先出FIFO', avg: '均价成本AVG', collapseCreate: '收起新建', createAccount: '新建账户', deleteAccount: '删除账户', deletingAccount: '删除中...', deleteAccountTitle: '删除持仓账户', deleteAccountConfirm: '确认删除', deleteAccountMessage: '确认删除账户 {name}#{id})吗?删除后该账户会从默认列表、快照、风险和录入入口隐藏;历史流水不会物理删除。', refreshing: '刷新中...', refreshData: '刷新数据', noAccounts: '还没有可用账户,请先创建账户后再录入交易或导入 CSV。', riskDegraded: '风险模块降级', operationHint: '操作提示', analysisTask: '分析任务', totalEquity: '总权益', totalMarketValue: '总市值', totalCash: '总现金', fxStatus: '汇率状态', refreshFx: '刷新汇率', stale: '过期', latest: '最新', fxRefreshResult: '汇率刷新结果', positionsTitle: '持仓明细', countItems: '共 {count} 项', noPositionsTitle: '当前无持仓数据', noPositionsDescription: '录入交易或导入 CSV 后,这里会展示按账户汇总的持仓明细。', account: '账户', code: '代码', quantity: '数量', avgCost: '均价', lastPrice: '现价', marketValue: '市值', unrealizedPnl: '未实现盈亏', returnPct: '收益率', action: '操作', submitting: '提交中', analyze: '分析', sectorConcentration: '行业集中度分布', positionConcentrationFallback: '行业数据暂不可用,当前展示个股集中度', noConcentrationTitle: '暂无集中度数据', noConcentrationDescription: '风险模块完成计算后,这里会展示行业或个股维度的集中度分布。', displayScope: '展示口径', sectorDimension: '行业维度', positionDimensionFallback: '个股维度(降级显示)', sectorAlert: '板块集中度告警', topWeight: 'Top1 权重', yes: '是', no: '否', writeBlocked: '当前处于“全部账户”视图。为避免误写,请先选择一个具体账户后再进行手工录入或 CSV 提交。', drawdownMonitor: '回撤监控', maxDrawdown: '最大回撤', currentDrawdown: '当前回撤', alert: '告警', stopLossWarning: '止损接近预警', triggeredCount: '触发数', nearCount: '接近数', scope: '口径', accountCount: '账户数', currency: '计价币种', costMethodShort: '成本法' },
en: { documentTitle: 'Portfolio Analysis - DSA', title: 'Portfolio management', description: 'Portfolio snapshots, manual entries, CSV import, and risk analysis with full-portfolio or single-account views', accountView: 'Account view', allAccounts: 'All accounts', costMethod: 'Cost method', fifo: 'FIFO', avg: 'Average cost', collapseCreate: 'Collapse', createAccount: 'New account', deleteAccount: 'Delete account', deletingAccount: 'Deleting...', deleteAccountTitle: 'Delete portfolio account', deleteAccountConfirm: 'Delete account', deleteAccountMessage: 'Delete account {name} (#{id})? It will be hidden from default lists, snapshots, risk views, and entry forms; historical ledger rows are not physically deleted.', refreshing: 'Refreshing...', refreshData: 'Refresh data', noAccounts: 'No accounts are available. Create an account before entering trades or importing CSV files.', riskDegraded: 'Risk module degraded', operationHint: 'Operation hint', analysisTask: 'Analysis task', totalEquity: 'Total equity', totalMarketValue: 'Total market value', totalCash: 'Total cash', fxStatus: 'FX status', refreshFx: 'Refresh FX', stale: 'Stale', latest: 'Current', fxRefreshResult: 'FX refresh result', positionsTitle: 'Positions', countItems: '{count} items', noPositionsTitle: 'No positions', noPositionsDescription: 'After you enter trades or import CSV data, account-level positions appear here.', account: 'Account', code: 'Code', quantity: 'Quantity', avgCost: 'Avg cost', lastPrice: 'Last price', marketValue: 'Market value', unrealizedPnl: 'Unrealized P/L', returnPct: 'Return', action: 'Action', submitting: 'Submitting', analyze: 'Analyze', sectorConcentration: 'Sector concentration', positionConcentrationFallback: 'Sector data unavailable; showing position concentration', noConcentrationTitle: 'No concentration data', noConcentrationDescription: 'Sector or position concentration appears after the risk module finishes.', displayScope: 'Display scope', sectorDimension: 'Sector', positionDimensionFallback: 'Position fallback', sectorAlert: 'Sector concentration alert', topWeight: 'Top1 weight', yes: 'Yes', no: 'No', writeBlocked: 'You are viewing all accounts. Select a specific account before manual entry or CSV submission to avoid writing to the wrong scope.', drawdownMonitor: 'Drawdown monitor', maxDrawdown: 'Max drawdown', currentDrawdown: 'Current drawdown', alert: 'Alert', stopLossWarning: 'Stop-loss proximity warning', triggeredCount: 'Triggered', nearCount: 'Near', scope: 'Scope', accountCount: 'Accounts', currency: 'Quote currency', costMethodShort: 'Cost method' },
zh: { documentTitle: '持仓分析 - DSA', title: '持仓管理', description: '组合快照、手工录入、CSV 导入与风险分析(支持全组合 / 单账户切换)', accountView: '账户视图', allAccounts: '全部账户', costMethod: '成本口径', fifo: '先进先出FIFO', avg: '均价成本AVG', collapseCreate: '收起新建', createAccount: '新建账户', deleteAccount: '删除账户', deletingAccount: '删除中...', deleteAccountTitle: '删除持仓账户', deleteAccountConfirm: '确认删除', deleteAccountMessage: '确认删除账户 {name}#{id})吗?删除后该账户会从默认列表、快照、风险和录入入口隐藏;历史流水不会物理删除。', refreshing: '刷新中...', refreshData: '刷新数据', noAccounts: '还没有可用账户,请先创建账户后再录入交易或导入 CSV。', riskDegraded: '风险模块降级', operationHint: '操作提示', analysisTask: '分析任务', totalEquity: '总权益', totalMarketValue: '总市值', totalCash: '总现金', fxStatus: '汇率状态', refreshFx: '刷新汇率', stale: '过期', latest: '最新', fxRefreshResult: '汇率刷新结果', positionsTitle: '持仓明细', countItems: '共 {count} 项', noPositionsTitle: '当前无持仓数据', noPositionsDescription: '录入交易或导入 CSV 后,这里会展示按账户汇总的持仓明细。', account: '账户', code: '代码', quantity: '数量', avgCost: '均价', lastPrice: '现价', marketValue: '市值', unrealizedPnl: '未实现盈亏', returnPct: '收益率', action: '操作', submitting: '提交中', analyze: '分析', sectorConcentration: '行业集中度分布', positionConcentrationFallback: '行业数据暂不可用,当前展示个股集中度', noConcentrationTitle: '暂无集中度数据', noConcentrationDescription: '风险模块完成计算后,这里会展示行业或个股维度的集中度分布。', displayScope: '展示口径', sectorDimension: '行业维度', positionDimensionFallback: '个股维度(降级显示)', sectorAlert: '板块集中度告警', topWeight: 'Top1 权重', yes: '是', no: '否', writeBlocked: '当前处于“全部账户”视图。为避免误写,请先选择一个具体账户后再进行手工录入或 CSV 提交。', drawdownMonitor: '回撤监控', maxDrawdown: '最大回撤', currentDrawdown: '当前回撤', alert: '告警', stopLossWarning: '止损接近预警', triggeredCount: '触发数', nearCount: '接近数', scope: '口径', accountCount: '账户数', currency: '计价币种', costMethodShort: '成本法', aiRiskSignals: 'AI 风险信号', aiRiskUnavailable: '信号风险暂不可用', aiRiskTotal: '风险信号', sellSignals: '卖出', reduceSignals: '减仓', alertSignals: '预警', noAiRiskSignals: '暂无防御型信号' },
en: { documentTitle: 'Portfolio Analysis - DSA', title: 'Portfolio management', description: 'Portfolio snapshots, manual entries, CSV import, and risk analysis with full-portfolio or single-account views', accountView: 'Account view', allAccounts: 'All accounts', costMethod: 'Cost method', fifo: 'FIFO', avg: 'Average cost', collapseCreate: 'Collapse', createAccount: 'New account', deleteAccount: 'Delete account', deletingAccount: 'Deleting...', deleteAccountTitle: 'Delete portfolio account', deleteAccountConfirm: 'Delete account', deleteAccountMessage: 'Delete account {name} (#{id})? It will be hidden from default lists, snapshots, risk views, and entry forms; historical ledger rows are not physically deleted.', refreshing: 'Refreshing...', refreshData: 'Refresh data', noAccounts: 'No accounts are available. Create an account before entering trades or importing CSV files.', riskDegraded: 'Risk module degraded', operationHint: 'Operation hint', analysisTask: 'Analysis task', totalEquity: 'Total equity', totalMarketValue: 'Total market value', totalCash: 'Total cash', fxStatus: 'FX status', refreshFx: 'Refresh FX', stale: 'Stale', latest: 'Current', fxRefreshResult: 'FX refresh result', positionsTitle: 'Positions', countItems: '{count} items', noPositionsTitle: 'No positions', noPositionsDescription: 'After you enter trades or import CSV data, account-level positions appear here.', account: 'Account', code: 'Code', quantity: 'Quantity', avgCost: 'Avg cost', lastPrice: 'Last price', marketValue: 'Market value', unrealizedPnl: 'Unrealized P/L', returnPct: 'Return', action: 'Action', submitting: 'Submitting', analyze: 'Analyze', sectorConcentration: 'Sector concentration', positionConcentrationFallback: 'Sector data unavailable; showing position concentration', noConcentrationTitle: 'No concentration data', noConcentrationDescription: 'Sector or position concentration appears after the risk module finishes.', displayScope: 'Display scope', sectorDimension: 'Sector', positionDimensionFallback: 'Position fallback', sectorAlert: 'Sector concentration alert', topWeight: 'Top1 weight', yes: 'Yes', no: 'No', writeBlocked: 'You are viewing all accounts. Select a specific account before manual entry or CSV submission to avoid writing to the wrong scope.', drawdownMonitor: 'Drawdown monitor', maxDrawdown: 'Max drawdown', currentDrawdown: 'Current drawdown', alert: 'Alert', stopLossWarning: 'Stop-loss proximity warning', triggeredCount: 'Triggered', nearCount: 'Near', scope: 'Scope', accountCount: 'Accounts', currency: 'Quote currency', costMethodShort: 'Cost method', aiRiskSignals: 'AI risk signals', aiRiskUnavailable: 'Signal risk unavailable', aiRiskTotal: 'Risk signals', sellSignals: 'Sell', reduceSignals: 'Reduce', alertSignals: 'Alert', noAiRiskSignals: 'No defensive signals' },
} as const;
export const PORTFOLIO_SIDE_LABELS: Record<UiLanguage, Record<PortfolioSide, string>> = { zh: { buy: '买入', sell: '卖出' }, en: { buy: 'Buy', sell: 'Sell' } };
export const PORTFOLIO_CASH_DIRECTION_LABELS: Record<UiLanguage, Record<PortfolioCashDirection, string>> = { zh: { in: '流入', out: '流出' }, en: { in: 'Inflow', out: 'Outflow' } };

View File

@@ -912,6 +912,8 @@ const PortfolioPage: React.FC = () => {
}
};
const decisionSignalRiskPreviewItems = (risk?.decisionSignalRisk?.items ?? []).slice(0, 3);
return (
<div className="portfolio-page min-h-screen space-y-4 p-4 md:p-6">
<section className="space-y-3">
@@ -1261,7 +1263,7 @@ const PortfolioPage: React.FC = () => {
/>
) : null}
<section className="grid grid-cols-1 md:grid-cols-3 gap-3">
<section className="grid grid-cols-1 md:grid-cols-2 xl:grid-cols-4 gap-3">
<Card padding="md">
<h3 className="text-sm font-semibold text-foreground mb-2">{text.drawdownMonitor}</h3>
<div className="text-xs text-secondary space-y-1">
@@ -1286,6 +1288,32 @@ const PortfolioPage: React.FC = () => {
<div>{text.costMethodShort}: {(snapshot?.costMethod || costMethod).toUpperCase()}</div>
</div>
</Card>
<Card padding="md">
<h3 className="text-sm font-semibold text-foreground mb-2">{text.aiRiskSignals}</h3>
<div className="text-xs text-secondary space-y-1">
{risk?.decisionSignalRisk?.available === false ? (
<div className="text-warning">{text.aiRiskUnavailable}</div>
) : (
<>
<div>{text.aiRiskTotal}: {risk?.decisionSignalRisk?.total ?? 0}</div>
<div>
{text.sellSignals}: {risk?.decisionSignalRisk?.actions?.sell ?? 0} · {text.reduceSignals}: {risk?.decisionSignalRisk?.actions?.reduce ?? 0} · {text.alertSignals}: {risk?.decisionSignalRisk?.actions?.alert ?? 0}
</div>
{decisionSignalRiskPreviewItems.length > 0 ? (
<div className="space-y-1 pt-1">
{decisionSignalRiskPreviewItems.map((item) => (
<div key={`${item.accountId ?? 'all'}-${item.market}-${item.symbol}-${item.signal.id ?? item.signal.action}`} className="truncate text-foreground">
{item.symbol} · {item.signal.actionLabel || item.signal.action || text.alert}
</div>
))}
</div>
) : (
<div>{text.noAiRiskSignals}</div>
)}
</>
)}
</div>
</Card>
</section>
<section className="grid grid-cols-1 xl:grid-cols-3 gap-3">

View File

@@ -182,7 +182,7 @@ function makePosition(overrides: Record<string, unknown> = {}) {
};
}
function makeRisk() {
function makeRisk(overrides: Record<string, unknown> = {}) {
return {
asOf: '2026-03-19',
accountId: null,
@@ -216,6 +216,13 @@ function makeRisk() {
nearCount: 0,
items: [],
},
decisionSignalRisk: {
available: true,
total: 0,
actions: { sell: 0, reduce: 0, alert: 0 },
items: [],
},
...overrides,
};
}
@@ -354,9 +361,62 @@ describe('PortfolioPage FX refresh', () => {
expect(screen.getByText(/Current drawdown:/)).toBeInTheDocument();
expect(screen.getByText('Stop-loss proximity warning')).toBeInTheDocument();
expect(screen.getByText('Scope')).toBeInTheDocument();
expect(screen.getByText('AI risk signals')).toBeInTheDocument();
expect(screen.getByText('No defensive signals')).toBeInTheDocument();
expect(screen.queryByText('回撤监控')).not.toBeInTheDocument();
});
it('renders portfolio decision signal risk summary', async () => {
getRisk.mockResolvedValueOnce(makeRisk({
decisionSignalRisk: {
available: true,
total: 2,
actions: { sell: 1, reduce: 0, alert: 1 },
items: [
{
accountId: 1,
symbol: '600519',
market: 'cn',
signal: makeDecisionSignal({ id: 201, action: 'sell', actionLabel: '卖出' }),
},
{
accountId: 1,
symbol: '300750',
market: 'cn',
signal: makeDecisionSignal({ id: 202, stockCode: '300750', action: 'alert', actionLabel: '预警' }),
},
],
},
}));
render(<PortfolioPage />);
await waitForInitialLoad();
expect(screen.getByText('AI 风险信号')).toBeInTheDocument();
expect(screen.getByText(/风险信号: 2/)).toBeInTheDocument();
expect(screen.getByText(/卖出: 1 · 减仓: 0 · 预警: 1/)).toBeInTheDocument();
expect(screen.getByText('600519 · 卖出')).toBeInTheDocument();
expect(screen.getByText('300750 · 预警')).toBeInTheDocument();
});
it('renders portfolio decision signal risk fail-open state', async () => {
getRisk.mockResolvedValueOnce(makeRisk({
decisionSignalRisk: {
available: false,
total: 0,
actions: { sell: 0, reduce: 0, alert: 0 },
items: [],
},
}));
render(<PortfolioPage />);
await waitForInitialLoad();
expect(screen.getByText('信号风险暂不可用')).toBeInTheDocument();
});
it('refreshes FX for a single selected account and only reloads snapshot/risk', async () => {
getSnapshot
.mockResolvedValueOnce(makeSnapshot({ fxStale: true }))

View File

@@ -1,4 +1,5 @@
import type { AnalysisContextPackOverview, MarketPhaseSummary } from './analysis';
import type { DecisionSignalItem } from './decisionSignals';
export type AlertType =
| 'price_cross'
@@ -122,6 +123,7 @@ export interface AlertTriggerItem {
marketPhaseSummary?: MarketPhaseSummary | null;
analysisContextPackOverview?: AnalysisContextPackOverview | null;
analysisVisibilitySource?: string | null;
decisionSignalSummary?: Partial<DecisionSignalItem> | null;
}
export interface AlertTriggerListResponse {

View File

@@ -1,3 +1,5 @@
import type { DecisionSignalItem } from './decisionSignals';
export type PortfolioCostMethod = 'fifo' | 'avg';
export type PortfolioSide = 'buy' | 'sell';
export type PortfolioCashDirection = 'in' | 'out';
@@ -121,6 +123,25 @@ export interface PortfolioStopLossItem {
isTriggered: boolean;
}
export interface PortfolioDecisionSignalRiskItem {
accountId?: number | null;
symbol: string;
market: string;
signal: Partial<DecisionSignalItem>;
}
export interface PortfolioDecisionSignalRiskBlock {
available: boolean;
total: number;
actions: {
sell?: number;
reduce?: number;
alert?: number;
[key: string]: number | undefined;
};
items: PortfolioDecisionSignalRiskItem[];
}
export interface PortfolioRiskResponse {
asOf: string;
accountId?: number | null;
@@ -148,6 +169,7 @@ export interface PortfolioRiskResponse {
nearCount: number;
items: PortfolioStopLossItem[];
};
decisionSignalRisk?: PortfolioDecisionSignalRiskBlock;
}
export interface PortfolioTradeCreateRequest {

View File

@@ -9,6 +9,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
## [Unreleased]
- [新功能] #1390 P6 将 DecisionSignal 复用到告警、通知和组合风险:告警触发关联 latest active 信号或创建最小 alert 信号,通知追加低敏信号摘要,持仓风险聚合 active sell/reduce/alert 信号并保持 fail-open。
- [新功能] #1707 资讯源新增 `newsnow` 类型、`NEWSNOW_BASE_URL` 配置和 `/api/v1/intelligence/sources/defaults` 默认源初始化接口,内置财联社热门、雪球热门股票、华尔街见闻快讯、金十数据和格隆汇事件等财经热点源,可直接拉取落库并进入既有分析证据链路;官方 NewsNow 部署指南见 https://github.com/qqhann/newsnow生产环境建议自建实例而非使用公开示例。
- [修复] AlphaSift 热点题材刷新在 EastMoney 瞬断且无缓存时返回友好空态,并让桌面更新保留 AlphaSift 热点缓存。

View File

@@ -342,6 +342,7 @@ P6 不做:
- `analysis_context_pack_overview` 只来自 evaluator 已带 overview 或最近 30 天内的历史 snapshot。最近历史查询复用历史服务的代码变体候选并以 best-effort + 批内短缓存方式执行;缺失或解析失败返回 `null`,不伪造 pack。
- 告警通知只输出公开摘要阶段标签、trigger source、partial-bar warning、数据质量等级和前两条 limitations。通知不得输出 raw context pack、Prompt、新闻正文、完整 diagnostics JSON、webhook URL、token 或持仓敏感细节。
- Web 告警历史展示 phase badge、数据质量等级和 limitations 空态;旧触发记录缺少公开摘要时不影响列表读取。
- #1390 P6 进一步复用 `DecisionSignal`:股票级真实触发会优先关联同标的 latest active 信号,并把低敏 `decision_signal_summary` 写入 diagnostics无 active 信号时只创建最小 `source_type=alert/action=alert` 信号。`trace_id=alert-rule-<hash>` 只用于同源重试的 best-effort 幂等去重,不覆盖 active 信号;新建告警信号不写 `market_phase`,避免同一规则跨阶段重复创建。`market``portfolio_account`、overflow 或无法解析为具体股票的触发不会创建个股信号。
#1386 P7 的用户边界:告警联动只解释触发时已经可公开的阶段和数据质量摘要,不会自动发起轻量 LLM 盘中分析,也不会新增告警表、规则类型、环境变量或 migration。需要阶段化分析时仍应通过分析 API / Web 手动分析入口触发告警通知只保留阶段标签、trigger source、partial-bar warning、数据质量等级和前两条 limitations。

File diff suppressed because it is too large Load Diff

View File

@@ -1342,6 +1342,8 @@ P5 在 Web `/decision-signals` 页面筛选区下方展示当前 outcome engine
持仓页会把 AI 建议作为非阻断增强异步加载:组合快照和风险模块先按原逻辑渲染,随后按当前快照中的唯一持仓调用 `GET /api/v1/decision-signals/latest/{stock_code}?market=<market>&limit=1` 查询 latest active 信号;不再通过 `holding_only=true` 通用列表分页扫描,也不存在固定页数截断。单个持仓 latest 查询失败时,页面保留其他已加载信号并显示可见降级提示;无匹配信号时持仓行显示空占位。匹配逻辑复用 Web 端股票代码等价规则,覆盖 A 股 `600519/SH600519/600519.SH`、港股 `00700/HK00700/00700.HK` 和美股大小写 ticker。
#1390 P6 将 `DecisionSignal` 复用到告警、通知和组合风险,不新增表、迁移或配置。真实股票级告警触发会优先关联同标的 latest active 信号,并把低敏 `decision_signal_summary` 写入 `alert_triggers.diagnostics`;没有 active 信号时worker 只创建最小 `source_type=alert`、`action=alert` 信号,`trace_id=alert-rule-<hash>` 仅用于同源重试的 best-effort 幂等去重,不覆盖 active 信号本体,且不写 `market_phase` 避免跨阶段重复。告警通知和分析通知只引用摘要中的 `action/horizon/reason/watch_conditions/risk_summary/source_report_id` 等公开字段,通知失败不影响 trigger 或信号写入。`GET /api/v1/portfolio/risk` 追加 `decision_signal_risk` 聚合块,只统计当前持仓中的 active `sell/reduce/alert` 信号,明确排除 `avoid/buy/add/hold/watch`;信号查询失败时风险接口 fail-openWeb 风险区显示降级状态。
普通个股历史报告详情会在策略区后展示该报告提取出的 `source_type=analysis` 信号,查询条件为 `source_report_id=<recordId>`;无 `recordId`、大盘复盘或其他非普通个股报告不会发起该查询。空结果显示“本报告暂无决策信号”,加载失败只影响该卡片,不影响报告主体、资讯、运行诊断或透明度区展示。
## 回测功能

View File

@@ -1177,6 +1177,8 @@ P5 extends the existing Web `/decision-signals` page instead of adding a new nav
The portfolio page loads AI signals as a non-blocking enhancement: portfolio snapshots and risk cards render first, then the page calls `GET /api/v1/decision-signals/latest/{stock_code}?market=<market>&limit=1` for each unique holding in the current snapshot to read the latest active signal. It no longer scans the generic `holding_only=true` list endpoint and has no fixed page-count cutoff. If a single latest lookup fails, the page keeps other loaded signals and shows a visible degradation warning; rows without a matching signal show an empty placeholder. Matching reuses the Web stock-code equivalence rules for CN variants such as `600519/SH600519/600519.SH`, HK variants such as `00700/HK00700/00700.HK`, and case-insensitive US tickers.
#1390 P6 reuses `DecisionSignal` across alerts, notifications, and portfolio risk without adding tables, migrations, or configuration. Real stock-level alert triggers first link the latest active signal for the same symbol and write a low-sensitive `decision_signal_summary` into `alert_triggers.diagnostics`; when no active signal exists, the worker creates only a minimal `source_type=alert`, `action=alert` signal. Its `trace_id=alert-rule-<hash>` is for best-effort retry de-duplication, not active-signal overwrites, and the payload intentionally omits `market_phase` to avoid cross-phase duplicates. Alert and analysis notifications reference only public summary fields such as `action/horizon/reason/watch_conditions/risk_summary/source_report_id`, and notification failure does not block trigger or signal writes. `GET /api/v1/portfolio/risk` now includes a `decision_signal_risk` block that counts active `sell/reduce/alert` signals for current holdings, explicitly excluding `avoid/buy/add/hold/watch`; if signal lookup fails, the risk endpoint fails open and the Web risk card shows a degraded state.
Regular stock history reports show `source_type=analysis` signals extracted from that report after the strategy block, using `source_report_id=<recordId>` as the query. Reports without a `recordId`, market reviews, and other non-regular stock reports do not issue this query. Empty results show a report-level empty state, and loading failures affect only that card, not the report body, news, diagnostics, or transparency sections.
## Backtesting

View File

@@ -222,6 +222,8 @@ P5 强化聚合报告通知路径的失败边界:`_send_notifications()` 在 r
- 任一静态渠道发送成功时P4 降噪 reservation 会写入正式记录;全部静态渠道失败或抛异常时,会释放 reservation。
- `send_to_context()` 仍独立于静态渠道 route 和降噪记录,用于回复触发任务的 Bot 会话上下文。
#1390 P6 的决策信号摘要沿用同一失败隔离边界:分析报告通知和告警通知只追加低敏 `decision_signal_summary` 摘要(动作、周期、理由、观察条件、风险和来源报告),不会输出 signal `metadata`、`evidence`、raw diagnostics 或 webhook/token。告警通知发送失败只记录通知尝试或 dispatch fallback不回滚已经写入的 trigger 或 DecisionSignal。
## 通知降噪机制
P4 新增进程内降噪,只影响静态配置渠道,不影响 `send_to_context()` 的机器人触发会话回执。默认所有配置关闭,未设置时保持旧行为。

View File

@@ -72,6 +72,7 @@ from src.services.run_diagnostics import (
sanitize_diagnostic_text,
)
from src.services.decision_signal_extractor import extract_and_persist_from_analysis_result
from src.services.decision_signal_summary import summarize_decision_signal
from src.enums import ReportType
from src.stock_analyzer import StockTrendAnalyzer, TrendAnalysisResult
from src.core.trading_calendar import (
@@ -2219,7 +2220,7 @@ class StockAnalysisPipeline:
or getattr(self, "trace_id", None)
or query_id
)
extract_and_persist_from_analysis_result(
signal_result = extract_and_persist_from_analysis_result(
result,
context_snapshot=context_snapshot,
source_report_id=source_report_id,
@@ -2228,6 +2229,10 @@ class StockAnalysisPipeline:
report_type=report_type,
portfolio_context=portfolio_context,
)
if isinstance(signal_result, dict):
summary = summarize_decision_signal(signal_result.get("item"))
if summary:
setattr(result, "decision_signal_summary", summary)
except Exception as exc:
logger.warning(
"Decision signal extraction skipped after history save: query_id=%s stock_code=%s error=%s",

View File

@@ -26,6 +26,7 @@ from enum import Enum
from src.config import Config, get_config
from src.enums import ReportType
from src.market_phase_summary import format_public_market_status_line, format_public_phase_pack_excerpt
from src.services.decision_signal_summary import format_decision_signal_excerpt
from src.notification_routing import (
get_notification_route_config,
split_notification_route_channels,
@@ -346,6 +347,12 @@ class NotificationService(
report_language=report_language,
)
def _decision_signal_excerpt(self, result: AnalysisResult, report_language: str) -> str:
return format_decision_signal_excerpt(
getattr(result, "decision_signal_summary", None),
report_language=report_language,
)
def _public_market_status_line(self, results: List[AnalysisResult], report_language: str) -> str:
for result in results or []:
line = format_public_market_status_line(
@@ -1631,6 +1638,9 @@ class NotificationService(
f"{localize_operation_advice(r.operation_advice, report_language)} | "
f"{labels['score_label']} {r.sentiment_score} | {one}"
)
signal_excerpt = self._decision_signal_excerpt(r, report_language)
if signal_excerpt:
lines.append(signal_excerpt)
lines.append("")
lines.append(f"*{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}*")
models = self._collect_models_used(results)
@@ -1673,6 +1683,10 @@ class NotificationService(
if excerpt:
lines.extend([excerpt, ""])
signal_excerpt = self._decision_signal_excerpt(result, report_language)
if signal_excerpt:
lines.extend([signal_excerpt, ""])
self._append_market_snapshot(lines, result)
# 核心决策(一句话)

View File

@@ -56,6 +56,7 @@ from src.services.market_light_alerts import (
normalize_market_alert_parameters,
)
from src.services.market_light_service import normalize_market_region
from src.services.decision_signal_summary import summarize_decision_signal
from src.analysis_context_pack_overview import (
ANALYSIS_CONTEXT_PACK_OVERVIEW_KEY,
extract_analysis_context_pack_overview,
@@ -1226,6 +1227,7 @@ class AlertService:
"market_phase_summary": visibility.get("market_phase_summary"),
"analysis_context_pack_overview": visibility.get("analysis_context_pack_overview"),
"analysis_visibility_source": visibility.get("analysis_visibility_source"),
"decision_signal_summary": visibility.get("decision_signal_summary"),
}
@staticmethod
@@ -1234,6 +1236,7 @@ class AlertService:
"market_phase_summary": None,
"analysis_context_pack_overview": None,
"analysis_visibility_source": None,
"decision_signal_summary": None,
}
if not diagnostics:
return result
@@ -1245,6 +1248,7 @@ class AlertService:
if not isinstance(parsed, dict):
result["analysis_visibility_source"] = "legacy_text"
return result
result["decision_signal_summary"] = summarize_decision_signal(parsed.get("decision_signal_summary"))
visibility = parsed.get("analysis_visibility")
if not isinstance(visibility, dict):
return result

View File

@@ -4,6 +4,7 @@
from __future__ import annotations
import asyncio
import hashlib
import json
import logging
import time
@@ -20,6 +21,7 @@ from src.agent.events import (
validate_event_alert_rule,
)
from data_provider.base import normalize_stock_code
from data_provider.us_index_mapping import is_us_index_code
from src.analysis_context_pack_overview import (
ANALYSIS_CONTEXT_PACK_OVERVIEW_KEY,
extract_analysis_context_pack_overview,
@@ -30,6 +32,11 @@ from src.market_phase_summary import (
render_market_phase_summary,
)
from src.services.alert_service import AlertService
from src.services.decision_signal_service import DecisionSignalService
from src.services.decision_signal_summary import (
format_decision_signal_excerpt,
summarize_decision_signal,
)
from src.services.history_service import HistoryService
from src.services.market_light_service import normalize_market_region
@@ -76,12 +83,14 @@ class AlertWorker:
*,
config_provider: Optional[Callable[[], Any]] = None,
service: Optional[AlertService] = None,
decision_signal_service: Optional[DecisionSignalService] = None,
notifier: Optional[Any] = None,
now_provider: Optional[Callable[[], float]] = None,
fingerprint_ttl_seconds: int = ALERT_WORKER_FINGERPRINT_TTL_SECONDS,
) -> None:
self.config_provider = config_provider or self._default_config_provider
self.service = service or AlertService()
self.decision_signal_service = decision_signal_service or DecisionSignalService()
self.notifier = notifier
self.now_provider = now_provider or time.time
self.fingerprint_ttl_seconds = max(1, int(fingerprint_ttl_seconds))
@@ -152,6 +161,8 @@ class AlertWorker:
}
record_status = result.get("record_status")
if record_status == "triggered":
self._attach_decision_signal_summary_safely(runtime_rule, result)
if record_status in WRITABLE_TRIGGER_STATUSES:
trigger_write = self._record_trigger_safely(runtime_rule, result, record_status)
trigger_id = trigger_write.trigger_id
@@ -361,6 +372,154 @@ class AlertWorker:
return dict(parsed) if isinstance(parsed, dict) else {"legacy_diagnostics": value}
return {}
def _attach_decision_signal_summary_safely(
self,
runtime_rule: RuntimeAlertRule,
result: Dict[str, Any],
) -> None:
try:
summary = self._resolve_decision_signal_summary(runtime_rule, result)
if not summary:
return
payload = self._diagnostics_payload(result.get("diagnostics"))
payload["decision_signal_summary"] = summary
result["diagnostics"] = payload
except Exception as exc:
logger.debug(
"[AlertWorker] decision signal summary unavailable for %s: %s",
self._display_target(runtime_rule),
self.service._sanitize_text(str(exc) or "decision signal summary failed"),
)
def _resolve_decision_signal_summary(
self,
runtime_rule: RuntimeAlertRule,
result: Dict[str, Any],
) -> Optional[Dict[str, Any]]:
identity = self._symbol_identity_for_decision_signal(runtime_rule)
if identity is None:
return None
stock_code, market = identity
latest = self.decision_signal_service.get_latest_active(
stock_code=stock_code,
market=market,
limit=1,
)
items = latest.get("items") if isinstance(latest, dict) else None
if items:
return summarize_decision_signal(items[0])
created = self.decision_signal_service.create_signal(
self._alert_decision_signal_payload(
runtime_rule,
result,
stock_code=stock_code,
market=market,
)
)
item = created.get("item") if isinstance(created, dict) else None
return summarize_decision_signal(item)
def _symbol_identity_for_decision_signal(self, runtime_rule: RuntimeAlertRule) -> Optional[Tuple[str, str]]:
rule = getattr(runtime_rule, "rule", runtime_rule)
metadata = getattr(rule, "metadata", None)
if not isinstance(metadata, dict):
metadata = {}
target_scope = str(
getattr(rule, "target_scope", None)
or metadata.get("target_scope")
or ""
).strip()
if target_scope in {"market", "portfolio_account"}:
return None
target = str(
metadata.get("effective_target")
or runtime_rule.effective_target
or getattr(rule, "stock_code", "")
or ""
).strip()
if not target or ":" in target:
return None
stock_code = normalize_stock_code(target)
if is_us_index_code(stock_code):
return None
market = get_market_for_stock(stock_code)
if market not in {"cn", "hk", "us"}:
return None
return stock_code, market
def _alert_decision_signal_payload(
self,
runtime_rule: RuntimeAlertRule,
result: Dict[str, Any],
*,
stock_code: str,
market: str,
) -> Dict[str, Any]:
rule = getattr(runtime_rule, "rule", runtime_rule)
metadata = getattr(rule, "metadata", None)
if not isinstance(metadata, dict):
metadata = {}
alert_type = self._public_alert_type(getattr(rule, "alert_type", None) or result.get("alert_type"))
key_hash = hashlib.sha1(str(runtime_rule.key or "").encode("utf-8")).hexdigest()
return {
"stock_code": stock_code,
"stock_name": getattr(rule, "stock_name", None),
"market": market,
"source_type": "alert",
"source_agent": "alert_worker",
"trace_id": f"alert-rule-{key_hash[:32]}",
"trigger_source": "alert",
"action": "alert",
"reason": result.get("reason") or result.get("message") or getattr(rule, "description", None),
"watch_conditions": self._alert_watch_conditions(runtime_rule, result, alert_type),
"risk_summary": self._alert_risk_summary(runtime_rule, result),
"metadata": {
"rule_id": self.service._runtime_rule_id(rule),
"alert_type": alert_type,
"severity": runtime_rule.severity,
"observed_value": result.get("observed_value"),
"threshold": result.get("threshold"),
"data_source": result.get("data_source"),
"data_timestamp": self._iso_or_text(result.get("data_timestamp")),
"rule_key_hash": key_hash,
},
}
@staticmethod
def _public_alert_type(value: Any) -> str:
raw = getattr(value, "value", value)
return str(raw or "").strip()[:64]
def _alert_watch_conditions(
self,
runtime_rule: RuntimeAlertRule,
result: Dict[str, Any],
alert_type: str,
) -> str:
threshold = result.get("threshold")
observed = result.get("observed_value")
target = self._display_target(runtime_rule)
parts = [part for part in (target, alert_type) if part]
if threshold not in (None, ""):
parts.append(f"threshold={threshold}")
if observed not in (None, ""):
parts.append(f"observed={observed}")
return " | ".join(str(part) for part in parts)
def _alert_risk_summary(self, runtime_rule: RuntimeAlertRule, result: Dict[str, Any]) -> str:
severity = str(runtime_rule.severity or "warning")
reason = result.get("reason") or result.get("message") or "Alert triggered"
return f"{severity}: {reason}"
@staticmethod
def _iso_or_text(value: Any) -> Optional[str]:
if value in (None, ""):
return None
if hasattr(value, "isoformat"):
return value.isoformat()
return str(value)
def _build_analysis_visibility(
self,
runtime_rule: RuntimeAlertRule,
@@ -492,6 +651,9 @@ class AlertWorker:
)
if excerpt:
content = f"{content}\n\n{excerpt}"
signal_excerpt = format_decision_signal_excerpt(diagnostics.get("decision_signal_summary"))
if signal_excerpt:
content = f"{content}\n\n{signal_excerpt}"
alert_text = NotificationBuilder.build_simple_alert(title=title, content=content, alert_type="warning")
return notification_service.send_with_results(alert_text, route_type="alert")

View File

@@ -483,6 +483,12 @@ class DecisionSignalService:
hk_normalized = cls._normalize_hk_stock_code(str(stock_code).strip())
return list(dict.fromkeys([normalized, hk_normalized]))
@classmethod
def normalize_stock_code_for_signal(cls, value: Any, *, market: Optional[str] = None) -> str:
"""Normalize a stock code for DecisionSignal identity matching."""
return cls._normalize_stock_code(value, market=market)
@classmethod
def _normalize_stock_code(cls, value: Any, *, market: Optional[str] = None) -> str:
raw = str(value or "").strip()

View File

@@ -0,0 +1,111 @@
# -*- coding: utf-8 -*-
"""Low-sensitive DecisionSignal summaries for notifications and risk views."""
from __future__ import annotations
from typing import Any, Dict, Optional
from src.utils.sanitize import sanitize_decision_signal_payload, sanitize_decision_signal_text
SUMMARY_FIELDS = (
"id",
"stock_code",
"stock_name",
"market",
"action",
"action_label",
"horizon",
"status",
"source_type",
"source_report_id",
"reason",
"watch_conditions",
"risk_summary",
"created_at",
"expires_at",
)
def summarize_decision_signal(item: Any) -> Optional[Dict[str, Any]]:
"""Return a low-sensitive summary from a serialized DecisionSignal item."""
if not isinstance(item, dict):
return None
summary: Dict[str, Any] = {}
for field_name in SUMMARY_FIELDS:
value = item.get(field_name)
if value in (None, "", [], {}):
continue
summary[field_name] = sanitize_decision_signal_payload(value)
return summary or None
def format_decision_signal_excerpt(summary: Any, report_language: str = "zh") -> str:
"""Format a compact public DecisionSignal excerpt for notification text."""
if not isinstance(summary, dict) or not summary:
return ""
language = "en" if str(report_language or "").lower().startswith("en") else "zh"
labels = {
"zh": {
"heading": "AI 决策信号",
"action": "动作",
"horizon": "周期",
"reason": "理由",
"watch_conditions": "观察条件",
"risk_summary": "风险",
"source_report_id": "报告",
},
"en": {
"heading": "AI decision signal",
"action": "Action",
"horizon": "Horizon",
"reason": "Reason",
"watch_conditions": "Watch",
"risk_summary": "Risk",
"source_report_id": "Report",
},
}[language]
parts = []
action_label = _public_scalar(summary.get("action_label") or summary.get("action"), max_length=32)
if action_label:
parts.append(f"{labels['action']}: {action_label}")
horizon = _public_scalar(summary.get("horizon"), max_length=16)
if horizon:
parts.append(f"{labels['horizon']}: {horizon}")
source_report_id = _public_scalar(summary.get("source_report_id"), max_length=24)
if source_report_id:
parts.append(f"{labels['source_report_id']}: #{source_report_id}")
lines = [f"**{labels['heading']}**"]
if parts:
lines.append(" | ".join(parts))
for key in ("reason", "watch_conditions", "risk_summary"):
text = _public_text(summary.get(key), max_length=120)
if text:
lines.append(f"- {labels[key]}: {text}")
return "\n".join(lines)
def _public_scalar(value: Any, *, max_length: int) -> str:
if value in (None, ""):
return ""
return sanitize_decision_signal_text(value)[:max_length]
def _public_text(value: Any, *, max_length: int) -> str:
if value in (None, "", [], {}):
return ""
if isinstance(value, (list, tuple)):
text = "".join(str(item).strip() for item in value if str(item or "").strip())
elif isinstance(value, dict):
text = "".join(
f"{key}: {item}"
for key, item in value.items()
if str(key or "").strip() and str(item or "").strip()
)
else:
text = str(value).strip()
return sanitize_decision_signal_text(text)[:max_length]

View File

@@ -3,13 +3,20 @@
from __future__ import annotations
import logging
from datetime import date, timedelta
from typing import Any, Dict, List, Optional, Tuple
from src.config import Config, get_config
from src.repositories.portfolio_repo import PortfolioRepository
from src.services.decision_signal_service import DecisionSignalService
from src.services.decision_signal_summary import summarize_decision_signal
from src.services.portfolio_service import PortfolioService
logger = logging.getLogger(__name__)
DEFENSIVE_DECISION_SIGNAL_ACTIONS = ("sell", "reduce", "alert")
class PortfolioRiskService:
"""Compute portfolio risk blocks on top of replayed snapshot data."""
@@ -19,10 +26,12 @@ class PortfolioRiskService:
*,
repo: Optional[PortfolioRepository] = None,
portfolio_service: Optional[PortfolioService] = None,
decision_signal_service: Optional[DecisionSignalService] = None,
config: Optional[Config] = None,
):
self.repo = repo or PortfolioRepository()
self.portfolio_service = portfolio_service or PortfolioService(repo=self.repo)
self.decision_signal_service = decision_signal_service or DecisionSignalService(portfolio_repo=self.repo)
self.config = config or get_config()
self._data_manager = None
self._data_manager_init_error = ""
@@ -73,6 +82,7 @@ class PortfolioRiskService:
lookback_days=thresholds["lookback_days"],
)
stop_loss = self._build_stop_loss(snapshot, thresholds)
decision_signal_risk = self._build_decision_signal_risk(snapshot, account_id=account_id)
return {
"as_of": as_of_date.isoformat(),
@@ -84,8 +94,108 @@ class PortfolioRiskService:
"sector_concentration": sector_concentration,
"drawdown": drawdown,
"stop_loss": stop_loss,
"decision_signal_risk": decision_signal_risk,
}
def _build_decision_signal_risk(
self,
snapshot: Dict[str, Any],
*,
account_id: Optional[int],
) -> Dict[str, Any]:
try:
held_positions = self._held_position_identities(snapshot)
if not held_positions:
return self._empty_decision_signal_risk(available=True)
defensive_actions = set(DEFENSIVE_DECISION_SIGNAL_ACTIONS)
latest_by_identity: Dict[Tuple[str, str], Dict[str, Any]] = {}
page = 1
while True:
response = self.decision_signal_service.list_signals(
holding_only=True,
account_id=account_id,
status="active",
page=page,
page_size=100,
)
items = response.get("items", []) if isinstance(response, dict) else []
for item in items:
if str(item.get("action") or "") not in defensive_actions:
continue
key = (
str(item.get("market") or "").strip().lower(),
str(item.get("stock_code") or "").strip().upper(),
)
if key[0] and key[1] and key not in latest_by_identity:
latest_by_identity[key] = item
total = int(response.get("total", 0) or 0) if isinstance(response, dict) else 0
if page * 100 >= total or not items:
break
page += 1
risk_items: List[Dict[str, Any]] = []
action_counts = {action: 0 for action in DEFENSIVE_DECISION_SIGNAL_ACTIONS}
seen: set[Tuple[Optional[int], str, str, int]] = set()
for position in held_positions:
signal = latest_by_identity.get((position["market"], position["signal_stock_code"]))
summary = summarize_decision_signal(signal)
if not summary:
continue
action = str(summary.get("action") or "")
if action not in action_counts:
continue
signal_id = int(summary.get("id") or 0)
dedupe_key = (position["account_id"], position["market"], position["signal_stock_code"], signal_id)
if dedupe_key in seen:
continue
seen.add(dedupe_key)
action_counts[action] += 1
risk_items.append({
"account_id": position["account_id"],
"symbol": position["symbol"],
"market": position["market"],
"signal": summary,
})
return {
"available": True,
"total": len(risk_items),
"actions": action_counts,
"items": risk_items,
}
except Exception:
logger.exception("[PortfolioRiskService] Decision signal risk unavailable")
return self._empty_decision_signal_risk(available=False)
@staticmethod
def _empty_decision_signal_risk(*, available: bool) -> Dict[str, Any]:
return {
"available": available,
"total": 0,
"actions": {action: 0 for action in DEFENSIVE_DECISION_SIGNAL_ACTIONS},
"items": [],
}
@staticmethod
def _held_position_identities(snapshot: Dict[str, Any]) -> List[Dict[str, Any]]:
positions: List[Dict[str, Any]] = []
for account in snapshot.get("accounts", []) or []:
account_id = account.get("account_id")
for pos in account.get("positions", []) or []:
symbol = str(pos.get("symbol") or "").strip().upper()
market = str(pos.get("market") or "").strip().lower()
if not symbol or market not in {"cn", "hk", "us"}:
continue
signal_stock_code = DecisionSignalService.normalize_stock_code_for_signal(symbol, market=market)
positions.append({
"account_id": account_id,
"symbol": symbol,
"market": market,
"signal_stock_code": signal_stock_code,
})
return positions
def _ensure_drawdown_snapshot_window(
self,
*,

View File

@@ -16,6 +16,7 @@ from typing import Any, Dict, List, Optional
from src.analyzer import AnalysisResult
from src.config import get_config
from src.market_phase_summary import format_public_market_status_line, format_public_phase_pack_excerpt
from src.services.decision_signal_summary import format_decision_signal_excerpt
from src.report_language import (
get_localized_stock_name,
get_report_labels,
@@ -157,6 +158,12 @@ def render(
report_language=report_language,
)
def decision_signal_excerpt(result: AnalysisResult) -> str:
return format_decision_signal_excerpt(
getattr(result, "decision_signal_summary", None),
report_language=report_language,
)
def market_status_line() -> str:
for source_results in (results or [], sorted_results):
for result in source_results:
@@ -186,6 +193,7 @@ def render(
"clean_sniper": _clean_sniper_value,
"failed_checks": failed_checks,
"phase_pack_excerpt": phase_pack_excerpt,
"decision_signal_excerpt": decision_signal_excerpt,
"history_by_code": {},
"get_chip_unavailable_reason": get_chip_unavailable_reason,
"is_chip_structure_unavailable": is_chip_structure_unavailable,

View File

@@ -10,6 +10,10 @@
{% set core = (dash.get('core_conclusion') or {}) if dash else {} %}
{% set one = (core.get('one_sentence') or e.result.analysis_summary or '')[:60] %}
**{{ e.stock_name }}({{ e.result.code }})** {{ e.signal_emoji }} {{ e.localized_operation_advice }} | {{ labels.score_label }} {{ e.result.sentiment_score }} | {{ one }}
{% set signal_excerpt = decision_signal_excerpt(e.result) %}
{% if signal_excerpt %}
{{ signal_excerpt }}
{% endif %}
{% endfor %}
*{{ report_timestamp }}*

View File

@@ -25,6 +25,7 @@ from src.services.alert_indicators import (
)
from src.services.alert_service import AlertService
from src.services.alert_worker import AlertWorker
from src.services.decision_signal_service import DecisionSignalService
from src.storage import DatabaseManager
@@ -315,6 +316,181 @@ class AlertWorkerTestCase(unittest.TestCase):
notifier.send_with_results.side_effect = list(results)
return notifier
def test_p6_triggered_stock_alert_links_latest_active_decision_signal(self) -> None:
self._create_rule(target="600519")
signal_service = DecisionSignalService()
signal_service.create_signal({
"stock_code": "600519",
"stock_name": "贵州茅台",
"market": "cn",
"source_type": "analysis",
"source_report_id": 1390,
"trace_id": "analysis-1390",
"trigger_source": "api",
"action": "sell",
"reason": "跌破关键支撑",
"watch_conditions": "观察能否收回均线",
"risk_summary": "下行风险扩大",
})
notifier = self._notifier()
worker = AlertWorker(
config_provider=lambda: self._config(),
service=self.service,
decision_signal_service=signal_service,
notifier=notifier,
)
with patch(
"src.agent.events.EventMonitor._get_realtime_quote",
new=AsyncMock(return_value=SimpleNamespace(price=1810.0)),
):
stats = worker.run_once()
self.assertEqual(stats["triggered"], 1)
all_signals = signal_service.list_signals(stock_code="600519", market="cn", page_size=10)["items"]
self.assertEqual(len(all_signals), 1)
self.assertEqual(all_signals[0]["source_type"], "analysis")
trigger = self._triggers(status="triggered")[0]
summary = trigger["decision_signal_summary"]
self.assertEqual(summary["id"], all_signals[0]["id"])
self.assertEqual(summary["action"], "sell")
alert_text = notifier.send_with_results.call_args.args[0]
self.assertIn("AI 决策信号", alert_text)
self.assertIn("跌破关键支撑", alert_text)
self.assertIn("观察能否收回均线", alert_text)
def test_p6_triggered_stock_alert_creates_alert_signal_when_no_active_signal(self) -> None:
self._create_rule(target="600519")
signal_service = DecisionSignalService()
notifier = self._notifier()
worker = AlertWorker(
config_provider=lambda: self._config(),
service=self.service,
decision_signal_service=signal_service,
notifier=notifier,
)
with patch(
"src.agent.events.EventMonitor._get_realtime_quote",
new=AsyncMock(return_value=SimpleNamespace(price=1810.0)),
):
stats = worker.run_once()
self.assertEqual(stats["triggered"], 1)
signals = signal_service.list_signals(source_type="alert", stock_code="600519", market="cn", page_size=10)[
"items"
]
self.assertEqual(len(signals), 1)
item = signals[0]
self.assertEqual(item["action"], "alert")
self.assertEqual(item["trigger_source"], "alert")
self.assertEqual(item["source_agent"], "alert_worker")
self.assertIsNone(item["market_phase"])
self.assertTrue(str(item["trace_id"]).startswith("alert-rule-"))
self.assertEqual(item["metadata"]["rule_id"], 1)
self.assertEqual(item["metadata"]["alert_type"], "price_cross")
self.assertEqual(self._triggers(status="triggered")[0]["decision_signal_summary"]["id"], item["id"])
def test_p6_alert_signal_trace_id_is_idempotent_for_same_rule(self) -> None:
self._create_rule(target="600519")
signal_service = DecisionSignalService()
notifier = self._notifier(
self._dispatch_result(False, dispatched=True),
self._dispatch_result(False, dispatched=True),
)
worker = AlertWorker(
config_provider=lambda: self._config(),
service=self.service,
decision_signal_service=signal_service,
notifier=notifier,
)
with patch(
"src.agent.events.EventMonitor._get_realtime_quote",
new=AsyncMock(return_value=SimpleNamespace(price=1810.0)),
):
worker.run_once()
worker.run_once()
signals = signal_service.list_signals(source_type="alert", stock_code="600519", market="cn", page_size=10)[
"items"
]
self.assertEqual(len(signals), 1)
self.assertIsNone(signals[0]["market_phase"])
self.assertEqual(len(self._triggers(status="triggered")), 2)
def test_p6_non_stock_alert_targets_skip_decision_signal_write(self) -> None:
signal_service = MagicMock()
worker = AlertWorker(
config_provider=lambda: self._config(),
service=self.service,
decision_signal_service=signal_service,
)
for runtime_rule in (
SimpleNamespace(
key="market:cn",
rule=SimpleNamespace(target_scope="market", target="cn", metadata={}),
source="db",
severity="warning",
effective_target="cn",
display_target="cn",
),
SimpleNamespace(
key="portfolio_account:all",
rule=SimpleNamespace(target_scope="portfolio_account", target="all", metadata={}),
source="db",
severity="warning",
effective_target="portfolio_account:all",
display_target="all accounts",
),
SimpleNamespace(
key="single_symbol:SPX",
rule=SimpleNamespace(
target_scope="single_symbol",
stock_code="SPX",
metadata={},
),
source="db",
severity="warning",
effective_target="SPX",
display_target="SPX",
),
):
worker._attach_decision_signal_summary_safely(runtime_rule, {"record_status": "triggered"})
signal_service.get_latest_active.assert_not_called()
signal_service.create_signal.assert_not_called()
def test_p6_notification_failure_does_not_block_alert_signal_write(self) -> None:
self._create_rule(target="600519")
signal_service = DecisionSignalService()
class FailingNotifier:
def send_with_results(self, *_args, **_kwargs):
raise RuntimeError("webhook token=secret failed")
worker = AlertWorker(
config_provider=lambda: self._config(),
service=self.service,
decision_signal_service=signal_service,
notifier=FailingNotifier(),
)
with patch(
"src.agent.events.EventMonitor._get_realtime_quote",
new=AsyncMock(return_value=SimpleNamespace(price=1810.0)),
):
stats = worker.run_once()
self.assertEqual(stats["triggered"], 1)
self.assertEqual(stats["recorded"], 1)
self.assertEqual(stats["notified"], 0)
self.assertEqual(len(self._triggers(status="triggered")), 1)
signals = signal_service.list_signals(source_type="alert", stock_code="600519", market="cn", page_size=10)[
"items"
]
self.assertEqual(len(signals), 1)
def test_triggered_diagnostics_merge_visibility_and_market_scope_uses_region(self) -> None:
worker = AlertWorker(config_provider=lambda: self._config(), service=self.service)
runtime_rule = SimpleNamespace(

View File

@@ -2,6 +2,7 @@
import json
from pathlib import Path
from typing import Any
from api.app import create_app
from api.v1.router import router as api_v1_router
@@ -37,6 +38,31 @@ DECISION_SIGNAL_SCHEMAS = (
"DecisionSignalOutcomeStatsResponse",
"DecisionSignalStatusUpdateRequest",
)
P6_SIGNAL_LINKED_PATHS = (
"/api/v1/alerts/triggers",
"/api/v1/portfolio/risk",
)
P6_SIGNAL_LINKED_SCHEMAS = (
"AlertTriggerItem",
"AlertTriggerListResponse",
"PortfolioDecisionSignalRiskBlock",
"PortfolioDecisionSignalRiskItem",
"PortfolioRiskResponse",
)
def _collect_component_schema_refs(node: Any) -> set[str]:
refs: set[str] = set()
if isinstance(node, dict):
ref = node.get("$ref")
if isinstance(ref, str) and ref.startswith("#/components/schemas/"):
refs.add(ref.rsplit("/", 1)[-1])
for value in node.values():
refs.update(_collect_component_schema_refs(value))
elif isinstance(node, list):
for value in node:
refs.update(_collect_component_schema_refs(value))
return refs
def test_schema_examples_remain_in_openapi_schema() -> None:
@@ -119,6 +145,14 @@ def test_decision_signal_static_api_spec_matches_runtime_paths() -> None:
for schema_name in DECISION_SIGNAL_SCHEMAS:
assert static_spec["components"]["schemas"][schema_name] == runtime_spec["components"]["schemas"][schema_name]
for path in P6_SIGNAL_LINKED_PATHS:
assert static_spec["paths"][path] == runtime_spec["paths"][path]
for schema_name in P6_SIGNAL_LINKED_SCHEMAS:
assert static_spec["components"]["schemas"][schema_name] == runtime_spec["components"]["schemas"][schema_name]
schema_refs = _collect_component_schema_refs(static_spec)
missing_schema_refs = sorted(schema_refs - set(static_spec["components"]["schemas"]))
assert missing_schema_refs == []
status_schema = static_spec["components"]["schemas"]["DecisionSignalStatusUpdateRequest"]["properties"]["status"]
assert status_schema["enum"] == ["active", "expired", "invalidated", "closed", "archived"]

View File

@@ -173,6 +173,13 @@ def test_build_payload_records_empty_holding_state_from_explicit_portfolio_conte
assert payload["metadata"]["holding_state"] == "empty"
def test_runtime_decision_signal_summary_is_not_serialized_by_analysis_result_to_dict() -> None:
result = _result()
setattr(result, "decision_signal_summary", {"action": "sell", "reason": "risk"})
assert "decision_signal_summary" not in result.to_dict()
def test_build_payload_maps_secondary_only_entry_to_entry_high() -> None:
result = _result()
result.dashboard = {

View File

@@ -0,0 +1,101 @@
# -*- coding: utf-8 -*-
"""Tests for low-sensitive DecisionSignal summary helpers."""
from __future__ import annotations
from src.services.decision_signal_summary import (
format_decision_signal_excerpt,
summarize_decision_signal,
)
def test_summarize_decision_signal_keeps_only_low_sensitive_fields() -> None:
summary = summarize_decision_signal({
"id": 42,
"stock_code": "600519",
"stock_name": "贵州茅台",
"market": "cn",
"action": "sell",
"action_label": "卖出",
"horizon": "3d",
"status": "active",
"source_type": "alert",
"source_agent": "alert_worker",
"source_report_id": 88,
"reason": "token=secret-value 触发止损",
"watch_conditions": ["观察量能", "password=hidden"],
"risk_summary": {"drawdown": "webhook=https://hooks.slack.com/services/T/B/C"},
"created_at": "2026-06-18T10:00:00+08:00",
"expires_at": "2026-06-25T10:00:00+08:00",
"metadata": {"webhook_url": "https://hooks.slack.com/services/T/B/C"},
"evidence": {"secret": "raw"},
"diagnostics": "authorization=Bearer raw",
})
assert summary is not None
assert set(summary) == {
"id",
"stock_code",
"stock_name",
"market",
"action",
"action_label",
"horizon",
"status",
"source_type",
"source_report_id",
"reason",
"watch_conditions",
"risk_summary",
"created_at",
"expires_at",
}
assert summary["reason"] == "token=[REDACTED] 触发止损"
assert summary["watch_conditions"] == ["观察量能", "password=[REDACTED]"]
assert summary["risk_summary"] == {"drawdown": "webhook=[REDACTED_URL]"}
def test_summarize_decision_signal_rejects_non_dict_and_empty_payload() -> None:
assert summarize_decision_signal(None) is None
assert summarize_decision_signal(["not", "a", "dict"]) is None
assert summarize_decision_signal({"metadata": {"token": "secret"}, "evidence": {"raw": True}}) is None
assert summarize_decision_signal({"stock_code": "", "reason": None}) is None
def test_format_decision_signal_excerpt_formats_chinese_list_and_dict_fields() -> None:
excerpt = format_decision_signal_excerpt({
"action_label": "卖出",
"horizon": "3d",
"source_report_id": 88,
"reason": "跌破止损线",
"watch_conditions": ["观察 1660 支撑", "等待成交量收缩"],
"risk_summary": {"drawdown": "组合回撤扩大"},
})
assert excerpt.startswith("**AI 决策信号**")
assert "动作: 卖出 | 周期: 3d | 报告: #88" in excerpt
assert "- 理由: 跌破止损线" in excerpt
assert "- 观察条件: 观察 1660 支撑;等待成交量收缩" in excerpt
assert "- 风险: drawdown: 组合回撤扩大" in excerpt
def test_format_decision_signal_excerpt_formats_english_and_redacts_text() -> None:
excerpt = format_decision_signal_excerpt({
"action": "alert",
"horizon": "5d",
"reason": "authorization: Bearer raw-token",
"watch_conditions": "Check price",
"risk_summary": "token=hidden",
}, report_language="en")
assert excerpt.startswith("**AI decision signal**")
assert "Action: alert | Horizon: 5d" in excerpt
assert "- Reason: authorization: [REDACTED]" in excerpt
assert "- Watch: Check price" in excerpt
assert "- Risk: token=[REDACTED]" in excerpt
def test_format_decision_signal_excerpt_returns_empty_for_invalid_input() -> None:
assert format_decision_signal_excerpt(None) == ""
assert format_decision_signal_excerpt({}) == ""
assert format_decision_signal_excerpt(["not", "a", "dict"]) == ""

View File

@@ -23,6 +23,7 @@ except ModuleNotFoundError:
import src.auth as auth
from api.app import create_app
from src.config import Config
from src.services.decision_signal_service import DecisionSignalService
from src.services.portfolio_import_service import PortfolioImportService
from src.services.portfolio_risk_service import PortfolioRiskService
from src.services.portfolio_service import PortfolioBusyError, PortfolioService
@@ -102,6 +103,34 @@ class PortfolioPr2TestCase(unittest.TestCase):
)
self.db.save_daily_data(df, code=symbol, data_source="portfolio-pr2-test")
def _create_position(self, account_id: int, symbol: str, price: float = 100.0, *, market: str = "cn") -> None:
self.service.record_trade(
account_id=account_id,
symbol=symbol,
trade_date=date(2026, 1, 1),
side="buy",
quantity=10,
price=price,
market=market,
currency="CNY" if market == "cn" else "USD",
)
self._save_close(symbol, date(2026, 1, 1), price)
def _create_signal(self, symbol: str, action: str, **overrides) -> dict:
payload = {
"stock_code": symbol,
"stock_name": symbol,
"market": "cn",
"source_type": "manual",
"trace_id": f"signal-{symbol}-{action}-{len(overrides)}",
"trigger_source": "api",
"action": action,
"reason": f"{symbol} {action} reason",
"status": "active",
}
payload.update(overrides)
return DecisionSignalService().create_signal(payload)["item"]
@staticmethod
def _csv_bytes(with_trade_uid: bool = True) -> bytes:
if with_trade_uid:
@@ -495,6 +524,93 @@ class PortfolioPr2TestCase(unittest.TestCase):
self.assertTrue(len(sectors) >= 1)
self.assertEqual(sectors[0]["sector"], "白酒")
def test_risk_report_aggregates_active_defensive_decision_signals_for_holdings(self) -> None:
account = self.service.create_account(name="Main", broker="Demo", market="cn", base_currency="CNY")
aid = account["id"]
self.service.record_cash_ledger(
account_id=aid,
event_date=date(2026, 1, 1),
direction="in",
amount=100000.0,
currency="CNY",
)
for symbol in ("SH600519", "300750", "000001", "000002", "000003"):
self._create_position(aid, symbol)
self._create_signal("600519.SH", "sell")
self._create_signal("300750", "reduce")
self._create_signal("000001", "alert")
self._create_signal("000002", "buy")
self._create_signal("000003", "watch")
self._create_signal("002000", "sell")
self._create_signal("600519", "alert", status="expired", trace_id="expired-alert-600519")
report = self.risk_service.get_risk_report(account_id=aid, as_of=date(2026, 1, 1), cost_method="fifo")
block = report["decision_signal_risk"]
self.assertTrue(block["available"])
self.assertEqual(block["total"], 3)
self.assertEqual(block["actions"], {"sell": 1, "reduce": 1, "alert": 1})
symbols = {item["symbol"] for item in block["items"]}
self.assertEqual(symbols, {"SH600519", "300750", "000001"})
signal_actions = {item["symbol"]: item["signal"]["action"] for item in block["items"]}
self.assertNotIn("000002", signal_actions)
self.assertNotIn("000003", signal_actions)
self.assertNotIn("002000", signal_actions)
self.assertEqual(signal_actions["SH600519"], "sell")
self.assertEqual(signal_actions["300750"], "reduce")
self.assertEqual(signal_actions["000001"], "alert")
def test_risk_report_decision_signal_fail_open(self) -> None:
account = self.service.create_account(name="Main", broker="Demo", market="cn", base_currency="CNY")
aid = account["id"]
self.service.record_cash_ledger(
account_id=aid,
event_date=date(2026, 1, 1),
direction="in",
amount=10000.0,
currency="CNY",
)
self._create_position(aid, "600519")
class BrokenSignalService:
def list_signals(self, **_kwargs):
raise RuntimeError("decision signal unavailable")
risk_service = PortfolioRiskService(
portfolio_service=self.service,
decision_signal_service=BrokenSignalService(),
)
report = risk_service.get_risk_report(account_id=aid, as_of=date(2026, 1, 1), cost_method="fifo")
self.assertFalse(report["decision_signal_risk"]["available"])
self.assertEqual(report["decision_signal_risk"]["total"], 0)
self.assertEqual(report["decision_signal_risk"]["items"], [])
def test_portfolio_risk_endpoint_returns_defensive_decision_signals(self) -> None:
account = self.service.create_account(name="Main", broker="Demo", market="cn", base_currency="CNY")
aid = account["id"]
self.service.record_cash_ledger(
account_id=aid,
event_date=date(2026, 1, 1),
direction="in",
amount=10000.0,
currency="CNY",
)
self._create_position(aid, "600519")
self._create_signal("600519", "sell")
response = self.client.get(
"/api/v1/portfolio/risk",
params={"account_id": aid, "as_of": "2026-01-01", "cost_method": "fifo"},
)
self.assertEqual(response.status_code, 200)
payload = response.json()
self.assertIn("decision_signal_risk", payload)
self.assertEqual(payload["decision_signal_risk"]["total"], 1)
self.assertEqual(payload["decision_signal_risk"]["items"][0]["signal"]["action"], "sell")
def test_snapshot_does_not_trigger_online_fx_refresh(self) -> None:
account = self.service.create_account(name="US", broker="Demo", market="us", base_currency="CNY")
aid = account["id"]