Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion .dockerignore
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ data/
*.sqlite

# Build
frontend/dist/
static/

# Docs
Expand Down
24 changes: 5 additions & 19 deletions Dockerfile
Original file line number Diff line number Diff line change
@@ -1,23 +1,9 @@
# PanWatch Dockerfile
# 多阶段构建,减小最终镜像大小

# ===== Stage 1: 前端构建 =====
FROM node:20-alpine AS frontend-builder

WORKDIR /app/frontend

# 安装 pnpm
RUN npm install -g pnpm

# 复制依赖文件
COPY frontend/package.json frontend/pnpm-lock.yaml ./

# 安装依赖
RUN pnpm install --frozen-lockfile

# 复制源码并构建
COPY frontend/ ./
RUN pnpm build
# ===== Stage 1: 前端静态文件 =====
# 前端已在本机构建(frontend/dist/),直接复制到运行阶段
# 如需在 Docker 内重建,可回退至:使用 node:20-alpine 执行 pnpm install && pnpm build


# ===== Stage 2: Python 运行环境 =====
Expand Down Expand Up @@ -93,8 +79,8 @@ COPY prompts/ ./prompts/
# 写入版本号
RUN echo "${VERSION}" > VERSION

# 从前端构建阶段复制静态文件
COPY --from=frontend-builder /app/frontend/dist ./static/
# 复制本地已构建的前端静态文件
COPY frontend/dist ./static/

