diff --git a/application/rebalance_service.py b/application/rebalance_service.py index e28beb2..c5e74b2 100644 --- a/application/rebalance_service.py +++ b/application/rebalance_service.py @@ -40,6 +40,12 @@ _DETAIL_FIELD_SPLIT_RE = re.compile(r"\s+(?=[^\s=::]+[=::])") DRY_RUN_BYPASS_EXECUTION_MARKER_ENV = "DRY_RUN_BYPASS_EXECUTION_MARKER" +_LIVE_COMMAND_BINDING_ERRORS = frozenset( + { + "invalid live execution command", + "live execution command strategy release is invalid", + } +) def _env_flag_enabled(name: str) -> bool: @@ -613,8 +619,42 @@ def fetch_replanned_state(): if not getattr(config, "durable_live_execution_session_authorized", False): raise RuntimeError("durable live execution requires an open exchange session") matching_commands = _matching_live_commands(config=config) + terminal_states = { + ExecutionCommandState.FILLED, + ExecutionCommandState.CANCELLED, + ExecutionCommandState.REJECTED, + } + unresolved_matching_commands = [] for command in matching_commands: - _validate_live_command_binding(command=command, config=config) + if config.execution_command_store.current_state(command) in terminal_states: + continue + try: + _validate_live_command_binding(command=command, config=config) + 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) + 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, + }, + }, + allocation={}, + logs=(), + skip_logs=(), + note_logs=(message,), + action_done=False, + ) + unresolved_matching_commands.append(command) + for command in unresolved_matching_commands: _reconcile_live_command( command=command, store=config.execution_command_store, trade_context=trade_context, fetch_order_status=runtime.fetch_order_status, diff --git a/main.py b/main.py index 61d772e..843b7e1 100644 --- a/main.py +++ b/main.py @@ -387,6 +387,13 @@ def _summarize_cycle_result_for_report(cycle_result, *, dry_run: bool) -> dict: "quotes": [dict(snapshot) for snapshot in quote_snapshots], } execution = dict(getattr(cycle_result, "execution", {}) or {}) + if execution.get("execution_status") == "blocked": + summary["execution_status"] = "blocked" + summary["blocked_reason"] = ( + "durable_live_execution_command_binding_invalid" + if execution.get("blocked_reason") == "durable_live_execution_command_binding_invalid" + else "unknown" + ) for field in ( "signal_date", "effective_date", diff --git a/tests/test_rebalance_service.py b/tests/test_rebalance_service.py index c51e7a2..2821a2a 100644 --- a/tests/test_rebalance_service.py +++ b/tests/test_rebalance_service.py @@ -1700,6 +1700,178 @@ def frozen_plan(*, allocation, execution, snapshot): self.assertEqual(command_store.current_state(read_error_command), ExecutionCommandState.CLAIMED) self.assertEqual(len(orders), 2) + def test_terminal_stale_live_command_is_ignored_before_current_cycle(self): + from application.durable_execution_commands import build_live_execution_command + + plan = _build_plan( + strategy_symbols=("SOXL",), risk_symbols=("SOXL",), + targets={"SOXL": 400.0}, market_values={"SOXL": 0.0}, + sellable_quantities={"SOXL": 0}, quantities={"SOXL": 0}, + current_min_trade=10.0, trade_threshold_value=10.0, + investable_cash=500.0, available_cash=500.0, total_strategy_equity=500.0, + market_status="Risk on", deploy_ratio_text="70.0%", income_ratio_text="0.0%", + income_locked_ratio_text="0.0%", signal_message="SOXL target", + portfolio_rows=(("SOXL",),), signal_date="2026-04-21", effective_date="2026-04-22", + ) + command_store = ExecutionCommandStore(local_dir=self.enterContext(TemporaryDirectory())) + stale = build_live_execution_command( + platform="longbridge", account_scope="SG", strategy_profile="soxl_soxx_trend_income", + physical_account_id="lb-sg-001", runtime_identity_digest="b" * 64, + execution={**plan["execution"], "signal_date": "2026-04-20", "effective_date": "2026-04-21"}, + allocation=plan["allocation"], + ) + self.assertTrue(command_store.enqueue(stale)) + command_store.append_event(stale, next_state=ExecutionCommandState.CANCELLED) + orders = [] + runtime = LongBridgeRebalanceRuntime( + bootstrap=lambda: ("quote", "trade", {"trend": "ok"}), + resolve_rebalance_plan=lambda **_kwargs: plan, + market_data_port_factory=lambda _context: CallableMarketDataPort( + quote_loader=lambda symbol: QuoteSnapshot(symbol=symbol, as_of="2026-04-21", last_price=100.0)), + estimate_max_purchase_quantity=lambda *_args, **_kwargs: 5, + notifications=CallableNotificationPort(lambda _message: None), + notify_issue=lambda *_args: None, + portfolio_port_factory=lambda *_contexts: CallablePortfolioPort(lambda: _build_snapshot(plan)), + execution_port_factory=lambda _context: CallableExecutionPort(lambda order: orders.append(order)), + resolve_frozen_rebalance_plan=lambda **_kwargs: plan, + ) + config = LongBridgeRebalanceConfig( + limit_sell_discount=0.995, limit_buy_premium=1.005, separator="-", + translator=build_translator("en"), with_prefix=lambda message: message, + strategy_profile="soxl_soxx_trend_income", execution_state_account_scope="SG", + physical_account_id="lb-sg-001", dry_run_only=False, notify_no_trade_cycles=False, + durable_execution_command_live_enabled=True, execution_command_store=command_store, + durable_live_execution_session_authorized=True, + durable_execution_runtime_identity_digest="a" * 64, + ) + + result = rebalance_service.run_strategy(runtime=runtime, config=config) + + self.assertFalse(result.action_done) + self.assertEqual(orders, []) + self.assertIs(command_store.current_state(stale), ExecutionCommandState.CANCELLED) + self.assertEqual(len(rebalance_service.list_live_execution_commands(command_store)), 2) + + + def test_same_day_terminal_live_command_prevents_duplicate_command_or_submit(self): + from application.durable_execution_commands import build_live_execution_command + + plan = _build_plan( + strategy_symbols=("SOXL",), risk_symbols=("SOXL",), + targets={"SOXL": 400.0}, market_values={"SOXL": 0.0}, + sellable_quantities={"SOXL": 0}, quantities={"SOXL": 0}, + current_min_trade=10.0, trade_threshold_value=10.0, + investable_cash=500.0, available_cash=500.0, total_strategy_equity=500.0, + market_status="Risk on", deploy_ratio_text="70.0%", income_ratio_text="0.0%", + income_locked_ratio_text="0.0%", signal_message="SOXL target", + portfolio_rows=(("SOXL",),), signal_date="2026-04-21", effective_date="2026-04-22", + ) + command_store = ExecutionCommandStore(local_dir=self.enterContext(TemporaryDirectory())) + terminal = build_live_execution_command( + platform="longbridge", account_scope="SG", strategy_profile="soxl_soxx_trend_income", + physical_account_id="lb-sg-001", runtime_identity_digest="b" * 64, + execution=plan["execution"], allocation=plan["allocation"], + ) + self.assertTrue(command_store.enqueue(terminal)) + command_store.append_event(terminal, next_state=ExecutionCommandState.CANCELLED) + orders = [] + runtime = LongBridgeRebalanceRuntime( + bootstrap=lambda: ("quote", "trade", {"trend": "ok"}), + resolve_rebalance_plan=lambda **_kwargs: plan, + market_data_port_factory=lambda _context: CallableMarketDataPort( + quote_loader=lambda symbol: QuoteSnapshot(symbol=symbol, as_of="2026-04-21", last_price=100.0)), + estimate_max_purchase_quantity=lambda *_args, **_kwargs: 5, + notifications=CallableNotificationPort(lambda _message: None), + notify_issue=lambda *_args: None, + portfolio_port_factory=lambda *_contexts: CallablePortfolioPort(lambda: _build_snapshot(plan)), + execution_port_factory=lambda _context: CallableExecutionPort(lambda order: orders.append(order)), + resolve_frozen_rebalance_plan=lambda **_kwargs: plan, + ) + config = LongBridgeRebalanceConfig( + limit_sell_discount=0.995, limit_buy_premium=1.005, separator="-", + translator=build_translator("en"), with_prefix=lambda message: message, + strategy_profile="soxl_soxx_trend_income", execution_state_account_scope="SG", + physical_account_id="lb-sg-001", dry_run_only=False, notify_no_trade_cycles=False, + durable_execution_command_live_enabled=True, execution_command_store=command_store, + durable_live_execution_session_authorized=True, + durable_execution_runtime_identity_digest="a" * 64, + ) + + result = rebalance_service.run_strategy(runtime=runtime, config=config) + + self.assertFalse(result.action_done) + self.assertTrue(result.execution["direct_live_routing_blocked"]) + self.assertEqual(orders, []) + self.assertEqual(rebalance_service.list_live_execution_commands(command_store), (terminal,)) + self.assertIs(command_store.current_state(terminal), ExecutionCommandState.CANCELLED) + + def test_queued_stale_live_command_blocks_without_store_or_broker_writes(self): + from application.durable_execution_commands import build_live_execution_command + + plan = _build_plan( + strategy_symbols=("SOXL",), risk_symbols=("SOXL",), + targets={"SOXL": 400.0}, market_values={"SOXL": 0.0}, + sellable_quantities={"SOXL": 0}, quantities={"SOXL": 0}, + current_min_trade=10.0, trade_threshold_value=10.0, + investable_cash=500.0, available_cash=500.0, total_strategy_equity=500.0, + market_status="Risk on", deploy_ratio_text="70.0%", income_ratio_text="0.0%", + income_locked_ratio_text="0.0%", signal_message="SOXL target", + portfolio_rows=(("SOXL",),), signal_date="2026-04-21", effective_date="2026-04-22", + ) + command_store = ExecutionCommandStore(local_dir=self.enterContext(TemporaryDirectory())) + stale = build_live_execution_command( + platform="longbridge", account_scope="SG", strategy_profile="soxl_soxx_trend_income", + physical_account_id="lb-sg-001", runtime_identity_digest="b" * 64, + execution=plan["execution"], allocation=plan["allocation"], + ) + valid = build_live_execution_command( + platform="longbridge", account_scope="SG", strategy_profile="soxl_soxx_trend_income", + physical_account_id="lb-sg-001", runtime_identity_digest="a" * 64, + execution=plan["execution"], allocation=plan["allocation"], + ) + self.assertTrue(command_store.enqueue(valid)) + self.assertTrue(command_store.enqueue(stale)) + issues = [] + runtime = LongBridgeRebalanceRuntime( + bootstrap=lambda: ("quote", "trade", {"trend": "ok"}), + resolve_rebalance_plan=lambda **_kwargs: (_ for _ in ()).throw(AssertionError("must not evaluate")), + market_data_port_factory=lambda _context: CallableMarketDataPort( + quote_loader=lambda _symbol: (_ for _ in ()).throw(AssertionError("must not load quote"))), + estimate_max_purchase_quantity=lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not estimate")), + notifications=CallableNotificationPort(lambda _message: None), + notify_issue=lambda title, detail: issues.append((title, detail)), + portfolio_port_factory=lambda *_contexts: CallablePortfolioPort( + lambda: (_ for _ in ()).throw(AssertionError("must not load snapshot"))), + execution_port_factory=lambda _context: CallableExecutionPort( + lambda _order: (_ for _ in ()).throw(AssertionError("must not submit"))), + ) + config = LongBridgeRebalanceConfig( + limit_sell_discount=0.995, limit_buy_premium=1.005, separator="-", + translator=build_translator("en"), with_prefix=lambda message: message, + strategy_profile="soxl_soxx_trend_income", execution_state_account_scope="SG", + physical_account_id="lb-sg-001", dry_run_only=False, notify_no_trade_cycles=False, + durable_execution_command_live_enabled=True, execution_command_store=command_store, + durable_live_execution_session_authorized=True, + durable_execution_runtime_identity_digest="a" * 64, + ) + + with patch.object(rebalance_service, "list_live_execution_commands", return_value=(valid, stale)), patch.object( + rebalance_service, "_reconcile_live_command" + ) as reconcile: + result = rebalance_service.run_strategy(runtime=runtime, config=config) + + self.assertFalse(result.action_done) + self.assertEqual(result.execution["execution_status"], "blocked") + self.assertEqual(result.execution["blocked_reason"], "durable_live_execution_command_binding_invalid") + self.assertEqual(result.execution["durable_live_execution_command"]["status"], "BLOCKED_INVALID_BINDING") + self.assertEqual(issues, [("Durable live execution blocked", "Durable live execution command binding is invalid; broker orders blocked")]) + reconcile.assert_not_called() + self.assertIs(command_store.current_state(valid), ExecutionCommandState.QUEUED) + self.assertIs(command_store.current_state(stale), ExecutionCommandState.QUEUED) + self.assertEqual(command_store.events(valid), ()) + self.assertEqual(command_store.events(stale), ()) + self.assertEqual(len(rebalance_service.list_live_execution_commands(command_store)), 2) + def test_live_order_detail_normalizes_real_sdk_enum_and_checks_identity(self): from application.longbridge_execution import fetch_live_order_status from longport.openapi import OrderStatus diff --git a/tests/test_request_handling.py b/tests/test_request_handling.py index 15ef1a9..b0691d3 100644 --- a/tests/test_request_handling.py +++ b/tests/test_request_handling.py @@ -1374,6 +1374,28 @@ def test_cycle_result_summary_keeps_broker_submission_pending_until_reconciled(s self.assertEqual(summary["orders_pending"][0]["broker_order_id"], "lb-order-pending") self.assertEqual(summary["durable_live_execution_command"]["status"], "QUEUED") + def test_cycle_result_summary_projects_fixed_live_command_binding_block(self): + module = load_module() + cycle_result = types.SimpleNamespace( + logs=(), skip_logs=(), note_logs=("blocked",), action_done=False, + execution={ + "execution_status": "blocked", + "blocked_reason": "durable_live_execution_command_binding_invalid", + "durable_live_execution_command": {"status": "BLOCKED_INVALID_BINDING"}, + }, + dry_run_orders=(), pending_orders=(), quote_snapshots=(), + ) + + summary = module._summarize_cycle_result_for_report(cycle_result, dry_run=False) + + self.assertEqual(summary["execution_status"], "blocked") + self.assertEqual(summary["blocked_reason"], "durable_live_execution_command_binding_invalid") + cycle_result.execution["blocked_reason"] = "SYNTHETIC_SECRET_DO_NOT_EMIT" + self.assertEqual( + module._summarize_cycle_result_for_report(cycle_result, dry_run=False)["blocked_reason"], + "unknown", + ) + def test_notification_delivery_log_summary_records_sent_dry_run_without_raw_text(self): module = load_module()