From 4db3ff4b5c78b871354078f1d41fad80805cdf2d Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Sat, 27 Jun 2026 15:43:18 +0800 Subject: [PATCH] refactor: unify Cloud Run env sync with build_cloud_run_env_sync_plan.py Co-authored-by: Cursor --- .github/workflows/sync-cloud-run-env.yml | 270 +++-------- scripts/build_cloud_run_env_sync_plan.py | 502 ++++++++++++++++++++ tests/test_build_cloud_run_env_sync_plan.py | 153 ++++++ tests/test_sync_cloud_run_env_workflow.py | 37 +- 4 files changed, 746 insertions(+), 216 deletions(-) create mode 100644 scripts/build_cloud_run_env_sync_plan.py create mode 100644 tests/test_build_cloud_run_env_sync_plan.py diff --git a/.github/workflows/sync-cloud-run-env.yml b/.github/workflows/sync-cloud-run-env.yml index f3441cb..99a94e1 100644 --- a/.github/workflows/sync-cloud-run-env.yml +++ b/.github/workflows/sync-cloud-run-env.yml @@ -269,81 +269,28 @@ jobs: python -m pip install --upgrade pip python -m pip install -r requirements.txt - - name: Resolve selected strategy runtime requirements + - name: Resolve Cloud Run sync targets id: strategy_requirements if: steps.env_sync_config.outputs.enabled == 'true' run: | set -euo pipefail - python - <<'PY' - import json - import os - import subprocess - import sys - from us_equity_strategies import resolve_canonical_profile - - raw_runtime_target = os.environ.get("RUNTIME_TARGET_JSON", "").strip() - if not raw_runtime_target: - raise SystemExit("RUNTIME_TARGET_JSON is required") - runtime_target = json.loads(raw_runtime_target) - profile = str(runtime_target.get("strategy_profile") or "").strip().lower() - if not profile: - raise SystemExit("RUNTIME_TARGET_JSON.strategy_profile is required") - canonical_profile = resolve_canonical_profile(profile) - runtime_target["strategy_profile"] = canonical_profile - - expected_service = os.environ.get("CLOUD_RUN_SERVICE", "").strip() - configured_service = str(runtime_target.get("service_name") or "").strip() - if configured_service and expected_service and configured_service != expected_service: - raise SystemExit( - "RUNTIME_TARGET_JSON.service_name does not match CLOUD_RUN_SERVICE: " - f"{configured_service!r} != {expected_service!r}" - ) - - raw_status = subprocess.check_output( - [sys.executable, "scripts/print_strategy_profile_status.py", "--json"], - text=True, - ) - rows = json.loads(raw_status) - selected = next((row for row in rows if row["canonical_profile"] == canonical_profile), None) - if selected is None: - supported = ", ".join(sorted(row["canonical_profile"] for row in rows)) - raise SystemExit(f"Unsupported STRATEGY_PROFILE={profile!r}; supported: {supported}") - if not selected.get("eligible") or not selected.get("enabled"): - raise SystemExit(f"STRATEGY_PROFILE={profile!r} is not eligible/enabled: {selected}") - - output_path = os.environ["GITHUB_OUTPUT"] - with open(output_path, "a", encoding="utf-8") as output: - output.write(f"strategy_profile={canonical_profile}\n") - output.write( - f"requires_snapshot_artifacts={str(bool(selected.get('requires_snapshot_artifacts'))).lower()}\n" - ) - output.write( - f"requires_snapshot_manifest_path={str(bool(selected.get('requires_snapshot_manifest_path'))).lower()}\n" - ) - output.write( - f"requires_strategy_config_path={str(bool(selected.get('requires_strategy_config_path'))).lower()}\n" - ) - output.write( - f"config_source_policy={str(selected.get('config_source_policy') or 'none')}\n" - ) - output.write(f"runtime_target_json={json.dumps(runtime_target, sort_keys=True)}\n") - PY + sync_plan_json="$(python scripts/build_cloud_run_env_sync_plan.py --json)" + { + echo "sync_plan_json<<__SYNC_PLAN_JSON__" + printf '%s\n' "${sync_plan_json}" + echo "__SYNC_PLAN_JSON__" + } >> "$GITHUB_OUTPUT" - name: Validate env sync inputs if: steps.env_sync_config.outputs.enabled == 'true' env: - REQUIRES_SNAPSHOT_ARTIFACTS: ${{ steps.strategy_requirements.outputs.requires_snapshot_artifacts }} - REQUIRES_SNAPSHOT_MANIFEST_PATH: ${{ steps.strategy_requirements.outputs.requires_snapshot_manifest_path }} - REQUIRES_STRATEGY_CONFIG_PATH: ${{ steps.strategy_requirements.outputs.requires_strategy_config_path }} - CONFIG_SOURCE_POLICY: ${{ steps.strategy_requirements.outputs.config_source_policy }} - RUNTIME_TARGET_JSON: ${{ steps.strategy_requirements.outputs.runtime_target_json }} + SYNC_PLAN_JSON: ${{ steps.strategy_requirements.outputs.sync_plan_json }} run: | set -euo pipefail required_vars=( CLOUD_RUN_REGION CLOUD_RUN_SERVICE - RUNTIME_TARGET_JSON ) if [ -z "${NOTIFY_LANG:-}" ]; then @@ -370,20 +317,6 @@ jobs: required_vars+=(GLOBAL_TELEGRAM_CHAT_ID) fi - if [ "${REQUIRES_SNAPSHOT_ARTIFACTS:-}" = "true" ] && [ -z "${FIRSTRADE_FEATURE_SNAPSHOT_PATH:-}" ]; then - required_vars+=(FIRSTRADE_FEATURE_SNAPSHOT_PATH) - fi - - if [ "${REQUIRES_SNAPSHOT_MANIFEST_PATH:-}" = "true" ] && [ -z "${FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH:-}" ]; then - required_vars+=(FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH) - fi - - if [ "${REQUIRES_STRATEGY_CONFIG_PATH:-}" = "true" ] \ - && [ "${CONFIG_SOURCE_POLICY:-}" = "env_only" ] \ - && [ -z "${FIRSTRADE_STRATEGY_CONFIG_PATH:-}" ]; then - required_vars+=(FIRSTRADE_STRATEGY_CONFIG_PATH) - fi - missing_vars=() for var_name in "${required_vars[@]}"; do if [[ "${var_name}" == *" or "* ]]; then @@ -393,6 +326,22 @@ jobs: fi done + python - <<'PY' + import json + import os + + plan = json.loads(os.environ["SYNC_PLAN_JSON"]) + targets = plan.get("targets") or [] + if not targets: + raise SystemExit("Cloud Run env sync did not resolve any targets") + for target in targets: + service_name = str(target.get("service_name") or "").strip() + if not service_name: + raise SystemExit("Cloud Run sync target is missing service_name") + if not isinstance(target.get("env"), dict): + raise SystemExit(f"Cloud Run sync target {service_name} is missing env") + PY + if [ "${#missing_vars[@]}" -gt 0 ]; then echo "Cloud Run env sync is enabled, but these values are missing:" >&2 printf ' - %s\n' "${missing_vars[@]}" >&2 @@ -426,11 +375,31 @@ jobs: - name: Sync Cloud Run environment if: steps.env_sync_config.outputs.enabled == 'true' env: - STRATEGY_PROFILE: ${{ steps.strategy_requirements.outputs.strategy_profile }} - RUNTIME_TARGET_JSON: ${{ steps.strategy_requirements.outputs.runtime_target_json }} + SYNC_PLAN_JSON: ${{ steps.strategy_requirements.outputs.sync_plan_json }} run: | set -euo pipefail + mapfile -t target_env_pairs < <(python - <<'PY' + import json + import os + + plan = json.loads(os.environ["SYNC_PLAN_JSON"]) + target = plan["targets"][0] + for key, value in sorted(target["env"].items()): + print(f"{key}={value}") + PY + ) + mapfile -t target_remove_env_vars < <(python - <<'PY' + import json + import os + + plan = json.loads(os.environ["SYNC_PLAN_JSON"]) + target = plan["targets"][0] + for key in sorted(target.get("remove_env_vars") or []): + print(key) + PY + ) + join_by_delimiter() { local delimiter="$1" shift @@ -446,13 +415,11 @@ jobs: printf '%s' "${output}" } - env_pairs=( - "GOOGLE_CLOUD_PROJECT=${GCP_PROJECT_ID}" - "RUNTIME_TARGET_JSON=${RUNTIME_TARGET_JSON}" - "STRATEGY_PROFILE=${STRATEGY_PROFILE}" - ) + env_pairs=("${target_env_pairs[@]}") + env_pairs+=("GOOGLE_CLOUD_PROJECT=${GCP_PROJECT_ID}") secret_pairs=() - remove_env_vars=( + remove_env_vars=("${target_remove_env_vars[@]}") + remove_env_vars+=( "TELEGRAM_CHAT_ID" "CRISIS_ALERT_GOOGLE_VOICE_TO" "CRISIS_ALERT_GOOGLE_VOICE_GATEWAY" @@ -513,16 +480,6 @@ jobs: "CRISIS_ALERT_TELEGRAM_BOT_TOKEN" ) - add_optional_env() { - local name="$1" - local value="${!name:-}" - if [ -n "${value}" ]; then - env_pairs+=("${name}=${value}") - else - remove_env_vars+=("${name}") - fi - } - add_optional_secret() { local env_name="$1" local secret_var_name="$2" @@ -541,89 +498,6 @@ jobs: fi } - add_optional_env ACCOUNT_PREFIX - add_optional_env ACCOUNT_REGION - add_optional_env FIRSTRADE_ACCOUNT - add_optional_env FIRSTRADE_COOKIE_DIR - add_optional_env FIRSTRADE_DRY_RUN_ONLY - add_optional_env FIRSTRADE_REUSE_SESSION - add_optional_env FIRSTRADE_SESSION_CACHE_TTL_SECONDS - add_optional_env FIRSTRADE_PERSIST_SESSION_CACHE - add_optional_env FIRSTRADE_GCS_STATE_BUCKET - add_optional_env FIRSTRADE_STATE_PREFIX - add_optional_env FIRSTRADE_PERSIST_ACCOUNT_SNAPSHOT - add_optional_env FIRSTRADE_PERSIST_STRATEGY_RUNS - add_optional_env FIRSTRADE_ENABLE_LIVE_TRADING - add_optional_env FIRSTRADE_RUN_SMOKE_ON_HTTP - add_optional_env FIRSTRADE_RUN_SESSION_CHECK_ON_HTTP - add_optional_env FIRSTRADE_SESSION_CHECK_INCLUDE_POSITIONS - add_optional_env FIRSTRADE_RUN_STRATEGY_ON_HTTP - add_optional_env FIRSTRADE_LIVE_ORDER_ACK - add_optional_env FIRSTRADE_MAX_ORDER_NOTIONAL_USD - add_optional_env FIRSTRADE_MIN_RESERVED_CASH_USD - add_optional_env FIRSTRADE_RESERVED_CASH_RATIO - add_optional_env FIRSTRADE_SAFE_HAVEN_CASH_SUBSTITUTE_THRESHOLD_USD - add_optional_env FIRSTRADE_SMOKE_SYMBOL - add_optional_env FIRSTRADE_FEATURE_SNAPSHOT_PATH - add_optional_env FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH - add_optional_env FIRSTRADE_FEATURE_SNAPSHOT_FALLBACK_MODE - add_optional_env FIRSTRADE_FEATURE_SNAPSHOT_FALLBACK_CACHE_DIR - add_optional_env FIRSTRADE_FEATURE_SNAPSHOT_MAX_STALE_DAYS - add_optional_env FIRSTRADE_STRATEGY_CONFIG_PATH - add_optional_env FIRSTRADE_STRATEGY_PLUGIN_MOUNTS_JSON - add_optional_env FIRSTRADE_MARKET_SIGNAL_HANDOFF_INDEX_URI - add_optional_env FIRSTRADE_MARKET_SIGNAL_HANDOFF_MANIFEST_URI - add_optional_env FIRSTRADE_MARKET_SIGNAL_CONSUMPTION_AUDIT_URI - add_optional_env FIRSTRADE_MARKET_SIGNAL_CACHE_DIR - add_optional_env FIRSTRADE_MARKET_SIGNAL_REQUIRED - add_optional_env FIRSTRADE_MARKET_SIGNAL_FALLBACK_MODE - add_optional_env FIRSTRADE_MARKET_SIGNAL_MAX_STALE_DAYS - add_optional_env STRATEGY_PLUGIN_ALERT_CHANNELS - add_optional_env STRATEGY_PLUGIN_ALERT_EMAIL_RECIPIENTS - add_optional_env STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_EMAIL - add_optional_env STRATEGY_PLUGIN_ALERT_EMAIL_SMTP_HOST - add_optional_env STRATEGY_PLUGIN_ALERT_EMAIL_SMTP_PORT - add_optional_env STRATEGY_PLUGIN_ALERT_EMAIL_SMTP_SECURITY - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_RECIPIENTS - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_PROVIDER - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_ACCOUNT_ID - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_SENDER - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_MESSAGING_SERVICE_ID - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_API_BASE_URL - add_optional_env STRATEGY_PLUGIN_ALERT_SMS_BODY_MAX_CHARS - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_RECIPIENTS - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_PROVIDER - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_API_BASE_URL - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_DEVICE - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_PRIORITY - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_TAGS - add_optional_env STRATEGY_PLUGIN_ALERT_PUSH_BODY_MAX_CHARS - add_optional_env STRATEGY_PLUGIN_ALERT_TELEGRAM_CHAT_IDS - add_optional_env STRATEGY_PLUGIN_ALERT_TELEGRAM_API_BASE_URL - add_optional_env STRATEGY_PLUGIN_ALERT_TELEGRAM_PARSE_MODE - add_optional_env STRATEGY_PLUGIN_ALERT_TELEGRAM_DISABLE_WEB_PAGE_PREVIEW - add_optional_env STRATEGY_PLUGIN_ALERT_TELEGRAM_BODY_MAX_CHARS - add_optional_env FIRSTRADE_RUNTIME_EXECUTION_WINDOW_TRADING_DAYS - add_optional_env FIRSTRADE_TECH_RUNTIME_EXECUTION_WINDOW_TRADING_DAYS - add_optional_env INCOME_THRESHOLD_USD - add_optional_env QQQI_INCOME_RATIO - add_optional_env INCOME_LAYER_ENABLED - add_optional_env INCOME_LAYER_START_USD - add_optional_env INCOME_LAYER_MAX_RATIO - add_optional_env CASH_ONLY_EXECUTION - add_optional_env DCA_MODE - add_optional_env DCA_BASE_INVESTMENT_USD - add_optional_env IBIT_ZSCORE_EXIT_ENABLED - add_optional_env IBIT_ZSCORE_EXIT_MODE - add_optional_env IBIT_ZSCORE_EXIT_PARKING_SYMBOL - add_optional_env IBIT_ZSCORE_EXIT_RISK_REDUCED_EXPOSURE - add_optional_env IBIT_ZSCORE_EXIT_RISK_OFF_EXPOSURE - add_optional_env IBIT_ZSCORE_EXIT_ALLOW_OUTSIDE_EXECUTION_WINDOW - add_optional_env RUNTIME_TARGET_ENABLED - add_optional_env EXECUTION_REPORT_GCS_URI - add_optional_env GLOBAL_TELEGRAM_CHAT_ID - add_optional_env NOTIFY_LANG - add_optional_secret TELEGRAM_TOKEN TELEGRAM_TOKEN_SECRET_NAME TELEGRAM_TOKEN add_optional_secret STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_PASSWORD STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_PASSWORD_SECRET_NAME STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_PASSWORD add_optional_secret STRATEGY_PLUGIN_ALERT_SMS_AUTH_TOKEN STRATEGY_PLUGIN_ALERT_SMS_AUTH_TOKEN_SECRET_NAME STRATEGY_PLUGIN_ALERT_SMS_AUTH_TOKEN @@ -651,16 +525,19 @@ jobs: import json import os - runtime_target = json.loads(os.environ.get("RUNTIME_TARGET_JSON") or "{}") - target = { - "service_name": os.environ.get("CLOUD_RUN_SERVICE"), + plan = json.loads(os.environ["SYNC_PLAN_JSON"]) + target = plan["targets"][0] + env = target.get("env") or {} + runtime_target = json.loads(env.get("RUNTIME_TARGET_JSON") or "{}") + payload = { + "service_name": target.get("service_name") or os.environ.get("CLOUD_RUN_SERVICE"), "service_url": os.environ.get("SERVICE_URL"), - "strategy_profile": runtime_target.get("strategy_profile"), + "strategy_profile": target.get("strategy_profile") or env.get("STRATEGY_PROFILE"), "account_scope": runtime_target.get("account_scope"), - "runtime_target_enabled": os.environ.get("RUNTIME_TARGET_ENABLED", "true"), - "scheduler": runtime_target.get("scheduler") if isinstance(runtime_target.get("scheduler"), dict) else {}, + "runtime_target_enabled": env.get("RUNTIME_TARGET_ENABLED", "true"), + "scheduler": target.get("scheduler") or {}, } - print(json.dumps({"targets": [target]}, separators=(",", ":"))) + print(json.dumps({"targets": [payload]}, separators=(",", ":"))) PY )" env_pairs+=("MONITOR_DISPATCH_TARGETS_JSON=${monitor_targets_json}") @@ -685,7 +562,7 @@ jobs: - name: Sync Cloud Scheduler schedule if: steps.env_sync_config.outputs.enabled == 'true' env: - RUNTIME_TARGET_JSON: ${{ steps.strategy_requirements.outputs.runtime_target_json }} + SYNC_PLAN_JSON: ${{ steps.strategy_requirements.outputs.sync_plan_json }} run: | set -euo pipefail @@ -699,24 +576,13 @@ jobs: import json import os - raw_runtime_target = os.environ.get("RUNTIME_TARGET_JSON", "").strip() - runtime_scheduler = {} - if raw_runtime_target: - try: - runtime_target = json.loads(raw_runtime_target) - except json.JSONDecodeError: - runtime_target = {} - scheduler = runtime_target.get("scheduler") if isinstance(runtime_target, dict) else {} - if isinstance(scheduler, dict): - runtime_scheduler = scheduler - - def configured_time(key: str, name: str, default: str) -> str: - return str(runtime_scheduler.get(key) or os.environ.get(name, "").strip() or default) - - print(str(runtime_scheduler.get("timezone") or "America/New_York").strip()) - print(configured_time("main_time", "CLOUD_SCHEDULER_MAIN_TIME", "45 15")) - print(configured_time("probe_time", "CLOUD_SCHEDULER_PROBE_TIME", "35 9,15")) - print(configured_time("precheck_time", "CLOUD_SCHEDULER_PRECHECK_TIME", "45 9")) + plan = json.loads(os.environ["SYNC_PLAN_JSON"]) + target = plan["targets"][0] + scheduler = target.get("scheduler") or {} + print(str(scheduler.get("timezone") or "America/New_York").strip()) + print(str(scheduler.get("main_time") or os.environ.get("CLOUD_SCHEDULER_MAIN_TIME", "").strip() or "45 15")) + print(str(scheduler.get("probe_time") or os.environ.get("CLOUD_SCHEDULER_PROBE_TIME", "").strip() or "35 9,15")) + print(str(scheduler.get("precheck_time") or os.environ.get("CLOUD_SCHEDULER_PRECHECK_TIME", "").strip() or "45 9")) PY ) market_timezone="${scheduler_config[0]}" diff --git a/scripts/build_cloud_run_env_sync_plan.py b/scripts/build_cloud_run_env_sync_plan.py new file mode 100644 index 0000000..f11f53c --- /dev/null +++ b/scripts/build_cloud_run_env_sync_plan.py @@ -0,0 +1,502 @@ +from __future__ import annotations + +import argparse +import json +import os +import sys +from collections.abc import Mapping, Sequence +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[1] +QPK_SRC = ROOT.parent / "QuantPlatformKit" / "src" +UES_SRC = ROOT.parent / "UsEquityStrategies" / "src" + + +def _has_catalog_marker(candidate: Path, package_name: str, marker: str) -> bool: + catalog_path = candidate / package_name / "catalog.py" + if not catalog_path.exists(): + return False + return marker in catalog_path.read_text(encoding="utf-8") + + +def _should_add_local_src(candidate: Path) -> bool: + if candidate == QPK_SRC: + return (candidate / "quant_platform_kit" / "common" / "runtime_target.py").exists() + if candidate == UES_SRC: + return _has_catalog_marker( + candidate, + "us_equity_strategies", + "global_etf_rotation", + ) + return True + + +for candidate in (ROOT, QPK_SRC, UES_SRC): + if not _should_add_local_src(candidate): + continue + candidate_str = str(candidate) + if candidate_str not in sys.path: + sys.path.insert(0, candidate_str) + +from strategy_registry import ( # noqa: E402 + FIRSTRADE_PLATFORM, + get_platform_profile_status_matrix, + resolve_strategy_definition, +) +from us_equity_strategies.runtime_adapters import ( # noqa: E402 + describe_platform_runtime_requirements, +) + + +TARGETS_JSON_ENV = "CLOUD_RUN_SERVICE_TARGETS_JSON" +SHARED_TARGET_FALLBACK_ENV = frozenset( + { + "GLOBAL_TELEGRAM_CHAT_ID", + "NOTIFY_LANG", + "ACCOUNT_PREFIX", + "ACCOUNT_REGION", + "EXECUTION_REPORT_GCS_URI", + } +) +REQUIRED_ENV = ("NOTIFY_LANG",) +OPTIONAL_TARGET_ENV = ( + "GLOBAL_TELEGRAM_CHAT_ID", + "ACCOUNT_PREFIX", + "ACCOUNT_REGION", + "FIRSTRADE_ACCOUNT", + "FIRSTRADE_COOKIE_DIR", + "FIRSTRADE_DRY_RUN_ONLY", + "FIRSTRADE_REUSE_SESSION", + "FIRSTRADE_SESSION_CACHE_TTL_SECONDS", + "FIRSTRADE_ENABLE_LIVE_TRADING", + "FIRSTRADE_RUN_SMOKE_ON_HTTP", + "FIRSTRADE_RUN_STRATEGY_ON_HTTP", + "FIRSTRADE_LIVE_ORDER_ACK", + "FIRSTRADE_MAX_ORDER_NOTIONAL_USD", + "FIRSTRADE_MIN_RESERVED_CASH_USD", + "FIRSTRADE_RESERVED_CASH_RATIO", + "FIRSTRADE_SAFE_HAVEN_CASH_SUBSTITUTE_THRESHOLD_USD", + "FIRSTRADE_SMOKE_SYMBOL", + "FIRSTRADE_FEATURE_SNAPSHOT_PATH", + "FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH", + "FIRSTRADE_FEATURE_SNAPSHOT_FALLBACK_MODE", + "FIRSTRADE_FEATURE_SNAPSHOT_FALLBACK_CACHE_DIR", + "FIRSTRADE_FEATURE_SNAPSHOT_MAX_STALE_DAYS", + "FIRSTRADE_GCS_STATE_BUCKET", + "FIRSTRADE_PERSIST_ACCOUNT_SNAPSHOT", + "FIRSTRADE_PERSIST_STRATEGY_RUNS", + "FIRSTRADE_PERSIST_SESSION_CACHE", + "FIRSTRADE_RUN_SESSION_CHECK_ON_HTTP", + "FIRSTRADE_SESSION_CHECK_INCLUDE_POSITIONS", + "FIRSTRADE_STATE_PREFIX", + "FIRSTRADE_STRATEGY_CONFIG_PATH", + "FIRSTRADE_STRATEGY_PLUGIN_MOUNTS_JSON", + "FIRSTRADE_MARKET_SIGNAL_HANDOFF_INDEX_URI", + "FIRSTRADE_MARKET_SIGNAL_HANDOFF_MANIFEST_URI", + "FIRSTRADE_MARKET_SIGNAL_CONSUMPTION_AUDIT_URI", + "FIRSTRADE_MARKET_SIGNAL_CACHE_DIR", + "FIRSTRADE_MARKET_SIGNAL_REQUIRED", + "FIRSTRADE_MARKET_SIGNAL_FALLBACK_MODE", + "FIRSTRADE_MARKET_SIGNAL_MAX_STALE_DAYS", + "STRATEGY_PLUGIN_ALERT_CHANNELS", + "STRATEGY_PLUGIN_ALERT_EMAIL_RECIPIENTS", + "STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_EMAIL", + "STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_PASSWORD_SECRET_NAME", + "STRATEGY_PLUGIN_ALERT_EMAIL_SMTP_HOST", + "STRATEGY_PLUGIN_ALERT_EMAIL_SMTP_PORT", + "STRATEGY_PLUGIN_ALERT_EMAIL_SMTP_SECURITY", + "STRATEGY_PLUGIN_ALERT_SMS_RECIPIENTS", + "STRATEGY_PLUGIN_ALERT_SMS_PROVIDER", + "STRATEGY_PLUGIN_ALERT_SMS_ACCOUNT_ID", + "STRATEGY_PLUGIN_ALERT_SMS_AUTH_TOKEN_SECRET_NAME", + "STRATEGY_PLUGIN_ALERT_SMS_SENDER", + "STRATEGY_PLUGIN_ALERT_SMS_MESSAGING_SERVICE_ID", + "STRATEGY_PLUGIN_ALERT_SMS_API_BASE_URL", + "STRATEGY_PLUGIN_ALERT_SMS_BODY_MAX_CHARS", + "STRATEGY_PLUGIN_ALERT_PUSH_RECIPIENTS", + "STRATEGY_PLUGIN_ALERT_PUSH_PROVIDER", + "STRATEGY_PLUGIN_ALERT_PUSH_APP_TOKEN_SECRET_NAME", + "STRATEGY_PLUGIN_ALERT_PUSH_ACCESS_TOKEN_SECRET_NAME", + "STRATEGY_PLUGIN_ALERT_PUSH_API_BASE_URL", + "STRATEGY_PLUGIN_ALERT_PUSH_DEVICE", + "STRATEGY_PLUGIN_ALERT_PUSH_PRIORITY", + "STRATEGY_PLUGIN_ALERT_PUSH_TAGS", + "STRATEGY_PLUGIN_ALERT_PUSH_BODY_MAX_CHARS", + "STRATEGY_PLUGIN_ALERT_TELEGRAM_CHAT_IDS", + "STRATEGY_PLUGIN_ALERT_TELEGRAM_BOT_TOKEN_SECRET_NAME", + "STRATEGY_PLUGIN_ALERT_TELEGRAM_API_BASE_URL", + "STRATEGY_PLUGIN_ALERT_TELEGRAM_PARSE_MODE", + "STRATEGY_PLUGIN_ALERT_TELEGRAM_DISABLE_WEB_PAGE_PREVIEW", + "STRATEGY_PLUGIN_ALERT_TELEGRAM_BODY_MAX_CHARS", + "FIRSTRADE_RUNTIME_EXECUTION_WINDOW_TRADING_DAYS", + "FIRSTRADE_TECH_RUNTIME_EXECUTION_WINDOW_TRADING_DAYS", + "INCOME_THRESHOLD_USD", + "QQQI_INCOME_RATIO", + "INCOME_LAYER_ENABLED", + "INCOME_LAYER_START_USD", + "INCOME_LAYER_MAX_RATIO", + "CASH_ONLY_EXECUTION", + "DCA_MODE", + "DCA_BASE_INVESTMENT_USD", + "IBIT_ZSCORE_EXIT_ENABLED", + "IBIT_ZSCORE_EXIT_MODE", + "IBIT_ZSCORE_EXIT_PARKING_SYMBOL", + "IBIT_ZSCORE_EXIT_RISK_REDUCED_EXPOSURE", + "IBIT_ZSCORE_EXIT_RISK_OFF_EXPOSURE", + "IBIT_ZSCORE_EXIT_ALLOW_OUTSIDE_EXECUTION_WINDOW", + "RUNTIME_TARGET_ENABLED", + "EXECUTION_REPORT_GCS_URI", +) +SCHEDULER_TIME_DEFAULTS = { + "main_time": "45 15", + "probe_time": "35 9,15", + "precheck_time": "45 9", +} +SCHEDULER_TIME_ENV = { + "main_time": "CLOUD_SCHEDULER_MAIN_TIME", + "probe_time": "CLOUD_SCHEDULER_PROBE_TIME", + "precheck_time": "CLOUD_SCHEDULER_PRECHECK_TIME", +} + + +def build_sync_plan(env: Mapping[str, str] = os.environ) -> dict[str, object]: + target_entries, defaults, per_service_mode = _load_target_entries(env) + status_rows = { + str(row["canonical_profile"]): { + **row, + **describe_platform_runtime_requirements( + str(row["canonical_profile"]), + platform_id=FIRSTRADE_PLATFORM, + ), + } + for row in get_platform_profile_status_matrix() + } + planned_targets = [ + _build_target_plan( + target=target, + defaults=defaults, + env=env, + status_rows=status_rows, + per_service_mode=per_service_mode, + ) + for target in target_entries + ] + if not planned_targets: + raise ValueError( + f"{TARGETS_JSON_ENV}, CLOUD_RUN_SERVICES, or CLOUD_RUN_SERVICE is required" + ) + return { + "mode": "per_service" if per_service_mode else "legacy", + "targets": planned_targets, + } + + +def _load_target_entries( + env: Mapping[str, str], +) -> tuple[list[dict[str, object]], Mapping[str, object], bool]: + raw_targets = str(env.get(TARGETS_JSON_ENV, "") or "").strip() + if raw_targets: + payload = json.loads(raw_targets) + if isinstance(payload, Mapping): + raw_entries = payload.get("targets") + defaults = _coerce_mapping(payload.get("defaults") or {}) + else: + raw_entries = payload + defaults = {} + if not isinstance(raw_entries, Sequence) or isinstance(raw_entries, (str, bytes)): + raise ValueError(f"{TARGETS_JSON_ENV} must be a JSON array or object with targets") + entries = [_coerce_mapping(item) for item in raw_entries] + return entries, defaults, True + + raw_services = str(env.get("CLOUD_RUN_SERVICES", "") or "").strip() + if not raw_services: + raw_services = str(env.get("CLOUD_RUN_SERVICE", "") or "").strip() + services = [ + item.strip() + for chunk in raw_services.replace(";", ",").replace("\n", ",").split(",") + for item in [chunk] + if item.strip() + ] + return [{"service": service} for service in services], {}, False + + +def _build_target_plan( + *, + target: Mapping[str, object], + defaults: Mapping[str, object], + env: Mapping[str, str], + status_rows: Mapping[str, Mapping[str, object]], + per_service_mode: bool, +) -> dict[str, object]: + service_name = _first_non_empty( + _target_field(target, defaults, "service"), + _target_field(target, defaults, "service_name"), + _target_field(target, defaults, "cloud_run_service"), + ) + runtime_target = _resolve_runtime_target(target, defaults, env, per_service_mode) + if not service_name: + service_name = str(runtime_target.get("service_name") or "").strip() + if not service_name: + raise ValueError("Each Cloud Run sync target requires service/service_name") + + runtime_target_service = str(runtime_target.get("service_name") or "").strip() + if runtime_target_service and runtime_target_service != service_name: + raise ValueError( + f"Target {service_name} runtime_target.service_name={runtime_target_service!r} does not match" + ) + runtime_target["service_name"] = service_name + + raw_profile = str(runtime_target.get("strategy_profile") or "").strip() + if not raw_profile: + raise ValueError(f"Target {service_name} runtime_target.strategy_profile is required") + definition = resolve_strategy_definition(raw_profile, platform_id=FIRSTRADE_PLATFORM) + canonical_profile = definition.profile + runtime_target["strategy_profile"] = canonical_profile + + status = status_rows.get(canonical_profile) + if status is None: + supported = ", ".join(sorted(status_rows)) + raise ValueError( + f"Unsupported STRATEGY_PROFILE={raw_profile!r} for {service_name}; supported: {supported}" + ) + if not status.get("eligible") or not status.get("enabled"): + raise ValueError( + f"STRATEGY_PROFILE={raw_profile!r} is not eligible/enabled for {service_name}: {status}" + ) + + env_values: dict[str, str] = {} + missing: list[str] = [] + for name in REQUIRED_ENV: + value = _target_env_value( + target, + defaults, + env, + name, + per_service_mode=per_service_mode, + allow_shared_fallback=name in SHARED_TARGET_FALLBACK_ENV, + ) + if value is None: + missing.append(f"{service_name}:{name}") + else: + env_values[name] = value + + env_values["STRATEGY_PROFILE"] = canonical_profile + env_values["RUNTIME_TARGET_JSON"] = json.dumps( + runtime_target, + separators=(",", ":"), + sort_keys=True, + ) + + remove_env_vars: list[str] = [] + for name in OPTIONAL_TARGET_ENV: + value = _target_env_value( + target, + defaults, + env, + name, + per_service_mode=per_service_mode, + allow_shared_fallback=name in SHARED_TARGET_FALLBACK_ENV, + ) + if value is None and name == "FIRSTRADE_DRY_RUN_ONLY": + dry_run_value = runtime_target.get("dry_run_only") + if dry_run_value is not None: + value = _coerce_env_value(dry_run_value) + if value is None: + remove_env_vars.append(name) + else: + env_values[name] = value + + if _runtime_target_enabled(env_values): + _validate_profile_inputs( + service_name=service_name, + env_values=env_values, + status=status, + missing=missing, + ) + if missing: + raise ValueError( + "Cloud Run env sync target values are missing:\n" + + "\n".join(f" - {item}" for item in missing) + ) + + return { + "service_name": service_name, + "strategy_profile": canonical_profile, + "env": env_values, + "scheduler": _build_scheduler_plan( + runtime_target=runtime_target, + target=target, + defaults=defaults, + env=env, + per_service_mode=per_service_mode, + ), + "remove_env_vars": sorted(set(remove_env_vars) - set(env_values)), + } + + +def _build_scheduler_plan( + *, + runtime_target: Mapping[str, object], + target: Mapping[str, object], + defaults: Mapping[str, object], + env: Mapping[str, str], + per_service_mode: bool, +) -> dict[str, str]: + runtime_scheduler = runtime_target.get("scheduler") if isinstance(runtime_target, Mapping) else {} + if not isinstance(runtime_scheduler, Mapping): + runtime_scheduler = {} + timezone = str(runtime_scheduler.get("timezone") or "").strip() + if not timezone: + timezone = "America/New_York" + + scheduler = {"timezone": timezone} + for key, env_name in SCHEDULER_TIME_ENV.items(): + configured_value = _target_env_value( + target, + defaults, + env, + env_name, + per_service_mode=per_service_mode, + allow_shared_fallback=True, + ) + scheduler[key] = str(runtime_scheduler.get(key) or configured_value or SCHEDULER_TIME_DEFAULTS[key]) + return scheduler + + +def _validate_profile_inputs( + *, + service_name: str, + env_values: Mapping[str, str], + status: Mapping[str, object], + missing: list[str], +) -> None: + if bool(status.get("requires_snapshot_artifacts")) and not env_values.get( + "FIRSTRADE_FEATURE_SNAPSHOT_PATH" + ): + missing.append(f"{service_name}:FIRSTRADE_FEATURE_SNAPSHOT_PATH") + if bool(status.get("requires_snapshot_manifest_path")) and not env_values.get( + "FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH" + ): + missing.append(f"{service_name}:FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH") + if ( + bool(status.get("requires_strategy_config_path")) + and str(status.get("config_source_policy") or "none") == "env_only" + and not env_values.get("FIRSTRADE_STRATEGY_CONFIG_PATH") + ): + missing.append(f"{service_name}:FIRSTRADE_STRATEGY_CONFIG_PATH") + + +def _runtime_target_enabled(env_values: Mapping[str, str]) -> bool: + raw = str(env_values.get("RUNTIME_TARGET_ENABLED") or "").strip().lower() + if not raw: + return True + if raw in {"1", "true", "yes", "on"}: + return True + if raw in {"0", "false", "no", "off"}: + return False + raise ValueError("RUNTIME_TARGET_ENABLED must be true or false") + + +def _resolve_runtime_target( + target: Mapping[str, object], + defaults: Mapping[str, object], + env: Mapping[str, str], + per_service_mode: bool, +) -> dict[str, object]: + raw = _first_non_empty( + _target_field(target, defaults, "runtime_target"), + _target_field(target, defaults, "runtime_target_json"), + ) + if raw is None and not per_service_mode: + raw = str(env.get("RUNTIME_TARGET_JSON", "") or "").strip() or None + if raw is None: + raise ValueError( + f"{TARGETS_JSON_ENV}.targets[].runtime_target is required in per-service mode" + ) + if isinstance(raw, Mapping): + return dict(raw) + if isinstance(raw, str): + loaded = json.loads(raw) + if not isinstance(loaded, Mapping): + raise ValueError("runtime_target_json must decode to a JSON object") + return dict(loaded) + raise ValueError("runtime_target must be a JSON object or JSON object string") + + +def _target_env_value( + target: Mapping[str, object], + defaults: Mapping[str, object], + env: Mapping[str, str], + env_name: str, + *, + per_service_mode: bool, + allow_shared_fallback: bool, +) -> str | None: + raw = _first_non_empty( + _target_field(target, defaults, env_name), + _target_field(target, defaults, env_name.lower()), + ) + if raw is None and (not per_service_mode or allow_shared_fallback): + raw = str(env.get(env_name, "") or "").strip() or None + return _coerce_env_value(raw) + + +def _target_field( + target: Mapping[str, object], + defaults: Mapping[str, object], + name: str, +) -> object | None: + for source in (target, _coerce_mapping(target.get("env") or {}), defaults, _coerce_mapping(defaults.get("env") or {})): + if name in source: + return source[name] + return None + + +def _coerce_mapping(value: object) -> Mapping[str, object]: + if not isinstance(value, Mapping): + raise ValueError("Expected a JSON object") + return value + + +def _coerce_env_value(value: object) -> str | None: + if value is None: + return None + if isinstance(value, bool): + return "true" if value else "false" + if isinstance(value, (Mapping, list, tuple)): + return json.dumps(value, separators=(",", ":"), sort_keys=True) + text = str(value).strip() + return text or None + + +def _first_non_empty(*values: object | None) -> object | None: + for value in values: + if value is None: + continue + if isinstance(value, str) and not value.strip(): + continue + return value + return None + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--json", action="store_true", help="Print compact JSON.") + args = parser.parse_args() + + try: + plan = build_sync_plan() + except Exception as exc: + print(str(exc), file=sys.stderr) + return 1 + + if args.json: + print(json.dumps(plan, separators=(",", ":"), sort_keys=True)) + else: + print(json.dumps(plan, indent=2, sort_keys=True)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_build_cloud_run_env_sync_plan.py b/tests/test_build_cloud_run_env_sync_plan.py new file mode 100644 index 0000000..f3b7e6d --- /dev/null +++ b/tests/test_build_cloud_run_env_sync_plan.py @@ -0,0 +1,153 @@ +from __future__ import annotations + +import json +import os +import subprocess +import sys +from pathlib import Path + + +SYNC_PLAN_SCRIPT_PATH = ( + Path(__file__).resolve().parents[1] / "scripts" / "build_cloud_run_env_sync_plan.py" +) + + +def runtime_target_json( + strategy_profile: str, + *, + dry_run_only: bool = True, + platform_id: str = "firstrade", + service_name: str | None = None, +) -> str: + payload: dict[str, object] = { + "platform_id": platform_id, + "strategy_profile": strategy_profile, + "dry_run_only": dry_run_only, + "execution_mode": "paper" if dry_run_only else "live", + } + if service_name is not None: + payload["service_name"] = service_name + return json.dumps(payload, separators=(",", ":")) + + +def test_build_cloud_run_env_sync_plan_legacy_mode_without_telegram(): + env = { + **os.environ, + "CLOUD_RUN_SERVICE": "firstrade-platform-service", + "NOTIFY_LANG": "zh", + "RUNTIME_TARGET_JSON": runtime_target_json( + "tqqq_growth_income", + service_name="firstrade-platform-service", + ), + "FIRSTRADE_MIN_RESERVED_CASH_USD": "250", + "GLOBAL_TELEGRAM_CHAT_ID": "", + } + + result = subprocess.run( + [sys.executable, str(SYNC_PLAN_SCRIPT_PATH), "--json"], + check=True, + capture_output=True, + text=True, + env=env, + ) + + plan = json.loads(result.stdout) + assert plan["mode"] == "legacy" + target = plan["targets"][0] + assert target["service_name"] == "firstrade-platform-service" + assert target["strategy_profile"] == "tqqq_growth_income" + assert target["env"]["NOTIFY_LANG"] == "zh" + assert target["env"]["STRATEGY_PROFILE"] == "tqqq_growth_income" + assert target["env"]["FIRSTRADE_DRY_RUN_ONLY"] == "true" + assert target["env"]["FIRSTRADE_MIN_RESERVED_CASH_USD"] == "250" + assert "GLOBAL_TELEGRAM_CHAT_ID" not in target["env"] + assert "GLOBAL_TELEGRAM_CHAT_ID" in target["remove_env_vars"] + assert "FIRSTRADE_FEATURE_SNAPSHOT_PATH" in target["remove_env_vars"] + assert target["scheduler"]["timezone"] == "America/New_York" + + +def test_build_cloud_run_env_sync_plan_requires_snapshot_for_snapshot_backed_profile(): + env = { + **os.environ, + "CLOUD_RUN_SERVICE": "firstrade-platform-service", + "NOTIFY_LANG": "en", + "RUNTIME_TARGET_JSON": runtime_target_json( + "global_etf_rotation", + service_name="firstrade-platform-service", + ), + "FIRSTRADE_FEATURE_SNAPSHOT_PATH": "gs://stale-paper/snapshot.csv", + } + + result = subprocess.run( + [sys.executable, str(SYNC_PLAN_SCRIPT_PATH), "--json"], + capture_output=True, + text=True, + env=env, + ) + + assert result.returncode != 0 + assert "firstrade-platform-service:FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH" in result.stderr + assert "gs://stale-paper/snapshot.csv" not in result.stderr + + +def test_build_cloud_run_env_sync_plan_skips_snapshot_requirements_when_disabled(): + payload = { + "defaults": {"NOTIFY_LANG": "en"}, + "targets": [ + { + "service": "firstrade-platform-service", + "runtime_target_enabled": "false", + "runtime_target": json.loads( + runtime_target_json( + "global_etf_rotation", + service_name="firstrade-platform-service", + ) + ), + } + ], + } + env = { + **os.environ, + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps(payload), + "FIRSTRADE_FEATURE_SNAPSHOT_PATH": "gs://stale-paper/snapshot.csv", + "FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH": "gs://stale-paper/snapshot.csv.manifest.json", + } + + result = subprocess.run( + [sys.executable, str(SYNC_PLAN_SCRIPT_PATH), "--json"], + check=True, + capture_output=True, + text=True, + env=env, + ) + + plan = json.loads(result.stdout) + assert plan["mode"] == "per_service" + target = plan["targets"][0] + assert target["env"]["RUNTIME_TARGET_ENABLED"] == "false" + assert "FIRSTRADE_FEATURE_SNAPSHOT_PATH" not in target["env"] + assert "FIRSTRADE_FEATURE_SNAPSHOT_MANIFEST_PATH" not in target["env"] + assert "FIRSTRADE_FEATURE_SNAPSHOT_PATH" in target["remove_env_vars"] + assert "gs://stale-paper/snapshot.csv" not in result.stdout + + +def test_build_cloud_run_env_sync_plan_requires_notify_lang(): + env = { + **os.environ, + "CLOUD_RUN_SERVICE": "firstrade-platform-service", + "RUNTIME_TARGET_JSON": runtime_target_json( + "tqqq_growth_income", + service_name="firstrade-platform-service", + ), + } + env.pop("NOTIFY_LANG", None) + + result = subprocess.run( + [sys.executable, str(SYNC_PLAN_SCRIPT_PATH), "--json"], + capture_output=True, + text=True, + env=env, + ) + + assert result.returncode != 0 + assert "firstrade-platform-service:NOTIFY_LANG" in result.stderr diff --git a/tests/test_sync_cloud_run_env_workflow.py b/tests/test_sync_cloud_run_env_workflow.py index 38d68f3..66e35fa 100644 --- a/tests/test_sync_cloud_run_env_workflow.py +++ b/tests/test_sync_cloud_run_env_workflow.py @@ -3,10 +3,18 @@ from pathlib import Path -def test_sync_cloud_run_env_workflow_syncs_strategy_plugin_alert_settings(): +def test_sync_cloud_run_env_workflow_uses_sync_plan_script(): workflow_path = Path(__file__).resolve().parents[1] / ".github/workflows/sync-cloud-run-env.yml" workflow = workflow_path.read_text(encoding="utf-8") + assert "Resolve Cloud Run sync targets" in workflow + assert "scripts/build_cloud_run_env_sync_plan.py --json" in workflow + assert "sync_plan_json<<__SYNC_PLAN_JSON__" in workflow + assert "SYNC_PLAN_JSON: ${{ steps.strategy_requirements.outputs.sync_plan_json }}" in workflow + assert "Cloud Run env sync did not resolve any targets" in workflow + assert "Cloud Run sync target is missing service_name" in workflow + assert "Cloud Run sync target {service_name} is missing env" in workflow + for name in ( "CLOUD_SCHEDULER_LOCATION", "CLOUD_SCHEDULER_MAIN_TIME", @@ -41,11 +49,6 @@ def test_sync_cloud_run_env_workflow_syncs_strategy_plugin_alert_settings(): "STRATEGY_PLUGIN_ALERT_TELEGRAM_PARSE_MODE", "STRATEGY_PLUGIN_ALERT_TELEGRAM_DISABLE_WEB_PAGE_PREVIEW", "STRATEGY_PLUGIN_ALERT_TELEGRAM_BODY_MAX_CHARS", - ): - assert f"{name}: ${{{{ vars.{name} }}}}" in workflow - assert f"add_optional_env {name}" in workflow - - for name in ( "INCOME_LAYER_ENABLED", "INCOME_LAYER_START_USD", "INCOME_LAYER_MAX_RATIO", @@ -69,7 +72,6 @@ def test_sync_cloud_run_env_workflow_syncs_strategy_plugin_alert_settings(): "FIRSTRADE_FEATURE_SNAPSHOT_MAX_STALE_DAYS", ): assert f"{name}: ${{{{ vars.{name} }}}}" in workflow - assert f"add_optional_env {name}" in workflow assert ( "STRATEGY_PLUGIN_ALERT_EMAIL_SENDER_PASSWORD_SECRET_NAME: " @@ -132,9 +134,16 @@ def test_sync_cloud_run_env_workflow_syncs_strategy_plugin_alert_settings(): assert '"CRISIS_ALERT_GOOGLE_VOICE_SMTP_PORT"' in workflow assert '"CRISIS_ALERT_GOOGLE_VOICE_SMTP_SECURITY"' in workflow assert '"CRISIS_ALERT_SMTP_HOST"' in workflow + assert 'env_pairs+=("GOOGLE_CLOUD_PROJECT=${GCP_PROJECT_ID}")' in workflow + assert "for key, value in sorted(target[\"env\"].items()):" in workflow + assert "target.get(\"remove_env_vars\")" in workflow + + assert "add_optional_env " not in workflow + assert "requires_snapshot_artifacts=" not in workflow + assert "Resolve selected strategy runtime requirements" not in workflow -def test_sync_cloud_run_env_workflow_syncs_scheduler_from_runtime_target(): +def test_sync_cloud_run_env_workflow_syncs_scheduler_from_sync_plan(): workflow_path = Path(__file__).resolve().parents[1] / ".github/workflows/sync-cloud-run-env.yml" workflow = workflow_path.read_text(encoding="utf-8") @@ -142,12 +151,12 @@ def test_sync_cloud_run_env_workflow_syncs_scheduler_from_runtime_target(): assert "GCP_SCHEDULER_SERVICE_ACCOUNT: firstrade-platform-scheduler@firstradequant.iam.gserviceaccount.com" in workflow assert "MONITOR_DISPATCH_TARGETS_JSON=${monitor_targets_json}" in workflow assert 'scheduler_location="${CLOUD_SCHEDULER_LOCATION:-${CLOUD_RUN_REGION}}"' in workflow - assert 'raw_runtime_target = os.environ.get("RUNTIME_TARGET_JSON", "").strip()' in workflow - assert 'scheduler = runtime_target.get("scheduler") if isinstance(runtime_target, dict) else {}' in workflow - assert 'print(str(runtime_scheduler.get("timezone") or "America/New_York").strip())' in workflow - assert 'configured_time("main_time", "CLOUD_SCHEDULER_MAIN_TIME", "45 15")' in workflow - assert 'configured_time("probe_time", "CLOUD_SCHEDULER_PROBE_TIME", "35 9,15")' in workflow - assert 'configured_time("precheck_time", "CLOUD_SCHEDULER_PRECHECK_TIME", "45 9")' in workflow + assert 'plan = json.loads(os.environ["SYNC_PLAN_JSON"])' in workflow + assert 'scheduler = target.get("scheduler") or {}' in workflow + assert 'print(str(scheduler.get("timezone") or "America/New_York").strip())' in workflow + assert 'scheduler.get("main_time") or os.environ.get("CLOUD_SCHEDULER_MAIN_TIME"' in workflow + assert 'scheduler.get("probe_time") or os.environ.get("CLOUD_SCHEDULER_PROBE_TIME"' in workflow + assert 'scheduler.get("precheck_time") or os.environ.get("CLOUD_SCHEDULER_PRECHECK_TIME"' in workflow assert 'scheduler_job_candidates=("${CLOUD_RUN_SERVICE}-scheduler")' in workflow assert 'scheduler_job_candidates+=("${CLOUD_RUN_SERVICE%-service}-scheduler")' in workflow assert 'if len(time_fields) == 5:' in workflow