#!/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 <