e2e-verify-workflow-terminate.sh 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458
  1. #!/usr/bin/env bash
  2. set -euo pipefail
  3. ROOT_DIR="$(cd "$(dirname "$0")/.." && pwd)"
  4. cd "$ROOT_DIR"
  5. BACKEND_URL="${E2E_BACKEND_URL:-http://127.0.0.1:8080}"
  6. ORCH_URL="${E2E_ORCH_URL:-http://127.0.0.1:18080}"
  7. POLL_INTERVAL_SECONDS="${E2E_POLL_INTERVAL_SECONDS:-1}"
  8. log() {
  9. printf '[INFO] %s\n' "$*"
  10. }
  11. fail() {
  12. printf '[FAIL] %s\n' "$*" >&2
  13. exit 1
  14. }
  15. require_cmd() {
  16. if ! command -v "$1" >/dev/null 2>&1; then
  17. fail "command not found: $1"
  18. fi
  19. }
  20. require_cmd docker
  21. require_cmd kubectl
  22. require_cmd curl
  23. require_cmd jq
  24. require_cmd awk
  25. BACKEND_ENV_DUMP="$(docker inspect ws-backend --format '{{range .Config.Env}}{{println .}}{{end}}' 2>/dev/null || true)"
  26. if [[ -z "$BACKEND_ENV_DUMP" ]]; then
  27. fail "cannot read backend container env; ensure container ws-backend is running"
  28. fi
  29. backend_env_get() {
  30. local key="$1"
  31. printf '%s\n' "$BACKEND_ENV_DUMP" | awk -F= -v k="$key" '$1==k {sub(/^[^=]*=/, "", $0); print; exit}'
  32. }
  33. MYSQL_URL="$(backend_env_get MYSQL_URL)"
  34. MYSQL_DB="$(backend_env_get MYSQL_DATABASE)"
  35. MYSQL_USER="$(backend_env_get MYSQL_USER)"
  36. MYSQL_PASSWORD="$(backend_env_get MYSQL_PASSWORD)"
  37. if [[ -z "$MYSQL_URL" || -z "$MYSQL_DB" || -z "$MYSQL_USER" ]]; then
  38. fail "backend mysql env is incomplete (MYSQL_URL/MYSQL_DATABASE/MYSQL_USER)"
  39. fi
  40. if [[ "$MYSQL_URL" == *:* ]]; then
  41. DB_HOST="${MYSQL_URL%:*}"
  42. DB_PORT="${MYSQL_URL##*:}"
  43. else
  44. DB_HOST="$MYSQL_URL"
  45. DB_PORT="3306"
  46. fi
  47. mysql_exec() {
  48. local sql="$1"
  49. docker compose exec -T mysql sh -lc \
  50. "mysql -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB' -e \"$sql\""
  51. }
  52. mysql_query() {
  53. local sql="$1"
  54. docker compose exec -T mysql sh -lc \
  55. "mysql -N -B -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB' -e \"$sql\""
  56. }
  57. api_register_token() {
  58. local user="e2e_workflow_term_$(date +%s)_$RANDOM"
  59. local pass='E2eWorkflow@123456'
  60. local payload
  61. payload=$(jq -cn \
  62. --arg username "$user" \
  63. --arg password "$pass" \
  64. --arg email "$user@example.com" \
  65. '{username:$username,password:$password,email:$email,role:"OPS"}')
  66. local resp
  67. resp=$(curl -fsS -X POST "$BACKEND_URL/api/auth/register" \
  68. -H 'Content-Type: application/json' \
  69. -d "$payload")
  70. local token
  71. token=$(echo "$resp" | jq -r '.token // empty')
  72. if [[ -z "$token" ]]; then
  73. fail "register token failed: $resp"
  74. fi
  75. printf '%s' "$token"
  76. }
  77. API_TOKEN="$(api_register_token)"
  78. AUTH_HEADER="Authorization: Bearer $API_TOKEN"
  79. api_get() {
  80. local path="$1"
  81. curl -fsS "$BACKEND_URL$path" -H "$AUTH_HEADER"
  82. }
  83. api_post() {
  84. local path="$1"
  85. local body="${2:-{}}"
  86. curl -fsS -X POST "$BACKEND_URL$path" \
  87. -H "$AUTH_HEADER" \
  88. -H 'Content-Type: application/json' \
  89. -d "$body"
  90. }
  91. terminate_workflow() {
  92. local workflow_instance_id="$1"
  93. local code
  94. code=$(curl -sS -o /tmp/e2e_terminate_resp.$$ -w '%{http_code}' \
  95. -X POST "$BACKEND_URL/api/task-exec/workflows/instances/$workflow_instance_id/terminate" \
  96. -H "$AUTH_HEADER")
  97. if [[ "$code" != "204" ]]; then
  98. local body
  99. body="$(cat /tmp/e2e_terminate_resp.$$ 2>/dev/null || true)"
  100. fail "terminate api failed, workflowInstanceId=$workflow_instance_id, status=$code, body=$body"
  101. fi
  102. }
  103. get_running_cluster_id() {
  104. local component="$1"
  105. local id
  106. id="$(mysql_query "select cluster_id from cluster where component_type='$component' and status='RUNNING' order by cluster_id desc limit 1;" | head -n1)"
  107. if [[ -z "$id" ]]; then
  108. fail "no RUNNING cluster found for component_type=$component"
  109. fi
  110. printf '%s' "$id"
  111. }
  112. wait_task_running_with_engine_id() {
  113. local workflow_instance_id="$1"
  114. local task_id="$2"
  115. local timeout_seconds="$3"
  116. local deadline=$((SECONDS + timeout_seconds))
  117. while (( SECONDS < deadline )); do
  118. local row
  119. row="$(mysql_query "select state,coalesce(engine_task_id,'') from task_instance where workflow_instance_id=$workflow_instance_id and task_id=$task_id limit 1;" | head -n1 || true)"
  120. if [[ -n "$row" ]]; then
  121. local state
  122. local engine_id
  123. state="$(echo "$row" | cut -f1)"
  124. engine_id="$(echo "$row" | cut -f2)"
  125. if [[ "$state" == "RUNNING" && -n "$engine_id" ]]; then
  126. printf '%s' "$engine_id"
  127. return 0
  128. fi
  129. fi
  130. sleep "$POLL_INTERVAL_SECONDS"
  131. done
  132. fail "task did not reach RUNNING with engine id in time, workflowInstanceId=$workflow_instance_id, taskId=$task_id"
  133. }
  134. wait_terminated_and_states() {
  135. local workflow_instance_id="$1"
  136. local running_task_id="$2"
  137. local downstream_task_id="$3"
  138. local timeout_seconds="$4"
  139. local deadline=$((SECONDS + timeout_seconds))
  140. while (( SECONDS < deadline )); do
  141. local status_json
  142. status_json="$(api_get "/api/task-exec/workflows/instances/$workflow_instance_id")"
  143. local wf_state
  144. local running_state
  145. local downstream_state
  146. wf_state="$(echo "$status_json" | jq -r '.state // empty')"
  147. running_state="$(echo "$status_json" | jq -r --argjson taskId "$running_task_id" '.taskInstances[] | select(.taskId==$taskId) | .state // empty')"
  148. downstream_state="$(echo "$status_json" | jq -r --argjson taskId "$downstream_task_id" '.taskInstances[] | select(.taskId==$taskId) | .state // empty')"
  149. if [[ "$wf_state" == "TERMINATED" && "$running_state" == "KILLED" && "$downstream_state" == "SKIPPED" ]]; then
  150. return 0
  151. fi
  152. sleep "$POLL_INTERVAL_SECONDS"
  153. done
  154. fail "workflow/task states mismatch after terminate, workflowInstanceId=$workflow_instance_id"
  155. }
  156. wait_spark_job_ref() {
  157. local operation_id="$1"
  158. local timeout_seconds="$2"
  159. local deadline=$((SECONDS + timeout_seconds))
  160. while (( SECONDS < deadline )); do
  161. local resp
  162. resp="$(curl -fsS "$ORCH_URL/v1/spark-jobs/operations/$operation_id")"
  163. local ns
  164. local app
  165. ns="$(echo "$resp" | jq -r '.namespace // empty')"
  166. app="$(echo "$resp" | jq -r '.application // empty')"
  167. if [[ -n "$ns" && -n "$app" ]]; then
  168. printf '%s\t%s' "$ns" "$app"
  169. return 0
  170. fi
  171. sleep "$POLL_INTERVAL_SECONDS"
  172. done
  173. fail "spark job namespace/application not ready, operationId=$operation_id"
  174. }
  175. wait_spark_app_exists() {
  176. local namespace="$1"
  177. local app="$2"
  178. local timeout_seconds="$3"
  179. local deadline=$((SECONDS + timeout_seconds))
  180. while (( SECONDS < deadline )); do
  181. if kubectl -n "$namespace" get sparkapplications.sparkoperator.k8s.io "$app" >/dev/null 2>&1; then
  182. return 0
  183. fi
  184. sleep "$POLL_INTERVAL_SECONDS"
  185. done
  186. fail "spark application not found before terminate, namespace=$namespace, app=$app"
  187. }
  188. wait_spark_app_deleted() {
  189. local namespace="$1"
  190. local app="$2"
  191. local timeout_seconds="$3"
  192. local deadline=$((SECONDS + timeout_seconds))
  193. while (( SECONDS < deadline )); do
  194. if kubectl -n "$namespace" get sparkapplications.sparkoperator.k8s.io "$app" >/dev/null 2>&1; then
  195. sleep "$POLL_INTERVAL_SECONDS"
  196. continue
  197. fi
  198. local pod_count
  199. pod_count="$(kubectl -n "$namespace" get pods --no-headers 2>/dev/null | awk -v app="$app" '$1 ~ app {c++} END {print c+0}')"
  200. if [[ "$pod_count" == "0" ]]; then
  201. return 0
  202. fi
  203. sleep "$POLL_INTERVAL_SECONDS"
  204. done
  205. fail "spark application resources still exist, namespace=$namespace, app=$app"
  206. }
  207. run_starrocks_sql() {
  208. local sql="$1"
  209. docker compose exec -T \
  210. -e E2E_SR_HOST="$STARROCKS_HOST" \
  211. -e E2E_SR_PORT="$STARROCKS_PORT" \
  212. -e E2E_SR_USER="$STARROCKS_USER" \
  213. -e E2E_SR_PASSWORD="$STARROCKS_PASSWORD" \
  214. -e E2E_SR_SQL="$sql" \
  215. mysql sh -lc '
  216. if [ -n "$E2E_SR_PASSWORD" ]; then
  217. mysql --connect-timeout=10 --protocol=TCP \
  218. --host="$E2E_SR_HOST" \
  219. --port="$E2E_SR_PORT" \
  220. --user="$E2E_SR_USER" \
  221. --password="$E2E_SR_PASSWORD" \
  222. --batch --raw \
  223. -e "$E2E_SR_SQL"
  224. else
  225. mysql --connect-timeout=10 --protocol=TCP \
  226. --host="$E2E_SR_HOST" \
  227. --port="$E2E_SR_PORT" \
  228. --user="$E2E_SR_USER" \
  229. --batch --raw \
  230. -e "$E2E_SR_SQL"
  231. fi
  232. '
  233. }
  234. wait_starrocks_query_present() {
  235. local tag="$1"
  236. local timeout_seconds="$2"
  237. local deadline=$((SECONDS + timeout_seconds))
  238. while (( SECONDS < deadline )); do
  239. local processlist
  240. processlist="$(run_starrocks_sql "SHOW PROCESSLIST;" || true)"
  241. if echo "$processlist" | awk -F'\t' -v tag="$tag" 'NR>1 && index($10, tag)>0 {found=1} END {exit found?0:1}'; then
  242. return 0
  243. fi
  244. sleep "$POLL_INTERVAL_SECONDS"
  245. done
  246. fail "starrocks target query not found in processlist, tag=$tag"
  247. }
  248. wait_starrocks_query_absent() {
  249. local tag="$1"
  250. local timeout_seconds="$2"
  251. local deadline=$((SECONDS + timeout_seconds))
  252. while (( SECONDS < deadline )); do
  253. local processlist
  254. processlist="$(run_starrocks_sql "SHOW PROCESSLIST;" || true)"
  255. if echo "$processlist" | awk -F'\t' -v tag="$tag" 'NR>1 && index($10, tag)>0 {found=1} END {exit found?1:0}'; then
  256. return 0
  257. fi
  258. sleep "$POLL_INTERVAL_SECONDS"
  259. done
  260. fail "starrocks query still present after terminate, tag=$tag"
  261. }
  262. main() {
  263. log "checking backend and orchestrator reachability"
  264. curl -fsS "$BACKEND_URL/api/auth/login" -X OPTIONS >/dev/null 2>&1 || true
  265. curl -fsS "$ORCH_URL/v1/operations/non-exist" >/dev/null 2>&1 || true
  266. SPARK_CLUSTER_ID="$(get_running_cluster_id SPARK)"
  267. STARROCKS_CLUSTER_ID="$(get_running_cluster_id STARROCKS)"
  268. log "using clusters: spark=$SPARK_CLUSTER_ID, starrocks=$STARROCKS_CLUSTER_ID"
  269. local sr_cluster_name="starrocks-$STARROCKS_CLUSTER_ID"
  270. local sr_endpoint_json
  271. sr_endpoint_json="$(curl -fsS "$ORCH_URL/v1/starrocks-clusters/$sr_cluster_name/direct-endpoint")"
  272. STARROCKS_HOST="$(echo "$sr_endpoint_json" | jq -r '.host // empty')"
  273. STARROCKS_PORT="$(echo "$sr_endpoint_json" | jq -r '.port // empty')"
  274. STARROCKS_POD_NAMESPACE="$(echo "$sr_endpoint_json" | jq -r '.namespace // empty')"
  275. STARROCKS_POD_NAME="$(echo "$sr_endpoint_json" | jq -r '.podName // empty')"
  276. STARROCKS_USER="$(backend_env_get STARROCKS_SQL_RUNNER_USERNAME)"
  277. STARROCKS_PASSWORD="$(backend_env_get STARROCKS_SQL_RUNNER_PASSWORD)"
  278. if [[ -z "$STARROCKS_HOST" || -z "$STARROCKS_PORT" || -z "$STARROCKS_USER" ]]; then
  279. fail "starrocks direct endpoint or auth config is incomplete"
  280. fi
  281. local run_epoch
  282. run_epoch="$(date +%s)"
  283. local base_id=$((1900000000 + (run_epoch % 1000000) * 1000 + (RANDOM % 1000)))
  284. local spark_wf_id=$((base_id + 1))
  285. local spark_task_running=$((base_id + 11))
  286. local spark_task_downstream=$((base_id + 12))
  287. local sr_wf_id=$((base_id + 101))
  288. local sr_task_running=$((base_id + 111))
  289. local sr_task_downstream=$((base_id + 112))
  290. local sr_tag="e2e_sr_${run_epoch}_${RANDOM}"
  291. log "preparing spark terminate e2e workflow"
  292. cat <<SQL | docker compose exec -T mysql sh -lc "mysql -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB'"
  293. INSERT INTO workflow_definition(workflow_id,workflow_name,description,dag_json,timeout_seconds,create_time,failure_strategy)
  294. VALUES (
  295. $spark_wf_id,
  296. 'E2E_SPARK_TERMINATE',
  297. 'spark terminate e2e',
  298. '{"nodes":[{"taskId":$spark_task_running},{"taskId":$spark_task_downstream}],"edges":[{"fromTaskId":$spark_task_running,"toTaskId":$spark_task_downstream}]}',
  299. 1800,
  300. NOW(),
  301. 'TERMINATE'
  302. )
  303. ON DUPLICATE KEY UPDATE
  304. workflow_name = VALUES(workflow_name),
  305. description = VALUES(description),
  306. dag_json = VALUES(dag_json),
  307. timeout_seconds = VALUES(timeout_seconds),
  308. failure_strategy = VALUES(failure_strategy);
  309. INSERT INTO task_definition(task_id,task_name,workflow_id,task_type,task_content,exector_id,timeout_seconds,retry_times)
  310. VALUES
  311. ($spark_task_running,'e2e_spark_running',$spark_wf_id,'SPARK_SQL','SELECT COUNT(*) AS c FROM range(1, 1000000000000);',$SPARK_CLUSTER_ID,1800,0),
  312. ($spark_task_downstream,'e2e_spark_downstream',$spark_wf_id,'SPARK_SQL','SELECT 1;',$SPARK_CLUSTER_ID,1800,0)
  313. ON DUPLICATE KEY UPDATE
  314. task_name = VALUES(task_name),
  315. task_type = VALUES(task_type),
  316. task_content = VALUES(task_content),
  317. exector_id = VALUES(exector_id),
  318. timeout_seconds = VALUES(timeout_seconds),
  319. retry_times = VALUES(retry_times);
  320. SQL
  321. local spark_submit_resp
  322. spark_submit_resp="$(api_post "/api/task-exec/workflows/$spark_wf_id/submit" '{}')"
  323. local spark_wfi
  324. spark_wfi="$(echo "$spark_submit_resp" | jq -r '.workflowInstanceId // empty')"
  325. if [[ -z "$spark_wfi" ]]; then
  326. fail "spark submit failed: $spark_submit_resp"
  327. fi
  328. log "spark workflow submitted: workflowInstanceId=$spark_wfi"
  329. local spark_operation_id
  330. spark_operation_id="$(wait_task_running_with_engine_id "$spark_wfi" "$spark_task_running" 180)"
  331. log "spark task is RUNNING: operationId=$spark_operation_id"
  332. local spark_ref
  333. spark_ref="$(wait_spark_job_ref "$spark_operation_id" 120)"
  334. local spark_namespace
  335. local spark_app
  336. IFS=$'\t' read -r spark_namespace spark_app <<<"$spark_ref"
  337. wait_spark_app_exists "$spark_namespace" "$spark_app" 120
  338. log "spark application observed in k8s: namespace=$spark_namespace, app=$spark_app"
  339. terminate_workflow "$spark_wfi"
  340. wait_terminated_and_states "$spark_wfi" "$spark_task_running" "$spark_task_downstream" 60
  341. wait_spark_app_deleted "$spark_namespace" "$spark_app" 120
  342. log "spark terminate verified: running task killed, downstream skipped, spark resources released"
  343. log "preparing starrocks terminate e2e workflow"
  344. cat <<SQL | docker compose exec -T mysql sh -lc "mysql -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB'"
  345. INSERT INTO workflow_definition(workflow_id,workflow_name,description,dag_json,timeout_seconds,create_time,failure_strategy)
  346. VALUES (
  347. $sr_wf_id,
  348. 'E2E_STARROCKS_TERMINATE',
  349. 'starrocks terminate e2e',
  350. '{"nodes":[{"taskId":$sr_task_running},{"taskId":$sr_task_downstream}],"edges":[{"fromTaskId":$sr_task_running,"toTaskId":$sr_task_downstream}]}',
  351. 1800,
  352. NOW(),
  353. 'TERMINATE'
  354. )
  355. ON DUPLICATE KEY UPDATE
  356. workflow_name = VALUES(workflow_name),
  357. description = VALUES(description),
  358. dag_json = VALUES(dag_json),
  359. timeout_seconds = VALUES(timeout_seconds),
  360. failure_strategy = VALUES(failure_strategy);
  361. INSERT INTO task_definition(task_id,task_name,workflow_id,task_type,task_content,exector_id,timeout_seconds,retry_times)
  362. VALUES
  363. ($sr_task_running,'e2e_sr_running',$sr_wf_id,'STARROCKS_SQL','SELECT /*$sr_tag*/ SLEEP(300);',$STARROCKS_CLUSTER_ID,1800,0),
  364. ($sr_task_downstream,'e2e_sr_downstream',$sr_wf_id,'STARROCKS_SQL','SELECT 1;',$STARROCKS_CLUSTER_ID,1800,0)
  365. ON DUPLICATE KEY UPDATE
  366. task_name = VALUES(task_name),
  367. task_type = VALUES(task_type),
  368. task_content = VALUES(task_content),
  369. exector_id = VALUES(exector_id),
  370. timeout_seconds = VALUES(timeout_seconds),
  371. retry_times = VALUES(retry_times);
  372. SQL
  373. local sr_submit_resp
  374. sr_submit_resp="$(api_post "/api/task-exec/workflows/$sr_wf_id/submit" '{}')"
  375. local sr_wfi
  376. sr_wfi="$(echo "$sr_submit_resp" | jq -r '.workflowInstanceId // empty')"
  377. if [[ -z "$sr_wfi" ]]; then
  378. fail "starrocks submit failed: $sr_submit_resp"
  379. fi
  380. log "starrocks workflow submitted: workflowInstanceId=$sr_wfi, tag=$sr_tag"
  381. local sr_engine_task_id
  382. sr_engine_task_id="$(wait_task_running_with_engine_id "$sr_wfi" "$sr_task_running" 90)"
  383. log "starrocks task is RUNNING: engineTaskId=$sr_engine_task_id"
  384. wait_starrocks_query_present "$sr_tag" 60
  385. local sr_pod_phase_before
  386. sr_pod_phase_before="$(kubectl -n "$STARROCKS_POD_NAMESPACE" get pod "$STARROCKS_POD_NAME" -o jsonpath='{.status.phase}')"
  387. if [[ "$sr_pod_phase_before" != "Running" ]]; then
  388. fail "starrocks pod is not running before terminate: $STARROCKS_POD_NAMESPACE/$STARROCKS_POD_NAME phase=$sr_pod_phase_before"
  389. fi
  390. log "starrocks processlist observed target query and pod is running"
  391. terminate_workflow "$sr_wfi"
  392. wait_terminated_and_states "$sr_wfi" "$sr_task_running" "$sr_task_downstream" 60
  393. wait_starrocks_query_absent "$sr_tag" 60
  394. local sr_pod_phase_after
  395. sr_pod_phase_after="$(kubectl -n "$STARROCKS_POD_NAMESPACE" get pod "$STARROCKS_POD_NAME" -o jsonpath='{.status.phase}')"
  396. if [[ "$sr_pod_phase_after" != "Running" ]]; then
  397. fail "starrocks pod should remain running after terminate: $STARROCKS_POD_NAMESPACE/$STARROCKS_POD_NAME phase=$sr_pod_phase_after"
  398. fi
  399. log "starrocks terminate verified: running query killed, downstream skipped, pod kept running"
  400. printf '[PASS] workflow terminate e2e verified for Spark + StarRocks\n'
  401. }
  402. main "$@"