diff --git a/application/rebalance_service.py b/application/rebalance_service.py index 7287ea4..d5e27fb 100644 --- a/application/rebalance_service.py +++ b/application/rebalance_service.py @@ -39,7 +39,7 @@ resolve_strategy_run_period, ) from decision_mapper import map_strategy_decision_to_plan -from notifications.telegram import build_sender, build_translator, render_cycle_summary +from notifications.telegram import build_sender, build_translator, render_cycle_notification from quant_platform_kit.common.execution_outcomes import ( DEFAULT_EXECUTION_BLOCKING_SKIP_REASONS, filter_execution_blocking_skips, @@ -57,7 +57,7 @@ load_configured_strategy_plugin_signals, parse_strategy_plugin_mounts, ) -from quant_platform_kit.notifications.events import NotificationPublisher, RenderedNotification +from quant_platform_kit.notifications.events import NotificationPublisher from quant_platform_kit.notifications.strategy_plugin_alerts import ( StrategyPluginAlertStateSettings, build_strategy_plugin_alert_context_label as build_alert_context_label, @@ -238,7 +238,7 @@ def _publish_cycle_notification( if not settings.tg_token or not settings.tg_chat_id: return False sender = build_sender(settings.tg_token, settings.tg_chat_id) - message = render_cycle_summary(result, lang=settings.notify_lang) + notification = render_cycle_notification(result, lang=settings.notify_lang) def publish_log(text: str) -> None: try: log_message(text, flush=True) @@ -256,7 +256,7 @@ def send_and_capture(text: str) -> bool | None: NotificationPublisher( log_message=publish_log, send_message=send_and_capture, - ).publish(RenderedNotification(detailed_text=message, compact_text=message)) + ).publish(notification) return delivery_sent diff --git a/notifications/compact_adapter.py b/notifications/compact_adapter.py new file mode 100644 index 0000000..0e978a0 --- /dev/null +++ b/notifications/compact_adapter.py @@ -0,0 +1,87 @@ +"""Adapter for compact, user-facing notification sections.""" + +from __future__ import annotations + +import re +from collections.abc import Iterable + +_HOLDINGS_HEADERS = { + "💼 持仓", + "💼 策略持仓", + "💼 当前持仓", + "💼 Holdings", + "💼 Current Holdings", + "💼 Strategy Holdings", + "💼 Strategy holdings", +} +_NUMBER_RE = re.compile(r"[-+]?\d[\d,]*(?:\.\d+)?") + + +def _contains_nonzero_number(text: str) -> bool: + for match in _NUMBER_RE.finditer(text): + try: + if abs(float(match.group(0).replace(",", ""))) > 1e-12: + return True + except ValueError: + continue + return False + + +def _localize_holding_detail(detail: str, *, locale: str) -> str: + if str(locale).lower().startswith("zh"): + return re.sub(r"\s+shares?\b", "股", detail, flags=re.IGNORECASE) + + def replace_share(match: re.Match[str]) -> str: + quantity = match.group(1) + try: + unit = "share" if abs(float(quantity.replace(",", ""))) == 1 else "shares" + except ValueError: + unit = "shares" + return f"{quantity} {unit}" + + return re.sub(r"([-+]?\d[\d,]*(?:\.\d+)?)\s*股", replace_share, detail) + + +def adapt_compact_sections( + dashboard_text: str, + *, + locale: str, + supplemental_lines: Iterable[object] = (), +) -> tuple[str, ...]: + """Return normalized non-zero holdings followed by explicit supplements. + + Holdings are read from the already-rendered same-cycle dashboard so each + platform keeps ownership of broker-specific valuation and quantity rules. + Supplemental lines must already be localized by the platform translator. + """ + holdings: list[str] = [] + in_holdings = False + for raw_line in str(dashboard_text or "").splitlines(): + line = raw_line.strip() + if line in _HOLDINGS_HEADERS: + in_holdings = True + continue + if not in_holdings or not line: + continue + if line.startswith("━") or line.startswith(("📌", "💵", "📊", "🎯", "🧾", "⏱", "🧩")): + break + normalized = line.lstrip("-• ").strip() + if ":" not in normalized and ":" not in normalized: + continue + separator = ":" if ":" in normalized else ":" + symbol, detail = (part.strip() for part in normalized.split(separator, 1)) + if not symbol or not detail or not _contains_nonzero_number(detail): + continue + detail = _localize_holding_detail(detail, locale=locale) + holdings.append(f"- {symbol}: {detail}") + + lines: list[str] = [] + if holdings: + lines.append("💼 持仓" if str(locale).lower().startswith("zh") else "💼 Holdings") + lines.extend(holdings) + + for raw_line in supplemental_lines: + line = str(raw_line or "").strip() + if line and line not in lines: + lines.append(line) + return tuple(lines) diff --git a/notifications/telegram.py b/notifications/telegram.py index df76c0c..0e9b616 100644 --- a/notifications/telegram.py +++ b/notifications/telegram.py @@ -25,6 +25,8 @@ localize_price_source_label as _localize_price_source_label, present as _present, ) +from notifications.compact_adapter import adapt_compact_sections +from quant_platform_kit.notifications.events import RenderedNotification try: from quant_platform_kit.common.notification_localization import ( @@ -1162,3 +1164,125 @@ def render_cycle_summary(result: Mapping[str, Any], *, lang: str = "en") -> str: else: lines.append(translator("no_rebalance_needed")) return "\n".join(str(line) for line in lines if str(line).strip()) + + +def _compact_cycle_total_assets_line( + result: Mapping[str, Any], + *, + translator: Callable[..., str], + has_submitted_orders: bool, +) -> str: + if not has_submitted_orders: + snapshot = result.get("heartbeat_account_snapshot") + snapshot = snapshot if isinstance(snapshot, Mapping) else {} + observed_at = str(snapshot.get("observed_at") or "").strip() + amount = snapshot.get("net_assets") + if ( + isinstance(amount, (int, float)) + and not isinstance(amount, bool) + and math.isfinite(amount) + and amount > 0 + and observed_at + ): + return translator("heartbeat_account_equity", value=f"USD {amount:,.2f}") + + portfolio = result.get("portfolio") + portfolio = portfolio if isinstance(portfolio, Mapping) else {} + total_equity = _safe_float(portfolio.get("total_equity")) + if total_equity is None: + return "" + execution = result.get("execution") + execution = execution if isinstance(execution, Mapping) else {} + cash_only = bool(execution.get("cash_only_execution", portfolio.get("cash_only_execution", True))) + label = translator("total_assets" if cash_only else "total_assets_margin") + return f"💰 {label}: {_format_money(total_equity)}" + + +def _compact_strategy_name(result: Mapping[str, Any], *, lang: str, translator) -> str: + strategy_profile = str(result.get("strategy_profile") or "").strip() + strategy_name = str(result.get("strategy_display_name") or strategy_profile).strip() + strategy_metadata = result.get("strategy_metadata") + if strategy_metadata is not None: + from quant_platform_kit.common.notification_localization import resolve_strategy_display_name + + return resolve_strategy_display_name(lang, strategy_metadata, translator=translator) + translated = translator(f"strategy_name_{strategy_profile}") if strategy_profile else "" + if translated and translated != f"strategy_name_{strategy_profile}": + return translated + return strategy_name + + +def render_cycle_notification(result: Mapping[str, Any], *, lang: str = "en") -> RenderedNotification: + """Render full technical detail for logs and a concise user-facing message.""" + translator = build_translator(lang) + submitted = list(result.get("submitted_orders") or ()) + skipped = list(result.get("skipped_orders") or ()) + allocation = dict(result.get("allocation") or {}) + portfolio = dict(result.get("portfolio") or {}) + execution = dict(result.get("execution") or {}) + target_diff_lines = _format_target_diff_lines(allocation, portfolio, translator=translator) + meaningful_skipped = [ + item for item in skipped if str(item.get("reason") or "") != "below_trade_threshold" + ] + has_rebalance_attempt = bool(submitted or target_diff_lines or meaningful_skipped) + + lines = [translator("rebalance_title" if has_rebalance_attempt else "heartbeat_title")] + strategy_name = _compact_strategy_name(result, lang=lang, translator=translator) + if strategy_name: + lines.append(translator("strategy_label", name=strategy_name)) + total_assets_line = _compact_cycle_total_assets_line( + result, + translator=translator, + has_submitted_orders=bool(submitted), + ) + if total_assets_line: + lines.append(total_assets_line) + if bool(result.get("dry_run_only")): + lines.append(translator("dry_run_banner")) + compact_dashboard = "\n".join( + _format_dashboard_lines(portfolio, execution, translator=translator) + ) + lines.extend( + adapt_compact_sections( + compact_dashboard, + locale=lang, + supplemental_lines=result.get("compact_supplemental_lines", ()), + ) + ) + + if bool(result.get("execution_blocked")): + blocked = list(result.get("execution_blocking_skips") or skipped) + reason = _format_skipped_reason(blocked, translator=translator) + if bool(result.get("funding_blocked")): + banner_key = "funding_blocked_banner" + elif bool(result.get("execution_block_retryable")): + banner_key = "execution_blocked_retryable_banner" + else: + banner_key = "execution_blocked_banner" + lines.append(translator(banner_key, reason=reason)) + + if submitted: + lines.extend( + _format_order_lines( + submitted, + dry_run_only=bool(result.get("dry_run_only")), + translator=translator, + ) + ) + hard_skips = [ + item + for item in meaningful_skipped + if str(item.get("reason") or "") + not in {"buy_quantity_zero", "sell_quantity_zero", "quantity_zero", "min_notional"} + ] + if hard_skips: + lines.append(translator("no_order_submitted", reason=_format_skipped_reason(hard_skips, translator=translator))) + elif skipped and has_rebalance_attempt: + lines.append(translator("no_order_submitted", reason=_format_skipped_reason(skipped, translator=translator))) + else: + lines.append(translator("no_rebalance_needed")) + + return RenderedNotification( + detailed_text=render_cycle_summary(result, lang=lang), + compact_text="\n".join(str(line) for line in lines if str(line).strip()), + ) diff --git a/scripts/reconcile_cloud_runtime.py b/scripts/reconcile_cloud_runtime.py index f82ec28..5160a11 100644 --- a/scripts/reconcile_cloud_runtime.py +++ b/scripts/reconcile_cloud_runtime.py @@ -179,6 +179,62 @@ def _describe_revision( ) +def _list_revisions( + ctx: RuntimeContext, + run_gcloud: RunGcloud = _run_gcloud, +) -> list[dict[str, Any]]: + command = [ + "run", + "revisions", + "list", + "--service", + ctx.service_name, + "--project", + ctx.project_id, + "--region", + ctx.region, + "--format=json", + ] + result = run_gcloud(command) + if result.returncode != 0: + stderr = (result.stderr or "").strip() + raise RuntimeError( + f"gcloud {' '.join(command)} failed with exit code {result.returncode}" + + (f": {stderr}" if stderr else "") + ) + payload = json.loads((result.stdout or "[]").strip() or "[]") + if not isinstance(payload, list): + raise RuntimeError("Cloud Run revision list did not return a JSON array") + return [item for item in payload if isinstance(item, dict)] + + +def _ready_revision_for_commit( + ctx: RuntimeContext, + target_sha: str, + *, + run_gcloud: RunGcloud = _run_gcloud, +) -> str: + matches = [] + for revision in _list_revisions(ctx, run_gcloud=run_gcloud): + metadata = revision.get("metadata", {}) + if str(metadata.get("labels", {}).get("commit-sha") or "").strip() != target_sha: + continue + conditions = revision.get("status", {}).get("conditions", []) or [] + if not any( + isinstance(condition, Mapping) + and condition.get("type") == "Ready" + and str(condition.get("status") or "").lower() == "true" + for condition in conditions + ): + continue + name = str(metadata.get("name") or "").strip() + if name: + matches.append((str(metadata.get("creationTimestamp") or ""), name)) + if not matches: + return "" + return max(matches)[1] + + def _wait_for_latest_ready_revision( ctx: RuntimeContext, *, @@ -204,7 +260,7 @@ def _wait_for_latest_ready_revision( time.sleep(poll_seconds) -def _traffic_is_reconciled(service: Mapping[str, Any], latest_revision: str) -> bool: +def _traffic_is_reconciled(service: Mapping[str, Any], target_revision: str) -> bool: traffic = service.get("status", {}).get("traffic", []) or [] positive = [ entry @@ -214,9 +270,15 @@ def _traffic_is_reconciled(service: Mapping[str, Any], latest_revision: str) -> if len(positive) != 1: return False entry = positive[0] - return int(entry.get("percent", 0) or 0) == 100 and ( - entry.get("latestRevision") is True or str(entry.get("revisionName") or "") == latest_revision - ) + if int(entry.get("percent", 0) or 0) != 100: + return False + revision_name = str(entry.get("revisionName") or "").strip() + if revision_name: + return revision_name == target_revision + latest_ready = str( + service.get("status", {}).get("latestReadyRevisionName") or "" + ).strip() + return entry.get("latestRevision") is True and latest_ready == target_revision def reconcile_traffic( @@ -230,17 +292,25 @@ def reconcile_traffic( raise ValueError("GITHUB_SHA is required for traffic reconciliation") service, latest_revision = _wait_for_latest_ready_revision(ctx, run_gcloud=run_gcloud) + target_revision = latest_revision revision = _describe_revision(ctx, latest_revision, run_gcloud=run_gcloud) revision_sha = str( revision.get("metadata", {}).get("labels", {}).get("commit-sha") or "" ).strip() if revision_sha != target_sha: - raise RuntimeError( - f"Latest ready revision {latest_revision} on {ctx.service_name} has commit-sha " - f"{revision_sha or ''}, expected {target_sha}" + target_revision = _ready_revision_for_commit( + ctx, + target_sha, + run_gcloud=run_gcloud, ) + if not target_revision: + raise RuntimeError( + f"Latest ready revision {latest_revision} on {ctx.service_name} has commit-sha " + f"{revision_sha or ''}, and no Ready revision matches expected " + f"commit {target_sha}" + ) - if not _traffic_is_reconciled(service, latest_revision): + if not _traffic_is_reconciled(service, target_revision): _run_gcloud_ok( [ "run", @@ -252,7 +322,7 @@ def reconcile_traffic( "--region", ctx.region, "--to-revisions", - f"{latest_revision}=100", + f"{target_revision}=100", "--quiet", ], run_gcloud=run_gcloud, @@ -261,14 +331,14 @@ def reconcile_traffic( deadline = time.monotonic() + 300 while True: service = _describe_service(ctx, run_gcloud=run_gcloud) - if _traffic_is_reconciled(service, latest_revision): + if _traffic_is_reconciled(service, target_revision): print( - f"Cloud Run service {ctx.service_name} traffic reconciled to {latest_revision}." + f"Cloud Run service {ctx.service_name} traffic reconciled to {target_revision}." ) return if time.monotonic() >= deadline: raise RuntimeError( - f"Timed out waiting for {ctx.service_name} traffic to converge on {latest_revision}" + f"Timed out waiting for {ctx.service_name} traffic to converge on {target_revision}" ) time.sleep(5) diff --git a/scripts/send_paper_notification_preview.py b/scripts/send_paper_notification_preview.py index 86b7a97..b6c5f1c 100644 --- a/scripts/send_paper_notification_preview.py +++ b/scripts/send_paper_notification_preview.py @@ -24,7 +24,7 @@ build_sender, build_strategy_display_name, build_translator, - render_cycle_summary, + render_cycle_notification, ) _MAX_PREVIEW_MESSAGES = 6 @@ -117,14 +117,14 @@ def build_preview_messages(*, locale: str | None = None) -> list[str]: ) account_line = translator("account_label", account="PAPER") - heartbeat = render_cycle_summary( + heartbeat = render_cycle_notification( { **_base_result(dry_run_only=True, strategy_display_name=strategy_name), }, lang=resolved_locale, - ) + ).compact_text - dry_run = render_cycle_summary( + dry_run = render_cycle_notification( { **_base_result(dry_run_only=True, strategy_display_name=strategy_name), "portfolio": { @@ -146,9 +146,9 @@ def build_preview_messages(*, locale: str | None = None) -> list[str]: ], }, lang=resolved_locale, - ) + ).compact_text - pending = render_cycle_summary( + pending = render_cycle_notification( { **_base_result(dry_run_only=False, strategy_display_name=strategy_name), "allocation": {"targets": {_SYNTHETIC_SYMBOL: 100.0}}, @@ -164,7 +164,7 @@ def build_preview_messages(*, locale: str | None = None) -> list[str]: ], }, lang=resolved_locale, - ) + ).compact_text pending = "\n".join((pending, "synthetic PREVIEW pending confirmation / 订单待确认")) filled_order = ( @@ -176,15 +176,13 @@ def build_preview_messages(*, locale: str | None = None) -> list[str]: ( translator("rebalance_title"), translator("strategy_label", name=strategy_name), - account_line, translator("dry_run_banner"), - translator("order_logs_title"), filled_order, "synthetic PREVIEW filled / 成交确认", ) ) - rejected = render_cycle_summary( + rejected = render_cycle_notification( { **_base_result(dry_run_only=True, strategy_display_name=strategy_name), "allocation": {"targets": {_SYNTHETIC_SYMBOL: 100.0}}, @@ -193,7 +191,7 @@ def build_preview_messages(*, locale: str | None = None) -> list[str]: ], }, lang=resolved_locale, - ) + ).compact_text rejected = "\n".join((rejected, "synthetic PREVIEW reject / 拒单异常")) unknown_status = "\n".join( diff --git a/tests/test_compact_notification_adapter.py b/tests/test_compact_notification_adapter.py new file mode 100644 index 0000000..fe23fd2 --- /dev/null +++ b/tests/test_compact_notification_adapter.py @@ -0,0 +1,41 @@ +from notifications.compact_adapter import adapt_compact_sections + + +def test_adapter_keeps_only_nonzero_holdings_and_appends_supplements(): + dashboard = """📌 策略账户概览 +- 总资产(策略净值): $581.59 +💼 策略持仓 +- SOXL: $151.80 / 1股 +- SOXX: $0.00 / 0股 +- BOXX: $0.00 / 2股 +━━━━━━━━━━━━━━━━━━ +📊 市场状态: 观察""" + + assert adapt_compact_sections( + dashboard, + locale="zh", + supplemental_lines=("⚠️ 订单结果待确认", "", "⚠️ 订单结果待确认"), + ) == ( + "💼 持仓", + "- SOXL: $151.80 / 1股", + "- BOXX: $0.00 / 2股", + "⚠️ 订单结果待确认", + ) + + +def test_adapter_omits_empty_holdings_section(): + assert adapt_compact_sections( + "💼 Strategy Holdings\n- TQQQ: $0.00 / 0 shares", + locale="en", + ) == () + + +def test_adapter_localizes_holding_units_without_mixed_language(): + assert adapt_compact_sections( + "💼 策略持仓\n- TQQQ: $80.18 / 1股", + locale="en", + ) == ("💼 Holdings", "- TQQQ: $80.18 / 1 share") + assert adapt_compact_sections( + "💼 Strategy Holdings\n- SOXL: $151.80 / 2 shares", + locale="zh", + ) == ("💼 持仓", "- SOXL: $151.80 / 2股") diff --git a/tests/test_notifications_telegram.py b/tests/test_notifications_telegram.py index 0b4f89f..d714a59 100644 --- a/tests/test_notifications_telegram.py +++ b/tests/test_notifications_telegram.py @@ -1,6 +1,6 @@ from __future__ import annotations -from notifications.telegram import render_cycle_summary +from notifications.telegram import render_cycle_notification, render_cycle_summary def test_no_trade_heartbeat_shows_broker_account_values_in_both_locales(): @@ -25,6 +25,76 @@ def test_no_trade_heartbeat_shows_broker_account_values_in_both_locales(): assert "Total account equity: USD 1,234.56" in en +def test_compact_heartbeat_keeps_nonzero_holdings_and_result(): + notification = render_cycle_notification({ + "account": "****1234", + "strategy_profile": "soxl_soxx_trend_income", + "portfolio": { + "total_equity": 1234.56, + "liquid_cash": 100.0, + "portfolio_rows": (("SOXL",),), + "market_values": {"SOXL": 1134.56}, + "quantities": {"SOXL": 1}, + }, + "allocation": {"targets": {"SOXL": 1134.56}}, + "execution": {"signal_display": "hold"}, + "strategy_plugin_error_lines": ("🧩 插件本次影响:仅通知复核",), + "submitted_orders": [], + "skipped_orders": [], + "heartbeat_account_snapshot": { + "available_cash": 100.0, + "net_assets": 1234.56, + "observed_at": "2026-09-25T19:45:00+00:00", + }, + }, lang="zh") + + assert notification.compact_text == ( + "💓 【心跳检测】\n" + "🧭 策略: SOXL/SOXX 半导体趋势收益\n" + "💰 账户总权益: USD 1,234.56\n" + "💼 持仓\n" + "- SOXL: $1,134.56 / 1股\n" + "✅ 无需调仓" + ) + + +def test_compact_trade_keeps_nonzero_holdings_and_order_result(): + notification = render_cycle_notification({ + "account": "****1234", + "strategy_profile": "soxl_soxx_trend_income", + "portfolio": { + "total_equity": 1234.56, + "liquid_cash": 100.0, + "portfolio_rows": (("BOXX", "QQQM"),), + "market_values": {"BOXX": 500.0, "QQQM": 0.0}, + "quantities": {"BOXX": 5, "QQQM": 0}, + }, + "execution": {"cash_only_execution": True}, + "compact_supplemental_lines": ("⚠️ 订单仍待券商确认",), + "allocation": {"targets": {"SOXL": 500.0}}, + "submitted_orders": [{ + "side": "buy", + "symbol": "SOXL", + "quantity": 1, + "order_type": "limit", + "limit_price": 151.8, + "broker_order_id": "21", + }], + "skipped_orders": [{"symbol": "QQQM", "reason": "buy_quantity_zero"}], + }, lang="zh") + + assert notification.compact_text == ( + "🔔 【调仓指令】\n" + "🧭 策略: SOXL/SOXX 半导体趋势收益\n" + "💰 总资产(策略标的+现金): $1,234.56\n" + "💼 持仓\n" + "- BOXX: $500.00 / 5股\n" + "⚠️ 订单仍待券商确认\n" + "📈 已提交限价买入 SOXL: 1股 @ $151.80(订单号: 21)" + "(尚未确认成交;限价单可能未成交或取消)" + ) + + def test_no_trade_heartbeat_marks_missing_account_values_unverified(): message = render_cycle_summary({ "account": "****1234", diff --git a/tests/test_rebalance_service.py b/tests/test_rebalance_service.py index 3a1cd35..92e9524 100644 --- a/tests/test_rebalance_service.py +++ b/tests/test_rebalance_service.py @@ -391,10 +391,10 @@ def fake_client_factory(*args, **kwargs): assert result["notification_sent"] is True assert "🔔 【Rebalance Instruction】" in messages[0] assert "🧭 Strategy: TQQQ Growth Income" in messages[0] - assert "🆔 Account: 12345678" in messages[0] - assert "📌 Strategy Account" in messages[0] - assert "Target changes: AAA +50.00 USD" in messages[0] - assert "🧾 Execution details" in messages[0] + assert "🆔 Account: 12345678" not in messages[0] + assert "📌 Strategy Account" not in messages[0] + assert "Target changes: AAA +50.00 USD" not in messages[0] + assert "🧾 Execution details" not in messages[0] assert "🧪 Dry-run limit buy AAA: 2 shares @ $10.05" in messages[0] @@ -589,7 +589,7 @@ def evaluate(self, **inputs): assert len(messages) == 1 assert "Heartbeat" in messages[0] assert "Total assets: $0.00" in messages[0] - assert "Available cash: $0.00" in messages[0] + assert "Available cash: $0.00" not in messages[0] def test_run_strategy_cycle_loads_strategy_plugin_report_and_sends_email( diff --git a/tests/test_reconcile_cloud_runtime.py b/tests/test_reconcile_cloud_runtime.py index 63e0ffe..2aafd90 100644 --- a/tests/test_reconcile_cloud_runtime.py +++ b/tests/test_reconcile_cloud_runtime.py @@ -121,11 +121,78 @@ def fake_run_gcloud(command: list[str]): return _completed(command, stdout=json.dumps(service_state)) if command[:4] == ["run", "revisions", "describe", "firstrade-platform-service-00002"]: return _completed(command, stdout=json.dumps(revision_state)) + if command[:4] == ["run", "revisions", "list", "--service"]: + return _completed(command, stdout="[]") raise AssertionError(f"Unexpected gcloud command: {command}") with self.assertRaisesRegex(RuntimeError, "commit-sha"): reconciler.reconcile_traffic(env, run_gcloud=fake_run_gcloud) + def test_reconcile_traffic_selects_ready_revision_for_target_commit(self): + env = { + "GCP_PROJECT_ID": "firstradequant", + "CLOUD_RUN_SERVICE": "firstrade-platform-service", + "CLOUD_RUN_REGION": "us-central1", + "GITHUB_SHA": "abc123", + "SYNC_PLAN_JSON": json.dumps( + {"targets": [{"service_name": "firstrade-platform-service"}]} + ), + } + service_state = { + "status": { + "latestReadyRevisionName": "firstrade-platform-service-00003", + "traffic": [ + {"revisionName": "firstrade-platform-service-00001", "percent": 100} + ], + } + } + revisions = [ + { + "metadata": { + "name": "firstrade-platform-service-00002", + "creationTimestamp": "2026-09-26T15:03:00Z", + "labels": {"commit-sha": "abc123"}, + }, + "status": {"conditions": [{"type": "Ready", "status": "True"}]}, + }, + { + "metadata": { + "name": "firstrade-platform-service-00003", + "creationTimestamp": "2026-09-26T15:04:00Z", + "labels": {"commit-sha": "oldsha"}, + }, + "status": {"conditions": [{"type": "Ready", "status": "True"}]}, + }, + ] + calls: list[list[str]] = [] + + def fake_run_gcloud(command: list[str]): + calls.append(command) + if command[:4] == ["run", "services", "describe", "firstrade-platform-service"]: + return _completed(command, stdout=json.dumps(service_state)) + if command[:4] == ["run", "revisions", "describe", "firstrade-platform-service-00003"]: + return _completed( + command, + stdout=json.dumps({"metadata": {"labels": {"commit-sha": "oldsha"}}}), + ) + if command[:4] == ["run", "revisions", "list", "--service"]: + return _completed(command, stdout=json.dumps(revisions)) + if command[:4] == ["run", "services", "update-traffic", "firstrade-platform-service"]: + service_state["status"]["traffic"] = [ + {"revisionName": "firstrade-platform-service-00002", "percent": 100} + ] + return _completed(command) + raise AssertionError(f"Unexpected gcloud command: {command}") + + reconciler.reconcile_traffic(env, run_gcloud=fake_run_gcloud) + + update_call = next( + command + for command in calls + if command[:3] == ["run", "services", "update-traffic"] + ) + self.assertIn("firstrade-platform-service-00002=100", update_call) + def test_cleanup_legacy_scheduler_jobs_preserves_canonical_direct_jobs(self): env = { "GCP_PROJECT_ID": "firstradequant",