jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124#!/bin/bash# Experiment 6 check. Reports every rate with its discriminators attached, so a# short-window artifact cannot be presented as a trend. See memory# feedback_monitoring_window_discipline.set -uIP=37.27.43.39K=$HOME/.ssh/waow_ed25519SC="${SC:-$HOME/.stream-exp6}"; mkdir -p "$SC"STATE=$SC/check6.stateSTART=1785214907DEADLINE=1785862907
S() { ssh -i "$K" -o IdentitiesOnly=yes -o ConnectTimeout=10 root@$IP "$1" 2>/dev/null; }Q() { S "cd /opt/stream-experiment; docker compose exec -T prometheus wget -qO- 'http://127.0.0.1:9090/api/v1/query?query=$1' 2>/dev/null"; }scalar() { python3 -c "import json,systry: r=json.load(sys.stdin)['data']['result'] print(r[0]['value'][1] if r else '')except Exception: print('')"; }
[ -s "$SC/EVENTS.log" ] && { echo '*** EVENTS.log NON-EMPTY — MERGE MAY HAVE STARTED ***'; cat "$SC/EVENTS.log"; }
# Control read first. NONE means mid-restart; every other value is unusable.phase=""for i in $(seq 1 10); do phase=$(Q jetstream_orchestrator_phase | scalar) [ -n "$phase" ] && break sleep 25done[ -z "$phase" ] && { echo "phase=NONE after 10 attempts (~4min) — mid-restart, no other value is trustworthy"; exit 0; }echo "phase=$phase"
now=$(S 'date +%s')comp=$(Q 'stream_backfill_repos_durable%7Bstatus%3D%22complete%22%7D' | scalar)disk=$(S 'df -B1 --output=used /data | tail -1' | tr -dc '0-9')notstarted=$(Q 'stream_backfill_repos_durable%7Bstatus%3D%22not_started%22%7D' | scalar)failed=$(Q 'stream_backfill_repos_durable%7Bstatus%3D%22failed%22%7D' | scalar)[ -z "$comp" ] && { echo "complete unreadable — treat as mid-restart"; exit 0; }
echo "complete=$comp not_started=$notstarted failed=$failed"
# Measured delta since the last run of this script: immune to window artifacts.if [ -f "$STATE" ]; then read -r p_ts p_comp p_failed p_disk < "$STATE" dt=$((now - p_ts)); dc=$((comp - p_comp)); df=$((failed - p_failed)) [ "$dt" -gt 0 ] && printf 'MEASURED since last check (%.2fh): %.1f repos/s, %d new failures (%.0f/hr)\n' \ "$(python3 -c "print($dt/3600)")" "$(python3 -c "print($dc/$dt)")" "$df" "$(python3 -c "print($df/$dt*3600)")" if [ -n "${p_disk:-}" ] && [ "$dc" -gt 0 ] && [ -n "$disk" ]; then python3 - "$((disk - p_disk))" "$dc" "$disk" "$comp" <<'PYD'import sysdd,dc,disk,comp=[int(x) for x in sys.argv[1:5]]per=dd/dcprint(f"MARGINAL disk: {dd/1e9:+.2f} GB over {dc:,} repos = {per/1024:+.0f} KiB/repo")if per <= 0: print(" (<=0 means compaction reclaimed more than the window wrote; needs a longer baseline)")else: for tgt,name in ((19_400_000,'19.4M'),(22_845_221,'22.85M')): proj=disk+per*(tgt-comp); free=3_219_914_752_000-proj warn=" <<< EXCEEDS VOLUME" if free<0 else (" <<< below 0.4TB floor" if free<0.4e12 else "") print(f" projected at {name}: {proj/1e12:.2f} TB used, {free/1e12:+.2f} TB free{warn}")PYD fifiecho "$now $comp $failed $disk" > "$STATE"
# The backfill-era numbers below (crawl rate, per-repo cost, batch troughs) are# meaningless once bootstrap is over. Report merge/serving state instead, or a# reader sees "1.3 repos/s" in phase 2 and thinks the thing has died.if [ "$phase" != "1" ]; then echo "--- phase $phase: bootstrap is over, backfill rates below do not apply ---" for m in merge_events_kept_total merge_events_dropped_total merge_segments_consumed_total; do printf ' %s = %s\n' "$m" "$(Q "jetstream_orchestrator_$m" | scalar)" done printf ' compaction_watermark_seq = %s (0 = compaction has not started)\n' \ "$(Q stream_compaction_watermark_seq | scalar)" printf ' pending (repos left in the pending pass) = %s\n' \ "$(Q 'stream_backfill_repos_durable%7Bstatus%3D%22pending%22%7D' | scalar)" code=$(S 'cd /opt/stream-experiment; docker compose exec -T prometheus wget -S -qO- http://172.18.0.2:8080/xrpc/jetstream.listSegments 2>&1 | grep -m1 -oE "HTTP/1.1 [0-9]+"') printf ' listSegments -> %s (503 until steady_state; 200 is the win)\n' "${code:-unreadable}" printf ' backfill tree = %s (GONE means merge cleanup finished)\n' \ "$(S 'du -sh /data/stream/backfill 2>/dev/null | cut -f1 || echo GONE')" echo " replay window: run 'ssh root@37.27.43.39 /usr/local/bin/wsprobe 32408688748'" echo " 101 + a binary frame = still replayable"fi
# Nested windows. A trend requires these to move together. BOOTSTRAP ONLY.[ "$phase" = "1" ] && printf 'rate windows:'if [ "$phase" = "1" ]; thenfor w in 1h 6h 12h; do r=$(Q "rate(stream_backfill_repos_durable%7Bstatus%3D%22complete%22%7D%5B$w%5D)" | scalar) printf ' %s=%s' "$w" "$(python3 -c "print(round(float('${r:-0}'),1))")"done; echofi
# Queue state qualifies the 1h number: near-zero not_started == batch trough.if [ "$phase" = "1" ] && [ -n "$notstarted" ] && [ "${notstarted%.*}" -lt 1000 ]; then echo " ^ NOTE: not_started=$notstarted (batch trough) — a low 1h rate here is EXPECTED, not a fault"fi
# Per-repo cost: rising means big repos, falling means the tail is thinning.pr=$(Q 'rate(jetstream_backfill_handle_repo_duration_seconds_sum%5B2h%5D)/rate(jetstream_backfill_handle_repo_duration_seconds_count%5B2h%5D)' | scalar)[ "$phase" = "1" ] && [ -n "$pr" ] && printf 'per-repo: %.1f ms (falling => tail thinning, NOT degradation; bytes/s falls with it)\n' \ "$(python3 -c "print(float('$pr')*1000)")"
# Per-process counter: it RESETS on every restart, so only increase() is# meaningful. Reading the raw value and dividing is how you get a number that# is off by two orders of magnitude.rec=$(Q "increase(jetstream_ingest_events_appended_total%5B24h%5D)" | scalar)[ -n "$rec" ] && printf 'records appended (24h): %.2f B\n' "$(python3 -c "print(float('$rec')/1e9)")"
rss=$(Q process_resident_memory_bytes | scalar)[ -n "$rss" ] && printf 'RSS %.1f GB (blooms; ~46-50 GB is normal)\n' "$(python3 -c "print(float('$rss')/1e9)")"S 'df -h /data | tail -1 docker inspect -f "restarts={{.RestartCount}}" stream-experiment-stream-1 echo "oom_kills=$(dmesg 2>/dev/null | grep -c "Out of memory: Killed process")" echo "recovery_resumed=$(docker logs --tail 3000 stream-experiment-stream-1 2>&1 | grep -c "recovery: resumed active segment")"'
python3 - "$now" "$comp" <<'PY'import sys, datetime as dtnow=int(sys.argv[1]); comp=int(float(sys.argv[2])); start=1785214907el=(now-start)/3600print(f"elapsed {el:.1f}h of 180h {comp/19_400_000*100:.1f}% of 19.4M / {comp/22_845_221*100:.1f}% of 22.85M")PY