From a2de3ffb026ff262f633a59db811b3b2b46b1759 Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Tue, 21 Apr 2026 05:37:02 +0800 Subject: [PATCH] Render strategy dashboard in IBKR notifications --- application/rebalance_service.py | 36 +++++++++- decision_mapper.py | 25 +++++++ requirements.txt | 2 +- strategy_runtime.py | 113 +++++++++++++++++++++++++++---- tests/test_decision_mapper.py | 3 + tests/test_rebalance_service.py | 13 ++++ tests/test_strategy_runtime.py | 13 +++- 7 files changed, 188 insertions(+), 17 deletions(-) diff --git a/application/rebalance_service.py b/application/rebalance_service.py index c88afd5..c33a931 100644 --- a/application/rebalance_service.py +++ b/application/rebalance_service.py @@ -2,6 +2,7 @@ from __future__ import annotations +from collections.abc import Mapping import json import re @@ -297,6 +298,27 @@ def _resolve_weight_allocation(signal_metadata, *, required: bool) -> dict: } +def _format_dashboard_text(text) -> str: + lines = [line.rstrip() for line in str(text or "").strip().splitlines()] + while lines and not lines[0].strip(): + lines.pop(0) + while lines and not lines[-1].strip(): + lines.pop() + return "\n".join(lines) + + +def _strategy_dashboard_text(signal_metadata) -> str: + metadata = signal_metadata if isinstance(signal_metadata, Mapping) else {} + raw_annotations = metadata.get("execution_annotations") + annotations = raw_annotations if isinstance(raw_annotations, Mapping) else {} + return _format_dashboard_text( + annotations.get("dashboard_text") + or metadata.get("dashboard_text") + or metadata.get("dashboard") + or "" + ) + + def build_dashboard( positions, account_values, @@ -311,6 +333,10 @@ def build_dashboard( separator, status_icon="🐤", ): + signal_metadata = signal_metadata or {} + strategy_dashboard = _strategy_dashboard_text(signal_metadata) + if strategy_dashboard: + return strategy_dashboard equity = account_values.get("equity", 0) buying_power = account_values.get("buying_power", 0) position_lines = [] @@ -320,7 +346,6 @@ def build_dashboard( market_value = qty * avg position_lines.append(f" - {symbol}: {qty}股 | ${market_value:,.2f}") position_text = "\n".join(position_lines) if position_lines else translator("empty_positions") - signal_metadata = signal_metadata or {} allocation = _resolve_weight_allocation(signal_metadata, required=False) target_lines = [] for symbol, weight in sorted(allocation.get("targets", {}).items(), key=lambda item: (-item[1], item[0])): @@ -395,10 +420,15 @@ def _build_compact_message( translator, separator: str, body_lines, + dashboard_text: str = "", ) -> str: lines = [title] strategy_name = _format_text(strategy_display_name, fallback="") lines.append(translator("strategy_label", name=strategy_name)) + dashboard = _format_dashboard_text(dashboard_text) + if dashboard: + lines.append(separator) + lines.extend(dashboard.splitlines()) status_line = _first_prefixed_line(status_icon, status_desc, translator=translator) if status_line: lines.append(status_line) @@ -452,6 +482,7 @@ def run_strategy_core( separator=separator, status_icon=signal_metadata.get("status_icon", "🐤"), ) + strategy_dashboard = _strategy_dashboard_text(signal_metadata) if target_weights is None: decision = signal_metadata.get("snapshot_guard_decision") @@ -491,6 +522,7 @@ def run_strategy_core( translator=translator, separator=separator, body_lines=[no_op_text], + dashboard_text=strategy_dashboard, ) print(detailed_message, flush=True) send_tg_message(compact_message) @@ -566,6 +598,7 @@ def run_strategy_core( translator=translator, separator=separator, body_lines=notification_trade_lines, + dashboard_text=strategy_dashboard, ) else: detailed_message = f"{translator('heartbeat_title')}\n{dashboard}\n{separator}\n{translator('no_trades')}" @@ -578,6 +611,7 @@ def run_strategy_core( translator=translator, separator=separator, body_lines=[translator("no_trades")], + dashboard_text=strategy_dashboard, ) print(detailed_message, flush=True) diff --git a/decision_mapper.py b/decision_mapper.py index 356eb88..e5b0888 100644 --- a/decision_mapper.py +++ b/decision_mapper.py @@ -112,6 +112,20 @@ def _derive_status_description( return _derive_signal_description(decision, runtime_metadata) +def _derive_execution_annotations( + diagnostics: Mapping[str, Any], + runtime_metadata: Mapping[str, Any], +) -> dict[str, Any]: + annotations: dict[str, Any] = {} + raw_runtime_annotations = runtime_metadata.get("execution_annotations") + if isinstance(raw_runtime_annotations, Mapping): + annotations.update(raw_runtime_annotations) + raw_diagnostic_annotations = diagnostics.get("execution_annotations") + if isinstance(raw_diagnostic_annotations, Mapping): + annotations.update(raw_diagnostic_annotations) + return annotations + + def map_strategy_decision( decision: StrategyDecision, *, @@ -152,5 +166,16 @@ def map_strategy_decision( metadata.setdefault("actionable", not no_execute) if allocation_payload: metadata.setdefault("allocation", allocation_payload) + execution_annotations = _derive_execution_annotations(diagnostics, runtime_metadata) + if execution_annotations: + metadata.setdefault("execution_annotations", execution_annotations) + dashboard_text = str( + execution_annotations.get("dashboard_text") + or diagnostics.get("dashboard") + or metadata.get("dashboard_text") + or "" + ).strip() + if dashboard_text: + metadata.setdefault("dashboard_text", dashboard_text) return target_weights, signal_desc, is_emergency, status_desc, metadata diff --git a/requirements.txt b/requirements.txt index d950166..0b4163f 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,7 +1,7 @@ flask gunicorn quant-platform-kit @ git+https://github.com/QuantStrategyLab/QuantPlatformKit.git@v0.7.19 -us-equity-strategies @ git+https://github.com/QuantStrategyLab/UsEquityStrategies.git@main +us-equity-strategies @ git+https://github.com/QuantStrategyLab/UsEquityStrategies.git@v0.7.31 pandas numpy requests diff --git a/strategy_runtime.py b/strategy_runtime.py index bd3a3b7..f79d916 100644 --- a/strategy_runtime.py +++ b/strategy_runtime.py @@ -1,7 +1,7 @@ from __future__ import annotations from collections.abc import Callable, Mapping -from dataclasses import dataclass, field +from dataclasses import dataclass, field, replace from typing import Any import pandas as pd @@ -23,6 +23,8 @@ StrategyDecision, StrategyEntrypoint, StrategyRuntimeAdapter, + build_strategy_context_from_available_inputs, + build_strategy_evaluation_inputs, ) from runtime_config_support import PlatformRuntimeSettings from strategy_loader import ( @@ -65,6 +67,68 @@ def profile(self) -> str: def required_inputs(self) -> frozenset[str]: return frozenset(self.entrypoint.manifest.required_inputs) + def _runtime_adapter_with_portfolio( + self, + runtime_adapter: StrategyRuntimeAdapter, + portfolio_snapshot: Any | None, + ) -> StrategyRuntimeAdapter: + if portfolio_snapshot is None or runtime_adapter.portfolio_input_name: + return runtime_adapter + available_inputs = set(runtime_adapter.available_inputs or self.required_inputs) + available_inputs.update(self.required_inputs) + available_inputs.add(_PORTFOLIO_SNAPSHOT_INPUT) + return replace( + runtime_adapter, + available_inputs=frozenset(available_inputs), + portfolio_input_name=_PORTFOLIO_SNAPSHOT_INPUT, + ) + + def _fetch_portfolio_snapshot_for_context(self, ib, *, required: bool) -> Any | None: + if ib is None and not required: + return None + if required: + return fetch_portfolio_snapshot(ib) + try: + return fetch_portfolio_snapshot(ib) + except Exception as exc: + self.logger( + "strategy_dashboard_portfolio_snapshot_failed | " + f"profile={self.profile} error_type={type(exc).__name__} error={exc}" + ) + return None + + def _build_strategy_context( + self, + *, + runtime_adapter: StrategyRuntimeAdapter, + as_of: pd.Timestamp, + market_inputs: Mapping[str, Any], + portfolio_snapshot: Any | None, + runtime_config: Mapping[str, Any], + current_holdings, + ib, + ): + context_adapter = self._runtime_adapter_with_portfolio(runtime_adapter, portfolio_snapshot) + available_inputs = set(context_adapter.available_inputs or self.required_inputs) + available_inputs.update(self.required_inputs) + evaluation_inputs = build_strategy_evaluation_inputs( + available_inputs=available_inputs, + market_inputs=market_inputs, + portfolio_snapshot=portfolio_snapshot, + ) + capabilities = {} + if ib is not None: + capabilities["broker_client"] = ib + return build_strategy_context_from_available_inputs( + entrypoint=self.entrypoint, + runtime_adapter=context_adapter, + as_of=as_of, + available_inputs=evaluation_inputs, + runtime_config=runtime_config, + state={"current_holdings": tuple(current_holdings)}, + capabilities=capabilities, + ) + def evaluate( self, *, @@ -129,11 +193,12 @@ def _evaluate_market_data_strategy( runtime_config = dict(self.runtime_config) runtime_config.setdefault("translator", translator) runtime_config.setdefault("pacing_sec", float(pacing_sec)) - ctx = build_ibkr_strategy_context( - entrypoint=self.entrypoint, + portfolio_snapshot = self._fetch_portfolio_snapshot_for_context(ib, required=False) + ctx = self._build_strategy_context( runtime_adapter=self.runtime_adapter, as_of=run_as_of, market_inputs=build_market_history_inputs(historical_close_loader), + portfolio_snapshot=portfolio_snapshot, runtime_config=runtime_config, current_holdings=current_holdings, ib=ib, @@ -151,6 +216,8 @@ def _evaluate_market_data_strategy( "status_icon": self.status_icon, "dry_run_only": self.runtime_settings.dry_run_only, } + if portfolio_snapshot is not None: + metadata["portfolio_total_equity"] = float(getattr(portfolio_snapshot, "total_equity", 0.0) or 0.0) if safe_haven_symbol: metadata["safe_haven_symbol"] = safe_haven_symbol return StrategyEvaluationResult(decision=decision, metadata=metadata) @@ -248,14 +315,24 @@ def _evaluate_feature_snapshot_strategy( translator: Callable[[str], str], pacing_sec: float, ) -> StrategyEvaluationResult: - del translator, pacing_sec + del pacing_sec runtime_config_path = self.merged_runtime_config.get("runtime_config_path") or self.runtime_settings.strategy_config_path benchmark_symbol = str(self.merged_runtime_config.get("benchmark_symbol") or "SPY").strip().upper() portfolio_snapshot_holder: dict[str, Any] = {} + runtime_config = dict(self.runtime_config) + runtime_config.setdefault("translator", translator) def build_available_inputs(feature_snapshot) -> Mapping[str, Any]: - if _PORTFOLIO_SNAPSHOT_INPUT in self.required_inputs: - portfolio_snapshot_holder["portfolio_snapshot"] = fetch_portfolio_snapshot(ib) + requires_portfolio = ( + _PORTFOLIO_SNAPSHOT_INPUT in self.required_inputs + or self.runtime_adapter.portfolio_input_name == _PORTFOLIO_SNAPSHOT_INPUT + ) + portfolio_snapshot = self._fetch_portfolio_snapshot_for_context( + ib, + required=requires_portfolio, + ) + if portfolio_snapshot is not None: + portfolio_snapshot_holder["portfolio_snapshot"] = portfolio_snapshot market_inputs: dict[str, Any] = {_FEATURE_SNAPSHOT_INPUT: feature_snapshot} if _MARKET_HISTORY_INPUT in self.required_inputs: market_inputs.update(build_market_history_inputs(historical_close_loader)) @@ -274,15 +351,25 @@ def build_available_inputs(feature_snapshot) -> Mapping[str, Any]: return market_inputs def build_context(request: FeatureSnapshotContextRequest): - return build_ibkr_strategy_context( + portfolio_snapshot = portfolio_snapshot_holder.get("portfolio_snapshot") + runtime_adapter = self._runtime_adapter_with_portfolio( + request.runtime_adapter, + portfolio_snapshot, + ) + available_inputs = dict(request.available_inputs) + if portfolio_snapshot is not None: + available_inputs[_PORTFOLIO_SNAPSHOT_INPUT] = portfolio_snapshot + capabilities = {} + if ib is not None: + capabilities["broker_client"] = ib + return build_strategy_context_from_available_inputs( entrypoint=request.entrypoint, - runtime_adapter=request.runtime_adapter, + runtime_adapter=runtime_adapter, as_of=request.as_of, - market_inputs=request.available_inputs, - portfolio_snapshot=portfolio_snapshot_holder.get("portfolio_snapshot"), + available_inputs=available_inputs, runtime_config=request.runtime_config, - current_holdings=current_holdings, - ib=ib, + state={"current_holdings": tuple(current_holdings)}, + capabilities=capabilities, ) def log_guard_metadata(guard_metadata: Mapping[str, Any]) -> None: @@ -323,7 +410,7 @@ def build_extra_metadata( strategy_config_source=self.runtime_settings.strategy_config_source, dry_run_only=self.runtime_settings.dry_run_only, ), - runtime_config=dict(self.runtime_config), + runtime_config=runtime_config, merged_runtime_config=self.merged_runtime_config, as_of=run_as_of, base_managed_symbols=(), diff --git a/tests/test_decision_mapper.py b/tests/test_decision_mapper.py index d3052a8..1977725 100644 --- a/tests/test_decision_mapper.py +++ b/tests/test_decision_mapper.py @@ -12,6 +12,7 @@ def test_map_strategy_decision_maps_weight_positions_and_safe_haven(): diagnostics={ "signal_description": "risk on", "status_description": "breadth=60.0%", + "execution_annotations": {"dashboard_text": "strategy dashboard"}, }, ) @@ -32,6 +33,8 @@ def test_map_strategy_decision_maps_weight_positions_and_safe_haven(): assert metadata["allocation"]["strategy_symbols"] == ("AAA", "BOXX") assert metadata["allocation"]["targets"] == {"AAA": 0.6, "BOXX": 0.4} assert metadata["allocation"]["positions"][1]["role"] == "safe_haven" + assert metadata["execution_annotations"]["dashboard_text"] == "strategy dashboard" + assert metadata["dashboard_text"] == "strategy dashboard" assert "target_mode" not in metadata diff --git a/tests/test_rebalance_service.py b/tests/test_rebalance_service.py index 58e019e..f2bcbd6 100644 --- a/tests/test_rebalance_service.py +++ b/tests/test_rebalance_service.py @@ -202,6 +202,16 @@ def fake_execute_rebalance( { "managed_symbols": ("AAA", "BOXX"), "status_icon": "📏", + "execution_annotations": { + "dashboard_text": ( + "📌 Strategy portfolio\n" + " - Total assets (strategy symbols + cash): $1,000.00\n" + " - Buying power: $500.00\n" + "💼 Strategy holdings\n" + " - AAA: $0.00 / 0 shares\n" + " - BOXX: $0.00 / 0 shares" + ) + }, "allocation": _weight_allocation( {"AAA": 0.9, "BOXX": 0.1}, risk_symbols=("AAA",), @@ -223,6 +233,9 @@ def fake_execute_rebalance( assert "Account Summary" not in observed["messages"][0] assert "Current Positions" not in observed["messages"][0] assert "Execution Summary" not in observed["messages"][0] + assert "📌 Strategy portfolio" in observed["messages"][0] + assert "Total assets (strategy symbols + cash): $1,000.00" in observed["messages"][0] + assert "💼 Strategy holdings" in observed["messages"][0] assert "📏 breadth=60.0%" in observed["messages"][0] assert "Target Weights" not in observed["messages"][0] diff --git a/tests/test_strategy_runtime.py b/tests/test_strategy_runtime.py index 5678413..377cc5c 100644 --- a/tests/test_strategy_runtime.py +++ b/tests/test_strategy_runtime.py @@ -172,6 +172,7 @@ class FakeEntrypoint: def evaluate(self, ctx): captured["market_data"] = dict(ctx.market_data) + captured["portfolio"] = ctx.portfolio captured["runtime_config"] = dict(ctx.runtime_config) return StrategyDecision() @@ -206,9 +207,11 @@ def fake_guard(path, **kwargs): ) monkeypatch.setattr(strategy_runtime_module, "load_feature_snapshot_guarded", fake_guard) + portfolio_snapshot = SimpleNamespace(total_equity=25000.0) + monkeypatch.setattr(strategy_runtime_module, "fetch_portfolio_snapshot", lambda _ib: portfolio_snapshot) result = runtime.evaluate( - ib=None, + ib="fake-ib", current_holdings={"AAPL"}, historical_close_loader=lambda *_args, **_kwargs: None, run_as_of=strategy_runtime_module.pd.Timestamp("2026-04-01"), @@ -222,6 +225,8 @@ def fake_guard(path, **kwargs): assert captured["require_manifest"] is True assert captured["expected_contract_version"] == "adapter.contract" assert captured["market_data"]["feature_snapshot"] == [{"as_of": "2026-03-31", "symbol": "AAPL", "close": 1.0}] + assert captured["portfolio"] is portfolio_snapshot + assert captured["runtime_config"]["translator"]("equity") == "equity" assert "pacing_sec" not in captured["runtime_config"] assert result.metadata["managed_symbols"] == ("AAPL", "BOXX") @@ -314,7 +319,7 @@ def candle_loader(_ib, symbol, duration="2 Y", bar_size="1 day"): assert result.metadata["managed_symbols"] == ("NVDL", "BOXX") -def test_market_history_runtime_uses_canonical_market_history_key(): +def test_market_history_runtime_uses_canonical_market_history_key(monkeypatch): captured = {} class FakeEntrypoint: @@ -329,6 +334,7 @@ class FakeEntrypoint: def evaluate(self, ctx): captured["market_data"] = dict(ctx.market_data) + captured["portfolio"] = ctx.portfolio return StrategyDecision() def loader(*_args, **_kwargs): @@ -343,6 +349,8 @@ def loader(*_args, **_kwargs): status_icon="🐤", logger=lambda _message: None, ) + portfolio_snapshot = SimpleNamespace(total_equity=1200.0) + monkeypatch.setattr(strategy_runtime_module, "fetch_portfolio_snapshot", lambda _ib: portfolio_snapshot) runtime.evaluate( ib="fake-ib", @@ -354,6 +362,7 @@ def loader(*_args, **_kwargs): ) assert captured["market_data"]["market_history"] is loader + assert captured["portfolio"] is portfolio_snapshot assert "historical_close_loader" not in captured["market_data"]