Skip to content

Latest commit

 

History

History
313 lines (233 loc) · 14.6 KB

File metadata and controls

313 lines (233 loc) · 14.6 KB

api/ — FastAPI REST API

职责

FastAPI 应用工厂,挂载路由、中间件、依赖注入。提供 9 个 API 端点模块,含 SSE 实时任务流。

文件清单

根目录

文件 职责
app.py create_app() 应用工厂;lifespan 创建 long_pool + io_pool 两个独立线程池;CORS / 认证中间件挂载
deps.py 依赖注入(Config / Session / DataManager / Services)
__init__.py 导出

middlewares/

文件 职责
auth.py 认证中间件(ADMIN_AUTH_ENABLED=true 时启用)
error_handler.py 全局异常处理

v1/endpoints/(10 个端点模块)

文件 职责
analysis.py 分析任务(提交/状态/结果/SSE 实时流)
agent.py Agent 分析 + 多轮对话(流式)
history.py 分析历史(查询/详情/删除/统计)
portfolio.py 组合管理(账户/交易/持仓/CSV 导入)
backtest.py 回测(运行/结果/统计)
system_config.py 系统配置(读取/修改/校验/LLM 状态)
stocks.py 股票数据(行情/搜索)
sectors.py 行业列表 / 热力图 / 美股指数(首页 dashboard)
auth.py 可选认证(登录/登出/会话检查)
usage.py LLM 用量审计

v1/schemas/

每个 endpoints 都有对应的 Pydantic schema,分组为 analysis / backtest / common / history / portfolio / stocks / system_config / usage

应用工厂

# api/app.py
@asynccontextmanager
async def app_lifespan(app: FastAPI):
    app.state.system_config_service = SystemConfigService()

    # 长任务池(LLM/Agent 流),默认 4
    app.state.long_pool = ThreadPoolExecutor(
        max_workers=int(os.environ.get("API_LONG_POOL_WORKERS", "4")),
        thread_name_prefix="api-long",
    )
    # 轻 IO 池,默认 16
    app.state.io_pool = ThreadPoolExecutor(
        max_workers=int(os.environ.get("API_IO_POOL_WORKERS", "16")),
        thread_name_prefix="api-io",
    )
    yield
    # cleanup ...

注意:lifespan 关闭时会 shutdown(wait=False) 两个池。

def create_app() -> FastAPI:
    app = FastAPI(...)
    # CORS
    # add_auth_middleware(app)
    # app.include_router(api_v1_router)
    # add_error_handlers(app)
    # @app.get("/")  -> JSON 描述(service / version / docs / health)
    # @app.get("/api/health") -> HealthResponse
    return app

app = create_app()  # 模块级实例

依赖注入(deps.py)

def get_config() -> Config: ...
def get_session() -> Session: ...
def get_data_manager() -> DataFetcherManager: ...
def get_task_queue() -> AnalysisTaskQueue: ...
def get_history_service() -> HistoryService: ...
def get_portfolio_service() -> PortfolioService: ...
def get_current_user(request: Request) -> Optional[dict]: ...

主要端点

分析(analysis.py)

方法 路径 说明
POST /api/v1/analysis/submit 提交分析任务(异步,via TaskQueue)
GET /api/v1/analysis/status/{task_id} 查询任务状态
GET /api/v1/analysis/result/{task_id} 获取任务结果
GET /api/v1/analysis/stream/{task_id} SSE 实时任务流
POST /api/v1/analysis/cancel/{task_id} 取消任务

Agent(agent.py)

方法 路径 说明
POST /api/v1/agent/analyze 提交 Agent 分析(用 long_pool)
GET /api/v1/agent/status/{task_id} 状态
POST /api/v1/agent/chat 多轮对话
GET /api/v1/agent/chat/stream 流式对话(SSE)
GET /api/v1/agent/models 可用 Agent 模型

历史(history.py)

方法 路径 说明
GET /api/v1/history 列表
GET /api/v1/history/{id} 详情
DELETE /api/v1/history/{id} 删除
GET /api/v1/history/statistics 统计

组合(portfolio.py)

账户 / 交易 / 持仓 / CSV 导入 / 风险概览。详见源码 router。

2026-05 重构:所有端点统一接入本地 @_portfolio_safe 装饰器,端点 body 只剩业务路径。装饰器映射:HTTPException 透传 / PortfolioBusyError409 portfolio_busy / PortfolioOversellError409 portfolio_oversell / PortfolioConflictError409 conflict / ValueError400 validation_error / Exception500 internal_error。手写的 HTTPException(404, ...) 全部改为 api_errors.not_found(...)。比 @safe_endpoint 多 3 个域异常映射,所以是局部装饰器;语义稳定后可考虑合并到 safe_endpoint(extra_handlers=...)

回测(backtest.py)

run / results / statistics

配置(system_config.py)

读 / 写 / 批量 / 字段元数据 / LLM 状态。

股票(stocks.py)

行情 / 搜索 / 名称解析。

认证(auth.py)

可选;ADMIN_AUTH_ENABLED=true 启用,提供 login / logout / check。

