Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 19 additions & 2 deletions cmd/errand/gc_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -175,20 +175,37 @@ func TestGCCachePreviewReportsRunnerPolicies(t *testing.T) {
func TestGCChangesSkipsDeletedWorkspaceWithoutFailing(t *testing.T) {
t.Setenv("XDG_STATE_HOME", t.TempDir())
root := t.TempDir()
if err := os.WriteFile(filepath.Join(root, ".errandignore"), nil, 0600); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(root, "data"), []byte("data"), 0600); err != nil {
t.Fatal(err)
}
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost || strings.HasSuffix(r.URL.Path, "/snapshot/diff") {
if r.Method != http.MethodPost {
http.NotFound(w, r)
return
}
if strings.HasSuffix(r.URL.Path, "/snapshot/diff") {
var request proto.SnapshotDiffRequest
if err := json.NewDecoder(r.Body).Decode(&request); err != nil {
t.Error(err)
http.Error(w, "invalid snapshot negotiation", http.StatusBadRequest)
return
}
missing := make([]string, 0, len(request.Blobs))
for _, blob := range request.Blobs {
missing = append(missing, blob.SHA256)
}
json.NewEncoder(w).Encode(proto.SnapshotDiffResponse{Missing: missing})
return
}
io.Copy(io.Discard, r.Body)
w.WriteHeader(http.StatusCreated)
json.NewEncoder(w).Encode(proto.Workspace{ID: strings.TrimPrefix(r.URL.Path, "/v0/workspaces/")})
}))
defer server.Close()
if _, err := client.CreateWorkspace(client.RunOptions{PeerURL: server.URL, Root: root, IncludeAll: true}, "gone"); err != nil {
if _, err := client.CreateWorkspace(client.RunOptions{PeerURL: server.URL, Root: root}, "gone"); err != nil {
t.Fatal(err)
}
if err := os.RemoveAll(root); err != nil {
Expand Down
4 changes: 1 addition & 3 deletions cmd/errand/placement.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,10 +144,8 @@ func probeCandidates(ctx context.Context, candidates []config.RunCandidate, wher
}
r.info = &info
missing := q.Missing(info.Facts)
r.matched = info.Placement && len(missing) == 0
r.matched = len(missing) == 0

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Retain the admission-validation capability check

When a configured runner predates admission-time requirement validation but still reports protocol 0, it can now be selected solely because its probe-time facts match. Version differences are advisory and the protocol number was not bumped, so this occurs during ordinary staggered upgrades; the older runner then admits the job without rechecking where, allowing facts that changed after probing—or requirements measured differently for the job environment—to run on an unsuitable machine. Keep the Placement capability gate or reject runners whose exact protocol/version cannot guarantee admission validation.

Useful? React with 👍 / 👎.

switch {
case !info.Placement:
r.reason = "runner does not support requirement validation; upgrade it"
case info.Busy:
r.reason = "runner is full or unavailable"
case info.MaxJobs <= 0 || info.MaxQueued < 0 || info.RunningJobs < 0 || info.StartingJobs < 0 || info.StagingJobs < 0 || info.QueuedJobs < 0:
Expand Down
12 changes: 6 additions & 6 deletions cmd/errand/placement_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,20 +31,20 @@ func TestCLIWhereFallsBackAfterCapacityRace(t *testing.T) {
var rejections, admissions atomic.Int32
first := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/v0/info" {
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, Placement: true, MaxJobs: 2})
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, MaxJobs: 2})
return
}
if r.Method == http.MethodPut {
rejections.Add(1)
http.Error(w, "runner filled after probe", 429)
return
}
http.NotFound(w, r)
d.Handler().ServeHTTP(w, r)
}))
defer first.Close()
second := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/v0/info" {
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, Placement: true, MaxJobs: 2, RunningJobs: 1})
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, MaxJobs: 2, RunningJobs: 1})
return
}
if r.Method == http.MethodPut {
Expand Down Expand Up @@ -87,15 +87,15 @@ func TestWorkspaceWhereCreationReportsSelectedPeer(t *testing.T) {
defer d.Close()
first := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/v0/info" {
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, Placement: true, MaxJobs: 2})
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, MaxJobs: 2})
return
}
http.Error(w, "requirements changed", 412)
}))
defer first.Close()
second := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/v0/info" {
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, Placement: true, MaxJobs: 2, RunningJobs: 1})
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, MaxJobs: 2, RunningJobs: 1})
return
}
d.Handler().ServeHTTP(w, r)
Expand Down Expand Up @@ -128,7 +128,7 @@ func TestWorkspaceWhereCreationReportsSelectedPeer(t *testing.T) {
func TestWorkspaceWhereFinalRejectionPrintedOnce(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/v0/info" {
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, Placement: true, MaxJobs: 2})
json.NewEncoder(w).Encode(proto.Info{Proto: proto.ProtoVersion, Version: version, MaxJobs: 2})
return
}
http.Error(w, "requirements changed", 412)
Expand Down
18 changes: 8 additions & 10 deletions cmd/errand/placement_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,18 +16,16 @@ import (

func TestWhereSelectionFiltersAndBalances(t *testing.T) {
e := config.EffectiveRun{Where: "os=linux,go"}
for _, name := range []string{"offline", "wrong-os", "unsupported", "busy", "loaded", "available"} {
for _, name := range []string{"offline", "wrong-os", "busy", "loaded", "available"} {
e.Candidates = append(e.Candidates, config.RunCandidate{Name: name, URL: name})
}
selection, err := chooseRunners(context.Background(), e, func(_ context.Context, url, where string, _ time.Duration) (proto.Info, error) {
i := proto.Info{Placement: true, MaxJobs: 4, Facts: proto.Facts{OS: "linux", Tools: map[string]string{"go": "/bin/go"}}}
i := proto.Info{MaxJobs: 4, Facts: proto.Facts{OS: "linux", Tools: map[string]string{"go": "/bin/go"}}}
switch url {
case "offline":
return i, fmt.Errorf("unreachable")
case "wrong-os":
i.Facts.OS = "darwin"
case "unsupported":
i.Placement = false
case "busy":
i.Busy = true
case "loaded":
Expand All @@ -44,7 +42,7 @@ func TestWhereSelectionFiltersAndBalances(t *testing.T) {
}
var errOut bytes.Buffer
selection.printExcluded(&errOut)
for _, s := range []string{"offline", "wrong-os", "unsupported", "busy"} {
for _, s := range []string{"offline", "wrong-os", "busy"} {
if !strings.Contains(errOut.String(), s) {
t.Fatalf("missing exclusion %s: %s", s, &errOut)
}
Expand Down Expand Up @@ -113,7 +111,7 @@ func TestDoctorWhereReportsChoiceWithoutSubmitting(t *testing.T) {
if target != "http://runner.invalid" || where != "go" {
t.Errorf("probe %s %s", target, where)
}
return proto.Info{Placement: true, Version: version, MaxJobs: 1, Facts: proto.Facts{Tools: map[string]string{"go": "/bin/go"}}}, nil
return proto.Info{Version: version, MaxJobs: 1, Facts: proto.Facts{Tools: map[string]string{"go": "/bin/go"}}}, nil
}, probe: func(context.Context, string) (proto.Info, error) {
t.Fatal("redundant probe")
return proto.Info{}, nil
Expand All @@ -131,7 +129,7 @@ func TestDoctorWhereKeepsSSHDisplayURL(t *testing.T) {
if target == "ssh://host" {
t.Error("custom SSH transport was not registered")
}
return proto.Info{Placement: true, Version: version, MaxJobs: 1}, nil
return proto.Info{Version: version, MaxJobs: 1}, nil
}, ssh: func(context.Context, string) error {
t.Fatal("repeated SSH inspection after successful info probe")
return nil
Expand All @@ -145,7 +143,7 @@ func TestDoctorWhereKeepsSSHDisplayURL(t *testing.T) {
// so only your own config, a profile you chose or --where may lease.
func TestWorkspaceWhereNeverLeases(t *testing.T) {
broker := func(_ context.Context, _, _ string, _ time.Duration) (proto.Info, error) {
return proto.Info{Placement: true, MaxJobs: 1, Facts: proto.Facts{OS: "linux"}, Offers: []proto.Offer{{Name: "h100", Facts: proto.Facts{OS: "linux", GPUs: []proto.GPU{{Name: "H100", MemoryMiB: 81920}}}}}}, nil
return proto.Info{MaxJobs: 1, Facts: proto.Facts{OS: "linux"}, Offers: []proto.Offer{{Name: "h100", Facts: proto.Facts{OS: "linux", GPUs: []proto.GPU{{Name: "H100", MemoryMiB: 81920}}}}}}, nil
}
e := config.EffectiveRun{Where: "gpu", Candidates: []config.RunCandidate{{Name: "cloud", URL: "cloud"}}}
selection, err := chooseRunners(context.Background(), e, broker)
Expand All @@ -168,7 +166,7 @@ func TestReadyLeasesAreReusedThroughTheirCloudPeer(t *testing.T) {
var probed []string
probe := func(_ context.Context, target, _ string, _ time.Duration) (proto.Info, error) {
probed = append(probed, target)
info := proto.Info{Placement: true, MaxJobs: 1, Facts: proto.Facts{OS: "linux"}}
info := proto.Info{MaxJobs: 1, Facts: proto.Facts{OS: "linux"}}
if target == "http://cloud:7443" {
info.Offers = []proto.Offer{{Name: "h100", Facts: h100}}
info.Leases = []proto.Lease{{ID: proto.NewULID(), Offer: "h100", State: proto.LeaseReady, Target: &proto.LeaseTarget{URL: "http://box:7443"}, Facts: &h100}}
Expand All @@ -193,7 +191,7 @@ func TestLeaseSuppliersRankedByOffer(t *testing.T) {
return proto.Info{}, errors.New("unreachable")
}
busy := map[string]int{"cabal": 1}[target]
return proto.Info{Placement: true, MaxJobs: 1, RunningJobs: busy, Facts: proto.Facts{OS: "linux"}, Offers: offers[target]}, nil
return proto.Info{MaxJobs: 1, RunningJobs: busy, Facts: proto.Facts{OS: "linux"}, Offers: offers[target]}, nil
}
order := func(peers ...string) []string {
e := config.EffectiveRun{Where: "gpu=h100", WhereMayLease: true}
Expand Down
26 changes: 13 additions & 13 deletions cmd/errand/push_batch_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,18 +15,18 @@ import (
"github.com/lydakis/errand/internal/proto"
)

func TestPushBatchUsesDeltaAndFailsClosedAfterRollback(t *testing.T) {
var rollback atomic.Bool
var legacyUploads, negotiations atomic.Int32
func TestPushBatchUsesDeltaAndDoesNotRetryMissingEndpoint(t *testing.T) {
var missingEndpoint atomic.Bool
var uploads, negotiations atomic.Int32
root, destination, peer, ws := watchFixtureHandler(t, func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/push") {
legacyUploads.Add(1)
uploads.Add(1)
}
if strings.HasSuffix(r.URL.Path, "/push/diff") {
negotiations.Add(1)
}
if strings.HasSuffix(r.URL.Path, "/push/delta-v1") && rollback.Load() {
if strings.HasSuffix(r.URL.Path, "/push") && missingEndpoint.Load() {
http.NotFound(w, r)
return
}
Expand All @@ -49,24 +49,24 @@ func TestPushBatchUsesDeltaAndFailsClosedAfterRollback(t *testing.T) {
if _, err := client.PushChanges(opts); err != nil {
t.Fatal(err)
}
if stats.TransferredBytes > 16<<10 || negotiations.Load() != 0 || legacyUploads.Load() != 0 {
t.Fatalf("small batch failed delta path: %+v negotiations=%d legacy=%d", stats, negotiations.Load(), legacyUploads.Load())
if stats.TransferredBytes > 16<<10 || negotiations.Load() != 0 || uploads.Load() != 2 {
t.Fatalf("small batch failed delta path: %+v negotiations=%d uploads=%d", stats, negotiations.Load(), uploads.Load())
}
rollback.Store(true)
missingEndpoint.Store(true)
if err := os.Remove(filepath.Join(root, "value")); err != nil {
t.Fatal(err)
}
if _, err := client.PushChanges(opts); err == nil {
t.Fatal("rollback unexpectedly accepted delta")
t.Fatal("missing endpoint unexpectedly accepted delta")
}
if _, err := os.Stat(filepath.Join(destination, "keep")); err != nil {
t.Fatalf("unrelated file lost: %v", err)
}
if _, err := os.Stat(filepath.Join(destination, "value")); err != nil {
t.Fatalf("rollback applied deletion: %v", err)
t.Fatalf("missing endpoint applied deletion: %v", err)
}
if legacyUploads.Load() != 0 {
t.Fatal("delta fell back to unsafe legacy endpoint")
if uploads.Load() != 3 {
t.Fatal("push retried a missing endpoint")
}
}

Expand All @@ -80,7 +80,7 @@ func TestPushBatchRebuildsOnlyTypedUnstagedCheckpointRejection(t *testing.T) {
if strings.HasSuffix(r.URL.Path, "/push/base") {
bases.Add(1)
}
if strings.HasSuffix(r.URL.Path, "/push/delta-v1") {
if strings.HasSuffix(r.URL.Path, "/push") {
uploads.Add(1)
if !rejected.Swap(true) {
w.WriteHeader(http.StatusConflict)
Expand Down
2 changes: 1 addition & 1 deletion cmd/errand/push_benchmark_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,7 @@ func benchmarkPushWorkload(b *testing.B, watch bool, workload watchWorkload) {
switch {
case strings.HasSuffix(r.URL.Path, "/push/diff"):
negotiation.Add(int64(time.Since(started)))
case (strings.HasSuffix(r.URL.Path, "/push") || strings.HasSuffix(r.URL.Path, "/push/delta-v1")):
case strings.HasSuffix(r.URL.Path, "/push"):
stage.Add(int64(time.Since(started)))
case strings.HasSuffix(r.URL.Path, "/apply"):
apply.Add(int64(time.Since(started)))
Expand Down
8 changes: 6 additions & 2 deletions cmd/errand/push_record_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ func TestPushDeltaRecordOmitsInventoryAndResumes(t *testing.T) {
t.Fatal(err)
}
var record struct{ Request proto.PushRequest }
var fields struct{ Request map[string]json.RawMessage }
found := false
err = filepath.WalkDir(os.Getenv("XDG_STATE_HOME"), func(name string, entry fs.DirEntry, err error) error {
if err != nil || entry.Name() != "push.json" {
Expand All @@ -34,13 +35,16 @@ func TestPushDeltaRecordOmitsInventoryAndResumes(t *testing.T) {
if err != nil {
return err
}
if err := json.Unmarshal(raw, &fields); err != nil {
return err
}
return json.Unmarshal(raw, &record)
})
if err != nil || !found || record.Request.Delta == nil || record.Request.SourceRoot == "" {
t.Fatalf("missing recoverable delta record: found=%v err=%v", found, err)
}
if len(record.Request.Manifest.Entries) != 0 {
t.Fatalf("delta recovery record redundantly contains %d full-inventory entries", len(record.Request.Manifest.Entries))
if _, exists := fields.Request["manifest"]; exists {
t.Fatal("delta recovery record redundantly contains the full inventory")
}
opts.Apply = true
applied, err := client.PushChanges(opts)
Expand Down
Loading
Loading