diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 00000000..faa69b95 --- /dev/null +++ b/.github/workflows/ci.yml @@ -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 diff --git a/backend/app/api/events.py b/backend/app/api/events.py new file mode 100644 index 00000000..7cef159f --- /dev/null +++ b/backend/app/api/events.py @@ -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"}, + ) diff --git a/backend/app/main.py b/backend/app/main.py index b18fa28f..d0a3335d 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -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) diff --git a/backend/app/services/alert_store.py b/backend/app/services/alert_store.py index 55ea67c0..b3ce026b 100644 --- a/backend/app/services/alert_store.py +++ b/backend/app/services/alert_store.py @@ -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( diff --git a/backend/app/services/api_gateway.py b/backend/app/services/api_gateway.py index 50c666bc..3e498e44 100644 --- a/backend/app/services/api_gateway.py +++ b/backend/app/services/api_gateway.py @@ -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', [])})", diff --git a/backend/app/services/api_tokens.py b/backend/app/services/api_tokens.py index 18fe6ff1..260779d5 100644 --- a/backend/app/services/api_tokens.py +++ b/backend/app/services/api_tokens.py @@ -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_" diff --git a/backend/app/services/event_tickets.py b/backend/app/services/event_tickets.py new file mode 100644 index 00000000..2a2f1a9e --- /dev/null +++ b/backend/app/services/event_tickets.py @@ -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] diff --git a/backend/app/services/events.py b/backend/app/services/events.py new file mode 100644 index 00000000..75e5a75e --- /dev/null +++ b/backend/app/services/events.py @@ -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() diff --git a/backend/tests/test_api_gateway.py b/backend/tests/test_api_gateway.py index ad8cdcfd..13207598 100644 --- a/backend/tests/test_api_gateway.py +++ b/backend/tests/test_api_gateway.py @@ -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 diff --git a/backend/tests/test_events.py b/backend/tests/test_events.py new file mode 100644 index 00000000..1a279db8 --- /dev/null +++ b/backend/tests/test_events.py @@ -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 diff --git a/backend/tests/test_openapi_contract.py b/backend/tests/test_openapi_contract.py new file mode 100644 index 00000000..150eed19 --- /dev/null +++ b/backend/tests/test_openapi_contract.py @@ -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 中说明。", + ) diff --git a/docs/features.md b/docs/features.md index 6090cc34..0827b274 100644 --- a/docs/features.md +++ b/docs/features.md @@ -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 会拦下无意识的契约变更。 ### 桌面客户端版本清单 diff --git a/docs/open-platform-plan.md b/docs/open-platform-plan.md index 53e6e488..4c2e4f4b 100644 --- a/docs/open-platform-plan.md +++ b/docs/open-platform-plan.md @@ -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 钩子契约 + 模拟盘拆分试点 diff --git a/examples/open-api/01_pull_market.py b/examples/open-api/01_pull_market.py new file mode 100644 index 00000000..1fa2a5b6 --- /dev/null +++ b/examples/open-api/01_pull_market.py @@ -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"]) diff --git a/examples/open-api/02_write_ext_data.py b/examples/open-api/02_write_ext_data.py new file mode 100644 index 00000000..0b005b38 --- /dev/null +++ b/examples/open-api/02_write_ext_data.py @@ -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"]) diff --git a/examples/open-api/03_run_backtest.py b/examples/open-api/03_run_backtest.py new file mode 100644 index 00000000..9dd2c9c4 --- /dev/null +++ b/examples/open-api/03_run_backtest.py @@ -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]}") diff --git a/examples/open-api/04_events_stream.py b/examples/open-api/04_events_stream.py new file mode 100644 index 00000000..1b8fdd25 --- /dev/null +++ b/examples/open-api/04_events_stream.py @@ -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 帧后退出; 生产环境保持连接持续接收)") diff --git a/examples/open-api/README.md b/examples/open-api/README.md new file mode 100644 index 00000000..d40641d7 --- /dev/null +++ b/examples/open-api/README.md @@ -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` diff --git a/examples/open-api/common.py b/examples/open-api/common.py new file mode 100644 index 00000000..5d723a4b --- /dev/null +++ b/examples/open-api/common.py @@ -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 diff --git a/frontend/src/pages/settings/ApiTokens.tsx b/frontend/src/pages/settings/ApiTokens.tsx index 5e05490f..cb5df817 100644 --- a/frontend/src/pages/settings/ApiTokens.tsx +++ b/frontend/src/pages/settings/ApiTokens.tsx @@ -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() { ))}
- 管理接口 (数据同步 / 扩展表写入 / 设置) 永不开放给 Token。调用方式: Authorization: Bearer tsp_..., 默认限流 120 次/分钟, 契约文档 /api/openapi.json?tier=a + 扩展表的结构配置/上传/拉取、数据同步、设置等管理接口永不开放给 Token (行数据写入走 write:ext)。调用方式: Authorization: Bearer tsp_..., 默认限流 120 次/分钟, 契约文档 /api/openapi.json?tier=a
diff --git a/frontend/src/pages/settings/MenuSettings.tsx b/frontend/src/pages/settings/MenuSettings.tsx index bebc5c06..72b9b4dc 100644 --- a/frontend/src/pages/settings/MenuSettings.tsx +++ b/frontend/src/pages/settings/MenuSettings.tsx @@ -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 = { + '/': { 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 {entry.label} + {(() => { + const arch = archOf(entry) + return ( + + {arch.kind === 'core' ? '核心' : '扩展'} + + ) + })()} {hidden && ( 已隐藏 )} @@ -295,6 +341,17 @@ export function SettingsMenuSettingsPanel() {

拖动左侧手柄调整菜单排列顺序,点击眼睛图标控制菜单在侧边栏中的显示或隐藏。

+

+ 架构归属 (悬停看原因): + + 核心 + 页面承载核心域能力 (数据底座 / 订阅锚点 / 策略回测 / 环境生产) + + + 扩展 + 核心数据的消费页 — 二次开发可用「扩展页面」平行实现替换, 不伤核心 + +