Skip to content

Commit bfabb1e

Browse files
authored
Add shared monitor dispatcher scheduler
1 parent 201e6c4 commit bfabb1e

6 files changed

Lines changed: 602 additions & 26 deletions

File tree

.github/workflows/sync-cloud-run-env.yml

Lines changed: 206 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@ env:
5353
GCP_WORKLOAD_IDENTITY_PROVIDER: projects/252919773759/locations/global/workloadIdentityPools/github-actions/providers/github-main
5454
GCP_WORKLOAD_IDENTITY_SERVICE_ACCOUNT: longbridge-platform-deploy@longbridgequant.iam.gserviceaccount.com
5555
GCP_RUNTIME_SERVICE_ACCOUNT: longbridge-platform-runtime@longbridgequant.iam.gserviceaccount.com
56+
GCP_SCHEDULER_SERVICE_ACCOUNT: longbridge-platform-scheduler@longbridgequant.iam.gserviceaccount.com
5657
GCP_ARTIFACT_REGISTRY_REPOSITORY: cloud-run-source-deploy
5758

5859
concurrency:
@@ -99,11 +100,13 @@ jobs:
99100
# Set CLOUD_RUN_REGION per Environment so paper/HK/SG can target different regions.
100101
CLOUD_RUN_REGION: ${{ vars.CLOUD_RUN_REGION }}
101102
CLOUD_RUN_SERVICE: ${{ vars.CLOUD_RUN_SERVICE }}
103+
CLOUD_RUN_SERVICE_TARGETS_JSON: ${{ vars.CLOUD_RUN_SERVICE_TARGETS_JSON }}
102104
CLOUD_RUN_ENV_SYNC_WAIT_FOR_COMMIT: ${{ vars.CLOUD_RUN_ENV_SYNC_WAIT_FOR_COMMIT }}
103105
CLOUD_SCHEDULER_LOCATION: ${{ vars.CLOUD_SCHEDULER_LOCATION }}
104106
CLOUD_SCHEDULER_MAIN_TIME: ${{ vars.CLOUD_SCHEDULER_MAIN_TIME }}
105107
CLOUD_SCHEDULER_PROBE_TIME: ${{ vars.CLOUD_SCHEDULER_PROBE_TIME }}
106108
CLOUD_SCHEDULER_PRECHECK_TIME: ${{ vars.CLOUD_SCHEDULER_PRECHECK_TIME }}
109+
MONITOR_DISPATCHER_OWNER_LABEL: ${{ vars.MONITOR_DISPATCHER_OWNER_LABEL || 'SG' }}
107110
ACCOUNT_PREFIX: ${{ vars.ACCOUNT_PREFIX }}
108111
TELEGRAM_TOKEN_SECRET_NAME: ${{ vars.TELEGRAM_TOKEN_SECRET_NAME }}
109112
LONGPORT_APP_KEY_SECRET_NAME: ${{ vars.LONGPORT_APP_KEY_SECRET_NAME }}
@@ -961,6 +964,89 @@ jobs:
961964
remove_env_vars+=("RUNTIME_TARGET_ENABLED")
962965
fi
963966
967+
monitor_targets_json="$(
968+
python - <<'PY'
969+
import json
970+
import os
971+
import subprocess
972+
973+
def decode_json(raw, fallback):
974+
raw = str(raw or "").strip()
975+
if not raw:
976+
return fallback
977+
try:
978+
return json.loads(raw)
979+
except json.JSONDecodeError:
980+
return fallback
981+
982+
def runtime_target_from(source):
983+
runtime_target = source.get("runtime_target") or source.get("runtime_target_json")
984+
if isinstance(runtime_target, str):
985+
runtime_target = decode_json(runtime_target, {})
986+
return runtime_target if isinstance(runtime_target, dict) else {}
987+
988+
raw_targets = decode_json(os.environ.get("CLOUD_RUN_SERVICE_TARGETS_JSON"), None)
989+
if isinstance(raw_targets, dict):
990+
source_targets = raw_targets.get("targets")
991+
else:
992+
source_targets = raw_targets
993+
if not isinstance(source_targets, list) or not source_targets:
994+
source_targets = [
995+
{
996+
"service": os.environ.get("CLOUD_RUN_SERVICE"),
997+
"region": os.environ.get("CLOUD_RUN_REGION"),
998+
"runtime_target": decode_json(os.environ.get("RUNTIME_TARGET_JSON"), {}),
999+
"runtime_target_enabled": os.environ.get("RUNTIME_TARGET_ENABLED", "true"),
1000+
}
1001+
]
1002+
1003+
project = os.environ.get("GCP_PROJECT_ID", "").strip()
1004+
output = []
1005+
for source in source_targets:
1006+
if not isinstance(source, dict):
1007+
continue
1008+
runtime_target = runtime_target_from(source)
1009+
service_name = (
1010+
source.get("service")
1011+
or source.get("service_name")
1012+
or source.get("cloud_run_service")
1013+
or runtime_target.get("service_name")
1014+
)
1015+
region = source.get("region") or source.get("cloud_run_region") or os.environ.get("CLOUD_RUN_REGION")
1016+
service_name = str(service_name or "").strip()
1017+
region = str(region or "").strip()
1018+
if not service_name or not region:
1019+
continue
1020+
service_url = subprocess.check_output(
1021+
[
1022+
"gcloud",
1023+
"run",
1024+
"services",
1025+
"describe",
1026+
service_name,
1027+
f"--project={project}",
1028+
f"--region={region}",
1029+
"--format=value(status.url)",
1030+
],
1031+
text=True,
1032+
).strip()
1033+
if not service_url:
1034+
continue
1035+
output.append(
1036+
{
1037+
"service_name": service_name,
1038+
"service_url": service_url,
1039+
"strategy_profile": runtime_target.get("strategy_profile") or source.get("STRATEGY_PROFILE"),
1040+
"account_scope": runtime_target.get("account_scope") or source.get("ACCOUNT_REGION") or source.get("account_scope"),
1041+
"runtime_target_enabled": source.get("runtime_target_enabled", source.get("RUNTIME_TARGET_ENABLED", True)),
1042+
"scheduler": runtime_target.get("scheduler") if isinstance(runtime_target.get("scheduler"), dict) else {},
1043+
}
1044+
)
1045+
print(json.dumps({"targets": output}, separators=(",", ":")))
1046+
PY
1047+
)"
1048+
env_pairs+=("MONITOR_DISPATCH_TARGETS_JSON=${monitor_targets_json}")
1049+
9641050
gcloud_args=(
9651051
run services update "${CLOUD_RUN_SERVICE}"
9661052
--region "${CLOUD_RUN_REGION}"
@@ -1042,33 +1128,42 @@ jobs:
10421128
exit 1
10431129
fi
10441130
1045-
for suffix in scheduler probe-scheduler precheck-scheduler; do
1046-
job_name="${CLOUD_RUN_SERVICE}-${suffix}"
1047-
case "${suffix}" in
1048-
scheduler)
1049-
schedule_time="${main_time}"
1050-
scheduler_path="/run"
1051-
;;
1052-
probe-scheduler)
1053-
schedule_time="${probe_time}"
1054-
scheduler_path="/probe"
1055-
;;
1056-
precheck-scheduler)
1057-
schedule_time="${precheck_time}"
1058-
scheduler_path="/dry-run"
1059-
;;
1060-
esac
1061-
1062-
current_schedule="$(gcloud scheduler jobs describe "${job_name}" \
1131+
gcloud run services add-iam-policy-binding "${CLOUD_RUN_SERVICE}" \
1132+
--project="${GCP_PROJECT_ID}" \
1133+
--region="${CLOUD_RUN_REGION}" \
1134+
--member="serviceAccount:${GCP_SCHEDULER_SERVICE_ACCOUNT}" \
1135+
--role="roles/run.invoker" \
1136+
--quiet
1137+
gcloud run services add-iam-policy-binding "${CLOUD_RUN_SERVICE}" \
1138+
--project="${GCP_PROJECT_ID}" \
1139+
--region="${CLOUD_RUN_REGION}" \
1140+
--member="serviceAccount:${GCP_RUNTIME_SERVICE_ACCOUNT}" \
1141+
--role="roles/run.invoker" \
1142+
--quiet
1143+
1144+
scheduler_job_candidates=("${CLOUD_RUN_SERVICE}-scheduler")
1145+
if [[ "${CLOUD_RUN_SERVICE}" == *-service ]]; then
1146+
scheduler_job_candidates+=("${CLOUD_RUN_SERVICE%-service}-scheduler")
1147+
fi
1148+
1149+
job_name=""
1150+
current_schedule=""
1151+
for candidate_job in "${scheduler_job_candidates[@]}"; do
1152+
current_schedule="$(gcloud scheduler jobs describe "${candidate_job}" \
10631153
--project="${GCP_PROJECT_ID}" \
10641154
--location="${scheduler_location}" \
10651155
--format='value(schedule)' 2>/dev/null || true)"
1066-
if [ -z "${current_schedule}" ]; then
1067-
echo "Cloud Scheduler job ${job_name} was not found in ${scheduler_location}; skipping schedule sync."
1068-
continue
1156+
if [ -n "${current_schedule}" ]; then
1157+
job_name="${candidate_job}"
1158+
break
10691159
fi
1160+
done
1161+
if [ -z "${job_name}" ]; then
1162+
job_name="${scheduler_job_candidates[0]}"
1163+
fi
10701164
1071-
desired_schedule="$(CURRENT_SCHEDULE="${current_schedule}" SCHEDULE_TIME="${schedule_time}" python - <<'PY'
1165+
if [ -n "${current_schedule}" ]; then
1166+
desired_schedule="$(CURRENT_SCHEDULE="${current_schedule}" SCHEDULE_TIME="${main_time}" python - <<'PY'
10721167
import os
10731168
10741169
current_fields = os.environ["CURRENT_SCHEDULE"].split()
@@ -1085,8 +1180,25 @@ jobs:
10851180
)
10861181
PY
10871182
)"
1183+
else
1184+
desired_schedule="$(SCHEDULE_TIME="${main_time}" python - <<'PY'
1185+
import os
1186+
1187+
fields = os.environ["SCHEDULE_TIME"].split()
1188+
if len(fields) == 5:
1189+
print(" ".join(fields))
1190+
elif len(fields) == 2:
1191+
print(" ".join([*fields, "*", "*", "*"]))
1192+
else:
1193+
raise SystemExit(
1194+
f"Cloud Scheduler override must have 2 time fields or 5 cron fields: {os.environ['SCHEDULE_TIME']!r}"
1195+
)
1196+
PY
1197+
)"
1198+
fi
10881199
1089-
scheduler_uri="${service_url}${scheduler_path}"
1200+
scheduler_uri="${service_url}/run"
1201+
if [ -n "${current_schedule}" ]; then
10901202
echo "Updating Cloud Scheduler job ${job_name} schedule to ${desired_schedule}, timezone to ${market_timezone}, and URI to ${scheduler_uri}."
10911203
gcloud scheduler jobs update http "${job_name}" \
10921204
--project="${GCP_PROJECT_ID}" \
@@ -1095,6 +1207,77 @@ jobs:
10951207
--schedule="${desired_schedule}" \
10961208
--time-zone="${market_timezone}" \
10971209
--quiet
1210+
else
1211+
echo "Creating Cloud Scheduler job ${job_name} schedule ${desired_schedule}, timezone ${market_timezone}, and URI ${scheduler_uri}."
1212+
gcloud scheduler jobs create http "${job_name}" \
1213+
--project="${GCP_PROJECT_ID}" \
1214+
--location="${scheduler_location}" \
1215+
--uri="${scheduler_uri}" \
1216+
--http-method=POST \
1217+
--oidc-service-account-email="${GCP_SCHEDULER_SERVICE_ACCOUNT}" \
1218+
--oidc-token-audience="${service_url}" \
1219+
--schedule="${desired_schedule}" \
1220+
--time-zone="${market_timezone}" \
1221+
--attempt-deadline=600s \
1222+
--quiet
1223+
fi
1224+
1225+
if [ "${DEPLOYMENT_LABEL:-}" = "${MONITOR_DISPATCHER_OWNER_LABEL:-SG}" ]; then
1226+
monitor_job_name="longbridge-monitor-dispatcher-scheduler"
1227+
monitor_uri="${service_url}/monitor-dispatch"
1228+
if gcloud scheduler jobs describe "${monitor_job_name}" \
1229+
--project="${GCP_PROJECT_ID}" \
1230+
--location="${scheduler_location}" >/dev/null 2>&1; then
1231+
echo "Updating Cloud Scheduler job ${monitor_job_name} to ${monitor_uri}."
1232+
gcloud scheduler jobs update http "${monitor_job_name}" \
1233+
--project="${GCP_PROJECT_ID}" \
1234+
--location="${scheduler_location}" \
1235+
--uri="${monitor_uri}" \
1236+
--http-method=POST \
1237+
--oidc-service-account-email="${GCP_SCHEDULER_SERVICE_ACCOUNT}" \
1238+
--oidc-token-audience="${service_url}" \
1239+
--schedule="*/5 * * * *" \
1240+
--time-zone="UTC" \
1241+
--attempt-deadline=180s \
1242+
--quiet
1243+
else
1244+
echo "Creating Cloud Scheduler job ${monitor_job_name} at ${monitor_uri}."
1245+
gcloud scheduler jobs create http "${monitor_job_name}" \
1246+
--project="${GCP_PROJECT_ID}" \
1247+
--location="${scheduler_location}" \
1248+
--uri="${monitor_uri}" \
1249+
--http-method=POST \
1250+
--oidc-service-account-email="${GCP_SCHEDULER_SERVICE_ACCOUNT}" \
1251+
--oidc-token-audience="${service_url}" \
1252+
--schedule="*/5 * * * *" \
1253+
--time-zone="UTC" \
1254+
--attempt-deadline=180s \
1255+
--quiet
1256+
fi
1257+
else
1258+
echo "Skipping shared LongBridge monitor dispatcher scheduler from ${DEPLOYMENT_LABEL}; owner is ${MONITOR_DISPATCHER_OWNER_LABEL:-SG}."
1259+
fi
1260+
1261+
legacy_jobs=(
1262+
"${CLOUD_RUN_SERVICE}-probe-scheduler"
1263+
"${CLOUD_RUN_SERVICE}-precheck-scheduler"
1264+
)
1265+
if [[ "${CLOUD_RUN_SERVICE}" == *-service ]]; then
1266+
legacy_jobs+=(
1267+
"${CLOUD_RUN_SERVICE%-service}-probe-scheduler"
1268+
"${CLOUD_RUN_SERVICE%-service}-precheck-scheduler"
1269+
)
1270+
fi
1271+
for legacy_job in "${legacy_jobs[@]}"; do
1272+
if gcloud scheduler jobs describe "${legacy_job}" \
1273+
--project="${GCP_PROJECT_ID}" \
1274+
--location="${scheduler_location}" >/dev/null 2>&1; then
1275+
echo "Deleting legacy Cloud Scheduler job ${legacy_job}; monitor dispatcher now owns probe/precheck."
1276+
gcloud scheduler jobs delete "${legacy_job}" \
1277+
--project="${GCP_PROJECT_ID}" \
1278+
--location="${scheduler_location}" \
1279+
--quiet
1280+
fi
10981281
done
10991282
11001283
- name: Prune old Cloud Run revisions

0 commit comments

Comments
 (0)