diff --git a/.gitignore b/.gitignore index f5195eb..4125a86 100644 --- a/.gitignore +++ b/.gitignore @@ -5,6 +5,7 @@ Cargo.lock # Go /api/bin/ +/api/api-bin *.exe *.test *.out diff --git a/README.md b/README.md index ffde36d..b91242b 100644 --- a/README.md +++ b/README.md @@ -269,6 +269,36 @@ Acceptance bar for engine correctness (M1): `σ²_30d`** — roughly one part per million in published BVOL. Enforced by `cargo test -p volx-engine` against snapshot fixtures. +### End-to-end smoke + +`scripts/e2e-smoke.sh` boots every M1 service in dependency order, waits +two engine snapshots, and asserts that fresh data has reached the API +surface. The M1 close gate (issue #66). + +```bash +./scripts/e2e-smoke.sh +``` + +Requirements on `PATH`: `docker`, `cargo`, `go`, `curl`, `python3` with +the `websockets` package (`pip install websockets`). + +Asserts, in order: + +1. `options_ticks` has ≥ 1 fresh row in the last 1 minute (ingestion + + normalizer reaching ClickHouse). +2. `index_ticks` has ≥ 1 fresh row in the last 2 minutes (engine writing). +3. `GET /v1/index/bvol/latest` returns 200 with `value > 0` and + `age < 150 s` (two engine cycles + slack). +4. `GET /v1/index/bvol/history?interval=5m&limit=12` returns 200 with + `bars ≥ 1`. +5. `ws://…/v1/stream` delivers at least one `type:tick` frame for both + `bvol` and `evol` inside a 75 s window (via + `scripts/e2e-ws-client.py`). + +Exits 0 on success with a stage-timing summary; non-zero with the name +of the failed assertion. Compose teardown runs in the exit trap so the +script is idempotent across repeated runs. + --- ## Repo layout diff --git a/scripts/e2e-smoke.sh b/scripts/e2e-smoke.sh new file mode 100755 index 0000000..81a77ad --- /dev/null +++ b/scripts/e2e-smoke.sh @@ -0,0 +1,282 @@ +#!/usr/bin/env bash +# End-to-end smoke for the VolX local pipeline (issue #66). +# +# Brings every M1 service online in order, waits long enough for at least +# two engine snapshots, and asserts that fresh data has propagated all the +# way to the API surface + WebSocket client. Exits non-zero with the name +# of the failed check so a regression bisects to a single hop. +# +# Pipeline: +# +# Deribit WS → ingestion → normalizer → ClickHouse + Redis +# ↓ +# engine (60 s) +# ↓ +# API REST + WS +# ↓ +# Python WS client +# +# Idempotent: tears down on exit, even on failure. +# Requirements on PATH: docker, cargo, go, curl, python3 (with `websockets`). +# +# Usage: ./scripts/e2e-smoke.sh + +set -euo pipefail + +# ------- config -------------------------------------------------------------- + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +COMPOSE_FILE="${ROOT_DIR}/docker/docker-compose.yml" +COMPOSE_PROJECT="volx-local" + +CLICKHOUSE_DB="${CLICKHOUSE_DB:-volx}" +CLICKHOUSE_HOST="${CLICKHOUSE_HOST:-127.0.0.1}" +CLICKHOUSE_HTTP_PORT="${CLICKHOUSE_HTTP_PORT:-8123}" + +API_HOST="${API_HOST:-localhost}" +API_PORT="${API_PORT:-8080}" +API_BASE="http://${API_HOST}:${API_PORT}" + +# Two engine snapshots + cold-start safety margin. Deribit instrument +# enumeration on first connect can take 30-60 s, which pushes snapshot +# #1 to t=60-90 and snapshot #2 to t=120-150. 240 s leaves room for a +# slow first connect on top of the two-cycle wait. Override with the +# env var to tune on faster hardware / warm caches. +ENGINE_WAIT_S="${ENGINE_WAIT_S:-240}" +WS_TIMEOUT_S="${WS_TIMEOUT_S:-75}" +# Preserve clickhouse + redis volumes across runs. Off by default — the +# smoke wipes data so stale `index_ticks` rows from prior sessions can +# not silently satisfy "fresh row" assertions. Set to "1" to skip the +# wipe when iterating on the script without losing local history. +PRESERVE_VOLUMES="${PRESERVE_VOLUMES:-0}" + +# Prefer the research venv (already has `websockets` from M0 work); +# fall back to system python3 otherwise. Override via env if needed. +PYTHON_BIN="${PYTHON_BIN:-}" +if [ -z "${PYTHON_BIN}" ]; then + if [ -x "${ROOT_DIR}/research/.venv/bin/python3" ]; then + PYTHON_BIN="${ROOT_DIR}/research/.venv/bin/python3" + else + PYTHON_BIN="python3" + fi +fi + +LOG_DIR="$(mktemp -d -t volx-e2e-XXXXXX)" +INGESTION_LOG="${LOG_DIR}/ingestion.log" +ENGINE_LOG="${LOG_DIR}/engine.log" +API_LOG="${LOG_DIR}/api.log" + +declare -a CHILD_PIDS=() + +# ------- helpers ------------------------------------------------------------- + +t_start_total=$(date +%s) +# Parallel indexed arrays — associative arrays would force bash 4+ and +# macOS ships bash 3.2 by default. +STAGE_NAMES=() +STAGE_START_T=() +STAGE_END_T=() + +stage_begin() { + STAGE_NAMES+=("$1") + STAGE_START_T+=("$(date +%s)") + STAGE_END_T+=("0") + echo "==> $1" >&2 +} +stage_end() { + local i=$(( ${#STAGE_NAMES[@]} - 1 )) + STAGE_END_T[${i}]=$(date +%s) +} + +elapsed_total() { echo $(( $(date +%s) - t_start_total )); } + +teardown() { + local rc=$? + set +e + echo "==> teardown" >&2 + # `cargo run` parents fork-exec the compiled binary; killing the cargo + # PID does not always reap the child. Hit the binaries by name first, + # then fall back to the recorded parent PIDs. + pkill -f "target/release/volx-ingestion" 2>/dev/null + pkill -f "target/release/volx-normalizer" 2>/dev/null + pkill -f "target/release/volx-engine" 2>/dev/null + pkill -f "api/api-bin" 2>/dev/null + pkill -f "exe/api" 2>/dev/null # historical: go run uses tmp + for pid in "${CHILD_PIDS[@]:-}"; do + [ -n "${pid:-}" ] && kill "${pid}" 2>/dev/null + done + # Catch any stragglers still bound to the well-known ports. + lsof -ti ":${API_PORT}" 2>/dev/null | xargs -r kill -9 2>/dev/null + docker compose -p "${COMPOSE_PROJECT}" -f "${COMPOSE_FILE}" down >/dev/null 2>&1 + echo "logs preserved in: ${LOG_DIR}" >&2 + echo "total runtime: $(elapsed_total)s" >&2 + exit "${rc}" +} +trap teardown EXIT INT TERM + +fail() { + echo "FAIL: $*" >&2 + exit 1 +} + +# Polls until `cond` succeeds or `timeout` seconds elapse. `cond` is the +# command line passed in $@, evaluated as a shell expression. +wait_until() { + local label="$1"; shift + local timeout="$1"; shift + local deadline=$(( $(date +%s) + timeout )) + while [ "$(date +%s)" -lt "${deadline}" ]; do + if eval "$@" >/dev/null 2>&1; then + return 0 + fi + sleep 2 + done + fail "${label} did not become ready within ${timeout}s" +} + +# ClickHouse HTTP query helper. The `-G` flag is critical — without it +# curl POSTs the form-encoded query as the request body and ClickHouse +# tries to parse the literal "query=SELECT…" string as SQL. +ch_query() { + curl -sS -G --fail \ + --data-urlencode "query=$1" \ + "http://${CLICKHOUSE_HOST}:${CLICKHOUSE_HTTP_PORT}/?database=${CLICKHOUSE_DB}" +} + +start_service() { + local label="$1" logfile="$2"; shift 2 + stage_begin "${label}" + ( cd "${ROOT_DIR}" && "$@" ) >"${logfile}" 2>&1 & + CHILD_PIDS+=("$!") + stage_end "${label}" +} + +# Path layout of the built binaries. +ING_BIN="${ROOT_DIR}/target/release/volx-ingestion" +ENG_BIN="${ROOT_DIR}/target/release/volx-engine" +API_BIN="${ROOT_DIR}/api/api-bin" + +# ------- stages -------------------------------------------------------------- + +stage_begin "compose-up" +# DESTRUCTIVE by default: wipe ClickHouse + Redis volumes so stale +# `index_ticks` rows from prior sessions can not satisfy "fresh row" +# assertions before the new pipeline produces anything. Set +# PRESERVE_VOLUMES=1 in the environment to keep data across runs (e.g. +# while iterating on the script without losing a day's history). +if [ "${PRESERVE_VOLUMES}" = "1" ]; then + echo " PRESERVE_VOLUMES=1 — keeping clickhouse + redis data" >&2 + docker compose -p "${COMPOSE_PROJECT}" -f "${COMPOSE_FILE}" down >/dev/null 2>&1 || true +else + docker compose -p "${COMPOSE_PROJECT}" -f "${COMPOSE_FILE}" down --volumes >/dev/null 2>&1 || true +fi +docker compose -p "${COMPOSE_PROJECT}" -f "${COMPOSE_FILE}" up -d +stage_end "compose-up" + +stage_begin "compose-healthy" +wait_until "clickhouse healthy" 60 \ + "curl -sS --fail http://${CLICKHOUSE_HOST}:${CLICKHOUSE_HTTP_PORT}/ping" +wait_until "redis healthy" 30 \ + "docker exec volx-redis redis-cli ping | grep -q PONG" +stage_end "compose-healthy" + +stage_begin "build-rust" +# `volx-normalizer` is a library that lives inside the ingestion process — +# only ingestion + engine need binaries. +( cd "${ROOT_DIR}" && cargo build --release -p volx-ingestion -p volx-engine ) \ + || fail "cargo build --release failed (see ${LOG_DIR} for prior logs)" +stage_end "build-rust" + +stage_begin "build-go" +( cd "${ROOT_DIR}/api" && go build -o "${API_BIN}" ./cmd/api ) \ + || fail "go build ./cmd/api failed" +stage_end "build-go" + +start_service "ingestion" "${INGESTION_LOG}" "${ING_BIN}" +start_service "engine" "${ENGINE_LOG}" "${ENG_BIN}" +start_service "api" "${API_LOG}" "${API_BIN}" + +stage_begin "api-ready" +wait_until "api /v1/health" 60 "curl -sS --fail ${API_BASE}/v1/health" +stage_end "api-ready" + +stage_begin "engine-cycles" +echo " polling for ≥ 2 fresh engine snapshots (timeout ${ENGINE_WAIT_S}s)" >&2 +# Wait for at least TWO snapshots so the second-cycle path (engine +# warm) is exercised before assertions run. One snapshot would race +# against the `/latest` age threshold below — by the time the four +# assertions finish, the single tick is already ~60-100 s old. +wait_until "engine ≥ 2 snapshots" "${ENGINE_WAIT_S}" \ + "[ \$(curl -sS -G --data-urlencode 'query=SELECT count(DISTINCT ts) FROM index_ticks' 'http://${CLICKHOUSE_HOST}:${CLICKHOUSE_HTTP_PORT}/?database=${CLICKHOUSE_DB}' | tr -d '[:space:]') -ge 2 ]" +stage_end "engine-cycles" + +# ------- assertions --------------------------------------------------------- + +stage_begin "assert-options_ticks" +opt_rows=$(ch_query "SELECT count() FROM options_ticks WHERE ts > now() - INTERVAL 1 MINUTE" | tr -d '[:space:]') +if [ -z "${opt_rows}" ] || [ "${opt_rows}" -lt 1 ]; then + fail "options_ticks had ${opt_rows:-0} fresh rows (expected ≥ 1) in the last 60 s" +fi +echo " options_ticks fresh rows: ${opt_rows}" >&2 +stage_end "assert-options_ticks" + +stage_begin "assert-index_ticks" +idx_rows=$(ch_query "SELECT count() FROM index_ticks WHERE ts > now() - INTERVAL 2 MINUTE" | tr -d '[:space:]') +if [ -z "${idx_rows}" ] || [ "${idx_rows}" -lt 1 ]; then + fail "index_ticks had ${idx_rows:-0} fresh rows (expected ≥ 1) in the last 120 s" +fi +echo " index_ticks fresh rows: ${idx_rows}" >&2 +stage_end "assert-index_ticks" + +stage_begin "assert-rest-latest" +latest_body=$(curl -sS --fail "${API_BASE}/v1/index/bvol/latest") +latest_value=$(echo "${latest_body}" | "${PYTHON_BIN}" -c "import sys,json;print(json.load(sys.stdin)['value'])") +latest_age=$(echo "${latest_body}" | "${PYTHON_BIN}" -c " +import sys, json, datetime +j = json.load(sys.stdin) +ts = j['ts'] +dt = datetime.datetime.fromisoformat(ts.replace('Z', '+00:00')) +now = datetime.datetime.now(datetime.timezone.utc) +print(int((now - dt).total_seconds())) +") +echo " /latest bvol value=${latest_value} age=${latest_age}s" >&2 +"${PYTHON_BIN}" -c "import sys; sys.exit(0 if float('${latest_value}') > 0 else 1)" \ + || fail "/latest value (${latest_value}) is not > 0" +# Engine publishes every 60 s; with the post-second-snapshot delay +# burned by the ClickHouse assertions + curl latency, a healthy tick +# can reach ~70-90 s of age before this line. 150 s = two cycles + a +# safety margin so a brief engine stall doesn't false-fail the smoke. +[ "${latest_age}" -lt 150 ] || fail "/latest age (${latest_age}s) ≥ 150s" +stage_end "assert-rest-latest" + +stage_begin "assert-rest-history" +hist_count=$(curl -sS --fail "${API_BASE}/v1/index/bvol/history?interval=5m&limit=12" \ + | "${PYTHON_BIN}" -c "import sys,json;print(len(json.load(sys.stdin)['bars']))") +echo " /history bvol bars=${hist_count}" >&2 +[ "${hist_count}" -ge 1 ] || fail "/history bars=${hist_count} (expected ≥ 1)" +stage_end "assert-rest-history" + +stage_begin "assert-ws-stream" +"${PYTHON_BIN}" "${ROOT_DIR}/scripts/e2e-ws-client.py" \ + --url "ws://${API_HOST}:${API_PORT}/v1/stream" \ + --timeout "${WS_TIMEOUT_S}" \ + || fail "ws stream did not deliver one tick per channel inside ${WS_TIMEOUT_S}s" +stage_end "assert-ws-stream" + +# ------- timing table ------------------------------------------------------- + +echo "" +echo "==================================== summary ====================================" +printf "%-26s %8s\n" "stage" "seconds" +printf "%-26s %8s\n" "--------------------------" "--------" +for i in "${!STAGE_NAMES[@]}"; do + start=${STAGE_START_T[$i]} + end=${STAGE_END_T[$i]} + if [ "${end}" -gt 0 ]; then + printf "%-26s %8d\n" "${STAGE_NAMES[$i]}" "$(( end - start ))" + fi +done +printf "%-26s %8s\n" "--------------------------" "--------" +printf "%-26s %8d\n" "TOTAL" "$(elapsed_total)" +echo "=================================================================================" +echo "OK" diff --git a/scripts/e2e-ws-client.py b/scripts/e2e-ws-client.py new file mode 100755 index 0000000..f9c135b --- /dev/null +++ b/scripts/e2e-ws-client.py @@ -0,0 +1,111 @@ +#!/usr/bin/env python3 +"""End-to-end WebSocket client used by `e2e-smoke.sh` (issue #66). + +Connects to `ws://localhost:8080/v1/stream`, subscribes to both `bvol` and +`evol`, and asserts that at least one tick of each channel arrives inside +a fixed window. + +Exit codes: + 0 — both channels delivered at least one tick. + 1 — at least one channel was silent (or the WS handshake failed). + +The frame contract is fixed by `api/internal/stream/hub.go`: + + {"type": "tick", "channel": "bvol", "value": 37.37, + "ts": 1779782592444, "confidence": 1.0} + +Usage: + python3 scripts/e2e-ws-client.py [--url URL] [--timeout SECONDS] +""" +from __future__ import annotations + +import argparse +import asyncio +import json +import sys +from typing import Any + +import websockets + +REQUIRED_CHANNELS = {"bvol", "evol"} +DEFAULT_URL = "ws://localhost:8080/v1/stream" +DEFAULT_TIMEOUT = 75.0 + + +def _validate(frame: dict[str, Any]) -> bool: + """Strict wire-shape check against PRD §6.""" + if frame.get("type") != "tick": + return False + if frame.get("channel") not in REQUIRED_CHANNELS: + return False + for k in ("value", "confidence"): + if not isinstance(frame.get(k), (int, float)): + return False + if not isinstance(frame.get("ts"), int): + return False + return True + + +async def run(url: str, timeout: float) -> int: + seen: dict[str, dict[str, Any]] = {} + print(f"[ws] connect {url}", flush=True) + try: + async with websockets.connect(url, ping_interval=20, ping_timeout=15) as ws: + await ws.send(json.dumps({"action": "subscribe", "channels": list(REQUIRED_CHANNELS)})) + print("[ws] subscribed; awaiting ticks", flush=True) + + async def consume() -> None: + async for raw in ws: + try: + frame = json.loads(raw) + except json.JSONDecodeError: + print(f"[ws] non-json frame: {raw!r}", flush=True) + continue + if frame.get("type") == "error": + print(f"[ws] server error: {frame}", flush=True) + continue + if not _validate(frame): + print(f"[ws] malformed tick: {frame}", flush=True) + continue + ch = frame["channel"] + if ch not in seen: + seen[ch] = frame + print( + f"[ws] tick {ch} value={frame['value']:.4f} " + f"ts={frame['ts']} confidence={frame['confidence']}", + flush=True, + ) + if REQUIRED_CHANNELS.issubset(seen.keys()): + return + + try: + await asyncio.wait_for(consume(), timeout=timeout) + except asyncio.TimeoutError: + pass + except (OSError, websockets.exceptions.WebSocketException) as e: + print(f"[ws] handshake/transport failure: {e}", file=sys.stderr, flush=True) + return 1 + + missing = REQUIRED_CHANNELS - seen.keys() + if missing: + print( + f"[ws] FAIL — channels with no tick in {timeout:.0f}s window: {sorted(missing)}", + file=sys.stderr, + flush=True, + ) + return 1 + + print(f"[ws] OK — both channels delivered (saw {sorted(seen.keys())})", flush=True) + return 0 + + +def main() -> int: + ap = argparse.ArgumentParser() + ap.add_argument("--url", default=DEFAULT_URL) + ap.add_argument("--timeout", type=float, default=DEFAULT_TIMEOUT) + args = ap.parse_args() + return asyncio.run(run(args.url, args.timeout)) + + +if __name__ == "__main__": + sys.exit(main())