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
14 changes: 13 additions & 1 deletion application/longbridge_portfolio.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,7 @@ def warn(message: str) -> None:
filter_enabled = bool(assets)

position_rows: list[tuple[str, str, Any, Any]] = []
position_currency_by_symbol: dict[str, str | None] = {}
try:
positions_response = t_ctx.stock_positions()
except Exception as exc:
Expand Down Expand Up @@ -176,6 +177,13 @@ def warn(message: str) -> None:
raise RuntimeError("LongBridge position quantity missing")
if raw_available_quantity is None:
raw_available_quantity = raw_quantity
raw_currency = str(getattr(position, "currency", "") or "").strip().upper()
position_currency = raw_currency if len(raw_currency) == 3 and raw_currency.isalpha() else None
if root_symbol in position_currency_by_symbol:
if position_currency_by_symbol[root_symbol] != position_currency:
position_currency_by_symbol[root_symbol] = None
else:
position_currency_by_symbol[root_symbol] = position_currency
if position_log_fn is not None:
position_log_fn(
"[position_snapshot] raw "
Expand All @@ -185,7 +193,10 @@ def warn(message: str) -> None:

position_rows.append((root_symbol, full_symbol, raw_quantity, raw_available_quantity))

prices = _fetch_last_prices(q_ctx, [full_symbol for _root_symbol, full_symbol, _quantity, _available in position_rows])
prices = _fetch_last_prices(
q_ctx,
[full_symbol for _root_symbol, full_symbol, _quantity, _available in position_rows],
)
for root_symbol, full_symbol, raw_quantity, raw_available_quantity in position_rows:
try:
quantity = float(raw_quantity)
Expand Down Expand Up @@ -215,6 +226,7 @@ def warn(message: str) -> None:
return {
"broker_capital": broker_capital,
"heartbeat_account_snapshot": _heartbeat_account_snapshot(account_balance, trading_currency, observed_at),
"position_currency_by_symbol": position_currency_by_symbol,
"available_cash": available_cash,
"cash_by_currency": cash_by_currency,
"market_values": market_values,
Expand Down
59 changes: 46 additions & 13 deletions application/rebalance_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -657,20 +657,37 @@ def fetch_replanned_state():
except ValueError as exc:
if str(exc) not in _LIVE_COMMAND_BINDING_ERRORS:
raise
message = "Durable live execution command binding is invalid; broker orders blocked"
runtime.notify_issue("Durable live execution blocked", message)
message = config.translator("issue_durable_execution_binding_invalid")
runtime.notify_issue(
config.translator("issue_durable_execution_blocked_title"), message
)
blocked_execution = {
"execution_status": "blocked",
"blocked_reason": "durable_live_execution_command_binding_invalid",
"heartbeat_execution_state": "blocked",
"durable_live_execution_command": {
"command_id": command.command_id,
"status": "BLOCKED_INVALID_BINDING",
"effective_date": command.effective_date,
},
}
notification_publisher.publish(
notification_renderers.render_heartbeat_notification(
execution=blocked_execution,
skip_logs=(),
note_logs=(),
translator=config.translator,
separator=config.separator,
strategy_display_name=config.strategy_display_name,
dry_run_only=config.dry_run_only,
extra_notification_lines=config.extra_notification_lines,
title_key=config.notification_title_key or "heartbeat_title",
)
)
return ExecutionCycleResult(
plan={},
portfolio={},
execution={
"execution_status": "blocked",
"blocked_reason": "durable_live_execution_command_binding_invalid",
"durable_live_execution_command": {
"command_id": command.command_id,
"status": "BLOCKED_INVALID_BINDING",
"effective_date": command.effective_date,
},
},
execution=blocked_execution,
allocation={},
logs=(),
skip_logs=(),
Expand Down Expand Up @@ -870,6 +887,15 @@ def fetch_replanned_state():
if direct_live_routing_blocked:
execution["direct_live_routing_blocked"] = True
execution["direct_live_routing_block_reason"] = "durable_execution_command_required"
if (
account_identity_blocked
or (direct_live_routing_blocked and not live_command_waiting)
or (live_command_blocked and not live_command_waiting)
or str(execution.get("execution_status") or "").strip().lower() == "blocked"
):
execution["heartbeat_execution_state"] = "blocked"
elif live_command_waiting:
execution["heartbeat_execution_state"] = "waiting_window"
execution_already_recorded = (
direct_live_routing_blocked or account_identity_blocked or live_command_blocked
)
Expand Down Expand Up @@ -934,8 +960,10 @@ def fetch_replanned_state():
elif live_command_waiting:
message = "Durable live execution command queued; waiting for its effective trading session"
elif live_command_blocked:
message = "Durable live execution command is unresolved; broker orders blocked"
runtime.notify_issue("Durable live execution blocked", message)
message = config.translator("issue_durable_execution_unresolved")
runtime.notify_issue(
config.translator("issue_durable_execution_blocked_title"), message
)
elif direct_live_routing_blocked:
message = _durable_command_required_message(execution=execution)
runtime.notify_issue("Next-session execution blocked", message)
Expand Down Expand Up @@ -1107,6 +1135,11 @@ def submit_claimed_order(order_intent):
execution.pop("heartbeat_account_snapshot", None)
if isinstance(account_snapshot, dict):
execution["heartbeat_account_snapshot"] = dict(account_snapshot)
position_currency_by_symbol = (getattr(initial_snapshot, "metadata", {}) or {}).get(
"position_currency_by_symbol"
)
if isinstance(position_currency_by_symbol, dict):
execution["position_currency_by_symbol"] = dict(position_currency_by_symbol)

if pending_orders:
try:
Expand Down
1 change: 1 addition & 0 deletions application/runtime_broker_adapters.py
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,7 @@ def build_portfolio_snapshot_from_account_state(self, account_state):
"account_hash": self.account_hash,
"broker_capital": account_state.get("broker_capital"),
"heartbeat_account_snapshot": account_state.get("heartbeat_account_snapshot"),
"position_currency_by_symbol": account_state.get("position_currency_by_symbol"),
},
)

Expand Down
37 changes: 37 additions & 0 deletions notifications/compact_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@
"💼 Strategy holdings",
}
_NUMBER_RE = re.compile(r"[-+]?\d[\d,]*(?:\.\d+)?")
_ISO_CURRENCY_RE = re.compile(r"^[A-Z]{3}$")
_DOLLAR_AMOUNT_RE = re.compile(r"\$\s*(?=[-+]?\d)")


