#!/usr/bin/env bash # Stonks Oracle — Deep Pipeline Health Check # Verifies every component is ACTUALLY working, not just "running" # # Usage: ./scripts/pipeline_health_check.sh [namespace] # Default namespace: stonks-oracle set -uo pipefail NS="${1:-stonks-oracle}" PASS=0 WARN=0 FAIL=0 DETAILS="" # Colors RED='\033[0;31m' GREEN='\033[0;32m' YELLOW='\033[1;33m' BLUE='\033[0;34m' NC='\033[0m' pass() { ((PASS++)); echo -e " ${GREEN}✓${NC} $1"; } warn() { ((WARN++)); echo -e " ${YELLOW}⚠${NC} $1"; } fail() { ((FAIL++)); echo -e " ${RED}✗${NC} $1"; } header() { echo -e "\n${BLUE}━━━ $1 ━━━${NC}"; } # ─── 1. POD HEALTH ───────────────────────────────────────────── header "Pod Health (namespace: $NS)" # Check all deployments are at desired replicas while IFS= read -r line; do name=$(echo "$line" | awk '{print $1}') ready=$(echo "$line" | awk '{print $2}') desired=$(echo "$ready" | cut -d/ -f2) actual=$(echo "$ready" | cut -d/ -f1) if [[ "$desired" == "0" ]]; then continue # intentionally scaled down fi if [[ "$actual" == "$desired" ]]; then pass "$name: $ready" else fail "$name: $ready (not fully ready)" fi done < <(kubectl get deployments -n "$NS" --no-headers 2>/dev/null) # Check for restarts in the last hour RESTART_PODS=$(kubectl get pods -n "$NS" --no-headers -o custom-columns=NAME:.metadata.name,RESTARTS:.status.containerStatuses[0].restartCount 2>/dev/null | awk '$2 > 0 {print $1 "(" $2 ")"}') if [[ -z "$RESTART_PODS" ]]; then pass "No pod restarts" else for p in $RESTART_PODS; do warn "Restart detected: $p" done fi # Check for OOMKilled in last state OOM_PODS=$(kubectl get pods -n "$NS" -o json 2>/dev/null | python3 -c " import json, sys data = json.load(sys.stdin) for pod in data.get('items', []): for cs in pod.get('status', {}).get('containerStatuses', []): last = cs.get('lastState', {}).get('terminated', {}) if last.get('reason') == 'OOMKilled': print(pod['metadata']['name']) " 2>/dev/null) if [[ -z "$OOM_PODS" ]]; then pass "No OOMKilled containers" else for p in $OOM_PODS; do fail "OOMKilled: $p" done fi # ─── 2. DATABASE CONNECTIVITY & HEALTH ──────────────────────── header "PostgreSQL" DB_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c " import asyncio, asyncpg, os, json async def check(): pool = await asyncpg.create_pool(dsn=f\"postgresql://{os.environ['POSTGRES_USER']}:{os.environ['POSTGRES_PASSWORD']}@{os.environ['POSTGRES_HOST']}:{os.environ.get('POSTGRES_PORT','5432')}/{os.environ['POSTGRES_DB']}\", min_size=1, max_size=2) result = {} result['connected'] = True # DB size row = await pool.fetchrow('SELECT pg_database_size(current_database()) as size') result['db_size_mb'] = round(row['size'] / 1024 / 1024) # Table count result['table_count'] = await pool.fetchval(\"SELECT count(*) FROM pg_tables WHERE schemaname='public'\") # Active connections result['connections'] = await pool.fetchval('SELECT count(*) FROM pg_stat_activity WHERE datname = current_database()') await pool.close() print(json.dumps(result)) asyncio.run(check()) " 2>/dev/null) || DB_CHECK='{"connected": false}' if echo "$DB_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin); exit(0 if d.get('connected') else 1)" 2>/dev/null; then DB_SIZE=$(echo "$DB_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin)['db_size_mb'])") DB_TABLES=$(echo "$DB_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin)['table_count'])") DB_CONNS=$(echo "$DB_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin)['connections'])") pass "Connected (${DB_SIZE}MB, ${DB_TABLES} tables, ${DB_CONNS} connections)" else fail "Cannot connect to PostgreSQL" fi # ─── 3. REDIS CONNECTIVITY & QUEUES ─────────────────────────── header "Redis" REDIS_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c " import redis, os, json r = redis.from_url(f'redis://:{os.environ.get(\"REDIS_PASSWORD\",\"\")}@{os.environ.get(\"REDIS_HOST\",\"redis-master.redis-service.svc.cluster.local\")}:6379/{os.environ.get(\"REDIS_DB\",\"1\")}') result = {} result['connected'] = True result['memory'] = r.info('memory')['used_memory_human'] result['keys'] = r.dbsize() # Check queues for backlog queues = {} for q in ['ingestion','parsing','extraction','aggregation','recommendation','trading','broker_orders','macro_classification']: stage = os.environ.get('DEPLOY_STAGE', '') prefix = f'stonks:{stage}' if stage else 'stonks' key = f'{prefix}:queue:{q}' length = r.llen(key) if length > 0: queues[q] = length result['queue_backlog'] = queues print(json.dumps(result)) " 2>/dev/null) || REDIS_CHECK='{"connected": false}' if echo "$REDIS_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin); exit(0 if d.get('connected') else 1)" 2>/dev/null; then R_MEM=$(echo "$REDIS_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin)['memory'])") R_KEYS=$(echo "$REDIS_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin)['keys'])") pass "Connected (${R_MEM}, ${R_KEYS} keys)" BACKLOG=$(echo "$REDIS_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin)['queue_backlog']; print(' '.join(f'{k}={v}' for k,v in d.items()) if d else '')") if [[ -z "$BACKLOG" ]]; then pass "All queues drained (no backlog)" else warn "Queue backlog: $BACKLOG" fi else fail "Cannot connect to Redis" fi # ─── 4. MINIO CONNECTIVITY ──────────────────────────────────── header "MinIO (S3)" MINIO_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c " from minio import Minio import os, json client = Minio( os.environ.get('MINIO_ENDPOINT', 'minio.minio-service.svc.cluster.local:80'), access_key=os.environ.get('MINIO_ACCESS_KEY', ''), secret_key=os.environ.get('MINIO_SECRET_KEY', ''), secure=False ) buckets = client.list_buckets() print(json.dumps({'connected': True, 'bucket_count': len(buckets)})) " 2>/dev/null) || MINIO_CHECK='{"connected": false}' if echo "$MINIO_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin); exit(0 if d.get('connected') else 1)" 2>/dev/null; then BUCKETS=$(echo "$MINIO_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin)['bucket_count'])") pass "Connected ($BUCKETS buckets)" else fail "Cannot connect to MinIO" fi # ─── 5. vLLM / AI MODEL ─────────────────────────────────────── header "vLLM (NuExtract3)" VLLM_CHECK=$(kubectl exec -n "$NS" deployment/extractor -- python3 -c " import httpx, json, os url = os.environ.get('VLLM_BASE_URL', os.environ.get('OLLAMA_BASE_URL', '')) try: resp = httpx.get(f'{url}/v1/models', timeout=10) models = resp.json().get('data', []) names = [m['id'] for m in models] print(json.dumps({'connected': True, 'models': names})) except Exception as e: print(json.dumps({'connected': False, 'error': str(e)})) " 2>/dev/null) || VLLM_CHECK='{"connected": false}' if echo "$VLLM_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin); exit(0 if d.get('connected') else 1)" 2>/dev/null; then MODELS=$(echo "$VLLM_CHECK" | python3 -c "import json,sys; print(', '.join(json.load(sys.stdin)['models']))") pass "Connected (models: $MODELS)" else fail "Cannot reach vLLM" fi # ─── 6. BROKER CONNECTIVITY ─────────────────────────────────── header "Alpaca Broker" BROKER_CHECK=$(kubectl exec -n "$NS" deployment/broker-adapter -- python3 -c " import httpx, json, os key = os.environ.get('BROKER_API_KEY', '') secret = os.environ.get('BROKER_API_SECRET', '') base = os.environ.get('BROKER_BASE_URL', 'https://paper-api.alpaca.markets') try: resp = httpx.get(f'{base}/v2/account', headers={'APCA-API-KEY-ID': key, 'APCA-API-SECRET-KEY': secret}, timeout=10) if resp.status_code == 200: acct = resp.json() print(json.dumps({'connected': True, 'equity': acct.get('equity'), 'buying_power': acct.get('buying_power'), 'status': acct.get('status')})) else: print(json.dumps({'connected': False, 'status_code': resp.status_code})) except Exception as e: print(json.dumps({'connected': False, 'error': str(e)})) " 2>/dev/null) || BROKER_CHECK='{"connected": false}' if echo "$BROKER_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin); exit(0 if d.get('connected') else 1)" 2>/dev/null; then EQUITY=$(echo "$BROKER_CHECK" | python3 -c "import json,sys; d=json.load(sys.stdin); print(f'equity=\${float(d[\"equity\"]):,.0f} buying_power=\${float(d[\"buying_power\"]):,.0f}')") pass "Connected ($EQUITY)" else fail "Cannot reach Alpaca API" fi # ─── 7. PIPELINE FLOW (the real test) ───────────────────────── header "Pipeline Flow (last 6 hours)" FLOW_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c " import asyncio, asyncpg, os, json from datetime import datetime, timezone async def check(): pool = await asyncpg.create_pool(dsn=f\"postgresql://{os.environ['POSTGRES_USER']}:{os.environ['POSTGRES_PASSWORD']}@{os.environ['POSTGRES_HOST']}:{os.environ.get('POSTGRES_PORT','5432')}/{os.environ['POSTGRES_DB']}\", min_size=1, max_size=2) r = {} now = datetime.now(tz=timezone.utc) # Ingestion: new docs created r['docs_ingested_6h'] = await pool.fetchval(\"SELECT count(*) FROM documents WHERE created_at > now() - interval '6 hours'\") # Extraction: docs extracted recently r['docs_extracted_6h'] = await pool.fetchval(\"SELECT count(*) FROM documents WHERE status = 'extracted' AND updated_at > now() - interval '6 hours'\") # Stuck docs r['stuck_parsed'] = await pool.fetchval(\"SELECT count(*) FROM documents WHERE status = 'parsed' AND updated_at < now() - interval '1 hour'\") r['stuck_ingested'] = await pool.fetchval(\"SELECT count(*) FROM documents WHERE status = 'ingested' AND updated_at < now() - interval '1 hour'\") # Aggregation: trend history entries r['trends_6h'] = await pool.fetchval(\"SELECT count(*) FROM trend_history WHERE generated_at > now() - interval '6 hours'\") # Recommendations generated r['recs_6h'] = await pool.fetchval(\"SELECT count(*) FROM recommendations WHERE created_at > now() - interval '6 hours'\") # Trading decisions r['decisions_6h'] = await pool.fetchval(\"SELECT count(*) FROM trading_decisions WHERE created_at > now() - interval '6 hours'\") # Global events (macro pipeline) r['global_events_6h'] = await pool.fetchval(\"SELECT count(*) FROM global_events WHERE created_at > now() - interval '6 hours'\") # Prediction snapshots r['snapshots_6h'] = await pool.fetchval(\"SELECT count(*) FROM prediction_snapshots WHERE created_at > now() - interval '6 hours'\") # Market data freshness latest_market = await pool.fetchval('SELECT max(captured_at) FROM market_snapshots') if latest_market: r['market_data_age_min'] = round((now - latest_market).total_seconds() / 60) else: r['market_data_age_min'] = -1 await pool.close() print(json.dumps(r)) asyncio.run(check()) " 2>/dev/null) || FLOW_CHECK='{}' if [[ -n "$FLOW_CHECK" && "$FLOW_CHECK" != "{}" ]]; then DOCS_IN=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('docs_ingested_6h', 0))") DOCS_EX=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('docs_extracted_6h', 0))") STUCK_P=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('stuck_parsed', 0))") STUCK_I=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('stuck_ingested', 0))") TRENDS=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('trends_6h', 0))") RECS=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('recs_6h', 0))") DECISIONS=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('decisions_6h', 0))") EVENTS=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('global_events_6h', 0))") SNAPS=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('snapshots_6h', 0))") MKT_AGE=$(echo "$FLOW_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('market_data_age_min', -1))") # Ingestion flowing? if [[ "$DOCS_IN" -gt 0 ]]; then pass "Ingestion: $DOCS_IN docs ingested" else fail "Ingestion: ZERO docs in last 6h" fi # Extraction flowing? if [[ "$DOCS_EX" -gt 0 ]]; then pass "Extraction: $DOCS_EX docs extracted" else fail "Extraction: ZERO docs extracted in last 6h" fi # Stuck documents? if [[ "$STUCK_P" -eq 0 && "$STUCK_I" -eq 0 ]]; then pass "No stuck documents" else [[ "$STUCK_P" -gt 0 ]] && warn "Stuck in 'parsed': $STUCK_P docs (>1h old)" [[ "$STUCK_I" -gt 0 ]] && warn "Stuck in 'ingested': $STUCK_I docs (>1h old)" fi # Aggregation producing trends? if [[ "$TRENDS" -gt 0 ]]; then pass "Aggregation: $TRENDS trend entries" else warn "Aggregation: 0 trends (expected on weekends/off-hours)" fi # Recommendations flowing? if [[ "$RECS" -gt 0 ]]; then pass "Recommendations: $RECS generated" else fail "Recommendations: ZERO in last 6h" fi # Trading decisions happening? if [[ "$DECISIONS" -gt 0 ]]; then pass "Trading engine: $DECISIONS decisions" else warn "Trading engine: 0 decisions (market may be closed)" fi # Macro events? if [[ "$EVENTS" -gt 0 ]]; then pass "Macro pipeline: $EVENTS global events classified" else warn "Macro pipeline: 0 events (may be low news volume)" fi # Prediction snapshots? if [[ "$SNAPS" -gt 0 ]]; then pass "Validation: $SNAPS prediction snapshots" else warn "Validation: 0 snapshots in 6h" fi # Market data freshness if [[ "$MKT_AGE" -ge 0 && "$MKT_AGE" -lt 30 ]]; then pass "Market data: ${MKT_AGE}min old" elif [[ "$MKT_AGE" -ge 30 && "$MKT_AGE" -lt 120 ]]; then warn "Market data: ${MKT_AGE}min old (stale)" elif [[ "$MKT_AGE" -ge 120 ]]; then # Could be weekend/after hours HOURS=$((MKT_AGE / 60)) warn "Market data: ${HOURS}h old (weekend/after-hours?)" fi else fail "Could not query pipeline flow data" fi # ─── 8. QUALITY GATE STATUS ─────────────────────────────────── header "Model Quality Gate" QG_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c " import asyncio, asyncpg, os, json from datetime import datetime, timezone async def check(): pool = await asyncpg.create_pool(dsn=f\"postgresql://{os.environ['POSTGRES_USER']}:{os.environ['POSTGRES_PASSWORD']}@{os.environ['POSTGRES_HOST']}:{os.environ.get('POSTGRES_PORT','5432')}/{os.environ['POSTGRES_DB']}\", min_size=1, max_size=2) row = await pool.fetchrow(\"\"\" SELECT generated_at, prediction_count, win_rate, directional_accuracy, brier_score FROM model_metric_snapshots WHERE lookback_window = '30d' AND horizon = '7d' ORDER BY generated_at DESC LIMIT 1 \"\"\") r = {} if row: age_h = (datetime.now(tz=timezone.utc) - row['generated_at']).total_seconds() / 3600 r['age_hours'] = round(age_h, 1) r['prediction_count'] = row['prediction_count'] r['win_rate'] = float(row['win_rate']) if row['win_rate'] else 0 r['directional_accuracy'] = float(row['directional_accuracy']) if row['directional_accuracy'] else 0 r['brier_score'] = float(row['brier_score']) if row['brier_score'] else 0 else: r['age_hours'] = -1 await pool.close() print(json.dumps(r)) asyncio.run(check()) " 2>/dev/null) || QG_CHECK='{}' if [[ -n "$QG_CHECK" && "$QG_CHECK" != "{}" ]]; then QG_AGE=$(echo "$QG_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('age_hours', -1))") QG_PREDS=$(echo "$QG_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('prediction_count', 0))") QG_WR=$(echo "$QG_CHECK" | python3 -c "import json,sys; print(f\"{json.load(sys.stdin).get('win_rate', 0)*100:.1f}%\")") QG_DA=$(echo "$QG_CHECK" | python3 -c "import json,sys; print(f\"{json.load(sys.stdin).get('directional_accuracy', 0)*100:.1f}%\")") if python3 -c "exit(0 if $QG_AGE >= 0 and $QG_AGE <= 48 else 1)" 2>/dev/null; then pass "Metric snapshot: ${QG_AGE}h old (predictions=$QG_PREDS, win_rate=$QG_WR, accuracy=$QG_DA)" elif python3 -c "exit(0 if $QG_AGE > 48 else 1)" 2>/dev/null; then fail "Metric snapshot STALE: ${QG_AGE}h old (max 48h) — quality gate failing" else fail "No metric snapshot exists" fi fi # ─── 9. PORTFOLIO STATUS ────────────────────────────────────── header "Paper Trading Portfolio" PORT_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c " import asyncio, asyncpg, os, json async def check(): pool = await asyncpg.create_pool(dsn=f\"postgresql://{os.environ['POSTGRES_USER']}:{os.environ['POSTGRES_PASSWORD']}@{os.environ['POSTGRES_HOST']}:{os.environ.get('POSTGRES_PORT','5432')}/{os.environ['POSTGRES_DB']}\", min_size=1, max_size=2) r = {} positions = await pool.fetch('SELECT ticker, quantity, avg_entry_price, current_price, unrealized_pnl FROM positions') r['position_count'] = len(positions) r['total_invested'] = sum(float(p['avg_entry_price']) * float(p['quantity']) for p in positions) r['unrealized_pnl'] = sum(float(p['unrealized_pnl'] or 0) for p in positions) r['tickers'] = [p['ticker'] for p in positions] # Recent orders r['orders_24h'] = await pool.fetchval(\"SELECT count(*) FROM orders WHERE created_at > now() - interval '24 hours'\") r['total_orders'] = await pool.fetchval('SELECT count(*) FROM orders') await pool.close() print(json.dumps(r)) asyncio.run(check()) " 2>/dev/null) || PORT_CHECK='{}' if [[ -n "$PORT_CHECK" && "$PORT_CHECK" != "{}" ]]; then POS_CNT=$(echo "$PORT_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('position_count', 0))") INVESTED=$(echo "$PORT_CHECK" | python3 -c "import json,sys; print(f\"\${json.load(sys.stdin).get('total_invested', 0):,.0f}\")") PNL=$(echo "$PORT_CHECK" | python3 -c "import json,sys; print(f\"\${json.load(sys.stdin).get('unrealized_pnl', 0):,.0f}\")") TICKERS=$(echo "$PORT_CHECK" | python3 -c "import json,sys; print(', '.join(json.load(sys.stdin).get('tickers', [])))") ORDERS_24=$(echo "$PORT_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('orders_24h', 0))") TOTAL_ORD=$(echo "$PORT_CHECK" | python3 -c "import json,sys; print(json.load(sys.stdin).get('total_orders', 0))") pass "$POS_CNT positions ($TICKERS)" pass "Invested: $INVESTED | Unrealized P&L: $PNL" pass "Orders: $ORDERS_24 today, $TOTAL_ORD lifetime" fi # ─── 10. API ENDPOINTS ──────────────────────────────────────── header "API Endpoints" # Query API health QA_STATUS=$(kubectl exec -n "$NS" deployment/query-api -- curl -sf http://localhost:8000/health 2>/dev/null) && pass "Query API: healthy" || fail "Query API: unreachable" # Symbol Registry health SR_STATUS=$(kubectl exec -n "$NS" deployment/symbol-registry -- curl -sf http://localhost:8000/health 2>/dev/null) && pass "Symbol Registry: healthy" || fail "Symbol Registry: unreachable" # Trading Engine health TE_STATUS=$(kubectl exec -n "$NS" deployment/trading-engine -- curl -sf http://localhost:8000/health 2>/dev/null) && pass "Trading Engine: healthy" || fail "Trading Engine: unreachable" # Risk Engine health RISK_STATUS=$(kubectl exec -n "$NS" deployment/risk -- curl -sf http://localhost:8000/health 2>/dev/null) && pass "Risk Engine: healthy" || fail "Risk Engine: unreachable" # Specialist NER service health SPEC_STATUS=$(kubectl exec -n "$NS" deployment/specialist -- curl -sf http://localhost:8000/health 2>/dev/null) && pass "Specialist: healthy" || fail "Specialist: unreachable" # ─── 11. RESOURCE USAGE ─────────────────────────────────────── header "Resource Usage" TOTAL_MEM=$(kubectl top pods -n "$NS" --no-headers 2>/dev/null | awk '{sum += $3} END {print sum}') TOTAL_CPU=$(kubectl top pods -n "$NS" --no-headers 2>/dev/null | awk '{sum += $2} END {print sum}') pass "Total: ${TOTAL_CPU}m CPU, ${TOTAL_MEM}Mi memory" # Any pod over 80% of its limit? HIGH_MEM=$(kubectl top pods -n "$NS" --no-headers 2>/dev/null | awk '$3+0 > 200 {print $1 "=" $3}') if [[ -n "$HIGH_MEM" ]]; then warn "High memory pods: $HIGH_MEM" fi # ─── 12. END-TO-END LATENCY CHECK ───────────────────────────── header "End-to-End Pipeline Latency" # Copy helper script to pod and run it LATENCY_CHECK=$(kubectl exec -n "$NS" deployment/query-api -- python3 -c ' import asyncio, asyncpg, os, json from datetime import datetime, timezone, timedelta async def check(): dsn = "postgresql://{}:{}@{}:{}/{}".format(os.environ["POSTGRES_USER"], os.environ["POSTGRES_PASSWORD"], os.environ["POSTGRES_HOST"], os.environ.get("POSTGRES_PORT","5432"), os.environ["POSTGRES_DB"]) pool = await asyncpg.create_pool(dsn=dsn, min_size=1, max_size=2) r = {} now = datetime.now(tz=timezone.utc) cutoff = now - timedelta(hours=24) row = await pool.fetchrow("SELECT avg(extract(epoch from (updated_at - created_at))) as avg_sec, max(extract(epoch from (updated_at - created_at))) as max_sec FROM documents WHERE status = $1 AND updated_at > $2 AND updated_at > created_at", "extracted", cutoff) if row and row["avg_sec"]: r["avg_extraction_sec"] = round(float(row["avg_sec"]), 1) r["max_extraction_sec"] = round(float(row["max_sec"]), 1) latest_trend = await pool.fetchval("SELECT max(generated_at) FROM trend_windows") if latest_trend: r["trend_age_min"] = round((now - latest_trend).total_seconds() / 60) latest_rec = await pool.fetchval("SELECT max(created_at) FROM recommendations") if latest_rec: r["rec_age_min"] = round((now - latest_rec).total_seconds() / 60) await pool.close() print(json.dumps(r)) asyncio.run(check()) ' 2>/dev/null) || LATENCY_CHECK='{}' if [[ -n "$LATENCY_CHECK" && "$LATENCY_CHECK" != "{}" ]]; then AVG_EXT=$(echo "$LATENCY_CHECK" | python3 -c "import json,sys; v=json.load(sys.stdin).get('avg_extraction_sec'); print(f'{v}s' if v else 'N/A')") MAX_EXT=$(echo "$LATENCY_CHECK" | python3 -c "import json,sys; v=json.load(sys.stdin).get('max_extraction_sec'); print(f'{v}s' if v else 'N/A')") TREND_AGE=$(echo "$LATENCY_CHECK" | python3 -c "import json,sys; v=json.load(sys.stdin).get('trend_age_min'); print(f'{v}min' if v else 'N/A')") REC_AGE=$(echo "$LATENCY_CHECK" | python3 -c "import json,sys; v=json.load(sys.stdin).get('rec_age_min'); print(f'{v}min' if v else 'N/A')") pass "Doc ingestion→extraction: avg=${AVG_EXT}, max=${MAX_EXT}" if echo "$LATENCY_CHECK" | python3 -c "import json,sys; v=json.load(sys.stdin).get('trend_age_min',999); exit(0 if v < 30 else 1)" 2>/dev/null; then pass "Latest trend: ${TREND_AGE} ago" else warn "Latest trend: ${TREND_AGE} ago" fi if echo "$LATENCY_CHECK" | python3 -c "import json,sys; v=json.load(sys.stdin).get('rec_age_min',999); exit(0 if v < 30 else 1)" 2>/dev/null; then pass "Latest recommendation: ${REC_AGE} ago" else warn "Latest recommendation: ${REC_AGE} ago" fi else warn "Could not measure pipeline latency" fi # ─── SUMMARY ────────────────────────────────────────────────── header "Summary" echo -e " ${GREEN}PASS: $PASS${NC} ${YELLOW}WARN: $WARN${NC} ${RED}FAIL: $FAIL${NC}" echo "" if [[ "$FAIL" -gt 0 ]]; then echo -e "${RED}Pipeline has failures that need attention.${NC}" exit 1 elif [[ "$WARN" -gt 3 ]]; then echo -e "${YELLOW}Pipeline is running but has multiple warnings.${NC}" exit 0 else echo -e "${GREEN}Pipeline is healthy and flowing.${NC}" exit 0 fi