# 创建数据目录
RUN mkdir -p /app/data
Expand Down
33 changes: 31 additions & 2 deletions server.py
Original file line number Diff line number Diff line change
Expand Up @@ -360,7 +360,8 @@ def seed_data_sources():
"provider": "xueqiu",
"config": {
"cookies": "",
"description": "雪球个股新闻聚合,需要登录 cookie",
"auto_refresh_waf": True,
"description": "雪球个股新闻聚合,需要登录 cookie;超时/风控时自动重新过 WAF",
},
"enabled": False,
"priority": 0,
Expand Down Expand Up @@ -398,6 +399,21 @@ def seed_data_sources():
"supports_batch": False,
"test_symbols": ["601127", "600519", "300750"],
},
{
"name": "通达信K线",
"type": "kline",
"provider": "tdx",
"config": {
"host": "124.71.187.122:7709",
"binary": "/app/tools/tdx-kline/tdx-kline",
"timeout_sec": 15,
"description": "基于 injoyai/tdx 的通达信日K数据。仅 A 股;需编译 helper 二进制。",
},
"enabled": False,
"priority": 5,
"supports_batch": False,
"test_symbols": ["600519", "000001"],
},
{
"name": "Tushare K线",
"type": "kline",
Expand Down Expand Up @@ -1120,8 +1136,21 @@ async def trigger_agent_for_stock(
market=market,
)

# 从数据库获取 Stock 记录(可能不在自选表中)
stock_id = 0
db_lookup = SessionLocal()
try:
db_stock = db_lookup.query(Stock).filter(
Stock.symbol == stock.symbol,
Stock.market == stock.market,
).first()
if db_stock:
stock_id = db_stock.id
finally:
db_lookup.close()

# 加载该股票的持仓信息
portfolio = load_portfolio_for_stock(stock.id)
portfolio = load_portfolio_for_stock(stock_id)

model, service = resolve_ai_model(agent_name, stock_agent_id)
channels = [] if suppress_notify else resolve_notify_channels(agent_name, stock_agent_id)
Expand Down
19 changes: 14 additions & 5 deletions src/collectors/kline_collector.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
logger = logging.getLogger(__name__)

# 腾讯日K线 API
TENCENT_KLINE_URL = "http://web.ifzq.gtimg.cn/appstock/app/fqkline/get"
TENCENT_KLINE_URL = "http://ifzq.gtimg.cn/appstock/app/fqkline/get"
EASTMONEY_KLINE_URL = "https://push2his.eastmoney.com/api/qt/stock/kline/get"


Expand Down Expand Up @@ -661,7 +661,10 @@ def _fetch_tencent_klines(
klines = _parse_tencent_kline_text(text, tencent_sym)
if klines:
break
last_err = "空响应" # gtimg 突发限流常回空 body,退避后重试
if resp.status_code != 200:
last_err = f"HTTP {resp.status_code}: {text[:100]}"
else:
last_err = "空响应"
except Exception as e:
last_err = e
if attempt < 2:
Expand Down Expand Up @@ -958,9 +961,15 @@ def get_technical_indicators(
kline_pattern=kline_pattern,
)

def get_kline_summary(self, symbol: str) -> dict:
"""获取 K 线摘要(用于 prompt 和前端展示)"""
klines = self.get_klines(symbol, days=120)
def get_kline_summary(self, symbol: str = "", klines: list[KlineData] | None = None) -> dict:
"""获取 K 线摘要(用于 prompt 和前端展示)

Args:
symbol: 股票代码(仅当 klines 为空时用于联网获取)
klines: 外部传入的 K 线数据(非空则跳过联网获取)
"""
if klines is None:
klines = self.get_klines(symbol, days=120)
if not klines:
return {"error": "无K线数据"}
indicators = self.get_technical_indicators(klines=klines)
Expand Down
212 changes: 212 additions & 0 deletions src/core/providers/kline/tdx.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,212 @@
"""通达信 K 线 Provider — 通过 Go helper 调用 injoyai/tdx。

通过编译后的 Go 二进制(tools/tdx-kline/tdx-kline)直连通达信行情服务器。

约束:
- 仅 A 股(CN 市场),港股/美股不支持
- 软依赖:未编译 Go binary 时,fetch 返回明确错误
- 默认使用公开行情服务器 124.71.187.122:7709,可通过 config 自定义 host
"""

from __future__ import annotations

import asyncio
import json
import logging
import os
import subprocess

from src.collectors.kline_collector import KlineData
from src.core.providers.base import KlineProvider, ProviderRequest, ProviderResponse
from src.core.cn_symbol import get_cn_exchange

logger = logging.getLogger(__name__)

# 默认二进制路径(容器内路径)
_DEFAULT_BINARY = "tools/tdx-kline/tdx-kline"

# 默认通达信行情服务器
_DEFAULT_HOST = "124.71.187.122:7709"

# 单次最大拉取条数
_MAX_DAYS = 800

# Go binary 默认超时(秒)
_DEFAULT_TIMEOUT = 5


def _tdx_symbol(symbol: str) -> str:
"""将纯 A 股代码转为通达信前缀格式。

sh600519 / sz000001 / bj...
"""
exchange = get_cn_exchange(symbol)
return f"{exchange.lower()}{symbol}"


class TdxKlineProvider(KlineProvider):
name = "tdx"
supports_markets = {"CN"} # 仅 A 股

def __init__(self, config: dict | None = None):
super().__init__(config=config)
self._binary = (self.config or {}).get("binary", _DEFAULT_BINARY)
self._host = (self.config or {}).get("host", _DEFAULT_HOST)
self._timeout = int((self.config or {}).get("timeout_sec", _DEFAULT_TIMEOUT))
self._init_error = ""

# 检查二进制是否存在
if not os.path.isfile(self._binary):
self._init_error = (
f"通达信 helper 不存在({self._binary}),"
f"请先编译 tools/tdx-kline: "
f"cd tools/tdx-kline && GOOS=linux GOARCH=amd64 go build -o tdx-kline ."
)

def _days(self, req: ProviderRequest) -> int:
for k, v in req.extra:
if k == "days":
try:
return max(1, min(int(v), _MAX_DAYS))
except Exception:
return 60
return 60

def _ktype(self, req: ProviderRequest) -> str:
"""从 extra 中提取 K 线类型,默认 day"""
for k, v in req.extra:
if k == "type":
if v in ("5m", "15m", "30m", "60m"):
return v
return "day"

async def fetch(self, req: ProviderRequest) -> ProviderResponse:
if self._init_error:
return ProviderResponse(success=False, error=self._init_error)
if not req.symbols:
return ProviderResponse(success=True, data=[])
if len(req.symbols) > 1:
return ProviderResponse(
success=False, error="TdxKlineProvider 仅支持单 symbol"
)
if req.market != "CN":
return ProviderResponse(
success=False, error="TdxKlineProvider 仅支持 CN 市场"
)

symbol = req.symbols[0]
days = self._days(req)
ktype = self._ktype(req)
tdx_sym = _tdx_symbol(symbol)
logger.info(
"TdxKlineProvider 调用 helper: binary=%s symbol=%s tdx_sym=%s host=%s days=%d type=%s",
self._binary, symbol, tdx_sym, self._host, days, ktype,
)

def _run_helper() -> ProviderResponse:
"""同步函数: 在 asyncio.to_thread 中执行。"""
try:
proc = subprocess.run(
[
self._binary,
"--symbol", tdx_sym,
"--days", str(days),
"--host", self._host,
"--timeout", str(self._timeout),
"--type", ktype,
],
capture_output=True,
text=True,
timeout=self._timeout + 2, # 额外给 2s 缓冲
)
except FileNotFoundError:
return ProviderResponse(
success=False,
error=f"通达信 helper 不存在({self._binary})",
)
except subprocess.TimeoutExpired:
return ProviderResponse(
success=False,
error=f"通达信接口超时({self._host}, {self._timeout}s)",
)
except Exception as e:
return ProviderResponse(
success=False, error=f"通达信 helper 执行失败: {e}"
)

if proc.returncode != 0:
err_msg = proc.stderr.strip() or "未知错误"
return ProviderResponse(
success=False, error=f"通达信接口调用失败: {err_msg}"
)

# 解析 JSON: Go helper 内部日志也输出到 stdout(含 ANSI 颜色码)
stdout = proc.stdout
# 找内容为 "[{" 开始的 JSON 数组(可能前有 ANSI 码)
idx = stdout.find('[{"date"')
if idx < 0:
idx = stdout.rfind("[")
if idx >= 0:
# 只取 JSON 部分(从 [ 到匹配的 ]),避免尾部日志干扰
decoder = json.JSONDecoder()
try:
raw, _ = decoder.raw_decode(stdout, idx)
except json.JSONDecodeError as e:
return ProviderResponse(
success=False,
error=f"通达信返回数据解析失败: {e}",
)
else:
return ProviderResponse(
success=False,
error=f"通达信未返回有效 JSON: {stdout[:200]}",
)

if not raw:
return ProviderResponse(
success=False,
error=f"通达信未返回数据(symbol={symbol})",
)

klines: list[KlineData] = []
for item in raw:
try:
klines.append(
KlineData(
date=str(item["date"]),
open=float(item["open"]),
close=float(item["close"]),
high=float(item["high"]),
low=float(item["low"]),
volume=float(item["volume"]),
)
)
except (KeyError, ValueError, TypeError) as e:
logger.debug(f"通达信 JSON 条目解析失败: {e}, item={item}")
continue

if not klines:
return ProviderResponse(
success=False,
error=f"通达信数据解析后为空(symbol={symbol})",
)

# 按日期升序排列
klines.sort(key=lambda k: k.date)

return ProviderResponse(success=True, data=klines[-days:])

return await asyncio.to_thread(_run_helper)

async def health_check(self) -> bool:
if self._init_error:
return False
try:
resp = await self.fetch(
ProviderRequest(
symbols=("600519",), market="CN", extra=(("days", 5),)
)
)
return resp.success and not resp.is_empty
except Exception:
return False
2 changes: 2 additions & 0 deletions src/core/providers/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -299,9 +299,11 @@ def get_kline_orchestrator() -> KlineOrchestrator:
orch = KlineOrchestrator()
from src.core.providers.kline.tencent import TencentKlineProvider
from src.core.providers.kline.tushare import TushareKlineProvider
from src.core.providers.kline.tdx import TdxKlineProvider
from src.core.providers.kline.yfinance import YFinanceKlineProvider

orch.register("tencent", lambda cfg: TencentKlineProvider(config=cfg))
orch.register("tdx", lambda cfg: TdxKlineProvider(config=cfg))
orch.register("tushare", lambda cfg: TushareKlineProvider(config=cfg))
orch.register("yfinance", lambda cfg: YFinanceKlineProvider(config=cfg))
_kline_orchestrator = orch
Expand Down
Loading