def _contains_nonzero_number(text: str) -> bool:
Expand Down Expand Up @@ -42,6 +44,41 @@ def replace_share(match: re.Match[str]) -> str:
return re.sub(r"([-+]?\d[\d,]*(?:\.\d+)?)\s*股", replace_share, detail)


def label_holding_currencies(
dashboard_text: str,
*,
position_currency_by_symbol: object,
unknown_currency_label: str,
) -> str:
"""Label dollar-denominated holding amounts only from same-cycle position metadata."""
currencies = position_currency_by_symbol if isinstance(position_currency_by_symbol, dict) else {}
normalized_currencies = {
str(symbol).strip().upper(): str(currency).strip().upper()
for symbol, currency in currencies.items()
if isinstance(currency, str) and _ISO_CURRENCY_RE.fullmatch(currency.strip().upper())
}
output: list[str] = []
in_holdings = False
for raw_line in str(dashboard_text or "").splitlines():
line = raw_line
stripped = line.strip()
if stripped in _HOLDINGS_HEADERS:
in_holdings = True
output.append(line)
continue
if in_holdings and (not stripped or stripped.startswith(("━", "📌", "💵", "📊", "🎯", "🧾", "⏱", "🧩"))):
in_holdings = False
if in_holdings and (":" in line or ":" in line):
separator = ":" if ":" in line else ":"
prefix, remainder = line.split(separator, 1)
symbol = prefix.strip().lstrip("-• ").strip().upper()
currency = normalized_currencies.get(symbol, unknown_currency_label)
remainder = _DOLLAR_AMOUNT_RE.sub(f"{currency} ", remainder)
line = f"{prefix}{separator}{remainder}"
output.append(line)
return "\n".join(output)


