chore(coordinator): add API smoke script and request collection
Two ways to exercise every endpoint, both living in the repo rather than in a personal Postman workspace: - scripts/smoke.sh walks the full lifecycle and asserts each status, exiting non-zero on the first surprise, so it works in CI as well as by hand; - api/requests.http drives the same calls from an editor's REST client, with later requests reusing ids captured from earlier responses. It doubles as API documentation for the worker author. The script claims until it sees its own job's chunks instead of assuming an empty queue: a shared development database usually holds pending tasks from earlier runs, and it takes the attempt number from the claim response, since a task requeued after an expired lease comes back with attempt 2 or 3. Note for whoever extends the validation cases: Go matches JSON field names case-insensitively, so "worker_ID" is accepted as "worker_id". Only a genuinely unknown key trips DisallowUnknownFields.
This commit is contained in:
@@ -63,3 +63,10 @@ rebuild:
|
||||
|
||||
psql:
|
||||
docker compose exec postgres psql -U scimesh -d scimesh
|
||||
|
||||
# --- api ------------------------------------------------------------------
|
||||
# Exercises every endpoint against a running coordinator; exits non-zero on the
|
||||
# first unexpected status. See also api/requests.http for clicking through them
|
||||
# one at a time in an editor.
|
||||
smoke:
|
||||
./scripts/smoke.sh
|
||||
|
||||
@@ -111,6 +111,18 @@ See `.env.example`; only `DATABASE_URL` is required.
|
||||
| GET | `/jobs/{job_id}` | Aggregate job progress |
|
||||
| GET | `/health` | Liveness (unauthenticated) |
|
||||
|
||||
## Poking the API
|
||||
|
||||
Two ways, both checked in:
|
||||
|
||||
```sh
|
||||
make smoke # every endpoint, asserted; non-zero exit on failure
|
||||
```
|
||||
|
||||
`api/requests.http` runs the same calls one at a time from an editor with a REST
|
||||
client (VSCodium/VS Code "REST Client", JetBrains HTTP Client). Later requests
|
||||
reuse ids captured from earlier responses, so it doubles as API documentation.
|
||||
|
||||
## Status
|
||||
|
||||
The queue works end to end: a job can be submitted, split into tasks, leased to
|
||||
|
||||
@@ -0,0 +1,157 @@
|
||||
# SciMesh Coordinator — API requests
|
||||
#
|
||||
# Runnable from any editor with a REST client (VSCodium/VS Code "REST Client",
|
||||
# JetBrains HTTP Client). Click "Send Request" above each block, top to bottom:
|
||||
# later requests reuse ids captured from earlier responses.
|
||||
#
|
||||
# Start the stack first: docker compose up -d
|
||||
|
||||
@host = http://localhost:8080
|
||||
@token = change-me
|
||||
@worker = worker-1
|
||||
|
||||
### Health — the only unauthenticated endpoint
|
||||
GET {{host}}/health
|
||||
|
||||
### Auth check — no token must be rejected with 401
|
||||
POST {{host}}/tasks/claim
|
||||
Content-Type: application/json
|
||||
|
||||
{ "worker_id": "{{worker}}" }
|
||||
|
||||
### 1. Create a job and its chunks (201)
|
||||
# The coordinator splits the submission into one task per chunk, transactionally.
|
||||
# @name createJob
|
||||
POST {{host}}/jobs
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"workload": "similarity_search",
|
||||
"input_uri": "s3://chembl/full.sdf",
|
||||
"parameters": { "top_k": 10 },
|
||||
"chunks": [
|
||||
{ "chunk_index": 0, "input_uri": "s3://chembl/shard-0.sdf", "input_sha256": "aaa", "max_attempts": 3 },
|
||||
{ "chunk_index": 1, "input_uri": "s3://chembl/shard-1.sdf", "input_sha256": "bbb", "max_attempts": 3 },
|
||||
{ "chunk_index": 2, "input_uri": "s3://chembl/shard-2.sdf", "input_sha256": "ccc", "max_attempts": 3 }
|
||||
]
|
||||
}
|
||||
|
||||
@jobId = {{createJob.response.body.id}}
|
||||
|
||||
### 2. Claim a task (200, or 204 when the queue is empty)
|
||||
# Each call leases a different task; run it repeatedly to see chunk_index advance.
|
||||
# @name claim
|
||||
POST {{host}}/tasks/claim
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "{{worker}}",
|
||||
"workloads": ["similarity_search"]
|
||||
}
|
||||
|
||||
@taskId = {{claim.response.body.task_id}}
|
||||
@attempt = {{claim.response.body.attempt}}
|
||||
|
||||
### 3. Heartbeat — renew the lease while the task is still running (200)
|
||||
POST {{host}}/tasks/{{taskId}}/heartbeat
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "{{worker}}",
|
||||
"attempt": {{attempt}}
|
||||
}
|
||||
|
||||
### 4. Submit the result (200)
|
||||
POST {{host}}/tasks/{{taskId}}/result
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "{{worker}}",
|
||||
"attempt": {{attempt}},
|
||||
"result_uri": "s3://results/shard-0.csv",
|
||||
"result_sha256": "r0sha",
|
||||
"metrics": { "elapsed_ms": 1234, "candidates": 50000 }
|
||||
}
|
||||
|
||||
### 4a. Replay the same result — must be idempotent (200, not 409)
|
||||
POST {{host}}/tasks/{{taskId}}/result
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "{{worker}}",
|
||||
"attempt": {{attempt}},
|
||||
"result_uri": "s3://results/shard-0.csv",
|
||||
"result_sha256": "r0sha",
|
||||
"metrics": { "elapsed_ms": 1234, "candidates": 50000 }
|
||||
}
|
||||
|
||||
### 4b. A different result for the same task — conflict (409)
|
||||
POST {{host}}/tasks/{{taskId}}/result
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "{{worker}}",
|
||||
"attempt": {{attempt}},
|
||||
"result_uri": "s3://results/SOMETHING-ELSE.csv",
|
||||
"result_sha256": "different"
|
||||
}
|
||||
|
||||
### 4c. Another worker submitting for this task — conflict (409)
|
||||
POST {{host}}/tasks/{{taskId}}/result
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "impostor",
|
||||
"attempt": {{attempt}},
|
||||
"result_uri": "s3://results/x.csv",
|
||||
"result_sha256": "x"
|
||||
}
|
||||
|
||||
### 5. Report a failure instead (200)
|
||||
# retryable=true returns the task to the queue while attempts remain;
|
||||
# retryable=false fails it terminally.
|
||||
POST {{host}}/tasks/{{taskId}}/failure
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{
|
||||
"worker_id": "{{worker}}",
|
||||
"attempt": {{attempt}},
|
||||
"error_code": "download_failed",
|
||||
"error_message": "checksum mismatch on shard",
|
||||
"retryable": true
|
||||
}
|
||||
|
||||
### 6. Job progress (200)
|
||||
GET {{host}}/jobs/{{jobId}}
|
||||
Authorization: Bearer {{token}}
|
||||
|
||||
### --- error cases -------------------------------------------------------
|
||||
|
||||
### Malformed UUID in the path (400)
|
||||
POST {{host}}/tasks/not-a-uuid/result
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{ "worker_id": "{{worker}}", "attempt": 1, "result_uri": "s3://x", "result_sha256": "x" }
|
||||
|
||||
### Unknown field in the body (400) — a misspelled key must not pass silently
|
||||
POST {{host}}/tasks/claim
|
||||
Authorization: Bearer {{token}}
|
||||
Content-Type: application/json
|
||||
|
||||
{ "worker_ID": "{{worker}}" }
|
||||
|
||||
### Unknown job (404)
|
||||
GET {{host}}/jobs/00000000-0000-0000-0000-000000000000
|
||||
Authorization: Bearer {{token}}
|
||||
|
||||
### Stitching is not implemented yet (501)
|
||||
# Any endpoint whose use case is still a stub answers 501.
|
||||
Executable
+122
@@ -0,0 +1,122 @@
|
||||
#!/usr/bin/env bash
|
||||
#
|
||||
# End-to-end smoke test against a running coordinator.
|
||||
#
|
||||
# ./scripts/smoke.sh # localhost:8080, token from .env
|
||||
# HOST=http://1.2.3.4:8080 TOKEN=x ./scripts/smoke.sh
|
||||
#
|
||||
# Exits non-zero on the first unexpected status, so it is usable in CI.
|
||||
|
||||
set -uo pipefail
|
||||
|
||||
HOST="${HOST:-http://localhost:8080}"
|
||||
TOKEN="${TOKEN:-$(grep -s '^WORKER_AUTH_TOKEN=' .env | cut -d= -f2- || echo change-me)}"
|
||||
|
||||
pass=0
|
||||
fail=0
|
||||
|
||||
# check <label> <expected-status> <curl args...>
|
||||
check() {
|
||||
local label="$1" want="$2"
|
||||
shift 2
|
||||
local body status
|
||||
body=$(curl -sS -w '\n%{http_code}' "$@" 2>&1)
|
||||
status=$(printf '%s' "$body" | tail -n1)
|
||||
|
||||
if [[ "$status" == "$want" ]]; then
|
||||
printf ' \033[32m✓\033[0m %-46s %s\n' "$label" "$status"
|
||||
pass=$((pass + 1))
|
||||
else
|
||||
printf ' \033[31m✗\033[0m %-46s got %s, want %s\n' "$label" "$status" "$want"
|
||||
printf ' %s\n' "$(printf '%s' "$body" | head -n-1)"
|
||||
fail=$((fail + 1))
|
||||
fi
|
||||
}
|
||||
|
||||
json() { printf '%s' "$1" | head -n-1; }
|
||||
|
||||
auth=(-H "Authorization: Bearer ${TOKEN}" -H 'Content-Type: application/json')
|
||||
|
||||
echo "coordinator: ${HOST}"
|
||||
echo
|
||||
|
||||
echo "health & auth"
|
||||
check "GET /health" 200 "${HOST}/health"
|
||||
check "claim without a token → 401" 401 -X POST "${HOST}/tasks/claim" \
|
||||
-H 'Content-Type: application/json' -d '{"worker_id":"w1"}'
|
||||
|
||||
echo
|
||||
echo "job lifecycle"
|
||||
job=$(curl -sS "${auth[@]}" -X POST "${HOST}/jobs" -d '{
|
||||
"workload":"similarity_search","input_uri":"s3://chembl","parameters":{"top_k":10},
|
||||
"chunks":[{"chunk_index":0,"input_uri":"s3://c0","input_sha256":"aaa"},
|
||||
{"chunk_index":1,"input_uri":"s3://c1","input_sha256":"bbb"}]}')
|
||||
job_id=$(printf '%s' "$job" | python3 -c 'import json,sys;print(json.load(sys.stdin)["id"])' 2>/dev/null)
|
||||
|
||||
if [[ -z "${job_id:-}" ]]; then
|
||||
echo " ✗ could not create a job: $job"
|
||||
exit 1
|
||||
fi
|
||||
printf ' \033[32m✓\033[0m %-46s %s\n' "POST /jobs" "$job_id"
|
||||
pass=$((pass + 1))
|
||||
|
||||
# The database may hold pending tasks from earlier runs, so claim until we have
|
||||
# both of *our* chunks rather than assuming the queue starts empty. The attempt
|
||||
# number comes from the response too: a task requeued by an expired lease is
|
||||
# handed out with attempt 2 or 3, and hard-coding 1 would fail the lease check.
|
||||
declare -A our_chunks
|
||||
task_id=""
|
||||
attempt=""
|
||||
for _ in $(seq 1 40); do
|
||||
claim=$(curl -sS "${auth[@]}" -X POST "${HOST}/tasks/claim" -d '{"worker_id":"w1"}')
|
||||
[[ -z "$claim" ]] && break # 204: queue drained
|
||||
|
||||
read -r c_job c_task c_chunk c_attempt < <(printf '%s' "$claim" |
|
||||
python3 -c 'import json,sys;d=json.load(sys.stdin);print(d["job_id"],d["task_id"],d["chunk_index"],d["attempt"])' 2>/dev/null)
|
||||
[[ "$c_job" != "$job_id" ]] && continue # someone else's leftover task
|
||||
|
||||
our_chunks["$c_chunk"]=1
|
||||
if [[ -z "$task_id" ]]; then
|
||||
task_id="$c_task"
|
||||
attempt="$c_attempt"
|
||||
fi
|
||||
[[ "${#our_chunks[@]}" -eq 2 ]] && break
|
||||
done
|
||||
|
||||
if [[ "${#our_chunks[@]}" -eq 2 ]]; then
|
||||
printf ' \033[32m✓\033[0m %-46s chunks %s\n' "POST /tasks/claim × 2 (distinct)" "${!our_chunks[*]}"
|
||||
pass=$((pass + 1))
|
||||
else
|
||||
printf ' \033[31m✗\033[0m %-46s got %d distinct chunks, want 2\n' "claim" "${#our_chunks[@]}"
|
||||
fail=$((fail + 1))
|
||||
exit 1
|
||||
fi
|
||||
|
||||
check "heartbeat" 200 -X POST "${HOST}/tasks/${task_id}/heartbeat" "${auth[@]}" \
|
||||
-d "{\"worker_id\":\"w1\",\"attempt\":${attempt}}"
|
||||
check "foreign worker submits → 409" 409 -X POST "${HOST}/tasks/${task_id}/result" "${auth[@]}" \
|
||||
-d "{\"worker_id\":\"impostor\",\"attempt\":${attempt},\"result_uri\":\"s3://x\",\"result_sha256\":\"x\"}"
|
||||
check "submit result" 200 -X POST "${HOST}/tasks/${task_id}/result" "${auth[@]}" \
|
||||
-d "{\"worker_id\":\"w1\",\"attempt\":${attempt},\"result_uri\":\"s3://r0\",\"result_sha256\":\"rrr\"}"
|
||||
check "replay same result → idempotent" 200 -X POST "${HOST}/tasks/${task_id}/result" "${auth[@]}" \
|
||||
-d "{\"worker_id\":\"w1\",\"attempt\":${attempt},\"result_uri\":\"s3://r0\",\"result_sha256\":\"rrr\"}"
|
||||
check "different result → 409" 409 -X POST "${HOST}/tasks/${task_id}/result" "${auth[@]}" \
|
||||
-d "{\"worker_id\":\"w1\",\"attempt\":${attempt},\"result_uri\":\"s3://other\",\"result_sha256\":\"zzz\"}"
|
||||
check "GET /jobs/{id}" 200 "${HOST}/jobs/${job_id}" "${auth[@]}"
|
||||
|
||||
echo
|
||||
echo "input validation"
|
||||
check "malformed uuid → 400" 400 -X POST "${HOST}/tasks/not-a-uuid/result" "${auth[@]}" \
|
||||
-d '{"worker_id":"w1","attempt":1,"result_uri":"s3://x","result_sha256":"x"}'
|
||||
# Note: Go's encoding/json matches field names case-insensitively, so
|
||||
# "worker_ID" would be accepted as "worker_id". Only a genuinely unknown key
|
||||
# trips DisallowUnknownFields.
|
||||
check "unknown json field → 400" 400 -X POST "${HOST}/tasks/claim" "${auth[@]}" \
|
||||
-d '{"worker_id":"w1","totally_unknown":1}'
|
||||
check "unknown job → 404" 404 "${HOST}/jobs/00000000-0000-0000-0000-000000000000" "${auth[@]}"
|
||||
|
||||
echo
|
||||
curl -sS "${HOST}/jobs/${job_id}" "${auth[@]}"
|
||||
echo
|
||||
printf '\n%d passed, %d failed\n' "$pass" "$fail"
|
||||
[[ "$fail" -eq 0 ]]
|
||||
Reference in New Issue
Block a user