用量(usage.py)

LLM 用量审计统计。

行业 / 热力图(sectors.py)

驱动首页 dashboard 的热力图(treemap)+ 美股大盘指数横排卡片。

方法 路径 说明
GET /api/v1/sectors 行业列表(名称 / 股票数 / 总市值)
GET /api/v1/sectors/heatmap 个股热力图(默认 top_n=30/父行业;支持 ?parent= / ?sector= 下钻)
GET /api/v1/sectors/heatmap-etf 板块 ETF 热力图(?level=parent|sub
GET /api/v1/sectors/indices 美股 4 大指数(道琼斯 / 纳斯达克 / 标普500 / 罗素2000)。2026-05 扩展:返回字段新增 open / high / low / pre_close / volume / data_date,驱动首页指数卡 hover tooltip(HmTooltip 'index' schema)。2026-05 freshness gateexpected_datecompute_batch_as_of(["us"]) 决定;quote 的 trade_date 不匹配 → 整个源被 reject,endpoint 返回该指数 price=None,前端渲染 --(防御 Xueqiu silent-staleness)。
GET /api/v1/sectors/{name:path} 单行业详情

缓存策略(kvcache namespaces)

  • 统一 30 min TTL(盘后复盘工具,无需盘中切换)
  • dashboard_indices namespace(SWR):单 key strip 装整张指数响应;TTL 后第一个用户立刻拿到 stale 值,单次后台 refresh 在 kvcache 共享 SWR 池里跑(无需手写 _INDICES_REFRESH_LOCK
  • dashboard_heatmap namespace(SWR,max_entries=50):以 parent + sector + top_n 组合为 key;同样的 SWR 语义自动合并并发 miss
  • 上游慢调用合并到 kvcache namespace akshare_us_index_daily(1h TTL,stale_while_revalidate)—— akshare index_us_stock_sina 单次 1-3s,缓存命中后 ~250ms
  • as_completed(timeout=4.0s) + ThreadPoolExecutor.shutdown(wait=False, cancel_futures=True):单 index 慢源(如 .RUT 走 yfinance 被限流 ~7s)不阻塞响应,超时即返回 --
  • _enrich_stocks_with_quotesstock_quote_snapshot 表读取最近一次入库的报价;symbols without snapshot 保留 catalog 提供的 market_cap_bprice / change_pct 留空
  • 手动刷新(POST /api/v1/sectors/refresh-quotes)会清空两个 namespace(_heatmap_ns().clear() / _indices_ns().clear()),下一个请求重新走 loader

字段:每只股票返回 code / name / sector / parent_sector / sub_sector / market_cap_b / price / change_pctchange_pct 为最近一次入库报价的涨跌幅(百分数,例如 -1.55)。 数据源:tickbridge.DataProvider.get_realtime_quote()(雪球 P1 优先)。

双 Executor 设计

  • app.state.long_pool — 长 LLM/Agent 调用走这里,4 个 worker(API_LONG_POOL_WORKERS
  • app.state.io_pool — 轻 IO 走这里,16 个 worker(API_IO_POOL_WORKERS
  • 避免长任务饿死短任务(之前都用默认 executor 共享)

配置项

变量 默认 说明
ADMIN_AUTH_ENABLED false 启用认证
ADMIN_PASSWORD 管理员密码
CORS_ORIGINS (内置 5173/3000) 额外允许 origin(逗号分隔;生产部署必填)
CORS_ALLOW_ALL false demo / 紧急逃生口;开启后会强制关闭 allow_credentials(CORS RFC 不允许 wildcard + credentials),admin session cookie 失效(G4.1)

G1 数据保留 endpoints(2026-05)

Method + Path 描述
GET /api/v1/system/retention Dry-run:返回每张表多少行可以删,不删任何数据
POST /api/v1/system/retention/execute 实际执行删除(破坏性,logs WARNING)

返回结构(参见 RetentionReport.to_dict()):

{
  "executed": false,
  "started_at": "...",
  "finished_at": "...",
  "total_deletable": 12345,
  "total_deleted": 0,
  "plans": [
    {"table": "news_intel", "age_column": "published_date", "days": 90,
     "cutoff_date": "2026-02-21", "deletable_rows": 1000, "total_rows": 50000,
     "deleted_rows": 0, "category": "news"}
  ],
  "skipped_tables": []
}

| API_LONG_POOL_WORKERS | 4 | 长任务池 | | API_IO_POOL_WORKERS | 16 | IO 池 | | --port / --host | 8000 / 0.0.0.0 | CLI 参数 |

依赖关系

  • 依赖:fastapi, uvicorn, stocklens/*, tickbridge/*, stocklens.contracts
  • 被依赖:main.py(start_api_server)

注意事项

  • 删除前端代码后,根路由 / 返回简单 JSON(含 docs 链接),不再做 SPA fallback
  • CORS(G4.1,2026-05):默认允许 localhost:5173 / localhost:3000;生产环境必须 CORS_ORIGINS=https://...,https://... 显式配置。CORS_ALLOW_ALL=true 是 demo 开关,会触发启动 WARNING(WEBUI_HOST=0.0.0.0 时警告更醒目);该模式下 allow_credentials 被强制关掉以满足 CORS RFC,session-based 认证 API 会失效。
  • SSE 流使用 text/event-stream
  • 分析任务通过 AnalysisTaskQueue 异步执行
  • 长任务池/IO 池在 lifespan 中初始化与关闭

错误响应统一化(2026-05)

新增 api/v1/errors.py,提供统一的 HTTPException 工厂:

from api.v1 import errors as api_errors

raise api_errors.bad_request("Invalid stock code")
raise api_errors.not_found("Stock", stock_code)
raise api_errors.internal_error("Failed to fetch quote", exc)
raise api_errors.upstream_error("Tavily API down", source="tavily")

错误响应统一形如 {"detail": {"error": "<code>", "message": "<text>", ...}}。 错误码:validation_error / bad_request / unauthorized / forbidden / not_found / conflict / rate_limited / internal_error / upstream_error / service_unavailable

2026-05portfolio.py 原本的 _bad_request / _internal_error / _conflict_error 三个本地 helper 已删除;端点统一接入 @_portfolio_safe 装饰器(详见上一节)。

M9 — @safe_endpoint 装饰器(2026-05)

api/v1/errors.py 新增 safe_endpoint(*, fallback_message="Internal server error", log_label=None) 装饰器。包住 endpoint 后的语义:

抛出的异常 转成什么
HTTPException 原样冒泡(保留显式状态码)
ValueError 400 validation_error
其它 Exception 500 internal_error(自动 logger.error(... exc_info=True)

支持 sync / async 两种 endpoint。这样 endpoint body 只剩业务路径 + raise not_found(...) / raise validation_error(...) 之类的显式 raise,不再需要 try: ... except HTTPException: raise; except Exception as e: raise HTTPException(500, ...) 的样板。

from api.v1.errors import safe_endpoint, not_found

@router.get("/{stock_code}")
@safe_endpoint(fallback_message="Failed to fetch stock")
def get_stock(stock_code: str):
    stock = repo.get(stock_code)
    if not stock:
        raise not_found("Stock", stock_code)
    return stock

当前为新增工具,未强制迁移现有 endpoint。 旧的 try/except/raise 模式仍然 OK;新写或大改的 endpoint 推荐用 @safe_endpoint

F4-F6 — 三个高频模块迁移到 @safe_endpoint(2026-05)

模块 迁移的 endpoint 备注
stocks.py get_stock_quote / get_stock_history / get_stock_logo get_stock_history 内部仍显式 raise HTTPException(422, "unsupported_period") —— 比装饰器默认的 400 validation_error 更精确
history.py (369→287 行) get_history_list / delete_history_records / get_history_detail / get_history_news / get_history_markdown get_history_markdown 显式 raise 500 generation_failed(区别于通用 internal_error),客户端可单独处理
analysis.py _handle_sync_analysis / get_analysis_status _handle_sync_analysis 显式 raise 500 analysis_failedget_analysis_statusapi_errors.not_found("Task", ...)

回归保护:tests/unit/test_safe_endpoint.py 11 用例(HTTPException 透传、ValueError → 400、Exception → 500、async 支持、functools.wraps metadata 保留、args/kwargs 传递)。

未迁移的 endpoint(admin / auth / agent / backtest / dashboard / metrics / portfolio / sectors / system_config / usage)签名保持不变;后续随触改触迁移即可。

H2 — 剩余端点全部接入 @safe_endpoint(2026-05)

H 轮完成了 F4-F6 留下的尾巴。继续迁移:

模块 迁移的 endpoint
agent.py agent_chat
backtest.py run_backtest / get_backtest_results / get_overall_performance / get_stock_performance
dashboard.py dashboard_summary
system_config.py get_system_config / update_system_config / validate_system_config / test_llm_channel / get_system_config_schema / get_retention_dry_run / execute_retention
sectors.py 3 个 404 raise 改为 api_errors.not_found(...)
admin.py revoke_api_key 的 404 改为 api_errors.not_found(...)

剩余的非 @safe_endpoint 项是 intentional(域定制错误码):

  • stocks.get_stock_history 422 unsupported_period(更精确)
  • analysis._handle_sync_analysis 500 analysis_failed(与 internal_error 区分)
  • history.get_history_markdown 500 generation_failed(客户端可单独处理)
  • analysis._handle_async_analysis_batch 409 duplicate_task(业务码)
  • system_config.update_system_config 400 validation_failed + 409 config_version_conflict(带 issues 列表 / 当前版本号)
  • system_config.test_llm_channel 422 validation_error(验证器抛 ValueError
  • sectors.get_etf_heatmap 400 bad_level + 503 etf_catalog_empty(域错误码)
  • portfolio.py 整体使用本地 @_portfolio_safe 装饰器(多 3 个 Portfolio*Error 域异常映射),形态等价 @safe_endpoint

auth.py 内部用 fallback 链 + 显式 try/except 是设计正确的(多源回退),不属于 boilerplate。