Files
cloud/cli/gpu_queue_test.go
antje 25d296a6f5
Hanzo CI/CD / cicd (push) Successful in 20s
CI/CD / gate (push) Successful in 20s
CI/CD / containment (push) Successful in 1m40s
CI/CD / image (push) Skipped
CI/CD / rollout (push) Skipped
CI/CD / reach (push) Skipped
CI/CD / fanout (push) Skipped
CI/CD / receipt (push) Skipped
fix: a refused credential is not a blip — stop retrying it
The register retry added earlier exists so a rolling control plane cannot kill a
worker mid-render: api.hanzo.ai answers 503 for a few seconds while its pod is
replaced, and riding that out saves whatever was sampling.

401/403 is the opposite kind of failure. It says this node's token is not accepted,
and waiting never changes that. Retried, it cost 30s per boot inside a systemd restart
loop — found at restart counter 10 on a node whose credential had expired — and buried
the one line naming the cause under five that said "retrying".

It now fails on the first refusal and says what to do: run `hanzo login` on that node.
Same lesson as the held spool, one layer over: a permanent condition wearing a
retryable shape is worse than an error, because it looks like progress.
2026-08-01 18:26:35 -07:00

453 lines
17 KiB
Go

package cli
// gpu_queue_test.go — the per-GPU claim contract: a job pinned to THIS machine's
// lane ("gpu:<identity>") is claimed BEFORE the shared any-GPU lane ("gpu-jobs"),
// both within the gpu-jobs namespace. A stub cloud records the taskQueue of every
// claim so the ORDER is the assertion.
import (
"context"
"encoding/json"
"io"
"net"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
)
// stubJobsCloud records the taskQueue of every claim and serves the queued job (if
// any) on the matching lane, once. A lane with no job (or already drained) answers
// 204; complete/fail/heartbeat answer 200.
func stubJobsCloud(t *testing.T, claims *[]string, jobsByLane map[string]*claimedActivity) *httptest.Server {
t.Helper()
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/activities/claim") {
var body struct {
TaskQueue string `json:"taskQueue"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
*claims = append(*claims, body.TaskQueue)
if job := jobsByLane[body.TaskQueue]; job != nil {
jobsByLane[body.TaskQueue] = nil // deliver once
_ = json.NewEncoder(w).Encode(job)
return
}
w.WriteHeader(http.StatusNoContent)
return
}
w.WriteHeader(http.StatusOK) // complete / fail / heartbeat
}))
}
// echoJob is a GPU-free job the worker's echo handler runs to completion instantly,
// so a claim test never needs a real render backend.
func echoJob(id string) *claimedActivity {
a := &claimedActivity{Input: json.RawMessage(`{}`)}
a.Execution.WorkflowId, a.Execution.RunId = id, id
a.Type.Name = "echo"
return a
}
func testWorker(t *testing.T, url string) *worker {
t.Helper()
t.Setenv("HANZO_TOKEN", "test-token") // ensureToken honors this; no `hanzo login`
return &worker{
env: &Env{CloudURL: url},
http: &http.Client{Timeout: 5 * time.Second},
baseURL: url,
identity: "spark",
hostname: "spark",
jobsNS: "gpu-jobs",
handlers: map[string]jobHandler{"echo": echoHandler},
studioReady: true, // a render-capable node (preflight passed)
}
}
func TestGPUQueueLaneName(t *testing.T) {
w := &worker{identity: "spark"}
if got := w.gpuQueue(); got != "gpu:spark" {
t.Fatalf("gpuQueue() = %q, want gpu:spark", got)
}
}
// A job pinned to THIS GPU's lane is claimed first — the shared lane is not even
// polled that cycle.
func TestClaimPrefersTargetedLane(t *testing.T) {
var claims []string
srv := stubJobsCloud(t, &claims, map[string]*claimedActivity{"gpu:spark": echoJob("j1")})
defer srv.Close()
if err := testWorker(t, srv.URL).claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("claimAndRun: %v", err)
}
if len(claims) != 1 || claims[0] != "gpu:spark" {
t.Fatalf("claims = %v, want exactly [gpu:spark] (targeted first; shared not polled)", claims)
}
}
// With its own lane empty, the worker falls through to the shared any-GPU lane.
func TestClaimFallsBackToSharedLane(t *testing.T) {
var claims []string
srv := stubJobsCloud(t, &claims, map[string]*claimedActivity{"gpu-jobs": echoJob("j2")})
defer srv.Close()
if err := testWorker(t, srv.URL).claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("claimAndRun: %v", err)
}
if len(claims) != 2 || claims[0] != "gpu:spark" || claims[1] != "gpu-jobs" {
t.Fatalf("claims = %v, want [gpu:spark gpu-jobs] (targeted, then shared)", claims)
}
}
// Both lanes empty polls targeted THEN shared and runs nothing.
func TestClaimBothLanesEmpty(t *testing.T) {
var claims []string
srv := stubJobsCloud(t, &claims, map[string]*claimedActivity{})
defer srv.Close()
if err := testWorker(t, srv.URL).claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("claimAndRun: %v", err)
}
if len(claims) != 2 || claims[0] != "gpu:spark" || claims[1] != "gpu-jobs" {
t.Fatalf("claims = %v, want [gpu:spark gpu-jobs]", claims)
}
}
// The render submit seam is the gated worker-mode execute path, not the open /prompt
// — the shared contract with the studio's --worker-mode gate.
func TestWorkerExecuteSeamIsGated(t *testing.T) {
if localWorkerExecute != "http://127.0.0.1:8188/v1/worker/execute" {
t.Fatalf("localWorkerExecute = %q, want the gated /v1/worker/execute seam", localWorkerExecute)
}
}
// A HUNG nvidia-smi (blocks until its bounded context fires) must NOT block the
// caller: reportSample detaches the probe+POST onto a goroutine and returns at once,
// so the worker's select loop keeps heartbeating and claiming. Guards the
// worker-wedge regression (a synchronous probe that stalls the loop under GPU/driver
// pressure → the machine flaps offline mid-render).
func TestReportSampleNeverBlocksLoop(t *testing.T) {
orig := nvidiaSmi
defer func() { nvidiaSmi = orig }()
probing := make(chan struct{})
release := make(chan struct{})
nvidiaSmi = func(ctx context.Context) ([]byte, error) {
close(probing) // entered the probe (the nvidiaSmi var read already happened)
select {
case <-release: // the test lets us finish
case <-ctx.Done(): // or the bounded probe timeout fires
}
return nil, ctx.Err()
}
t.Setenv("HANZO_TOKEN", "t")
w := &worker{identity: "spark", hostname: "spark", http: &http.Client{Timeout: time.Second}, env: &Env{}, baseURL: "http://127.0.0.1:0"}
done := make(chan struct{})
go func() { w.reportSample(context.Background()); close(done) }()
select {
case <-done: // returned immediately — the select loop is never wedged
case <-time.After(500 * time.Millisecond):
t.Fatal("reportSample blocked the caller — a hung sampler would wedge the worker loop")
}
select {
case <-probing: // the probe really ran, on the detached goroutine (off the critical path)
case <-time.After(2 * time.Second):
t.Fatal("probe never started")
}
close(release)
}
// A node that can't serve renders (preflight failed) claims NOTHING — it must never
// pull a render job onto a box that will only refuse it on the gated seam (poison
// loop). It still heartbeats presence; it just stays idle.
func TestNotStudioReadyClaimsNothing(t *testing.T) {
var claims []string
srv := stubJobsCloud(t, &claims, map[string]*claimedActivity{"gpu:spark": echoJob("j1")})
defer srv.Close()
w := testWorker(t, srv.URL)
w.studioReady = false
if err := w.claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("claimAndRun: %v", err)
}
if len(claims) != 0 {
t.Fatalf("a not-ready node claimed %v; want zero claims", claims)
}
}
// Declining a render this node cannot serve must report NOTHING — the claim's lease
// lapses and the engine returns the job to pending for a worker that can run it.
//
// It used to POST `fail`, which is terminal: the job was not handed to anyone, it was
// destroyed. Two workers that had both just restarted each claimed one render and each
// failed it, so the person who asked for that render got nothing back. `fail` is the
// only verb here that can end a job, and a decline is precisely the case where it must
// not be used.
func TestDecliningARenderReturnsItInsteadOfFailingIt(t *testing.T) {
var claims []string
var reported []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case strings.HasSuffix(r.URL.Path, "/activities/claim"):
var body struct {
TaskQueue string `json:"taskQueue"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
claims = append(claims, body.TaskQueue)
if body.TaskQueue == "gpu:spark" {
a := &claimedActivity{Input: json.RawMessage(`{}`)}
a.Execution.WorkflowId, a.Execution.RunId = "wf1", "wf1"
a.Type.Name = studioCap // a RENDER, on a node whose studio is down
_ = json.NewEncoder(w).Encode(a)
return
}
w.WriteHeader(http.StatusNoContent)
default:
reported = append(reported, r.URL.Path) // complete / fail / heartbeat
w.WriteHeader(http.StatusOK)
}
}))
defer srv.Close()
// A node with a non-render lane keeps claiming even when its studio is down, so
// it is the one that meets a render it cannot serve.
w := testWorker(t, srv.URL)
w.studioReady = false
w.handlers["fn.run"] = echoHandler
if err := w.claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("claimAndRun: %v", err)
}
if len(reported) != 0 {
t.Fatalf("a declined render was reported terminal via %v; it must be left to the lease", reported)
}
}
// The terminal report must hit the RIGHT activity — namespace gpu-jobs, the CLAIMED
// workflow+run ids, the correct verb. A stub that 200s every path lets an ns/id
// routing regression pass, so assert the exact paths for both complete and fail.
func TestTerminalReportsHitCorrectActivityPath(t *testing.T) {
t.Setenv("HANZO_TOKEN", "t")
var mu sync.Mutex
terminal := map[string]string{}
served := 0
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case strings.HasSuffix(r.URL.Path, "/activities/claim"):
var body struct {
TaskQueue string `json:"taskQueue"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
if body.TaskQueue != "gpu:spark" {
w.WriteHeader(http.StatusNoContent)
return
}
mu.Lock()
n := served
served++
mu.Unlock()
switch n {
case 0:
_ = json.NewEncoder(w).Encode(echoJob("wf1")) // echo handler → complete
case 1:
j := echoJob("wf2")
j.Type.Name = "nope" // no handler → fail
_ = json.NewEncoder(w).Encode(j)
default:
w.WriteHeader(http.StatusNoContent)
}
case strings.HasSuffix(r.URL.Path, "/complete"):
mu.Lock()
terminal["complete"] = r.URL.Path
mu.Unlock()
w.WriteHeader(http.StatusOK)
case strings.HasSuffix(r.URL.Path, "/fail"):
mu.Lock()
terminal["fail"] = r.URL.Path
mu.Unlock()
w.WriteHeader(http.StatusOK)
default:
w.WriteHeader(http.StatusOK)
}
}))
defer srv.Close()
w := testWorker(t, srv.URL)
if err := w.claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("run1 (echo→complete): %v", err)
}
if err := w.claimAndRun(context.Background(), io.Discard); err != nil {
t.Fatalf("run2 (unknown→fail): %v", err)
}
if got := terminal["complete"]; got != "/v1/tasks/namespaces/gpu-jobs/activities/wf1/wf1/complete" {
t.Fatalf("complete path = %q, want the claimed activity's exact ns+ids", got)
}
if got := terminal["fail"]; got != "/v1/tasks/namespaces/gpu-jobs/activities/wf2/wf2/fail" {
t.Fatalf("fail path = %q, want the claimed activity's exact ns+ids", got)
}
}
// SharePolicy.reject's fallback: an ABSENT or unparseable input field skips its gate
// (permissive), never a hard error, so a policy only ever narrows on fields it can read.
func TestSharePolicyRejectFallback(t *testing.T) {
var nilp *SharePolicy
if r := nilp.reject("studio.render", nil); r != "" {
t.Fatalf("nil policy must allow everything: %q", r)
}
p := &SharePolicy{AllowedJobTypes: []string{"studio.render"}}
if p.reject("echo", nil) == "" {
t.Fatal("a disallowed job type must be rejected")
}
if r := p.reject("studio.render", nil); r != "" {
t.Fatalf("an allowed job type must pass: %q", r)
}
p2 := &SharePolicy{AllowedOrgs: []string{"acme"}}
if r := p2.reject("studio.render", json.RawMessage(`{}`)); r != "" {
t.Fatalf("absent org must SKIP the org gate (fallback), not reject: %q", r)
}
if r := p2.reject("studio.render", json.RawMessage(`{"org":"acme"}`)); r != "" {
t.Fatalf("a matching org must pass: %q", r)
}
if p2.reject("studio.render", json.RawMessage(`{"org":"other"}`)) == "" {
t.Fatal("a non-allowed org must be rejected")
}
if r := p2.reject("studio.render", json.RawMessage(`not json`)); r != "" {
t.Fatalf("unparseable input must skip input gates, not reject: %q", r)
}
}
// studioCap is advertised ONLY when the node can actually render: a missing worker
// token is never ready (the gated seam would 403), a studio that does not answer is
// never ready EVEN ON A NODE THAT LAUNCHES ITS OWN, and the block reason is explicit.
//
// That middle clause is the regression this pins. Launching a studio was once taken
// as proof of readiness, so a --studio-dir node advertised studio.render and claimed
// render jobs from its first instant — before it had started the studio at all, and
// long before that studio bound its port and loaded models. It won every job it could
// reach while cold and failed each one on arrival.
func TestStudioReadyGatesCapability(t *testing.T) {
// The probe is against a fixed loopback address, so a real studio already
// serving there would make "not reachable" untestable in this process.
ln, err := net.Listen("tcp", strings.TrimPrefix(localComfyUI, "http://"))
if err != nil {
t.Skipf("%s is already in use; this test owns that address: %v", localComfyUI, err)
}
_ = ln.Close()
t.Setenv("STUDIO_WORKER_TOKEN", "")
w := &worker{launchesStudio: true, http: &http.Client{}}
w.refreshStudioReady(context.Background())
if w.studioReady {
t.Fatal("no token must never be studio-ready")
}
if contains(w.capabilities(), studioCap) {
t.Fatalf("studioCap advertised without a token: %v", w.capabilities())
}
if w.studioBlockReason() == "" {
t.Fatal("a not-ready node must explain why it won't render")
}
// A token alone is not readiness: we launch the studio, but nothing answers yet.
t.Setenv("STUDIO_WORKER_TOKEN", "tok")
w.refreshStudioReady(context.Background())
if w.studioReady {
t.Fatal("a node that launches its own studio claimed readiness before that studio answered")
}
if contains(w.capabilities(), studioCap) {
t.Fatalf("studioCap advertised by a cold node: %v", w.capabilities())
}
// Once the studio answers, readiness flips and the capability is advertised —
// the same recovery path a studio that died and came back travels.
srv := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = io.WriteString(w, `{"queue_running":[],"queue_pending":[]}`)
})}
up, err := net.Listen("tcp", strings.TrimPrefix(localComfyUI, "http://"))
if err != nil {
t.Fatalf("bind %s: %v", localComfyUI, err)
}
go func() { _ = srv.Serve(up) }()
defer func() { _ = srv.Close() }()
if changed := w.refreshStudioReady(context.Background()); !changed {
t.Fatal("a studio that started answering should flip readiness")
}
if !w.studioReady || !contains(w.capabilities(), studioCap) {
t.Fatalf("a reachable studio must be ready + advertise studioCap: ready=%v caps=%v", w.studioReady, w.capabilities())
}
if w.studioBlockReason() != "" {
t.Fatalf("a ready node must give no block reason: %q", w.studioBlockReason())
}
}
// A control plane that is ROLLING must not kill this worker. The presence write is
// the call whose error ends the process, and systemd's restart kills the studio the
// worker supervises — taking any render that was sampling with it. A cloud deploy did
// exactly that: api.hanzo.ai answered 503 for the seconds its pod was replaced, and a
// compose that had run for minutes was gone.
//
// The retry is what separates a blip from an outage, so it is asserted on the shape
// that bit: fail, fail, then succeed — one call's worth of error must not be terminal.
func TestRegisterRidesOutARollingControlPlane(t *testing.T) {
t.Setenv("HANZO_TOKEN", "t")
var mu sync.Mutex
presence := 0
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !strings.HasSuffix(r.URL.Path, "/fleet/activities") {
w.WriteHeader(http.StatusOK) // namespace ensure — already best-effort
return
}
mu.Lock()
presence++
n := presence
mu.Unlock()
if n < 3 {
http.Error(w, "no available server", http.StatusServiceUnavailable)
return
}
w.WriteHeader(http.StatusOK)
}))
defer srv.Close()
w := testWorker(t, srv.URL)
if err := w.register(context.Background()); err != nil {
t.Fatalf("register gave up on a rolling control plane: %v", err)
}
if presence != 3 {
t.Fatalf("presence attempts = %d, want 3 (two 503s ridden out, then success)", presence)
}
}
// A credential the cloud REFUSES is not a blip, and must not be retried.
//
// The register retry exists so a rolling control plane (503 for a few seconds) cannot
// kill a worker mid-render. 401/403 is the opposite kind of failure: it says this
// node's token is not accepted, which waiting never fixes. Retried, it burned 30s per
// boot inside a systemd restart loop and buried the one useful line under five
// "retrying" ones.
func TestRegisterDoesNotRetryARefusedCredential(t *testing.T) {
t.Setenv("HANZO_TOKEN", "t")
var mu sync.Mutex
presence := 0
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !strings.HasSuffix(r.URL.Path, "/fleet/activities") {
w.WriteHeader(http.StatusOK)
return
}
mu.Lock()
presence++
mu.Unlock()
http.Error(w, "identity required", http.StatusForbidden)
}))
defer srv.Close()
w := testWorker(t, srv.URL)
err := w.register(context.Background())
if err == nil {
t.Fatal("a refused credential must be reported, not swallowed")
}
if presence != 1 {
t.Fatalf("presence attempts = %d, want 1 (a refusal is never retried)", presence)
}
if !strings.Contains(err.Error(), "hanzo login") {
t.Fatalf("the error must say what to do about it: %v", err)
}
}