| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458 |
- #!/usr/bin/env bash
- set -euo pipefail
- ROOT_DIR="$(cd "$(dirname "$0")/.." && pwd)"
- cd "$ROOT_DIR"
- BACKEND_URL="${E2E_BACKEND_URL:-http://127.0.0.1:8080}"
- ORCH_URL="${E2E_ORCH_URL:-http://127.0.0.1:18080}"
- POLL_INTERVAL_SECONDS="${E2E_POLL_INTERVAL_SECONDS:-1}"
- log() {
- printf '[INFO] %s\n' "$*"
- }
- fail() {
- printf '[FAIL] %s\n' "$*" >&2
- exit 1
- }
- require_cmd() {
- if ! command -v "$1" >/dev/null 2>&1; then
- fail "command not found: $1"
- fi
- }
- require_cmd docker
- require_cmd kubectl
- require_cmd curl
- require_cmd jq
- require_cmd awk
- BACKEND_ENV_DUMP="$(docker inspect ws-backend --format '{{range .Config.Env}}{{println .}}{{end}}' 2>/dev/null || true)"
- if [[ -z "$BACKEND_ENV_DUMP" ]]; then
- fail "cannot read backend container env; ensure container ws-backend is running"
- fi
- backend_env_get() {
- local key="$1"
- printf '%s\n' "$BACKEND_ENV_DUMP" | awk -F= -v k="$key" '$1==k {sub(/^[^=]*=/, "", $0); print; exit}'
- }
- MYSQL_URL="$(backend_env_get MYSQL_URL)"
- MYSQL_DB="$(backend_env_get MYSQL_DATABASE)"
- MYSQL_USER="$(backend_env_get MYSQL_USER)"
- MYSQL_PASSWORD="$(backend_env_get MYSQL_PASSWORD)"
- if [[ -z "$MYSQL_URL" || -z "$MYSQL_DB" || -z "$MYSQL_USER" ]]; then
- fail "backend mysql env is incomplete (MYSQL_URL/MYSQL_DATABASE/MYSQL_USER)"
- fi
- if [[ "$MYSQL_URL" == *:* ]]; then
- DB_HOST="${MYSQL_URL%:*}"
- DB_PORT="${MYSQL_URL##*:}"
- else
- DB_HOST="$MYSQL_URL"
- DB_PORT="3306"
- fi
- mysql_exec() {
- local sql="$1"
- docker compose exec -T mysql sh -lc \
- "mysql -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB' -e \"$sql\""
- }
- mysql_query() {
- local sql="$1"
- docker compose exec -T mysql sh -lc \
- "mysql -N -B -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB' -e \"$sql\""
- }
- api_register_token() {
- local user="e2e_workflow_term_$(date +%s)_$RANDOM"
- local pass='E2eWorkflow@123456'
- local payload
- payload=$(jq -cn \
- --arg username "$user" \
- --arg password "$pass" \
- --arg email "$user@example.com" \
- '{username:$username,password:$password,email:$email,role:"OPS"}')
- local resp
- resp=$(curl -fsS -X POST "$BACKEND_URL/api/auth/register" \
- -H 'Content-Type: application/json' \
- -d "$payload")
- local token
- token=$(echo "$resp" | jq -r '.token // empty')
- if [[ -z "$token" ]]; then
- fail "register token failed: $resp"
- fi
- printf '%s' "$token"
- }
- API_TOKEN="$(api_register_token)"
- AUTH_HEADER="Authorization: Bearer $API_TOKEN"
- api_get() {
- local path="$1"
- curl -fsS "$BACKEND_URL$path" -H "$AUTH_HEADER"
- }
- api_post() {
- local path="$1"
- local body="${2:-{}}"
- curl -fsS -X POST "$BACKEND_URL$path" \
- -H "$AUTH_HEADER" \
- -H 'Content-Type: application/json' \
- -d "$body"
- }
- terminate_workflow() {
- local workflow_instance_id="$1"
- local code
- code=$(curl -sS -o /tmp/e2e_terminate_resp.$$ -w '%{http_code}' \
- -X POST "$BACKEND_URL/api/task-exec/workflows/instances/$workflow_instance_id/terminate" \
- -H "$AUTH_HEADER")
- if [[ "$code" != "204" ]]; then
- local body
- body="$(cat /tmp/e2e_terminate_resp.$$ 2>/dev/null || true)"
- fail "terminate api failed, workflowInstanceId=$workflow_instance_id, status=$code, body=$body"
- fi
- }
- get_running_cluster_id() {
- local component="$1"
- local id
- id="$(mysql_query "select cluster_id from cluster where component_type='$component' and status='RUNNING' order by cluster_id desc limit 1;" | head -n1)"
- if [[ -z "$id" ]]; then
- fail "no RUNNING cluster found for component_type=$component"
- fi
- printf '%s' "$id"
- }
- wait_task_running_with_engine_id() {
- local workflow_instance_id="$1"
- local task_id="$2"
- local timeout_seconds="$3"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- local row
- 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)"
- if [[ -n "$row" ]]; then
- local state
- local engine_id
- state="$(echo "$row" | cut -f1)"
- engine_id="$(echo "$row" | cut -f2)"
- if [[ "$state" == "RUNNING" && -n "$engine_id" ]]; then
- printf '%s' "$engine_id"
- return 0
- fi
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "task did not reach RUNNING with engine id in time, workflowInstanceId=$workflow_instance_id, taskId=$task_id"
- }
- wait_terminated_and_states() {
- local workflow_instance_id="$1"
- local running_task_id="$2"
- local downstream_task_id="$3"
- local timeout_seconds="$4"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- local status_json
- status_json="$(api_get "/api/task-exec/workflows/instances/$workflow_instance_id")"
- local wf_state
- local running_state
- local downstream_state
- wf_state="$(echo "$status_json" | jq -r '.state // empty')"
- running_state="$(echo "$status_json" | jq -r --argjson taskId "$running_task_id" '.taskInstances[] | select(.taskId==$taskId) | .state // empty')"
- downstream_state="$(echo "$status_json" | jq -r --argjson taskId "$downstream_task_id" '.taskInstances[] | select(.taskId==$taskId) | .state // empty')"
- if [[ "$wf_state" == "TERMINATED" && "$running_state" == "KILLED" && "$downstream_state" == "SKIPPED" ]]; then
- return 0
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "workflow/task states mismatch after terminate, workflowInstanceId=$workflow_instance_id"
- }
- wait_spark_job_ref() {
- local operation_id="$1"
- local timeout_seconds="$2"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- local resp
- resp="$(curl -fsS "$ORCH_URL/v1/spark-jobs/operations/$operation_id")"
- local ns
- local app
- ns="$(echo "$resp" | jq -r '.namespace // empty')"
- app="$(echo "$resp" | jq -r '.application // empty')"
- if [[ -n "$ns" && -n "$app" ]]; then
- printf '%s\t%s' "$ns" "$app"
- return 0
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "spark job namespace/application not ready, operationId=$operation_id"
- }
- wait_spark_app_exists() {
- local namespace="$1"
- local app="$2"
- local timeout_seconds="$3"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- if kubectl -n "$namespace" get sparkapplications.sparkoperator.k8s.io "$app" >/dev/null 2>&1; then
- return 0
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "spark application not found before terminate, namespace=$namespace, app=$app"
- }
- wait_spark_app_deleted() {
- local namespace="$1"
- local app="$2"
- local timeout_seconds="$3"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- if kubectl -n "$namespace" get sparkapplications.sparkoperator.k8s.io "$app" >/dev/null 2>&1; then
- sleep "$POLL_INTERVAL_SECONDS"
- continue
- fi
- local pod_count
- pod_count="$(kubectl -n "$namespace" get pods --no-headers 2>/dev/null | awk -v app="$app" '$1 ~ app {c++} END {print c+0}')"
- if [[ "$pod_count" == "0" ]]; then
- return 0
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "spark application resources still exist, namespace=$namespace, app=$app"
- }
- run_starrocks_sql() {
- local sql="$1"
- docker compose exec -T \
- -e E2E_SR_HOST="$STARROCKS_HOST" \
- -e E2E_SR_PORT="$STARROCKS_PORT" \
- -e E2E_SR_USER="$STARROCKS_USER" \
- -e E2E_SR_PASSWORD="$STARROCKS_PASSWORD" \
- -e E2E_SR_SQL="$sql" \
- mysql sh -lc '
- if [ -n "$E2E_SR_PASSWORD" ]; then
- mysql --connect-timeout=10 --protocol=TCP \
- --host="$E2E_SR_HOST" \
- --port="$E2E_SR_PORT" \
- --user="$E2E_SR_USER" \
- --password="$E2E_SR_PASSWORD" \
- --batch --raw \
- -e "$E2E_SR_SQL"
- else
- mysql --connect-timeout=10 --protocol=TCP \
- --host="$E2E_SR_HOST" \
- --port="$E2E_SR_PORT" \
- --user="$E2E_SR_USER" \
- --batch --raw \
- -e "$E2E_SR_SQL"
- fi
- '
- }
- wait_starrocks_query_present() {
- local tag="$1"
- local timeout_seconds="$2"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- local processlist
- processlist="$(run_starrocks_sql "SHOW PROCESSLIST;" || true)"
- if echo "$processlist" | awk -F'\t' -v tag="$tag" 'NR>1 && index($10, tag)>0 {found=1} END {exit found?0:1}'; then
- return 0
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "starrocks target query not found in processlist, tag=$tag"
- }
- wait_starrocks_query_absent() {
- local tag="$1"
- local timeout_seconds="$2"
- local deadline=$((SECONDS + timeout_seconds))
- while (( SECONDS < deadline )); do
- local processlist
- processlist="$(run_starrocks_sql "SHOW PROCESSLIST;" || true)"
- if echo "$processlist" | awk -F'\t' -v tag="$tag" 'NR>1 && index($10, tag)>0 {found=1} END {exit found?1:0}'; then
- return 0
- fi
- sleep "$POLL_INTERVAL_SECONDS"
- done
- fail "starrocks query still present after terminate, tag=$tag"
- }
- main() {
- log "checking backend and orchestrator reachability"
- curl -fsS "$BACKEND_URL/api/auth/login" -X OPTIONS >/dev/null 2>&1 || true
- curl -fsS "$ORCH_URL/v1/operations/non-exist" >/dev/null 2>&1 || true
- SPARK_CLUSTER_ID="$(get_running_cluster_id SPARK)"
- STARROCKS_CLUSTER_ID="$(get_running_cluster_id STARROCKS)"
- log "using clusters: spark=$SPARK_CLUSTER_ID, starrocks=$STARROCKS_CLUSTER_ID"
- local sr_cluster_name="starrocks-$STARROCKS_CLUSTER_ID"
- local sr_endpoint_json
- sr_endpoint_json="$(curl -fsS "$ORCH_URL/v1/starrocks-clusters/$sr_cluster_name/direct-endpoint")"
- STARROCKS_HOST="$(echo "$sr_endpoint_json" | jq -r '.host // empty')"
- STARROCKS_PORT="$(echo "$sr_endpoint_json" | jq -r '.port // empty')"
- STARROCKS_POD_NAMESPACE="$(echo "$sr_endpoint_json" | jq -r '.namespace // empty')"
- STARROCKS_POD_NAME="$(echo "$sr_endpoint_json" | jq -r '.podName // empty')"
- STARROCKS_USER="$(backend_env_get STARROCKS_SQL_RUNNER_USERNAME)"
- STARROCKS_PASSWORD="$(backend_env_get STARROCKS_SQL_RUNNER_PASSWORD)"
- if [[ -z "$STARROCKS_HOST" || -z "$STARROCKS_PORT" || -z "$STARROCKS_USER" ]]; then
- fail "starrocks direct endpoint or auth config is incomplete"
- fi
- local run_epoch
- run_epoch="$(date +%s)"
- local base_id=$((1900000000 + (run_epoch % 1000000) * 1000 + (RANDOM % 1000)))
- local spark_wf_id=$((base_id + 1))
- local spark_task_running=$((base_id + 11))
- local spark_task_downstream=$((base_id + 12))
- local sr_wf_id=$((base_id + 101))
- local sr_task_running=$((base_id + 111))
- local sr_task_downstream=$((base_id + 112))
- local sr_tag="e2e_sr_${run_epoch}_${RANDOM}"
- log "preparing spark terminate e2e workflow"
- cat <<SQL | docker compose exec -T mysql sh -lc "mysql -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB'"
- INSERT INTO workflow_definition(workflow_id,workflow_name,description,dag_json,timeout_seconds,create_time,failure_strategy)
- VALUES (
- $spark_wf_id,
- 'E2E_SPARK_TERMINATE',
- 'spark terminate e2e',
- '{"nodes":[{"taskId":$spark_task_running},{"taskId":$spark_task_downstream}],"edges":[{"fromTaskId":$spark_task_running,"toTaskId":$spark_task_downstream}]}',
- 1800,
- NOW(),
- 'TERMINATE'
- )
- ON DUPLICATE KEY UPDATE
- workflow_name = VALUES(workflow_name),
- description = VALUES(description),
- dag_json = VALUES(dag_json),
- timeout_seconds = VALUES(timeout_seconds),
- failure_strategy = VALUES(failure_strategy);
- INSERT INTO task_definition(task_id,task_name,workflow_id,task_type,task_content,exector_id,timeout_seconds,retry_times)
- VALUES
- ($spark_task_running,'e2e_spark_running',$spark_wf_id,'SPARK_SQL','SELECT COUNT(*) AS c FROM range(1, 1000000000000);',$SPARK_CLUSTER_ID,1800,0),
- ($spark_task_downstream,'e2e_spark_downstream',$spark_wf_id,'SPARK_SQL','SELECT 1;',$SPARK_CLUSTER_ID,1800,0)
- ON DUPLICATE KEY UPDATE
- task_name = VALUES(task_name),
- task_type = VALUES(task_type),
- task_content = VALUES(task_content),
- exector_id = VALUES(exector_id),
- timeout_seconds = VALUES(timeout_seconds),
- retry_times = VALUES(retry_times);
- SQL
- local spark_submit_resp
- spark_submit_resp="$(api_post "/api/task-exec/workflows/$spark_wf_id/submit" '{}')"
- local spark_wfi
- spark_wfi="$(echo "$spark_submit_resp" | jq -r '.workflowInstanceId // empty')"
- if [[ -z "$spark_wfi" ]]; then
- fail "spark submit failed: $spark_submit_resp"
- fi
- log "spark workflow submitted: workflowInstanceId=$spark_wfi"
- local spark_operation_id
- spark_operation_id="$(wait_task_running_with_engine_id "$spark_wfi" "$spark_task_running" 180)"
- log "spark task is RUNNING: operationId=$spark_operation_id"
- local spark_ref
- spark_ref="$(wait_spark_job_ref "$spark_operation_id" 120)"
- local spark_namespace
- local spark_app
- IFS=$'\t' read -r spark_namespace spark_app <<<"$spark_ref"
- wait_spark_app_exists "$spark_namespace" "$spark_app" 120
- log "spark application observed in k8s: namespace=$spark_namespace, app=$spark_app"
- terminate_workflow "$spark_wfi"
- wait_terminated_and_states "$spark_wfi" "$spark_task_running" "$spark_task_downstream" 60
- wait_spark_app_deleted "$spark_namespace" "$spark_app" 120
- log "spark terminate verified: running task killed, downstream skipped, spark resources released"
- log "preparing starrocks terminate e2e workflow"
- cat <<SQL | docker compose exec -T mysql sh -lc "mysql -h'$DB_HOST' -P'$DB_PORT' -u'$MYSQL_USER' -p'$MYSQL_PASSWORD' '$MYSQL_DB'"
- INSERT INTO workflow_definition(workflow_id,workflow_name,description,dag_json,timeout_seconds,create_time,failure_strategy)
- VALUES (
- $sr_wf_id,
- 'E2E_STARROCKS_TERMINATE',
- 'starrocks terminate e2e',
- '{"nodes":[{"taskId":$sr_task_running},{"taskId":$sr_task_downstream}],"edges":[{"fromTaskId":$sr_task_running,"toTaskId":$sr_task_downstream}]}',
- 1800,
- NOW(),
- 'TERMINATE'
- )
- ON DUPLICATE KEY UPDATE
- workflow_name = VALUES(workflow_name),
- description = VALUES(description),
- dag_json = VALUES(dag_json),
- timeout_seconds = VALUES(timeout_seconds),
- failure_strategy = VALUES(failure_strategy);
- INSERT INTO task_definition(task_id,task_name,workflow_id,task_type,task_content,exector_id,timeout_seconds,retry_times)
- VALUES
- ($sr_task_running,'e2e_sr_running',$sr_wf_id,'STARROCKS_SQL','SELECT /*$sr_tag*/ SLEEP(300);',$STARROCKS_CLUSTER_ID,1800,0),
- ($sr_task_downstream,'e2e_sr_downstream',$sr_wf_id,'STARROCKS_SQL','SELECT 1;',$STARROCKS_CLUSTER_ID,1800,0)
- ON DUPLICATE KEY UPDATE
- task_name = VALUES(task_name),
- task_type = VALUES(task_type),
- task_content = VALUES(task_content),
- exector_id = VALUES(exector_id),
- timeout_seconds = VALUES(timeout_seconds),
- retry_times = VALUES(retry_times);
- SQL
- local sr_submit_resp
- sr_submit_resp="$(api_post "/api/task-exec/workflows/$sr_wf_id/submit" '{}')"
- local sr_wfi
- sr_wfi="$(echo "$sr_submit_resp" | jq -r '.workflowInstanceId // empty')"
- if [[ -z "$sr_wfi" ]]; then
- fail "starrocks submit failed: $sr_submit_resp"
- fi
- log "starrocks workflow submitted: workflowInstanceId=$sr_wfi, tag=$sr_tag"
- local sr_engine_task_id
- sr_engine_task_id="$(wait_task_running_with_engine_id "$sr_wfi" "$sr_task_running" 90)"
- log "starrocks task is RUNNING: engineTaskId=$sr_engine_task_id"
- wait_starrocks_query_present "$sr_tag" 60
- local sr_pod_phase_before
- sr_pod_phase_before="$(kubectl -n "$STARROCKS_POD_NAMESPACE" get pod "$STARROCKS_POD_NAME" -o jsonpath='{.status.phase}')"
- if [[ "$sr_pod_phase_before" != "Running" ]]; then
- fail "starrocks pod is not running before terminate: $STARROCKS_POD_NAMESPACE/$STARROCKS_POD_NAME phase=$sr_pod_phase_before"
- fi
- log "starrocks processlist observed target query and pod is running"
- terminate_workflow "$sr_wfi"
- wait_terminated_and_states "$sr_wfi" "$sr_task_running" "$sr_task_downstream" 60
- wait_starrocks_query_absent "$sr_tag" 60
- local sr_pod_phase_after
- sr_pod_phase_after="$(kubectl -n "$STARROCKS_POD_NAMESPACE" get pod "$STARROCKS_POD_NAME" -o jsonpath='{.status.phase}')"
- if [[ "$sr_pod_phase_after" != "Running" ]]; then
- fail "starrocks pod should remain running after terminate: $STARROCKS_POD_NAMESPACE/$STARROCKS_POD_NAME phase=$sr_pod_phase_after"
- fi
- log "starrocks terminate verified: running query killed, downstream skipped, pod kept running"
- printf '[PASS] workflow terminate e2e verified for Spark + StarRocks\n'
- }
- main "$@"
|