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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
181 changes: 155 additions & 26 deletions application/execution_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
import json
import tempfile
import time
from collections.abc import Mapping
from dataclasses import dataclass
from datetime import datetime
from pathlib import Path
from typing import Any
Expand Down Expand Up @@ -42,14 +44,22 @@ def should_sell_cash_sweep_to_fund_whole_share_buy(
return False
try:
from quant_platform_kit.common.small_account_compatibility import (
project_unbuyable_value_targets_to_cash,
apply_small_account_cash_compatibility,
format_small_account_cash_substitution_notes,
)
except ImportError: # pragma: no cover - compatibility with older pinned shared wheels
def project_unbuyable_value_targets_to_cash(
@dataclass(frozen=True)
class _SmallAccountCashCompatibilityResult:
targets: dict[str, float]
whole_share_substituted_symbols: tuple[str, ...]
safe_haven_cash_substituted_symbols: tuple[str, ...]
cash_substitution_notes: tuple[dict[str, Any], ...]

def _project_unbuyable_value_targets_to_cash(
target_values,
prices,
*,
symbols=None,
candidate_symbols=None,
quantity_step=1.0,
):
adjusted = {
Expand All @@ -59,23 +69,140 @@ def project_unbuyable_value_targets_to_cash(
step = max(0.0, float(quantity_step or 0.0))
if step <= 0.0:
return adjusted, ()
candidate_symbols = (
normalized_candidates = (
tuple(adjusted)
if symbols is None
else tuple(dict.fromkeys(str(symbol or "").strip().upper() for symbol in symbols))
if candidate_symbols is None
else tuple(dict.fromkeys(str(symbol or "").strip().upper() for symbol in candidate_symbols))
)
normalized_prices = {
str(symbol or "").strip().upper(): float(price or 0.0)
for symbol, price in dict(prices or {}).items()
}
substituted = []
for symbol in candidate_symbols:
for symbol in normalized_candidates:
target_value = max(0.0, float(adjusted.get(symbol, 0.0) or 0.0))
price = max(0.0, float(normalized_prices.get(symbol, 0.0) or 0.0))
if price > 0.0 and 0.0 < target_value < (price * step):
adjusted[symbol] = 0.0
substituted.append(symbol)
return adjusted, tuple(dict.fromkeys(substituted))

def apply_small_account_cash_compatibility(
target_values,
prices,
*,
candidate_symbols=None,
safe_haven_cash_symbols=(),
quantity_step=1.0,
cash_substitute_limit_usd=2000.0,
):
adjusted_targets, substituted = _project_unbuyable_value_targets_to_cash(
target_values,
prices,
candidate_symbols=candidate_symbols,
quantity_step=quantity_step,
)
normalized_candidates = (
tuple(adjusted_targets)
if candidate_symbols is None
else tuple(dict.fromkeys(str(symbol or "").strip().upper() for symbol in candidate_symbols))
)
remaining_non_safe_targets = [
symbol
for symbol in normalized_candidates
if float(adjusted_targets.get(str(symbol or "").strip().upper(), 0.0) or 0.0) > 0.0
]
safe_haven_symbols = tuple(
dict.fromkeys(
str(symbol or "").strip().upper()
for symbol in safe_haven_cash_symbols
if str(symbol or "").strip()
)
)
safe_haven_substituted = []
if (
substituted
and not remaining_non_safe_targets
and _positive_target_total(adjusted_targets) <= max(0.0, float(cash_substitute_limit_usd or 0.0))
):
for symbol in safe_haven_symbols:
if float(adjusted_targets.get(symbol, 0.0) or 0.0) > 0.0:
adjusted_targets[symbol] = 0.0
safe_haven_substituted.append(symbol)
normalized_targets = {
str(symbol or "").strip().upper(): float(value or 0.0)
for symbol, value in dict(target_values or {}).items()
}
normalized_prices = {
str(symbol or "").strip().upper(): float(price or 0.0)
for symbol, price in dict(prices or {}).items()
}
notes = []
if safe_haven_substituted:
for symbol in substituted:
target_value = max(0.0, float(normalized_targets.get(symbol, 0.0) or 0.0))
price = max(0.0, float(normalized_prices.get(symbol, 0.0) or 0.0))
if target_value <= 0.0 or price <= 0.0:
continue
notes.append(
{
"symbol": symbol,
"target_value": target_value,
"price": price,
"cash_symbols": tuple(safe_haven_substituted),
}
)
return _SmallAccountCashCompatibilityResult(
targets=adjusted_targets,
whole_share_substituted_symbols=substituted,
safe_haven_cash_substituted_symbols=tuple(safe_haven_substituted),
cash_substitution_notes=tuple(notes),
)

def format_small_account_cash_substitution_notes(
notes,
*,
translator,
wrapper_key="buy_deferred",
detail_key="buy_deferred_small_account_cash_substitution",
cash_label_key="cash_label",
symbol_suffix=".US",
):
messages = []
seen_keys = set()
for note in tuple(notes or ()):
if not isinstance(note, Mapping):
continue
symbol = str(note.get("symbol") or "").strip().upper()
if not symbol:
continue
target_value = max(0.0, float(note.get("target_value") or 0.0))
price = max(0.0, float(note.get("price") or 0.0))
if target_value <= 0.0 or price <= 0.0:
continue
cash_symbols = tuple(
dict.fromkeys(
str(cash_symbol or "").strip().upper()
for cash_symbol in tuple(note.get("cash_symbols") or ())
if str(cash_symbol or "").strip()
)
)
cash_symbols_text = ", ".join(f"{cash_symbol}{symbol_suffix}" for cash_symbol in cash_symbols)
if not cash_symbols_text:
cash_symbols_text = translator(cash_label_key)
note_key = (symbol, f"{target_value:.2f}", cash_symbols_text)
if note_key in seen_keys:
continue
seen_keys.add(note_key)
detail = translator(
detail_key,
symbol=f"{symbol}{symbol_suffix}",
diff=f"{target_value:.2f}",
price=f"{price:.2f}",
cash_symbols=cash_symbols_text,
)
messages.append(translator(wrapper_key, detail=detail))
return tuple(messages)
from quant_platform_kit.common.quantity import (
floor_to_quantity_step,
format_quantity,
Expand Down Expand Up @@ -663,6 +790,7 @@ def execute_rebalance(
"safe_haven_cash_substituted_symbols": [],
"small_account_whole_share_substituted_symbols": [],
"small_account_safe_haven_cash_substituted_symbols": [],
"small_account_whole_share_cash_notes": [],
"residual_cash_estimate": float(account_values.get("buying_power", 0.0) or 0.0),
"current_stock_weight": 0.0,
"current_safe_haven_weight": 0.0,
Expand Down Expand Up @@ -750,31 +878,23 @@ def execute_rebalance(
for symbol in target_mv
if str(symbol or "").strip().upper() not in safe_haven_symbols
)
target_mv, small_account_substituted_symbols = project_unbuyable_value_targets_to_cash(
small_account_compatibility = apply_small_account_cash_compatibility(
target_mv,
prices,
symbols=small_account_candidate_symbols,
candidate_symbols=small_account_candidate_symbols,
safe_haven_cash_symbols=safe_haven_symbols,
quantity_step=order_quantity_step,
cash_substitute_limit_usd=SMALL_ACCOUNT_SAFE_HAVEN_CASH_SUBSTITUTE_LIMIT_USD,
)
target_mv = small_account_compatibility.targets
small_account_substituted_symbols = small_account_compatibility.whole_share_substituted_symbols
for symbol in small_account_substituted_symbols:
target_weights[symbol] = 0.0
remaining_non_safe_targets = [
symbol
for symbol in small_account_candidate_symbols
if float(target_mv.get(str(symbol or "").strip().upper(), 0.0) or 0.0) > 0.0
]
small_account_safe_haven_cash_substituted_symbols: list[str] = []
if (
small_account_substituted_symbols
and not remaining_non_safe_targets
and _positive_target_total(target_mv) <= SMALL_ACCOUNT_SAFE_HAVEN_CASH_SUBSTITUTE_LIMIT_USD
):
for symbol in safe_haven_symbols:
normalized_symbol = str(symbol or "").strip().upper()
if float(target_mv.get(normalized_symbol, 0.0) or 0.0) > 0.0:
target_mv[normalized_symbol] = 0.0
target_weights[normalized_symbol] = 0.0
small_account_safe_haven_cash_substituted_symbols.append(normalized_symbol)
small_account_safe_haven_cash_substituted_symbols = list(
small_account_compatibility.safe_haven_cash_substituted_symbols
)
for symbol in small_account_safe_haven_cash_substituted_symbols:
target_weights[symbol] = 0.0
if safe_haven_symbols:
execution_summary["realized_safe_haven_weight"] = float(
sum(
Expand All @@ -788,7 +908,16 @@ def execute_rebalance(
execution_summary["small_account_safe_haven_cash_substituted_symbols"] = (
small_account_safe_haven_cash_substituted_symbols
)
execution_summary["small_account_whole_share_cash_notes"] = list(
small_account_compatibility.cash_substitution_notes
)
trade_logs = []
trade_logs.extend(
format_small_account_cash_substitution_notes(
small_account_compatibility.cash_substitution_notes,
translator=translator,
)
)
target_hash = _build_target_hash(target_weights)
execution_summary["target_vs_current"] = _build_target_diff_rows(target_weights, current_mv, equity)
if equity > 0:
Expand Down
6 changes: 6 additions & 0 deletions notifications/telegram.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
"buying_power": "购买力",
"reserved_cash": "预留现金",
"investable_cash": "可投资现金",
"cash_label": "现金",
"empty_positions": " (空仓)",
"empty_target_weights": " (无目标持仓)",
"account_summary_title": "📊 账户摘要",
Expand All @@ -40,6 +41,8 @@
"price_fallback_prices": "📌 执行估价: 使用最近收盘价 {count}个标的 ({symbols})",
"target_diff_summary": "调仓变化: {details}",
"no_order_plan_reason": "未下单: {reason}",
"buy_deferred": "ℹ️ [买入说明] {detail}",
"buy_deferred_small_account_cash_substitution": "{symbol} 目标金额 ${diff} 低于 1 股价格 ${price};为避免超过目标仓位,小账户本轮保留现金,不回补 {cash_symbols}",
"trade_date_detail": "交易日={value}",
"target_diff": "目标差异 {symbol}: 当前={current} 目标={target} 变化={delta}",
"pending_orders_detected": "检测到未完成订单: 策略={profile} 标的={symbols}",
Expand Down Expand Up @@ -151,6 +154,7 @@
"buying_power": "Buying Power",
"reserved_cash": "Reserved Cash",
"investable_cash": "Investable Cash",
"cash_label": "Cash",
"empty_positions": " (No positions)",
"empty_target_weights": " (No target positions)",
"account_summary_title": "📊 Account Summary",
Expand All @@ -174,6 +178,8 @@
"price_fallback_prices": "📌 execution pricing: recent close for {count} symbols ({symbols})",
"target_diff_summary": "Target changes: {details}",
"no_order_plan_reason": "No order submitted: {reason}",
"buy_deferred": "ℹ️ [Buy note] {detail}",
"buy_deferred_small_account_cash_substitution": "{symbol} target ${diff} is below the 1-share price ${price}; to avoid exceeding the target allocation, this small account keeps cash this cycle and does not rebuy {cash_symbols}",
"trade_date_detail": "trade_date={value}",
"target_diff": "target_diff {symbol}: current={current} target={target} delta={delta}",
"pending_orders_detected": "pending_orders_detected profile={profile} symbols={symbols}",
Expand Down
4 changes: 2 additions & 2 deletions requirements.txt
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
flask
gunicorn
quant-platform-kit @ git+https://github.com/QuantStrategyLab/QuantPlatformKit.git@190edb21c5fe0c8e0efba66bfd8bacfdd00b3f77
us-equity-strategies @ git+https://github.com/QuantStrategyLab/UsEquityStrategies.git@ead0c075f5b47ec1d05d3a45cfcdda58217fe086
quant-platform-kit @ git+https://github.com/QuantStrategyLab/QuantPlatformKit.git@ceb84a366ed1bf9a53292ff2c73e06b4baac05e2
us-equity-strategies @ git+https://github.com/QuantStrategyLab/UsEquityStrategies.git@f2ebae8aacd8c70292c5b6115a80c6657e64ad1f
pandas
numpy
requests
Expand Down
15 changes: 14 additions & 1 deletion tests/test_execution_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,12 @@ def translate(key, **kwargs):
"dry_run_snapshot_prices": "dry_run_snapshot_prices count={count} symbols={symbols}",
"price_fallback_prices": "price_fallback_prices count={count} symbols={symbols}",
"no_equity": "❌ No equity",
"cash_label": "现金",
"buy_deferred": "ℹ️ [买入说明] {detail}",
"buy_deferred_small_account_cash_substitution": (
"{symbol} 目标金额 ${diff} 低于 1 股价格 ${price};"
"为避免超过目标仓位,小账户本轮保留现金,不回补 {cash_symbols}"
),
}
template = templates[key]
return template.format(**kwargs) if kwargs else template
Expand Down Expand Up @@ -238,7 +244,7 @@ def accountValues(self):
submitted = []
monkeypatch.setattr("application.execution_service.time.sleep", lambda _seconds: None)

_trade_logs, summary = execute_rebalance(
trade_logs, summary = execute_rebalance(
FakeIB(),
{},
{},
Expand Down Expand Up @@ -273,6 +279,13 @@ def accountValues(self):
assert submitted == []
assert summary["small_account_whole_share_substituted_symbols"] == ["SOXX"]
assert summary["small_account_safe_haven_cash_substituted_symbols"] == ["BOXX"]
assert len(summary["small_account_whole_share_cash_notes"]) == 1
cash_note = summary["small_account_whole_share_cash_notes"][0]
assert cash_note["symbol"] == "SOXX"
assert round(cash_note["target_value"], 3) == 188.277
assert cash_note["price"] == 525.0
assert cash_note["cash_symbols"] == ("BOXX",)
assert any("SOXX.US 目标金额 $188.28 低于 1 股价格 $525.00" in log for log in trade_logs)
assert summary["realized_safe_haven_weight"] == 0.0
boxx_row = next(row for row in summary["target_vs_current"] if row["symbol"] == "BOXX")
assert boxx_row["target_weight"] == 0.0
Expand Down