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
36 changes: 35 additions & 1 deletion application/rebalance_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

from __future__ import annotations

from collections.abc import Mapping
import json
import re

Expand Down Expand Up @@ -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,
Expand All @@ -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 = []
Expand All @@ -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])):
Expand Down Expand Up @@ -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="<unknown>")
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)
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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')}"
Expand All @@ -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)
Expand Down
25 changes: 25 additions & 0 deletions decision_mapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
*,
Expand Down Expand Up @@ -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
2 changes: 1 addition & 1 deletion 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@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
Expand Down
113 changes: 100 additions & 13 deletions strategy_runtime.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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 (
Expand Down Expand Up @@ -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,
*,
Expand Down Expand Up @@ -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,
Expand All @@ -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)
Expand Down Expand Up @@ -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))
Expand All @@ -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:
Expand Down Expand Up @@ -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=(),
Expand Down
3 changes: 3 additions & 0 deletions tests/test_decision_mapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"},
},
)

Expand All @@ -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


Expand Down
13 changes: 13 additions & 0 deletions tests/test_rebalance_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",),
Expand All @@ -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]

Expand Down
Loading