GPU fleet: Vast client, box setup (0.3.12 node + patched SP1 server + cuda host), the phase-1 matrix, the Ember ladder, the Linux segment prover, the fleet page builder
This commit is contained in:
parent
f11b02eddd
commit
627a2256c7
9 changed files with 1249 additions and 0 deletions
86
tools/fleet/box-ember.sh
Executable file
86
tools/fleet/box-ember.sh
Executable file
|
|
@ -0,0 +1,86 @@
|
|||
#!/usr/bin/env bash
|
||||
# Ember Tune's two-knob ladder on a rented NVIDIA card (docs/plans/ember-tune.md, branch ember-tune): the miner runs
|
||||
# throughout; the power ladder 100, 90, 80, 70, 60, 50% of the default limit at the unlocked clock (clamped at the
|
||||
# card's reported minimum), then the clock ladder 90, 80, 70, 60% of the maximum graphics clock at the power the
|
||||
# first ladder chose; 15 s settle and 60 s hold per step; the choice is the best MH/W among the steps whose rate is
|
||||
# within 1% of the fastest. Linux root: nvidia-smi -pl and -lgc, no prompt. Every step prints a RESULT line; the
|
||||
# ladder and the choice go to /root/fleet/out/ember.json as a TUNE-shaped record (relay/lib/ember.mjs parseRecords).
|
||||
set -uo pipefail
|
||||
F=/root/fleet; OUT=$F/out; LOG=$OUT/ember.log; B=/opt/igneum/pkg/bin
|
||||
LABEL="${LABEL:-box}"; WALLET="${WALLET:-0x1919191919191919191919191919191919191919}"
|
||||
SETTLE="${SETTLE:-15}"; HOLD="${HOLD:-60}"
|
||||
mkdir -p $OUT $F/mine/packs
|
||||
exec > >(tee -a $LOG) 2>&1
|
||||
stamp() { date -u +%Y-%m-%dT%H:%M:%SZ; }
|
||||
say() { echo "$(stamp) $*"; }
|
||||
q() { nvidia-smi --query-gpu="$1" --format=csv,noheader,nounits -i 0 2>/dev/null | head -1 | tr -d ' '; }
|
||||
NAME="$(q name)"; DRV="$(q driver_version)"; PDEF="$(q power.default_limit)"; PMIN="$(q power.min_limit)"; PMAX="$(q power.max_limit)"; CMAX="$(q clocks.max.graphics)"
|
||||
pkill -f sp1-gpu-server 2>/dev/null; rm -f /tmp/sp1-cuda-*.sock
|
||||
echo "RESULT start $(stamp) card=$NAME driver=$DRV power_default_w=$PDEF min_w=$PMIN max_w=$PMAX clock_max_mhz=$CMAX"
|
||||
# can we set anything?
|
||||
nvidia-smi -i 0 -pl "$PDEF" >/dev/null 2>&1 && PL_OK=1 || PL_OK=0
|
||||
nvidia-smi -i 0 -lgc 0,"$CMAX" >/dev/null 2>&1 && LGC_OK=1 || LGC_OK=0
|
||||
nvidia-smi -i 0 -rgc >/dev/null 2>&1
|
||||
echo "RESULT knobs power_limit_settable=$PL_OK clock_cap_settable=$LGC_OK"
|
||||
cd $F/mine
|
||||
rm -rf packs/devnet; $B/igneum-miner export-pack grpc://127.0.0.1:26610 packs/devnet > $OUT/ember-export.log 2>&1
|
||||
nohup $B/igneum-miner mine grpc://127.0.0.1:26610 1 100000000 "$LABEL" --worker $B/igneum-worker-cuda --worker-args "--device 0 --pack packs/devnet" \
|
||||
--prepare-packs packs/prepare --exit-on-seed-change --evm-address "$WALLET" --payout-label "$LABEL" --status-secs 10 > $OUT/ember-miner.log 2>&1 &
|
||||
MPID=$!; cd $F
|
||||
cleanup() { nvidia-smi -i 0 -rgc >/dev/null 2>&1; [ "$PL_OK" = 1 ] && nvidia-smi -i 0 -pl "$PDEF" >/dev/null 2>&1; kill $MPID 2>/dev/null; pkill -f igneum-worker-cuda 2>/dev/null; }
|
||||
trap cleanup EXIT
|
||||
say "miner warming 90 s"; sleep 90
|
||||
STEPS=$OUT/ember-steps.jsonl; : > $STEPS
|
||||
step() { # <label> <power_pct> <limit_w> <clock_cap_mhz or 0>
|
||||
local lab="$1" pct="$2" lim="$3" cap="$4"
|
||||
if [ "$PL_OK" = 1 ]; then nvidia-smi -i 0 -pl "$lim" >/dev/null 2>&1 || say "could not set -pl $lim"; fi
|
||||
if [ "$cap" != 0 ]; then nvidia-smi -i 0 -lgc 0,"$cap" >/dev/null 2>&1 || say "could not set -lgc $cap"; else nvidia-smi -i 0 -rgc >/dev/null 2>&1; fi
|
||||
sleep "$SETTLE"
|
||||
local n0; n0="$(grep -c 'now=' $OUT/ember-miner.log)"
|
||||
nvidia-smi --query-gpu=power.draw,clocks.sm,clocks.mem,temperature.gpu --format=csv,noheader,nounits -i 0 -l 1 > $OUT/ember-samp-$lab.csv 2>/dev/null & local sp=$!
|
||||
sleep "$HOLD"; kill $sp 2>/dev/null; wait $sp 2>/dev/null
|
||||
local n1; n1="$(grep -c 'now=' $OUT/ember-miner.log)"
|
||||
local mhs; mhs="$(grep -o 'now=[0-9.]*' $OUT/ember-miner.log | tail -n $((n1 - n0 > 0 ? n1 - n0 : 1)) | cut -d= -f2 | awk '{s+=$1; n++} END {if (n) printf "%.2f", s/n; else print 0}')"
|
||||
local w gclk mclk tmax; read -r w gclk mclk tmax <<< "$(awk -F', *' '{w+=$1; g+=$2; m+=$3; if ($4+0 > t) t=$4+0; n++} END {if (n) printf "%.1f %.0f %.0f %d", w/n, g/n, m/n, t; else print "0 0 0 0"}' $OUT/ember-samp-$lab.csv)"
|
||||
local eff; eff="$(awk -v a="$mhs" -v b="$w" 'BEGIN {if (b > 0) printf "%.4f", a / b; else print 0}')"
|
||||
echo "RESULT step $lab power_pct=$pct limit_w=$lim clock_cap_mhz=$cap watts=$w mhs=$mhs mhw=$eff gclk=$gclk mclk=$mclk tmax=$tmax"
|
||||
printf '%s\n' "{\"label\":\"$lab\",\"clock_mhz\":$cap,\"power_pct\":$pct,\"limit_w\":$lim,\"watts\":$w,\"mhs\":$mhs,\"eff\":$eff,\"gclk\":$gclk,\"mclk\":$mclk,\"tmax\":$tmax,\"faults\":0,\"mark\":\"ok\"}" >> $STEPS
|
||||
}
|
||||
# the power ladder
|
||||
last=""
|
||||
for pct in 100 90 80 70 60 50; do
|
||||
lim="$(awk -v d="$PDEF" -v p="$pct" -v mn="$PMIN" 'BEGIN {l = d * p / 100; if (l < mn) l = mn; printf "%d", l}')"
|
||||
[ "$lim" = "$last" ] && { say "power $pct% clamps to the same $lim W; skipped"; continue; }
|
||||
last="$lim"; step "p$pct" "$pct" "$lim" 0
|
||||
[ "$PL_OK" = 0 ] && { say "power limit not settable here; one baseline step only"; break; }
|
||||
done
|
||||
# choose the power point: best eff within 1% of the top rate
|
||||
read -r BEST_PCT BEST_LIM <<< "$(python3 - $STEPS <<'PY'
|
||||
import json, sys
|
||||
s = [json.loads(l) for l in open(sys.argv[1]) if l.strip()]
|
||||
top = max(x["mhs"] for x in s); ok = [x for x in s if x["mhs"] >= top * 0.99]
|
||||
b = max(ok, key=lambda x: x["eff"]); print(b["power_pct"], b["limit_w"])
|
||||
PY
|
||||
)"
|
||||
echo "RESULT power_choice pct=$BEST_PCT limit_w=$BEST_LIM"
|
||||
if [ "$LGC_OK" = 1 ]; then
|
||||
for cpct in 90 80 70 60; do
|
||||
cap="$(awk -v c="$CMAX" -v p="$cpct" 'BEGIN {printf "%d", int(c * p / 100 / 10) * 10}')"
|
||||
step "c$cpct" "$BEST_PCT" "$BEST_LIM" "$cap"
|
||||
done
|
||||
fi
|
||||
python3 - $STEPS $OUT/ember.json "$NAME" "$DRV" "$LABEL" "$PDEF" "$CMAX" "$PL_OK" "$LGC_OK" <<'PY'
|
||||
import json, sys, datetime
|
||||
steps = [json.loads(l) for l in open(sys.argv[1]) if l.strip()]
|
||||
name, drv, label, pdef, cmax, pl, lgc = sys.argv[3:]
|
||||
top = max(x["mhs"] for x in steps); ok = [x for x in steps if x["mhs"] >= top * 0.99]
|
||||
chosen = max(ok, key=lambda x: x["eff"]); base = steps[0]
|
||||
rec = {"ts": datetime.datetime.utcnow().strftime("%Y-%m-%dT%H:%M:%SZ"), "machine": "fleet-" + label, "app": "fleet-ladder", "os": "linux", "card": name, "vendor": "nvidia",
|
||||
"driver": drv, "driver_major": drv.split(".")[0], "class": "v3", "key": f"nvidia|{name}|{drv.split('.')[0]}|v3",
|
||||
"plan": "full" if pl == "1" else "baseline", "steps": steps, "chosen": chosen, "before": base, "eff": chosen["eff"], "mhs": chosen["mhs"], "watts": chosen["watts"],
|
||||
"power_default_w": float(pdef or 0), "clock_max_mhz": float(cmax or 0), "knobs": {"power": pl == "1", "clock": lgc == "1"}}
|
||||
json.dump(rec, open(sys.argv[2], "w"), indent=1)
|
||||
print("TUNE " + json.dumps(rec))
|
||||
print(f"RESULT choice clock_cap_mhz={chosen['clock_mhz']} power_pct={chosen['power_pct']} limit_w={chosen['limit_w']} watts={chosen['watts']} mhs={chosen['mhs']} mhw={chosen['eff']} baseline_mhs={base['mhs']} baseline_watts={base['watts']} baseline_mhw={base['eff']} steps={len(steps)}")
|
||||
PY
|
||||
echo "RESULT ember_done $(stamp)"
|
||||
124
tools/fleet/box-matrix.sh
Executable file
124
tools/fleet/box-matrix.sh
Executable file
|
|
@ -0,0 +1,124 @@
|
|||
#!/usr/bin/env bash
|
||||
# Phase 1 on one rented card: the memory matrix of docs/analysis/prover-floor.md on the real card. Every point is one
|
||||
# igneum-prove-host run against a fresh sp1-gpu-server (killed and its socket unlinked around every point), with an
|
||||
# nvidia-smi sampler at 1 s (memory.used, power.draw, utilization); peak = max memory.used, base = the reading just
|
||||
# before the point (the card's idle, or the miner's resident set when the miner runs beside the prover), own = peak
|
||||
# minus base. Every proof is verified by the host's own (unpatched) verifier: the VERIFIED word in the RESULT line.
|
||||
# Writes /root/fleet/out/matrix.json (card facts, the miner row, every point) and the raw logs beside it; the last log
|
||||
# line is "RESULT matrix_done" or "RESULT matrix_failed <why>".
|
||||
set -uo pipefail
|
||||
F=/root/fleet; OUT=$F/out; LOG=$OUT/matrix.log; B=/opt/igneum/pkg/bin; FLOOR=/opt/igneum-floor
|
||||
HOST=$FLOOR/bin/igneum-prove-host; FIX=$FLOOR/prove/proving/fixtures
|
||||
V1=$FIX/fees-v1-shards2.json; EMPTY=$FIX/block-72854-empty-block-first.json
|
||||
LABEL="${LABEL:-box}"; WALLET="${WALLET:-0x1919191919191919191919191919191919191919}"
|
||||
MINER_SECS="${MINER_SECS:-150}"; POINT_TIMEOUT="${POINT_TIMEOUT:-1500}"
|
||||
mkdir -p $OUT $F/mine/packs
|
||||
exec > >(tee -a $LOG) 2>&1
|
||||
stamp() { date -u +%Y-%m-%dT%H:%M:%SZ; }
|
||||
say() { echo "$(stamp) $*"; }
|
||||
q() { nvidia-smi --query-gpu="$1" --format=csv,noheader,nounits -i 0 2>/dev/null | head -1 | tr -d ' '; }
|
||||
[ -x "$HOST" ] || { echo "RESULT matrix_failed no host"; exit 2; }
|
||||
[ -x "$FLOOR/bin/sp1-gpu-server" ] || { echo "RESULT matrix_failed no patched server"; exit 2; }
|
||||
NAME="$(q name)"; TOTAL="$(q memory.total)"; DRV="$(q driver_version)"; PLIM="$(q power.limit)"
|
||||
echo "RESULT start $(stamp) label=$LABEL card=$NAME total_mib=$TOTAL driver=$DRV power_limit_w=$PLIM"
|
||||
ROWS=$OUT/rows.jsonl; : > $ROWS
|
||||
|
||||
# the sampler: one csv per point, "ts,mem_mib,power_w,util_pct"
|
||||
SAMP=""
|
||||
sampler_start() { nvidia-smi --query-gpu=timestamp,memory.used,power.draw,utilization.gpu --format=csv,noheader,nounits -i 0 -l 1 > "$1" 2>/dev/null & SAMP=$!; }
|
||||
sampler_stop() { [ -n "$SAMP" ] && kill $SAMP 2>/dev/null; wait $SAMP 2>/dev/null; SAMP=""; }
|
||||
peak_of() { awk -F', *' 'NR>0 {if ($2+0 > m) m=$2+0} END {print m+0}' "$1"; }
|
||||
mean_col() { awk -F', *' -v c="$2" '{s+=$c; n++} END {if (n) printf "%.1f", s/n; else print 0}' "$1"; }
|
||||
kill_server() { pkill -f sp1-gpu-server 2>/dev/null; sleep 2; pkill -9 -f sp1-gpu-server 2>/dev/null; rm -f /tmp/sp1-cuda-*.sock; }
|
||||
idle_mib() { sleep 3; q memory.used; }
|
||||
|
||||
# 1. idle
|
||||
kill_server
|
||||
sleep 5
|
||||
IDLE_MIB="$(q memory.used)"; IDLE_W="$(q power.draw)"
|
||||
echo "RESULT idle mem_mib=$IDLE_MIB power_w=$IDLE_W"
|
||||
|
||||
# 2. the node's state (the miner mines IBD templates when not synced; the rate is the same, the pack may move)
|
||||
NODE_LINE="$($B/igneum-miner watch 1 grpc://127.0.0.1:26610 2>/dev/null | grep -o 'blocks=[0-9]*.*synced=[a-z]*' | tail -1)"
|
||||
waited=0; while [[ "$NODE_LINE" != *synced=true* && $waited -lt 900 ]]; do sleep 30; waited=$((waited+30)); NODE_LINE="$($B/igneum-miner watch 1 grpc://127.0.0.1:26610 2>/dev/null | grep -o 'blocks=[0-9]*.*synced=[a-z]*' | tail -1)"; done
|
||||
echo "RESULT node $NODE_LINE waited_s=$waited"
|
||||
|
||||
# 3. the miner alone: rate, watts, working set
|
||||
MPID=""
|
||||
miner_start() {
|
||||
cd $F/mine
|
||||
rm -rf packs/devnet; $B/igneum-miner export-pack grpc://127.0.0.1:26610 packs/devnet > $OUT/export-pack.log 2>&1 || say "export-pack failed (see export-pack.log)"
|
||||
nohup $B/igneum-miner mine grpc://127.0.0.1:26610 1 100000000 "$LABEL" --worker $B/igneum-worker-cuda --worker-args "--device 0 --pack packs/devnet" \
|
||||
--prepare-packs packs/prepare --exit-on-seed-change --evm-address "$WALLET" --payout-label "$LABEL" --status-secs 10 > $OUT/miner.log 2>&1 &
|
||||
MPID=$!; cd $F
|
||||
}
|
||||
miner_stop() { [ -n "$MPID" ] && { kill $MPID 2>/dev/null; sleep 2; kill -9 $MPID 2>/dev/null; }; pkill -f igneum-worker-cuda 2>/dev/null; pkill -f "igneum-miner mine" 2>/dev/null; MPID=""; sleep 3; }
|
||||
miner_rate() { # mean of the STATUS now= values in the last N lines
|
||||
grep -o 'now=[0-9.]*' $OUT/miner.log | tail -n "${1:-10}" | cut -d= -f2 | awk '{s+=$1; n++} END {if (n) printf "%.2f", s/n; else print 0}'
|
||||
}
|
||||
miner_start
|
||||
say "miner warming 60 s"; sleep 60
|
||||
sampler_start $OUT/samp-miner.csv; sleep "$MINER_SECS"; sampler_stop
|
||||
MINER_MHS="$(miner_rate 12)"; MINER_W="$(mean_col $OUT/samp-miner.csv 3)"; MINER_UTIL="$(mean_col $OUT/samp-miner.csv 4)"; MINER_MIB="$(peak_of $OUT/samp-miner.csv)"
|
||||
MINER_ACC="$(grep -o 'accepted=[0-9]*' $OUT/miner.log | tail -1 | cut -d= -f2)"
|
||||
MINER_STATUS="$(grep STATUS $OUT/miner.log | tail -1 | cut -c1-200)"
|
||||
echo "RESULT miner mhs=$MINER_MHS watts=$MINER_W util=$MINER_UTIL mem_mib=$MINER_MIB own_mib=$((MINER_MIB - IDLE_MIB)) accepted=${MINER_ACC:-?} status=\"$MINER_STATUS\""
|
||||
printf '%s\n' "{\"row\":\"miner\",\"mhs\":$MINER_MHS,\"watts\":$MINER_W,\"util\":$MINER_UTIL,\"mem_mib\":$MINER_MIB,\"own_mib\":$((MINER_MIB - IDLE_MIB)),\"accepted\":\"${MINER_ACC:-}\"}" >> $ROWS
|
||||
miner_stop
|
||||
|
||||
# 4. one point: name, home (stock or patched), mode, threshold env, fixture, beside
|
||||
point() {
|
||||
local name="$1" home="$2" mode="$3" thr="$4" fixture="$5" beside="$6"
|
||||
kill_server; sleep 2
|
||||
local base; base="$(q memory.used)"
|
||||
local env=(HOME="$home" SP1_PROVER=cuda RUST_LOG=off SP1_GPU_FLOOR_LOG=1)
|
||||
[ -n "$thr" ] && env+=(SP1_GPU_ELEMENT_THRESHOLD="$thr")
|
||||
local log=$OUT/point-$name.log res=$OUT/point-$name.json samp=$OUT/samp-$name.csv
|
||||
say "point $name: mode $mode threshold ${thr:-default} fixture $(basename $fixture) beside_miner=$beside base_mib=$base"
|
||||
sampler_start "$samp"
|
||||
local t0=$(date +%s)
|
||||
env "${env[@]}" timeout "$POINT_TIMEOUT" "$HOST" "$fixture" --mode "$mode" --shard 0 --prover "$WALLET" --out "$res" > "$log" 2>&1
|
||||
local rc=$?; local wall=$(( $(date +%s) - t0 ))
|
||||
sampler_stop; kill_server
|
||||
local peak; peak="$(peak_of "$samp")"
|
||||
local rline; rline="$(grep -E "^RESULT (core|compressed) shard" "$log" | tail -1)"
|
||||
local secs; secs="$(printf '%s' "$rline" | grep -o 'prove [0-9.]* s' | grep -o '[0-9.]*' | head -1)"
|
||||
local bytes; bytes="$(printf '%s' "$rline" | grep -o 'proof [0-9]* bytes' | grep -o '[0-9]*' | head -1)"
|
||||
local ver="no"; printf '%s' "$rline" | grep -q 'VERIFIED' && ver="yes"; printf '%s' "$rline" | grep -q 'VERIFY FAILED' && ver="FAILED"
|
||||
local err; err="$(grep -m1 -E 'Unsupported GPU memory|out of memory|OutOfMemory|CUDA_ERROR|panicked|Error:|error:' "$log" | cut -c1-200 | tr '"' "'")"
|
||||
local setup; setup="$(grep -o 'RESULT setup: [0-9.]* s' "$log" | grep -o '[0-9.]* s' | head -1)"
|
||||
echo "RESULT point name=$name mode=$mode thr=${thr:-default} beside=$beside rc=$rc wall_s=$wall prove_s=${secs:-} proof_bytes=${bytes:-} verified=$ver peak_mib=$peak base_mib=$base own_mib=$((peak - base)) setup=\"${setup:-}\" err=\"${err:-}\""
|
||||
printf '%s\n' "{\"row\":\"point\",\"name\":\"$name\",\"home\":\"$home\",\"mode\":\"$mode\",\"threshold\":\"${thr:-default}\",\"fixture\":\"$(basename $fixture)\",\"beside_miner\":$beside,\"rc\":$rc,\"wall_s\":$wall,\"prove_s\":${secs:-null},\"proof_bytes\":${bytes:-null},\"verified\":\"$ver\",\"peak_mib\":$peak,\"base_mib\":$base,\"own_mib\":$((peak - base)),\"err\":\"${err:-}\",\"result\":\"$(printf '%s' "$rline" | cut -c1-160 | tr '"' "'")\"}" >> $ROWS
|
||||
}
|
||||
STOCK=/root; PATCHED=$FLOOR/home
|
||||
T25=33554432; T26=67108864; T27=134217728; T24=16777216
|
||||
SMALL=0; [ "$TOTAL" -lt 11000 ] && SMALL=1
|
||||
|
||||
# 5. the stock server, alone (the SDK downloads it into /root/.sp1/bin on the first run)
|
||||
point stock-comp-v1 $STOCK compressed "" $V1 false
|
||||
# 6. the patched server, alone
|
||||
point alone-comp-26-v1 $PATCHED compressed $T26 $V1 false
|
||||
point alone-comp-27-v1 $PATCHED compressed $T27 $V1 false
|
||||
point alone-comp-26-empty $PATCHED compressed $T26 $EMPTY false
|
||||
point alone-core-25-v1 $PATCHED core $T25 $V1 false
|
||||
point alone-core-26-v1 $PATCHED core $T26 $V1 false
|
||||
if [ $SMALL = 1 ]; then point alone-comp-25-v1 $PATCHED compressed $T25 $V1 false; point alone-core-24-v1 $PATCHED core $T24 $V1 false; fi
|
||||
# 7. beside the miner
|
||||
miner_start; say "miner warming 45 s for the beside rows"; sleep 45
|
||||
MINER_RES="$(q memory.used)"; echo "RESULT miner_resident mem_mib=$MINER_RES mhs=$(miner_rate 4)"
|
||||
point miner-comp-26-v1 $PATCHED compressed $T26 $V1 true
|
||||
[ "$TOTAL" -ge 15000 ] && point miner-comp-27-v1 $PATCHED compressed $T27 $V1 true
|
||||
point miner-core-25-v1 $PATCHED core $T25 $V1 true
|
||||
point miner-core-26-v1 $PATCHED core $T26 $V1 true
|
||||
if [ $SMALL = 1 ]; then point miner-core-24-v1 $PATCHED core $T24 $V1 true; fi
|
||||
MINER_MHS_BESIDE="$(miner_rate 30)"
|
||||
miner_stop
|
||||
echo "RESULT miner_beside mhs_mean_during_points=$MINER_MHS_BESIDE"
|
||||
|
||||
python3 - "$OUT" "$LABEL" "$NAME" "$TOTAL" "$DRV" "$IDLE_MIB" "$IDLE_W" "$PLIM" <<'PY'
|
||||
import json, sys
|
||||
out, label, name, total, drv, idle, idle_w, plim = sys.argv[1:]
|
||||
rows = [json.loads(l) for l in open(f"{out}/rows.jsonl") if l.strip()]
|
||||
json.dump({"label": label, "card": name, "total_mib": int(total), "driver": drv, "idle_mib": int(idle), "idle_w": float(idle_w or 0), "power_limit_w": float(plim or 0), "rows": rows}, open(f"{out}/matrix.json", "w"), indent=1)
|
||||
PY
|
||||
echo "RESULT matrix_done $(stamp) points=$(grep -c '"row":"point"' $ROWS)"
|
||||
246
tools/fleet/box-prover.py
Normal file
246
tools/fleet/box-prover.py
Normal file
|
|
@ -0,0 +1,246 @@
|
|||
#!/usr/bin/env python3
|
||||
"""The fleet's segment-aligned prover for a Linux box (phase 2): a port of app/igneum-app/src/prover.rs (branch
|
||||
proving-v1, 272b025) and tools/proving-v1/pc2-segments.ps1 around the four binaries. One pass every 15 s:
|
||||
|
||||
1. igneum_getProvingStatus (start, n, unproven, tip) and igneum_getAssignedShards [[keyHash], 600], grouped into
|
||||
whole untouched segments (every block present, every shard open, unpaid, not in this pool) whose deadline
|
||||
(last DAA + unproven) is at least 240 DAA (or 1.5x the last segment's time) past the tip; ranked by FNV-1a of
|
||||
(first, key hash) so a fleet of provers spreads over the segments instead of racing for one.
|
||||
2. igneum_getSegmentStatement [first]: executed and pending; chained with --prev when the previous segment's proof is
|
||||
in this pool (igneum_getSegmentProofBytes), else fresh. A fresh record refused with "pending until" is held and
|
||||
offered again every pass until the segment's deadline (the 272b025 behaviour).
|
||||
3. one export (igneum_exportSegments 0..last), one fixture per block (igneum-prove-export), one host run
|
||||
(igneum-prove-host --mode chain --chain ... --save-shards [--prev]) on the patched server (HOME=/opt/igneum-floor/home,
|
||||
SP1_GPU_ELEMENT_THRESHOLD from THRESHOLD), the miner paused for the run when MINER=pause (prove-alone cards).
|
||||
4. every shard record signed (igneum-miner sign-record) and submitted (igneum_submitProofRecord); the segment record
|
||||
(sign-segment-record, igneum_submitSegmentRecord) once every shard is accepted and the statement equals the node's.
|
||||
5. the paid state of every submitted segment polled each pass (igneum_getSegmentRecords); a state file for the
|
||||
collector: /root/fleet/out/prover-state.json; every event a RESULT line in /root/fleet/out/prover.log.
|
||||
|
||||
Env: LABEL (the key label, kept for the box's life), WALLET (payout), THRESHOLD (element threshold or empty),
|
||||
MINER (keep|pause), RUN_HOURS (default 9).
|
||||
"""
|
||||
import json, os, sys, time, subprocess, datetime, binascii, urllib.request, signal
|
||||
F = "/root/fleet"; OUT = f"{F}/out"; B = "/opt/igneum/pkg/bin"; FLOOR = "/opt/igneum-floor"
|
||||
HOST = f"{FLOOR}/bin/igneum-prove-host"; EXPORT = f"{FLOOR}/bin/igneum-prove-export"
|
||||
EVM = "http://127.0.0.1:26790"; GRPC = "grpc://127.0.0.1:26610"; CHAIN = "igneum-devnet"
|
||||
LABEL = os.environ.get("LABEL", "box"); WALLET = os.environ.get("WALLET", "0x" + "19" * 20)
|
||||
THRESHOLD = os.environ.get("THRESHOLD", ""); MINER = os.environ.get("MINER", "keep"); RUN_HOURS = float(os.environ.get("RUN_HOURS", "9"))
|
||||
os.makedirs(f"{OUT}/segs", exist_ok=True); os.makedirs(f"{F}/mine/packs", exist_ok=True)
|
||||
LOG = open(f"{OUT}/prover.log", "a")
|
||||
def stamp(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
def say(s): LOG.write(f"{s}\n"); LOG.flush(); print(s, flush=True)
|
||||
def hexi(v): return int(v, 16) if isinstance(v, str) and v.startswith("0x") else int(v or 0)
|
||||
def rpc(method, params, timeout=60):
|
||||
body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": method, "params": params}).encode()
|
||||
try:
|
||||
with urllib.request.urlopen(urllib.request.Request(EVM, data=body, headers={"Content-Type": "application/json"}), timeout=timeout) as r:
|
||||
d = json.loads(r.read())
|
||||
if d.get("error"): say(f"RESULT rpc_error {stamp()} {method}: {str(d['error'])[:200]}"); return None
|
||||
return d.get("result")
|
||||
except Exception as e:
|
||||
say(f"RESULT rpc_fail {stamp()} {method}: {str(e)[:120]}"); return None
|
||||
def fnv1a(s):
|
||||
h = 0xcbf29ce484222325
|
||||
for c in s.encode(): h ^= c; h = (h * 0x100000001b3) & 0xffffffffffffffff
|
||||
return h
|
||||
def kill_server(): subprocess.run("pkill -f sp1-gpu-server; sleep 1; rm -f /tmp/sp1-cuda-*.sock", shell=True)
|
||||
MPROC = None
|
||||
def miner_start():
|
||||
global MPROC
|
||||
if MPROC and MPROC.poll() is None: return
|
||||
subprocess.run(f"cd {F}/mine && rm -rf packs/devnet && {B}/igneum-miner export-pack {GRPC} packs/devnet > {OUT}/prover-export-pack.log 2>&1", shell=True)
|
||||
MPROC = subprocess.Popen([f"{B}/igneum-miner", "mine", GRPC, "1", "100000000", LABEL, "--worker", f"{B}/igneum-worker-cuda", "--worker-args", "--device 0 --pack packs/devnet",
|
||||
"--prepare-packs", "packs/prepare", "--exit-on-seed-change", "--evm-address", WALLET, "--payout-label", LABEL, "--status-secs", "30"],
|
||||
cwd=f"{F}/mine", stdout=open(f"{OUT}/prover-miner.log", "a"), stderr=subprocess.STDOUT)
|
||||
say(f"RESULT miner_start {stamp()} pid={MPROC.pid}")
|
||||
def miner_stop():
|
||||
global MPROC
|
||||
if MPROC: MPROC.terminate(); time.sleep(2); MPROC.kill(); MPROC = None
|
||||
subprocess.run("pkill -f igneum-worker-cuda", shell=True)
|
||||
def miner_rate(n=6):
|
||||
try:
|
||||
vals = [float(l.split("now=")[1].split()[0]) for l in open(f"{OUT}/prover-miner.log").read().split("\n") if "now=" in l][-n:]
|
||||
return round(sum(vals) / len(vals), 2) if vals else 0
|
||||
except Exception: return 0
|
||||
# the key
|
||||
kh = subprocess.run([f"{B}/igneum-miner", "key-hash", LABEL], capture_output=True, text=True).stdout.strip().split("\n")[-1].strip()
|
||||
if len(kh) == 64: kh = "0x" + kh
|
||||
if len(kh) != 66: say(f"RESULT prover_failed key-hash gave '{kh[:40]}'"); sys.exit(2)
|
||||
say(f"RESULT start {stamp()} label={LABEL} key={kh[:18]} wallet={WALLET[:10]} threshold={THRESHOLD or 'default'} miner={MINER} host={os.path.exists(HOST)}")
|
||||
# wait for the node to be synced
|
||||
for _ in range(120):
|
||||
w = subprocess.run(f"{B}/igneum-miner watch 1 {GRPC} 2>/dev/null | grep -o 'synced=[a-z]*' | tail -1", shell=True, capture_output=True, text=True).stdout.strip()
|
||||
if w == "synced=true": break
|
||||
time.sleep(15)
|
||||
say(f"RESULT node {stamp()} {w}")
|
||||
miner_start()
|
||||
state = {"passes": 0, "claimed": 0, "submitted": 0, "paid": 0, "paid_wei": 0, "shards_accepted": 0, "shards_refused": 0, "segment_refused": 0, "held": 0,
|
||||
"last_segment_s": 0, "segments": [], "started": stamp(), "label": LABEL, "wallet": WALLET, "key": kh}
|
||||
attempted = set(); submitted = {}; held = {} # held: first -> {body file, deadline, last}
|
||||
last_seg_secs = 0; t_run0 = time.time()
|
||||
def save_state():
|
||||
state["miner_mhs"] = miner_rate(); state["updated"] = stamp(); state["run_min"] = round((time.time() - t_run0) / 60, 1)
|
||||
json.dump(state, open(f"{OUT}/prover-state.json.tmp", "w"), indent=1); os.replace(f"{OUT}/prover-state.json.tmp", f"{OUT}/prover-state.json")
|
||||
def submit(method, record, proof_file):
|
||||
proof = "0x" + binascii.hexlify(open(proof_file, "rb").read()).decode()
|
||||
return rpc(method, [{"record": record, "proof": proof}], timeout=180)
|
||||
def candidates(st):
|
||||
start = hexi(st["v1"]["start"]); n = max(1, hexi(st["v1"]["segmentBlocks"])); unproven = hexi(st["v1"]["unprovenDaa"]); tip = hexi(st["tipDaa"])
|
||||
work = rpc("igneum_getAssignedShards", [[kh], 600]) or []
|
||||
by = {}
|
||||
for w in work:
|
||||
num = hexi(w.get("number"))
|
||||
if num < start: continue
|
||||
by.setdefault(num, []).append(w)
|
||||
need = max(240, int(last_seg_secs * 1.5) + 1)
|
||||
segs = []; seen = set()
|
||||
for num in sorted(by):
|
||||
k = (num - start) // n; first = start + k * n; last = first + n - 1
|
||||
if first in seen or first in attempted: continue
|
||||
seen.add(first)
|
||||
whole = True; shards = []; last_daa = 0
|
||||
for b in range(first, last + 1):
|
||||
if b not in by: whole = False; break
|
||||
es = {int(e.get("shard", 0)): e for e in by[b]}
|
||||
for si in sorted(es):
|
||||
e = es[si]
|
||||
if not e.get("open") or e.get("paid") is not None or (e.get("pool") and len(e["pool"]) > 0): whole = False; break
|
||||
shards.append({"number": b, "hash": e.get("hash"), "shard": si})
|
||||
if b == last: last_daa = hexi(e.get("daaScore"))
|
||||
if not whole: break
|
||||
if whole and shards:
|
||||
deadline = last_daa + unproven
|
||||
if deadline >= tip + 1 + need: segs.append({"first": first, "last": last, "last_daa": last_daa, "deadline": deadline, "shards": shards, "margin": deadline - tip - 1})
|
||||
segs.sort(key=lambda s: fnv1a(f"{s['first']}:{kh}"))
|
||||
return segs, start, n, tip, len(work)
|
||||
while (time.time() - t_run0) / 3600 < RUN_HOURS:
|
||||
state["passes"] += 1; p = state["passes"]
|
||||
if MPROC and MPROC.poll() is not None and MINER == "keep": say(f"RESULT miner_exit {stamp()} rc={MPROC.returncode}; restarting"); MPROC = None; miner_start()
|
||||
st = rpc("igneum_getProvingStatus", [])
|
||||
if not st or not st.get("v1") or not st["v1"].get("active"): say(f"RESULT pass {p} {stamp()} v1 not active or no status"); time.sleep(20); continue
|
||||
tip = hexi(st["tipDaa"])
|
||||
# paid state
|
||||
for first in list(submitted):
|
||||
rec = rpc("igneum_getSegmentRecords", [hex(first)])
|
||||
if rec and rec.get("paid"):
|
||||
wei = hexi(rec["paid"]["wei"]); state["paid"] += 1; state["paid_wei"] += wei
|
||||
say(f"RESULT paid {stamp()} segment {first}..{submitted[first]['last']} wei={wei} ({wei/1e18:.4f} IGN) carrier={hexi(rec['paid'].get('carrierNumber'))} after {int(time.time()-submitted[first]['at'])} s")
|
||||
for s in state["segments"]:
|
||||
if s["first"] == first: s["paid_wei"] = wei; s["paid_at"] = stamp()
|
||||
del submitted[first]
|
||||
elif rec is not None and tip > submitted[first]["deadline"] + 50:
|
||||
say(f"RESULT unpaid {stamp()} segment {first} past its deadline unpaid; carried={len(rec.get('carried') or [])} pool={len(rec.get('pool') or [])}"); del submitted[first]
|
||||
# held fresh records offered again
|
||||
for first in list(held):
|
||||
h = held[first]
|
||||
if tip > h["deadline"]: say(f"RESULT held_expired {stamp()} segment {first} deadline passed"); del held[first]; continue
|
||||
rr = submit("igneum_submitSegmentRecord", h["record"], h["proof_file"])
|
||||
if rr and rr.get("accepted"):
|
||||
state["submitted"] += 1; submitted[first] = {"last": h["last"], "at": time.time(), "deadline": h["deadline"]}; say(f"RESULT submitted {stamp()} segment {first}..{h['last']} record accepted on retry (held {int(time.time()-h['since'])} s)"); del held[first]
|
||||
state["held"] = len(held)
|
||||
cands, start, n, tip, entries = candidates(st)
|
||||
save_state()
|
||||
if not cands:
|
||||
if p % 4 == 1: say(f"RESULT pass {p} {stamp()} no whole segment inside the margin (worklist {entries} entries, tip {tip}, mhs {miner_rate()}); waiting")
|
||||
time.sleep(15); continue
|
||||
picked = None; prev_file = None; expected = ""
|
||||
for c in cands[:3]:
|
||||
stmt = rpc("igneum_getSegmentStatement", [hex(c["first"])])
|
||||
if not stmt or not stmt.get("executed") or (stmt.get("status") or {}).get("status") != "pending":
|
||||
attempted.add(c["first"]); say(f"RESULT skip {stamp()} segment {c['first']}: executed={stmt and stmt.get('executed')} status={(stmt or {}).get('status')}"); continue
|
||||
if stmt.get("previous") is None:
|
||||
if not st["v1"].get("freshRuleActive") and c["first"] >= start + n:
|
||||
pr = rpc("igneum_getSegmentRecords", [hex(c["first"] - n)])
|
||||
waiting = any(e.get("verified") and e.get("includedIn") is None for e in (pr or {}).get("pool") or [])
|
||||
if waiting or (pr and pr.get("paid")): say(f"RESULT skip {stamp()} segment {c['first']}: previous has a record waiting or paid, a fresh chain would be refused (fresh rule off)"); continue
|
||||
picked = c; expected = stmt.get("publicValuesFresh") or ""; break
|
||||
if not stmt["previous"].get("proofInPool"): say(f"RESULT skip {stamp()} segment {c['first']}: previous paid, its proof not in this pool"); continue
|
||||
got = rpc("igneum_getSegmentProofBytes", [stmt["previous"]["first"], stmt["previous"]["keyHash"]])
|
||||
if not got or not got.get("proof"): continue
|
||||
prev_file = f"{OUT}/segs/prev-{c['first']}.bin"; h = got["proof"]; open(prev_file, "wb").write(binascii.unhexlify(h[2:] if h.startswith("0x") else h))
|
||||
picked = c; expected = stmt.get("publicValuesContinuing") or ""; break
|
||||
if not picked: say(f"RESULT pass {p} {stamp()} {len(cands)} candidates, none usable; waiting"); time.sleep(15); continue
|
||||
first, last = picked["first"], picked["last"]; attempted.add(first); state["claimed"] += 1
|
||||
say(f"RESULT claim {stamp()} segment {first}..{last} ({len(picked['shards'])} shards, {'continuing' if prev_file else 'fresh'}) margin={picked['margin']} tip={tip} candidates={len(cands)} rank_by=fnv")
|
||||
seg = {"first": first, "last": last, "claimed_at": stamp(), "shards": len(picked["shards"]), "fresh": prev_file is None}; state["segments"].append(seg)
|
||||
d = f"{OUT}/segs/seg-{first}"; os.makedirs(d, exist_ok=True); t_seg0 = time.time()
|
||||
# export
|
||||
t = time.time(); body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": "igneum_exportSegments", "params": ["0x0", hex(last)]})
|
||||
r = subprocess.run(["curl", "-s", "-m", "600", "-X", "POST", EVM, "-H", "Content-Type: application/json", "--data-binary", body, "-o", f"{d}/seq.json"])
|
||||
try: json.dump(json.load(open(f"{d}/seq.json"))["result"], open(f"{d}/export.json", "w"))
|
||||
except Exception as e: say(f"RESULT seg {first} export FAILED {str(e)[:100]}"); continue
|
||||
seg["export_s"] = round(time.time() - t, 1); os.remove(f"{d}/seq.json")
|
||||
# cut
|
||||
t = time.time(); fixtures = []
|
||||
ok = True
|
||||
for b in range(first, last + 1):
|
||||
rr = subprocess.run([EXPORT, f"{d}/export.json", str(b), f"{d}/block-{b}.json", "--source", f"fleet {LABEL} live devnet, segment-aligned prover"], capture_output=True, text=True, timeout=600)
|
||||
if rr.returncode != 0: say(f"RESULT seg {first} cut {b} FAILED: {(rr.stdout + rr.stderr)[-200:]}"); ok = False; break
|
||||
fixtures.append(f"{d}/block-{b}.json")
|
||||
if not ok: continue
|
||||
seg["cut_s"] = round(time.time() - t, 1)
|
||||
# chain
|
||||
if MINER == "pause": miner_stop()
|
||||
kill_server()
|
||||
env = dict(os.environ, HOME=f"{FLOOR}/home", SP1_PROVER="cuda", RUST_LOG="off")
|
||||
if THRESHOLD: env["SP1_GPU_ELEMENT_THRESHOLD"] = THRESHOLD
|
||||
args = [HOST, "--mode", "chain", "--chain", ",".join(fixtures), "--prover", WALLET, "--save-shards", "--out", f"{d}/chain-results.json"]
|
||||
if prev_file: args += ["--prev", prev_file]
|
||||
samp = subprocess.Popen("nvidia-smi --query-gpu=memory.used,utilization.gpu,power.draw --format=csv,noheader,nounits -l 1", shell=True, stdout=open(f"{d}/smi.csv", "w"), stderr=subprocess.DEVNULL)
|
||||
t = time.time()
|
||||
try: rr = subprocess.run(args, env=env, capture_output=True, text=True, timeout=3600)
|
||||
except subprocess.TimeoutExpired: rr = None
|
||||
seg["chain_s"] = round(time.time() - t, 1); samp.terminate(); kill_server()
|
||||
if MINER == "pause": miner_start()
|
||||
try: peak = max(float(l.split(",")[0]) for l in open(f"{d}/smi.csv") if l.strip())
|
||||
except Exception: peak = 0
|
||||
seg["peak_mib"] = peak
|
||||
open(f"{d}/chain.log", "w").write((rr.stdout if rr else "") + "\n" + (rr.stderr if rr else "TIMEOUT"))
|
||||
if not rr or rr.returncode != 0 or not os.path.exists(f"{d}/chain-results.json"):
|
||||
say(f"RESULT seg {first} chain FAILED {stamp()} rc={rr.returncode if rr else 'timeout'} wall={seg['chain_s']} s: {((rr.stderr if rr else '') or '')[-200:].strip()}"); seg["failed"] = "chain"; continue
|
||||
res = json.load(open(f"{d}/chain-results.json"))
|
||||
recs = [s for blk in res.get("blocks", []) for s in blk.get("shard_records", [])]
|
||||
say(f"RESULT seg {first} chain {stamp()} {len(recs)} shard records, chain_len {res.get('segment_chain_len')}, proof {res.get('segment_proof_bytes')} bytes, shards {res.get('shard_prove_seconds_total', 0):.1f} s, aggregation {res.get('aggregate_prove_seconds_total', 0):.1f} s, wall {seg['chain_s']} s, peak {peak:.0f} MiB")
|
||||
# shard records
|
||||
ok_shards = 0
|
||||
for rcd in recs:
|
||||
sg = subprocess.run([f"{B}/igneum-miner", "sign-record", LABEL, CHAIN, rcd["block_hash"], str(rcd["number"]), str(rcd["shard"]), WALLET, rcd["statement"], rcd["proof_sha256"]], capture_output=True, text=True).stdout.strip().split("\n")[-1]
|
||||
try: record = json.loads(sg).get("record")
|
||||
except Exception: record = None
|
||||
if not record: say(f"RESULT seg {first} shard {rcd['number']}/{rcd['shard']} sign FAILED: {sg[:120]}"); state["shards_refused"] += 1; continue
|
||||
reply = submit("igneum_submitProofRecord", record, rcd["proof_file"])
|
||||
if reply and reply.get("accepted"): ok_shards += 1; state["shards_accepted"] += 1
|
||||
else: state["shards_refused"] += 1; say(f"RESULT seg {first} shard {rcd['number']}/{rcd['shard']} refused: {(reply or {}).get('reason', reply)}")
|
||||
seg["shards_accepted"] = ok_shards
|
||||
say(f"RESULT seg {first} shards {stamp()} accepted {ok_shards} of {len(recs)}")
|
||||
if ok_shards != len(recs): seg["failed"] = "shards"; continue
|
||||
pv = res.get("segment_public_values", "")
|
||||
strip = lambda h: (h[2:] if h.startswith("0x") else h); strip2 = lambda h: (strip(h)[:472] + strip(h)[536:]) if len(strip(h)) == 680 else strip(h)
|
||||
if strip2(pv) != strip2(expected): say(f"RESULT seg {first} FAILED: statement differs from the node's; ours {pv[:34]} node {expected[:34]}"); seg["failed"] = "statement"; continue
|
||||
last_hash = next(s["hash"] for s in picked["shards"] if s["number"] == last)
|
||||
sg = subprocess.run([f"{B}/igneum-miner", "sign-segment-record", LABEL, CHAIN, str(first), str(last), last_hash, WALLET, pv, res["segment_proof_sha256"]], capture_output=True, text=True).stdout.strip().split("\n")[-1]
|
||||
try: record = json.loads(sg).get("record")
|
||||
except Exception: record = None
|
||||
if not record: say(f"RESULT seg {first} segment sign FAILED: {sg[:120]}"); seg["failed"] = "sign"; continue
|
||||
reply = submit("igneum_submitSegmentRecord", record, res["segment_proof_file"])
|
||||
seg_s = round(time.time() - t_seg0, 1); seg["end_to_end_s"] = seg_s
|
||||
if reply and reply.get("accepted"):
|
||||
state["submitted"] += 1; last_seg_secs = seg_s; state["last_segment_s"] = seg_s
|
||||
stmt2 = rpc("igneum_getSegmentStatement", [hex(first)]) or {}
|
||||
submitted[first] = {"last": last, "at": time.time(), "deadline": picked["deadline"]}; seg["submitted_at"] = stamp()
|
||||
say(f"RESULT submitted {stamp()} segment {first}..{last} record accepted (new={reply.get('new')}, chain_len {res.get('segment_chain_len')}), aggregator share {hexi(stmt2.get('aggregatorWei'))/1e18:.4f} IGN, end to end {seg_s} s, mhs {miner_rate()}")
|
||||
else:
|
||||
reason = str((reply or {}).get("reason", reply))[:200]; state["segment_refused"] += 1; seg["refused"] = reason
|
||||
say(f"RESULT segment_refused {stamp()} segment {first}..{last}: {reason}; end to end {seg_s} s")
|
||||
if "pending until" in reason or "does not chain" in reason:
|
||||
held[first] = {"record": record, "proof_file": res["segment_proof_file"], "last": last, "deadline": picked["deadline"], "since": time.time()}; say(f"RESULT held {stamp()} segment {first} held for retry until DAA {picked['deadline']}")
|
||||
for b in range(first, last + 1):
|
||||
try: os.remove(f"{d}/block-{b}.json")
|
||||
except OSError: pass
|
||||
try: os.remove(f"{d}/export.json")
|
||||
except OSError: pass
|
||||
save_state()
|
||||
miner_stop(); kill_server(); save_state()
|
||||
say(f"RESULT summary {stamp()} passes={state['passes']} claimed={state['claimed']} submitted={state['submitted']} paid={state['paid']} paid_wei={state['paid_wei']} shards_accepted={state['shards_accepted']} shards_refused={state['shards_refused']} segment_refused={state['segment_refused']}")
|
||||
say(f"RESULT prover_done {stamp()}")
|
||||
111
tools/fleet/box-setup.sh
Executable file
111
tools/fleet/box-setup.sh
Executable file
|
|
@ -0,0 +1,111 @@
|
|||
#!/usr/bin/env bash
|
||||
# Igneum GPU fleet, on-box setup (6 October 2026). Runs as root inside a rented Vast or RunPod container built from
|
||||
# nvidia/cuda:12.8.1-devel-ubuntu24.04 (or RunPod's CUDA 12.8 image). Idempotent; every step prints a RESULT line to
|
||||
# /root/fleet/setup.log and the last line is "RESULT setup_done" or "RESULT setup_failed <step>".
|
||||
#
|
||||
# What it does, in order: tools (apt), the 0.3.12 Linux package (sha256 checked) unpacked to /opt/igneum/pkg, the node
|
||||
# started at once on the devnet with the ten-field override (so it syncs while the builds run), rustup, Go 1.27.1 (pinned
|
||||
# sha256), SP1 v6.8.1 cloned and patched with proving/prover-floor/sp1-gpu-6.8.1-floor.patch, the patched server built
|
||||
# with CUDA_ARCHS=$ARCHS into /opt/igneum-floor (the SDK's stock server stays under /root/.sp1 for the stock rows), and
|
||||
# igneum-prove-host + igneum-prove-export built with the cuda feature from the prover-floor host bundle.
|
||||
#
|
||||
# Inputs in /root/fleet/in: igneum-prove-wsl2-floor.zip (the host sources, pinned ELFs, fixtures), floor.patch,
|
||||
# override.json. Env: ARCHS (CUDA arch list, default 86,89,120), LABEL (this box's name), WALLET (0x payout, throwaway).
|
||||
set -uo pipefail
|
||||
export DEBIAN_FRONTEND=noninteractive
|
||||
F=/root/fleet; IN=$F/in; OUT=$F/out; LOG=$F/setup.log
|
||||
mkdir -p $F $IN $OUT /opt/igneum /opt/igneum-floor/bin /opt/igneum-floor/home/.sp1/bin /opt/igneum-floor/logs
|
||||
exec > >(tee -a $LOG) 2>&1
|
||||
stamp() { date -u +%Y-%m-%dT%H:%M:%SZ; }
|
||||
say() { echo "$(stamp) $*"; }
|
||||
fail() { echo "RESULT setup_failed $1 $(stamp)"; exit 2; }
|
||||
ARCHS="${ARCHS:-86,89,120}"; LABEL="${LABEL:-box}"; WALLET="${WALLET:-}"
|
||||
PKG_URL=https://dl.igneum.network/dl/public/igneum-hive-0.3.12.tar.gz
|
||||
PKG_SHA=7972af92e7cd9a032303eca4d95b533f53e0e68d1b9cae5bfe406a5b7c30a454
|
||||
PATCH_SHA=e81cb0d03b291f9fd4bf0a109d6da2d7c897795c9ffd7f797c0ddce723eee2b1
|
||||
SEED=188.245.5.161:26611
|
||||
echo "RESULT start $(stamp) label=$LABEL archs=$ARCHS host=$(hostname) nproc=$(nproc) ram_gb=$(( $(awk '/MemTotal/{print $2}' /proc/meminfo) / 1048576 )) disk_avail=$(df -BG /root | awk 'NR==2{print $4}')"
|
||||
echo "RESULT gpu $(nvidia-smi --query-gpu=name,memory.total,driver_version,pci.bus_id,power.limit,power.min_limit,power.max_limit,clocks.max.sm --format=csv,noheader 2>&1 | tr '\n' ';')"
|
||||
echo "RESULT os $(. /etc/os-release; echo "$ID $VERSION_ID") glibc $(ldd --version | head -1 | awk '{print $NF}') nvcc $(nvcc --version 2>/dev/null | grep -o 'release [0-9.]*' || echo none)"
|
||||
|
||||
# 1. tools
|
||||
if ! command -v protoc >/dev/null 2>&1 || ! command -v cmake >/dev/null 2>&1 || ! command -v clang >/dev/null 2>&1; then
|
||||
say "apt"
|
||||
apt-get update -qq >/dev/null 2>&1
|
||||
apt-get install -y -qq build-essential clang cmake pkg-config libssl-dev git curl ca-certificates python3 wget unzip protobuf-compiler jq rsync bc pciutils >/dev/null 2>&1 || fail apt
|
||||
fi
|
||||
echo "RESULT tools protoc=$(protoc --version 2>/dev/null) cmake=$(cmake --version | head -1) clang=$(clang --version | head -1 | cut -c1-40)"
|
||||
|
||||
# 2. the 0.3.12 Linux package and the node
|
||||
if [ ! -x /opt/igneum/pkg/bin/igneumd ]; then
|
||||
say "package"
|
||||
curl -fsSL -o $F/pkg.tgz "$PKG_URL" || fail package_download
|
||||
echo "$PKG_SHA $F/pkg.tgz" | sha256sum -c - >/dev/null || fail package_sha256
|
||||
rm -rf /opt/igneum/pkg.new && mkdir -p /opt/igneum/pkg.new && tar -C /opt/igneum/pkg.new --strip-components=1 -xzf $F/pkg.tgz && mv /opt/igneum/pkg.new /opt/igneum/pkg
|
||||
fi
|
||||
B=/opt/igneum/pkg/bin
|
||||
echo "RESULT package igneumd=$(sha256sum $B/igneumd | cut -c1-16) miner=$(sha256sum $B/igneum-miner | cut -c1-16) cuda_worker=$(sha256sum $B/igneum-worker-cuda | cut -c1-16) version=$($B/igneumd --version 2>&1 | head -1)"
|
||||
cp $IN/override.json $F/override.json
|
||||
if ! pgrep -x igneumd >/dev/null; then
|
||||
say "node"
|
||||
mkdir -p $F/node
|
||||
nohup $B/igneumd --devnet --appdir=$F/node --rpclisten=127.0.0.1:26610 --evm-rpclisten=127.0.0.1:26790 --listen=0.0.0.0:26611 \
|
||||
--addpeer=$SEED --override-params-file=$F/override.json --nodnsseed --disable-upnp --nologfiles --yes > $F/node.log 2>&1 &
|
||||
sleep 8
|
||||
fi
|
||||
echo "RESULT node pid=$(pgrep -x igneumd | head -1) digest=$(grep -o 'digest: [0-9a-f]*' $F/node.log | head -1 | awk '{print $2}') fresh=$(grep -c 'fresh-record rule' $F/node.log)"
|
||||
|
||||
# 3. rust and go
|
||||
export PATH="$HOME/.cargo/bin:/opt/igneum-floor/go/bin:$PATH"
|
||||
if [ ! -x "$HOME/.cargo/bin/cargo" ]; then say "rustup"; curl -sSf https://sh.rustup.rs | sh -s -- -y --profile minimal >/dev/null 2>&1 || fail rustup; fi
|
||||
if [ ! -x /opt/igneum-floor/go/bin/go ]; then
|
||||
say "go"
|
||||
curl -sSL -o /opt/igneum-floor/go.tgz https://go.dev/dl/go1.27.1.linux-amd64.tar.gz || fail go_download
|
||||
echo "63d339f0da5ab53635a56f2490a7984dfe12dfcff22ad749f63edaf590168445 /opt/igneum-floor/go.tgz" | sha256sum -c - >/dev/null || fail go_sha256
|
||||
tar -xzf /opt/igneum-floor/go.tgz -C /opt/igneum-floor && rm -f /opt/igneum-floor/go.tgz
|
||||
fi
|
||||
export GOPATH=/opt/igneum-floor/gopath GOCACHE=/opt/igneum-floor/gocache GOFLAGS=-mod=mod
|
||||
echo "RESULT toolchain cargo=$(cargo --version) go=$(go version | awk '{print $3}')"
|
||||
|
||||
# 4. the patched SP1 GPU server
|
||||
SRC=/opt/igneum-floor/sp1
|
||||
echo "RESULT patch_sha256 $(sha256sum $IN/floor.patch | cut -c1-64) expected $PATCH_SHA"
|
||||
[ "$(sha256sum $IN/floor.patch | cut -c1-64)" = "$PATCH_SHA" ] || fail patch_sha256
|
||||
if [ ! -x /opt/igneum-floor/bin/sp1-gpu-server ]; then
|
||||
if [ ! -d "$SRC/.git" ]; then say "clone sp1"; git clone -q --depth 1 --branch v6.8.1 https://github.com/succinctlabs/sp1 "$SRC" || fail clone; fi
|
||||
cd "$SRC" && git checkout -q -- . && git clean -qfd sp1-gpu/crates >/dev/null 2>&1
|
||||
echo "RESULT source $(git describe --tags --always) $(git rev-parse HEAD)"
|
||||
git apply $IN/floor.patch || fail patch_apply
|
||||
touch sp1-gpu/crates/prover_components/src/builder.rs sp1-gpu/crates/jagged_tracegen/src/lib.rs sp1-gpu/crates/server/src/server.rs sp1-gpu/crates/cuda/src/task.rs
|
||||
echo "RESULT patched $(git diff --stat | tail -1)"
|
||||
export CUDA_ARCHS="$ARCHS" CARGO_TARGET_DIR=/opt/igneum-floor/target CUDA_PATH=/usr/local/cuda CUDACXX=/usr/local/cuda/bin/nvcc
|
||||
J=$(nproc); [ "$J" -gt 16 ] && J=16
|
||||
say "build server jobs=$J archs=$ARCHS"; t0=$(date +%s)
|
||||
cargo build --release --bin sp1-gpu-server -j $J > /opt/igneum-floor/logs/build-server.log 2>&1 &
|
||||
BP=$!; while kill -0 $BP 2>/dev/null; do sleep 60; echo "STAGE server $(stamp) $(( ($(date +%s) - t0) / 60 )) min $(grep -c '^ Compiling' /opt/igneum-floor/logs/build-server.log) crates"; done
|
||||
wait $BP; rc=$?
|
||||
echo "RESULT server_build_exit $rc time_s=$(( $(date +%s) - t0 ))"
|
||||
[ $rc -eq 0 ] || { grep -n -A8 '^error' /opt/igneum-floor/logs/build-server.log | head -60; fail server_build; }
|
||||
cp /opt/igneum-floor/target/release/sp1-gpu-server /opt/igneum-floor/bin/ && cp /opt/igneum-floor/bin/sp1-gpu-server /opt/igneum-floor/home/.sp1/bin/ && chmod +x /opt/igneum-floor/bin/sp1-gpu-server /opt/igneum-floor/home/.sp1/bin/sp1-gpu-server
|
||||
fi
|
||||
echo "RESULT server bytes=$(stat -c %s /opt/igneum-floor/bin/sp1-gpu-server) sha256=$(sha256sum /opt/igneum-floor/bin/sp1-gpu-server | cut -c1-64) version=$(/opt/igneum-floor/bin/sp1-gpu-server --version 2>&1 | head -1) elf=$(cuobjdump --list-elf /opt/igneum-floor/bin/sp1-gpu-server 2>/dev/null | grep -o 'sm_[0-9]*' | sort -u | tr '\n' ' ')"
|
||||
|
||||
# 5. the host and the exporter (cuda feature), from the prover-floor bundle
|
||||
H=/opt/igneum-floor/prove
|
||||
if [ ! -x $H/proving/igneum-prove/target/release/igneum-prove-host ]; then
|
||||
say "host"
|
||||
rm -rf $H && mkdir -p $H && unzip -q -o $IN/igneum-prove-wsl2-floor.zip -d $F/unz && cp -r $F/unz/igneum-prove-wsl2/package/. $H/ && find $H -type f -exec touch {} +
|
||||
cd $H/proving/igneum-prove || fail host_src
|
||||
t0=$(date +%s)
|
||||
cargo build --release -p igneum-prove-export -p igneum-prove-host --features igneum-prove-host/cuda -j $(( $(nproc) > 16 ? 16 : $(nproc) )) > /opt/igneum-floor/logs/build-host.log 2>&1 &
|
||||
BP=$!; while kill -0 $BP 2>/dev/null; do sleep 60; echo "STAGE host $(stamp) $(( ($(date +%s) - t0) / 60 )) min $(grep -c '^ Compiling' /opt/igneum-floor/logs/build-host.log) crates"; done
|
||||
wait $BP; rc=$?
|
||||
echo "RESULT host_build_exit $rc time_s=$(( $(date +%s) - t0 ))"
|
||||
[ $rc -eq 0 ] || { grep -n -A8 '^error' /opt/igneum-floor/logs/build-host.log | head -60; fail host_build; }
|
||||
fi
|
||||
HOST=$H/proving/igneum-prove/target/release/igneum-prove-host; EXPORT=$H/proving/igneum-prove/target/release/igneum-prove-export
|
||||
ln -sfn $HOST /opt/igneum-floor/bin/igneum-prove-host; ln -sfn $EXPORT /opt/igneum-floor/bin/igneum-prove-export
|
||||
echo "RESULT host sha256=$(sha256sum $HOST | cut -c1-16) export=$(sha256sum $EXPORT | cut -c1-16) ids=$($HOST --mode id 2>&1 | grep -o '0x[0-9a-f]*' | tr '\n' ' ')"
|
||||
echo "RESULT fixtures $(ls $H/proving/fixtures | wc -l) files, v1=$(sha256sum $H/proving/fixtures/fees-v1-shards2.json | cut -c1-16) empty=$(sha256sum $H/proving/fixtures/block-72854-empty-block-first.json | cut -c1-16)"
|
||||
echo "RESULT node_now $($B/igneum-miner watch 1 grpc://127.0.0.1:26610 2>/dev/null | grep -o 'blocks=[0-9]*.*synced=[a-z]*' | tail -1)"
|
||||
echo "RESULT setup_done $(stamp)"
|
||||
153
tools/fleet/fleet.py
Executable file
153
tools/fleet/fleet.py
Executable file
|
|
@ -0,0 +1,153 @@
|
|||
#!/usr/bin/env python3
|
||||
"""The fleet orchestrator (Mac side). Registry: ~/Desktop/fleet/boxes.json ({instance: {...}}); per-box raw logs in
|
||||
~/Desktop/fleet/<instance>/. Uses vast.py for the provider and ssh with ~/.ssh/igneum-fleet for the boxes.
|
||||
|
||||
fleet.py rent-card "<gpu name>" <label> <archs> [--min-ram MB] [--max-ram MB] [--disk 60] [--phase 1]
|
||||
fleet.py wait [labels...] until ssh answers on every (named) box
|
||||
fleet.py setup [labels...] push the inputs and start box-setup.sh (nohup) on every (named) box
|
||||
fleet.py status [labels...] the last RESULT or STAGE line of setup.log and the node's sync line
|
||||
fleet.py run <script> [labels...] push tools/fleet/<script> and start it under nohup (box-matrix.sh, box-ember.sh, ...)
|
||||
fleet.py tail <label> [file] the last 30 lines of a box's log
|
||||
fleet.py sh <label> "<cmd>" run one command on a box
|
||||
fleet.py pull [labels...] rsync /root/fleet/out and the logs into ~/Desktop/fleet/<instance>/
|
||||
fleet.py destroy <labels...> pull, then destroy, then mark the registry
|
||||
"""
|
||||
import json, os, sys, subprocess, time, secrets, datetime
|
||||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
||||
import vast
|
||||
|
||||
HERE = os.path.dirname(os.path.abspath(__file__))
|
||||
ROOT = os.path.expanduser("~/Desktop/fleet"); REG = os.path.join(ROOT, "boxes.json")
|
||||
SSH_OPTS = ["-i", os.path.expanduser("~/.ssh/igneum-fleet"), "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null",
|
||||
"-o", "LogLevel=ERROR", "-o", "ConnectTimeout=20", "-o", "ServerAliveInterval=15"]
|
||||
INPUTS = [os.path.expanduser("~/Desktop/igneum-prove-wsl2-floor.zip"), os.path.join(HERE, "floor.patch"), os.path.join(HERE, "override.json")]
|
||||
|
||||
def now(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
def load(): return json.load(open(REG)) if os.path.exists(REG) else {}
|
||||
def save(reg): os.makedirs(ROOT, exist_ok=True); json.dump(reg, open(REG, "w"), indent=1)
|
||||
def boxes(labels, reg=None):
|
||||
reg = reg or load()
|
||||
sel = {k: v for k, v in reg.items() if v.get("state") != "destroyed" and (not labels or v["label"] in labels)}
|
||||
if labels and len(sel) != len(set(labels)): print("unknown or destroyed labels:", set(labels) - {v["label"] for v in sel.values()}, file=sys.stderr)
|
||||
return sel
|
||||
|
||||
def refresh_ssh(reg):
|
||||
live = {str(i.get("id")): i for i in vast.instances()}
|
||||
for iid, b in reg.items():
|
||||
if b.get("provider", "vast") == "vast" and iid in live:
|
||||
i = live[iid]; b["ssh_host"] = i.get("ssh_host"); b["ssh_port"] = i.get("ssh_port"); b["actual_status"] = i.get("actual_status"); b["status_msg"] = (i.get("status_msg") or "")[:120]
|
||||
save(reg); return reg
|
||||
|
||||
def ssh(b, cmd, timeout=120, capture=True):
|
||||
try:
|
||||
r = subprocess.run(["ssh"] + SSH_OPTS + ["-p", str(b["ssh_port"]), f"root@{b['ssh_host']}", cmd], capture_output=capture, text=True, timeout=timeout, stdin=subprocess.DEVNULL)
|
||||
except subprocess.TimeoutExpired: return 124, "", "timeout"
|
||||
return r.returncode, (r.stdout or ""), (r.stderr or "")
|
||||
|
||||
def scp(b, files, dest):
|
||||
return subprocess.run(["scp"] + SSH_OPTS + ["-P", str(b["ssh_port"])] + files + [f"root@{b['ssh_host']}:{dest}"], capture_output=True, text=True, timeout=900).returncode
|
||||
|
||||
def rent_card(gpu, label, archs, min_ram=None, max_ram=None, disk=60, phase="1", min_cores=4, gpus=1):
|
||||
offers = vast.search(gpu, n=12, min_cores=min_cores, gpus=gpus)
|
||||
if min_ram: offers = [o for o in offers if (o.get("gpu_ram") or 0) >= min_ram]
|
||||
if max_ram: offers = [o for o in offers if (o.get("gpu_ram") or 0) <= max_ram]
|
||||
if not offers: print(f"{label}: no offer for {gpu}"); return None
|
||||
# prefer a few more cores for the build when the price is close: score = price + 0.01 per missing core under 8
|
||||
offers.sort(key=lambda o: o["dph_total"] + 0.01 * max(0, 8 - (o.get("cpu_cores_effective") or 0)))
|
||||
o = offers[0]; print(label, "->", vast.fmt_offer(o))
|
||||
iid = vast.rent(o["id"], label, disk, offer=o)
|
||||
reg = load()
|
||||
reg[str(iid)] = {"label": label, "card": gpu, "vram_mb": o.get("gpu_ram"), "archs": archs, "offer": o["id"], "dph": o["dph_total"], "num_gpus": o.get("num_gpus"),
|
||||
"cores": o.get("cpu_cores_effective"), "ram_gb": round((o.get("cpu_ram") or 0) / 1024), "geo": o.get("geolocation"), "driver": o.get("driver_version"),
|
||||
"provider": "vast", "phase": phase, "state": "renting", "rented_at": now(), "wallet": "0x" + secrets.token_hex(20)}
|
||||
save(reg); os.makedirs(os.path.join(ROOT, str(iid)), exist_ok=True); return iid
|
||||
|
||||
def wait(labels, limit=1200):
|
||||
t0 = time.time()
|
||||
while True:
|
||||
reg = refresh_ssh(load()); pending = []
|
||||
for iid, b in boxes(labels, reg).items():
|
||||
if b.get("state") in ("destroyed",): continue
|
||||
if b.get("ssh_ok"): continue
|
||||
if not b.get("ssh_host") or not b.get("ssh_port"): pending.append((b["label"], b.get("actual_status"), b.get("status_msg"))); continue
|
||||
try: rc, out, err = ssh(b, "nvidia-smi --query-gpu=name,memory.total --format=csv,noheader", timeout=40)
|
||||
except subprocess.TimeoutExpired: rc, out, err = 1, "", "timeout"
|
||||
if rc == 0 and out.strip(): b["ssh_ok"] = True; b["state"] = "installing"; b["gpu_seen"] = out.strip().replace("\n", ";"); print(f"{b['label']}: ssh ok, {b['gpu_seen']}")
|
||||
else: pending.append((b["label"], b.get("actual_status"), (err or out).strip()[:80]))
|
||||
save(reg)
|
||||
if not pending: print("all up"); return
|
||||
if time.time() - t0 > limit: print("still pending:", pending); return
|
||||
print(f"waiting ({int(time.time()-t0)} s): " + "; ".join(f"{l} {s} {m}" for l, s, m in pending)); time.sleep(20)
|
||||
|
||||
def setup(labels):
|
||||
for iid, b in boxes(labels).items():
|
||||
if not b.get("ssh_ok"): print(b["label"], "no ssh yet"); continue
|
||||
ssh(b, "mkdir -p /root/fleet/in /root/fleet/out")
|
||||
if scp(b, INPUTS + [os.path.join(HERE, "box-setup.sh")], "/root/fleet/in/") != 0: print(b["label"], "scp failed"); continue
|
||||
rc, out, err = ssh(b, f"cd /root/fleet && chmod +x in/box-setup.sh && if pgrep -f in/box-setup.sh >/dev/null; then echo already; else ARCHS={b['archs']} LABEL={b['label']} WALLET={b['wallet']} setsid nohup in/box-setup.sh </dev/null >/dev/null 2>&1 & echo started; fi", timeout=60)
|
||||
print(b["label"], out.strip()[:40])
|
||||
b["setup_started"] = now(); print(b["label"], "setup started")
|
||||
save(load() | {k: v for k, v in boxes(labels).items()})
|
||||
|
||||
def status(labels):
|
||||
reg = load()
|
||||
for iid, b in boxes(labels, reg).items():
|
||||
if not b.get("ssh_ok"): print(f"{b['label']:<14} {iid} no ssh"); continue
|
||||
try:
|
||||
rc, out, err = ssh(b, "grep -E '^(RESULT|STAGE)' /root/fleet/setup.log 2>/dev/null | tail -1; grep -c '^RESULT setup_done' /root/fleet/setup.log 2>/dev/null; /opt/igneum/pkg/bin/igneum-miner watch 1 grpc://127.0.0.1:26610 2>/dev/null | grep -o 'blocks=[0-9]*.*synced=[a-z]*' | tail -1; grep -E '^(RESULT|STAGE)' /root/fleet/out/matrix.log 2>/dev/null | tail -1", timeout=60)
|
||||
lines = out.strip().split("\n"); print(f"{b['label']:<14} {iid} | " + " | ".join(l[:110] for l in lines))
|
||||
if any(l.startswith("1") for l in lines[1:2]) and b.get("state") == "installing": b["state"] = "running"; save(reg)
|
||||
except subprocess.TimeoutExpired: print(f"{b['label']:<14} {iid} ssh timeout")
|
||||
|
||||
def run(script, labels, env=""):
|
||||
for iid, b in boxes(labels).items():
|
||||
if scp(b, [os.path.join(HERE, script)], "/root/fleet/in/") != 0: print(b["label"], "scp failed"); continue
|
||||
name = os.path.basename(script)
|
||||
ssh(b, f"cd /root/fleet && chmod +x in/{name} && LABEL={b['label']} WALLET={b['wallet']} ARCHS={b['archs']} {env} setsid nohup in/{name} </dev/null >/dev/null 2>&1 & echo ok", timeout=60)
|
||||
print(b["label"], name, "started")
|
||||
|
||||
def pull(labels):
|
||||
for iid, b in boxes(labels).items():
|
||||
d = os.path.join(ROOT, iid); os.makedirs(d, exist_ok=True)
|
||||
r = subprocess.run(["rsync", "-az", "-e", "ssh " + " ".join(SSH_OPTS) + f" -p {b['ssh_port']}", f"root@{b['ssh_host']}:/root/fleet/out/", f"root@{b['ssh_host']}:/root/fleet/setup.log", f"root@{b['ssh_host']}:/root/fleet/node.log", d + "/"], capture_output=True, text=True, timeout=600)
|
||||
print(b["label"], "pulled" if r.returncode == 0 else f"pull failed: {r.stderr[:200]}")
|
||||
|
||||
def destroy(labels):
|
||||
reg = load()
|
||||
for iid, b in boxes(labels, reg).items():
|
||||
if b.get("ssh_ok"):
|
||||
try: pull([b["label"]])
|
||||
except Exception as e: print("pull failed", e)
|
||||
vast.destroy(iid); b["state"] = "destroyed"; b["destroyed_at"] = now()
|
||||
t0 = datetime.datetime.strptime(b["rented_at"], "%Y-%m-%dT%H:%M:%SZ"); h = (datetime.datetime.strptime(b["destroyed_at"], "%Y-%m-%dT%H:%M:%SZ") - t0).total_seconds() / 3600
|
||||
b["hours"] = round(h, 2); b["cost_usd"] = round(h * b["dph"], 3); print(b["label"], f"destroyed after {h:.2f} h, USD {b['cost_usd']:.2f}")
|
||||
save(reg)
|
||||
|
||||
if __name__ == "__main__":
|
||||
a = sys.argv[1:]
|
||||
if not a: print(__doc__); sys.exit(1)
|
||||
c = a[0]
|
||||
if c == "rent-card":
|
||||
kw = {}; pos = []
|
||||
i = 1
|
||||
while i < len(a):
|
||||
if a[i] == "--min-ram": kw["min_ram"] = int(a[i+1]); i += 2
|
||||
elif a[i] == "--max-ram": kw["max_ram"] = int(a[i+1]); i += 2
|
||||
elif a[i] == "--disk": kw["disk"] = int(a[i+1]); i += 2
|
||||
elif a[i] == "--phase": kw["phase"] = a[i+1]; i += 2
|
||||
elif a[i] == "--min-cores": kw["min_cores"] = int(a[i+1]); i += 2
|
||||
elif a[i] == "--gpus": kw["gpus"] = int(a[i+1]); i += 2
|
||||
else: pos.append(a[i]); i += 1
|
||||
rent_card(pos[0], pos[1], pos[2], **kw)
|
||||
elif c == "wait": wait(a[1:])
|
||||
elif c == "setup": setup(a[1:])
|
||||
elif c == "status": status(a[1:])
|
||||
elif c == "run": run(a[1], a[2:])
|
||||
elif c == "pull": pull(a[1:])
|
||||
elif c == "destroy": destroy(a[1:])
|
||||
elif c == "tail":
|
||||
b = [v for v in boxes([a[1]]).values()][0]; f = a[2] if len(a) > 2 else "/root/fleet/setup.log"; print(ssh(b, f"tail -n 30 {f}")[1])
|
||||
elif c == "sh":
|
||||
b = [v for v in boxes([a[1]]).values()][0]; rc, out, err = ssh(b, a[2], timeout=600); print(out, err)
|
||||
elif c == "list":
|
||||
for iid, b in load().items(): print(iid, b["label"], b["card"], b.get("state"), f"${b['dph']:.3f}/h", b.get("ssh_host"), b.get("ssh_port"))
|
||||
345
tools/fleet/floor.patch
Normal file
345
tools/fleet/floor.patch
Normal file
|
|
@ -0,0 +1,345 @@
|
|||
diff --git a/sp1-gpu/crates/cuda/src/task.rs b/sp1-gpu/crates/cuda/src/task.rs
|
||||
index a503a86..813016b 100644
|
||||
--- a/sp1-gpu/crates/cuda/src/task.rs
|
||||
+++ b/sp1-gpu/crates/cuda/src/task.rs
|
||||
@@ -149,7 +149,15 @@ pub enum GlobalTaskPoolBuildError {
|
||||
|
||||
impl TaskPoolBuilder {
|
||||
pub fn new() -> Self {
|
||||
- Self { capacity: None, device: CudaDevice(0), mem_release_threshold: u64::MAX }
|
||||
+ // Igneum prover-floor patch: upstream keeps every freed device allocation in the pool for the process's
|
||||
+ // life (threshold u64::MAX), so the prover holds its high-water mark between shards on a card it shares
|
||||
+ // with a miner. `SP1_GPU_MEM_RELEASE_THRESHOLD=<bytes>` sets the pool's release threshold (0 returns
|
||||
+ // freed memory to the driver at once); unset, upstream's behaviour.
|
||||
+ let mem_release_threshold = std::env::var("SP1_GPU_MEM_RELEASE_THRESHOLD")
|
||||
+ .ok()
|
||||
+ .and_then(|s| s.parse::<u64>().ok())
|
||||
+ .unwrap_or(u64::MAX);
|
||||
+ Self { capacity: None, device: CudaDevice(0), mem_release_threshold }
|
||||
}
|
||||
|
||||
pub fn num_tasks(mut self, num_tasks: usize) -> Self {
|
||||
diff --git a/sp1-gpu/crates/jagged_tracegen/src/lib.rs b/sp1-gpu/crates/jagged_tracegen/src/lib.rs
|
||||
index 579f70a..2264044 100644
|
||||
--- a/sp1-gpu/crates/jagged_tracegen/src/lib.rs
|
||||
+++ b/sp1-gpu/crates/jagged_tracegen/src/lib.rs
|
||||
@@ -481,6 +481,33 @@ async fn device_preprocessed_tracegen<A: CudaTracegenAir<Felt>>(
|
||||
named_traces
|
||||
}
|
||||
|
||||
+/// Igneum prover-floor patch: the dense elements a set of traces will occupy once `generate_jagged_traces`
|
||||
+/// has laid them out, that is the sum of their buffers padded to the next multiple of 2^log_stacking_height
|
||||
+/// (the "final padding" step below). Each phase (preprocessed, then main) is padded on its own.
|
||||
+pub fn padded_trace_elements(
|
||||
+ traces: &BTreeMap<String, Trace<TaskScope>>,
|
||||
+ log_stacking_height: u32,
|
||||
+) -> usize {
|
||||
+ let total: usize = traces
|
||||
+ .values()
|
||||
+ .map(|t| match t {
|
||||
+ Trace::Real(trace) => trace.guts().as_buffer().len(),
|
||||
+ Trace::Padding(_) => 0,
|
||||
+ })
|
||||
+ .sum();
|
||||
+ total.next_multiple_of(1 << log_stacking_height)
|
||||
+}
|
||||
+
|
||||
+/// Igneum prover-floor patch: the capacity to allocate for a trace set: the exact padded size plus one
|
||||
+/// stacking height of slack, never more than the prover's `max_trace_size`. `SP1_GPU_FLOOR_EXACT=0` restores
|
||||
+/// upstream's full-capacity allocation.
|
||||
+fn floor_capacity(max_trace_size: usize, needed: usize, log_stacking_height: u32) -> usize {
|
||||
+ if std::env::var("SP1_GPU_FLOOR_EXACT").map(|v| v == "0").unwrap_or(false) {
|
||||
+ return max_trace_size;
|
||||
+ }
|
||||
+ max_trace_size.min(needed + (1 << log_stacking_height))
|
||||
+}
|
||||
+
|
||||
async fn allocate_and_initialize_traces(
|
||||
preprocessed_traces: BTreeMap<String, Trace<TaskScope>>,
|
||||
max_trace_size: usize,
|
||||
@@ -494,6 +521,11 @@ async fn allocate_and_initialize_traces(
|
||||
|
||||
let total_gb = total_bytes as f64 / (1 << 30) as f64;
|
||||
tracing::debug!("Allocating {:?} GB of traces", total_gb);
|
||||
+ if std::env::var("SP1_GPU_FLOOR_LOG").is_ok() {
|
||||
+ eprintln!(
|
||||
+ "FLOOR tracegen alloc capacity_elements={max_trace_size} bytes={total_bytes} ({total_gb:.3} GB)"
|
||||
+ );
|
||||
+ }
|
||||
let mut dense_data: Buffer<Felt, TaskScope> =
|
||||
Buffer::with_capacity_in(max_trace_size, backend.clone());
|
||||
let mut col_index: Buffer<u32, TaskScope> =
|
||||
@@ -677,9 +709,14 @@ pub async fn setup_tracegen<A: CudaTracegenAir<Felt>>(
|
||||
let preprocessed_traces =
|
||||
device_preprocessed_tracegen(program, host_phase_tracegen, backend).await;
|
||||
|
||||
+ let capacity = floor_capacity(
|
||||
+ max_trace_size,
|
||||
+ padded_trace_elements(&preprocessed_traces, log_stacking_height),
|
||||
+ log_stacking_height,
|
||||
+ );
|
||||
let jagged_traces = allocate_and_initialize_traces(
|
||||
preprocessed_traces,
|
||||
- max_trace_size,
|
||||
+ capacity,
|
||||
log_stacking_height,
|
||||
max_log_row_count,
|
||||
backend,
|
||||
@@ -906,6 +943,11 @@ pub async fn main_tracegen<GC: IopCtx<F = Felt>, A: CudaTracegenAir<Felt>>(
|
||||
|
||||
log_chip_stats(machine, &chip_set, &traces);
|
||||
|
||||
+ // Igneum prover-floor patch: the key's buffer is sized to its preprocessed traces at setup (upstream sized it
|
||||
+ // for a whole shard), so grow it here to what this shard needs before the main traces are appended: a bigger
|
||||
+ // dense buffer and column index, the preprocessed region copied device to device, swapped into the key.
|
||||
+ grow_for_main(&mut jagged_traces.preprocessed_traces, &traces, log_stacking_height, backend);
|
||||
+
|
||||
copy_main_jagged_traces(
|
||||
traces,
|
||||
&mut jagged_traces.preprocessed_traces,
|
||||
@@ -918,6 +960,61 @@ pub async fn main_tracegen<GC: IopCtx<F = Felt>, A: CudaTracegenAir<Felt>>(
|
||||
(public_values, chip_set, permit)
|
||||
}
|
||||
|
||||
+/// Igneum prover-floor patch: see `main_tracegen`. The need is the preprocessed phase as laid out (its padded
|
||||
+/// end, `preprocessed_offset`) plus the main traces padded to the stacking height plus one stacking height of
|
||||
+/// slack; a buffer at least that big is left alone. The process aborts, loudly, if the copy cannot be made,
|
||||
+/// because a panic inside a prover task is what left sweep 2 hanging on the client's socket.
|
||||
+fn grow_for_main(
|
||||
+ jagged: &mut JaggedTraceMle<Felt, TaskScope>,
|
||||
+ main_traces: &BTreeMap<String, Trace<TaskScope>>,
|
||||
+ log_stacking_height: u32,
|
||||
+ backend: &TaskScope,
|
||||
+) {
|
||||
+ let pre_end = jagged.dense().preprocessed_offset;
|
||||
+ let needed = pre_end
|
||||
+ + padded_trace_elements(main_traces, log_stacking_height)
|
||||
+ + (1 << log_stacking_height);
|
||||
+ let have = jagged.dense().dense.capacity();
|
||||
+ if have >= needed {
|
||||
+ return;
|
||||
+ }
|
||||
+ let mut new_dense: Buffer<Felt, TaskScope> = Buffer::with_capacity_in(needed, backend.clone());
|
||||
+ let mut new_col_index: Buffer<u32, TaskScope> =
|
||||
+ Buffer::with_capacity_in(needed >> 1, backend.clone());
|
||||
+ unsafe {
|
||||
+ new_dense.assume_init();
|
||||
+ new_col_index.assume_init();
|
||||
+ }
|
||||
+ {
|
||||
+ let JaggedMle { dense_data, col_index, .. } = &mut **jagged;
|
||||
+ let src_dense: &Slice<_, _> = &dense_data.dense[..pre_end];
|
||||
+ let dst_dense: &mut Slice<_, _> = &mut new_dense[..pre_end];
|
||||
+ let src_col: &Slice<_, _> = &col_index[..pre_end >> 1];
|
||||
+ let dst_col: &mut Slice<_, _> = &mut new_col_index[..pre_end >> 1];
|
||||
+ unsafe {
|
||||
+ if dst_dense.copy_from_slice(src_dense, backend).is_err()
|
||||
+ || dst_col.copy_from_slice(src_col, backend).is_err()
|
||||
+ {
|
||||
+ eprintln!("FLOOR grow FAILED: could not copy the preprocessed region ({pre_end} elements) into the grown buffer ({needed} elements); aborting instead of hanging");
|
||||
+ std::process::abort();
|
||||
+ }
|
||||
+ }
|
||||
+ }
|
||||
+ if std::env::var("SP1_GPU_FLOOR_LOG").is_ok() {
|
||||
+ eprintln!(
|
||||
+ "FLOOR grow key buffer {have} -> {needed} elements (preprocessed {pre_end}, {} bytes)",
|
||||
+ needed * 6
|
||||
+ );
|
||||
+ }
|
||||
+ let JaggedMle { dense_data, col_index, .. } = &mut **jagged;
|
||||
+ dense_data.dense = new_dense;
|
||||
+ *col_index = new_col_index;
|
||||
+ unsafe {
|
||||
+ dense_data.dense.set_len(pre_end);
|
||||
+ col_index.set_len(pre_end >> 1);
|
||||
+ }
|
||||
+}
|
||||
+
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn main_tracegen_permit<GC: IopCtx<F = Felt>, A: CudaTracegenAir<Felt>>(
|
||||
machine: &Machine<Felt, A>,
|
||||
@@ -984,9 +1081,15 @@ pub async fn full_tracegen<A: CudaTracegenAir<Felt>>(
|
||||
|
||||
log_chip_stats(machine, &chip_set, &main_traces);
|
||||
|
||||
+ let capacity = floor_capacity(
|
||||
+ max_trace_size,
|
||||
+ padded_trace_elements(&preprocessed_traces, log_stacking_height)
|
||||
+ + padded_trace_elements(&main_traces, log_stacking_height),
|
||||
+ log_stacking_height,
|
||||
+ );
|
||||
let mut jagged_mle = allocate_and_initialize_traces(
|
||||
preprocessed_traces,
|
||||
- max_trace_size,
|
||||
+ capacity,
|
||||
log_stacking_height,
|
||||
max_log_row_count,
|
||||
backend,
|
||||
@@ -1002,6 +1105,18 @@ pub async fn full_tracegen<A: CudaTracegenAir<Felt>>(
|
||||
)
|
||||
.await;
|
||||
|
||||
+ if std::env::var("SP1_GPU_FLOOR_LOG").is_ok() {
|
||||
+ let dense = jagged_mle.dense();
|
||||
+ let (free, total) = sp1_gpu_cudart::cuda_memory_info().unwrap_or((0, 0));
|
||||
+ eprintln!(
|
||||
+ "FLOOR tracegen used preprocessed_elements={} main_elements={} dense_len={} capacity_elements={capacity} max_trace_size={max_trace_size} device_used_mib={}",
|
||||
+ dense.preprocessed_offset,
|
||||
+ dense.main_size(),
|
||||
+ dense.dense.len(),
|
||||
+ (total - free) >> 20
|
||||
+ );
|
||||
+ }
|
||||
+
|
||||
(public_values, jagged_mle, chip_set, permit)
|
||||
}
|
||||
|
||||
diff --git a/sp1-gpu/crates/prover_components/src/builder.rs b/sp1-gpu/crates/prover_components/src/builder.rs
|
||||
index 5dccd9d..574d4fa 100644
|
||||
--- a/sp1-gpu/crates/prover_components/src/builder.rs
|
||||
+++ b/sp1-gpu/crates/prover_components/src/builder.rs
|
||||
@@ -23,28 +23,75 @@ use crate::{
|
||||
SP1CudaProverComponents,
|
||||
};
|
||||
|
||||
+/// Igneum prover-floor patch (5 October 2026). Upstream sizes every device buffer for a 24 GB card or larger
|
||||
+/// and panics below that, whatever the shard. Here the card's memory (or `SP1_GPU_MEMORY_BUDGET_GB`) picks a
|
||||
+/// tier, and `SP1_GPU_ELEMENT_THRESHOLD` / `SP1_GPU_RECURSION_TRACE_ALLOCATION` set the two buffers directly.
|
||||
+/// The proof format, the verifier and the program ids do not change: the element threshold only decides where
|
||||
+/// the executor splits shards, as upstream's own 24 GB tier already does.
|
||||
+fn env_usize(name: &str) -> Option<usize> {
|
||||
+ std::env::var(name).ok().and_then(|s| s.parse::<usize>().ok())
|
||||
+}
|
||||
+
|
||||
+fn env_f64(name: &str) -> Option<f64> {
|
||||
+ std::env::var(name).ok().and_then(|s| s.parse::<f64>().ok())
|
||||
+}
|
||||
+
|
||||
+/// The core element threshold for a memory budget in GB (upstream's own figure for the budget, +4, as it
|
||||
+/// computed it: a 32 GB card is 36, a 24 GB card 28, a 16 GB card 20, a 12 GB card 16).
|
||||
+pub fn element_threshold_for_budget(gpu_memory_gb: usize, full_size_shards: bool) -> u64 {
|
||||
+ if gpu_memory_gb > 30 || (full_size_shards && gpu_memory_gb >= 24) {
|
||||
+ ELEMENT_THRESHOLD
|
||||
+ } else if gpu_memory_gb >= 24 {
|
||||
+ ELEMENT_THRESHOLD - (1 << 26) - (1 << 25) - (1 << 24)
|
||||
+ } else if gpu_memory_gb >= 18 {
|
||||
+ (1 << 27) + (1 << 26)
|
||||
+ } else {
|
||||
+ 1 << 27
|
||||
+ }
|
||||
+}
|
||||
+
|
||||
+/// The recursion trace allocation (elements) for a memory budget.
|
||||
+pub fn recursion_trace_allocation_for_budget(gpu_memory_gb: usize) -> usize {
|
||||
+ if gpu_memory_gb >= 24 {
|
||||
+ RECURSION_TRACE_ALLOCATION
|
||||
+ } else {
|
||||
+ RECURSION_TRACE_ALLOCATION
|
||||
+ }
|
||||
+}
|
||||
+
|
||||
+pub fn gpu_memory_gb() -> usize {
|
||||
+ let gb = 1024.0 * 1024.0 * 1024.0;
|
||||
+ match env_f64("SP1_GPU_MEMORY_BUDGET_GB") {
|
||||
+ Some(b) => (b.ceil() as usize) + 4,
|
||||
+ None => (((cuda_memory_info().unwrap().1 as f64) / gb).ceil() as usize) + 4,
|
||||
+ }
|
||||
+}
|
||||
+
|
||||
+pub fn recursion_trace_allocation() -> usize {
|
||||
+ env_usize("SP1_GPU_RECURSION_TRACE_ALLOCATION")
|
||||
+ .unwrap_or_else(|| recursion_trace_allocation_for_budget(gpu_memory_gb()))
|
||||
+}
|
||||
+
|
||||
pub fn local_gpu_opts() -> SP1CoreOpts {
|
||||
let mut opts = SP1CoreOpts::default();
|
||||
|
||||
let log2_shard_size = 24;
|
||||
opts.shard_size = 1 << log2_shard_size;
|
||||
|
||||
- let gb = 1024.0 * 1024.0 * 1024.0;
|
||||
-
|
||||
- // Get the amount of memory on the GPU.
|
||||
- let gpu_memory_gb: usize = (((cuda_memory_info().unwrap().1 as f64) / gb).ceil() as usize) + 4;
|
||||
-
|
||||
- if gpu_memory_gb < 24 {
|
||||
- panic!("Unsupported GPU memory: {gpu_memory_gb}, must be at least 24GB");
|
||||
- }
|
||||
+ // The card's memory plus 4, as upstream computed it (a 32 GB card reads 36), or the budget given.
|
||||
+ let gpu_memory_gb = gpu_memory_gb();
|
||||
|
||||
- let shard_threshold = if !opts.full_size_shards && gpu_memory_gb <= 30 {
|
||||
- ELEMENT_THRESHOLD - (1 << 26) - (1 << 25) - (1 << 24)
|
||||
- } else {
|
||||
- ELEMENT_THRESHOLD
|
||||
+ let shard_threshold = match env_usize("SP1_GPU_ELEMENT_THRESHOLD") {
|
||||
+ Some(t) => t as u64,
|
||||
+ None => element_threshold_for_budget(gpu_memory_gb, opts.full_size_shards),
|
||||
};
|
||||
+ let height_threshold = opts.sharding_threshold.height_threshold;
|
||||
|
||||
- tracing::debug!("Shard threshold: {shard_threshold}");
|
||||
+ eprintln!(
|
||||
+ "FLOOR opts gpu_memory_gb={gpu_memory_gb} element_threshold={shard_threshold} height_threshold={height_threshold} recursion_trace_allocation={} full_size_shards={}",
|
||||
+ recursion_trace_allocation(),
|
||||
+ opts.full_size_shards
|
||||
+ );
|
||||
opts.sharding_threshold.element_threshold = shard_threshold;
|
||||
|
||||
opts.global_dependencies_opt = true;
|
||||
@@ -92,7 +139,7 @@ pub async fn recursion_prover_and_verifier(
|
||||
) {
|
||||
let recursion_verifier = SP1CudaProverComponents::compress_verifier();
|
||||
(
|
||||
- new_cuda_prover(&recursion_verifier, RECURSION_TRACE_ALLOCATION, 4, false, false, scope)
|
||||
+ new_cuda_prover(&recursion_verifier, recursion_trace_allocation(), 4, false, false, scope)
|
||||
.await,
|
||||
recursion_verifier,
|
||||
)
|
||||
diff --git a/sp1-gpu/crates/server/src/server.rs b/sp1-gpu/crates/server/src/server.rs
|
||||
index 4035f1f..0d0d907 100644
|
||||
--- a/sp1-gpu/crates/server/src/server.rs
|
||||
+++ b/sp1-gpu/crates/server/src/server.rs
|
||||
@@ -157,6 +157,7 @@ impl Server {
|
||||
};
|
||||
let pk = CachedProgram { elf: Arc::new(Elf::Dynamic(elf.into())), vk: vk.clone() };
|
||||
ctx.pk_cache.insert(elf_hash, pk);
|
||||
+ floor_memory_line("after setup");
|
||||
Response::Setup { id: elf_hash, vk }
|
||||
}
|
||||
Request::Destroy { key } => {
|
||||
@@ -177,15 +178,31 @@ impl Server {
|
||||
);
|
||||
};
|
||||
let context = SP1Context::builder().proof_nonce(proof_nonce).build();
|
||||
- match prover.prove_with_mode(&cached.elf, stdin, context, mode).await {
|
||||
+ let started = std::time::Instant::now();
|
||||
+ let response = match prover.prove_with_mode(&cached.elf, stdin, context, mode).await {
|
||||
Ok(proof) => Response::Proof { proof },
|
||||
Err(e) => Response::ProverError(e.to_string()),
|
||||
- }
|
||||
+ };
|
||||
+ floor_memory_line(&format!("after prove {:?} in {:.1} s", mode, started.elapsed().as_secs_f64()));
|
||||
+ response
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+/// Igneum prover-floor patch: the device memory in use (total minus free, as the driver reports it) at the
|
||||
+/// points that bound a proof, so a run's log carries the terms of the peak without a sampler.
|
||||
+fn floor_memory_line(what: &str) {
|
||||
+ if let Ok((free, total)) = sp1_gpu_cudart::cuda_memory_info() {
|
||||
+ eprintln!(
|
||||
+ "FLOOR memory {what}: device_used_mib={} free_mib={} total_mib={}",
|
||||
+ (total - free) >> 20,
|
||||
+ free >> 20,
|
||||
+ total >> 20
|
||||
+ );
|
||||
+ }
|
||||
+}
|
||||
+
|
||||
fn sha256(data: &[u8]) -> [u8; 32] {
|
||||
use sha2::{Digest, Sha256};
|
||||
let mut hasher = Sha256::new();
|
||||
1
tools/fleet/override.json
Normal file
1
tools/fleet/override.json
Normal file
|
|
@ -0,0 +1 @@
|
|||
{"difficulty_v2_activation_daa":33000,"proving_v0_activation_daa":84100,"fees_v1_activation_daa":210000,"finality_v3_activation_daa":135200,"program_class_v3_activation_daa":154800,"proving_v1_activation_daa":154800,"proving_v1_segment_blocks":8,"proving_v1_unproven_daa":600,"proving_v1_aggregator_share_bps":1000,"proving_v1_fresh_rule_daa":198000}
|
||||
61
tools/fleet/page.py
Executable file
61
tools/fleet/page.py
Executable file
|
|
@ -0,0 +1,61 @@
|
|||
#!/usr/bin/env python3
|
||||
"""Builds ~/Desktop/fleet/fleet.json (the fleet page's single source) from the registry (boxes.json), the ledger, the
|
||||
result rows (results.json: {card_key: {...}}), the night rows (night.json: [..]), the phases (phases.json) and the log
|
||||
(log.jsonl), then publishes it with the master checkout's tools/fleet/publish-fleet.sh (at most once a minute).
|
||||
|
||||
page.py log "<text>" append a log line (and rebuild, no publish)
|
||||
page.py phase p1 state=running done=3 note="..."
|
||||
page.py build rebuild only
|
||||
page.py publish rebuild and publish (skipped when the last publish was under 60 s ago unless --force)
|
||||
"""
|
||||
import json, os, sys, time, datetime, subprocess
|
||||
ROOT = os.path.expanduser("~/Desktop/fleet"); OUT = os.path.join(ROOT, "fleet.json")
|
||||
PUBLISH = "/Users/joshm/Projects/igneum/tools/fleet/publish-fleet.sh"
|
||||
def now(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
def jl(name, default):
|
||||
p = os.path.join(ROOT, name)
|
||||
if not os.path.exists(p): return default
|
||||
if name.endswith(".jsonl"): return [json.loads(l) for l in open(p) if l.strip()]
|
||||
return json.load(open(p))
|
||||
def hours(a, b=None):
|
||||
t0 = datetime.datetime.strptime(a, "%Y-%m-%dT%H:%M:%SZ")
|
||||
t1 = datetime.datetime.strptime(b, "%Y-%m-%dT%H:%M:%SZ") if b else datetime.datetime.utcnow()
|
||||
return max((t1 - t0).total_seconds(), 0) / 3600
|
||||
def build():
|
||||
reg = jl("boxes.json", {}); res = jl("results.json", {}); night = jl("night.json", []); log = jl("log.jsonl", [])
|
||||
phases = jl("phases.json", [
|
||||
{"order": 1, "name": "Memory matrix on real cards", "state": "planned", "planned": 11, "done": 0, "note": ""},
|
||||
{"order": 2, "name": "Prover fleet on the devnet", "state": "planned", "planned": 25, "done": 0, "note": ""},
|
||||
{"order": 3, "name": "The 8x rigs", "state": "planned", "planned": 2, "done": 0, "note": ""},
|
||||
{"order": 4, "name": "AMD mining", "state": "planned", "planned": 3, "done": 0, "note": ""}])
|
||||
boxes = []; vast = 0.0; runpod = 0.0
|
||||
for iid, b in reg.items():
|
||||
h = b.get("hours") if b.get("state") == "destroyed" else hours(b["rented_at"])
|
||||
cost = round(h * b["dph"], 3)
|
||||
if b.get("provider", "vast") == "vast": vast += cost
|
||||
else: runpod += cost
|
||||
boxes.append({"id": iid, "label": b["label"], "card": b["card"], "vram_gb": round((b.get("vram_mb") or 0) / 1024), "provider": b.get("provider", "vast"),
|
||||
"rate_usd_h": b["dph"], "phase": b.get("phase", "1"), "state": "done" if b.get("state") == "destroyed" else b.get("state", "renting"),
|
||||
"started_at": b["rented_at"], "ended_at": b.get("destroyed_at"), "hours": round(h, 2), "cost_usd": cost, "doing": b.get("doing", "")})
|
||||
page = {"spend": {"vast_usd": round(vast, 2), "runpod_usd": round(runpod, 2), "cap_vast": 900, "cap_runpod": 450, "updated_at": now()},
|
||||
"phases": phases, "boxes": boxes, "results": list(res.values()), "night": night, "log": log[-200:]}
|
||||
json.dump(page, open(OUT, "w"), indent=1); return page
|
||||
def publish(force=False):
|
||||
stamp = os.path.join(ROOT, ".last-publish")
|
||||
if not force and os.path.exists(stamp) and time.time() - os.path.getmtime(stamp) < 60: print("published under a minute ago; skipped"); return
|
||||
build(); open(stamp, "w").write(now())
|
||||
r = subprocess.run([PUBLISH, OUT], capture_output=True, text=True, timeout=300); print((r.stdout + r.stderr).strip()[-300:])
|
||||
if __name__ == "__main__":
|
||||
c = sys.argv[1] if len(sys.argv) > 1 else "build"
|
||||
if c == "log":
|
||||
with open(os.path.join(ROOT, "log.jsonl"), "a") as f: f.write(json.dumps({"at": now(), "text": sys.argv[2]}) + "\n")
|
||||
build()
|
||||
elif c == "phase":
|
||||
ph = jl("phases.json", None) or build()["phases"]
|
||||
for p in ph:
|
||||
if f"p{p['order']}" == sys.argv[2]:
|
||||
for kv in sys.argv[3:]:
|
||||
k, v = kv.split("=", 1); p[k] = int(v) if v.isdigit() else v
|
||||
json.dump(ph, open(os.path.join(ROOT, "phases.json"), "w")); build()
|
||||
elif c == "publish": publish("--force" in sys.argv)
|
||||
else: build(); print(OUT)
|
||||
122
tools/fleet/vast.py
Executable file
122
tools/fleet/vast.py
Executable file
|
|
@ -0,0 +1,122 @@
|
|||
#!/usr/bin/env python3
|
||||
"""Vast.ai client for the Igneum GPU fleet (6 October 2026). The key is read from ~/.config/vast/credentials and never
|
||||
printed. Every rent and destroy is appended to ~/Desktop/fleet/ledger.jsonl with the offer's hourly price, so the spend
|
||||
can be summed without the provider's statement.
|
||||
|
||||
vast.py search "RTX 3090" [--n 8] [--min-cores 6] [--max-dph 0.5] [--gpus 1]
|
||||
vast.py rent <offer_id> --label <label> [--disk 60] [--image nvidia/cuda:12.8.1-devel-ubuntu24.04] [--onstart FILE]
|
||||
vast.py list every instance: id, label, gpu, state, ssh host and port, dph
|
||||
vast.py ssh <instance_id> prints the ssh command line (host, port) for the fleet key
|
||||
vast.py destroy <instance_id>... destroys and writes the ledger line
|
||||
vast.py spend the ledger's running total and the account's credit
|
||||
"""
|
||||
import json, os, sys, time, urllib.request, urllib.parse, urllib.error, argparse, datetime
|
||||
|
||||
KEY = open(os.path.expanduser("~/.config/vast/credentials")).read().strip()
|
||||
BASE = "https://console.vast.ai/api"
|
||||
LEDGER = os.path.expanduser("~/Desktop/fleet/ledger.jsonl")
|
||||
IMAGE = "nvidia/cuda:12.8.1-devel-ubuntu24.04"
|
||||
|
||||
def call(method, path, body=None, timeout=90):
|
||||
req = urllib.request.Request(BASE + path, data=(json.dumps(body).encode() if body is not None else None),
|
||||
headers={"Authorization": "Bearer " + KEY, "Content-Type": "application/json", "Accept": "application/json"},
|
||||
method=method)
|
||||
for attempt in range(4):
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=timeout) as r:
|
||||
return json.loads(r.read() or b"{}")
|
||||
except urllib.error.HTTPError as e:
|
||||
txt = e.read().decode(errors="replace")
|
||||
if e.code in (429, 502, 503, 504) and attempt < 3:
|
||||
time.sleep(3 * (attempt + 1)); continue
|
||||
raise SystemExit(f"{method} {path}: HTTP {e.code}: {txt[:400]}")
|
||||
except (urllib.error.URLError, TimeoutError) as e:
|
||||
if attempt < 3: time.sleep(3 * (attempt + 1)); continue
|
||||
raise
|
||||
|
||||
def now(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
|
||||
def ledger(row):
|
||||
os.makedirs(os.path.dirname(LEDGER), exist_ok=True)
|
||||
with open(LEDGER, "a") as f: f.write(json.dumps(row) + "\n")
|
||||
|
||||
def search(gpu, n=8, min_cores=4, max_dph=None, gpus=1, min_disk=70, min_rel=0.95, min_down=200, cuda=12.8):
|
||||
q = {"gpu_name": {"eq": gpu}, "rentable": {"eq": True}, "verified": {"eq": True}, "num_gpus": {"eq": gpus},
|
||||
"cuda_max_good": {"gte": cuda}, "disk_space": {"gte": min_disk}, "reliability2": {"gte": min_rel},
|
||||
"inet_down": {"gte": min_down}, "cpu_cores_effective": {"gte": min_cores},
|
||||
"order": [["dph_total", "asc"]], "type": "on-demand", "limit": n}
|
||||
if max_dph: q["dph_total"] = {"lte": max_dph}
|
||||
d = call("GET", "/v0/bundles/?q=" + urllib.parse.quote(json.dumps(q)))
|
||||
return d.get("offers", [])
|
||||
|
||||
def fmt_offer(o):
|
||||
return (f"id {o['id']} ${o['dph_total']:.3f}/h {o.get('num_gpus')}x {o.get('gpu_name')} {o.get('gpu_ram')}MB cuda{o.get('cuda_max_good')} "
|
||||
f"cores{o.get('cpu_cores_effective',0):.0f} ram{o.get('cpu_ram',0)/1024:.0f}G disk{o.get('disk_space',0):.0f} dl{o.get('inet_down',0):.0f} "
|
||||
f"ul{o.get('inet_up',0):.0f} rel{o.get('reliability2',0):.3f} {o.get('geolocation')} drv{o.get('driver_version')}")
|
||||
|
||||
def rent(offer_id, label, disk=60, image=IMAGE, onstart=None, offer=None):
|
||||
body = {"client_id": "me", "image": image, "disk": disk, "label": label, "runtype": "ssh", "image_login": None,
|
||||
"python_utf8": False, "lang_utf8": False, "use_jupyter_lab": False, "cancel_unavail": True}
|
||||
if onstart: body["onstart"] = open(onstart).read()
|
||||
# the offer's price for the ledger
|
||||
o = offer or {}
|
||||
d = call("PUT", f"/v0/asks/{offer_id}/", body)
|
||||
if not d.get("success"): raise SystemExit(f"rent failed: {d}")
|
||||
iid = d.get("new_contract")
|
||||
row = {"t": now(), "event": "rent", "instance": iid, "offer": int(offer_id), "label": label, "gpu": o.get("gpu_name"),
|
||||
"num_gpus": o.get("num_gpus"), "dph": o.get("dph_total"), "disk": disk, "geo": o.get("geolocation")}
|
||||
ledger(row); print(json.dumps(row)); return iid
|
||||
|
||||
def instances():
|
||||
try:
|
||||
d = call("GET", "/v0/instances/?owner=me")
|
||||
return d.get("instances", [])
|
||||
except SystemExit:
|
||||
d = call("GET", "/v1/instances/")
|
||||
return d.get("instances", d if isinstance(d, list) else [])
|
||||
|
||||
def destroy(iid):
|
||||
try: d = call("DELETE", f"/v0/instances/{iid}/", {})
|
||||
except SystemExit: d = call("DELETE", f"/v1/instances/{iid}/", {})
|
||||
row = {"t": now(), "event": "destroy", "instance": int(iid), "ok": bool(d.get("success", True))}
|
||||
ledger(row); print(json.dumps(row))
|
||||
|
||||
def spend():
|
||||
rows = [json.loads(l) for l in open(LEDGER)] if os.path.exists(LEDGER) else []
|
||||
starts = {r["instance"]: r for r in rows if r["event"] == "rent"}
|
||||
stops = {r["instance"]: r for r in rows if r["event"] == "destroy"}
|
||||
total = 0.0; lines = []
|
||||
tnow = datetime.datetime.now(datetime.timezone.utc)
|
||||
for iid, r in starts.items():
|
||||
t0 = datetime.datetime.strptime(r["t"], "%Y-%m-%dT%H:%M:%SZ").replace(tzinfo=datetime.timezone.utc)
|
||||
t1 = datetime.datetime.strptime(stops[iid]["t"], "%Y-%m-%dT%H:%M:%SZ").replace(tzinfo=datetime.timezone.utc) if iid in stops else tnow
|
||||
h = max((t1 - t0).total_seconds(), 60) / 3600
|
||||
cost = h * float(r.get("dph") or 0) + (r.get("disk") or 0) * 0.0002 * h # storage about USD 0.2 per TB-hour, approximate
|
||||
total += cost
|
||||
lines.append(f"{iid} {r['label']:<22} {r.get('gpu')!s:<16} {h:5.2f} h x ${r.get('dph') or 0:.3f} = ${cost:6.2f} {'running' if iid not in stops else 'destroyed'}")
|
||||
me = call("GET", "/v0/users/current/")
|
||||
print("\n".join(lines)); print(f"ledger total USD {total:.2f}; account credit USD {me.get('credit')}; balance {me.get('balance')}")
|
||||
|
||||
if __name__ == "__main__":
|
||||
ap = argparse.ArgumentParser(); sub = ap.add_subparsers(dest="cmd", required=True)
|
||||
s = sub.add_parser("search"); s.add_argument("gpu"); s.add_argument("--n", type=int, default=8); s.add_argument("--min-cores", type=int, default=4)
|
||||
s.add_argument("--max-dph", type=float); s.add_argument("--gpus", type=int, default=1); s.add_argument("--min-disk", type=int, default=70); s.add_argument("--json", action="store_true")
|
||||
r = sub.add_parser("rent"); r.add_argument("offer"); r.add_argument("--label", required=True); r.add_argument("--disk", type=int, default=60); r.add_argument("--image", default=IMAGE); r.add_argument("--onstart")
|
||||
sub.add_parser("list"); sub.add_parser("spend")
|
||||
d = sub.add_parser("destroy"); d.add_argument("ids", nargs="+")
|
||||
h = sub.add_parser("ssh"); h.add_argument("id")
|
||||
a = ap.parse_args()
|
||||
if a.cmd == "search":
|
||||
offs = search(a.gpu, a.n, a.min_cores, a.max_dph, a.gpus, a.min_disk)
|
||||
print(json.dumps(offs) if a.json else "\n".join(fmt_offer(o) for o in offs) or "no offers")
|
||||
elif a.cmd == "rent": rent(a.offer, a.label, a.disk, a.image, a.onstart)
|
||||
elif a.cmd == "list":
|
||||
for i in instances():
|
||||
print(f"{i.get('id')} {str(i.get('label')):<22} {i.get('num_gpus')}x {i.get('gpu_name')} {i.get('actual_status')}/{i.get('cur_state')} "
|
||||
f"ssh {i.get('ssh_host')}:{i.get('ssh_port')} ${i.get('dph_total',0):.3f}/h {i.get('geolocation')} {i.get('status_msg') or ''}"[:200])
|
||||
elif a.cmd == "ssh":
|
||||
for i in instances():
|
||||
if str(i.get("id")) == a.id: print(f"ssh -i ~/.ssh/igneum-fleet -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -p {i.get('ssh_port')} root@{i.get('ssh_host')}")
|
||||
elif a.cmd == "destroy":
|
||||
for i in a.ids: destroy(i)
|
||||
elif a.cmd == "spend": spend()
|
||||
Loading…
Reference in a new issue