def adapt_compact_sections(
dashboard_text: str,
*,
Expand Down
47 changes: 41 additions & 6 deletions notifications/renderers.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
from collections.abc import Mapping
import math

from notifications.compact_adapter import adapt_compact_sections
from notifications.compact_adapter import adapt_compact_sections, label_holding_currencies
from notifications.events import RenderedNotification
from quant_platform_kit.common.notification_localization import (
localize_notification_text as _base_localize_notification_text,
Expand Down Expand Up @@ -150,7 +150,13 @@ def _build_risk_control_lines(execution, *, translator):
)


def _format_dashboard_text(text, *, translator=None, cash_only_execution: bool = True) -> str:
def _format_dashboard_text(
text,
*,
translator=None,
cash_only_execution: bool = True,
position_currency_by_symbol=None,
) -> str:
lines = []
for raw_line in str(text or "").splitlines():
line = raw_line.rstrip()
Expand All @@ -160,6 +166,12 @@ def _format_dashboard_text(text, *, translator=None, cash_only_execution: bool =
line = _localize_notification_text(line, translator=translator)
lines.append(line)
result = "\n".join(lines)
if translator is not None:
result = label_holding_currencies(
result,
position_currency_by_symbol=position_currency_by_symbol,
unknown_currency_label=translator("holding_currency_unverified"),
)
if translator is not None:
result = _relabel_dashboard_cash_labels_shared(
result,
Expand All @@ -175,6 +187,7 @@ def _append_dashboard_block(lines, *, execution, separator, translator, compact:
execution.get("dashboard_text"),
translator=translator,
cash_only_execution=cash_only_execution,
position_currency_by_symbol=execution.get("position_currency_by_symbol"),
)
dashboard_lines = [
line for line in dashboard_text.splitlines()
Expand Down Expand Up @@ -316,6 +329,7 @@ def _compact_total_assets_line(execution, *, translator) -> str:
execution.get("dashboard_text"),
translator=translator,
cash_only_execution=bool(execution.get("cash_only_execution", True)),
position_currency_by_symbol=execution.get("position_currency_by_symbol"),
)
labels = (
"总资产",
Expand Down Expand Up @@ -397,6 +411,7 @@ def render_rebalance_notification(
execution.get("dashboard_text"),
translator=translator,
cash_only_execution=bool(execution.get("cash_only_execution", True)),
position_currency_by_symbol=execution.get("position_currency_by_symbol"),
)
compact_lines.extend(
adapt_compact_sections(
Expand Down Expand Up @@ -467,10 +482,11 @@ def render_heartbeat_notification(
translator=translator,
signal_key="heartbeat_signal",
)
outcome_key = _heartbeat_outcome_key(execution, skip_logs=skip_logs, note_logs=note_logs)
detailed_lines.extend(
[
separator,
translator("no_executable_orders") if (skip_logs or note_logs) else translator("no_trades"),
translator(outcome_key),
]
)
detailed_text = "\n".join(detailed_lines)
Expand Down Expand Up @@ -500,6 +516,7 @@ def render_heartbeat_notification(
execution.get("dashboard_text"),
translator=translator,
cash_only_execution=bool(execution.get("cash_only_execution", True)),
position_currency_by_symbol=execution.get("position_currency_by_symbol"),
)
compact_lines.extend(
adapt_compact_sections(
Expand All @@ -508,11 +525,29 @@ def render_heartbeat_notification(
supplemental_lines=execution.get("compact_supplemental_lines", ()),
)
)
compact_lines.append(
translator("no_executable_orders") if (skip_logs or note_logs) else translator("no_trades")
)
compact_lines.append(translator(outcome_key))

return RenderedNotification(
detailed_text=detailed_text,
compact_text="\n".join(compact_lines),
)


def _heartbeat_outcome_key(execution, *, skip_logs, note_logs) -> str:
state = str(execution.get("heartbeat_execution_state") or "").strip().lower()
durable = execution.get("durable_live_execution_command")
durable = durable if isinstance(durable, Mapping) else {}
is_hard_blocked = (
state == "blocked"
or str(execution.get("execution_status") or "").strip().lower() == "blocked"
or bool(execution.get("account_identity_blocked"))
or bool(execution.get("live_command_blocked"))
or str(durable.get("status") or "").strip().upper() == "BLOCKED_INVALID_BINDING"
)
if is_hard_blocked:
return "heartbeat_execution_blocked"
if state == "waiting_window":
return "heartbeat_waiting_window"
if bool(execution.get("direct_live_routing_blocked")):
return "heartbeat_execution_blocked"
return "no_executable_orders" if (skip_logs or note_logs) else "no_trades"
12 changes: 12 additions & 0 deletions notifications/telegram.py
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,12 @@ def _break_telegram_market_symbol_auto_links(value) -> str:
"heartbeat_account_equity": "💰 账户总权益: {value}",
"heartbeat_account_observed": "资金快照(本轮调仓前): {value}",
"heartbeat_unverified": "未核实",
"heartbeat_execution_blocked": "⚠️ 执行受阻,本轮未提交新订单",
"heartbeat_waiting_window": "⏳ 等待执行时段",
"holding_currency_unverified": "币种未核实",
"issue_durable_execution_blocked_title": "🛑 持久化实盘指令被阻断",
"issue_durable_execution_binding_invalid": "持久化实盘指令绑定无效;已阻止券商订单",
"issue_durable_execution_unresolved": "持久化实盘指令仍未解决;已阻止券商订单",
"precheck_title": "🧪 【策略演练】",
"dry_run_title": "🧪 【策略演练】",
"health_probe_title": "🔎 【连接探针】",
Expand Down Expand Up @@ -282,6 +288,12 @@ def _break_telegram_market_symbol_auto_links(value) -> str:
"heartbeat_account_equity": "💰 Total account equity: {value}",
"heartbeat_account_observed": "Account snapshot (before this cycle's rebalance): {value}",
"heartbeat_unverified": "Unverified",
"heartbeat_execution_blocked": "⚠️ Execution blocked; no new orders submitted this cycle",
"heartbeat_waiting_window": "⏳ Waiting for the execution window",
"holding_currency_unverified": "Currency unverified",
"issue_durable_execution_blocked_title": "Durable live execution blocked",
"issue_durable_execution_binding_invalid": "Durable live execution command binding is invalid; broker orders blocked",
"issue_durable_execution_unresolved": "Durable live execution command is unresolved; broker orders blocked",
"precheck_title": "🧪 【Strategy Dry Run】",
"dry_run_title": "🧪 【Strategy Dry Run】",
"health_probe_title": "🔎 【Health Probe】",
Expand Down
40 changes: 37 additions & 3 deletions tests/test_longbridge_local_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,12 @@ def quote(self, symbols):


class FakePosition:
def __init__(self, symbol, quantity, available_quantity=None):
def __init__(self, symbol, quantity, available_quantity=None, currency=None):
self.symbol = symbol
self.quantity = quantity
self.available_quantity = available_quantity if available_quantity is not None else quantity
if currency is not None:
self.currency = currency


class FakeChannel:
Expand All @@ -43,8 +45,8 @@ def __init__(self, positions):


class FakePositionsResponse:
def __init__(self):
self.channels = [FakeChannel([FakePosition("SOXL.US", 3), FakePosition("QQQI.US", 2, 1)])]
def __init__(self, positions=None):
self.channels = [FakeChannel(positions or [FakePosition("SOXL.US", 3), FakePosition("QQQI.US", 2, 1)])]


class LongBridgeLocalHelpersTests(unittest.TestCase):
Expand Down Expand Up @@ -110,6 +112,38 @@ def test_heartbeat_labels_cash_and_account_equity_in_their_own_currencies(self):
self.assertEqual(snapshot["equity_currency"], "SGD")
self.assertIn("observed_at", snapshot)

def test_position_currency_projection_uses_only_consistent_native_position_currency(self):
balance = types.SimpleNamespace(currency="SGD", net_assets="2500.50", cash_infos=[])
positions = FakePositionsResponse([
FakePosition("SOXL.US", 1, currency="USD"),
FakePosition("00700.HK", 1, currency="HKD"),
FakePosition("QQQI.US", 1),
])
trade = types.SimpleNamespace(
account_balance=lambda: [balance], stock_positions=lambda: positions,
)
state = fetch_strategy_account_state(
FakeQuoteContext(), trade, ["SOXL", "00700", "QQQI"],
)
self.assertEqual(
state["position_currency_by_symbol"],
{"SOXL": "USD", "00700": "HKD", "QQQI": None},
)
self.assertEqual(state["market_values"], {"SOXL": 50.0, "00700": 320.0, "QQQI": 20.0})

def test_conflicting_position_currencies_withhold_symbol_currency(self):
balance = types.SimpleNamespace(currency="SGD", net_assets="2500.50", cash_infos=[])
positions = FakePositionsResponse([
FakePosition("SOXL.US", 1, currency="USD"),
FakePosition("SOXL.US", 1, currency="SGD"),
])
trade = types.SimpleNamespace(
account_balance=lambda: [balance], stock_positions=lambda: positions,
)
state = fetch_strategy_account_state(FakeQuoteContext(), trade, ["SOXL"])
self.assertEqual(state["position_currency_by_symbol"], {"SOXL": None})
self.assertEqual(state["market_values"]["SOXL"], 100.0)

def test_fetch_strategy_account_state_rejects_account_balance_failure(self):
class BalanceFailingTradeContext:
def account_balance(self):
Expand Down
Loading
Loading