mirror of
https://github.com/shy3130/tick-stock-panel.git
synced 2026-10-02 01:35:07 +08:00
feat(open-api): 开放能力闭环 — write:ext 数据写入 + SSE 事件流 + 契约守护
二开闭环补齐三块: 数据进得来 (写入)、事件推得出 (SSE)、契约有人守 (快照+CI)。
- write:ext scope (第六档): POST /api/ext-data/{id}/ingest 对 Token 开放,
复用管理端同一校验/落盘/视图刷新路径; 结构配置/上传/拉取仍不开放,
外部只能写行数据, 不能改写表结构
- SSE 短期票据 (V2 收尾): POST /api/events/ticket 换 60s 一次性票据 →
GET /api/events?ticket=… 订阅事件流; 票据只继承 Token scope 不放大权限;
进程内事件总线 (services/events.py, 慢消费者丢旧保新), 告警落盘点已挂
(alert_store.append/append_many → bus), 事件按 scope 过滤双向验证
- 契约守护: test_openapi_contract.py 冻结 tier=a 开放面 (61 端点),
增删须显式更新快照; 新增 ci.yml (后端全量测试 + 前端构建) 把契约
变化拦在 PR 上
- examples/open-api: 四个零依赖可运行示例 (行情/写入/回测/事件流) + README
- 菜单设置架构标注: 每页核心/扩展徽标 + 悬停原因 + 图例 (按 open-platform-plan
§2 依赖方向分类, 核心 8 页/扩展 11 页); 顺带补上 MenuSettings 缺失的 模拟盘
- 测试: 网关 write矩阵/票据语义 + 票据生命周期/总线/SSE帧/告警联动 + 契约快照,
全量 2408 passed / 102s; 线上实测: ingest 200 / 票据一次性 / SSE hello+
alert 推送 / scope 过滤双向 / 契约 61 路径
This commit is contained in:
@@ -0,0 +1,64 @@
|
||||
name: CI
|
||||
|
||||
# 持续集成: 后端全量测试 + 前端构建。
|
||||
# 与 release.yml 的分工: 这里不打包, 只守质量门 —
|
||||
# 后端含 Tier A 开放契约快照测试 (test_openapi_contract.py),
|
||||
# 开放面的任何变化都会在这里被拦下, 须有意识地更新快照。
|
||||
on:
|
||||
push:
|
||||
branches: [main]
|
||||
pull_request:
|
||||
branches: [main]
|
||||
|
||||
concurrency:
|
||||
group: ci-${{ github.ref }}
|
||||
cancel-in-progress: true
|
||||
|
||||
jobs:
|
||||
backend:
|
||||
runs-on: ubuntu-latest
|
||||
defaults:
|
||||
run:
|
||||
working-directory: backend
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
- name: 安装 uv
|
||||
uses: astral-sh/setup-uv@v3
|
||||
with:
|
||||
enable-cache: true
|
||||
|
||||
- name: 设置 Python
|
||||
run: uv python install 3.12
|
||||
|
||||
- name: 安装依赖 (含 dev 测试链)
|
||||
run: uv sync --extra dev --frozen
|
||||
|
||||
- name: 后端全量测试
|
||||
run: uv run pytest tests -q
|
||||
|
||||
frontend:
|
||||
runs-on: ubuntu-latest
|
||||
defaults:
|
||||
run:
|
||||
working-directory: frontend
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
- name: 安装 pnpm
|
||||
uses: pnpm/action-setup@v4
|
||||
with:
|
||||
version: 9
|
||||
|
||||
- name: 设置 Node 20
|
||||
uses: actions/setup-node@v4
|
||||
with:
|
||||
node-version: '20'
|
||||
cache: 'pnpm'
|
||||
cache-dependency-path: frontend/pnpm-lock.yaml
|
||||
|
||||
- name: 安装依赖
|
||||
run: pnpm install --frozen-lockfile
|
||||
|
||||
- name: 构建 (tsc + vite)
|
||||
run: pnpm build
|
||||
@@ -0,0 +1,81 @@
|
||||
"""开放事件流 — SSE 短期票据 + /api/events 订阅端点。
|
||||
|
||||
EventSource 规范不能携带 Authorization 头, Token 调用方走两步:
|
||||
POST /api/events/ticket (Bearer, 任意 scope) → 一次性票据 (60s)
|
||||
GET /api/events?ticket=… (EventSource) → 按 scope 过滤的事件流
|
||||
|
||||
GET /api/events 在访问中间件白名单内 (票据即凭证, 不走密码会话);
|
||||
票据校验失败立即 401, 成功则流式推送: hello 帧 → 事件 + 15s 心跳。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import queue
|
||||
import time
|
||||
|
||||
from fastapi import APIRouter, HTTPException, Request
|
||||
from fastapi.responses import JSONResponse, StreamingResponse
|
||||
|
||||
from app.services import event_tickets
|
||||
from app.services.events import bus, sse_format
|
||||
|
||||
router = APIRouter(prefix="/api/events", tags=["events"])
|
||||
|
||||
_HEARTBEAT_S = 15.0
|
||||
_POLL_S = 0.5
|
||||
|
||||
|
||||
@router.post("/ticket")
|
||||
def issue_ticket(request: Request):
|
||||
"""Bearer Token (任意 scope) → 一次性 SSE 票据。"""
|
||||
# 认证与限流已在网关中间件完成 (scope 规则 "*"); 端点内再验一次
|
||||
# 明文是为了拿到 Token 记录、让票据继承 scope (不放大权限)。
|
||||
authz = request.headers.get("authorization", "")
|
||||
plaintext = authz[len("Bearer "):].strip() if authz.startswith("Bearer ") else ""
|
||||
|
||||
from app.services import api_tokens
|
||||
|
||||
record = api_tokens.verify_token(request.app.state.repo.store.data_dir, plaintext)
|
||||
if record is None:
|
||||
raise HTTPException(status_code=401, detail="API Token 无效或已吊销")
|
||||
try:
|
||||
ticket, ttl = event_tickets.issue(record.get("scopes", []))
|
||||
except RuntimeError as e:
|
||||
raise HTTPException(status_code=429, detail=str(e)) from e
|
||||
return {"ticket": ticket, "expires_in": int(ttl), "stream": "/api/events?ticket=" + ticket}
|
||||
|
||||
|
||||
@router.get("")
|
||||
async def event_stream(ticket: str = ""):
|
||||
"""SSE 事件流; 票据一次性, 重连需 POST /ticket 换新。"""
|
||||
scopes = event_tickets.consume(ticket) if ticket else None
|
||||
if scopes is None:
|
||||
return JSONResponse(status_code=401, content={"detail": "票据无效、已使用或过期"})
|
||||
|
||||
q = bus.subscribe()
|
||||
|
||||
async def gen():
|
||||
allowed = ("*", *scopes)
|
||||
try:
|
||||
yield sse_format({"type": "hello", "scope": "*", "data": {"scopes": scopes}, "ts": time.time()})
|
||||
last_send = time.monotonic()
|
||||
while True:
|
||||
try:
|
||||
event = q.get_nowait()
|
||||
except queue.Empty:
|
||||
if time.monotonic() - last_send >= _HEARTBEAT_S:
|
||||
yield ": ping\n\n" # SSE 注释帧 — 保活, 不产生 EventSource 事件
|
||||
last_send = time.monotonic()
|
||||
await asyncio.sleep(_POLL_S)
|
||||
continue
|
||||
if event["scope"] in allowed:
|
||||
yield sse_format(event)
|
||||
last_send = time.monotonic()
|
||||
finally:
|
||||
bus.unsubscribe(q)
|
||||
|
||||
return StreamingResponse(
|
||||
gen(),
|
||||
media_type="text/event-stream",
|
||||
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
|
||||
)
|
||||
+14
-1
@@ -19,6 +19,7 @@ from app.api import (
|
||||
backtest,
|
||||
data,
|
||||
ext_data,
|
||||
events,
|
||||
factors,
|
||||
financials,
|
||||
indices,
|
||||
@@ -431,7 +432,18 @@ app.add_middleware(
|
||||
# 3. 已设密码 → 检查 session, 无效则 401(前端跳登录)
|
||||
# 白名单: /api/auth/* (设密码/登录本身)、/health 等探活。
|
||||
_AUTH_WHITELIST_PREFIX = ("/api/auth/",)
|
||||
_AUTH_WHITELIST_EXACT = ("/health", "/api/health", "/openapi.json", "/api/openapi.json", "/docs", "/redoc")
|
||||
_AUTH_WHITELIST_EXACT = (
|
||||
"/health",
|
||||
"/api/health",
|
||||
"/openapi.json",
|
||||
"/api/openapi.json",
|
||||
"/docs",
|
||||
"/redoc",
|
||||
# SSE 事件流: EventSource 带不了 Authorization 头, 凭证即 query 里的
|
||||
# 一次性票据, 端点内校验 (api/events.py); POST /api/events/ticket 不在
|
||||
# 白名单, 仍走网关 Bearer 通道
|
||||
"/api/events",
|
||||
)
|
||||
|
||||
|
||||
@app.middleware("http")
|
||||
@@ -515,6 +527,7 @@ app.include_router(signals.router)
|
||||
app.include_router(monitor_rules.router)
|
||||
app.include_router(lots.router)
|
||||
app.include_router(alerts.router)
|
||||
app.include_router(events.router)
|
||||
app.include_router(rps.router)
|
||||
app.include_router(sector_rotation.router)
|
||||
|
||||
|
||||
@@ -47,6 +47,7 @@ def append(data_dir: Path, event: dict) -> None:
|
||||
if _write_count >= PRUNE_EVERY:
|
||||
_write_count = 0
|
||||
_prune_locked(p)
|
||||
_publish_events([event])
|
||||
|
||||
|
||||
def append_many(data_dir: Path, events: list[dict]) -> None:
|
||||
@@ -63,6 +64,18 @@ def append_many(data_dir: Path, events: list[dict]) -> None:
|
||||
if _write_count >= PRUNE_EVERY:
|
||||
_write_count = 0
|
||||
_prune_locked(p)
|
||||
_publish_events(events)
|
||||
|
||||
|
||||
def _publish_events(events: list[dict]) -> None:
|
||||
"""落盘成功后广播到事件总线 (SSE 出口); 广播失败不影响落盘。"""
|
||||
try:
|
||||
from app.services.events import bus
|
||||
|
||||
for ev in events:
|
||||
bus.publish("alert", ev, scope="read:analysis")
|
||||
except Exception: # noqa: BLE001 — 事件通道绝不反噬主流程
|
||||
logger.debug("alert 事件广播失败 (忽略)", exc_info=True)
|
||||
|
||||
|
||||
def list_recent(
|
||||
|
||||
@@ -55,6 +55,19 @@ _RULES: list[tuple[str, str, str]] = [
|
||||
("POST", "/api/backtest/strategy/run", "run:backtest"),
|
||||
# paper:trade — 模拟盘全部 (读+写一体, 单一 scope 简化心智)
|
||||
("*", "/api/paper", "paper:trade"),
|
||||
# events — 换取 SSE 短期票据: 任意有效 Token 即可 (票据只继承已有 scope,
|
||||
# 不放大权限, 故此处用 "*" 表示"仅需认证", 不做 scope 校验)
|
||||
("POST", "/api/events/ticket", "*"),
|
||||
# SSE 流本体: 网关只做标记 (契约可见); 真正凭证是 query 里的票据,
|
||||
# 由端点内校验 (EventSource 带不了 Authorization 头)
|
||||
("GET", "/api/events", "*"),
|
||||
]
|
||||
|
||||
# 后缀规则表 — 路径中段含动态段 (如 config_id), 前缀表表达不了"以 X 结尾"语义。
|
||||
# 限定在 /api/ext-data/ 下, 防止误伤其他路由的巧合后缀。
|
||||
_SUFFIX_RULES: list[tuple[str, str, str]] = [
|
||||
# write:ext — 程序化行数据写入 (结构配置/上传/拉取等管理端点仍不开放)
|
||||
("POST", "/ingest", "write:ext"),
|
||||
]
|
||||
|
||||
_RATE_WINDOW_S = 60.0
|
||||
@@ -70,13 +83,20 @@ def rate_limit_per_min() -> int:
|
||||
|
||||
|
||||
def required_scope(method: str, path: str) -> str | None:
|
||||
"""路径+方法 → 所需 scope; None = 该端点不对外开放。"""
|
||||
"""路径+方法 → 所需 scope; None = 该端点不对外开放。
|
||||
|
||||
返回 "*" 表示仅需有效 Token (不做特定 scope 校验)。
|
||||
"""
|
||||
for m, prefix, scope in _RULES:
|
||||
if (m == "*" or m == method) and (path == prefix or path.startswith(prefix.rstrip("/") + "/")):
|
||||
# read:ext 域内的拉取 Key 状态子路径属管理面 (脱敏也不外露)
|
||||
if scope == "read:ext" and path.endswith("/api-key"):
|
||||
return None
|
||||
return scope
|
||||
if path.startswith("/api/ext-data/"):
|
||||
for m, suffix, scope in _SUFFIX_RULES:
|
||||
if (m == "*" or m == method) and path.endswith(suffix):
|
||||
return scope
|
||||
return None
|
||||
|
||||
|
||||
@@ -112,7 +132,9 @@ def evaluate(data_dir: Path, method: str, path: str, plaintext: str) -> dict:
|
||||
"detail": f"该端点 ({method} {path}) 未对外开放; 开放清单见 /api/openapi.json?tier=a",
|
||||
"headers": {},
|
||||
}
|
||||
if scope not in record.get("scopes", []):
|
||||
if scope == "*":
|
||||
pass # 仅需有效 Token (票据签发): 已通过认证, 不校验具体 scope
|
||||
elif scope not in record.get("scopes", []):
|
||||
return {
|
||||
"status": 403,
|
||||
"detail": f"Token 缺少所需 scope: {scope} (持有: {record.get('scopes', [])})",
|
||||
|
||||
@@ -25,7 +25,14 @@ from app.services.fs_utils import atomic_write_text
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 对外开放的 scope 全集 (admin 刻意不存在 — 管理面永不开放给 Token)
|
||||
SCOPES = ("read:market", "read:ext", "read:analysis", "run:backtest", "paper:trade")
|
||||
SCOPES = (
|
||||
"read:market",
|
||||
"read:ext",
|
||||
"write:ext",
|
||||
"read:analysis",
|
||||
"run:backtest",
|
||||
"paper:trade",
|
||||
)
|
||||
|
||||
TOKEN_PREFIX = "tsp_"
|
||||
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
"""SSE 短期票据 — EventSource 无法携带 Authorization 头的替代通道。
|
||||
|
||||
流程 (docs/open-platform-plan.md §6 V2):
|
||||
1. POST /api/events/ticket (Bearer Token, 任意 scope) → 一次性票据 (tse_ 前缀)
|
||||
2. GET /api/events?ticket=... (EventSource) → 校验+消费票据 → 按票据
|
||||
继承的 scope 过滤事件流
|
||||
3. 票据 60 秒过期 / 消费即作废; 重连需重新走第 1 步换新票据
|
||||
|
||||
安全边界: 票据只继承 Token 已有 scope, 不放大权限; 存活窗口短,
|
||||
即使泄漏也只能订阅事件流, 不能调用任何 REST 端点。
|
||||
进程内存储 (单进程 uvicorn 场景), 重启清零无害。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import secrets
|
||||
import threading
|
||||
import time
|
||||
|
||||
TICKET_PREFIX = "tse_"
|
||||
TTL_S = 60.0
|
||||
_MAX_OUTSTANDING = 1000 # 防堆积: 远超正常并发, 超限先清理过期
|
||||
|
||||
_LOCK = threading.Lock()
|
||||
_TICKETS: dict[str, tuple[float, list[str]]] = {} # ticket → (过期时刻, scopes)
|
||||
|
||||
|
||||
def _prune(now: float) -> None:
|
||||
expired = [t for t, (exp, _) in _TICKETS.items() if exp <= now]
|
||||
for t in expired:
|
||||
_TICKETS.pop(t, None)
|
||||
|
||||
|
||||
def issue(scopes: list[str]) -> tuple[str, float]:
|
||||
"""签发一次性票据; 返回 (票据, 有效期秒)。"""
|
||||
now = time.monotonic()
|
||||
with _LOCK:
|
||||
_prune(now)
|
||||
if len(_TICKETS) >= _MAX_OUTSTANDING:
|
||||
raise RuntimeError("待使用票据过多, 稍后重试")
|
||||
ticket = TICKET_PREFIX + secrets.token_hex(16)
|
||||
_TICKETS[ticket] = (now + TTL_S, list(scopes))
|
||||
return ticket, TTL_S
|
||||
|
||||
|
||||
def consume(ticket: str) -> list[str] | None:
|
||||
"""校验并消费票据; 返回继承的 scopes, 无效/过期/已用 → None。"""
|
||||
now = time.monotonic()
|
||||
with _LOCK:
|
||||
_prune(now)
|
||||
hit = _TICKETS.pop(ticket, None)
|
||||
if hit is None or hit[0] <= now:
|
||||
return None
|
||||
return hit[1]
|
||||
@@ -0,0 +1,68 @@
|
||||
"""进程内事件总线 — SSE /api/events 的发布/订阅底座。
|
||||
|
||||
设计约束:
|
||||
- 只做进程内广播 (自托管单进程 uvicorn 场景); 不引入 Redis/MQ 依赖。
|
||||
- 订阅队列有界 (默认 200): 慢消费者丢最旧事件而不是拖垮发布方;
|
||||
SSE 断线由客户端换新票据重连补拉 (告警等持久数据另有 REST 出口)。
|
||||
- 事件带 required_scope: 流端点按票据 scope 过滤, 票据不放大权限。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import contextlib
|
||||
import json
|
||||
import queue
|
||||
import threading
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
_QUEUE_MAX = 200
|
||||
|
||||
|
||||
def _serialize(event_type: str, payload: dict[str, Any], scope: str) -> dict[str, Any]:
|
||||
return {
|
||||
"type": event_type,
|
||||
"scope": scope,
|
||||
"data": payload,
|
||||
"ts": time.time(),
|
||||
}
|
||||
|
||||
|
||||
class EventBus:
|
||||
def __init__(self) -> None:
|
||||
self._lock = threading.Lock()
|
||||
self._subscribers: list[queue.Queue[dict[str, Any]]] = []
|
||||
|
||||
def subscribe(self) -> queue.Queue[dict[str, Any]]:
|
||||
q: queue.Queue[dict[str, Any]] = queue.Queue(maxsize=_QUEUE_MAX)
|
||||
with self._lock:
|
||||
self._subscribers.append(q)
|
||||
return q
|
||||
|
||||
def unsubscribe(self, q: queue.Queue[dict[str, Any]]) -> None:
|
||||
with self._lock, contextlib.suppress(ValueError):
|
||||
self._subscribers.remove(q)
|
||||
|
||||
def publish(self, event_type: str, payload: dict[str, Any], scope: str) -> None:
|
||||
"""广播事件; scope = 订阅方票据需持有的最小 scope。"""
|
||||
event = _serialize(event_type, payload, scope)
|
||||
with self._lock:
|
||||
subscribers = list(self._subscribers)
|
||||
for q in subscribers:
|
||||
try:
|
||||
q.put_nowait(event)
|
||||
except queue.Full:
|
||||
# 慢消费者: 丢最旧, 保最新 (监控类事件新比全重要)
|
||||
try:
|
||||
q.get_nowait()
|
||||
q.put_nowait(event)
|
||||
except (queue.Empty, queue.Full):
|
||||
pass
|
||||
|
||||
|
||||
def sse_format(event: dict[str, Any]) -> str:
|
||||
"""事件 → SSE 帧 (id 用单调时间戳, EventSource lastEventId 可用于去重)。"""
|
||||
return f"id: {int(event['ts'] * 1000)}\nevent: {event['type']}\ndata: {json.dumps(event['data'], ensure_ascii=False)}\n\n"
|
||||
|
||||
|
||||
# 进程级单例 — 域模块与 API 层共用同一总线
|
||||
bus = EventBus()
|
||||
@@ -53,7 +53,14 @@ def test_required_scope_mapping():
|
||||
assert f("POST", "/api/kline/daily-batch") is None # 管理面
|
||||
assert f("GET", "/api/ext-data/ext_fuyao_hot/rows") == "read:ext"
|
||||
assert f("GET", "/api/ext-data/ext_fuyao_hot/api-key") is None # Key 状态不外露
|
||||
assert f("POST", "/api/ext-data/ext_fuyao_hot/ingest") is None # 写不开放
|
||||
assert f("POST", "/api/ext-data/ext_fuyao_hot/ingest") == "write:ext" # 程序化写入
|
||||
assert f("POST", "/api/ext-data/ext_fuyao_hot/upload") is None # 文件上传仍属管理面
|
||||
assert f("POST", "/api/ext-data/ext_fuyao_hot/pull/test") is None
|
||||
assert f("POST", "/api/ext-data") is None # 建表不开放
|
||||
assert f("PUT", "/api/ext-data/x/ingest") is None # 仅 POST
|
||||
assert f("GET", "/api/events/ticket") == "*" # 票据签发: 任意有效 Token
|
||||
assert f("POST", "/api/events/ticket") == "*"
|
||||
assert f("GET", "/api/events") == "*" # SSE 流本体 (票据认证)
|
||||
assert f("GET", "/api/backtest/candidates") == "read:analysis"
|
||||
assert f("POST", "/api/backtest/run") == "run:backtest"
|
||||
assert f("POST", "/api/backtest/strategy/run") == "run:backtest"
|
||||
@@ -78,19 +85,27 @@ def client(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> TestClient:
|
||||
app = FastAPI()
|
||||
|
||||
@app.get("/api/kline/daily")
|
||||
def market(): # noqa: ANN001
|
||||
def market():
|
||||
return {"ok": True}
|
||||
|
||||
@app.get("/api/ext-data/x/rows")
|
||||
def ext(): # noqa: ANN001
|
||||
def ext():
|
||||
return {"ok": True}
|
||||
|
||||
@app.post("/api/ext-data/x/ingest")
|
||||
def ingest():
|
||||
return {"status": "ok", "rows": 1}
|
||||
|
||||
@app.post("/api/events/ticket")
|
||||
def ticket():
|
||||
return {"ticket": "tse_x"}
|
||||
|
||||
@app.post("/api/paper/orders")
|
||||
def paper(): # noqa: ANN001
|
||||
def paper():
|
||||
return {"ok": True}
|
||||
|
||||
@app.get("/api/settings/api-tokens")
|
||||
def admin(): # noqa: ANN001
|
||||
def admin():
|
||||
return {"ok": True}
|
||||
|
||||
from app.config import Settings
|
||||
@@ -98,7 +113,7 @@ def client(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> TestClient:
|
||||
object.__setattr__(fake_settings, "data_dir", tmp_path)
|
||||
|
||||
@app.middleware("http")
|
||||
async def token_gate(request, call_next): # noqa: ANN001
|
||||
async def token_gate(request, call_next):
|
||||
path = request.url.path
|
||||
if not path.startswith("/api/"):
|
||||
return await call_next(request)
|
||||
@@ -124,17 +139,17 @@ def client(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> TestClient:
|
||||
|
||||
def test_gateway_401_403_200_matrix(client: TestClient, tmp_path: Path):
|
||||
_, market_token = api_tokens.create_token(tmp_path, "只读行情", ["read:market"])
|
||||
H = {"Authorization": f"Bearer {market_token}"}
|
||||
hdr = {"Authorization": f"Bearer {market_token}"}
|
||||
|
||||
# 401: 伪造 Token
|
||||
r = client.get("/api/kline/daily", headers={"Authorization": "Bearer tsp_deadbeef"})
|
||||
assert r.status_code == 401
|
||||
# 403: scope 不足 (read:market 调 ext)
|
||||
assert client.get("/api/ext-data/x/rows", headers=H).status_code == 403
|
||||
assert client.get("/api/ext-data/x/rows", headers=hdr).status_code == 403
|
||||
# 403: 未开放端点 (管理)
|
||||
assert client.get("/api/settings/api-tokens", headers=H).status_code == 403
|
||||
assert client.get("/api/settings/api-tokens", headers=hdr).status_code == 403
|
||||
# 200: 命中 scope + 限流头
|
||||
r = client.get("/api/kline/daily", headers=H)
|
||||
r = client.get("/api/kline/daily", headers=hdr)
|
||||
assert r.status_code == 200
|
||||
assert r.headers.get("X-RateLimit-Limit") == "120"
|
||||
assert int(r.headers["X-RateLimit-Remaining"]) < 120
|
||||
@@ -146,13 +161,41 @@ def test_gateway_scope_grants_endpoint(client: TestClient, tmp_path: Path):
|
||||
assert r.status_code == 200
|
||||
|
||||
|
||||
def test_gateway_write_ext_ingest(client: TestClient, tmp_path: Path):
|
||||
"""write:ext — 程序化写入扩展表; 其余写端点仍不开放。"""
|
||||
_, w = api_tokens.create_token(tmp_path, "数据写入", ["write:ext"])
|
||||
_, r = api_tokens.create_token(tmp_path, "只读", ["read:ext", "read:market"])
|
||||
hdr_w = {"Authorization": f"Bearer {w}"}
|
||||
hdr_r = {"Authorization": f"Bearer {r}"}
|
||||
|
||||
# 持有 write:ext → ingest 200
|
||||
assert client.post("/api/ext-data/x/ingest", headers=hdr_w).status_code == 200
|
||||
# 只读 Token → 403 (缺 scope)
|
||||
resp = client.post("/api/ext-data/x/ingest", headers=hdr_r)
|
||||
assert resp.status_code == 403
|
||||
assert "write:ext" in resp.json()["detail"]
|
||||
# write:ext 不能读 (rows 是 read:ext)
|
||||
assert client.get("/api/ext-data/x/rows", headers=hdr_w).status_code == 403
|
||||
# 管理面写端点对 write:ext 也不开 (结构/上传/拉取)
|
||||
assert client.post("/api/ext-data/x/upload", headers=hdr_w).status_code == 403
|
||||
|
||||
|
||||
def test_gateway_ticket_any_token(client: TestClient, tmp_path: Path):
|
||||
"""票据签发只要求 Token 有效, 不要求特定 scope (票据继承 scope, 不放大)。"""
|
||||
_, tok = api_tokens.create_token(tmp_path, "任意", ["read:market"])
|
||||
r = client.post("/api/events/ticket", headers={"Authorization": f"Bearer {tok}"})
|
||||
assert r.status_code == 200
|
||||
# 无效 Token → 401
|
||||
assert client.post("/api/events/ticket", headers={"Authorization": "Bearer tsp_bad"}).status_code == 401
|
||||
|
||||
|
||||
def test_gateway_rate_limit_429(client: TestClient, tmp_path: Path, monkeypatch: pytest.MonkeyPatch):
|
||||
monkeypatch.setattr(api_gateway, "rate_limit_per_min", lambda: 3)
|
||||
_, tok = api_tokens.create_token(tmp_path, "限流", ["read:market"])
|
||||
H = {"Authorization": f"Bearer {tok}"}
|
||||
hdr = {"Authorization": f"Bearer {tok}"}
|
||||
for _ in range(3):
|
||||
assert client.get("/api/kline/daily", headers=H).status_code == 200
|
||||
r = client.get("/api/kline/daily", headers=H)
|
||||
assert client.get("/api/kline/daily", headers=hdr).status_code == 200
|
||||
r = client.get("/api/kline/daily", headers=hdr)
|
||||
assert r.status_code == 429
|
||||
assert int(r.headers["Retry-After"]) >= 1
|
||||
# 无 Bearer 的 UI 会话路径不受 Token 桶影响
|
||||
@@ -161,7 +204,7 @@ def test_gateway_rate_limit_429(client: TestClient, tmp_path: Path, monkeypatch:
|
||||
|
||||
def test_revoked_token_rejected_immediately(client: TestClient, tmp_path: Path):
|
||||
record, tok = api_tokens.create_token(tmp_path, "待吊销", ["read:market"])
|
||||
H = {"Authorization": f"Bearer {tok}"}
|
||||
assert client.get("/api/kline/daily", headers=H).status_code == 200
|
||||
hdr = {"Authorization": f"Bearer {tok}"}
|
||||
assert client.get("/api/kline/daily", headers=hdr).status_code == 200
|
||||
api_tokens.revoke_token(tmp_path, record["id"])
|
||||
assert client.get("/api/kline/daily", headers=H).status_code == 401
|
||||
assert client.get("/api/kline/daily", headers=hdr).status_code == 401
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
"""开放事件流测试 — 票据生命周期 / 事件总线 / SSE 端点。"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import queue
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from app.services import alert_store, event_tickets
|
||||
from app.services.events import EventBus, bus, sse_format
|
||||
|
||||
# ── 票据 ──────────────────────────────────────────────
|
||||
|
||||
def test_ticket_roundtrip_single_use():
|
||||
ticket, ttl = event_tickets.issue(["read:analysis"])
|
||||
assert ticket.startswith("tse_") and ttl == 60
|
||||
assert event_tickets.consume(ticket) == ["read:analysis"]
|
||||
assert event_tickets.consume(ticket) is None # 一次性
|
||||
|
||||
|
||||
def test_ticket_invalid_and_expired(monkeypatch: pytest.MonkeyPatch):
|
||||
assert event_tickets.consume("tse_nonexistent") is None
|
||||
assert event_tickets.consume("") is None
|
||||
assert event_tickets.consume("not_a_ticket") is None
|
||||
|
||||
# 过期: 签发后把时钟推过 TTL
|
||||
ticket, _ = event_tickets.issue(["read:market"])
|
||||
real_monotonic = event_tickets.time.monotonic
|
||||
monkeypatch.setattr(event_tickets.time, "monotonic", lambda: real_monotonic() + 61)
|
||||
assert event_tickets.consume(ticket) is None
|
||||
|
||||
|
||||
def test_ticket_overflow_guard(monkeypatch: pytest.MonkeyPatch):
|
||||
# 上限调小 (避免 1000 张票据污染模块态, 影响后续测试)
|
||||
monkeypatch.setattr(event_tickets, "_MAX_OUTSTANDING", 2)
|
||||
event_tickets.issue(["read:market"])
|
||||
event_tickets.issue(["read:market"])
|
||||
with pytest.raises(RuntimeError):
|
||||
event_tickets.issue(["read:market"])
|
||||
|
||||
|
||||
# ── 总线 ──────────────────────────────────────────────
|
||||
|
||||
def test_bus_pubsub_and_scope_tag():
|
||||
b = EventBus()
|
||||
q = b.subscribe()
|
||||
b.publish("alert", {"rule_id": "r1"}, scope="read:analysis")
|
||||
ev = q.get_nowait()
|
||||
assert ev["type"] == "alert" and ev["scope"] == "read:analysis"
|
||||
assert ev["data"] == {"rule_id": "r1"}
|
||||
b.unsubscribe(q)
|
||||
b.publish("alert", {"x": 1}, scope="read:analysis")
|
||||
with pytest.raises(queue.Empty):
|
||||
q.get_nowait()
|
||||
|
||||
|
||||
def test_bus_slow_consumer_drops_oldest():
|
||||
b = EventBus()
|
||||
q = b.subscribe()
|
||||
for i in range(250): # 超过 _QUEUE_MAX=200, 前 50 个被挤掉
|
||||
b.publish("tick", {"i": i}, scope="*")
|
||||
first = q.get_nowait()["data"]["i"]
|
||||
assert first == 50
|
||||
assert q.get_nowait()["data"]["i"] == 51
|
||||
|
||||
|
||||
def test_sse_format_frame():
|
||||
frame = sse_format({"type": "alert", "scope": "read:analysis", "data": {"a": 1}, "ts": 1700000000.123})
|
||||
lines = frame.splitlines()
|
||||
assert lines[0] == "id: 1700000000123"
|
||||
assert lines[1] == "event: alert"
|
||||
assert json.loads(lines[2].removeprefix("data: ")) == {"a": 1}
|
||||
assert frame.endswith("\n\n")
|
||||
|
||||
|
||||
# ── 告警落盘 → 总线联动 ────────────────────────────────
|
||||
|
||||
def test_alert_append_publishes_to_bus(tmp_path: Path):
|
||||
q = bus.subscribe()
|
||||
try:
|
||||
alert_store.append(tmp_path, {"ts": 1, "rule_id": "r1", "source": "test"})
|
||||
ev = q.get_nowait()
|
||||
assert ev["type"] == "alert"
|
||||
assert ev["scope"] == "read:analysis"
|
||||
assert ev["data"]["rule_id"] == "r1"
|
||||
finally:
|
||||
bus.unsubscribe(q)
|
||||
|
||||
|
||||
def test_alert_append_many_publishes_each(tmp_path: Path):
|
||||
q = bus.subscribe()
|
||||
try:
|
||||
alert_store.append_many(tmp_path, [
|
||||
{"ts": 1, "rule_id": "a"}, {"ts": 2, "rule_id": "b"},
|
||||
])
|
||||
assert q.get_nowait()["data"]["rule_id"] == "a"
|
||||
assert q.get_nowait()["data"]["rule_id"] == "b"
|
||||
finally:
|
||||
bus.unsubscribe(q)
|
||||
|
||||
|
||||
# ── SSE 端点 (直接调函数, 消费前两帧后关闭) ────────────
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_event_stream_endpoint_frames(tmp_path: Path):
|
||||
from app.api import events as events_api
|
||||
|
||||
ticket, _ = event_tickets.issue(["read:analysis"])
|
||||
response = await events_api.event_stream(ticket=ticket)
|
||||
assert response.media_type == "text/event-stream"
|
||||
|
||||
agen = response.body_iterator
|
||||
# 帧 1: hello (回显票据 scope)
|
||||
hello = await asyncio.wait_for(agen.__anext__(), timeout=3)
|
||||
assert "event: hello" in hello
|
||||
assert "read:analysis" in hello
|
||||
# 帧 2: 总线事件, 票据 scope 覆盖 → 推送
|
||||
bus.publish("alert", {"rule_id": "r9"}, scope="read:analysis")
|
||||
# 帧 3: 票据 scope 不覆盖 → 不推送 (只发覆盖的)
|
||||
bus.publish("secret", {"x": 1}, scope="paper:trade")
|
||||
bus.publish("alert", {"rule_id": "r10"}, scope="read:analysis")
|
||||
second = await asyncio.wait_for(agen.__anext__(), timeout=3)
|
||||
assert "event: alert" in second and "r9" in second
|
||||
third = await asyncio.wait_for(agen.__anext__(), timeout=3)
|
||||
assert "r10" in third and "secret" not in third
|
||||
await agen.aclose()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_event_stream_rejects_bad_ticket():
|
||||
from app.api import events as events_api
|
||||
|
||||
response = await events_api.event_stream(ticket="tse_gone")
|
||||
assert response.status_code == 401
|
||||
response = await events_api.event_stream(ticket="")
|
||||
assert response.status_code == 401
|
||||
@@ -0,0 +1,111 @@
|
||||
"""Tier A 开放契约快照 — 开放面的任何变化必须是有意识的。
|
||||
|
||||
契约即承诺: /api/openapi.json?tier=a 列出的端点集合是二开方依赖的稳定面。
|
||||
本测试把该集合冻结成常量 — 新增/移除/改方法都会失败, 失败者须同步更新
|
||||
快照并在 PR 里说明契约变更 (对外破坏性变更需 major 版本)。
|
||||
|
||||
快照来源与 /api/openapi.json?tier=a 完全同源: app.openapi() 的 paths
|
||||
逐 (method, path) 过 api_gateway.required_scope — 规则表是唯一源。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from app.services import api_gateway
|
||||
|
||||
# 冻结的开放面 (method + OpenAPI 路径模板), 按字母序。
|
||||
# 变更此表 = 契约变更, 请在 commit message 与 changelog 中明示。
|
||||
EXPECTED_OPEN_ENDPOINTS: set[str] = {
|
||||
"GET /api/alerts",
|
||||
"GET /api/backtest/candidates",
|
||||
"POST /api/backtest/factor/batch",
|
||||
"GET /api/backtest/factor/columns",
|
||||
"POST /api/backtest/factor/run",
|
||||
"POST /api/backtest/run",
|
||||
"GET /api/backtest/status",
|
||||
"POST /api/backtest/strategy/run",
|
||||
"GET /api/events",
|
||||
"POST /api/events/ticket",
|
||||
"GET /api/ext-data",
|
||||
"POST /api/ext-data/{config_id}/ingest",
|
||||
"GET /api/ext-data/{config_id}/rows",
|
||||
"GET /api/ext-data/{config_id}/values",
|
||||
"GET /api/ext-data/schema-all",
|
||||
"GET /api/ext-data/schema/{config_id}",
|
||||
"GET /api/ext-data/{config_id}/dimension-intraday",
|
||||
"GET /api/ext-data/{config_id}/dimension-members",
|
||||
"GET /api/index/daily",
|
||||
"GET /api/index/minute",
|
||||
"GET /api/intraday/indices",
|
||||
"GET /api/intraday/status",
|
||||
"GET /api/kline/daily",
|
||||
"GET /api/kline/daily/latest",
|
||||
"GET /api/kline/instruments/search",
|
||||
"GET /api/kline/minute",
|
||||
"GET /api/kline/minute-range",
|
||||
"GET /api/overview/market",
|
||||
"GET /api/paper/account",
|
||||
"POST /api/paper/account",
|
||||
"GET /api/paper/accounts",
|
||||
"POST /api/paper/auto_rules",
|
||||
"DELETE /api/paper/auto_rules/{rule_id}",
|
||||
"GET /api/paper/auto_rules",
|
||||
"POST /api/paper/auto_rules/{rule_id}/enabled",
|
||||
"GET /api/paper/compare",
|
||||
"POST /api/paper/freeze",
|
||||
"GET /api/paper/nav",
|
||||
"POST /api/paper/orders",
|
||||
"DELETE /api/paper/orders/{order_id}",
|
||||
"GET /api/paper/orders",
|
||||
"GET /api/paper/overview",
|
||||
"GET /api/paper/positions",
|
||||
"POST /api/paper/rebuild",
|
||||
"POST /api/paper/settings",
|
||||
"GET /api/paper/stats",
|
||||
"GET /api/paper/trades",
|
||||
"GET /api/regime/coverage",
|
||||
"GET /api/regime/history",
|
||||
"GET /api/regime/latest",
|
||||
"GET /api/regime/mainline",
|
||||
"GET /api/regime/phases",
|
||||
"GET /api/regime/states",
|
||||
"GET /api/screener/cached",
|
||||
"GET /api/screener/limit-ladder",
|
||||
"GET /api/screener/market-snapshot",
|
||||
"POST /api/screener/run",
|
||||
"POST /api/screener/run_all",
|
||||
"POST /api/screener/run_preset",
|
||||
"GET /api/screener/strategies",
|
||||
"GET /api/strategies",
|
||||
"GET /api/strategies/ai/status",
|
||||
"GET /api/strategies/{strategy_id}",
|
||||
"GET /api/strategies/{strategy_id}/source",
|
||||
}
|
||||
|
||||
|
||||
def _collect_open_endpoints() -> set[str]:
|
||||
"""与 /api/openapi.json?tier=a 同源收集 (app.openapi paths 过 required_scope)。"""
|
||||
from app.main import app
|
||||
|
||||
spec = app.openapi()
|
||||
found: set[str] = set()
|
||||
for path, ops in spec.get("paths", {}).items():
|
||||
for method in ops:
|
||||
if method in ("get", "post", "put", "delete", "patch") and api_gateway.required_scope(
|
||||
method.upper(), path,
|
||||
):
|
||||
found.add(f"{method.upper()} {path}")
|
||||
return found
|
||||
|
||||
|
||||
def test_open_contract_snapshot():
|
||||
found = _collect_open_endpoints()
|
||||
added = found - EXPECTED_OPEN_ENDPOINTS
|
||||
removed = EXPECTED_OPEN_ENDPOINTS - found
|
||||
if added or removed:
|
||||
pytest.fail(
|
||||
"Tier A 开放契约发生变化 — 这是对外承诺面, 须有意识地更新快照:\n"
|
||||
f" 新增: {sorted(added) or '无'}\n"
|
||||
f" 移除: {sorted(removed) or '无'}\n"
|
||||
"确认无误后, 同步更新本文件的 EXPECTED_OPEN_ENDPOINTS 并在 commit 中说明。",
|
||||
)
|
||||
+22
-4
@@ -231,23 +231,24 @@ APScheduler 默认 15:35 CST 自动:拉日 K → 重算 enriched 表 → 跑监
|
||||
|
||||
## 🌐 开放接口(Open API · Tier A)
|
||||
|
||||
外部程序(自己的看板/脚本/量化服务)无需面板密码,凭 **API Token** 直接调用核心只读与任务接口 —— 这是「核心能力开放」的第一层出口。总体设计与后续 Tier 见 [open-platform-plan.md](./open-platform-plan.md)。
|
||||
外部程序(自己的看板/脚本/量化服务)无需面板密码,凭 **API Token** 直接调用核心读取、写入与任务接口 —— 「核心能力开放」的出口层。总体设计与后续 Tier 见 [open-platform-plan.md](./open-platform-plan.md),可运行示例见 [examples/open-api](../examples/open-api/README.md)。
|
||||
|
||||
### Token 管理
|
||||
|
||||
**设置 → 开放接口**: 新建 Token(名称 + 权限勾选),明文 `tsp_` 前缀只显示一次;可随时吊销,吊销立即生效。服务端只存 SHA-256 哈希(`data/user_data/api_tokens.json`,权限 0600)。
|
||||
|
||||
### 五个权限(scope)
|
||||
### 六个权限(scope)
|
||||
|
||||
| scope | 能力 |
|
||||
| :--- | :--- |
|
||||
| `read:market` | 标的搜索 / 日K / 分时 / 指数 / 市场快照 (默认勾选) |
|
||||
| `read:ext` | 扩展表 rows / values / schema 查询 |
|
||||
| `write:ext` | 向**已配置的**扩展表程序化写入行数据 (`POST /api/ext-data/{id}/ingest`) |
|
||||
| `read:analysis` | 策略清单与结果 / 回测报告 / 市场环境 / 告警 |
|
||||
| `run:backtest` | 提交回测 / 选股 / 因子检验任务 |
|
||||
| `paper:trade` | 模拟盘读取与下单/撤单 (最高敏感) |
|
||||
|
||||
管理接口(数据同步 / 扩展表写入 / 设置)**永不开放给 Token**;数据同步等写操作只能通过面板密码会话进行。
|
||||
管理面**永不开放给 Token**: 数据同步、扩展表的**结构**配置(建表/字段/拉取/上传)、设置等只能通过面板密码会话。`write:ext` 只写行数据 —— 外部程序能把数据喂进来,但不能改写表结构;写入的数据会进入策略/回测数据面,谨慎授予。
|
||||
|
||||
### 调用方式
|
||||
|
||||
@@ -261,13 +262,30 @@ curl -H "Authorization: Bearer tsp_xxxx" \
|
||||
- **错误语义**: 401 Token 无效或已吊销 / 403 权限不足或该接口未开放 / 429 超限;
|
||||
- **CORS 全开**(自托管场景),浏览器直连亦可。
|
||||
|
||||
### 事件流(SSE,实时推送)
|
||||
|
||||
EventSource 无法携带 Authorization 头,外部程序走**短期票据**两步接入:
|
||||
|
||||
```bash
|
||||
# ① Bearer 换一次性票据 (60 秒, 任意 scope 的 Token 皆可)
|
||||
curl -X POST -H "Authorization: Bearer tsp_xxxx" http://localhost:8398/api/events/ticket
|
||||
# → {"ticket":"tse_xxx","expires_in":60,"stream":"/api/events?ticket=tse_xxx"}
|
||||
|
||||
# ② 订阅事件流 (浏览器: new EventStream(stream); 脚本: 逐行读)
|
||||
curl -N "http://localhost:8398/api/events?ticket=tse_xxx"
|
||||
```
|
||||
|
||||
- 当前事件源:**告警触发**(`event: alert`,需票据含 `read:analysis`);心跳 15s;
|
||||
- 票据**一次性**且短时效 —— 泄漏最多损失一次事件订阅,不能调任何 REST 端点;断线重连需重新换票;
|
||||
- 事件总线是进程内广播(`app/services/events.py`),新的核心事件源在域模块里 `bus.publish(...)` 即可挂上。
|
||||
|
||||
### 契约文档(机器可读)
|
||||
|
||||
```bash
|
||||
curl "http://localhost:8398/api/openapi.json?tier=a"
|
||||
```
|
||||
|
||||
返回按网关规则表过滤后的 OpenAPI 3 规范(`x-tier: a`)—— 哪些路径对外开放、需要什么 scope,**规则表是唯一契约源**,生成代码 / Postman 导入即用。
|
||||
返回按网关规则表过滤后的 OpenAPI 3 规范(`x-tier: a`)—— 哪些路径对外开放、需要什么 scope,**规则表是唯一契约源**,生成代码 / Postman 导入即用。开放面有**契约快照测试**守护(`test_openapi_contract.py`):增删开放端点必须显式更新快照,CI 会拦下无意识的契约变更。
|
||||
|
||||
### 桌面客户端版本清单
|
||||
|
||||
|
||||
@@ -78,11 +78,12 @@
|
||||
| `read:market` | 行情/指数/标的搜索/分时/日K 读取 | ✅ 新 Token 默认 |
|
||||
| `read:analysis` | 策略列表与运行结果、回测报告与候选、市场环境、监控事件读取 | ❌ 显式勾选 |
|
||||
| `read:ext` | 扩展数据 rows/values/schema 读取 | ❌ 显式勾选 |
|
||||
| `write:ext` | 扩展表**行数据**程序化写入(`POST /api/ext-data/{id}/ingest`, 0.3.2 增) | ❌ 显式勾选 |
|
||||
| `run:backtest` | 触发回测/挖掘任务(受重活并发器约束) | ❌ 显式勾选 |
|
||||
| `paper:trade` | 模拟盘下单/撤单/建户(写操作, 最高敏感级) | ❌ 显式勾选 |
|
||||
| `admin` | 管理接口 | ❌ **永不签发给 Token**, 仅 UI 会话可用 |
|
||||
|
||||
写敏感端点(扩展表 upload/ingest/backfill/delete、数据同步、设置)一律不进 Token scope —— 外部只读 + 指定的两类写(回测任务、模拟盘)。
|
||||
结构管理端点(扩展表建表/字段/上传/拉取/回补/删除、数据同步、设置)一律不进 Token scope —— 外部可读 + 三类受控写(扩展表**行数据**、回测任务、模拟盘); 行数据写入复用管理端同一校验与落盘路径, 表结构不可被外部改写。
|
||||
|
||||
### 4.3 限流与 CORS
|
||||
|
||||
@@ -170,8 +171,8 @@
|
||||
### V2 契约化 + 示例仓库
|
||||
|
||||
- ✅ Tier A 已含策略/回测/监控事件/模拟盘; `GET /api/openapi.json?tier=a` 契约视图(规则表单源)。
|
||||
- 组织下建 `api-examples` 示例仓库: 「拉行情 → 跑策略 → 取结果 → 模拟盘跟单」三个可运行脚本 + README。
|
||||
- SSE 短期票据(query token)方案落地。
|
||||
- ✅ 示例脚本(0.3.2): `examples/open-api/` 四个零依赖可运行示例(行情 / 写入扩展数据 / 回测 / 事件流)+ README; 独立示例仓库待有真实用户再拆。
|
||||
- ✅ SSE 短期票据落地(0.3.2): `POST /api/events/ticket` 换 60s 一次性票据 → `GET /api/events?ticket=…`(`services/event_tickets.py` + `api/events.py`, 事件源已接告警触发)。
|
||||
|
||||
### V3 钩子契约 + 模拟盘拆分试点
|
||||
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
"""示例 1 — 读行情 (scope: read:market)。
|
||||
|
||||
标的搜索 → 拉日K → 市场总览。契约与参数详见 /api/openapi.json?tier=a。
|
||||
"""
|
||||
from common import call
|
||||
|
||||
# 1. 标的搜索 (支持代码/名称/拼音)
|
||||
hits = call("GET", "/api/kline/instruments/search?q=600519&limit=3")
|
||||
for h in hits["results"]:
|
||||
print(f"搜索命中: {h['symbol']} {h['name']}")
|
||||
|
||||
# 2. 贵州茅台 最近 10 根日K (前复权口径)
|
||||
daily = call("GET", "/api/kline/daily?symbol=600519.SH&limit=10")
|
||||
rows = daily["rows"]
|
||||
if rows:
|
||||
last = rows[-1]
|
||||
print(f"日K 最近 {len(rows)} 根, 最新 {last['date']} 收盘 {last['close']:.2f}")
|
||||
else:
|
||||
print("日K 无数据 — 先在面板「数据」页跑一次同步")
|
||||
|
||||
# 3. 市场总览 (快照日期 / 行情源状态)
|
||||
overview = call("GET", "/api/overview/market")
|
||||
print("市场总览 as_of:", overview.get("as_of"), "| 行情轮询:", overview["quote_status"]["enabled"])
|
||||
@@ -0,0 +1,30 @@
|
||||
"""示例 2 — 写入扩展数据 (scope: write:ext + read:ext)。
|
||||
|
||||
把自有数据程序化喂进扩展表, 与内置数据同台分析:
|
||||
列出扩展表 → ingest 写入一行 → /rows 读回验证。
|
||||
表结构 (字段/模式) 需先在面板「数据 → 扩展数据」配好 — 结构配置属管理面,
|
||||
Token 只能写行数据, 防止外部改写表结构污染策略数据面。
|
||||
"""
|
||||
from common import call
|
||||
|
||||
# 1. 列出已配置的扩展表 (items: [{id, label, mode, fields:[{name,dtype,...}]}])
|
||||
items = call("GET", "/api/ext-data")["items"]
|
||||
if not items:
|
||||
raise SystemExit("还没有扩展表 — 先在面板 数据 → 扩展数据 创建一张 (如 人气排行, 字段: symbol,rank)")
|
||||
|
||||
cfg = items[0]
|
||||
cid = cfg["id"]
|
||||
fields = [f["name"] for f in cfg["fields"]]
|
||||
print(f"目标表: {cfg.get('label', cid)} ({cid}) | 模式: {cfg.get('mode')} | 字段: {fields}")
|
||||
|
||||
# 2. 写入一行 (date 不传按北京当天落盘; 时序表自动进当日分区)
|
||||
row = {"symbol": "600519.SH"}
|
||||
for f in fields:
|
||||
if f not in row:
|
||||
row[f] = 1 # 示例值 — 实际接入时换成你的数据
|
||||
resp = call("POST", f"/api/ext-data/{cid}/ingest", {"rows": [row]})
|
||||
print("写入:", resp)
|
||||
|
||||
# 3. 读回验证 (等值过滤 + 分页)
|
||||
back = call("GET", f"/api/ext-data/{cid}/rows?filter=symbol:600519.SH&limit=3")
|
||||
print(f"读回 {len(back['rows'])} 行 (共 {back.get('total', '?')}):", back["rows"])
|
||||
@@ -0,0 +1,24 @@
|
||||
"""示例 3 — 策略与回测 (scope: read:analysis + run:backtest)。
|
||||
|
||||
列出策略清单 → 对单标的跑一次信号回测 → 读关键指标。
|
||||
可用信号名以面板「信号库」/ GET /api/strategies 为准;
|
||||
策略全周期回测走 POST /api/backtest/strategy/run (body 详见契约)。
|
||||
"""
|
||||
from common import call
|
||||
|
||||
# 1. 策略清单 (read:analysis)
|
||||
strategies = call("GET", "/api/strategies")["strategies"]
|
||||
for s in strategies[:5]:
|
||||
print(f"策略: {s['id']} {s['name']} [{','.join(s.get('tags', []))}]")
|
||||
if not strategies:
|
||||
print("(无策略)")
|
||||
|
||||
# 2. 信号回测 (run:backtest): 贵州茅台, 5/20 均线金叉入场
|
||||
result = call("POST", "/api/backtest/run", {
|
||||
"symbols": ["600519.SH"],
|
||||
"entries": ["signal_ma_golden_5_20"], # 内置信号; 更多见 信号库 页面
|
||||
"exits": [],
|
||||
})
|
||||
for key in ("total_return", "win_rate", "trades", "max_drawdown"):
|
||||
if key in result:
|
||||
print(f"回测 {key}: {result[key]}")
|
||||
@@ -0,0 +1,24 @@
|
||||
"""示例 4 — 事件流 SSE (任意有效 Token; 票据继承 scope)。
|
||||
|
||||
EventSource 带不了 Authorization 头, 走两步:
|
||||
POST /api/events/ticket (Bearer) → 一次性票据 (60 秒)
|
||||
GET /api/events?ticket=… → hello 帧 + 事件 (此处演示 alert 告警推送)
|
||||
|
||||
票据一次性: 断线重连需重新换票。浏览器端用 EventSource,
|
||||
脚本端像本例一样逐行读流即可。
|
||||
"""
|
||||
import itertools
|
||||
|
||||
from common import call, sse
|
||||
|
||||
# 1. 换票 (任意 scope 的 Token 都可; 流里只推票据 scope 覆盖的事件)
|
||||
ticket = call("POST", "/api/events/ticket")["ticket"]
|
||||
print("票据已签发 (60s 内一次性)")
|
||||
|
||||
# 2. 订阅: hello 帧之后, 面板触发任意监控告警即可看到 alert 事件
|
||||
for event in itertools.islice(sse(f"/api/events?ticket={ticket}"), 5):
|
||||
if event.get("event") == "hello":
|
||||
print("hello: scopes =", event.get("data", "").strip("{}"))
|
||||
continue
|
||||
print(f"[{event.get('event')}] {event.get('data', '')[:120]}")
|
||||
print("(收到 5 帧后退出; 生产环境保持连接持续接收)")
|
||||
@@ -0,0 +1,46 @@
|
||||
# 开放接口示例 (Open API Examples)
|
||||
|
||||
外部程序接入 tick-stock-panel 核心能力的可运行最小示例 — 纯 Python 标准库, 零第三方依赖。
|
||||
|
||||
## 准备
|
||||
|
||||
1. 启动面板 (或桌面客户端), 进入 **设置 → 开放接口**, 新建 Token 并勾选所需 scope:
|
||||
|
||||
| 脚本 | 所需 scope |
|
||||
| :--- | :--- |
|
||||
| `01_pull_market.py` | `read:market` |
|
||||
| `02_write_ext_data.py` | `write:ext` + `read:ext` |
|
||||
| `03_run_backtest.py` | `read:analysis` + `run:backtest` |
|
||||
| `04_events_stream.py` | 任意 (票据继承 scope) |
|
||||
|
||||
2. 设置环境变量后运行:
|
||||
|
||||
```bash
|
||||
export TSP_BASE=http://127.0.0.1:3018 # 面板实际地址
|
||||
export TSP_TOKEN=tsp_xxxx # 创建 Token 时只显示一次, 立即保存
|
||||
|
||||
python 01_pull_market.py
|
||||
```
|
||||
|
||||
## 脚本说明
|
||||
|
||||
| 脚本 | 演示能力 |
|
||||
| :--- | :--- |
|
||||
| `01_pull_market.py` | 标的搜索 (拼音可用) / 日K / 市场总览 |
|
||||
| `02_write_ext_data.py` | 程序化写入扩展表行数据 + 读回验证 (数据同台分析闭环) |
|
||||
| `03_run_backtest.py` | 策略清单 / 信号回测触发与结果读取 |
|
||||
| `04_events_stream.py` | SSE 短期票据 + 事件流 (告警推送) |
|
||||
|
||||
## 通用约定
|
||||
|
||||
- 认证: `Authorization: Bearer tsp_...` 头; 与面板密码会话并行, 互不影响
|
||||
- 限流: 每 Token 默认 120 次/分钟, 响应头 `X-RateLimit-Remaining` 实时可见, 超限 429 + `Retry-After`
|
||||
- 错误语义: 401 Token 无效或已吊销 / 403 scope 不足或端点未开放 / 429 超限
|
||||
- 机器可读契约: `GET /api/openapi.json?tier=a` — 可直接导入 Postman / 生成客户端
|
||||
- 完整文档: [docs/features.md → 开放接口](../../docs/features.md#-开放接口open-api--tier-a)
|
||||
|
||||
## 注意
|
||||
|
||||
- `write:ext` 只能向**已配置好的**扩展表写行数据; 表结构 (字段/模式/拉取) 属管理面, Token 不可触及
|
||||
- 写入的数据会进入策略信号 / 因子 / 回测的数据面, 谨慎授予该 scope
|
||||
- SSE 票据一次性且 60 秒过期, 断线重连需重新 `POST /api/events/ticket`
|
||||
@@ -0,0 +1,58 @@
|
||||
"""开放接口示例脚本 — 公共工具 (纯标准库, 无第三方依赖)。
|
||||
|
||||
使用前:
|
||||
1. 面板 → 设置 → 开放接口 → 新建 Token (按需勾选 scope)
|
||||
2. 设置环境变量:
|
||||
TSP_TOKEN=tsp_xxx (必填)
|
||||
TSP_BASE=http://127.0.0.1:3018 (默认本机开发端口, 按实际改)
|
||||
3. 逐个运行: python 01_pull_market.py
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
BASE = os.environ.get("TSP_BASE", "http://127.0.0.1:3018").rstrip("/")
|
||||
TOKEN = os.environ.get("TSP_TOKEN", "")
|
||||
|
||||
|
||||
def call(method: str, path: str, body: dict | None = None, headers: dict[str, str] | None = None):
|
||||
"""调开放接口; 非 2xx 抛 RuntimeError 并带出服务端 detail。"""
|
||||
if not TOKEN:
|
||||
raise SystemExit("请先设置环境变量 TSP_TOKEN (设置 → 开放接口 → 新建 Token)")
|
||||
req = urllib.request.Request(
|
||||
BASE + path,
|
||||
data=json.dumps(body).encode() if body is not None else None,
|
||||
method=method,
|
||||
headers={
|
||||
"Authorization": f"Bearer {TOKEN}",
|
||||
"Content-Type": "application/json",
|
||||
**(headers or {}),
|
||||
},
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30) as resp:
|
||||
return json.loads(resp.read().decode())
|
||||
except urllib.error.HTTPError as e:
|
||||
detail = e.read().decode()
|
||||
raise RuntimeError(f"{method} {path} → HTTP {e.code}: {detail}") from e
|
||||
|
||||
|
||||
def sse(path: str):
|
||||
"""GET 一个 SSE 流, 逐行 yield (demo 用; 生产建议浏览器 EventSource)。"""
|
||||
req = urllib.request.Request(BASE + path, headers={"Accept": "text/event-stream"})
|
||||
with urllib.request.urlopen(req, timeout=None) as resp:
|
||||
event: dict[str, str] = {}
|
||||
for raw in resp:
|
||||
line = raw.decode().rstrip("\n")
|
||||
if not line:
|
||||
if event:
|
||||
yield event
|
||||
event = {}
|
||||
continue
|
||||
if line.startswith(":"):
|
||||
continue # 心跳注释帧
|
||||
key, _, val = line.partition(": ")
|
||||
event[key] = val
|
||||
@@ -9,6 +9,7 @@ import { cn } from '@/lib/cn'
|
||||
const SCOPES: { key: string; label: string; desc: string; default: boolean }[] = [
|
||||
{ key: 'read:market', label: '行情读取', desc: '标的搜索 / 日K / 分时 / 指数 / 市场快照', default: true },
|
||||
{ key: 'read:ext', label: '扩展数据读取', desc: '扩展表 rows / values / schema 查询', default: false },
|
||||
{ key: 'write:ext', label: '扩展数据写入', desc: '向已配置的扩展表程序化写入行数据 (会进入策略/回测数据面, 谨慎授予)', default: false },
|
||||
{ key: 'read:analysis', label: '分析结果读取', desc: '策略清单与结果 / 回测报告 / 市场环境 / 告警', default: false },
|
||||
{ key: 'run:backtest', label: '触发回测', desc: '提交回测 / 选股 / 因子检验任务 (受并发约束)', default: false },
|
||||
{ key: 'paper:trade', label: '模拟盘交易', desc: '模拟盘读取与下单/撤单 (写操作, 最高敏感)', default: false },
|
||||
@@ -92,7 +93,7 @@ export function SettingsApiTokensPanel() {
|
||||
</label>
|
||||
))}
|
||||
<div className="text-[10px] text-muted">
|
||||
管理接口 (数据同步 / 扩展表写入 / 设置) 永不开放给 Token。调用方式: <span className="font-mono">Authorization: Bearer tsp_...</span>, 默认限流 120 次/分钟, 契约文档 <span className="font-mono text-accent/80">/api/openapi.json?tier=a</span>
|
||||
扩展表的结构配置/上传/拉取、数据同步、设置等管理接口永不开放给 Token (行数据写入走 write:ext)。调用方式: <span className="font-mono">Authorization: Bearer tsp_...</span>, 默认限流 120 次/分钟, 契约文档 <span className="font-mono text-accent/80">/api/openapi.json?tier=a</span>
|
||||
</div>
|
||||
</div>
|
||||
<div className="flex gap-2">
|
||||
|
||||
@@ -46,12 +46,43 @@ const BUILTIN_PAGES: NavEntry[] = [
|
||||
{ id: '/regime', label: '市场环境', type: 'builtin', visible: true },
|
||||
{ id: '/abnormal', label: '异动监控', type: 'builtin', visible: true },
|
||||
{ id: '/lots', label: '持仓提醒', type: 'builtin', visible: true },
|
||||
{ id: '/paper', label: '模拟盘', type: 'builtin', visible: true },
|
||||
{ id: '/signals', label: '信号库', type: 'builtin', visible: true },
|
||||
{ id: '/review', label: '复盘', type: 'builtin', visible: true },
|
||||
{ id: '/indices', label: '指数', type: 'builtin', visible: true },
|
||||
{ id: '/data', label: '数据', type: 'builtin', visible: true },
|
||||
]
|
||||
|
||||
// ── 架构归属标注 (docs/open-platform-plan.md §2 核心域/扩展域边界) ──
|
||||
// 核心 = 页面承载核心域能力 (数据底座/订阅锚点/策略回测/环境生产);
|
||||
// 扩展 = 核心数据的消费页, 二开可用「扩展页面」平行实现替换 UI 而不伤核心。
|
||||
// 与「类型」列的区别: 类型标注页面来源 (内置/自定义), 归属标注架构分层。
|
||||
type ArchKind = 'core' | 'ext'
|
||||
const ARCH_CLASS: Record<string, { kind: ArchKind; reason: string }> = {
|
||||
'/': { kind: 'core', reason: '总览页 — 聚合核心数据的多维视图' },
|
||||
'/watchlist': { kind: 'core', reason: '订阅锚点 — 监控/提醒/看板的数据范围基准' },
|
||||
'/screener': { kind: 'core', reason: '策略引擎 — 选股/评分/信号定义' },
|
||||
'/factors': { kind: 'core', reason: '因子平台 — 因子检验与回测联动' },
|
||||
'/backtest': { kind: 'core', reason: '回测引擎 — 矩阵回测与优化' },
|
||||
'/regime': { kind: 'core', reason: '环境数据生产 — 回测/挖掘/归因共用的市场状态口径' },
|
||||
'/indices': { kind: 'core', reason: '行情数据底座 — 指数日K/分时出口' },
|
||||
'/data': { kind: 'core', reason: '数据底座 — 同步管道与数据源路由' },
|
||||
'/stock-analysis': { kind: 'ext', reason: '个股分析视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/limit-ladder': { kind: 'ext', reason: '连板梯队视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/concept-analysis': { kind: 'ext', reason: '概念分析视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/industry-analysis': { kind: 'ext', reason: '行业分析视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/financials': { kind: 'ext', reason: '财务分析视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/monitor': { kind: 'ext', reason: '监控消费页 (规则引擎与推送管道属核心), 页面可替换' },
|
||||
'/abnormal': { kind: 'ext', reason: '异动监控视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/lots': { kind: 'ext', reason: '持仓提醒视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/paper': { kind: 'ext', reason: '模拟盘 — 官方插件化拆分候选 (V3)' },
|
||||
'/signals': { kind: 'ext', reason: '信号库视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
'/review': { kind: 'ext', reason: '大盘复盘视图 — 消费核心数据, 可由扩展页面替换' },
|
||||
}
|
||||
/** 自定义分析页天然属于扩展域 */
|
||||
const archOf = (entry: NavEntry): { kind: ArchKind; reason: string } =>
|
||||
ARCH_CLASS[entry.id] ?? { kind: 'ext', reason: '自定义分析页 — 扩展域' }
|
||||
|
||||
// ── Sortable row ──
|
||||
|
||||
function SortableItem({ entry, hidden, onToggleHidden, badgeEnabled, onToggleBadge }: {
|
||||
@@ -96,6 +127,21 @@ function SortableItem({ entry, hidden, onToggleHidden, badgeEnabled, onToggleBad
|
||||
<span className={`truncate text-sm font-medium ${!hidden ? 'text-foreground' : 'text-muted line-through'}`}>
|
||||
{entry.label}
|
||||
</span>
|
||||
{(() => {
|
||||
const arch = archOf(entry)
|
||||
return (
|
||||
<span
|
||||
title={arch.reason}
|
||||
className={`shrink-0 cursor-help rounded px-1 py-px text-[9px] leading-4 ${
|
||||
arch.kind === 'core'
|
||||
? 'bg-accent/10 text-accent'
|
||||
: 'border border-border text-muted'
|
||||
}`}
|
||||
>
|
||||
{arch.kind === 'core' ? '核心' : '扩展'}
|
||||
</span>
|
||||
)
|
||||
})()}
|
||||
{hidden && (
|
||||
<span className="rounded bg-elevated px-1.5 py-0.5 text-[10px] text-muted shrink-0">已隐藏</span>
|
||||
)}
|
||||
@@ -295,6 +341,17 @@ export function SettingsMenuSettingsPanel() {
|
||||
<p className="mt-2 max-w-3xl text-sm leading-6 text-secondary">
|
||||
拖动左侧手柄调整菜单排列顺序,点击眼睛图标控制菜单在侧边栏中的显示或隐藏。
|
||||
</p>
|
||||
<p className="mt-2 flex flex-wrap items-center gap-x-3 gap-y-1 text-[11px] text-muted">
|
||||
<span>架构归属 (悬停看原因):</span>
|
||||
<span className="inline-flex items-center gap-1">
|
||||
<span className="rounded bg-accent/10 px-1 py-px text-[9px] leading-4 text-accent">核心</span>
|
||||
页面承载核心域能力 (数据底座 / 订阅锚点 / 策略回测 / 环境生产)
|
||||
</span>
|
||||
<span className="inline-flex items-center gap-1">
|
||||
<span className="rounded border border-border px-1 py-px text-[9px] leading-4 text-muted">扩展</span>
|
||||
核心数据的消费页 — 二次开发可用「扩展页面」平行实现替换, 不伤核心
|
||||
</span>
|
||||
</p>
|
||||
</section>
|
||||
|
||||
<section className="rounded-card border border-border bg-surface overflow-hidden">
|
||||
|
||||
Reference in New Issue
Block a user