feat: adapt AlphaSift hotspot topics (#1635)

* feat: adapt AlphaSift hotspot topics

* fix(review-feedback-1635): preserve cached hotspots when refresh returns no rows and Avoid

* fix(review-feedback-1635): Reload detail when refreshing the selected hotspot and Fall back when

* fix(review-feedback-1635): update by the requested topic or a request token

* fix(review-feedback-1635): preserve analysis when optional serve startup fails and Fall back when

* fix(review-feedback-1635): 当前 CI 状态为 failure,且 backend-gate 为 cancelled、docker-build 为 skipped

* fix(review-feedback-1635): 让 backend-gate / docker-build 在当前 head 上完成并通过,或给出可接受的维护者级豁免说明

* fix(review-feedback-1635): 当前 CI 仍未闭环:结构化事实显示 backend-gate 被取消、docker-build skipped

* fix(review-feedback-1635): 修复 hotspot detail 的上游 schema 容错问题,并让当前 head 的阻断 CI 重新跑通

* fix(review-feedback-1635): docs/CHANGELOG.md:当前 diff 不是在 [Unreleased] 扁平格式下追加本 PR

* fix(review-feedback-1635): CI 闭环未完成:结构化事实显示当前 CI 状态为 failure,backend-gate cancelled、docker-build

* fix(review-feedback-1635): 让当前 head 的阻断型 CI 完整通过,并修正文档中的 AlphaSift commit hash 不一致

* fix(review-feedback-1635): CI 仍未在当前 head 完整通过:结构化事实显示

* fix(review-feedback-1635): 修复 docs/CHANGELOG.md 对主线既有 [Unreleased] 条目的删除,并补齐 AlphaSift

* fix(review-feedback-1635): Accept slash-containing hotspot topics

* fix(review-feedback-1635): Try industry constituents for industry hotspots

* fix(review-feedback-1635): Clear old hotspot details when loading a new topic and preserve the

* test: isolate main schedule environment

* chore: clean AlphaSift PR scope

* fix: respect AlphaSift hotspot fallbacks
This commit is contained in:
mumu
2026-06-12 22:38:42 +08:00
committed by GitHub
parent efa41e0ada
commit 4b3f679b2d
17 changed files with 2185 additions and 16 deletions

View File

@@ -31,9 +31,18 @@ ALPHASIFT_ENABLED=false
# - ALPHASIFT_ENABLED 只影响 AlphaSift 选股流程,不会改写/迁移/清理现有 LITELLM_*/LLM_* 配置。
# - 关闭/回退只需将 ALPHASIFT_ENABLED 置 false现有 provider/base_url/custom headers 与 fallback 语义保持原状。
# - 兼容依据运维核验LiteLLM providers 与 OpenAI-compatible 语义说明见下方注释;运行时默认使用 requirements.txt 中固定 litellm 版本。
ALPHASIFT_INSTALL_SPEC=git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344
ALPHASIFT_INSTALL_SPEC=git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386
# AlphaSift 全市场快照源优先级;未配置时 DSA 调用 AlphaSift 会优先使用更稳的东方财富数据中心源。
# SNAPSHOT_SOURCE_PRIORITY=em_datacenter,tushare,efinance,akshare_em
# AlphaSift 最新版支持 last-good 快照、日线历史和行业/概念 provider 缓存DSA 运行时默认使用 data/alphasift 下的隔离缓存目录。
# 如需自定义,可覆盖以下路径;留空不会改写既有 DSA 配置。
# ALPHASIFT_DATA_DIR=data/alphasift
# ALPHASIFT_FALLBACK_SNAPSHOT_PATH=data/alphasift/snapshot.last_good.json
# ALPHASIFT_DAILY_HISTORY_CACHE_DIR=data/alphasift/daily_history
# ALPHASIFT_INDUSTRY_PROVIDER_CACHE_DIR=data/alphasift/industry_provider_cache
# AlphaSift 行业/概念 provider 默认关闭;开启 akshare 后可为 theme_heat、行业热度和 LLM 上下文提供更多板块信息。
# INDUSTRY_PROVIDER=none
# INDUSTRY_PROVIDER_MAX_BOARDS=80
# 单次 LLM 请求超时秒数AlphaSift 选股会复用 DSA 的该配置,超时后降级返回非 LLM 排序结果。
# LLM_TIMEOUT_SEC=60

View File

@@ -6,7 +6,7 @@ from __future__ import annotations
import uuid
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, HTTPException, Request
from fastapi import APIRouter, Depends, HTTPException, Query, Request
from pydantic import BaseModel, Field
from api.deps import get_config_dep
@@ -82,6 +82,25 @@ def alphasift_strategies(
return _service(config).strategies()
@router.get("/hotspots")
def alphasift_hotspots(
provider: str = Query("", max_length=32),
top: int = Query(12, ge=1, le=50),
refresh: bool = Query(False),
config: Config = Depends(get_config_dep),
) -> Dict[str, Any]:
return _service(config).hotspots(provider=provider, top=top, refresh=refresh)
@router.get("/hotspots/{topic:path}")
def alphasift_hotspot_detail(
topic: str,
provider: str = Query("", max_length=32),
config: Config = Depends(get_config_dep),
) -> Dict[str, Any]:
return _service(config).hotspot_detail(topic=topic, provider=provider)
@router.post("/install")
def alphasift_install(
request: Request,

View File

@@ -129,6 +129,60 @@ describe('alphasiftApi', () => {
expect(result.strategies[0].marketScope).toEqual(['cn']);
});
it('loads hotspot themes from the AlphaSift API', async () => {
get.mockResolvedValueOnce({
data: {
enabled: true,
provider: 'akshare',
provider_used: 'akshare',
hotspots: [
{
topic: 'AI算力',
heat_score: 88,
trend_score: 12,
sample_stock_count: 8,
leaders: ['中际旭创'],
},
],
hotspot_count: 1,
},
});
const result = await alphasiftApi.getHotspots({ provider: 'akshare', top: 12, refresh: true });
expect(get).toHaveBeenCalledWith('/api/v1/alphasift/hotspots', {
params: { provider: 'akshare', top: 12, refresh: true },
timeout: 300000,
});
expect(result.providerUsed).toBe('akshare');
expect(result.hotspots[0].heatScore).toBe(88);
expect(result.hotspots[0].sampleStockCount).toBe(8);
});
it('loads hotspot detail for a concrete topic', async () => {
get.mockResolvedValueOnce({
data: {
enabled: true,
provider: 'akshare',
topic: '玻璃基板',
summary: '玻璃基板盘中发酵',
route: [{ title: '盘中发酵', description: '出现大笔买入' }],
stocks: [{ code: '920438', name: '戈碧迦', role: '异动核心' }],
stock_count: 1,
},
});
const result = await alphasiftApi.getHotspotDetail({ topic: '玻璃基板', provider: 'akshare' });
expect(get).toHaveBeenCalledWith('/api/v1/alphasift/hotspots/%E7%8E%BB%E7%92%83%E5%9F%BA%E6%9D%BF', {
params: { provider: 'akshare' },
timeout: 300000,
});
expect(result.topic).toBe('玻璃基板');
expect(result.stockCount).toBe(1);
expect(result.stocks[0].name).toBe('戈碧迦');
});
it('uses a long timeout for LLM-backed screening', async () => {
post.mockResolvedValueOnce({
data: {
@@ -183,6 +237,9 @@ describe('alphasiftApi', () => {
enabled: true,
candidates: [],
candidate_count: 0,
daily_enriched: true,
daily_enrich_count: 4,
post_analyzers: ['scorecard'],
},
},
});
@@ -192,5 +249,8 @@ describe('alphasiftApi', () => {
expect(get).toHaveBeenCalledWith('/api/v1/alphasift/screen/tasks/screen-task-1');
expect(result.taskId).toBe('screen-task-1');
expect(result.result?.candidateCount).toBe(0);
expect(result.result?.dailyEnriched).toBe(true);
expect(result.result?.dailyEnrichCount).toBe(4);
expect(result.result?.postAnalyzers).toEqual(['scorecard']);
});
});

View File

@@ -93,6 +93,73 @@ export type AlphaSiftStrategiesResponse = {
strategyCount: number;
};
export type AlphaSiftHotspot = {
topic: string;
name?: string;
source?: string;
rank?: number | null;
changePct?: number | null;
heatScore?: number | null;
trendScore?: number | null;
persistenceScore?: number | null;
coolingScore?: number | null;
observations?: number | null;
state?: string;
stage?: string;
sampleStockCount?: number | null;
leaders?: string[];
providerUsed?: string;
fallbackUsed?: boolean;
cacheUsed?: boolean;
cachedAt?: string | null;
sourceErrors?: string[];
stale?: boolean;
staleAgeHours?: number | null;
};
export type AlphaSiftHotspotRouteItem = {
title: string;
description: string;
source?: string;
};
export type AlphaSiftHotspotStock = {
code?: string;
name?: string;
changePct?: number | null;
amount?: number | null;
turnoverRate?: number | null;
volumeRatio?: number | null;
role?: string;
hotStockScore?: number | null;
};
export type AlphaSiftHotspotDetail = {
enabled: boolean;
provider: string;
topic: string;
name?: string;
summary?: string;
route: AlphaSiftHotspotRouteItem[];
stocks: AlphaSiftHotspotStock[];
stockCount: number;
sourceErrors?: string[];
};
export type AlphaSiftHotspotsResponse = {
enabled: boolean;
provider: string;
providerUsed?: string;
fallbackUsed?: boolean;
cacheUsed?: boolean;
cachedAt?: string | null;
sourceErrors?: string[];
stale?: boolean;
staleAgeHours?: number | null;
hotspots: AlphaSiftHotspot[];
hotspotCount: number;
};
export type AlphaSiftScreenResponse = {
enabled: boolean;
candidates: AlphaSiftCandidate[];
@@ -117,6 +184,13 @@ export type AlphaSiftScreenResponse = {
enrichedCount?: number;
warnings?: string[];
};
deepAnalysisRequested?: boolean | null;
postAnalyzers?: string[];
dailyEnriched?: boolean | null;
dailyEnrichCount?: number | null;
riskEnabled?: boolean | null;
portfolioDiversityEnabled?: boolean | null;
portfolioConcentrationNotes?: string[];
};
export type AlphaSiftScreenAccepted = {
@@ -193,6 +267,29 @@ export const alphasiftApi = {
return toCamelCase<AlphaSiftStrategiesResponse>(response.data);
},
async getHotspots(payload: { provider?: string; top?: number; refresh?: boolean } = {}): Promise<AlphaSiftHotspotsResponse> {
const response = await apiClient.get<Record<string, unknown>>('/api/v1/alphasift/hotspots', {
params: {
provider: payload.provider || 'akshare',
top: payload.top ?? 12,
refresh: payload.refresh ?? false,
},
timeout: ALPHASIFT_INSTALL_TIMEOUT_MS,
});
return toCamelCase<AlphaSiftHotspotsResponse>(response.data);
},
async getHotspotDetail(payload: { topic: string; provider?: string }): Promise<AlphaSiftHotspotDetail> {
const response = await apiClient.get<Record<string, unknown>>(
`/api/v1/alphasift/hotspots/${encodeURIComponent(payload.topic)}`,
{
params: { provider: payload.provider || 'akshare' },
timeout: ALPHASIFT_INSTALL_TIMEOUT_MS,
},
);
return toCamelCase<AlphaSiftHotspotDetail>(response.data);
},
async install(): Promise<AlphaSiftInstallResponse> {
const response = await apiClient.post<Record<string, unknown>>('/api/v1/alphasift/install', {}, { timeout: ALPHASIFT_INSTALL_TIMEOUT_MS });
return toCamelCase<AlphaSiftInstallResponse>(response.data);

View File

@@ -376,7 +376,7 @@ export function parseApiError(error: unknown): ParsedApiError {
if (errorCode === 'alphasift_unavailable' || includesAny(matchText, ['cannot import alphasift', 'alphasift.screen'])) {
return createParsedApiError({
title: 'AlphaSift 未就绪',
message: '当前 DSA 后端环境无法导入 alphasift。请先执行 pip install -r requirements.txt或重建 Docker/桌面后端产物。',
message: rawMessage,
rawMessage,
status,
category: 'http_error',

View File

@@ -1,9 +1,11 @@
import type React from 'react';
import { Fragment, useCallback, useEffect, useMemo, useState } from 'react';
import { CheckCircle2, CircleAlert, Play, PlusCircle, Search, SlidersHorizontal } from 'lucide-react';
import { Fragment, useCallback, useEffect, useMemo, useRef, useState } from 'react';
import { CheckCircle2, CircleAlert, Flame, Play, PlusCircle, RefreshCw, Search, SlidersHorizontal } from 'lucide-react';
import {
alphasiftApi,
type AlphaSiftCandidate,
type AlphaSiftHotspotDetail,
type AlphaSiftHotspot,
type AlphaSiftScreenResponse,
type AlphaSiftScreenTaskStatus,
type AlphaSiftStrategy,
@@ -273,6 +275,15 @@ const StockScreeningPage: React.FC = () => {
const [strategies, setStrategies] = useState<AlphaSiftStrategy[]>([]);
const [maxResults, setMaxResults] = useState(restoredTask?.maxResults || 3);
const [candidates, setCandidates] = useState<AlphaSiftCandidate[]>([]);
const [hotspots, setHotspots] = useState<AlphaSiftHotspot[]>([]);
const [selectedHotspotTopic, setSelectedHotspotTopic] = useState<string | null>(null);
const selectedHotspotTopicRef = useRef<string | null>(null);
const hotspotDetailRequestIdRef = useRef(0);
const [hotspotDetail, setHotspotDetail] = useState<AlphaSiftHotspotDetail | null>(null);
const [loadingHotspotDetail, setLoadingHotspotDetail] = useState(false);
const [hotspotDetailError, setHotspotDetailError] = useState('');
const [loadingHotspots, setLoadingHotspots] = useState(false);
const [hotspotError, setHotspotError] = useState('');
const [screenMeta, setScreenMeta] = useState<AlphaSiftScreenResponse | null>(null);
const [expandedCode, setExpandedCode] = useState<string | null>(null);
const [loading, setLoading] = useState(Boolean(restoredTask?.taskId));
@@ -311,6 +322,36 @@ const StockScreeningPage: React.FC = () => {
setExpandedCode(null);
};
const loadHotspotDetail = useCallback(async (topic: string) => {
if (!topic) {
return;
}
const requestId = hotspotDetailRequestIdRef.current + 1;
hotspotDetailRequestIdRef.current = requestId;
const isCurrentRequest = () => hotspotDetailRequestIdRef.current === requestId;
const canApplyRequest = () => isCurrentRequest() && selectedHotspotTopicRef.current === topic;
setLoadingHotspotDetail(true);
setHotspotDetail((currentDetail) => (currentDetail?.topic === topic ? currentDetail : null));
setHotspotDetailError('');
try {
const detail = await alphasiftApi.getHotspotDetail({ topic, provider: 'akshare' });
if (!canApplyRequest()) {
return;
}
setHotspotDetail(detail);
} catch (err) {
if (!canApplyRequest()) {
return;
}
setHotspotDetail(null);
setHotspotDetailError(toApiErrorMessage(err, '热点题材详情加载失败,请稍后重试。'));
} finally {
if (isCurrentRequest()) {
setLoadingHotspotDetail(false);
}
}
}, []);
const loadStrategies = useCallback(async () => {
setLoadingStrategies(true);
try {
@@ -331,6 +372,51 @@ const StockScreeningPage: React.FC = () => {
}
}, []);
const loadHotspots = useCallback(async (refresh = false) => {
setLoadingHotspots(true);
setHotspotError('');
try {
const result = await alphasiftApi.getHotspots({ provider: 'akshare', top: 12, refresh });
const nextHotspots = result.hotspots || [];
const currentTopic = selectedHotspotTopicRef.current;
const retainedTopic = Boolean(currentTopic && nextHotspots.some((item) => item.topic === currentTopic));
const nextTopic = retainedTopic ? currentTopic : nextHotspots[0]?.topic ?? null;
setHotspots(nextHotspots);
setSelectedHotspotTopic(nextTopic);
selectedHotspotTopicRef.current = nextTopic;
if (retainedTopic && refresh && nextTopic) {
void loadHotspotDetail(nextTopic);
} else if (!retainedTopic) {
setHotspotDetail(null);
}
setHotspotDetailError('');
if (nextHotspots.length === 0) {
const sourceError = result.sourceErrors?.[0];
setHotspotError(sourceError ? `热点题材暂未返回数据:${sourceError}` : '热点题材暂未返回数据');
}
} catch (err) {
setHotspotError(toApiErrorMessage(err, '热点题材加载失败,请稍后重试。'));
} finally {
setLoadingHotspots(false);
}
}, [loadHotspotDetail]);
const handleHotspotSelect = useCallback((topic: string) => {
selectedHotspotTopicRef.current = topic;
setSelectedHotspotTopic(topic);
}, []);
useEffect(() => {
selectedHotspotTopicRef.current = selectedHotspotTopic;
}, [selectedHotspotTopic]);
useEffect(() => {
if (!selectedHotspotTopic) {
return;
}
void loadHotspotDetail(selectedHotspotTopic);
}, [loadHotspotDetail, selectedHotspotTopic]);
useEffect(() => {
let active = true;
alphasiftApi
@@ -343,6 +429,7 @@ const StockScreeningPage: React.FC = () => {
setAvailable(status.available);
if (status.enabled && status.available) {
void loadStrategies();
void loadHotspots(false);
}
})
.catch(() => {
@@ -354,7 +441,7 @@ const StockScreeningPage: React.FC = () => {
return () => {
active = false;
};
}, [loadStrategies]);
}, [loadHotspots, loadStrategies]);
useEffect(() => {
if (!activeTaskId) {
@@ -566,6 +653,141 @@ const StockScreeningPage: React.FC = () => {
{error ? <InlineAlert variant="danger" title="调用失败" message={error} /> : null}
<section className="rounded-2xl border border-orange-300/60 bg-card/95 p-4 shadow-soft-card">
<div className="mb-4 flex flex-col gap-3 sm:flex-row sm:items-center sm:justify-between">
<div className="flex items-start gap-3">
<span className="grid h-8 w-8 place-items-center rounded-xl bg-orange-500/10 text-orange-500">
<Flame className="h-4 w-4" />
</span>
<div>
<h2 className="text-sm font-semibold text-foreground"></h2>
<p className="mt-1 text-xs leading-5 text-secondary-text">
AlphaSift hotspot capital_heatbalanced_alpha theme_heat
</p>
</div>
</div>
<Button
size="sm"
variant="secondary"
isLoading={loadingHotspots}
loadingText="刷新中..."
disabled={!isScreeningEnabled || loadingHotspots}
onClick={() => void loadHotspots(true)}
>
<RefreshCw className="h-4 w-4" />
</Button>
</div>
{hotspotError ? (
<p className="mb-3 rounded-xl border border-warning/30 bg-warning/10 px-3 py-2 text-xs text-warning">
{hotspotError}
</p>
) : null}
{hotspots.length === 0 ? (
<div className="rounded-xl border border-dashed border-border bg-surface/70 px-4 py-6 text-sm text-secondary-text">
/
</div>
) : (
<div className="grid gap-2 md:grid-cols-2 xl:grid-cols-4">
{hotspots.map((item) => {
const selected = selectedHotspotTopic === item.topic;
return (
<button
key={`${item.topic}-${item.rank ?? ''}`}
className={`rounded-xl border p-4 text-left transition-colors ${
selected
? 'border-orange-400 bg-orange-500/10 shadow-[0_0_0_1px_hsl(var(--warning)/0.18)]'
: 'border-border/80 bg-surface/70 hover:border-orange-300/80 hover:bg-orange-500/5'
}`}
type="button"
onClick={() => handleHotspotSelect(item.topic)}
>
<div className="flex items-start justify-between gap-3">
<div>
<p className="text-sm font-semibold text-foreground">{item.name || item.topic}</p>
<p className="mt-1 text-xs text-secondary-text">{item.stage || item.state || '阶段待观察'}</p>
</div>
<span className="rounded-full bg-orange-500/10 px-2 py-1 text-xs font-semibold text-orange-500">
{formatNumber(item.heatScore, 0)}
</span>
</div>
<div className="mt-3 grid gap-1 text-xs text-secondary-text">
<span> {formatNumber(item.changePct)}%</span>
<span> {formatNumber(item.trendScore, 1)} · {formatNumber(item.persistenceScore, 1)}</span>
<span> {item.sampleStockCount ?? 0} · {(item.leaders || []).slice(0, 2).join('、') || '-'}</span>
</div>
</button>
);
})}
</div>
)}
{selectedHotspotTopic ? (
<div className="mt-4 rounded-xl border border-border/80 bg-surface/80 p-4">
<div className="mb-3 flex flex-col gap-2 sm:flex-row sm:items-center sm:justify-between">
<div>
<h3 className="text-sm font-semibold text-foreground">
{hotspotDetail?.name || selectedHotspotTopic}
</h3>
<p className="mt-1 text-xs leading-5 text-secondary-text">
{loadingHotspotDetail ? '正在读取发酵路线与概念股...' : hotspotDetail?.summary || '点击题材查看发酵路线与概念股。'}
</p>
</div>
{hotspotDetail?.stockCount != null ? (
<span className="w-fit rounded-full bg-orange-500/10 px-3 py-1 text-xs font-semibold text-orange-500">
{hotspotDetail.stockCount}
</span>
) : null}
</div>
{hotspotDetailError ? (
<p className="mb-3 rounded-xl border border-warning/30 bg-warning/10 px-3 py-2 text-xs text-warning">
{hotspotDetailError}
</p>
) : null}
{hotspotDetail ? (
<div className="grid gap-4 lg:grid-cols-[1fr_1.3fr]">
<div>
<p className="mb-2 text-xs font-semibold text-secondary-text">线</p>
<div className="space-y-2">
{(hotspotDetail.route || []).map((item, index) => (
<div key={`${item.title}-${index}`} className="rounded-lg border border-border/70 bg-card/80 p-3">
<p className="text-xs font-semibold text-foreground">{item.title}</p>
<p className="mt-1 text-xs leading-5 text-secondary-text">{item.description}</p>
</div>
))}
</div>
</div>
<div>
<p className="mb-2 text-xs font-semibold text-secondary-text"></p>
<div className="grid gap-2 sm:grid-cols-2">
{(hotspotDetail.stocks || []).slice(0, 10).map((stock) => (
<div key={`${stock.code || stock.name}`} className="rounded-lg border border-border/70 bg-card/80 p-3">
<div className="flex items-start justify-between gap-2">
<div>
<p className="text-xs font-semibold text-foreground">{stock.name || stock.code || '-'}</p>
<p className="mt-1 text-[11px] text-secondary-text">{stock.code || '-'}</p>
</div>
<span className="rounded-full bg-cyan/10 px-2 py-1 text-[11px] font-semibold text-cyan">
{stock.role || '概念股'}
</span>
</div>
<p className="mt-2 text-[11px] text-secondary-text">
{formatNumber(stock.changePct)}% · {formatNumber(stock.hotStockScore, 0)}
</p>
</div>
))}
</div>
</div>
</div>
) : null}
</div>
) : null}
</section>
<section className="rounded-2xl border border-cyan/35 bg-card/95 p-4 shadow-soft-card">
<div className="mb-4 flex items-center justify-between gap-3">
<div>

View File

@@ -1,10 +1,12 @@
import { fireEvent, render, screen, waitFor } from '@testing-library/react';
import { act, fireEvent, render, screen, waitFor } from '@testing-library/react';
import { beforeEach, describe, expect, it, vi } from 'vitest';
import StockScreeningPage from '../StockScreeningPage';
const {
enableAlphaSift,
getAlphaSiftStatus,
getHotspotDetail,
getHotspots,
getStrategies,
getScreenTask,
resetLastScreenResult,
@@ -39,6 +41,8 @@ const {
return {
enableAlphaSift: vi.fn(),
getAlphaSiftStatus: vi.fn(),
getHotspotDetail: vi.fn(),
getHotspots: vi.fn(),
getStrategies: vi.fn(),
getScreenTask,
resetLastScreenResult: () => {
@@ -53,6 +57,8 @@ vi.mock('../../api/alphasift', () => ({
alphasiftApi: {
enable: () => enableAlphaSift(),
getStatus: () => getAlphaSiftStatus(),
getHotspotDetail: (payload: unknown) => getHotspotDetail(payload),
getHotspots: (payload: unknown) => getHotspots(payload),
getStrategies: () => getStrategies(),
getScreenTask: (taskId: string) => getScreenTask(taskId),
screen: (payload: unknown) => screenStocks(payload),
@@ -77,16 +83,39 @@ const mockStrategiesResponse = {
strategyCount: 1,
};
function createDeferred<T>() {
let resolve: (value: T) => void = () => {};
let reject: (reason?: unknown) => void = () => {};
const promise = new Promise<T>((resolvePromise, rejectPromise) => {
resolve = resolvePromise;
reject = rejectPromise;
});
return { promise, resolve, reject };
}
describe('StockScreeningPage', () => {
beforeEach(() => {
enableAlphaSift.mockReset();
getAlphaSiftStatus.mockReset();
getHotspotDetail.mockReset();
getHotspots.mockReset();
getStrategies.mockReset();
getScreenTask.mockClear();
resetLastScreenResult();
screenStocks.mockReset();
startScreenTask.mockClear();
getStrategies.mockResolvedValue(mockStrategiesResponse);
getHotspotDetail.mockResolvedValue({
enabled: true,
provider: 'akshare',
topic: 'AI算力',
name: 'AI算力',
summary: 'AI算力 盘中发酵。',
route: [{ title: '盘中发酵', description: '出现大笔买入。', source: 'eastmoney_board_change' }],
stocks: [{ code: '300000', name: '中际旭创', role: '核心龙头', hotStockScore: 88 }],
stockCount: 1,
});
getHotspots.mockResolvedValue({ enabled: true, provider: 'akshare', hotspots: [], hotspotCount: 0 });
window.sessionStorage.clear();
});
@@ -118,6 +147,382 @@ describe('StockScreeningPage', () => {
expect(screen.getByText('AlphaSift 适配层不可用。请执行 pip install -r requirements.txt')).toBeInTheDocument();
});
it('loads AlphaSift hotspot themes on demand', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,
available: true,
installSpecIsDefault: true,
});
getHotspots
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [],
hotspotCount: 0,
cacheUsed: true,
cachedAt: '2026-06-07T08:00:00Z',
})
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 88,
trendScore: 12,
persistenceScore: 66,
changePct: 4.2,
stage: '加速主升',
sampleStockCount: 8,
leaders: ['中际旭创', '工业富联'],
},
],
hotspotCount: 1,
});
render(<StockScreeningPage />);
expect(await screen.findByText('选股已开启')).toBeInTheDocument();
await waitFor(() => expect(getHotspots).toHaveBeenCalledWith({ provider: 'akshare', top: 12, refresh: false }));
fireEvent.click(screen.getByRole('button', { name: /刷新热点题材/ }));
await waitFor(() => expect(getHotspots).toHaveBeenCalledWith({ provider: 'akshare', top: 12, refresh: true }));
await waitFor(() => expect(getHotspotDetail).toHaveBeenCalledWith({ topic: 'AI算力', provider: 'akshare' }));
await waitFor(() => expect(screen.getAllByText('AI算力').length).toBeGreaterThan(0));
expect(screen.getByText('加速主升')).toBeInTheDocument();
expect(screen.getByText(/中际旭创、工业富联/)).toBeInTheDocument();
expect(await screen.findByText('发酵路线')).toBeInTheDocument();
expect(screen.getByText('盘中发酵')).toBeInTheDocument();
expect(screen.getByText('概念股')).toBeInTheDocument();
expect(screen.getByText('中际旭创')).toBeInTheDocument();
});
it('loads selected hotspot detail once when switching themes', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,
available: true,
installSpecIsDefault: true,
});
getHotspots.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 88,
stage: '加速主升',
},
{
topic: '机器人执行器',
name: '机器人执行器',
heatScore: 80,
stage: '轮动扩散',
},
],
hotspotCount: 2,
});
render(<StockScreeningPage />);
expect(await screen.findByText('选股已开启')).toBeInTheDocument();
await waitFor(() => expect(getHotspotDetail).toHaveBeenCalledWith({ topic: 'AI算力', provider: 'akshare' }));
expect(getHotspotDetail).toHaveBeenCalledTimes(1);
fireEvent.click(screen.getByRole('button', { name: /机器人执行器/ }));
await waitFor(() =>
expect(getHotspotDetail).toHaveBeenLastCalledWith({ topic: '机器人执行器', provider: 'akshare' }),
);
await new Promise((resolve) => window.setTimeout(resolve, 0));
expect(getHotspotDetail).toHaveBeenCalledTimes(2);
});
it('clears loaded hotspot detail while loading a different theme', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,
available: true,
installSpecIsDefault: true,
});
getHotspots.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 88,
stage: '加速主升',
},
{
topic: '机器人执行器',
name: '机器人执行器',
heatScore: 80,
stage: '轮动扩散',
},
],
hotspotCount: 2,
});
const robotDetail = createDeferred<unknown>();
getHotspotDetail
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
topic: 'AI算力',
name: 'AI算力',
summary: 'AI算力 盘中发酵。',
route: [{ title: '盘中发酵', description: '出现大笔买入。', source: 'eastmoney_board_change' }],
stocks: [{ code: '300000', name: '中际旭创', role: '核心龙头', hotStockScore: 88 }],
stockCount: 1,
})
.mockImplementationOnce(({ topic }: { topic: string }) => {
if (topic === '机器人执行器') {
return robotDetail.promise;
}
return Promise.reject(new Error(`unexpected topic: ${topic}`));
});
render(<StockScreeningPage />);
expect(await screen.findByText('盘中发酵')).toBeInTheDocument();
expect(screen.getByText('中际旭创')).toBeInTheDocument();
fireEvent.click(screen.getByRole('button', { name: /机器人执行器/ }));
await waitFor(() =>
expect(getHotspotDetail).toHaveBeenLastCalledWith({ topic: '机器人执行器', provider: 'akshare' }),
);
expect(screen.getAllByText('机器人执行器').length).toBeGreaterThan(0);
expect(screen.getByText('正在读取发酵路线与概念股...')).toBeInTheDocument();
expect(screen.queryByText('盘中发酵')).not.toBeInTheDocument();
expect(screen.queryByText('中际旭创')).not.toBeInTheDocument();
await act(async () => {
robotDetail.resolve({
enabled: true,
provider: 'akshare',
topic: '机器人执行器',
name: '机器人执行器',
summary: '机器人执行器 继续发酵。',
route: [{ title: '机器人发酵', description: '执行器链条扩散。', source: 'eastmoney_board_change' }],
stocks: [{ code: '300111', name: '机器人龙头', role: '核心龙头', hotStockScore: 86 }],
stockCount: 1,
});
});
expect(await screen.findByText('机器人发酵')).toBeInTheDocument();
expect(screen.getByText('机器人龙头')).toBeInTheDocument();
});
it('ignores stale hotspot detail responses when switching themes', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,
available: true,
installSpecIsDefault: true,
});
getHotspots.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 88,
stage: '加速主升',
},
{
topic: '机器人执行器',
name: '机器人执行器',
heatScore: 80,
stage: '轮动扩散',
},
],
hotspotCount: 2,
});
const aiDetail = createDeferred<unknown>();
const robotDetail = createDeferred<unknown>();
getHotspotDetail.mockImplementation(({ topic }: { topic: string }) => {
if (topic === 'AI算力') {
return aiDetail.promise;
}
if (topic === '机器人执行器') {
return robotDetail.promise;
}
return Promise.reject(new Error(`unexpected topic: ${topic}`));
});
render(<StockScreeningPage />);
expect(await screen.findByText('选股已开启')).toBeInTheDocument();
await waitFor(() => expect(getHotspotDetail).toHaveBeenCalledWith({ topic: 'AI算力', provider: 'akshare' }));
fireEvent.click(screen.getByRole('button', { name: /机器人执行器/ }));
await waitFor(() =>
expect(getHotspotDetail).toHaveBeenLastCalledWith({ topic: '机器人执行器', provider: 'akshare' }),
);
await act(async () => {
robotDetail.resolve({
enabled: true,
provider: 'akshare',
topic: '机器人执行器',
name: '机器人执行器',
summary: '机器人执行器 继续发酵。',
route: [{ title: '机器人发酵', description: '执行器链条扩散。', source: 'eastmoney_board_change' }],
stocks: [{ code: '300111', name: '机器人龙头', role: '核心龙头', hotStockScore: 86 }],
stockCount: 1,
});
});
expect(await screen.findByText('机器人发酵')).toBeInTheDocument();
await act(async () => {
aiDetail.resolve({
enabled: true,
provider: 'akshare',
topic: 'AI算力',
name: 'AI算力',
summary: 'AI算力 旧响应。',
route: [{ title: 'AI旧发酵', description: '旧请求晚到。', source: 'eastmoney_board_change' }],
stocks: [{ code: '300000', name: '中际旭创', role: '核心龙头', hotStockScore: 88 }],
stockCount: 1,
});
});
expect(screen.getByText('机器人发酵')).toBeInTheDocument();
expect(screen.getByText('机器人龙头')).toBeInTheDocument();
expect(screen.queryByText('AI旧发酵')).not.toBeInTheDocument();
expect(screen.queryByText('中际旭创')).not.toBeInTheDocument();
});
it('reloads selected hotspot detail when refreshed themes keep the same topic', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,
available: true,
installSpecIsDefault: true,
});
getHotspots
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 88,
stage: '加速主升',
},
{
topic: '机器人执行器',
name: '机器人执行器',
heatScore: 80,
stage: '轮动扩散',
},
],
hotspotCount: 2,
})
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 91,
stage: '高位发酵',
},
],
hotspotCount: 1,
});
getHotspotDetail
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
topic: 'AI算力',
name: 'AI算力',
summary: 'AI算力 盘中发酵。',
route: [{ title: '盘中发酵', description: '出现大笔买入。', source: 'eastmoney_board_change' }],
stocks: [{ code: '300000', name: '中际旭创', role: '核心龙头', hotStockScore: 88 }],
stockCount: 1,
})
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
topic: 'AI算力',
name: 'AI算力',
summary: 'AI算力 刷新后继续发酵。',
route: [{ title: '刷新发酵', description: '刷新后仍在榜内。', source: 'eastmoney_board_change' }],
stocks: [{ code: '601138', name: '工业富联', role: '核心龙头', hotStockScore: 90 }],
stockCount: 1,
});
render(<StockScreeningPage />);
expect(await screen.findByText('选股已开启')).toBeInTheDocument();
await waitFor(() => expect(getHotspotDetail).toHaveBeenCalledTimes(1));
fireEvent.click(screen.getByRole('button', { name: /刷新热点题材/ }));
await waitFor(() => expect(getHotspots).toHaveBeenCalledWith({ provider: 'akshare', top: 12, refresh: true }));
await waitFor(() => expect(getHotspotDetail).toHaveBeenCalledTimes(2));
expect(getHotspotDetail).toHaveBeenLastCalledWith({ topic: 'AI算力', provider: 'akshare' });
expect(await screen.findByText('刷新发酵')).toBeInTheDocument();
expect(screen.getByText('工业富联')).toBeInTheDocument();
});
it('keeps existing hotspot cards when manual refresh fails', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,
available: true,
installSpecIsDefault: true,
});
getHotspots
.mockResolvedValueOnce({
enabled: true,
provider: 'akshare',
providerUsed: 'akshare',
hotspots: [
{
topic: 'AI算力',
name: 'AI算力',
heatScore: 88,
trendScore: 12,
persistenceScore: 66,
changePct: 4.2,
stage: '加速主升',
sampleStockCount: 8,
leaders: ['中际旭创', '工业富联'],
},
],
hotspotCount: 1,
})
.mockRejectedValueOnce(new Error('manual refresh failed'));
render(<StockScreeningPage />);
expect(await screen.findByText('选股已开启')).toBeInTheDocument();
expect(await screen.findByText('加速主升')).toBeInTheDocument();
expect(screen.getByText(/中际旭创、工业富联/)).toBeInTheDocument();
fireEvent.click(screen.getByRole('button', { name: /刷新热点题材/ }));
await waitFor(() => expect(getHotspots).toHaveBeenCalledWith({ provider: 'akshare', top: 12, refresh: true }));
expect(await screen.findByText(/manual refresh failed/)).toBeInTheDocument();
expect(screen.getByText('加速主升')).toBeInTheDocument();
expect(screen.getByText(/中际旭创、工业富联/)).toBeInTheDocument();
expect(screen.queryByText(/点击刷新后会拉取热点概念/)).not.toBeInTheDocument();
});
it('shows input strategy when strategy is not in preset list', async () => {
getAlphaSiftStatus.mockResolvedValueOnce({
enabled: true,

View File

@@ -23,7 +23,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
- [修复] 修复运行流 live SSE 事件未复用快照层递归脱敏规则的问题避免本地路径、prompt/raw response、代理头等敏感诊断字段在 refetch 前短暂暴露。
<!-- 新条目格式:- [类型] 描述(类型取值:新功能/改进/修复/文档/测试/chore-->
<!-- 每条独立一行追加到本段末尾,无需分类标题,合并时冲突最小 -->
- [改进] AlphaSift 依赖锁定更新到 `de54ea0da367be85770d9589a5bf7ded4f62d386`,并为新版 last-good snapshot、日线历史、行业/概念 provider cache、hotspot 具体题材榜单、题材发酵路线、概念股详情、上次成功热点缓存与 post-analysis 元信息补齐 DSA 运行期和 Web 选股页适配;默认不启用 DSA deep-analysis 回调。
- [修复] 桌面发布打包改用冻结可执行文件运行时探针校验 `alphasift.dsa_adapter`,避免 macOS PyInstaller 将模块内嵌进可执行文件时被文件系统/zip 扫描误判为缺失。
- [新功能] 新增 AlphaSift 热点题材链路:后端新增 `/api/v1/alphasift/hotspots``/api/v1/alphasift/hotspots/{topic}` APIWeb 选股页新增“热点题材”区域并支持发酵路线与概念股查看。
- [改进] 新增 AlphaSift 热点题材读取与刷新策略:默认优先读取上次成功热点缓存,手动刷新才实时拉取并覆盖缓存,实时拉取失败时尽量回退旧缓存。
- [改进] 改造 `main.py --webui-only` 启动行为:若 FastAPI 监听端口已被占用,启动即 fail-fast 抛出明确错误并退出,避免 Windows 下 `WinError 10048`
- [文档] 补充 `docs/alphasift-integration.md`:明确 AlphaSift 锁定 commit 来源、Hotspot 契约边界、LLM/LiteLLM 兼容语义与关闭开关下回退路径。
- [修复] 为 THS 发酵路线补充列名兜底:当 `stock_board_concept_summary_ths` 返回缺列时仅跳过该来源富化,不影响 `/api/v1/alphasift/hotspots/{topic}` 返回。
- [测试] 新增/更新后端回归:`python -m pytest tests/test_alphasift_api.py -q``python -m pytest tests/test_docker_entrypoint.py -q``python -m pytest tests/test_main_schedule_mode.py -q -k "start_api_server_fails_before_thread_when_port_is_busy"`.
## [3.21.0] - 2026-06-07

View File

@@ -6,7 +6,7 @@ AlphaSift 作为独立仓库维护的选股引擎接入 DSA。DSA 默认不启
- 默认关闭:`ALPHASIFT_ENABLED=false`
- 启用入口:设置页或选股页点击开启,或在 `.env` 中配置 `ALPHASIFT_ENABLED=true`
- 依赖来源:`requirements.txt` 固定到已验证的 AlphaSift 适配层 commit`git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344#egg=alphasift`。该来源覆盖 `alphasift.dsa_adapter` 契约与 `screen/list_strategies/get_status` 调用。
- 依赖来源:`requirements.txt` 固定到已验证的 AlphaSift 适配层 commit`git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386#egg=alphasift`(对应提交 `https://github.com/ZhuLinsen/alphasift/commit/de54ea0da367be85770d9589a5bf7ded4f62d386`。该来源覆盖 `alphasift.dsa_adapter` 契约与 `screen/list_strategies/get_status` 调用。
- 修复安装来源:`ALPHASIFT_INSTALL_SPEC` 仍保留,默认等于同一个受信任 commit。它不再是策略列表或选股接口的运行时安装主路径只用于显式调用 `/api/v1/alphasift/install` 时做修复安装和来源校验。
- 缺失依赖边界:如果运行环境缺少 `alphasift.dsa_adapter``status` 返回 `available=false + diagnostics.reason=missing_module``strategies``screen` 返回 `424` 并提示执行 `pip install -r requirements.txt` 或重建 Docker/桌面后端产物,不会在业务请求中自动 `pip install`
- 运行异常边界:若适配层可导入但 `get_status()` 报错或返回 `available=false`DSA 返回 `424 + diagnostics`,保留故障诊断,防止用重装掩盖真实运行时错误。
@@ -15,6 +15,7 @@ AlphaSift 作为独立仓库维护的选股引擎接入 DSA。DSA 默认不启
- 日 K 线补特征DSA 调用 AlphaSift 时会优先复用 DSA 历史行情加载链路数据库缓存、Tushare、Efinance、Akshare、Pytdx、Baostock、Yfinance 等 fallback仅在 DSA 链路无可用数据时回退到 AlphaSift 原始日线数据源,减少单一上游超时拖垮选股。
- LLM 环境DSA 调用 AlphaSift 时会桥接 DSA 已解析的 `LITELLM_MODEL``LITELLM_FALLBACK_MODELS``LLM_CHANNELS``LLM_<NAME>_*``LITELLM_CONFIG`、渠道额外请求头和各模型密钥AlphaSift 独立运行时仍使用自己的 `.env`/环境变量。
- 快照源DSA 调用 AlphaSift 时,未显式配置 `SNAPSHOT_SOURCE_PRIORITY` 会默认优先使用 `em_datacenter`,减少 Tushare/东方财富行情接口在夜间或网络抖动时逐个失败造成的等待;显式配置的源顺序会原样保留。
- 最新 AlphaSift 能力:锁定 commit `de54ea0da367be85770d9589a5bf7ded4f62d386` 包含选股 pipeline 性能优化、last-good snapshot fallback、日线历史缓存、行业/概念 provider cache、热点/行业热度因子、hotspot 热点题材榜单与本地 scorecard/post-analysis 元信息。DSA 调用时会注入隔离缓存默认路径 `data/alphasift``data/alphasift/snapshot.last_good.json``data/alphasift/daily_history``data/alphasift/industry_provider_cache`Web 选股页提供“热点题材”手动刷新入口,请求 `/api/v1/alphasift/hotspots` 时会显式使用 `akshare` provider 优先拉取具体概念/题材异动(例如玻璃基板、机器人执行器、减速器),行业板块仅作为兜底;默认打开页面时优先读取上一次成功的热点题材缓存,点击刷新才实时拉取并覆盖缓存,实时拉取失败时会尽量回退旧缓存;点击题材会请求 `/api/v1/alphasift/hotspots/{topic}` 展示发酵路线与概念股;不会默认触发 AlphaSift 的 DSA deep-analysis 回调,避免无提示扩大递归调用面。
- 风险提示:前端设置页和选股页展示第三方来源与投资风险说明;不会弹窗打断用户。
## AlphaSift 适配层要求
@@ -94,7 +95,7 @@ context = {
AlphaSift 会在 L1 初筛后、LLM 重排前调用 `context["dsa"]` 中的 provider为有限 Top 候选补充 DSA 行情和基本面轻量上下文,并把 `dsa_context` 随候选返回。新闻搜索、完整摘要和缺失字段补全由 DSA API 在最终 Top 候选阶段执行若候选已经携带完整新闻上下文DSA API 返回阶段会复用这些字段,避免重复请求。
AlphaSift 侧已在 `ZhuLinsen/alphasift@1a0ed8c99b3615c0cb1076e6029827ffc6de2344` 提供 DSA provider context 支持、DSA adapter contract并支持复用 DSA 的 `LLM_TIMEOUT_SEC`
AlphaSift 侧已在 `ZhuLinsen/alphasift@de54ea0da367be85770d9589a5bf7ded4f62d386` 提供 DSA provider context 支持、DSA adapter contract并支持复用 DSA 的 `LLM_TIMEOUT_SEC`
## DSA 后端行为
@@ -108,13 +109,14 @@ AlphaSift 侧已在 `ZhuLinsen/alphasift@1a0ed8c99b3615c0cb1076e6029827ffc6de234
## 配置兼容边界LLM / LiteLLM / Base URL
- 兼容语义与版本证据(可追溯):
- 运行依赖约束:`requirements.txt` 中将 LiteLLM 固定到 `litellm>=1.80.10,!=1.82.7,!=1.82.8,<2.0.0`,并通过 `git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344` 安装 AlphaSift 适配层。
- 运行依赖约束:`requirements.txt` 中将 LiteLLM 固定到 `litellm>=1.80.10,!=1.82.7,!=1.82.8,<2.0.0`,并通过 `git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386` 安装 AlphaSift 适配层。
- 文档依据:
- LiteLLM Providers: https://docs.litellm.ai/docs/providers
- LiteLLM OpenAI-compatible: https://docs.litellm.ai/docs/providers/openai_compatible
- LiteLLM model_list/proxy 配置(含 `api_base``api_key``extra_headers`: https://docs.litellm.ai/docs/proxy/configs
- OpenAI 请求语义与授权头: https://platform.openai.com/docs/api-reference/making-requests、https://platform.openai.com/docs/api-reference/authentication
- 结构化检测澄清:本 PR 触及 `.env.example``requirements.txt``src/config.py` 与本文档,是因为 AlphaSift 依赖 pin 更新和调用期 runtime bridge 需要把既有 DSA LLM 配置透传给外部适配层;本 PR 没有升级 LiteLLM 主版本、没有新增或改名 provider 协议、没有修改 `LITELLM_MODEL`/`LITELLM_FALLBACK_MODELS`/`LLM_CHANNELS`/`LLM_<NAME>_*` 的持久化解析语义。
- LLM 运行时兼容边界AlphaSift 不改变主配置链路,只在调用期注入已解析的 `LITELLM_MODEL``LITELLM_FALLBACK_MODELS``LLM_CHANNELS``LLM_<NAME>_*` 到进程环境;受管 provider 的 fallback 过滤行为保持现有策略,不做历史配置的静默迁移。`ALPHASIFT_ENABLED` 是当前场景唯一新增持久化分支。
- 注意:本注入是**短时内存注入**,不会改写 `.env`、不会回写历史配置、不会静默迁移用户自定义 provider/model 路由;失败或未开启时,除了 AlphaSift 选股能力本身,其它 DSA 业务链路保持既有配置执行。
- 注入来源与回滚原则:
@@ -130,11 +132,17 @@ AlphaSift 侧已在 `ZhuLinsen/alphasift@1a0ed8c99b3615c0cb1076e6029827ffc6de234
- 官方 `model_list`/额外头依据LiteLLM config 文档([https://docs.litellm.ai/docs/proxy/configs](https://docs.litellm.ai/docs/proxy/configs))说明 `litellm_params` 支持 `model``api_base``api_key``extra_headers`。DSA 只把已声明渠道转换为同类结构传给 AlphaSift不新增模型路由映射不做 provider 模式迁移。
- 兼容头部语义依据OpenAI 调用约定([https://platform.openai.com/docs/api-reference/making-requests](https://platform.openai.com/docs/api-reference/making-requests))与鉴权约定([https://platform.openai.com/docs/api-reference/authentication](https://platform.openai.com/docs/api-reference/authentication))对应 `Authorization` 与自定义 header 传递行为,`extra_headers` 仅用于补充会话头,不改写模型路由。
- 回退路径为“设置页关闭 AlphaSift 或保留 `ALPHASIFT_ENABLED=false`”,并保持原有 `LITELLM_*``LLM_*` 配置,触发失败时可先核对 `status`/`screen``diagnostics` 后执行服务重启。
- 旧配置保留证据:
- `src/services/alphasift_service.py``_alphasift_runtime_env()` 会在调用前保存同名 `os.environ` 值,并在调用后逐项恢复或删除本次临时新增键;该路径不调用 `dotenv_values()` 写回,也不修改 `.env` 文件。
- `src/services/alphasift_service.py``_build_alphasift_runtime_env()` 只从当前 `Config` 生成临时 env dict未声明渠道不会生成 `LLM_<NAME>_*`,已有自定义 provider/model/base URL 不会被重命名或清理。
- `tests/test_alphasift_api.py` 覆盖 `test_screen_bridges_dsa_llm_config_into_alphasift_runtime``test_screen_bridges_legacy_openai_fields_into_alphasift_runtime_env``test_screen_injects_openai_compatible_model_headers_into_alphasift_litellm_calls``test_screen_disabled_preserves_existing_llm_env_state``test_screen_filters_undeclared_managed_fallbacks_for_dsa_routes`用于证明注入、OpenAI-compatible header/base URL、关闭状态和未声明 fallback 均不改写用户原始配置。
- 失败可见性:`status`/`screen` 接口返回明确错误码与 `message`,前端在设置页或选股页会将 `403/424/400/422` 等错误直接提示给用户,便于定位并回退到“关闭 AlphaSift + 保持原有 LLM 运行链路”。
## 兼容验收索引(发布前核验)
- 依赖与源码约束核验:`requirements.txt` 中的 `litellm` 约束与 `src/config.py`/`requirements.txt` 一致。
- Hotspot 契约兼容核验:`docs/alphasift-integration.md``api/v1/endpoints/alphasift.py``src/services/alphasift_service.py` 保持 `hotspots`/`hotspots/{topic}` 字段与 `tests/test_alphasift_api.py` 一致,调用前后默认使用 `snapshot.last_good` 缓存兜底。
- 外部版本来源:本次集成依赖来源为 `https://github.com/ZhuLinsen/alphasift/commit/de54ea0da367be85770d9589a5bf7ded4f62d386`,需在复验时按该 commit pin 回放导入与接口契约。
- 行为核验:`src/services/alphasift_service.py``_build_alphasift_runtime_env``_build_alphasift_context` 仅在调用期写入进程环境;`/api/v1/alphasift/screen``strategies``status` 在运行期不回写 `.env`
- 回退核验:关闭 `ALPHASIFT_ENABLED` 并重启配置链路后,系统恢复原始 `LITELLM_MODEL/FALLBACK_MODELS``LLM_CHANNELS``LLM_*` 运行语义,不执行迁移清理脚本。
- 语义来源核验LiteLLM 文档https://docs.litellm.ai/docs/providers、OpenAI-compatible 文档https://docs.litellm.ai/docs/providers/openai_compatible与 LiteLLM 配置文档https://docs.litellm.ai/docs/proxy/configs用于核对 provider/model/base_url/extra_headers 映射链路。
@@ -172,6 +180,7 @@ Docker 镜像与桌面发布包保持一致:`docker/Dockerfile` 会通过 `req
## 验证记录
- `python -m pytest tests/test_alphasift_api.py -q`
- `python -m pytest tests/test_main_schedule_mode.py -q -k "start_api_server_fails_before_thread_when_port_is_busy"`
- `python -m py_compile api/v1/endpoints/alphasift.py src/services/alphasift_service.py tests/test_alphasift_api.py src/config.py src/core/config_registry.py`
- `cd apps/dsa-web && npm run test -- alphasift.test.ts StockScreeningPage.test.tsx SettingsPage.test.tsx --run`
- `cd apps/dsa-web && npm run lint`

12
main.py
View File

@@ -713,9 +713,18 @@ def start_api_server(host: str, port: int, config: Config) -> None:
port: 监听端口
config: 配置对象
"""
import socket
import threading
import uvicorn
probe = socket.socket(socket.AF_INET6 if ":" in host else socket.AF_INET, socket.SOCK_STREAM)
try:
probe.bind((host, port))
except OSError as exc:
raise RuntimeError(f"FastAPI port is not available: {host}:{port}") from exc
finally:
probe.close()
def run_server():
level_name = (config.log_level or "INFO").lower()
uvicorn.run(
@@ -904,6 +913,9 @@ def main() -> int:
bot_clients_started = True
except Exception as e:
logger.error(f"启动 FastAPI 服务失败: {e}")
if args.serve_only:
return 1
start_serve = False
if bot_clients_started:
start_bot_stream_clients(config)

View File

@@ -19,7 +19,7 @@ yfinance>=0.2.0 # Priority 4: Yahoo Finance (Fallback)
longbridge>=0.2.77 # Priority 5: Longbridge OpenAPI fallback for US/HK stocks; OAuth capability checked at runtime
tickflow>=0.1.0 # TickFlow official SDK (Issue #632, market review enhancement)
# Built-in optional AlphaSift screening engine
git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344#egg=alphasift
git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386#egg=alphasift
# Feishu
lark-oapi>=1.0.0 # Feishu API

View File

@@ -40,7 +40,7 @@ from src.llm import generation_params as llm_generation_params
logger = logging.getLogger(__name__)
DEFAULT_ALPHASIFT_INSTALL_SPEC = (
"git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344"
"git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386"
)

View File

@@ -9,12 +9,14 @@ import json
import logging
import math
import os
import re
import subprocess
import sys
import threading
from contextvars import ContextVar
from contextlib import contextmanager
from dataclasses import asdict, is_dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Dict, Iterator, List, Optional, Tuple
@@ -37,6 +39,9 @@ DSA_PRE_RANK_CONTEXT_MAX_CANDIDATES = 3
DSA_ALPHASIFT_LLM_CANDIDATE_MULTIPLIER = 2
DSA_ALPHASIFT_LLM_MAX_CANDIDATES = 12
DSA_ALPHASIFT_SNAPSHOT_SOURCE_PRIORITY = "em_datacenter,tushare,efinance,akshare_em"
DSA_ALPHASIFT_DATA_DIR = Path("data") / "alphasift"
DSA_ALPHASIFT_HOTSPOT_CACHE_PATH = DSA_ALPHASIFT_DATA_DIR / "hotspots.json"
DSA_ALPHASIFT_HOTSPOT_HISTORY_PATH = DSA_ALPHASIFT_DATA_DIR / "hotspot.history.jsonl"
_DSA_FETCHER_MANAGER_LOCK = threading.RLock()
_DSA_FETCHER_MANAGER: Any = None
_FUNDAMENTAL_BLOCKS = ("valuation", "growth", "earnings", "institution", "capital_flow", "boards")
@@ -48,6 +53,88 @@ _ALPHASIFT_LITELLM_COMPLETION_ATTR = "_alphasift_litellm_completion_bridge"
_ALPHASIFT_LITELLM_COMPLETION_LOCK = threading.Lock()
def _safe_float(value: Any) -> Optional[float]:
try:
if value is None or value == "":
return None
number = float(value)
except (TypeError, ValueError):
return None
if not math.isfinite(number):
return None
return number
def _utc_now_iso() -> str:
return datetime.now(timezone.utc).isoformat(timespec="seconds").replace("+00:00", "Z")
def _resolve_alphasift_data_dir() -> Path:
configured = _env_text(os.getenv("ALPHASIFT_DATA_DIR"))
if configured:
return Path(configured)
return DSA_ALPHASIFT_DATA_DIR
def _alphasift_hotspot_cache_path() -> Path:
if _env_text(os.getenv("ALPHASIFT_DATA_DIR")):
return _resolve_alphasift_data_dir() / "hotspots.json"
return DSA_ALPHASIFT_HOTSPOT_CACHE_PATH
def _alphasift_hotspot_history_path() -> Path:
if _env_text(os.getenv("ALPHASIFT_DATA_DIR")):
return _resolve_alphasift_data_dir() / "hotspot.history.jsonl"
return DSA_ALPHASIFT_HOTSPOT_HISTORY_PATH
def _load_alphasift_hotspot_cache(*, provider: str, top: int) -> Optional[Dict[str, Any]]:
cache_path = _alphasift_hotspot_cache_path()
try:
raw = json.loads(cache_path.read_text(encoding="utf-8"))
except FileNotFoundError:
return None
except Exception as exc:
logger.warning("Failed to read AlphaSift hotspot cache from %s: %s", cache_path, exc)
return None
payload = raw.get("payload") if isinstance(raw, dict) else None
if not isinstance(payload, dict):
return None
hotspots = payload.get("hotspots")
if not isinstance(hotspots, list) or not hotspots:
return None
selected = hotspots[: max(1, min(int(top or 12), 50))]
cached = dict(payload)
cached.update({
"enabled": True,
"provider": provider or payload.get("provider") or "akshare",
"hotspots": selected,
"hotspot_count": len(selected),
"cache_used": True,
"cached_at": raw.get("cached_at") or payload.get("cached_at"),
})
cached["source_errors"] = list(cached.get("source_errors") or [])
return _remove_non_finite_json_values(cached)
def _write_alphasift_hotspot_cache(payload: Dict[str, Any]) -> None:
cache_path = _alphasift_hotspot_cache_path()
try:
cache_path.parent.mkdir(parents=True, exist_ok=True)
cached_at = _utc_now_iso()
cache_payload = dict(payload)
cache_payload["cache_used"] = False
cache_payload["cached_at"] = cached_at
cache_path.write_text(
json.dumps({"cached_at": cached_at, "payload": cache_payload}, ensure_ascii=False, indent=2),
encoding="utf-8",
)
except Exception as exc:
logger.warning("Failed to write AlphaSift hotspot cache to %s: %s", cache_path, exc)
class AlphaSiftStrategyResponse(BaseModel):
id: str
name: str = ""
@@ -95,6 +182,102 @@ class AlphaSiftService:
_ensure_alphasift_enabled(self.config)
return _install_alphasift(self.config)
def hotspots(self, *, provider: str = "", top: int = 12, refresh: bool = False) -> Dict[str, Any]:
_ensure_alphasift_enabled(self.config)
_ensure_alphasift_available_for_use()
hotspot_module = _import_alphasift_hotspot()
discover_hotspots = _get_adapter_callable(
hotspot_module,
"discover_hotspots",
"discover_hotspots() is not callable.",
)
provider_name, provider_arg = _resolve_hotspot_provider(provider)
top_count = max(1, min(int(top or 12), 50))
if not refresh:
cached = _load_alphasift_hotspot_cache(provider=provider_name, top=top_count)
if cached is not None:
return cached
try:
with _alphasift_runtime_env(self.config):
raw = discover_hotspots(
provider=provider_arg,
top=top_count,
history_path=_alphasift_hotspot_history_path(),
fallback_cache_path=_alphasift_hotspot_cache_path(),
)
except HTTPException:
raise
except Exception as exc:
cached = _load_alphasift_hotspot_cache(provider=provider_name, top=top_count)
if cached is not None:
errors = list(cached.get("source_errors") or [])
errors.append(f"live refresh failed: {exc}")
cached["source_errors"] = errors
cached["fallback_used"] = True
cached["cache_used"] = True
return cached
raise HTTPException(
status_code=424,
detail={"error": "alphasift_hotspots_failed", "message": f"AlphaSift hotspots failed: {exc}"},
) from exc
items = _remove_non_finite_json_values(_to_plain(raw))
if not isinstance(items, list):
items = []
selected = items[:top_count]
source_errors = list(getattr(raw, "source_errors", []) or [])
if not selected and source_errors:
cached = _load_alphasift_hotspot_cache(provider=provider_name, top=top_count)
if cached is not None:
errors = list(cached.get("source_errors") or [])
errors.extend(source_errors)
cached["source_errors"] = errors
cached["fallback_used"] = True
cached["cache_used"] = True
return cached
payload = {
"enabled": True,
"provider": provider_name,
"provider_used": str(getattr(raw, "provider_used", "")),
"fallback_used": bool(getattr(raw, "fallback_used", False)),
"cache_used": False,
"cached_at": None,
"source_errors": source_errors,
"stale": bool(getattr(raw, "stale", False)),
"stale_age_hours": getattr(raw, "stale_age_hours", None),
"hotspots": selected,
"hotspot_count": len(selected),
}
if selected:
_write_alphasift_hotspot_cache(payload)
return payload
def hotspot_detail(self, *, topic: str, provider: str = "") -> Dict[str, Any]:
_ensure_alphasift_enabled(self.config)
_ensure_alphasift_available_for_use()
topic_text = _env_text(topic)
if not topic_text:
raise HTTPException(
status_code=400,
detail={"error": "alphasift_hotspot_topic_required", "message": "热点题材名称不能为空。"},
)
provider_name, provider_arg = _resolve_hotspot_provider(provider)
if not isinstance(provider_arg, DsaEastMoneyHotspotProvider):
provider_arg = DsaEastMoneyHotspotProvider()
try:
with _alphasift_runtime_env(self.config):
detail = provider_arg.hotspot_detail(topic_text)
except Exception as exc:
raise HTTPException(
status_code=424,
detail={"error": "alphasift_hotspot_detail_failed", "message": f"AlphaSift hotspot detail failed: {exc}"},
) from exc
detail["enabled"] = True
detail["provider"] = provider_name
return _remove_non_finite_json_values(detail)
def screen(self, *, strategy: str, market: str, max_results: int) -> Dict[str, Any]:
_ensure_alphasift_enabled(self.config)
_ensure_alphasift_available_for_use()
@@ -150,6 +333,13 @@ class AlphaSiftService:
"warnings": raw_data.get("warnings") or [],
"source_errors": raw_data.get("source_errors") or [],
"dsa_enrichment": dsa_enrichment,
"deep_analysis_requested": raw_data.get("deep_analysis_requested"),
"post_analyzers": raw_data.get("post_analyzers") or [],
"daily_enriched": raw_data.get("daily_enriched"),
"daily_enrich_count": raw_data.get("daily_enrich_count"),
"risk_enabled": raw_data.get("risk_enabled"),
"portfolio_diversity_enabled": raw_data.get("portfolio_diversity_enabled"),
"portfolio_concentration_notes": raw_data.get("portfolio_concentration_notes") or [],
}
@@ -364,6 +554,35 @@ def _import_alphasift() -> Any:
) from exc
def _import_alphasift_hotspot() -> Any:
try:
_prepare_alphasift_runtime_env()
return importlib.import_module("alphasift.hotspot")
except ModuleNotFoundError as exc:
if getattr(exc, "name", None) in {"alphasift", "alphasift.hotspot"}:
diagnostics = {
"reason": "missing_module",
"stage": "import_hotspot",
"error_type": exc.__class__.__name__,
"module": str(getattr(exc, "name", "alphasift.hotspot")),
}
raise _alphasift_unavailable_exception(
f"AlphaSift hotspot module is unavailable: {exc}",
diagnostics=diagnostics,
) from exc
diagnostics = _log_unexpected_alphasift_exception("import_hotspot", exc)
raise _alphasift_unavailable_exception(
f"AlphaSift hotspot module import failed: {exc}",
diagnostics=diagnostics,
) from exc
except Exception as exc:
diagnostics = _log_unexpected_alphasift_exception("import_hotspot", exc)
raise _alphasift_unavailable_exception(
f"AlphaSift hotspot module import failed: {exc}",
diagnostics=diagnostics,
) from exc
def _prepare_alphasift_runtime_env() -> None:
if os.getenv("STRATEGIES_DIR"):
return
@@ -743,9 +962,542 @@ def _build_alphasift_runtime_env(config: Config, *, max_results: Optional[int] =
put_default("LLM_CANDIDATE_MULTIPLIER", str(DSA_ALPHASIFT_LLM_CANDIDATE_MULTIPLIER))
put_default("LLM_MAX_CANDIDATES", str(_resolve_dsa_llm_max_candidates(max_results)))
put_default("SNAPSHOT_SOURCE_PRIORITY", DSA_ALPHASIFT_SNAPSHOT_SOURCE_PRIORITY)
alphasift_data_dir = _resolve_alphasift_data_dir()
put_default("ALPHASIFT_DATA_DIR", str(alphasift_data_dir))
put_default("ALPHASIFT_FALLBACK_SNAPSHOT_PATH", str(alphasift_data_dir / "snapshot.last_good.json"))
put_default("ALPHASIFT_DAILY_HISTORY_CACHE_DIR", str(alphasift_data_dir / "daily_history"))
put_default("ALPHASIFT_INDUSTRY_PROVIDER_CACHE_DIR", str(alphasift_data_dir / "industry_provider_cache"))
return env
def _resolve_hotspot_provider(provider: str) -> Tuple[str, Any]:
requested = (provider or "").strip()
if requested.lower() == "akshare":
return requested, DsaEastMoneyHotspotProvider()
if requested:
return requested, requested
configured = (os.getenv("INDUSTRY_PROVIDER") or "").strip()
if configured.lower() == "akshare":
return configured, DsaEastMoneyHotspotProvider()
return configured or "none", configured or "none"
class DsaEastMoneyHotspotProvider:
"""Minimal EastMoney board provider for AlphaSift hotspot scoring."""
_BASE_URL = "https://79.push2.eastmoney.com/api/qt/clist/get"
_COMMON_PARAMS = {
"pn": "1",
"po": "1",
"np": "1",
"ut": "bd1d9ddb04089700cf9c27f6f7426281",
"fltt": "2",
"invt": "2",
"fid": "f12",
"fields": "f3,f12,f14",
}
_BROAD_BOARD_KEYWORDS = (
"融资融券",
"深股通",
"沪股通",
"创业板",
"昨日",
"机构重仓",
"富时罗素",
"MSCI",
"标普",
"上证",
"深证",
"中证",
"HS300",
"证金",
"QFII",
"基金",
"转融券",
"预增",
"预盈",
"亏损",
"低价",
"小盘股",
"中盘股",
"百元股",
"破发",
"破增发",
"趋势股",
"广东板块",
"江苏板块",
"浙江板块",
"上海板块",
"深圳特区",
"央国企",
"国企改革",
"专精特新",
"其他",
"",
"",
)
_CHANGE_EVENT_LABELS = {
4: "快速拉升",
8: "快速回落",
16: "大幅上涨",
32: "大幅下跌",
64: "有大笔买入",
128: "有大笔卖出",
8193: "火箭发射",
8194: "高台跳水",
8201: "大笔买入",
8202: "大笔卖出",
8203: "封涨停板",
8204: "打开涨停板",
8207: "有打开跌停板",
8208: "封跌停板",
8209: "向上缺口",
8210: "向下缺口",
8211: "60日新高",
8212: "60日新低",
8213: "60日大幅上涨",
8214: "60日大幅下跌",
8215: "竞价上涨",
8216: "竞价下跌",
8217: "高开",
8218: "低开",
8219: "放量",
8220: "缩量",
8221: "向上突破",
8222: "向下破位",
}
def stock_board_concept_name_em(self) -> Any:
frame = self._fetch_board_changes_with_fallback()
if frame is not None and not frame.empty:
return frame
frame = self._fetch_rankings_with_fallback("concept")
if frame is not None and not frame.empty:
return frame
return self._fetch_board_names(source_fs="m:90 t:3 f:!50")
def stock_board_industry_name_em(self) -> Any:
concept_frame = self._fetch_board_changes_with_fallback()
if concept_frame is not None and not concept_frame.empty:
import pandas as pd
return pd.DataFrame()
frame = self._fetch_rankings_with_fallback("industry")
if frame is not None and not frame.empty:
return frame
return self._fetch_board_names(source_fs="m:90 t:2 f:!50")
def stock_board_concept_cons_em(self, symbol: str = "") -> Any:
try:
frame = self._fetch_ths_constituents(symbol)
except Exception as exc:
logger.warning(
"AlphaSift THS constituent fetch failed for %s; falling back to alternative sources: %s",
symbol,
exc,
)
frame = None
if frame is not None and not frame.empty:
return frame
frame = self._fallback_constituents(symbol)
if frame is not None and not frame.empty:
return frame
return self._fetch_eastmoney_constituents(symbol, source="concept")
def stock_board_industry_cons_em(self, symbol: str = "") -> Any:
frame = self._fetch_eastmoney_constituents(symbol, source="industry")
if frame is not None and not frame.empty:
return frame
return self._fallback_constituents(symbol)
def hotspot_detail(self, topic: str) -> Dict[str, Any]:
try:
summary = self._find_board_change(topic)
except Exception as exc:
logger.warning(
"AlphaSift board-change summary fetch failed for %s; continuing without summary: %s",
topic,
exc,
)
summary = {}
if self._is_industry_hotspot(topic):
stocks = self._normalize_constituent_records(self.stock_board_industry_cons_em(topic))
else:
stocks = self._normalize_constituent_records(self.stock_board_concept_cons_em(topic))
route = self._build_hotspot_route(topic, summary)
info = self._fetch_ths_info(topic)
if info:
route.append({
"title": "同花顺板块概况",
"description": "".join(f"{key} {value}" for key, value in list(info.items())[:4]),
"source": "ths_info",
})
if not stocks and summary:
stock_code = _env_text(summary.get("板块异动最频繁个股及所属类型-股票代码"))
stock_name = _env_text(summary.get("板块异动最频繁个股及所属类型-股票名称"))
if stock_code or stock_name:
stocks.append({
"code": stock_code,
"name": stock_name,
"role": "异动核心",
"change_pct": None,
"hot_stock_score": 60.0,
})
return {
"topic": topic,
"name": topic,
"summary": self._build_hotspot_summary(topic, summary),
"route": route,
"stocks": stocks[:30],
"stock_count": len(stocks),
"source_errors": [],
}
def _fetch_board_changes(self) -> Any:
import akshare as ak
import pandas as pd
df = ak.stock_board_change_em()
if df is None or df.empty:
return pd.DataFrame()
rows = []
for index, row in df.iterrows():
topic = _env_text(row.get("板块名称"))
if not topic or self._is_broad_board(topic):
continue
change_pct = _safe_float(row.get("涨跌幅"))
event_count = int(_safe_float(row.get("板块异动总次数")) or 0)
leader = _env_text(row.get("板块异动最频繁个股及所属类型-股票名称"))
heat_score = min(99.0, max(1.0, event_count / 120.0 + max(change_pct or 0.0, 0.0) * 9.0))
rows.append({
"name": topic,
"change_pct": change_pct,
"rank": index + 1,
"heat_score": heat_score,
"leader": leader,
"event_count": event_count,
})
rows.sort(key=lambda item: (item.get("heat_score") or 0, item.get("event_count") or 0), reverse=True)
return pd.DataFrame(rows)
def _fetch_board_changes_with_fallback(self) -> Any:
import pandas as pd
try:
return self._fetch_board_changes()
except Exception as exc:
logger.warning("AlphaSift hotspot board-change fetch failed; falling back to ranking/board names: %s", exc)
return pd.DataFrame()
def _is_broad_board(self, name: str) -> bool:
return any(keyword in name for keyword in self._BROAD_BOARD_KEYWORDS)
def _fetch_rankings(self, source: str) -> Any:
import pandas as pd
manager = _get_dsa_fetcher_manager()
fetch = manager.get_concept_rankings if source == "concept" else manager.get_sector_rankings
top, _bottom = fetch(100)
rows = []
for index, item in enumerate(top or []):
name = _env_text((item or {}).get("name"))
if not name:
continue
rows.append({
"name": name,
"change_pct": (item or {}).get("change_pct"),
"rank": index + 1,
})
return pd.DataFrame(rows)
def _fetch_rankings_with_fallback(self, source: str) -> Any:
import pandas as pd
try:
return self._fetch_rankings(source)
except Exception as exc:
logger.warning("AlphaSift hotspot %s ranking fetch failed; falling back to board names: %s", source, exc)
return pd.DataFrame()
def _fetch_board_names(self, *, source_fs: str) -> Any:
import pandas as pd
import requests
params = dict(self._COMMON_PARAMS)
params.update({"pz": "100", "fs": source_fs})
session = requests.Session()
response = session.get(
self._BASE_URL,
params=params,
timeout=20,
headers={"User-Agent": "Mozilla/5.0", "Accept": "application/json,text/plain,*/*"},
)
response.raise_for_status()
payload = response.json()
rows = ((payload.get("data") or {}).get("diff") or []) if isinstance(payload, dict) else []
normalized = [
{
"板块名称": str(row.get("f14") or "").strip(),
"涨跌幅": row.get("f3"),
"序号": index + 1,
}
for index, row in enumerate(rows)
if str(row.get("f14") or "").strip()
]
return pd.DataFrame(normalized)
def _find_board_change(self, topic: str) -> Dict[str, Any]:
import akshare as ak
df = ak.stock_board_change_em()
if df is None or df.empty:
return {}
rows = df[df["板块名称"].astype(str) == topic]
if rows.empty:
rows = df[df["板块名称"].astype(str).str.contains(re.escape(topic), case=False, na=False)]
if rows.empty:
return {}
return rows.iloc[0].to_dict()
def _is_industry_hotspot(self, topic: str) -> bool:
try:
frame = self.stock_board_industry_name_em()
except Exception as exc:
logger.warning(
"AlphaSift industry hotspot source check failed for %s; using concept constituents: %s",
topic,
exc,
)
return False
return self._board_frame_contains_topic(frame, topic)
def _board_frame_contains_topic(self, frame: Any, topic: str) -> bool:
import pandas as pd
topic_text = _env_text(topic)
if not topic_text:
return False
df = pd.DataFrame(frame)
if df.empty:
return False
for column in ("name", "板块名称", "行业名称", "名称"):
if column not in df.columns:
continue
values = df[column].map(_env_text)
if bool((values == topic_text).any()):
return True
return False
def _build_hotspot_summary(self, topic: str, summary: Dict[str, Any]) -> str:
if not summary:
return f"{topic} 当前暂无可用的板块异动摘要。"
change_pct = _safe_float(summary.get("涨跌幅"))
event_count = int(_safe_float(summary.get("板块异动总次数")) or 0)
leader = _env_text(summary.get("板块异动最频繁个股及所属类型-股票名称"))
action = _env_text(summary.get("板块异动最频繁个股及所属类型-买卖方向"))
parts = [f"{topic} 当前涨跌幅 {change_pct:.2f}%" if change_pct is not None else f"{topic} 当前有异动记录"]
if event_count:
parts.append(f"盘中异动 {event_count}")
if leader:
parts.append(f"高频异动个股为 {leader}{f'{action}' if action else ''}")
return "".join(parts) + ""
def _build_hotspot_route(self, topic: str, summary: Dict[str, Any]) -> List[Dict[str, Any]]:
route: List[Dict[str, Any]] = []
ths_event = self._fetch_ths_summary_event(topic)
if ths_event:
route.append({
"title": "题材驱动",
"description": ths_event,
"source": "ths_summary",
})
if summary:
route.append({
"title": "盘中发酵",
"description": self._build_hotspot_summary(topic, summary),
"source": "eastmoney_board_change",
})
for item in self._parse_change_events(summary.get("板块具体异动类型列表及出现次数"))[:5]:
route.append({
"title": item["label"],
"description": f"出现 {item['count']} 次,反映题材内个股盘中联动。",
"source": "eastmoney_board_change",
})
if not route:
route.append({
"title": "等待发酵",
"description": "暂未获取到明确催化事件,可继续观察涨跌幅、成交额和核心个股联动。",
"source": "fallback",
})
return route
def _parse_change_events(self, raw: Any) -> List[Dict[str, Any]]:
if isinstance(raw, str):
try:
import ast
raw = ast.literal_eval(raw)
except Exception:
raw = []
events = []
for item in raw or []:
if not isinstance(item, dict):
continue
event_type = int(_safe_float(item.get("t")) or 0)
count = int(_safe_float(item.get("ct")) or 0)
if not count:
continue
events.append({
"type": event_type,
"label": self._CHANGE_EVENT_LABELS.get(event_type, f"异动类型 {event_type}"),
"count": count,
})
return sorted(events, key=lambda item: item["count"], reverse=True)
def _fetch_ths_summary_event(self, topic: str) -> str:
import akshare as ak
try:
df = ak.stock_board_concept_summary_ths()
except Exception:
return ""
if df is None or df.empty:
return ""
if "概念名称" not in df.columns:
logger.warning(
"AlphaSift THS summary missing required column '概念名称'; skip enrichment.",
)
return ""
rows = df[df["概念名称"].astype(str) == topic]
if rows.empty:
rows = df[df["概念名称"].astype(str).str.contains(re.escape(topic), case=False, na=False)]
if rows.empty:
return ""
row = rows.iloc[0]
date = _env_text(row.get("日期"))
event = _env_text(row.get("驱动事件"))
return f"{date}{event}" if date and event else event
def _fetch_ths_info(self, topic: str) -> Dict[str, str]:
import akshare as ak
try:
df = ak.stock_board_concept_info_ths(symbol=topic)
except Exception:
return {}
if df is None or df.empty or "项目" not in df.columns or "" not in df.columns:
return {}
return {
_env_text(row.get("项目")): _env_text(row.get(""))
for _, row in df.iterrows()
if _env_text(row.get("项目"))
}
def _fetch_eastmoney_constituents(self, topic: str, *, source: str) -> Any:
import akshare as ak
try:
if source == "industry":
return ak.stock_board_industry_cons_em(symbol=topic)
return ak.stock_board_concept_cons_em(symbol=topic)
except Exception:
return None
def _fetch_ths_constituents(self, topic: str) -> Any:
import pandas as pd
import requests
code = self._resolve_ths_concept_code(topic)
if not code:
return pd.DataFrame()
url = f"http://q.10jqka.com.cn/gn/detail/code/{code}/"
response = requests.get(
url,
headers={"User-Agent": "Mozilla/5.0", "Referer": "http://q.10jqka.com.cn/gn/"},
timeout=15,
)
response.raise_for_status()
html = response.content.decode("gbk", "ignore")
rows = []
seen = set()
for match in re.finditer(r">(\d{6})<.*?>([^<>\n]{2,12})<", html, re.S):
code_text = match.group(1)
name_text = re.sub(r"\s+", "", match.group(2))
if code_text in seen or not name_text or re.search(r"\d", name_text):
continue
seen.add(code_text)
rows.append({"code": code_text, "name": name_text})
if len(rows) >= 80:
break
return pd.DataFrame(rows)
def _resolve_ths_concept_code(self, topic: str) -> str:
import akshare as ak
try:
df = ak.stock_board_concept_name_ths()
except Exception:
return ""
if df is None or df.empty:
return ""
rows = df[df["name"].astype(str) == topic]
if rows.empty:
rows = df[df["name"].astype(str).str.contains(re.escape(topic), case=False, na=False)]
if rows.empty and topic.endswith("概念"):
base = topic[:-2]
rows = df[df["name"].astype(str).str.contains(re.escape(base), case=False, na=False)]
if rows.empty:
return ""
return _env_text(rows.iloc[0].get("code"))
def _fallback_constituents(self, topic: str) -> Any:
import pandas as pd
try:
summary = self._find_board_change(topic)
except Exception as exc:
logger.warning(
"AlphaSift board-change constituent fallback failed for %s; trying other sources: %s",
topic,
exc,
)
return pd.DataFrame()
code = _env_text(summary.get("板块异动最频繁个股及所属类型-股票代码"))
name = _env_text(summary.get("板块异动最频繁个股及所属类型-股票名称"))
if not code and not name:
return pd.DataFrame()
return pd.DataFrame([{
"code": code,
"name": name,
"change_pct": None,
"hot_stock_score": 60.0,
}])
def _normalize_constituent_records(self, frame: Any) -> List[Dict[str, Any]]:
import pandas as pd
df = pd.DataFrame(frame)
if df.empty:
return []
records = []
for _, row in df.iterrows():
code = _env_text(row.get("code") or row.get("代码") or row.get("证券代码"))
name = _env_text(row.get("name") or row.get("名称") or row.get("股票名称"))
if not code and not name:
continue
records.append({
"code": code,
"name": name,
"change_pct": _safe_float(row.get("change_pct") or row.get("涨跌幅") or row.get("涨幅")),
"amount": _safe_float(row.get("amount") or row.get("成交额") or row.get("成交金额")),
"turnover_rate": _safe_float(row.get("turnover_rate") or row.get("换手率")),
"volume_ratio": _safe_float(row.get("volume_ratio") or row.get("量比")),
"role": _env_text(row.get("role")) or "概念股",
"hot_stock_score": _safe_float(row.get("hot_stock_score")) or 0.0,
})
return records
def _build_alphasift_context(config: Config, *, max_results: Optional[int] = None) -> Dict[str, Any]:
# context.llm.model/fallback/model_list 与 LiteLLM 路由语义保持一致,
# 参见 https://docs.litellm.ai/docs/proxy/configs#the-model_list-key

View File

@@ -4,14 +4,18 @@
from __future__ import annotations
import os
import json
import sys
import tempfile
import unittest
from pathlib import Path
from types import ModuleType, SimpleNamespace
from typing import Any, Dict
from unittest.mock import ANY, MagicMock, patch
import threading
from fastapi import HTTPException
from fastapi import FastAPI, HTTPException
from fastapi.testclient import TestClient
try:
import litellm # noqa: F401
@@ -62,8 +66,11 @@ def _missing_alphasift_module_diagnostics() -> Dict[str, str]:
class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
def setUp(self) -> None:
Config.reset_instance()
self.env_patch = patch.dict(os.environ, {"ALPHASIFT_DATA_DIR": ""}, clear=False)
self.env_patch.start()
def tearDown(self) -> None:
self.env_patch.stop()
Config.reset_instance()
def _config(self, *, enabled: bool, install_spec: str = DEFAULT_ALPHASIFT_TEST_SPEC) -> Config:
@@ -102,6 +109,12 @@ class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
def _strategies(self, config: Config):
return alphasift_endpoint.alphasift_strategies(request=self._request(), config=config)
def _hotspots(self, config: Config, **kwargs):
return alphasift_endpoint.alphasift_hotspots(config=config, **kwargs)
def _hotspot_detail(self, config: Config, **kwargs):
return alphasift_endpoint.alphasift_hotspot_detail(config=config, **kwargs)
def test_default_install_spec_is_commit_pinned(self) -> None:
self.assertRegex(
DEFAULT_ALPHASIFT_TEST_SPEC,
@@ -254,6 +267,407 @@ class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
self.assertEqual(payload["strategies"][0]["name"], "双低选股")
self.assertEqual(payload["strategies"][1]["name"], "趋势质量")
def test_hotspots_returns_alphasift_hotspot_summaries(self) -> None:
config = self._config(enabled=True)
class HotspotRows(list):
provider_used = "akshare"
fallback_used = False
source_errors = []
stale = False
stale_age_hours = None
rows = HotspotRows([
{
"topic": "AI算力",
"name": "AI算力",
"heat_score": 88.0,
"stage": "加速主升",
"leaders": ["中际旭创"],
}
])
discover = MagicMock(return_value=rows)
with (
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._import_alphasift_hotspot", return_value=SimpleNamespace(discover_hotspots=discover)),
):
payload = self._hotspots(config=config, provider="akshare", top=6)
self.assertEqual(payload["enabled"], True)
self.assertEqual(payload["provider"], "akshare")
self.assertEqual(payload["provider_used"], "akshare")
self.assertEqual(payload["hotspot_count"], 1)
self.assertEqual(payload["hotspots"][0]["topic"], "AI算力")
self.assertEqual(payload["hotspots"][0]["heat_score"], 88.0)
discover.assert_called_once()
provider = discover.call_args.kwargs["provider"]
self.assertTrue(hasattr(provider, "stock_board_concept_name_em"))
self.assertTrue(hasattr(provider, "stock_board_industry_name_em"))
self.assertEqual(discover.call_args.kwargs["top"], 6)
def test_hotspots_uses_last_success_cache_by_default(self) -> None:
config = self._config(enabled=True)
with tempfile.TemporaryDirectory() as tmpdir:
cache_path = Path(tmpdir) / "hotspots.json"
cache_path.write_text(
json.dumps({
"cached_at": "2026-06-07T12:00:00Z",
"payload": {
"enabled": True,
"provider": "akshare",
"provider_used": "DsaEastMoneyHotspotProvider",
"fallback_used": False,
"cache_used": False,
"cached_at": "2026-06-07T12:00:00Z",
"source_errors": [],
"hotspots": [
{"topic": "玻璃基板", "heat_score": 88.0},
{"topic": "机器人执行器", "heat_score": 80.0},
],
"hotspot_count": 2,
},
}),
encoding="utf-8",
)
discover = MagicMock()
with (
patch("src.services.alphasift_service.DSA_ALPHASIFT_HOTSPOT_CACHE_PATH", cache_path),
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._import_alphasift_hotspot", return_value=SimpleNamespace(discover_hotspots=discover)),
):
payload = self._hotspots(config=config, provider="akshare", top=1, refresh=False)
self.assertEqual(payload["cache_used"], True)
self.assertEqual(payload["cached_at"], "2026-06-07T12:00:00Z")
self.assertEqual(payload["hotspot_count"], 1)
self.assertEqual(payload["hotspots"][0]["topic"], "玻璃基板")
discover.assert_not_called()
def test_hotspots_refresh_falls_back_to_cache_when_provider_returns_only_errors(self) -> None:
config = self._config(enabled=True)
class HotspotRows(list):
provider_used = "akshare"
fallback_used = False
source_errors = ["akshare returned no usable board rows"]
stale = False
stale_age_hours = None
rows = HotspotRows()
discover = MagicMock(return_value=rows)
with tempfile.TemporaryDirectory() as tmpdir:
cache_path = Path(tmpdir) / "hotspots.json"
cache_path.write_text(
json.dumps({
"cached_at": "2026-06-07T12:00:00Z",
"payload": {
"enabled": True,
"provider": "akshare",
"provider_used": "DsaEastMoneyHotspotProvider",
"fallback_used": False,
"cache_used": False,
"cached_at": "2026-06-07T12:00:00Z",
"source_errors": [],
"hotspots": [
{"topic": "MLCC", "heat_score": 91.0},
],
"hotspot_count": 1,
},
}),
encoding="utf-8",
)
with (
patch("src.services.alphasift_service.DSA_ALPHASIFT_HOTSPOT_CACHE_PATH", cache_path),
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._import_alphasift_hotspot", return_value=SimpleNamespace(discover_hotspots=discover)),
):
payload = self._hotspots(config=config, provider="akshare", top=12, refresh=True)
self.assertEqual(payload["cache_used"], True)
self.assertEqual(payload["fallback_used"], True)
self.assertEqual(payload["hotspot_count"], 1)
self.assertEqual(payload["hotspots"][0]["topic"], "MLCC")
self.assertIn("akshare returned no usable board rows", payload["source_errors"])
discover.assert_called_once()
def test_hotspots_respects_custom_alphasift_data_dir_for_cache_paths(self) -> None:
config = self._config(enabled=True)
class HotspotRows(list):
provider_used = "akshare"
fallback_used = False
source_errors = []
stale = False
stale_age_hours = None
rows = HotspotRows([
{"topic": "机器人执行器", "heat_score": 86.0},
])
captured: Dict[str, Any] = {}
def discover(**kwargs):
captured.update(kwargs)
return rows
with tempfile.TemporaryDirectory() as tmpdir:
data_dir = Path(tmpdir) / "persistent-alphasift"
cache_path = data_dir / "hotspots.json"
history_path = data_dir / "hotspot.history.jsonl"
with (
patch.dict(os.environ, {"ALPHASIFT_DATA_DIR": str(data_dir)}, clear=False),
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch(
"src.services.alphasift_service._import_alphasift_hotspot",
return_value=SimpleNamespace(discover_hotspots=discover),
),
):
payload = self._hotspots(config=config, provider="akshare", top=3, refresh=True)
self.assertEqual(payload["hotspots"][0]["topic"], "机器人执行器")
self.assertEqual(captured["history_path"], history_path)
self.assertEqual(captured["fallback_cache_path"], cache_path)
self.assertTrue(cache_path.exists())
discover_again = MagicMock()
with (
patch.dict(os.environ, {"ALPHASIFT_DATA_DIR": str(data_dir)}, clear=False),
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch(
"src.services.alphasift_service._import_alphasift_hotspot",
return_value=SimpleNamespace(discover_hotspots=discover_again),
),
):
cached = self._hotspots(config=config, provider="akshare", top=1, refresh=False)
self.assertEqual(cached["cache_used"], True)
self.assertEqual(cached["hotspots"][0]["topic"], "机器人执行器")
discover_again.assert_not_called()
def test_hotspot_detail_returns_route_and_concept_stocks(self) -> None:
config = self._config(enabled=True)
class FakeProvider(alphasift_service.DsaEastMoneyHotspotProvider):
def hotspot_detail(self, topic: str) -> Dict[str, Any]:
return {
"topic": topic,
"summary": f"{topic} 盘中发酵。",
"route": [{"title": "盘中发酵", "description": "出现大笔买入。"}],
"stocks": [{"code": "920438", "name": "戈碧迦", "role": "异动核心"}],
"stock_count": 1,
"source_errors": [],
}
with (
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._resolve_hotspot_provider", return_value=("akshare", FakeProvider())),
):
payload = self._hotspot_detail(config=config, provider="akshare", topic="玻璃基板")
self.assertEqual(payload["enabled"], True)
self.assertEqual(payload["provider"], "akshare")
self.assertEqual(payload["topic"], "玻璃基板")
self.assertEqual(payload["route"][0]["title"], "盘中发酵")
self.assertEqual(payload["stocks"][0]["name"], "戈碧迦")
def test_hotspot_detail_route_accepts_slash_containing_topic(self) -> None:
config = self._config(enabled=True)
app = FastAPI()
app.include_router(alphasift_endpoint.router, prefix="/api/v1/alphasift")
app.dependency_overrides[alphasift_endpoint.get_config_dep] = lambda: config
service = MagicMock()
service.hotspot_detail.return_value = {
"enabled": True,
"provider": "akshare",
"topic": "DRG/DIP",
"route": [],
"stocks": [],
"stock_count": 0,
}
with patch("api.v1.endpoints.alphasift._service", return_value=service):
response = TestClient(app).get("/api/v1/alphasift/hotspots/DRG%2FDIP?provider=akshare")
self.assertEqual(response.status_code, 200)
self.assertEqual(response.json()["topic"], "DRG/DIP")
service.hotspot_detail.assert_called_once_with(topic="DRG/DIP", provider="akshare")
def test_hotspot_detail_falls_back_when_ths_constituents_fail(self) -> None:
import pandas as pd
config = self._config(enabled=True)
class FakeProvider(alphasift_service.DsaEastMoneyHotspotProvider):
def _fetch_ths_constituents(self, topic: str) -> Any:
raise TimeoutError("ths timeout")
def _fallback_constituents(self, topic: str) -> Any:
return pd.DataFrame([{
"code": "300000",
"name": "中际旭创",
"change_pct": None,
"hot_stock_score": 60.0,
}])
def _fetch_eastmoney_constituents(self, topic: str, *, source: str) -> Any:
return pd.DataFrame()
def _find_board_change(self, topic: str) -> Dict[str, Any]:
return {}
def _build_hotspot_route(self, topic: str, summary: Dict[str, Any]) -> Any:
return [{"title": "fallback", "description": topic, "source": "test"}]
def _fetch_ths_info(self, topic: str) -> Dict[str, str]:
return {}
with (
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._resolve_hotspot_provider", return_value=("akshare", FakeProvider())),
):
payload = self._hotspot_detail(config=config, provider="akshare", topic="AI算力")
self.assertEqual(payload["enabled"], True)
self.assertEqual(payload["provider"], "akshare")
self.assertEqual(payload["topic"], "AI算力")
self.assertEqual(payload["stocks"][0]["name"], "中际旭创")
self.assertEqual(payload["route"][0]["title"], "fallback")
def test_hotspot_detail_uses_constituent_fallback_when_board_change_summary_fails(self) -> None:
import pandas as pd
config = self._config(enabled=True)
class FakeProvider(alphasift_service.DsaEastMoneyHotspotProvider):
def _find_board_change(self, topic: str) -> Dict[str, Any]:
raise TimeoutError("board change timeout")
def _fetch_ths_constituents(self, topic: str) -> Any:
return pd.DataFrame()
def _fetch_eastmoney_constituents(self, topic: str, *, source: str) -> Any:
return pd.DataFrame([{
"代码": "002138",
"名称": "顺络电子",
"涨跌幅": 3.2,
}])
def _fetch_ths_summary_event(self, topic: str) -> str:
return "需求升温"
def _fetch_ths_info(self, topic: str) -> Dict[str, str]:
return {}
with (
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._resolve_hotspot_provider", return_value=("akshare", FakeProvider())),
):
payload = self._hotspot_detail(config=config, provider="akshare", topic="MLCC")
self.assertEqual(payload["enabled"], True)
self.assertEqual(payload["topic"], "MLCC")
self.assertEqual(payload["summary"], "MLCC 当前暂无可用的板块异动摘要。")
self.assertEqual(payload["route"][0]["source"], "ths_summary")
self.assertEqual(payload["stocks"][0]["name"], "顺络电子")
def test_hotspot_detail_uses_industry_constituents_for_industry_hotspots(self) -> None:
import pandas as pd
config = self._config(enabled=True)
class FakeProvider(alphasift_service.DsaEastMoneyHotspotProvider):
def __init__(self) -> None:
self.constituent_sources = []
def stock_board_industry_name_em(self) -> Any:
return pd.DataFrame([{"name": "电池", "rank": 1}])
def _fetch_eastmoney_constituents(self, topic: str, *, source: str) -> Any:
self.constituent_sources.append(source)
if source == "industry":
return pd.DataFrame([{
"代码": "300750",
"名称": "宁德时代",
"涨跌幅": 2.6,
}])
return pd.DataFrame()
def _fetch_ths_constituents(self, topic: str) -> Any:
raise AssertionError("industry hotspots must not use concept constituents")
def _find_board_change(self, topic: str) -> Dict[str, Any]:
return {}
def _fetch_ths_summary_event(self, topic: str) -> str:
return ""
def _fetch_ths_info(self, topic: str) -> Dict[str, str]:
return {}
provider = FakeProvider()
with (
patch("src.services.alphasift_service._get_alphasift_status_snapshot", return_value=({}, True, {})),
patch("src.services.alphasift_service._resolve_hotspot_provider", return_value=("akshare", provider)),
):
payload = self._hotspot_detail(config=config, provider="akshare", topic="电池")
self.assertEqual(payload["enabled"], True)
self.assertEqual(payload["topic"], "电池")
self.assertEqual(payload["stocks"][0]["name"], "宁德时代")
self.assertEqual(provider.constituent_sources, ["industry"])
def test_hotspot_provider_uses_board_name_fallback_when_rankings_fail(self) -> None:
import pandas as pd
provider = alphasift_service.DsaEastMoneyHotspotProvider()
fallback = pd.DataFrame([{"板块名称": "玻璃基板", "涨跌幅": 1.8, "序号": 1}])
with (
patch.object(provider, "_fetch_board_changes", return_value=pd.DataFrame()),
patch.object(provider, "_fetch_rankings", side_effect=RuntimeError("ranking schema changed")),
patch.object(provider, "_fetch_board_names", return_value=fallback) as fetch_board_names,
):
concept = provider.stock_board_concept_name_em()
industry = provider.stock_board_industry_name_em()
self.assertEqual(concept.iloc[0]["板块名称"], "玻璃基板")
self.assertEqual(industry.iloc[0]["板块名称"], "玻璃基板")
fetch_board_names.assert_any_call(source_fs="m:90 t:3 f:!50")
fetch_board_names.assert_any_call(source_fs="m:90 t:2 f:!50")
def test_hotspot_provider_continues_fallback_when_board_change_fails(self) -> None:
import pandas as pd
provider = alphasift_service.DsaEastMoneyHotspotProvider()
rankings = pd.DataFrame([{"name": "减速器", "change_pct": 2.2, "rank": 1}])
with (
patch.object(provider, "_fetch_board_changes", side_effect=RuntimeError("akshare timeout")),
patch.object(provider, "_fetch_rankings", return_value=rankings) as fetch_rankings,
):
concept = provider.stock_board_concept_name_em()
self.assertEqual(concept.iloc[0]["name"], "减速器")
fetch_rankings.assert_called_once_with("concept")
def test_fetch_ths_summary_event_ignores_missing_concept_name_column(self) -> None:
import pandas as pd
provider = alphasift_service.DsaEastMoneyHotspotProvider()
summary = pd.DataFrame([
{"日期": "2026-06-07", "驱动事件": "行业政策利好"},
])
class _MockAkshare:
@staticmethod
def stock_board_concept_summary_ths():
return summary
with patch.dict("sys.modules", {"akshare": _MockAkshare()}):
text = provider._fetch_ths_summary_event("MLCC")
self.assertEqual(text, "")
def test_strategies_rejects_when_enabled_but_adapter_missing(self) -> None:
config = self._config(enabled=True)
@@ -541,6 +955,13 @@ class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
"llm_coverage": 1.0,
"warnings": ["fallback"],
"source_errors": [],
"deep_analysis_requested": False,
"post_analyzers": ["scorecard"],
"daily_enriched": True,
"daily_enrich_count": 12,
"risk_enabled": True,
"portfolio_diversity_enabled": True,
"portfolio_concentration_notes": ["sector concentration adjusted"],
"candidates": [
{
"code": "600519",
@@ -578,6 +999,10 @@ class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
self.assertEqual(payload["llm_coverage"], 1.0)
self.assertEqual(payload["warnings"], ["fallback"])
self.assertEqual(payload["candidate_count"], 1)
self.assertEqual(payload["post_analyzers"], ["scorecard"])
self.assertEqual(payload["daily_enriched"], True)
self.assertEqual(payload["daily_enrich_count"], 12)
self.assertEqual(payload["portfolio_concentration_notes"], ["sector concentration adjusted"])
self.assertEqual(payload["candidates"][0]["code"], "600519")
self.assertEqual(payload["candidates"][0]["llm_score"], 90.0)
self.assertEqual(payload["candidates"][0]["llm_thesis"], "LLM likes the setup")
@@ -912,6 +1337,10 @@ class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
"LLM_CANDIDATE_MULTIPLIER": alphasift_service.os.environ.get("LLM_CANDIDATE_MULTIPLIER"),
"LLM_MAX_CANDIDATES": alphasift_service.os.environ.get("LLM_MAX_CANDIDATES"),
"SNAPSHOT_SOURCE_PRIORITY": alphasift_service.os.environ.get("SNAPSHOT_SOURCE_PRIORITY"),
"ALPHASIFT_DATA_DIR": alphasift_service.os.environ.get("ALPHASIFT_DATA_DIR"),
"ALPHASIFT_FALLBACK_SNAPSHOT_PATH": alphasift_service.os.environ.get("ALPHASIFT_FALLBACK_SNAPSHOT_PATH"),
"ALPHASIFT_DAILY_HISTORY_CACHE_DIR": alphasift_service.os.environ.get("ALPHASIFT_DAILY_HISTORY_CACHE_DIR"),
"ALPHASIFT_INDUSTRY_PROVIDER_CACHE_DIR": alphasift_service.os.environ.get("ALPHASIFT_INDUSTRY_PROVIDER_CACHE_DIR"),
}
captured["context"] = kwargs.get("context")
return {"candidates": []}
@@ -948,6 +1377,19 @@ class AlphaSiftOpportunitiesApiTestCase(unittest.TestCase):
self.assertEqual(runtime_env["LLM_CANDIDATE_MULTIPLIER"], "2")
self.assertEqual(runtime_env["LLM_MAX_CANDIDATES"], "10")
self.assertEqual(runtime_env["SNAPSHOT_SOURCE_PRIORITY"], "em_datacenter,tushare,efinance,akshare_em")
self.assertEqual(runtime_env["ALPHASIFT_DATA_DIR"], str(alphasift_service.DSA_ALPHASIFT_DATA_DIR))
self.assertEqual(
runtime_env["ALPHASIFT_FALLBACK_SNAPSHOT_PATH"],
str(alphasift_service.DSA_ALPHASIFT_DATA_DIR / "snapshot.last_good.json"),
)
self.assertEqual(
runtime_env["ALPHASIFT_DAILY_HISTORY_CACHE_DIR"],
str(alphasift_service.DSA_ALPHASIFT_DATA_DIR / "daily_history"),
)
self.assertEqual(
runtime_env["ALPHASIFT_INDUSTRY_PROVIDER_CACHE_DIR"],
str(alphasift_service.DSA_ALPHASIFT_DATA_DIR / "industry_provider_cache"),
)
context = captured["context"]
self.assertIsInstance(context, dict)
self.assertEqual(context["llm"]["model"], "gemini/gemini-2.5-flash")

View File

@@ -29,7 +29,7 @@ def test_dockerfile_bundles_default_alphasift_adapter() -> None:
requirements = (REPO_ROOT / "requirements.txt").read_text(encoding="utf-8")
assert "git \\" in dockerfile
assert "git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344#egg=alphasift" in requirements
assert "git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386#egg=alphasift" in requirements
assert "pip install --no-cache-dir -r requirements.txt" in dockerfile
assert "import alphasift.dsa_adapter" in dockerfile

View File

@@ -3,6 +3,7 @@
import logging
import os
import socket
import tempfile
import unittest
from datetime import datetime
@@ -14,9 +15,17 @@ from tests.litellm_stub import ensure_litellm_stub
ensure_litellm_stub()
_ENV_BEFORE_MAIN_IMPORT = dict(os.environ)
import main
from src.config import Config
_MAIN_IMPORT_ENV_ADDITIONS = frozenset(set(os.environ) - set(_ENV_BEFORE_MAIN_IMPORT))
_MAIN_IMPORT_ENV_OVERRIDES = {
key: value
for key, value in _ENV_BEFORE_MAIN_IMPORT.items()
if os.environ.get(key) != value
}
class _DummyConfig(SimpleNamespace):
def validate(self):
@@ -51,6 +60,10 @@ class MainScheduleModeTestCase(unittest.TestCase):
os.chdir(self.original_cwd)
Config.reset_instance()
self.env_patch.stop()
for key in _MAIN_IMPORT_ENV_ADDITIONS:
os.environ.pop(key, None)
for key, value in _MAIN_IMPORT_ENV_OVERRIDES.items():
os.environ[key] = value
self.temp_dir.cleanup()
def _make_args(self, **overrides):
@@ -112,6 +125,25 @@ class MainScheduleModeTestCase(unittest.TestCase):
warning_log.assert_not_called()
def test_start_api_server_fails_before_thread_when_port_is_busy(self) -> None:
config = self._make_config(log_level="INFO")
class BusySocket:
def bind(self, address):
raise OSError("address already in use")
def close(self):
pass
with patch("socket.socket", return_value=BusySocket()) as socket_factory, \
patch("threading.Thread") as thread_cls:
with self.assertRaises(RuntimeError) as caught:
main.start_api_server("127.0.0.1", 8000, config)
socket_factory.assert_called_once_with(socket.AF_INET, socket.SOCK_STREAM)
self.assertIn("127.0.0.1:8000", str(caught.exception))
thread_cls.assert_not_called()
def test_schedule_mode_ignores_cli_stock_snapshot(self) -> None:
args = self._make_args(schedule=True, stocks="600519,000001")
config = self._make_config(schedule_enabled=False)
@@ -309,6 +341,99 @@ class MainScheduleModeTestCase(unittest.TestCase):
start_api_server.assert_not_called()
run_full_analysis.assert_not_called()
def test_serve_mode_exits_when_api_server_start_fails(self) -> None:
args = self._make_args(serve_only=True, host="127.0.0.1", port=8000)
config = self._make_config(webui_enabled=False)
with patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False), \
patch("main.parse_arguments", return_value=args), \
patch("main.get_config", return_value=config), \
patch("main.prepare_webui_frontend_assets", return_value=True), \
patch("main.start_api_server", side_effect=RuntimeError("port busy")), \
patch("main.start_bot_stream_clients") as start_bots, \
patch("main.logger.error") as error_log:
exit_code = main.main()
self.assertEqual(exit_code, 1)
start_bots.assert_not_called()
error_log.assert_called_once()
def test_webui_only_maps_to_serve_only_and_exits_when_api_server_start_fails(self) -> None:
args = self._make_args(webui_only=True, host="127.0.0.1", port=8000)
config = self._make_config(webui_enabled=False)
with patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False), \
patch("main.parse_arguments", return_value=args), \
patch("main.get_config", return_value=config), \
patch("main.prepare_webui_frontend_assets", return_value=True), \
patch("main.start_api_server", side_effect=RuntimeError("port busy")), \
patch("main.start_bot_stream_clients") as start_bots, \
patch("main.run_full_analysis") as run_full_analysis, \
patch("main.logger.error") as error_log:
exit_code = main.main()
self.assertEqual(exit_code, 1)
start_bots.assert_not_called()
run_full_analysis.assert_not_called()
error_log.assert_called_once()
def test_serve_mode_continues_single_analysis_when_api_server_start_fails(self) -> None:
args = self._make_args(serve=True, host="127.0.0.1", port=8000)
config = self._make_config(webui_enabled=False, run_immediately=True)
with patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False), \
patch("main.parse_arguments", return_value=args), \
patch("main.get_config", return_value=config), \
patch("main.prepare_webui_frontend_assets", return_value=True), \
patch("main.start_api_server", side_effect=RuntimeError("port busy")), \
patch("main.start_bot_stream_clients") as start_bots, \
patch("main.run_full_analysis") as run_full_analysis, \
patch("main.logger.error") as error_log:
exit_code = main.main()
self.assertEqual(exit_code, 0)
start_bots.assert_not_called()
run_full_analysis.assert_called_once_with(config, args, None)
error_log.assert_called_once()
def test_serve_schedule_mode_continues_scheduler_when_api_server_start_fails(self) -> None:
args = self._make_args(serve=True, schedule=True, host="127.0.0.1", port=8000)
config = self._make_config(webui_enabled=False, schedule_enabled=False)
scheduled_call = {}
def fake_run_with_schedule(
task,
schedule_time,
run_immediately,
background_tasks=None,
schedule_time_provider=None,
):
scheduled_call["schedule_time"] = schedule_time
scheduled_call["run_immediately"] = run_immediately
scheduled_call["background_tasks"] = background_tasks or []
task()
with patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False), \
patch("main.parse_arguments", return_value=args), \
patch("main.get_config", return_value=config), \
patch("main._reload_runtime_config", return_value=config), \
patch("main._build_schedule_time_provider", return_value=lambda: "18:00"), \
patch("main.prepare_webui_frontend_assets", return_value=True), \
patch("main.start_api_server", side_effect=RuntimeError("port busy")), \
patch("main.start_bot_stream_clients") as start_bots, \
patch("main.run_full_analysis") as run_full_analysis, \
patch("src.scheduler.run_with_schedule", side_effect=fake_run_with_schedule), \
patch("main.logger.error") as error_log:
exit_code = main.main()
self.assertEqual(exit_code, 0)
start_bots.assert_not_called()
run_full_analysis.assert_called_once_with(config, args, None)
self.assertEqual(scheduled_call["schedule_time"], "18:00")
self.assertEqual(scheduled_call["run_immediately"], True)
self.assertEqual(scheduled_call["background_tasks"], [])
error_log.assert_called_once()
def test_reload_runtime_config_preserves_process_env_overrides(self) -> None:
self.env_path.write_text(
"OPENAI_API_KEY=stale-file\nSCHEDULE_TIME=09:30\n",

View File

@@ -725,9 +725,14 @@ class SystemConfigServiceTestCase(unittest.TestCase):
"LITELLM_MODEL=openai/gpt-4o-mini",
"AGENT_LITELLM_MODEL=openai/gpt-4o",
"OPENAI_BASE_URL=https://api.openai.com/v1",
"LLM_CHANNELS=openai",
"LLM_OPENAI_PROTOCOL=openai",
"LLM_OPENAI_BASE_URL=https://api.openai.com/v1",
"LLM_OPENAI_API_KEYS=legacy-openai-secret",
"LLM_OPENAI_MODELS=openai/gpt-4o-mini,openai/gpt-4o",
"LITELLM_FALLBACK_MODELS=openai/gpt-4o-mini,openai/gpt-4o",
"ALPHASIFT_ENABLED=false",
"ALPHASIFT_INSTALL_SPEC=git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344",
"ALPHASIFT_INSTALL_SPEC=git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386",
"GEMINI_API_KEY=legacy-secret",
)
@@ -751,12 +756,17 @@ class SystemConfigServiceTestCase(unittest.TestCase):
self.assertEqual(current_map["ALPHASIFT_ENABLED"], "true")
self.assertEqual(
current_map["ALPHASIFT_INSTALL_SPEC"],
"git+https://github.com/ZhuLinsen/alphasift.git@1a0ed8c99b3615c0cb1076e6029827ffc6de2344",
"git+https://github.com/ZhuLinsen/alphasift.git@de54ea0da367be85770d9589a5bf7ded4f62d386",
)
self.assertEqual(current_map["GEMINI_API_KEY"], "legacy-secret")
self.assertEqual(current_map["LITELLM_MODEL"], "openai/gpt-4o-mini")
self.assertEqual(current_map["AGENT_LITELLM_MODEL"], "openai/gpt-4o")
self.assertEqual(current_map["OPENAI_BASE_URL"], "https://api.openai.com/v1")
self.assertEqual(current_map["LLM_CHANNELS"], "openai")
self.assertEqual(current_map["LLM_OPENAI_PROTOCOL"], "openai")
self.assertEqual(current_map["LLM_OPENAI_BASE_URL"], "https://api.openai.com/v1")
self.assertEqual(current_map["LLM_OPENAI_API_KEYS"], "legacy-openai-secret")
self.assertEqual(current_map["LLM_OPENAI_MODELS"], "openai/gpt-4o-mini,openai/gpt-4o")
self.assertEqual(current_map["LITELLM_FALLBACK_MODELS"], "openai/gpt-4o-mini,openai/gpt-4o")
def test_validate_reports_invalid_time(self) -> None: