mirror of
https://github.com/ZhuLinsen/daily_stock_analysis.git
synced 2026-10-06 14:33:11 +08:00
feat: 新增富途Futu真实持仓导入支持 (#2042)
* feat: add Futu portfolio import support * docs: condense Futu changelog entry * fix: tighten Futu portfolio import contracts
This commit is contained in:
@@ -8,6 +8,14 @@
|
||||
# 深市:000xxx, 002xxx, 300xxx
|
||||
STOCK_LIST=600519,300750,002594
|
||||
|
||||
# Futu OpenD(持仓导入)
|
||||
# --portfolio futu 只读取 ACTIVE REAL NORMAL / MASTER 账户中的沪深 A 股、港股、美股 LONG 正股持仓。
|
||||
# futu-api 10.8 仅支持 IPv4;Docker 连接宿主机 OpenD 时请勿使用容器内的 127.0.0.1,详见 docs/full-guide.md。
|
||||
# FUTU_OPEND_HOST=127.0.0.1
|
||||
# FUTU_OPEND_PORT=11111
|
||||
# FUTU_SECURITY_FIRM=NONE # 可选;默认由 OpenD 自动识别,也可显式指定券商
|
||||
# FUTU_ACC_ID= # 可选;正整数,指定后只读取该真实账户
|
||||
|
||||
# Anspire Open API Keys(支持多个,逗号分隔)
|
||||
# 获取: https://open.anspire.cn/?share_code=QFBC0FYC
|
||||
# 在未配置更高优先级 OpenAI-compatible 来源时,满足条件可复用该 key 给 Anspire 大模型网关与新闻搜索。
|
||||
|
||||
@@ -54,6 +54,7 @@ jobs:
|
||||
run: |
|
||||
pip install --upgrade pip
|
||||
pip install -r requirements.txt
|
||||
python -c "import futu; print('✅ futu SDK import OK')"
|
||||
|
||||
- name: 创建必要目录
|
||||
run: |
|
||||
|
||||
@@ -14,6 +14,7 @@ jobs:
|
||||
runs-on: ubuntu-latest
|
||||
outputs:
|
||||
frontend: ${{ steps.filter.outputs.frontend }}
|
||||
futu_packaging: ${{ steps.filter.outputs.futu_packaging }}
|
||||
steps:
|
||||
- name: 📥 Checkout
|
||||
uses: actions/checkout@v5
|
||||
@@ -24,6 +25,17 @@ jobs:
|
||||
filters: |
|
||||
frontend:
|
||||
- 'apps/dsa-web/**'
|
||||
futu_packaging:
|
||||
- 'requirements.txt'
|
||||
- 'main.py'
|
||||
- 'src/brokers/futu/**'
|
||||
- 'scripts/build-backend.ps1'
|
||||
- 'scripts/build-backend-macos.sh'
|
||||
- 'scripts/build-all.ps1'
|
||||
- 'scripts/build-all-macos.sh'
|
||||
- 'apps/dsa-desktop/**'
|
||||
- '.github/workflows/ci.yml'
|
||||
- '.github/workflows/desktop-release.yml'
|
||||
|
||||
ai-governance:
|
||||
name: ai-governance
|
||||
@@ -67,6 +79,8 @@ jobs:
|
||||
echo "Dependency install attempt ${attempt} failed, retrying in 15s..." >&2
|
||||
sleep 15
|
||||
done
|
||||
- name: 🔌 Verify global Futu SDK
|
||||
run: python -c "import futu; print('✅ futu SDK import OK')"
|
||||
- name: ✅ Python syntax check
|
||||
run: ./scripts/ci_gate.sh syntax
|
||||
- name: ✅ Flake8 critical checks
|
||||
@@ -107,9 +121,62 @@ jobs:
|
||||
from src.patches.eastmoney_patch import eastmoney_patch; print('✅ patch')
|
||||
from bot.dispatcher import CommandDispatcher; print('✅ bot')
|
||||
from api.app import app; print('✅ api')
|
||||
import futu; print('✅ futu')
|
||||
print('✅ All Docker imports OK')
|
||||
"
|
||||
|
||||
desktop-futu-package-windows:
|
||||
name: desktop-futu-package-windows
|
||||
runs-on: windows-latest
|
||||
needs: [changes, ai-governance]
|
||||
if: needs.changes.outputs.futu_packaging == 'true'
|
||||
steps:
|
||||
- name: 📥 Checkout
|
||||
uses: actions/checkout@v5
|
||||
- name: 🐍 Setup Python
|
||||
uses: actions/setup-python@v6
|
||||
with:
|
||||
python-version: '3.12'
|
||||
cache: 'pip'
|
||||
cache-dependency-path: requirements.txt
|
||||
- name: 🟢 Setup Node
|
||||
uses: actions/setup-node@v6
|
||||
with:
|
||||
node-version: '20'
|
||||
cache: 'npm'
|
||||
cache-dependency-path: apps/dsa-web/package-lock.json
|
||||
- name: 📦 Install Web dependencies
|
||||
shell: pwsh
|
||||
run: npm ci --prefix apps/dsa-web
|
||||
- name: 🧊 Build and verify frozen backend
|
||||
shell: pwsh
|
||||
run: powershell -ExecutionPolicy Bypass -File scripts/build-backend.ps1
|
||||
|
||||
desktop-futu-package-macos:
|
||||
name: desktop-futu-package-macos
|
||||
runs-on: macos-15
|
||||
needs: [changes, ai-governance]
|
||||
if: needs.changes.outputs.futu_packaging == 'true'
|
||||
steps:
|
||||
- name: 📥 Checkout
|
||||
uses: actions/checkout@v5
|
||||
- name: 🐍 Setup Python
|
||||
uses: actions/setup-python@v6
|
||||
with:
|
||||
python-version: '3.12'
|
||||
cache: 'pip'
|
||||
cache-dependency-path: requirements.txt
|
||||
- name: 🟢 Setup Node
|
||||
uses: actions/setup-node@v6
|
||||
with:
|
||||
node-version: '20'
|
||||
cache: 'npm'
|
||||
cache-dependency-path: apps/dsa-web/package-lock.json
|
||||
- name: 📦 Install Web dependencies
|
||||
run: npm ci --prefix apps/dsa-web
|
||||
- name: 🧊 Build and verify frozen backend
|
||||
run: bash scripts/build-backend-macos.sh
|
||||
|
||||
web-gate:
|
||||
name: web-gate
|
||||
runs-on: ubuntu-latest
|
||||
|
||||
@@ -136,6 +136,7 @@ jobs:
|
||||
from src.notification import NotificationService; print('ok-notification')
|
||||
from data_provider import DataFetcherManager; print('ok-data-provider')
|
||||
from src.analyzer import GeminiAnalyzer; print('ok-analyzer')
|
||||
import futu; print('ok-futu')
|
||||
print('release-smoke-ok')
|
||||
"
|
||||
|
||||
|
||||
@@ -96,6 +96,7 @@ jobs:
|
||||
from src.notification import NotificationService; print('ok-notification')
|
||||
from data_provider import DataFetcherManager; print('ok-data-provider')
|
||||
from src.analyzer import GeminiAnalyzer; print('ok-analyzer')
|
||||
import futu; print('ok-futu')
|
||||
print('manual-smoke-ok')
|
||||
"
|
||||
|
||||
|
||||
+4
-2
@@ -58,8 +58,10 @@ COPY strategies/ ./strategies/
|
||||
COPY --from=web-builder /app/static ./static/
|
||||
COPY docker/entrypoint.sh /usr/local/bin/docker-entrypoint.sh
|
||||
|
||||
# Verify the bundled AlphaSift adapter from requirements.
|
||||
RUN python -c "import alphasift.dsa_adapter"
|
||||
# Verify globally bundled runtime dependencies without leaving SDK logs in the image.
|
||||
RUN mkdir -p /tmp/futu-sdk-smoke-home && \
|
||||
HOME=/tmp/futu-sdk-smoke-home python -c "import alphasift.dsa_adapter; import futu" && \
|
||||
rm -rf /tmp/futu-sdk-smoke-home
|
||||
|
||||
# 确保数据目录存在并授权给 non-root 用户
|
||||
RUN mkdir -p /app/data /app/logs /app/reports && \
|
||||
|
||||
@@ -8,6 +8,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
|
||||
> For user-friendly release highlights, see the [GitHub Releases](https://github.com/ZhuLinsen/daily_stock_analysis/releases) page.
|
||||
|
||||
## [Unreleased]
|
||||
- [新功能] 新增 `--portfolio futu`,只读导入 Futu OpenD 真实账户的沪深 A 股、港股、美股 LONG 正股持仓作为分析列表。
|
||||
<!-- 新条目格式:- [类型] 描述(类型取值:新功能/改进/修复/文档/测试/chore)-->
|
||||
<!-- 每条独立一行追加到本段末尾,无需分类标题,合并时冲突最小 -->
|
||||
- [修复] #2026 外股代码映射到中文显示名时英文新闻相关性判定漏判:新增同源 STOCK_ENGLISH_NAME_MAP 单一真源、canonicalize_foreign_stock_code 规范化入口与 _foreign_english_query_terms 别名解析,使 AAPL/00700/BABA 等 ticker 即使 stock_name 为中文也能在查询构建、相关性打分与多维度情报路径上复用 canonical 英文名,并补齐 .US/.HK suffix / HK 前缀全形式的归类与回归用例;同时在 _score_news_relevance 对 alias 展开 term 做去重,避免 legal alias 展开短名与显式 short alias 重复计分。
|
||||
|
||||
@@ -8,6 +8,7 @@
|
||||
- Electron 启动时自动拉起后端服务,等待 `/api/health` 就绪后加载 UI
|
||||
- Windows 便携/安装模式下,用户配置文件 `.env` 和数据库放在 exe 同级目录;macOS 打包版使用 Electron 用户数据目录保存运行时配置
|
||||
- 桌面端会自动从本机 `8000-8100` 选择可用端口,并把实际选择的端口同步给内置后端;桌面端不依赖 `.env` 里的 `WEBUI_PORT` 来决定窗口连接地址,避免用户改端口后 Electron 仍等待旧端口导致启动超时
|
||||
- Desktop backend 默认随 `requirements.txt` 安装并冻结 `futu-api==10.8.6808`;Windows/macOS 构建脚本会在源码环境和 PyInstaller 产物中分别执行 `import futu`,防止发布包只安装但未携带 SDK。
|
||||
|
||||
## 本地开发
|
||||
|
||||
@@ -206,7 +207,7 @@ npm install
|
||||
npm run build
|
||||
```
|
||||
|
||||
2) 按现有脚本打包 Python 后端(脚本已内置 AlphaSift 与 AkShare 数据文件收集)
|
||||
2) 按现有脚本打包 Python 后端(脚本已内置 AlphaSift、Futu SDK 与 AkShare 数据文件收集)
|
||||
|
||||
- Windows:
|
||||
|
||||
@@ -220,7 +221,7 @@ powershell -ExecutionPolicy Bypass -File scripts\build-backend.ps1
|
||||
bash scripts/build-backend-macos.sh
|
||||
```
|
||||
|
||||
该脚本会在安装依赖后执行 `--collect-all alphasift` 和 `--collect-data akshare`。构建完成后会校验 `alphasift.dsa_adapter` 可导入,并确认 AkShare 的 `file_fold/calendar.json` 已进入冻结产物,避免发行包在热点题材或日线增强路径中因缺少 package data 降级。
|
||||
该脚本会在安装依赖后执行 `--collect-all alphasift`、`--collect-all futu` 和 `--collect-data akshare`。构建完成后会通过冻结可执行文件校验 `alphasift.dsa_adapter`、`futu`、`orjson` 均可导入,并确认 AkShare 的 `file_fold/calendar.json` 已进入冻结产物,避免发行包在热点题材、Futu 持仓导入或日线增强路径中因缺少依赖/package data 降级。PR 主 CI 在 `requirements.txt`、Futu broker、Desktop 打包入口或相关 workflow 变化时,会分别运行 `desktop-futu-package-windows` 与 `desktop-futu-package-macos` 阻断检查。
|
||||
|
||||
3) 打包 Electron 桌面应用
|
||||
|
||||
|
||||
@@ -389,6 +389,17 @@ daily_stock_analysis/
|
||||
|
||||
兼容与回退说明:该改动不新增/修改模型、provider、Base URL、LiteLLM route、配置清理或回写逻辑;若出现异常,只能通过回滚本次提交恢复旧排序行为,不涉及历史配置迁移。
|
||||
|
||||
### Futu 持仓导入配置
|
||||
|
||||
| 变量名 | 说明 | 默认值 | 必填 |
|
||||
|--------|------|--------|:----:|
|
||||
| `FUTU_OPEND_HOST` | OpenD 地址;锁定的 `futu-api==10.8.6808` 仅支持 IPv4 地址或可解析到 IPv4 的主机名。跨主机连接只应使用受信网络或本机端口转发。 | `127.0.0.1` | 可选 |
|
||||
| `FUTU_OPEND_PORT` | OpenD 端口,合法范围 `1-65535`。 | `11111` | 可选 |
|
||||
| `FUTU_SECURITY_FIRM` | Futu `SecurityFirm` 枚举名;`NONE` 表示使用 SDK 官方自动识别一次,也可显式指定券商。 | `NONE` | 可选 |
|
||||
| `FUTU_ACC_ID` | 指定一个符合条件的 REAL 账户 ID;留空时合并所有状态为 `ACTIVE` 的 `NORMAL`(普通)和 `MASTER`(主)证券账户。账户 ID 应按敏感配置处理,不要提交到仓库。 | 空 | 可选 |
|
||||
|
||||
`MASTER` 仅表示 Futu 的主账户角色,不表示账户具有只读属性。本集成的只读边界来自它只调用账户、持仓和证券信息查询接口,不调用交易解锁、下单、改单或撤单接口。
|
||||
|
||||
### 数据源配置
|
||||
|
||||
| 变量名 | 说明 | 默认值 | 必填 |
|
||||
@@ -680,6 +691,7 @@ python main.py # 完整分析(个股 + 大盘复盘)
|
||||
python main.py --market-review # 仅大盘复盘
|
||||
python main.py --no-market-review # 仅个股分析
|
||||
python main.py --stocks 600519,300750 # 指定股票
|
||||
python main.py --portfolio futu # 使用 Futu 真实 LONG 正股持仓(覆盖 --stocks/STOCK_LIST)
|
||||
python main.py --dry-run # 仅获取数据,不 AI 分析
|
||||
python main.py --no-notify # 不发送推送
|
||||
python main.py --schedule # 定时任务模式
|
||||
@@ -688,6 +700,25 @@ python main.py --debug # 调试模式(详细日志)
|
||||
python main.py --workers 5 # 指定并发数
|
||||
```
|
||||
|
||||
### Futu 真实持仓作为分析列表
|
||||
|
||||
标准源码安装(`pip install -r requirements.txt`)、官方 Docker 镜像和 Windows/macOS Desktop backend 已默认包含锁定的 `futu-api==10.8.6808`。仅在使用裁剪过的自定义 Python 环境时,才需要按 [Futu OpenAPI SDK 安装说明](https://openapi.futunn.com/futu-api-doc/en/intro/intro.html) 手动补装。启动并登录 Futu OpenD 后运行:
|
||||
|
||||
```bash
|
||||
# 仅裁剪过的自定义环境需要执行下一行
|
||||
pip install "futu-api==10.8.6808"
|
||||
# 所有标准安装均可直接运行
|
||||
python main.py --portfolio futu
|
||||
```
|
||||
|
||||
`--portfolio futu` 固定读取状态明确为 `ACTIVE` 的 `REAL` 真实证券账户,并在每次分析开始前用 `refresh_cache=True` 刷新持仓;状态缺失、`N/A`、未知或 `DISABLED` 的账户一律拒绝。未设置 `FUTU_ACC_ID` 时会合并所有可用的 `NORMAL`(普通)及 `MASTER`(主)证券账户并按代码去重;设置后只读取指定的正整数账户 ID。根据 [Futu `get_acc_list` 账户角色定义](https://openapi.futunn.com/futu-api-doc/trade/get-acc-list.html),`MASTER` 表示主账户而非只读属性,马来西亚 `IPO` 账户不属于本功能的持仓来源并会被跳过。本集成的只读边界来自它只调用查询接口。
|
||||
|
||||
只有持仓方向明确为 `LONG`、Futu 静态类型为 `STOCK` 且数量非零的正股持仓会进入分析;`SHORT`、方向未知、期权、ETF、窝轮、期货等持仓会被排除。Futu 持仓代码转换仅支持沪深 A 股、港股和美股;沪深 B 股、日股及其他 Futu 市场持仓会在日志中列出代码并跳过,这不改变手工股票列表的既有市场支持边界。如果可用账户 ID 无效,或 `LONG` 持仓数量无效、非零 `LONG` 持仓代码无效、静态类型缺失 / 未知,或已确认的正股代码无法转换为当前分析格式,整次持仓导入会明确失败,不会返回静默截断的部分结果。
|
||||
|
||||
OpenD 默认地址为 `127.0.0.1:11111`,可用 `FUTU_OPEND_HOST` / `FUTU_OPEND_PORT` 覆盖。锁定的 `futu-api==10.8.6808` 网络层使用 IPv4 socket,因此 `FUTU_OPEND_HOST` 应填写 IPv4 地址或可解析到 IPv4 的主机名,不支持 `::1` 等 IPv6 地址。在 Docker 容器中,`127.0.0.1` 指向容器自身;OpenD 运行在宿主机时,macOS / Windows 可设置 `FUTU_OPEND_HOST=host.docker.internal`,Linux 需要先为容器增加 `host.docker.internal:host-gateway` 映射后再使用该主机名。跨主机连接会传输真实账户与持仓信息;[Futu 官方建议实盘连接配置协议加密](https://openapi.futunn.com/futu-api-doc/en/ftapi/protocol.html)。本功能不修改进程级 SDK 加密配置,建议优先让 OpenD 与本程序同机,或使用受信网络 / 本机端口转发。未设置 `FUTU_SECURITY_FIRM` 时只使用 Futu SDK 官方的 `SecurityFirm.NONE` 自动识别一次,不会枚举多个券商或在部分探测失败后静默拼接结果;需要固定券商时可显式配置该变量。
|
||||
|
||||
若同时传入 `--stocks`,Futu 持仓优先;定时模式会在每轮执行前重新读取真实持仓,而不是复用启动时快照。若没有符合条件的 Futu 持仓,本轮会跳过个股分析且不会回退到 `STOCK_LIST`;已启用的大盘复盘仍按原配置执行,大盘复盘也未请求时不会刷新股票索引或构造分析管线,已启用的自动回测仍作为独立步骤执行。单次 CLI 仅在 SDK、OpenD、账户发现、持仓读取或证券分类等持仓解析边界失败时返回非零退出码;持仓解析成功后的交易日历、分析管线和报告异常仍沿用原分析流程的记录与容错语义。已启动服务与定时调度会记录持仓导入错误并继续运行。该能力只读取账户和持仓,不执行下单、改单、撤单或交易解锁。现有分析日志会记录本轮股票代码,但不会记录账户 ID、持仓数量、成本或资金;分享运行日志前请按需脱敏。
|
||||
|
||||
---
|
||||
|
||||
## 定时任务配置
|
||||
|
||||
@@ -327,6 +327,17 @@ For the notification baseline, diagnostics, and deployment notes, see [Notificat
|
||||
|
||||
> Behavior note: Search and social sentiment are optional enhancement services. If either service fails to initialize, the system logs a warning and degrades gracefully by skipping that stage without blocking the core analysis flow.
|
||||
|
||||
### Futu Portfolio Import Configuration
|
||||
|
||||
| Variable | Description | Default | Required |
|
||||
|--------|------|--------|:----:|
|
||||
| `FUTU_OPEND_HOST` | OpenD host. The pinned `futu-api==10.8.6808` accepts an IPv4 address or a hostname that resolves to IPv4. Cross-host connections should use only a trusted network or local port forwarding. | `127.0.0.1` | Optional |
|
||||
| `FUTU_OPEND_PORT` | OpenD port in the range `1-65535`. | `11111` | Optional |
|
||||
| `FUTU_SECURITY_FIRM` | Futu `SecurityFirm` enum name. `NONE` performs the SDK's official auto-detection once; set an explicit broker when required. | `NONE` | Optional |
|
||||
| `FUTU_ACC_ID` | Select one eligible REAL account ID. When empty, all explicitly `ACTIVE` `NORMAL` and `MASTER` securities accounts are merged. Treat account IDs as sensitive configuration and do not commit them. | empty | Optional |
|
||||
|
||||
`MASTER` is Futu's master-account role, not a read-only attribute. This integration is read-only because it calls only account, position, and security-information queries; it never unlocks trading or places, modifies, or cancels orders.
|
||||
|
||||
### Data Source Configuration
|
||||
|
||||
| Variable | Description | Default | Required |
|
||||
@@ -618,6 +629,7 @@ python main.py # Full analysis (stocks + market review)
|
||||
python main.py --market-review # Market review only
|
||||
python main.py --no-market-review # Stock analysis only
|
||||
python main.py --stocks 600519,300750 # Specify stocks
|
||||
python main.py --portfolio futu # Use real Futu LONG stock holdings (overrides --stocks/STOCK_LIST)
|
||||
python main.py --dry-run # Fetch data only, no AI analysis
|
||||
python main.py --no-notify # Don't send notifications
|
||||
python main.py --schedule # Scheduled task mode
|
||||
@@ -625,6 +637,25 @@ python main.py --debug # Debug mode (verbose logging)
|
||||
python main.py --workers 5 # Specify concurrency
|
||||
```
|
||||
|
||||
### Use real Futu holdings as the analysis list
|
||||
|
||||
Standard source installs (`pip install -r requirements.txt`), official Docker images, and Windows/macOS Desktop backends already include the pinned `futu-api==10.8.6808`. Install it manually from the [Futu OpenAPI SDK guide](https://openapi.futunn.com/futu-api-doc/en/intro/intro.html) only when using a reduced custom Python environment. After starting and signing in to Futu OpenD, run:
|
||||
|
||||
```bash
|
||||
# Only reduced custom environments need the next line
|
||||
pip install "futu-api==10.8.6808"
|
||||
# Standard installs can run the command directly
|
||||
python main.py --portfolio futu
|
||||
```
|
||||
|
||||
`--portfolio futu` only reads `REAL` securities accounts whose status is explicitly `ACTIVE`, and refreshes positions with `refresh_cache=True` before each analysis run. Accounts with a missing, `N/A`, unknown, or `DISABLED` status are rejected. Without `FUTU_ACC_ID`, it merges all usable `NORMAL` and `MASTER` securities accounts and deduplicates symbols; when set, only that positive integer account ID is read. Per the [Futu `get_acc_list` account-role contract](https://openapi.futunn.com/futu-api-doc/en/trade/get-acc-list.html), `MASTER` means the master-account role rather than a read-only attribute, and Malaysian `IPO` accounts are not portfolio sources. The integration is read-only because it calls only query APIs.
|
||||
|
||||
Only non-zero positions whose direction is explicitly `LONG` and whose Futu static type is `STOCK` are analyzed. `SHORT`, unknown-direction, option, ETF, warrant, futures, and other non-stock positions are excluded. Futu portfolio conversion is limited to Shanghai/Shenzhen A-shares, HK stocks, and US stocks. Shanghai/Shenzhen B-shares, JP holdings, and holdings from other Futu markets are logged with their codes and skipped; this does not change the market support for manually configured stock lists. An invalid eligible account ID, an invalid quantity on a `LONG` position, an invalid or missing code on a non-zero `LONG` position, a missing or unknown static type, or a confirmed stock code that cannot be converted to the analysis format fails the whole import instead of returning a silently truncated result.
|
||||
|
||||
OpenD defaults to `127.0.0.1:11111`; override it with `FUTU_OPEND_HOST` / `FUTU_OPEND_PORT`. The pinned `futu-api==10.8.6808` networking layer uses IPv4 sockets, so `FUTU_OPEND_HOST` must be an IPv4 address or a hostname that resolves to IPv4; IPv6 addresses such as `::1` are unsupported. Inside a Docker container, `127.0.0.1` refers to the container itself. When OpenD runs on the host, set `FUTU_OPEND_HOST=host.docker.internal` on macOS or Windows; on Linux, add a `host.docker.internal:host-gateway` mapping to the container before using that hostname. Cross-host connections carry real account and position data, and [Futu recommends protocol encryption for real-trading connections](https://openapi.futunn.com/futu-api-doc/en/ftapi/protocol.html). This integration does not modify process-wide SDK encryption settings; prefer running OpenD on the same host, or use a trusted network or local port forwarding. When `FUTU_SECURITY_FIRM` is unset, discovery makes one call with the Futu SDK's official `SecurityFirm.NONE` auto-detection mode; it does not enumerate brokers or silently combine partial probe results. Set the variable explicitly when a fixed broker is required.
|
||||
|
||||
If `--stocks` is also present, the Futu portfolio wins. Scheduled mode reloads real positions for every run instead of reusing a startup snapshot. If no Futu holdings qualify, stock analysis is skipped without falling back to `STOCK_LIST`; an enabled market review still runs according to its existing configuration. When no market review is requested either, the run does not refresh the stock index or construct the analysis pipeline; an enabled auto-backtest still runs as an independent step. A one-shot CLI exits non-zero only when SDK, OpenD, account discovery, position loading, or security classification fails inside the portfolio-resolution boundary. Trading-calendar, pipeline, and report failures after a successful portfolio import retain the existing analysis error semantics. An already running service or scheduler logs portfolio import errors and continues. This integration only reads accounts and positions; it does not place, modify, cancel, or unlock trades. Existing analysis logs include the stock symbols for the current run, but not account IDs, quantities, costs, or cash balances; redact those symbols as needed before sharing logs.
|
||||
|
||||
---
|
||||
|
||||
## Scheduled Task Configuration
|
||||
|
||||
@@ -72,6 +72,8 @@ from datetime import date, datetime, timezone, timedelta
|
||||
from src.webui_frontend import prepare_webui_frontend_assets
|
||||
from src.config import get_config, Config
|
||||
from src.logging_config import setup_logging
|
||||
from src.brokers.futu.portfolio import FutuPortfolioError
|
||||
from data_provider.base import canonical_stock_code
|
||||
from src.services.stock_list_parser import split_stock_list
|
||||
from src.services.stock_code_utils import resolve_index_stock_code_for_analysis
|
||||
|
||||
@@ -277,6 +279,7 @@ def parse_arguments() -> argparse.Namespace:
|
||||
python main.py --debug # 调试模式
|
||||
python main.py --dry-run # 仅获取数据,不进行 AI 分析
|
||||
python main.py --stocks 600519,000001 # 指定分析特定股票
|
||||
python main.py --portfolio futu # 使用 Futu 真实正股持仓(覆盖 --stocks)
|
||||
python main.py --no-notify # 不发送推送通知
|
||||
python main.py --check-notify # 检查通知配置,不发送通知
|
||||
python main.py --single-notify # 启用单股推送模式(每分析完一只立即推送)
|
||||
@@ -303,6 +306,13 @@ def parse_arguments() -> argparse.Namespace:
|
||||
help='指定要分析的股票代码,逗号分隔(覆盖配置文件)'
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
'--portfolio',
|
||||
type=str.lower,
|
||||
choices=('futu',),
|
||||
help='使用券商真实持仓作为股票列表;当前支持 futu,并覆盖 --stocks/STOCK_LIST'
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
'--no-notify',
|
||||
action='store_true',
|
||||
@@ -524,6 +534,25 @@ def _refresh_stock_index_cache_for_analysis(config: Config) -> None:
|
||||
logger.warning("[stock-index] 分析前刷新股票索引失败,继续执行分析: %s", exc)
|
||||
|
||||
|
||||
def _resolve_portfolio_stock_codes(args: argparse.Namespace) -> Optional[List[str]]:
|
||||
"""Resolve an optional broker portfolio into the analysis stock list."""
|
||||
portfolio = str(getattr(args, "portfolio", "") or "").strip().lower()
|
||||
if not portfolio:
|
||||
return None
|
||||
if portfolio != "futu": # argparse prevents this for CLI callers; keep API callers safe.
|
||||
raise ValueError(f"不支持的 portfolio: {portfolio}")
|
||||
|
||||
from src.brokers.futu.portfolio import load_futu_stock_codes
|
||||
|
||||
stock_codes = [
|
||||
canonical_stock_code(code)
|
||||
for code in load_futu_stock_codes()
|
||||
if (code or "").strip()
|
||||
]
|
||||
logger.info("portfolio=futu 已覆盖 stocks/STOCK_LIST,使用 %d 只真实正股", len(stock_codes))
|
||||
return stock_codes
|
||||
|
||||
|
||||
def _prime_daily_market_context(
|
||||
config: Config,
|
||||
pipeline: Any,
|
||||
@@ -660,6 +689,32 @@ def _save_reused_market_review_report(
|
||||
logger.warning("复用大盘上下文保存大盘复盘报告失败: %s", exc)
|
||||
|
||||
|
||||
def _run_auto_backtest(config: Config) -> None:
|
||||
"""Run the independently configured auto-backtest without failing analysis."""
|
||||
|
||||
try:
|
||||
if not getattr(config, 'backtest_enabled', False):
|
||||
return
|
||||
|
||||
from src.services.backtest_service import BacktestService
|
||||
|
||||
logger.info("开始自动回测...")
|
||||
service = BacktestService()
|
||||
stats = service.run_backtest(
|
||||
force=False,
|
||||
eval_window_days=getattr(config, 'backtest_eval_window_days', 10),
|
||||
min_age_days=getattr(config, 'backtest_min_age_days', 14),
|
||||
limit=200,
|
||||
)
|
||||
logger.info(
|
||||
f"自动回测完成: processed={stats.get('processed')} "
|
||||
f"saved={stats.get('saved')} completed={stats.get('completed')} "
|
||||
f"insufficient={stats.get('insufficient')} errors={stats.get('errors')}"
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.warning(f"自动回测失败(已忽略): {exc}")
|
||||
|
||||
|
||||
def run_full_analysis(
|
||||
config: Config,
|
||||
args: argparse.Namespace,
|
||||
@@ -670,8 +725,26 @@ def run_full_analysis(
|
||||
"""
|
||||
执行完整的分析流程(个股 + 大盘复盘)
|
||||
|
||||
这是定时任务调用的主函数
|
||||
这是定时任务调用的主函数。Futu 持仓解析失败始终传播给调用方;
|
||||
``raise_errors`` 只控制持仓解析成功后的分析流程异常语义。
|
||||
"""
|
||||
# Portfolio resolution is its own CLI contract boundary. A broker import
|
||||
# failure must reach the one-shot caller, while all later work keeps the
|
||||
# existing run_full_analysis return-value semantics.
|
||||
portfolio_stock_codes = _resolve_portfolio_stock_codes(args)
|
||||
portfolio_is_empty = portfolio_stock_codes == []
|
||||
market_review_requested = (
|
||||
getattr(config, 'market_review_enabled', False)
|
||||
and not getattr(args, 'no_market_review', False)
|
||||
)
|
||||
if portfolio_is_empty and not market_review_requested:
|
||||
logger.info(
|
||||
"真实账户中无符合条件的 Futu 持仓,"
|
||||
"本轮跳过个股分析和大盘复盘。"
|
||||
)
|
||||
_run_auto_backtest(config)
|
||||
return True
|
||||
|
||||
# Import pipeline modules outside the broad try/except so that import-time
|
||||
# failures propagate to the caller instead of being silently swallowed.
|
||||
from src.core.market_review import run_market_review
|
||||
@@ -679,9 +752,11 @@ def run_full_analysis(
|
||||
|
||||
try:
|
||||
_refresh_stock_index_cache_for_analysis(config)
|
||||
if portfolio_stock_codes is not None:
|
||||
stock_codes = portfolio_stock_codes
|
||||
|
||||
# Issue #529: Hot-reload STOCK_LIST from .env on each scheduled run
|
||||
if stock_codes is None:
|
||||
if stock_codes is None and portfolio_stock_codes is None:
|
||||
config.refresh_stock_list()
|
||||
|
||||
# Issue #373: Trading day filter (per-stock, per-market)
|
||||
@@ -690,14 +765,24 @@ def run_full_analysis(
|
||||
config, args, effective_codes
|
||||
)
|
||||
if should_skip:
|
||||
logger.info(
|
||||
"今日所有相关市场均为非交易日,跳过执行。可使用 --force-run 强制执行。"
|
||||
)
|
||||
if portfolio_is_empty:
|
||||
logger.info(
|
||||
"真实账户中无符合条件的 Futu 持仓,"
|
||||
"本轮无需执行个股分析或大盘复盘,跳过执行。"
|
||||
)
|
||||
else:
|
||||
logger.info(
|
||||
"今日所有相关市场均为非交易日,跳过执行。"
|
||||
"可使用 --force-run 强制执行。"
|
||||
)
|
||||
return True
|
||||
if set(filtered_codes) != set(effective_codes):
|
||||
skipped = set(effective_codes) - set(filtered_codes)
|
||||
logger.info("今日休市股票已跳过: %s", skipped)
|
||||
stock_codes = filtered_codes
|
||||
skip_futu_stock_analysis = (
|
||||
portfolio_stock_codes is not None and not stock_codes
|
||||
)
|
||||
|
||||
# 命令行参数 --single-notify 覆盖配置(#55)
|
||||
if getattr(args, 'single_notify', False):
|
||||
@@ -777,13 +862,20 @@ def run_full_analysis(
|
||||
)
|
||||
|
||||
# 1. 运行个股分析
|
||||
results = pipeline.run(
|
||||
stock_codes=stock_codes,
|
||||
dry_run=args.dry_run,
|
||||
send_notification=not args.no_notify,
|
||||
merge_notification=merge_notification,
|
||||
current_time=analysis_reference_time,
|
||||
)
|
||||
if skip_futu_stock_analysis:
|
||||
if portfolio_is_empty:
|
||||
logger.info("真实账户中无符合条件的 Futu 持仓,跳过个股分析。")
|
||||
else:
|
||||
logger.info("Futu 持仓经交易日过滤后无可分析股票,跳过个股分析。")
|
||||
results = []
|
||||
else:
|
||||
results = pipeline.run(
|
||||
stock_codes=stock_codes,
|
||||
dry_run=args.dry_run,
|
||||
send_notification=not args.no_notify,
|
||||
merge_notification=merge_notification,
|
||||
current_time=analysis_reference_time,
|
||||
)
|
||||
|
||||
if should_use_daily_market_context and not market_context_summary:
|
||||
(
|
||||
@@ -972,24 +1064,7 @@ def run_full_analysis(
|
||||
logger.error(f"飞书文档生成失败: {e}")
|
||||
|
||||
# === Auto backtest ===
|
||||
try:
|
||||
if getattr(config, 'backtest_enabled', False):
|
||||
from src.services.backtest_service import BacktestService
|
||||
|
||||
logger.info("开始自动回测...")
|
||||
service = BacktestService()
|
||||
stats = service.run_backtest(
|
||||
force=False,
|
||||
eval_window_days=getattr(config, 'backtest_eval_window_days', 10),
|
||||
min_age_days=getattr(config, 'backtest_min_age_days', 14),
|
||||
limit=200,
|
||||
)
|
||||
logger.info(
|
||||
f"自动回测完成: processed={stats.get('processed')} saved={stats.get('saved')} "
|
||||
f"completed={stats.get('completed')} insufficient={stats.get('insufficient')} errors={stats.get('errors')}"
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"自动回测失败(已忽略): {e}")
|
||||
_run_auto_backtest(config)
|
||||
|
||||
return True
|
||||
|
||||
@@ -1299,6 +1374,8 @@ def main() -> int:
|
||||
if (c or "").strip()
|
||||
]
|
||||
logger.info(f"使用命令行指定的股票列表: {stock_codes}")
|
||||
if getattr(args, "portfolio", None):
|
||||
logger.info("同时指定了 --portfolio;实际分析时 portfolio 将覆盖 --stocks")
|
||||
|
||||
# === 处理 --webui / --webui-only 参数,映射到 --serve / --serve-only ===
|
||||
if args.webui:
|
||||
@@ -1351,7 +1428,7 @@ def main() -> int:
|
||||
)
|
||||
else:
|
||||
os.environ.pop(RUNTIME_SCHEDULER_RUN_IMMEDIATELY_ENV, None)
|
||||
os.environ[RUNTIME_SCHEDULER_ARGS_ENV] = json.dumps({
|
||||
runtime_scheduler_args = {
|
||||
"no_notify": bool(getattr(args, "no_notify", False)),
|
||||
"no_market_review": bool(getattr(args, "no_market_review", False)),
|
||||
"dry_run": bool(getattr(args, "dry_run", False)),
|
||||
@@ -1359,7 +1436,10 @@ def main() -> int:
|
||||
"single_notify": bool(getattr(args, "single_notify", False)),
|
||||
"no_context_snapshot": bool(getattr(args, "no_context_snapshot", False)),
|
||||
"workers": getattr(args, "workers", None),
|
||||
})
|
||||
}
|
||||
if getattr(args, "portfolio", None):
|
||||
runtime_scheduler_args["portfolio"] = args.portfolio
|
||||
os.environ[RUNTIME_SCHEDULER_ARGS_ENV] = json.dumps(runtime_scheduler_args)
|
||||
if not prepare_webui_frontend_assets():
|
||||
logger.warning("前端静态资源未就绪,继续启动 FastAPI 服务(Web 页面可能不可用)")
|
||||
try:
|
||||
@@ -1511,7 +1591,15 @@ def main() -> int:
|
||||
|
||||
# 模式3: 正常单次运行
|
||||
if config.run_immediately:
|
||||
_run_analysis_with_runtime_scheduler_lock(config, args, stock_codes)
|
||||
try:
|
||||
_run_analysis_with_runtime_scheduler_lock(config, args, stock_codes)
|
||||
except FutuPortfolioError as exc:
|
||||
if not start_serve:
|
||||
raise
|
||||
logger.exception(
|
||||
"Futu 持仓导入失败,Web/API 服务继续运行: %s",
|
||||
exc,
|
||||
)
|
||||
else:
|
||||
logger.info("配置为不立即运行分析 (RUN_IMMEDIATELY=false)")
|
||||
|
||||
|
||||
@@ -22,6 +22,7 @@ longbridge==0.2.74; platform_system == "Linux" and python_version < "3.12"
|
||||
# Priority 5: Longbridge OpenAPI SDK with OAuth support.
|
||||
longbridge>=4.0.5,<5; platform_system != "Linux" or python_version >= "3.12"
|
||||
tickflow>=0.1.24 # TickFlow official SDK; 0.1.24 verified for klines/get batch quotes/universes
|
||||
futu-api==10.8.6808 # Futu OpenAPI SDK for read-only --portfolio futu holdings import
|
||||
# Built-in optional AlphaSift screening engine
|
||||
git+https://github.com/ZhuLinsen/alphasift.git@9f522747caafd3c0b1ddb7e14d5cf44c8580b6cf#egg=alphasift
|
||||
|
||||
|
||||
@@ -47,6 +47,9 @@ log "Checking python-multipart availability..."
|
||||
log "Checking AlphaSift adapter availability..."
|
||||
"${PYTHON_BIN}" -c "import alphasift.dsa_adapter"
|
||||
|
||||
log "Checking Futu SDK availability..."
|
||||
"${PYTHON_BIN}" -c "import futu"
|
||||
|
||||
log "Checking orjson availability..."
|
||||
"${PYTHON_BIN}" -c "import orjson"
|
||||
|
||||
@@ -116,6 +119,7 @@ done
|
||||
pushd "${ROOT_DIR}" >/dev/null
|
||||
cmd=("${PYTHON_BIN}" -m PyInstaller --name stock_analysis --onedir --noconfirm --noconsole --add-data "static:static" --add-data "strategies:strategies" --collect-data litellm --collect-data tiktoken --collect-data akshare)
|
||||
cmd+=("--collect-all" "alphasift")
|
||||
cmd+=("--collect-all" "futu")
|
||||
cmd+=("${hidden_import_args[@]}" "main.py")
|
||||
|
||||
echo "Running: ${cmd[*]}"
|
||||
@@ -124,7 +128,7 @@ popd >/dev/null
|
||||
|
||||
cp -R "${ROOT_DIR}/dist/stock_analysis" "${ROOT_DIR}/dist/backend/stock_analysis"
|
||||
|
||||
log "Verifying packaged AlphaSift importability..."
|
||||
log "Verifying packaged runtime imports..."
|
||||
packaged_root="${ROOT_DIR}/dist/backend/stock_analysis"
|
||||
|
||||
packaged_entry="${packaged_root}/stock_analysis"
|
||||
@@ -133,14 +137,14 @@ if [[ ! -x "${packaged_entry}" ]]; then
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# 先校验可执行文件可启动(不进入业务流程的参数),再检查冻结产物中是否携带 alphasift.
|
||||
if ! "${packaged_entry}" --help >/tmp/alphasift-packaged-help.log 2>&1; then
|
||||
# 先校验可执行文件可启动(不进入业务流程的参数),再检查冻结产物中的关键依赖。
|
||||
if ! "${packaged_entry}" --help >/tmp/dsa-packaged-help.log 2>&1; then
|
||||
echo "ERROR: packaged backend help startup check failed."
|
||||
cat /tmp/alphasift-packaged-help.log
|
||||
cat /tmp/dsa-packaged-help.log
|
||||
exit 1
|
||||
fi
|
||||
|
||||
for module in alphasift.dsa_adapter orjson; do
|
||||
for module in alphasift.dsa_adapter futu orjson; do
|
||||
if DSA_PACKAGED_IMPORT_PROBE="${module}" "${packaged_entry}" >/tmp/dsa-packaged-import.log 2>&1; then
|
||||
cat /tmp/dsa-packaged-import.log
|
||||
else
|
||||
|
||||
@@ -56,6 +56,11 @@ if (-not (Test-PythonCode -Python $pythonBin -Code "import alphasift.dsa_adapter
|
||||
throw 'alphasift.dsa_adapter is not importable after installing requirements.'
|
||||
}
|
||||
|
||||
Write-Host 'Checking Futu SDK availability...'
|
||||
if (-not (Test-PythonCode -Python $pythonBin -Code "import futu")) {
|
||||
throw 'futu is not importable after installing requirements.'
|
||||
}
|
||||
|
||||
Write-Host 'Checking orjson availability...'
|
||||
if (-not (Test-PythonCode -Python $pythonBin -Code "import orjson")) {
|
||||
throw 'orjson is not importable after installing requirements.'
|
||||
@@ -131,7 +136,8 @@ $pyInstallerArgs = @(
|
||||
'--collect-data', 'litellm',
|
||||
'--collect-data', 'tiktoken',
|
||||
'--collect-data', 'akshare',
|
||||
'--collect-all', 'alphasift'
|
||||
'--collect-all', 'alphasift',
|
||||
'--collect-all', 'futu'
|
||||
)
|
||||
$pyInstallerArgs += $hiddenImportArgs
|
||||
$pyInstallerArgs += 'main.py'
|
||||
@@ -155,7 +161,7 @@ if (-not (Test-Path $packagedEntry)) {
|
||||
}
|
||||
$previousProbe = $env:DSA_PACKAGED_IMPORT_PROBE
|
||||
try {
|
||||
foreach ($module in @('alphasift.dsa_adapter', 'orjson')) {
|
||||
foreach ($module in @('alphasift.dsa_adapter', 'futu', 'orjson')) {
|
||||
$env:DSA_PACKAGED_IMPORT_PROBE = $module
|
||||
$probeProcess = Start-Process -FilePath $packagedEntry -Wait -PassThru
|
||||
if ($probeProcess.ExitCode -ne 0) {
|
||||
|
||||
@@ -0,0 +1,484 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Read real stock holdings from a Futu OpenD instance."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import ipaddress
|
||||
import logging
|
||||
import math
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Iterable, List, Optional
|
||||
|
||||
from data_provider.us_index_mapping import is_us_stock_code
|
||||
from src.services.stock_code_utils import normalize_code
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class FutuPortfolioError(RuntimeError):
|
||||
"""Raised when a Futu portfolio cannot be resolved safely."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class _FutuAccount:
|
||||
"""Identify one usable real Futu securities account."""
|
||||
|
||||
acc_id: int
|
||||
security_firm: Any
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class _FutuApi:
|
||||
"""Hold the imported Futu SDK surface used by portfolio loading."""
|
||||
|
||||
OpenQuoteContext: Any
|
||||
OpenSecTradeContext: Any
|
||||
Market: Any
|
||||
RET_OK: Any
|
||||
SecurityFirm: Any
|
||||
SecurityType: Any
|
||||
TrdEnv: Any
|
||||
TrdMarket: Any
|
||||
|
||||
|
||||
_SUPPORTED_ACCOUNT_ROLES = frozenset({"NORMAL", "MASTER"})
|
||||
_SUPPORTED_ANALYSIS_MARKETS = frozenset({"US", "HK", "SH", "SZ"})
|
||||
_UNKNOWN_SECURITY_TYPES = frozenset({"", "N/A", "NONE", "UNKNOWN", "NAN"})
|
||||
_STATIC_INFO_BATCH_SIZE = 100
|
||||
|
||||
|
||||
def _load_futu_api() -> _FutuApi:
|
||||
"""Import the supported Futu SDK surface or raise an actionable error."""
|
||||
|
||||
try:
|
||||
from futu import (
|
||||
Market,
|
||||
OpenQuoteContext,
|
||||
OpenSecTradeContext,
|
||||
RET_OK,
|
||||
SecurityFirm,
|
||||
SecurityType,
|
||||
TrdEnv,
|
||||
TrdMarket,
|
||||
)
|
||||
except ImportError as exc:
|
||||
raise FutuPortfolioError(
|
||||
"未安装 Futu OpenAPI SDK;请先执行 "
|
||||
"`pip install \"futu-api==10.8.6808\"`。"
|
||||
) from exc
|
||||
except Exception as exc: # noqa: BLE001 - SDK import initializes its file logger
|
||||
raise FutuPortfolioError(f"加载 Futu OpenAPI SDK 失败: {exc}") from exc
|
||||
|
||||
return _FutuApi(
|
||||
OpenQuoteContext=OpenQuoteContext,
|
||||
OpenSecTradeContext=OpenSecTradeContext,
|
||||
Market=Market,
|
||||
RET_OK=RET_OK,
|
||||
SecurityFirm=SecurityFirm,
|
||||
SecurityType=SecurityType,
|
||||
TrdEnv=TrdEnv,
|
||||
TrdMarket=TrdMarket,
|
||||
)
|
||||
|
||||
|
||||
def _enum_text(value: Any) -> str:
|
||||
"""Normalize SDK enum-like values for stable comparisons."""
|
||||
|
||||
if value is None:
|
||||
return ""
|
||||
name = getattr(value, "name", None)
|
||||
return str(name if name is not None else value).strip().upper()
|
||||
|
||||
|
||||
def _iter_rows(data: Any, operation: str) -> Iterable[Any]:
|
||||
"""Iterate the pandas-style table returned by the pinned Futu SDK."""
|
||||
|
||||
iterrows = getattr(data, "iterrows", None)
|
||||
if not callable(iterrows):
|
||||
raise FutuPortfolioError(f"{operation}返回了非表格数据")
|
||||
return (row for _, row in iterrows())
|
||||
|
||||
|
||||
def _safe_close(context: Any) -> None:
|
||||
"""Close an SDK context without masking the primary operation result."""
|
||||
|
||||
if context is None:
|
||||
return
|
||||
try:
|
||||
context.close()
|
||||
except Exception: # pragma: no cover - closing is best effort
|
||||
logger.debug("关闭 Futu OpenD 连接失败", exc_info=True)
|
||||
|
||||
|
||||
def _connection_settings() -> tuple[str, int]:
|
||||
"""Return the validated IPv4 OpenD host and port from environment settings."""
|
||||
|
||||
host = (os.getenv("FUTU_OPEND_HOST") or "127.0.0.1").strip()
|
||||
raw_port = (os.getenv("FUTU_OPEND_PORT") or "11111").strip()
|
||||
try:
|
||||
port = int(raw_port)
|
||||
except ValueError as exc:
|
||||
raise FutuPortfolioError(f"FUTU_OPEND_PORT 不是有效端口: {raw_port!r}") from exc
|
||||
if not host or not 1 <= port <= 65535:
|
||||
raise FutuPortfolioError(f"Futu OpenD 地址无效: {host!r}:{port}")
|
||||
|
||||
address_text = host[1:-1] if host.startswith("[") and host.endswith("]") else host
|
||||
try:
|
||||
address = ipaddress.ip_address(address_text)
|
||||
except ValueError:
|
||||
address = None
|
||||
if address is not None and address.version != 4:
|
||||
raise FutuPortfolioError(
|
||||
"futu-api==10.8.6808 的网络层仅支持 IPv4;"
|
||||
f"FUTU_OPEND_HOST 当前为 {host!r},请改用 IPv4 地址或可解析到 IPv4 的主机名。"
|
||||
)
|
||||
return host, port
|
||||
|
||||
|
||||
def _configured_account_id() -> Optional[int]:
|
||||
"""Return the optional configured real account ID."""
|
||||
|
||||
value = (os.getenv("FUTU_ACC_ID") or "").strip()
|
||||
if not value:
|
||||
return None
|
||||
try:
|
||||
account_id = int(value)
|
||||
except ValueError as exc:
|
||||
raise FutuPortfolioError("FUTU_ACC_ID 必须是正整数账户 ID") from exc
|
||||
if account_id <= 0:
|
||||
raise FutuPortfolioError("FUTU_ACC_ID 必须是正整数账户 ID")
|
||||
return account_id
|
||||
|
||||
|
||||
def _configured_security_firm(api: _FutuApi) -> Any:
|
||||
"""Resolve one firm, defaulting to the SDK's official auto-detection mode."""
|
||||
|
||||
name = (os.getenv("FUTU_SECURITY_FIRM") or "NONE").strip().upper()
|
||||
firm = getattr(api.SecurityFirm, name, None)
|
||||
if firm is None:
|
||||
raise FutuPortfolioError(f"不支持的 FUTU_SECURITY_FIRM: {name}")
|
||||
return firm
|
||||
|
||||
|
||||
def _discover_real_accounts(api: _FutuApi, host: str, port: int) -> List[_FutuAccount]:
|
||||
"""Discover explicitly ACTIVE NORMAL or MASTER REAL accounts."""
|
||||
|
||||
accounts: List[_FutuAccount] = []
|
||||
seen_ids = set()
|
||||
requested_acc_id = _configured_account_id()
|
||||
security_firm = _configured_security_firm(api)
|
||||
context = None
|
||||
try:
|
||||
context = api.OpenSecTradeContext(
|
||||
host=host,
|
||||
port=port,
|
||||
filter_trdmarket=api.TrdMarket.NONE,
|
||||
security_firm=security_firm,
|
||||
)
|
||||
ret, data = context.get_acc_list()
|
||||
if ret != api.RET_OK:
|
||||
raise FutuPortfolioError(f"查询 Futu 真实账户失败: {data}")
|
||||
for row in _iter_rows(data, "Futu 账户查询"):
|
||||
if _enum_text(row.get("trd_env")) != "REAL":
|
||||
continue
|
||||
if _enum_text(row.get("acc_status")) != "ACTIVE":
|
||||
continue
|
||||
if _enum_text(row.get("acc_role")) not in _SUPPORTED_ACCOUNT_ROLES:
|
||||
continue
|
||||
raw_acc_id = row.get("acc_id")
|
||||
try:
|
||||
acc_id = int(raw_acc_id)
|
||||
exact_integer = isinstance(raw_acc_id, str) or bool(
|
||||
raw_acc_id == acc_id
|
||||
)
|
||||
except (TypeError, ValueError, OverflowError) as exc:
|
||||
raise FutuPortfolioError(
|
||||
"Futu 账户查询返回了无效账户 ID"
|
||||
) from exc
|
||||
if isinstance(raw_acc_id, bool) or not exact_integer or acc_id <= 0:
|
||||
raise FutuPortfolioError("Futu 账户查询返回了无效账户 ID")
|
||||
if acc_id in seen_ids:
|
||||
continue
|
||||
returned_firm_name = _enum_text(row.get("security_firm"))
|
||||
returned_firm = getattr(
|
||||
api.SecurityFirm,
|
||||
returned_firm_name,
|
||||
security_firm,
|
||||
)
|
||||
seen_ids.add(acc_id)
|
||||
accounts.append(_FutuAccount(acc_id=acc_id, security_firm=returned_firm))
|
||||
except FutuPortfolioError:
|
||||
raise
|
||||
except Exception as exc: # noqa: BLE001 - translate SDK/network failures
|
||||
raise FutuPortfolioError(f"查询 Futu 真实账户失败: {exc}") from exc
|
||||
finally:
|
||||
_safe_close(context)
|
||||
|
||||
if requested_acc_id is not None:
|
||||
accounts = [account for account in accounts if account.acc_id == requested_acc_id]
|
||||
if not accounts:
|
||||
raise FutuPortfolioError(
|
||||
"FUTU_ACC_ID 未匹配到可用的真实证券账户;请检查账户 ID、券商和 OpenD 登录状态。"
|
||||
)
|
||||
|
||||
if not accounts:
|
||||
raise FutuPortfolioError(
|
||||
"未找到状态为 ACTIVE 的 Futu REAL 普通或 MASTER 证券账户"
|
||||
)
|
||||
return accounts
|
||||
|
||||
|
||||
def _load_position_codes(
|
||||
api: _FutuApi,
|
||||
host: str,
|
||||
port: int,
|
||||
accounts: Iterable[_FutuAccount],
|
||||
) -> List[str]:
|
||||
"""Load deduplicated non-zero LONG position codes from selected accounts."""
|
||||
|
||||
codes: List[str] = []
|
||||
seen_codes = set()
|
||||
skipped_short_count = 0
|
||||
skipped_unknown_side_count = 0
|
||||
|
||||
for account in accounts:
|
||||
context = None
|
||||
try:
|
||||
context = api.OpenSecTradeContext(
|
||||
host=host,
|
||||
port=port,
|
||||
filter_trdmarket=api.TrdMarket.NONE,
|
||||
security_firm=account.security_firm,
|
||||
)
|
||||
ret, data = context.position_list_query(
|
||||
trd_env=api.TrdEnv.REAL,
|
||||
acc_id=account.acc_id,
|
||||
refresh_cache=True,
|
||||
)
|
||||
if ret != api.RET_OK:
|
||||
raise FutuPortfolioError(f"查询 Futu 真实持仓失败: {data}")
|
||||
for row in _iter_rows(data, "Futu 持仓查询"):
|
||||
position_side = _enum_text(row.get("position_side"))
|
||||
if position_side == "SHORT":
|
||||
skipped_short_count += 1
|
||||
continue
|
||||
if position_side != "LONG":
|
||||
skipped_unknown_side_count += 1
|
||||
continue
|
||||
raw_code = row.get("code")
|
||||
code = (
|
||||
raw_code.strip().upper()
|
||||
if isinstance(raw_code, str)
|
||||
else ""
|
||||
)
|
||||
raw_quantity = row.get("qty")
|
||||
try:
|
||||
if isinstance(raw_quantity, bool):
|
||||
raise TypeError("boolean quantity")
|
||||
quantity = float(raw_quantity)
|
||||
except (TypeError, ValueError) as exc:
|
||||
suffix = f": {code}" if code else ""
|
||||
raise FutuPortfolioError(f"Futu 持仓数量无效{suffix}") from exc
|
||||
if not math.isfinite(quantity):
|
||||
suffix = f": {code}" if code else ""
|
||||
raise FutuPortfolioError(f"Futu 持仓数量无效{suffix}")
|
||||
if quantity == 0:
|
||||
continue
|
||||
if not isinstance(raw_code, str):
|
||||
raise FutuPortfolioError("Futu 非零持仓返回了无效证券代码")
|
||||
if not code:
|
||||
raise FutuPortfolioError("Futu 非零持仓返回了空证券代码")
|
||||
market, separator, symbol = code.partition(".")
|
||||
if not separator or not market or not symbol:
|
||||
raise FutuPortfolioError(
|
||||
f"Futu 非零持仓返回了无效证券代码: {code}"
|
||||
)
|
||||
if code in seen_codes:
|
||||
continue
|
||||
seen_codes.add(code)
|
||||
codes.append(code)
|
||||
except FutuPortfolioError:
|
||||
raise
|
||||
except Exception as exc: # noqa: BLE001 - translate SDK/network errors for CLI callers
|
||||
raise FutuPortfolioError(f"查询 Futu 真实持仓失败: {exc}") from exc
|
||||
finally:
|
||||
_safe_close(context)
|
||||
|
||||
if skipped_short_count:
|
||||
logger.info("已跳过 %d 个 Futu SHORT 空头持仓", skipped_short_count)
|
||||
if skipped_unknown_side_count:
|
||||
logger.warning(
|
||||
"已跳过 %d 个持仓方向不是 LONG 的 Futu 持仓",
|
||||
skipped_unknown_side_count,
|
||||
)
|
||||
return codes
|
||||
|
||||
|
||||
def _market_prefix(code: str) -> str:
|
||||
"""Extract the Futu market prefix from a qualified security code."""
|
||||
|
||||
return code.split(".", 1)[0] if "." in code else ""
|
||||
|
||||
|
||||
def _is_cn_b_share_code(code: str) -> bool:
|
||||
"""Return whether a qualified Futu code is a Shanghai/Shenzhen B-share."""
|
||||
|
||||
prefix, separator, symbol = code.partition(".")
|
||||
if not separator or not (symbol.isdigit() and len(symbol) == 6):
|
||||
return False
|
||||
return (prefix == "SH" and symbol.startswith("900")) or (
|
||||
prefix == "SZ" and symbol.startswith("200")
|
||||
)
|
||||
|
||||
|
||||
def _to_analysis_code(futu_code: str) -> Optional[str]:
|
||||
"""Convert a supported Futu code into the analysis pipeline format."""
|
||||
|
||||
prefix, separator, symbol = futu_code.partition(".")
|
||||
if not separator or not symbol:
|
||||
return None
|
||||
prefix = prefix.upper()
|
||||
symbol = symbol.upper()
|
||||
if prefix == "US":
|
||||
normalized = normalize_code(symbol)
|
||||
if normalized == symbol and is_us_stock_code(normalized):
|
||||
return normalized
|
||||
return None
|
||||
if prefix == "HK":
|
||||
normalized = normalize_code(f"HK.{symbol}")
|
||||
return f"HK{normalized}" if normalized is not None else None
|
||||
if prefix in {"SH", "SZ"}:
|
||||
normalized = normalize_code(f"{prefix}.{symbol}")
|
||||
return normalized if normalized == symbol else None
|
||||
return None
|
||||
|
||||
|
||||
def _filter_stock_codes(
|
||||
api: _FutuApi,
|
||||
host: str,
|
||||
port: int,
|
||||
position_codes: List[str],
|
||||
) -> List[str]:
|
||||
"""Keep A/HK/US stocks and report unsupported Futu market codes."""
|
||||
|
||||
if not position_codes:
|
||||
return []
|
||||
|
||||
grouped: dict[str, List[str]] = {}
|
||||
unsupported_codes: List[str] = []
|
||||
for code in position_codes:
|
||||
prefix = _market_prefix(code)
|
||||
if prefix not in _SUPPORTED_ANALYSIS_MARKETS or _is_cn_b_share_code(code):
|
||||
unsupported_codes.append(code)
|
||||
continue
|
||||
grouped.setdefault(prefix, []).append(code)
|
||||
|
||||
if not grouped:
|
||||
logger.warning(
|
||||
"已跳过 %d 个当前分析流程不支持的 Futu 持仓: %s",
|
||||
len(unsupported_codes),
|
||||
", ".join(unsupported_codes),
|
||||
)
|
||||
return []
|
||||
|
||||
stock_codes = set()
|
||||
classified_codes = set()
|
||||
context = None
|
||||
try:
|
||||
context = api.OpenQuoteContext(host=host, port=port)
|
||||
for prefix, codes in grouped.items():
|
||||
market = getattr(api.Market, prefix, None)
|
||||
if market is None:
|
||||
unsupported_codes.extend(codes)
|
||||
continue
|
||||
for start in range(0, len(codes), _STATIC_INFO_BATCH_SIZE):
|
||||
batch = codes[start : start + _STATIC_INFO_BATCH_SIZE]
|
||||
ret, data = context.get_stock_basicinfo(
|
||||
market,
|
||||
stock_type=api.SecurityType.STOCK,
|
||||
code_list=batch,
|
||||
)
|
||||
if ret != api.RET_OK:
|
||||
raise FutuPortfolioError(
|
||||
f"查询 Futu 持仓证券类型失败({prefix}): {data}"
|
||||
)
|
||||
for row in _iter_rows(data, "Futu 证券类型查询"):
|
||||
code = str(row.get("code", "") or "").strip().upper()
|
||||
if not code:
|
||||
continue
|
||||
stock_type = _enum_text(row.get("stock_type"))
|
||||
if stock_type in _UNKNOWN_SECURITY_TYPES:
|
||||
continue
|
||||
classified_codes.add(code)
|
||||
if stock_type == "STOCK":
|
||||
stock_codes.add(code)
|
||||
except FutuPortfolioError:
|
||||
raise
|
||||
except Exception as exc: # noqa: BLE001 - translate SDK/network errors for CLI callers
|
||||
raise FutuPortfolioError(f"查询 Futu 持仓证券类型失败: {exc}") from exc
|
||||
finally:
|
||||
_safe_close(context)
|
||||
|
||||
missing_codes = [
|
||||
code
|
||||
for codes in grouped.values()
|
||||
for code in codes
|
||||
if code not in classified_codes
|
||||
]
|
||||
if unsupported_codes:
|
||||
logger.warning(
|
||||
"已跳过 %d 个当前分析流程不支持的 Futu 持仓: %s",
|
||||
len(unsupported_codes),
|
||||
", ".join(unsupported_codes),
|
||||
)
|
||||
if missing_codes:
|
||||
raise FutuPortfolioError(
|
||||
"无法确认证券类型的 Futu 持仓: " + ", ".join(missing_codes)
|
||||
)
|
||||
|
||||
result: List[str] = []
|
||||
conversion_failures: List[str] = []
|
||||
for futu_code in position_codes:
|
||||
if futu_code not in stock_codes:
|
||||
continue
|
||||
analysis_code = _to_analysis_code(futu_code)
|
||||
if not analysis_code:
|
||||
conversion_failures.append(futu_code)
|
||||
continue
|
||||
if analysis_code not in result:
|
||||
result.append(analysis_code)
|
||||
if conversion_failures:
|
||||
raise FutuPortfolioError(
|
||||
"无法转换已确认的 Futu 正股代码到当前分析格式: "
|
||||
+ ", ".join(conversion_failures)
|
||||
)
|
||||
return result
|
||||
|
||||
|
||||
def load_futu_stock_codes() -> List[str]:
|
||||
"""Return deduplicated analysis codes from all selected REAL Futu accounts.
|
||||
|
||||
Only explicitly ACTIVE REAL accounts and Futu ``SecurityType.STOCK`` LONG
|
||||
positions with non-zero quantity are kept. ``FUTU_ACC_ID`` can select one
|
||||
account; otherwise NORMAL and MASTER accounts are merged. ``MASTER`` is an
|
||||
account role, while read-only describes this integration's query-only API
|
||||
calls. Firm discovery uses the SDK's ``SecurityFirm.NONE`` auto-detection
|
||||
unless ``FUTU_SECURITY_FIRM`` is explicitly set. Position data is always
|
||||
refreshed. Symbol conversion is limited to A/HK/US stocks; holdings from
|
||||
other Futu markets are logged with their codes and skipped.
|
||||
"""
|
||||
api = _load_futu_api()
|
||||
host, port = _connection_settings()
|
||||
accounts = _discover_real_accounts(api, host, port)
|
||||
position_codes = _load_position_codes(api, host, port, accounts)
|
||||
stock_codes = _filter_stock_codes(api, host, port, position_codes)
|
||||
logger.info(
|
||||
"已从 Futu 真实账户加载 %d 只正股(账户数: %d,原始非零多头持仓数: %d): %s",
|
||||
len(stock_codes),
|
||||
len(accounts),
|
||||
len(position_codes),
|
||||
", ".join(stock_codes),
|
||||
)
|
||||
return stock_codes
|
||||
@@ -29,6 +29,7 @@ SCHEDULE_ARGS_OVERRIDE_KEYS = {
|
||||
"single_notify",
|
||||
"no_context_snapshot",
|
||||
"workers",
|
||||
"portfolio",
|
||||
}
|
||||
|
||||
|
||||
@@ -155,6 +156,7 @@ class RuntimeSchedulerService:
|
||||
"serve": False,
|
||||
"serve_only": True,
|
||||
"stocks": None,
|
||||
"portfolio": None,
|
||||
"workers": None,
|
||||
}
|
||||
defaults.update(self._schedule_args_overrides)
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Distribution contracts for the globally bundled Futu SDK."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
import yaml
|
||||
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parents[1]
|
||||
|
||||
|
||||
def _read(relative_path: str) -> str:
|
||||
return (REPO_ROOT / relative_path).read_text(encoding="utf-8")
|
||||
|
||||
|
||||
def _workflow(relative_path: str) -> dict:
|
||||
return yaml.load(_read(relative_path), Loader=yaml.BaseLoader)
|
||||
|
||||
|
||||
def _job_run_text(job: dict) -> str:
|
||||
return "\n".join(
|
||||
str(step.get("run", ""))
|
||||
for step in job.get("steps", [])
|
||||
if isinstance(step, dict)
|
||||
)
|
||||
|
||||
|
||||
def test_futu_sdk_is_pinned_and_verified_across_linux_distributions() -> None:
|
||||
requirements = _read("requirements.txt")
|
||||
dockerfile = _read("docker/Dockerfile")
|
||||
ci = _workflow(".github/workflows/ci.yml")
|
||||
daily = _workflow(".github/workflows/00-daily-analysis.yml")
|
||||
docker_publish = _workflow(".github/workflows/docker-publish.yml")
|
||||
manual_publish = _workflow(".github/workflows/ghcr-dockerhub.yml")
|
||||
|
||||
assert requirements.count("futu-api==10.8.6808") == 1
|
||||
assert 'python -c "import alphasift.dsa_adapter; import futu"' in dockerfile
|
||||
assert "import futu" in _job_run_text(ci["jobs"]["backend-gate"])
|
||||
assert "import futu" in _job_run_text(ci["jobs"]["docker-build"])
|
||||
assert "import futu" in _job_run_text(daily["jobs"]["analyze"])
|
||||
assert "import futu" in _job_run_text(docker_publish["jobs"]["build-and-push"])
|
||||
assert "import futu" in _job_run_text(manual_publish["jobs"]["build-and-push"])
|
||||
|
||||
|
||||
def test_futu_sdk_is_collected_and_probed_in_desktop_backends() -> None:
|
||||
ci = _workflow(".github/workflows/ci.yml")
|
||||
changes_job = ci["jobs"]["changes"]
|
||||
filter_step = next(
|
||||
step for step in changes_job["steps"] if step.get("id") == "filter"
|
||||
)
|
||||
filters = str(filter_step["with"]["filters"])
|
||||
jobs = ci["jobs"]
|
||||
macos_script = _read("scripts/build-backend-macos.sh")
|
||||
windows_script = _read("scripts/build-backend.ps1")
|
||||
|
||||
assert "futu_packaging:" in filters
|
||||
assert "requirements.txt" in filters
|
||||
assert "src/brokers/futu/**" in filters
|
||||
assert "scripts/build-backend-macos.sh" in filters
|
||||
assert "scripts/build-backend.ps1" in filters
|
||||
|
||||
for job_name in (
|
||||
"desktop-futu-package-windows",
|
||||
"desktop-futu-package-macos",
|
||||
):
|
||||
job = jobs[job_name]
|
||||
assert job["needs"] == ["changes", "ai-governance"]
|
||||
assert "needs.changes.outputs.futu_packaging == 'true'" in job["if"]
|
||||
assert "build-backend" in _job_run_text(job)
|
||||
|
||||
assert '"${PYTHON_BIN}" -c "import futu"' in macos_script
|
||||
assert 'cmd+=("--collect-all" "futu")' in macos_script
|
||||
assert "for module in alphasift.dsa_adapter futu orjson" in macos_script
|
||||
|
||||
assert 'Test-PythonCode -Python $pythonBin -Code "import futu"' in windows_script
|
||||
assert "'--collect-all', 'futu'" in windows_script
|
||||
assert "@('alphasift.dsa_adapter', 'futu', 'orjson')" in windows_script
|
||||
@@ -0,0 +1,988 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from types import SimpleNamespace
|
||||
import unittest
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pandas as pd
|
||||
|
||||
from src.brokers.futu import portfolio as service
|
||||
|
||||
|
||||
class _TradeContext:
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
filter_trdmarket,
|
||||
host,
|
||||
port,
|
||||
security_firm,
|
||||
accounts=None,
|
||||
positions_by_account=None,
|
||||
) -> None:
|
||||
self.closed = False
|
||||
self.position_queries = []
|
||||
self.accounts = accounts
|
||||
self.positions_by_account = positions_by_account
|
||||
self.open_arguments = {
|
||||
"filter_trdmarket": filter_trdmarket,
|
||||
"host": host,
|
||||
"port": port,
|
||||
"security_firm": security_firm,
|
||||
}
|
||||
|
||||
def get_acc_list(self):
|
||||
if self.accounts is not None:
|
||||
return 0, pd.DataFrame(self.accounts)
|
||||
return 0, pd.DataFrame([
|
||||
{
|
||||
"acc_id": 1001,
|
||||
"trd_env": "REAL",
|
||||
"acc_role": "NORMAL",
|
||||
"acc_status": "ACTIVE",
|
||||
"security_firm": "FUTUSECURITIES",
|
||||
},
|
||||
{
|
||||
"acc_id": 2002,
|
||||
"trd_env": "SIMULATE",
|
||||
"acc_role": "NORMAL",
|
||||
"acc_status": "ACTIVE",
|
||||
"security_firm": "FUTUSECURITIES",
|
||||
},
|
||||
])
|
||||
|
||||
def position_list_query(self, **kwargs):
|
||||
self.position_queries.append(kwargs)
|
||||
if self.positions_by_account is not None:
|
||||
return 0, pd.DataFrame(
|
||||
self.positions_by_account.get(kwargs["acc_id"], [])
|
||||
)
|
||||
return 0, pd.DataFrame([
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "US.DRAM", "qty": 3, "position_side": "LONG"},
|
||||
{
|
||||
"code": "US.AAPL261218C200000",
|
||||
"qty": -1,
|
||||
"position_side": "LONG",
|
||||
},
|
||||
{"code": "HK.00700", "qty": 20, "position_side": "LONG"},
|
||||
{"code": "SH.600519", "qty": 0, "position_side": "LONG"},
|
||||
{"code": "SZ.000001", "qty": 8, "position_side": "LONG"},
|
||||
])
|
||||
|
||||
def close(self) -> None:
|
||||
self.closed = True
|
||||
|
||||
|
||||
class _QuoteContext:
|
||||
def __init__(self, *, host, port, stock_types=None) -> None:
|
||||
self.closed = False
|
||||
self.open_arguments = {"host": host, "port": port}
|
||||
self.stock_types = {
|
||||
"US.AAPL": "STOCK",
|
||||
"US.DRAM": "ETF",
|
||||
"US.AAPL261218C200000": "DRVT",
|
||||
"HK.00700": "STOCK",
|
||||
"SZ.000001": "STOCK",
|
||||
"JP.7203": "STOCK",
|
||||
"JP.130A": "STOCK",
|
||||
} if stock_types is None else stock_types
|
||||
|
||||
def get_stock_basicinfo(self, market, *, stock_type, code_list):
|
||||
return 0, pd.DataFrame([
|
||||
{"code": code, "stock_type": self.stock_types[code]}
|
||||
for code in code_list
|
||||
if code in self.stock_types
|
||||
])
|
||||
|
||||
def close(self) -> None:
|
||||
self.closed = True
|
||||
|
||||
|
||||
def _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
*,
|
||||
accounts=None,
|
||||
positions_by_account=None,
|
||||
stock_types=None,
|
||||
):
|
||||
def open_trade_context(*, filter_trdmarket, host, port, security_firm):
|
||||
context = _TradeContext(
|
||||
filter_trdmarket=filter_trdmarket,
|
||||
host=host,
|
||||
port=port,
|
||||
security_firm=security_firm,
|
||||
accounts=accounts,
|
||||
positions_by_account=positions_by_account,
|
||||
)
|
||||
trade_contexts.append(context)
|
||||
return context
|
||||
|
||||
def open_quote_context(*, host, port):
|
||||
context = _QuoteContext(host=host, port=port, stock_types=stock_types)
|
||||
quote_contexts.append(context)
|
||||
return context
|
||||
|
||||
return service._FutuApi(
|
||||
OpenQuoteContext=open_quote_context,
|
||||
OpenSecTradeContext=open_trade_context,
|
||||
Market=SimpleNamespace(US="US", HK="HK", SH="SH", SZ="SZ", JP="JP"),
|
||||
RET_OK=0,
|
||||
SecurityFirm=SimpleNamespace(
|
||||
NONE="N/A",
|
||||
FUTUSECURITIES="FUTUSECURITIES",
|
||||
FUTUSG="FUTUSG",
|
||||
),
|
||||
SecurityType=SimpleNamespace(STOCK="STOCK"),
|
||||
TrdEnv=SimpleNamespace(REAL="REAL"),
|
||||
TrdMarket=SimpleNamespace(NONE="NONE"),
|
||||
)
|
||||
|
||||
|
||||
def _account(
|
||||
acc_id,
|
||||
acc_role,
|
||||
*,
|
||||
acc_status="ACTIVE",
|
||||
security_firm="FUTUSECURITIES",
|
||||
):
|
||||
return {
|
||||
"acc_id": acc_id,
|
||||
"trd_env": "REAL",
|
||||
"acc_role": acc_role,
|
||||
"acc_status": acc_status,
|
||||
"security_firm": security_firm,
|
||||
}
|
||||
|
||||
|
||||
def _load_codes_for_accounts(accounts, positions_by_account):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=accounts,
|
||||
positions_by_account=positions_by_account,
|
||||
)
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(service, "_load_futu_api", return_value=api):
|
||||
result = service.load_futu_stock_codes()
|
||||
return result, trade_contexts
|
||||
|
||||
|
||||
class FutuPortfolioServiceTest(unittest.TestCase):
|
||||
def test_missing_sdk_uses_actionable_install_error(self):
|
||||
with patch(
|
||||
"builtins.__import__",
|
||||
side_effect=ImportError("No module named 'futu'"),
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"未安装 Futu OpenAPI SDK",
|
||||
) as raised:
|
||||
service._load_futu_api()
|
||||
|
||||
self.assertIn('pip install "futu-api==10.8.6808"', str(raised.exception))
|
||||
|
||||
def test_sdk_initialization_failure_uses_portfolio_error_boundary(self):
|
||||
with patch(
|
||||
"builtins.__import__",
|
||||
side_effect=PermissionError("log directory denied"),
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"加载 Futu OpenAPI SDK 失败: log directory denied",
|
||||
):
|
||||
service._load_futu_api()
|
||||
|
||||
def test_load_futu_stock_codes_keeps_only_supported_a_hk_us_stocks(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(trade_contexts, quote_contexts)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(service, "_load_futu_api", return_value=api):
|
||||
result = service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(result, ["AAPL", "HK00700", "000001"])
|
||||
position_contexts = [ctx for ctx in trade_contexts if ctx.position_queries]
|
||||
self.assertEqual(len(position_contexts), 1)
|
||||
self.assertEqual(
|
||||
position_contexts[0].position_queries,
|
||||
[{"trd_env": "REAL", "acc_id": 1001, "refresh_cache": True}],
|
||||
)
|
||||
self.assertTrue(all(ctx.closed for ctx in trade_contexts))
|
||||
self.assertTrue(quote_contexts and all(ctx.closed for ctx in quote_contexts))
|
||||
|
||||
def test_load_futu_stock_codes_reports_unsupported_jp_holdings(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "JP.7203", "qty": 5, "position_side": "LONG"},
|
||||
{"code": "JP.130A", "qty": 3, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertLogs(service.logger, level="WARNING") as captured:
|
||||
result = service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(result, [])
|
||||
warning_text = "\n".join(captured.output)
|
||||
self.assertIn("JP.7203", warning_text)
|
||||
self.assertIn("JP.130A", warning_text)
|
||||
|
||||
def test_load_futu_stock_codes_keeps_supported_holdings_when_jp_is_present(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "JP.7203", "qty": 5, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertLogs(service.logger, level="WARNING") as captured:
|
||||
result = service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
self.assertIn("JP.7203", "\n".join(captured.output))
|
||||
|
||||
def test_load_futu_stock_codes_rejects_stock_code_outside_analysis_contract(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "HK.BAD", "qty": 3, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
stock_types={"HK.BAD": "STOCK"},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"无法转换.*HK.BAD",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
def test_load_futu_stock_codes_reports_unsupported_b_shares(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "SH.900901", "qty": 5, "position_side": "LONG"},
|
||||
{"code": "SZ.200012", "qty": 8, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
stock_types={
|
||||
"US.AAPL": "STOCK",
|
||||
"SH.900901": "STOCK",
|
||||
"SZ.200012": "STOCK",
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertLogs(
|
||||
service.logger,
|
||||
level="WARNING",
|
||||
) as captured:
|
||||
result = service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
warning_text = "\n".join(captured.output)
|
||||
self.assertIn("SH.900901", warning_text)
|
||||
self.assertIn("SZ.200012", warning_text)
|
||||
|
||||
def test_load_futu_stock_codes_rejects_partial_static_info_response(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "US.MSFT", "qty": 4, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
stock_types={"US.AAPL": "STOCK"},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"无法确认证券类型.*US.MSFT",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
def test_load_futu_stock_codes_rejects_unknown_static_security_type(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "US.MSFT", "qty": 4, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
stock_types={"US.AAPL": "STOCK", "US.MSFT": "N/A"},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"无法确认证券类型.*US.MSFT",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
def test_load_futu_stock_codes_rejects_invalid_eligible_account_id(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[
|
||||
_account("invalid", "NORMAL"),
|
||||
_account(1001, "NORMAL"),
|
||||
],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"账户查询返回了无效账户 ID",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertFalse(any(ctx.position_queries for ctx in trade_contexts))
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_load_futu_stock_codes_rejects_nonpositive_or_fractional_account_id(self):
|
||||
for invalid_acc_id in (0, -1, 1001.5, True):
|
||||
with self.subTest(acc_id=invalid_acc_id):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[
|
||||
_account(invalid_acc_id, "NORMAL"),
|
||||
_account(1001, "NORMAL"),
|
||||
],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{
|
||||
"code": "US.AAPL",
|
||||
"qty": 10,
|
||||
"position_side": "LONG",
|
||||
}
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"账户查询返回了无效账户 ID",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertFalse(any(ctx.position_queries for ctx in trade_contexts))
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_load_futu_stock_codes_rejects_invalid_position_quantity(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "US.MSFT", "qty": "bad", "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
stock_types={
|
||||
"US.AAPL": "STOCK",
|
||||
"US.MSFT": "STOCK",
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"持仓数量无效.*US.MSFT",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_load_futu_stock_codes_rejects_blank_nonzero_position_code(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[_account(1001, "NORMAL")],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "", "qty": 5, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"非零持仓返回了空证券代码",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_load_futu_stock_codes_rejects_missing_nonzero_long_code(self):
|
||||
with self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"非零持仓返回了无效证券代码",
|
||||
):
|
||||
_load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL")],
|
||||
{
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"qty": 5, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
def test_load_futu_stock_codes_rejects_unqualified_nonzero_long_code(self):
|
||||
with self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"非零持仓返回了无效证券代码",
|
||||
):
|
||||
_load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL")],
|
||||
{
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "AAPL", "qty": 5, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
def test_load_futu_stock_codes_rejects_non_string_nonzero_long_codes(self):
|
||||
for invalid_code in (True, 123, b"US.AAPL"):
|
||||
with self.subTest(code=invalid_code), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"非零持仓返回了无效证券代码",
|
||||
):
|
||||
_load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL")],
|
||||
{
|
||||
1001: [
|
||||
{
|
||||
"code": "US.AAPL",
|
||||
"qty": 10,
|
||||
"position_side": "LONG",
|
||||
},
|
||||
{
|
||||
"code": invalid_code,
|
||||
"qty": 5,
|
||||
"position_side": "LONG",
|
||||
},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
def test_default_firm_uses_one_official_none_discovery_context(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(trade_contexts, quote_contexts)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{},
|
||||
clear=True,
|
||||
), patch.object(service, "_load_futu_api", return_value=api):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(
|
||||
[context.open_arguments for context in trade_contexts],
|
||||
[
|
||||
{
|
||||
"filter_trdmarket": "NONE",
|
||||
"host": "127.0.0.1",
|
||||
"port": 11111,
|
||||
"security_firm": "N/A",
|
||||
},
|
||||
{
|
||||
"filter_trdmarket": "NONE",
|
||||
"host": "127.0.0.1",
|
||||
"port": 11111,
|
||||
"security_firm": "FUTUSECURITIES",
|
||||
},
|
||||
],
|
||||
)
|
||||
self.assertEqual(
|
||||
[context.open_arguments for context in quote_contexts],
|
||||
[{"host": "127.0.0.1", "port": 11111}],
|
||||
)
|
||||
|
||||
def test_configured_security_firm_replaces_auto_detection(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[
|
||||
_account(
|
||||
1001,
|
||||
"NORMAL",
|
||||
security_firm="FUTUSG",
|
||||
)
|
||||
],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"}
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{"FUTU_SECURITY_FIRM": "FUTUSG"},
|
||||
clear=True,
|
||||
), patch.object(service, "_load_futu_api", return_value=api):
|
||||
result = service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
self.assertEqual(
|
||||
trade_contexts[0].open_arguments["security_firm"],
|
||||
"FUTUSG",
|
||||
)
|
||||
|
||||
def test_unknown_security_firm_fails_before_opening_context(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(trade_contexts, quote_contexts)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{"FUTU_SECURITY_FIRM": "UNKNOWN"},
|
||||
clear=True,
|
||||
), patch.object(service, "_load_futu_api", return_value=api), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"不支持的 FUTU_SECURITY_FIRM: UNKNOWN",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(trade_contexts, [])
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_configured_account_id_must_be_a_positive_integer(self):
|
||||
for configured_acc_id in ("0", "-1", "1.5"):
|
||||
with self.subTest(acc_id=configured_acc_id):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(trade_contexts, quote_contexts)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{"FUTU_ACC_ID": configured_acc_id},
|
||||
clear=True,
|
||||
), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"FUTU_ACC_ID 必须是正整数账户 ID",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(trade_contexts, [])
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_account_discovery_failure_is_not_retried_or_partially_ignored(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
base_api = _fake_api(trade_contexts, quote_contexts)
|
||||
context = SimpleNamespace(
|
||||
get_acc_list=MagicMock(return_value=(1, "broker unavailable")),
|
||||
close=MagicMock(),
|
||||
)
|
||||
open_calls = []
|
||||
|
||||
def open_trade_context(**kwargs):
|
||||
open_calls.append(kwargs)
|
||||
return context
|
||||
|
||||
api = SimpleNamespace(**base_api.__dict__)
|
||||
api.OpenSecTradeContext = open_trade_context
|
||||
|
||||
with patch.dict("os.environ", {}, clear=True), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"查询 Futu 真实账户失败: broker unavailable",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertEqual(len(open_calls), 1)
|
||||
self.assertEqual(open_calls[0]["security_firm"], "N/A")
|
||||
context.close.assert_called_once_with()
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_real_accounts_require_explicit_active_status(self):
|
||||
cases = {
|
||||
"missing": None,
|
||||
"n/a": "N/A",
|
||||
"unknown": "UNKNOWN",
|
||||
"disabled": "DISABLED",
|
||||
}
|
||||
for label, status in cases.items():
|
||||
with self.subTest(status=label):
|
||||
account = _account(1001, "NORMAL", acc_status=status)
|
||||
if label == "missing":
|
||||
account.pop("acc_status")
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(
|
||||
trade_contexts,
|
||||
quote_contexts,
|
||||
accounts=[account],
|
||||
positions_by_account={
|
||||
1001: [
|
||||
{
|
||||
"code": "US.AAPL",
|
||||
"qty": 10,
|
||||
"position_side": "LONG",
|
||||
}
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {}, clear=True), patch.object(
|
||||
service,
|
||||
"_load_futu_api",
|
||||
return_value=api,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"未找到状态为 ACTIVE",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertFalse(any(ctx.position_queries for ctx in trade_contexts))
|
||||
self.assertEqual(quote_contexts, [])
|
||||
|
||||
def test_load_futu_stock_codes_keeps_active_master_account(self):
|
||||
result, trade_contexts = _load_codes_for_accounts(
|
||||
[_account(3003, "MASTER")],
|
||||
{
|
||||
3003: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"}
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
position_contexts = [ctx for ctx in trade_contexts if ctx.position_queries]
|
||||
self.assertEqual(len(position_contexts), 1)
|
||||
self.assertEqual(position_contexts[0].position_queries[0]["acc_id"], 3003)
|
||||
|
||||
def test_load_futu_stock_codes_merges_master_and_normal_accounts(self):
|
||||
result, trade_contexts = _load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL"), _account(3003, "MASTER")],
|
||||
{
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"}
|
||||
],
|
||||
3003: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
{"code": "HK.00700", "qty": 20, "position_side": "LONG"},
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
self.assertEqual(result, ["AAPL", "HK00700"])
|
||||
queried_account_ids = [
|
||||
context.position_queries[0]["acc_id"]
|
||||
for context in trade_contexts
|
||||
if context.position_queries
|
||||
]
|
||||
self.assertEqual(queried_account_ids, [1001, 3003])
|
||||
|
||||
def test_load_futu_stock_codes_skips_short_positions_before_deduplication(self):
|
||||
result, trade_contexts = _load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL"), _account(3003, "MASTER")],
|
||||
{
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "SHORT"},
|
||||
{"code": "HK.00700", "qty": 20, "position_side": "SHORT"},
|
||||
],
|
||||
3003: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"}
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
queried_account_ids = [
|
||||
context.position_queries[0]["acc_id"]
|
||||
for context in trade_contexts
|
||||
if context.position_queries
|
||||
]
|
||||
self.assertEqual(queried_account_ids, [1001, 3003])
|
||||
|
||||
def test_load_futu_stock_codes_skips_non_long_before_validating_fields(self):
|
||||
result, _ = _load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL")],
|
||||
{
|
||||
1001: [
|
||||
{"qty": "bad", "position_side": "SHORT"},
|
||||
{"qty": None, "position_side": "N/A"},
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"},
|
||||
]
|
||||
},
|
||||
)
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
|
||||
def test_load_futu_stock_codes_skips_unknown_position_sides(self):
|
||||
result, _ = _load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL")],
|
||||
{
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "N/A"},
|
||||
{"code": "HK.00700", "qty": 20},
|
||||
{"code": "JP.7203", "qty": 5, "position_side": "NONE"},
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
self.assertEqual(result, [])
|
||||
|
||||
def test_load_futu_stock_codes_rejects_non_finite_or_missing_quantities(self):
|
||||
for quantity in (float("nan"), float("inf"), float("-inf"), None, True):
|
||||
with self.subTest(quantity=quantity), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"持仓数量无效.*US.AAPL",
|
||||
):
|
||||
_load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL")],
|
||||
{
|
||||
1001: [
|
||||
{
|
||||
"code": "US.AAPL",
|
||||
"qty": quantity,
|
||||
"position_side": "LONG",
|
||||
}
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
def test_load_futu_stock_codes_skips_malaysian_ipo_accounts(self):
|
||||
result, trade_contexts = _load_codes_for_accounts(
|
||||
[_account(1001, "NORMAL"), _account(4004, "IPO")],
|
||||
{
|
||||
1001: [
|
||||
{"code": "US.AAPL", "qty": 10, "position_side": "LONG"}
|
||||
],
|
||||
4004: [
|
||||
{"code": "HK.00700", "qty": 20, "position_side": "LONG"}
|
||||
],
|
||||
},
|
||||
)
|
||||
|
||||
self.assertEqual(result, ["AAPL"])
|
||||
queried_account_ids = [
|
||||
context.position_queries[0]["acc_id"]
|
||||
for context in trade_contexts
|
||||
if context.position_queries
|
||||
]
|
||||
self.assertEqual(queried_account_ids, [1001])
|
||||
|
||||
def test_invalid_futu_account_id_fails_before_position_query(self):
|
||||
trade_contexts = []
|
||||
quote_contexts = []
|
||||
api = _fake_api(trade_contexts, quote_contexts)
|
||||
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{
|
||||
"FUTU_SECURITY_FIRM": "FUTUSECURITIES",
|
||||
"FUTU_ACC_ID": "9999",
|
||||
},
|
||||
clear=True,
|
||||
), patch.object(service, "_load_futu_api", return_value=api), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"FUTU_ACC_ID 未匹配",
|
||||
):
|
||||
service.load_futu_stock_codes()
|
||||
|
||||
self.assertFalse(any(ctx.position_queries for ctx in trade_contexts))
|
||||
|
||||
def test_to_analysis_code(self):
|
||||
cases = [
|
||||
("US.MSFT", "MSFT"),
|
||||
("US.BRK.B", "BRK.B"),
|
||||
("HK.01810", "HK01810"),
|
||||
("HK.700", "HK00700"),
|
||||
("SZ.000001", "000001"),
|
||||
("SH.600519", "600519"),
|
||||
("HK.123456", None),
|
||||
("HK.BAD", None),
|
||||
("SH.1", None),
|
||||
("SZ.1234567", None),
|
||||
("US.AAICPRC", None),
|
||||
("US.SPX", None),
|
||||
("JP.9984", None),
|
||||
("SG.D05", None),
|
||||
]
|
||||
for futu_code, expected in cases:
|
||||
with self.subTest(futu_code=futu_code):
|
||||
self.assertEqual(service._to_analysis_code(futu_code), expected)
|
||||
|
||||
def test_connection_settings_accepts_ipv4_and_hostnames(self):
|
||||
cases = {
|
||||
"default": (None, "127.0.0.1"),
|
||||
"explicit_ipv4": ("127.0.0.1", "127.0.0.1"),
|
||||
"remote_ipv4": ("192.168.1.10", "192.168.1.10"),
|
||||
"hostname": ("localhost", "localhost"),
|
||||
"remote_hostname": ("opend.internal", "opend.internal"),
|
||||
"padded": (" 127.0.0.1 ", "127.0.0.1"),
|
||||
}
|
||||
for label, (configured, expected_host) in cases.items():
|
||||
with self.subTest(host=label):
|
||||
env = {} if configured is None else {"FUTU_OPEND_HOST": configured}
|
||||
with patch.dict("os.environ", env, clear=True):
|
||||
host, port = service._connection_settings()
|
||||
self.assertEqual(host, expected_host)
|
||||
self.assertEqual(port, 11111)
|
||||
|
||||
def test_connection_settings_rejects_ipv6_literal(self):
|
||||
for host in ("::1", "[::1]", "2001:db8::1"):
|
||||
with self.subTest(host=host):
|
||||
with patch.dict(
|
||||
"os.environ",
|
||||
{"FUTU_OPEND_HOST": host},
|
||||
clear=True,
|
||||
), self.assertRaisesRegex(
|
||||
service.FutuPortfolioError,
|
||||
"网络层仅支持 IPv4",
|
||||
):
|
||||
service._connection_settings()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,395 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import builtins
|
||||
import sys
|
||||
from types import SimpleNamespace
|
||||
import unittest
|
||||
from unittest.mock import MagicMock, call, patch
|
||||
|
||||
import main
|
||||
from src.brokers.futu.portfolio import FutuPortfolioError
|
||||
from src.services.runtime_scheduler import RuntimeSchedulerService
|
||||
|
||||
|
||||
class MainPortfolioTest(unittest.TestCase):
|
||||
def test_parse_arguments_accepts_futu_portfolio(self):
|
||||
with patch.object(sys, "argv", ["main.py", "--portfolio", "FUTU"]):
|
||||
args = main.parse_arguments()
|
||||
|
||||
self.assertEqual(args.portfolio, "futu")
|
||||
|
||||
def test_resolve_portfolio_stock_codes_uses_futu_loader(self):
|
||||
args = SimpleNamespace(portfolio="futu")
|
||||
with patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=["aapl", "HK01810", "005930"],
|
||||
) as loader, patch.object(
|
||||
main,
|
||||
"resolve_index_stock_code_for_analysis",
|
||||
side_effect=AssertionError("broker codes must not be remapped"),
|
||||
):
|
||||
result = main._resolve_portfolio_stock_codes(args)
|
||||
|
||||
self.assertEqual(result, ["AAPL", "HK01810", "005930"])
|
||||
loader.assert_called_once_with()
|
||||
|
||||
def test_resolve_portfolio_stock_codes_returns_none_when_disabled(self):
|
||||
self.assertIsNone(main._resolve_portfolio_stock_codes(SimpleNamespace()))
|
||||
|
||||
def test_analysis_lock_propagates_requested_portfolio_failures(self):
|
||||
config = SimpleNamespace()
|
||||
args = SimpleNamespace(portfolio="futu")
|
||||
error = FutuPortfolioError("OpenD unavailable")
|
||||
|
||||
with patch.object(
|
||||
main,
|
||||
"run_full_analysis",
|
||||
side_effect=error,
|
||||
) as runner, self.assertRaisesRegex(FutuPortfolioError, "OpenD unavailable"):
|
||||
main._run_analysis_with_runtime_scheduler_lock(
|
||||
config,
|
||||
args,
|
||||
["600519"],
|
||||
)
|
||||
|
||||
runner.assert_called_once_with(config, args, ["600519"])
|
||||
|
||||
def test_run_full_analysis_propagates_futu_portfolio_load_failure(self):
|
||||
config = SimpleNamespace()
|
||||
args = SimpleNamespace(portfolio="futu")
|
||||
error = FutuPortfolioError("OpenD unavailable")
|
||||
|
||||
with patch.object(
|
||||
main,
|
||||
"_refresh_stock_index_cache_for_analysis",
|
||||
), patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
side_effect=error,
|
||||
), self.assertRaisesRegex(FutuPortfolioError, "OpenD unavailable"):
|
||||
main.run_full_analysis(config, args)
|
||||
|
||||
def test_run_full_analysis_keeps_downstream_failures_non_propagating(self):
|
||||
config = SimpleNamespace()
|
||||
args = SimpleNamespace(portfolio="futu")
|
||||
|
||||
with patch.object(
|
||||
main,
|
||||
"_refresh_stock_index_cache_for_analysis",
|
||||
), patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=["AAPL"],
|
||||
), patch.object(
|
||||
main,
|
||||
"_compute_trading_day_filter",
|
||||
side_effect=RuntimeError("calendar unavailable"),
|
||||
):
|
||||
result = main.run_full_analysis(config, args)
|
||||
|
||||
self.assertFalse(result)
|
||||
|
||||
def test_run_full_analysis_uses_futu_holdings_and_reloads_each_run(self):
|
||||
args = SimpleNamespace(
|
||||
portfolio="futu",
|
||||
single_notify=False,
|
||||
no_context_snapshot=True,
|
||||
no_market_review=True,
|
||||
workers=1,
|
||||
dry_run=True,
|
||||
no_notify=True,
|
||||
schedule=False,
|
||||
)
|
||||
config = SimpleNamespace(
|
||||
refresh_stock_list=MagicMock(),
|
||||
single_stock_notify=False,
|
||||
merge_email_notification=False,
|
||||
market_review_enabled=False,
|
||||
market_review_region="cn",
|
||||
daily_market_context_enabled=False,
|
||||
analysis_delay=0,
|
||||
backtest_enabled=False,
|
||||
)
|
||||
pipeline = MagicMock()
|
||||
pipeline.run.return_value = []
|
||||
trading_day_filter = MagicMock(
|
||||
return_value=(["AAPL", "HK00700"], "us,hk", False)
|
||||
)
|
||||
|
||||
with patch.object(main, "_refresh_stock_index_cache_for_analysis"), patch.object(
|
||||
main,
|
||||
"_compute_trading_day_filter",
|
||||
trading_day_filter,
|
||||
), patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=["AAPL", "HK00700"],
|
||||
) as loader, patch(
|
||||
"src.core.pipeline.StockAnalysisPipeline",
|
||||
return_value=pipeline,
|
||||
), patch(
|
||||
"src.core.market_review.run_market_review",
|
||||
), patch(
|
||||
"src.feishu_doc.FeishuDocManager",
|
||||
) as feishu_manager:
|
||||
feishu_manager.return_value.is_configured.return_value = False
|
||||
first_result = main.run_full_analysis(config, args, ["600519"])
|
||||
second_result = main.run_full_analysis(config, args, ["600519"])
|
||||
|
||||
self.assertTrue(first_result)
|
||||
self.assertTrue(second_result)
|
||||
self.assertEqual(loader.call_count, 2)
|
||||
config.refresh_stock_list.assert_not_called()
|
||||
self.assertEqual(trading_day_filter.call_count, 2)
|
||||
trading_day_filter.assert_has_calls(
|
||||
[
|
||||
call(config, args, ["AAPL", "HK00700"]),
|
||||
call(config, args, ["AAPL", "HK00700"]),
|
||||
]
|
||||
)
|
||||
self.assertEqual(pipeline.run.call_count, 2)
|
||||
for invocation in pipeline.run.call_args_list:
|
||||
self.assertEqual(invocation.kwargs["stock_codes"], ["AAPL", "HK00700"])
|
||||
|
||||
def test_run_full_analysis_skips_empty_futu_portfolio_without_fallback(self):
|
||||
args = SimpleNamespace(
|
||||
portfolio="futu",
|
||||
force_run=True,
|
||||
single_notify=False,
|
||||
no_context_snapshot=True,
|
||||
no_market_review=True,
|
||||
workers=1,
|
||||
dry_run=True,
|
||||
no_notify=True,
|
||||
schedule=False,
|
||||
)
|
||||
config = SimpleNamespace(
|
||||
refresh_stock_list=MagicMock(),
|
||||
single_stock_notify=False,
|
||||
merge_email_notification=False,
|
||||
market_review_enabled=True,
|
||||
market_review_region="cn",
|
||||
daily_market_context_enabled=False,
|
||||
analysis_delay=0,
|
||||
backtest_enabled=False,
|
||||
)
|
||||
real_import = builtins.__import__
|
||||
|
||||
def reject_pipeline_import(name, *args, **kwargs):
|
||||
if name in {"src.core.market_review", "src.core.pipeline"}:
|
||||
raise AssertionError(f"empty portfolio must not import {name}")
|
||||
return real_import(name, *args, **kwargs)
|
||||
|
||||
with patch.object(
|
||||
main,
|
||||
"_refresh_stock_index_cache_for_analysis",
|
||||
) as refresh_stock_index, patch.object(
|
||||
main,
|
||||
"_compute_trading_day_filter",
|
||||
return_value=([], None, False),
|
||||
) as trading_day_filter, patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=[],
|
||||
), patch.object(
|
||||
builtins,
|
||||
"__import__",
|
||||
side_effect=reject_pipeline_import,
|
||||
), self.assertLogs(main.logger, level="INFO") as captured:
|
||||
result = main.run_full_analysis(config, args, ["600519"])
|
||||
|
||||
self.assertTrue(result)
|
||||
config.refresh_stock_list.assert_not_called()
|
||||
refresh_stock_index.assert_not_called()
|
||||
trading_day_filter.assert_not_called()
|
||||
log_text = "\n".join(captured.output)
|
||||
self.assertIn("无符合条件的 Futu 持仓", log_text)
|
||||
self.assertNotIn("未配置自选股列表", log_text)
|
||||
|
||||
def test_empty_futu_portfolio_is_noop_when_trading_day_check_is_disabled(self):
|
||||
args = SimpleNamespace(
|
||||
portfolio="futu",
|
||||
force_run=False,
|
||||
no_market_review=False,
|
||||
)
|
||||
config = SimpleNamespace(
|
||||
refresh_stock_list=MagicMock(),
|
||||
market_review_enabled=False,
|
||||
trading_day_check_enabled=False,
|
||||
)
|
||||
|
||||
with patch.object(
|
||||
main,
|
||||
"_refresh_stock_index_cache_for_analysis",
|
||||
) as refresh_stock_index, patch.object(
|
||||
main,
|
||||
"_compute_trading_day_filter",
|
||||
) as trading_day_filter, patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=[],
|
||||
), self.assertLogs(main.logger, level="INFO") as captured:
|
||||
result = main.run_full_analysis(config, args, ["600519"])
|
||||
|
||||
self.assertTrue(result)
|
||||
config.refresh_stock_list.assert_not_called()
|
||||
refresh_stock_index.assert_not_called()
|
||||
trading_day_filter.assert_not_called()
|
||||
self.assertIn("无符合条件的 Futu 持仓", "\n".join(captured.output))
|
||||
|
||||
def test_empty_futu_portfolio_preserves_enabled_auto_backtest(self):
|
||||
args = SimpleNamespace(
|
||||
portfolio="futu",
|
||||
no_market_review=True,
|
||||
)
|
||||
config = SimpleNamespace(
|
||||
market_review_enabled=True,
|
||||
backtest_enabled=True,
|
||||
backtest_eval_window_days=10,
|
||||
backtest_min_age_days=14,
|
||||
)
|
||||
backtest_service = MagicMock()
|
||||
backtest_service.run_backtest.return_value = {
|
||||
"processed": 1,
|
||||
"saved": 1,
|
||||
"completed": 1,
|
||||
"insufficient": 0,
|
||||
"errors": 0,
|
||||
}
|
||||
real_import = builtins.__import__
|
||||
|
||||
def reject_pipeline_import(name, *args, **kwargs):
|
||||
if name in {"src.core.market_review", "src.core.pipeline"}:
|
||||
raise AssertionError(f"empty portfolio must not import {name}")
|
||||
return real_import(name, *args, **kwargs)
|
||||
|
||||
with patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=[],
|
||||
), patch(
|
||||
"src.services.backtest_service.BacktestService",
|
||||
return_value=backtest_service,
|
||||
) as backtest_class, patch.object(
|
||||
builtins,
|
||||
"__import__",
|
||||
side_effect=reject_pipeline_import,
|
||||
):
|
||||
result = main.run_full_analysis(config, args)
|
||||
|
||||
self.assertTrue(result)
|
||||
backtest_class.assert_called_once_with()
|
||||
backtest_service.run_backtest.assert_called_once_with(
|
||||
force=False,
|
||||
eval_window_days=10,
|
||||
min_age_days=14,
|
||||
limit=200,
|
||||
)
|
||||
|
||||
def test_empty_futu_portfolio_uses_an_accurate_skip_reason(self):
|
||||
args = SimpleNamespace(portfolio="futu")
|
||||
config = SimpleNamespace(refresh_stock_list=MagicMock())
|
||||
|
||||
with patch.object(main, "_refresh_stock_index_cache_for_analysis"), patch.object(
|
||||
main,
|
||||
"_compute_trading_day_filter",
|
||||
return_value=([], "", True),
|
||||
), patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=[],
|
||||
), patch(
|
||||
"src.core.pipeline.StockAnalysisPipeline",
|
||||
) as pipeline_class, patch(
|
||||
"src.core.market_review.run_market_review",
|
||||
), self.assertLogs(main.logger, level="INFO") as captured:
|
||||
result = main.run_full_analysis(config, args)
|
||||
|
||||
self.assertTrue(result)
|
||||
pipeline_class.assert_not_called()
|
||||
log_text = "\n".join(captured.output)
|
||||
self.assertIn("无符合条件的 Futu 持仓", log_text)
|
||||
self.assertNotIn("所有相关市场均为非交易日", log_text)
|
||||
|
||||
def test_futu_portfolio_without_effective_codes_still_runs_market_review(self):
|
||||
args = SimpleNamespace(
|
||||
portfolio="futu",
|
||||
single_notify=False,
|
||||
no_context_snapshot=True,
|
||||
no_market_review=False,
|
||||
workers=1,
|
||||
dry_run=True,
|
||||
no_notify=True,
|
||||
schedule=False,
|
||||
)
|
||||
config = SimpleNamespace(
|
||||
refresh_stock_list=MagicMock(),
|
||||
single_stock_notify=False,
|
||||
merge_email_notification=False,
|
||||
market_review_enabled=True,
|
||||
market_review_region="cn",
|
||||
daily_market_context_enabled=False,
|
||||
analysis_delay=0,
|
||||
backtest_enabled=False,
|
||||
)
|
||||
for holdings in ([], ["AAPL"]):
|
||||
with self.subTest(holdings=holdings):
|
||||
pipeline = MagicMock()
|
||||
run_market_review = MagicMock()
|
||||
|
||||
with patch.object(
|
||||
main,
|
||||
"_refresh_stock_index_cache_for_analysis",
|
||||
), patch.object(
|
||||
main,
|
||||
"_compute_trading_day_filter",
|
||||
return_value=([], "cn", False),
|
||||
), patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=holdings,
|
||||
), patch(
|
||||
"src.core.pipeline.StockAnalysisPipeline",
|
||||
return_value=pipeline,
|
||||
), patch(
|
||||
"src.core.market_review.run_market_review",
|
||||
run_market_review,
|
||||
), patch.object(
|
||||
main,
|
||||
"_run_market_review_with_shared_lock",
|
||||
return_value=SimpleNamespace(report="market review"),
|
||||
) as run_with_lock, patch(
|
||||
"src.feishu_doc.FeishuDocManager",
|
||||
) as feishu_manager:
|
||||
feishu_manager.return_value.is_configured.return_value = False
|
||||
result = main.run_full_analysis(config, args)
|
||||
|
||||
self.assertTrue(result)
|
||||
pipeline.run.assert_not_called()
|
||||
run_with_lock.assert_called_once()
|
||||
|
||||
def test_runtime_scheduler_preserves_futu_portfolio_override(self):
|
||||
scheduler = RuntimeSchedulerService(
|
||||
owns_schedule=False,
|
||||
schedule_args_overrides={"portfolio": "futu"},
|
||||
)
|
||||
|
||||
self.assertEqual(scheduler._make_schedule_args().portfolio, "futu")
|
||||
|
||||
def test_runtime_scheduler_records_futu_load_failure_and_keeps_running(self):
|
||||
config = SimpleNamespace(
|
||||
schedule_enabled=True,
|
||||
schedule_time="18:00",
|
||||
schedule_times=["18:00"],
|
||||
)
|
||||
error = FutuPortfolioError("OpenD unavailable")
|
||||
|
||||
def runner(config_arg, args, stock_codes):
|
||||
raise error
|
||||
|
||||
scheduler = RuntimeSchedulerService(
|
||||
config_provider=lambda: config,
|
||||
task_runner=runner,
|
||||
)
|
||||
scheduler._reload_config = lambda: config
|
||||
|
||||
self.assertTrue(scheduler._run_analysis_once())
|
||||
status = scheduler.status()
|
||||
self.assertIsNone(status["last_success_at"])
|
||||
self.assertEqual(status["last_error"], "OpenD unavailable")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -18,6 +18,7 @@ ensure_litellm_stub()
|
||||
|
||||
_ENV_BEFORE_MAIN_IMPORT = dict(os.environ)
|
||||
import main
|
||||
from src.brokers.futu.portfolio import FutuPortfolioError
|
||||
from src.config import Config
|
||||
|
||||
_MAIN_IMPORT_ENV_ADDITIONS = frozenset(set(os.environ) - set(_ENV_BEFORE_MAIN_IMPORT))
|
||||
@@ -88,6 +89,7 @@ class MainScheduleModeTestCase(unittest.TestCase):
|
||||
defaults = {
|
||||
"debug": False,
|
||||
"stocks": None,
|
||||
"portfolio": None,
|
||||
"webui": False,
|
||||
"webui_only": False,
|
||||
"serve": False,
|
||||
@@ -424,6 +426,72 @@ class MainScheduleModeTestCase(unittest.TestCase):
|
||||
_, _, stock_codes = run_full_analysis.call_args.args
|
||||
self.assertEqual(stock_codes, ["005930.KS"])
|
||||
|
||||
def test_standalone_futu_portfolio_failure_returns_nonzero(self) -> None:
|
||||
args = self._make_args(portfolio="futu")
|
||||
config = self._make_config(run_immediately=True)
|
||||
error = FutuPortfolioError("OpenD unavailable")
|
||||
|
||||
with (
|
||||
patch("main.parse_arguments", return_value=args),
|
||||
patch("main.get_config", return_value=config),
|
||||
patch("main.setup_logging"),
|
||||
patch("main._refresh_stock_index_cache_for_analysis"),
|
||||
patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
side_effect=error,
|
||||
) as loader,
|
||||
):
|
||||
exit_code = main.main()
|
||||
|
||||
self.assertEqual(exit_code, 1)
|
||||
loader.assert_called_once_with()
|
||||
|
||||
def test_standalone_futu_portfolio_success_returns_zero(self) -> None:
|
||||
args = self._make_args(portfolio="futu")
|
||||
config = self._make_config(run_immediately=True)
|
||||
|
||||
with (
|
||||
patch("main.parse_arguments", return_value=args),
|
||||
patch("main.get_config", return_value=config),
|
||||
patch("main.setup_logging"),
|
||||
patch("main._refresh_stock_index_cache_for_analysis"),
|
||||
patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=["AAPL"],
|
||||
) as loader,
|
||||
patch(
|
||||
"main._compute_trading_day_filter",
|
||||
return_value=([], "", True),
|
||||
),
|
||||
):
|
||||
exit_code = main.main()
|
||||
|
||||
self.assertEqual(exit_code, 0)
|
||||
loader.assert_called_once_with()
|
||||
|
||||
def test_standalone_futu_downstream_failure_keeps_existing_exit_semantics(self) -> None:
|
||||
args = self._make_args(portfolio="futu")
|
||||
config = self._make_config(run_immediately=True)
|
||||
|
||||
with (
|
||||
patch("main.parse_arguments", return_value=args),
|
||||
patch("main.get_config", return_value=config),
|
||||
patch("main.setup_logging"),
|
||||
patch("main._refresh_stock_index_cache_for_analysis"),
|
||||
patch(
|
||||
"src.brokers.futu.portfolio.load_futu_stock_codes",
|
||||
return_value=["AAPL"],
|
||||
) as loader,
|
||||
patch(
|
||||
"main._compute_trading_day_filter",
|
||||
side_effect=RuntimeError("calendar unavailable"),
|
||||
),
|
||||
):
|
||||
exit_code = main.main()
|
||||
|
||||
self.assertEqual(exit_code, 0)
|
||||
loader.assert_called_once_with()
|
||||
|
||||
def test_schedule_mode_reload_uses_latest_runtime_config(self) -> None:
|
||||
args = self._make_args(schedule=True)
|
||||
startup_config = self._make_config(schedule_enabled=True, schedule_time="18:00")
|
||||
@@ -731,18 +799,26 @@ class MainScheduleModeTestCase(unittest.TestCase):
|
||||
run_with_schedule.assert_not_called()
|
||||
|
||||
def test_serve_mode_uses_shared_analysis_lock_for_immediate_run_full_analysis(self) -> None:
|
||||
args = self._make_args(serve=True, schedule=False, host="127.0.0.1", port=8000)
|
||||
args = self._make_args(
|
||||
serve=True,
|
||||
schedule=False,
|
||||
portfolio="futu",
|
||||
host="127.0.0.1",
|
||||
port=8000,
|
||||
)
|
||||
config = self._make_config(webui_enabled=False, run_immediately=True)
|
||||
|
||||
with patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False), \
|
||||
patch("main.parse_arguments", return_value=args), \
|
||||
patch("main.get_config", return_value=config), \
|
||||
patch("main.prepare_webui_frontend_assets", return_value=True), \
|
||||
patch("main.start_api_server"), \
|
||||
patch("main.start_bot_stream_clients") as start_bots, \
|
||||
patch("main.time.sleep", side_effect=KeyboardInterrupt), \
|
||||
patch("main.run_full_analysis") as run_full_analysis, \
|
||||
patch("main._run_analysis_with_runtime_scheduler_lock") as run_with_lock:
|
||||
with (
|
||||
patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False),
|
||||
patch("main.parse_arguments", return_value=args),
|
||||
patch("main.get_config", return_value=config),
|
||||
patch("main.prepare_webui_frontend_assets", return_value=True),
|
||||
patch("main.start_api_server"),
|
||||
patch("main.start_bot_stream_clients") as start_bots,
|
||||
patch("main.time.sleep", side_effect=KeyboardInterrupt),
|
||||
patch("main.run_full_analysis") as run_full_analysis,
|
||||
patch("main._run_analysis_with_runtime_scheduler_lock") as run_with_lock,
|
||||
):
|
||||
exit_code = main.main()
|
||||
|
||||
self.assertEqual(exit_code, 0)
|
||||
@@ -751,6 +827,41 @@ class MainScheduleModeTestCase(unittest.TestCase):
|
||||
run_full_analysis.assert_not_called()
|
||||
start_bots.assert_called_once_with(config)
|
||||
|
||||
def test_serve_mode_keeps_running_after_futu_portfolio_load_failure(self) -> None:
|
||||
args = self._make_args(
|
||||
serve=True,
|
||||
schedule=False,
|
||||
portfolio="futu",
|
||||
host="127.0.0.1",
|
||||
port=8000,
|
||||
)
|
||||
config = self._make_config(webui_enabled=False, run_immediately=True)
|
||||
error = FutuPortfolioError("OpenD unavailable")
|
||||
|
||||
with (
|
||||
patch.dict(os.environ, {"GITHUB_ACTIONS": "false"}, clear=False),
|
||||
patch("main.parse_arguments", return_value=args),
|
||||
patch("main.get_config", return_value=config),
|
||||
patch("main.prepare_webui_frontend_assets", return_value=True),
|
||||
patch("main.start_api_server"),
|
||||
patch("main.start_bot_stream_clients") as start_bots,
|
||||
patch("main.time.sleep", side_effect=KeyboardInterrupt),
|
||||
patch(
|
||||
"main._run_analysis_with_runtime_scheduler_lock",
|
||||
side_effect=error,
|
||||
) as run_with_lock,
|
||||
patch("main.logger.exception") as exception_log,
|
||||
):
|
||||
exit_code = main.main()
|
||||
|
||||
self.assertEqual(exit_code, 0)
|
||||
run_with_lock.assert_called_once_with(config, args, None)
|
||||
start_bots.assert_called_once_with(config)
|
||||
exception_log.assert_any_call(
|
||||
"Futu 持仓导入失败,Web/API 服务继续运行: %s",
|
||||
error,
|
||||
)
|
||||
|
||||
def test_serve_schedule_flag_enables_api_runtime_scheduler(self) -> None:
|
||||
from src.services.runtime_scheduler import (
|
||||
CLI_SCHEDULER_OWNER_ENV,
|
||||
|
||||
Reference in New Issue
Block a user