diff --git a/cmd/errand/gc_test.go b/cmd/errand/gc_test.go index 0d47672..22f8b4e 100644 --- a/cmd/errand/gc_test.go +++ b/cmd/errand/gc_test.go @@ -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 { diff --git a/cmd/errand/placement.go b/cmd/errand/placement.go index f1e87ff..f0c3c54 100644 --- a/cmd/errand/placement.go +++ b/cmd/errand/placement.go @@ -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 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: diff --git a/cmd/errand/placement_integration_test.go b/cmd/errand/placement_integration_test.go index ca2d175..48efd76 100644 --- a/cmd/errand/placement_integration_test.go +++ b/cmd/errand/placement_integration_test.go @@ -31,7 +31,7 @@ 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 { @@ -39,12 +39,12 @@ func TestCLIWhereFallsBackAfterCapacityRace(t *testing.T) { 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 { @@ -87,7 +87,7 @@ 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) @@ -95,7 +95,7 @@ func TestWorkspaceWhereCreationReportsSelectedPeer(t *testing.T) { 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) @@ -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) diff --git a/cmd/errand/placement_test.go b/cmd/errand/placement_test.go index b432240..fab6b09 100644 --- a/cmd/errand/placement_test.go +++ b/cmd/errand/placement_test.go @@ -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": @@ -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) } @@ -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 @@ -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 @@ -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) @@ -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}} @@ -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} diff --git a/cmd/errand/push_batch_test.go b/cmd/errand/push_batch_test.go index 30994f9..3c1cbdd 100644 --- a/cmd/errand/push_batch_test.go +++ b/cmd/errand/push_batch_test.go @@ -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 } @@ -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") } } @@ -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) diff --git a/cmd/errand/push_benchmark_test.go b/cmd/errand/push_benchmark_test.go index 08c44ee..bcf28b2 100644 --- a/cmd/errand/push_benchmark_test.go +++ b/cmd/errand/push_benchmark_test.go @@ -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))) diff --git a/cmd/errand/push_record_test.go b/cmd/errand/push_record_test.go index f2f982b..852f57b 100644 --- a/cmd/errand/push_record_test.go +++ b/cmd/errand/push_record_test.go @@ -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" { @@ -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) diff --git a/cmd/errand/push_watch_delta_test.go b/cmd/errand/push_watch_delta_test.go index 97d7a61..7a8a954 100644 --- a/cmd/errand/push_watch_delta_test.go +++ b/cmd/errand/push_watch_delta_test.go @@ -14,69 +14,61 @@ import ( "github.com/lydakis/errand/internal/client" ) -func TestPushWatchStructuralEditsAndOldRunnerFallback(t *testing.T) { - for _, legacy := range []bool{false, true} { - t.Run(fmt.Sprintf("legacy=%v", legacy), func(t *testing.T) { - var 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 legacy && strings.HasSuffix(r.URL.Path, "/push/base") { - http.NotFound(w, r) - return - } - if strings.HasSuffix(r.URL.Path, "/push/diff") { - negotiations.Add(1) - } - next.ServeHTTP(w, r) - }) - }) - ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() - step := 0 - err := client.WatchPush(ctx, client.PushOptions{PeerURL: peer, Workspace: ws.Name, Root: root, Apply: true}, func(event client.PushWatchEvent) error { - if event.State != "watching" { - return nil - } - switch step { - case 0: - if err := os.MkdirAll(filepath.Join(root, "new", "nested"), 0700); err != nil { - return err - } - if err := os.WriteFile(filepath.Join(root, "new", "nested", "file"), []byte("created\n"), 0600); err != nil { - return err - } - case 1: - if body, err := os.ReadFile(filepath.Join(destination, "new", "nested", "file")); err != nil || string(body) != "created\n" { - return fmt.Errorf("new directory not delivered: %q %v", body, err) - } - if err := os.Rename(filepath.Join(root, "new"), filepath.Join(root, "renamed")); err != nil { - return err - } - if err := os.Remove(filepath.Join(root, "value")); err != nil { - return err - } - if err := os.Symlink("renamed/nested/file", filepath.Join(root, "value")); err != nil { - return err - } - case 2: - if _, err := os.Stat(filepath.Join(destination, "new")); !os.IsNotExist(err) { - return fmt.Errorf("renamed directory remains: %v", err) - } - if link, err := os.Readlink(filepath.Join(destination, "value")); err != nil || link != "renamed/nested/file" { - return fmt.Errorf("file replacement not delivered: %q %v", link, err) - } - cancel() - } - step++ - return nil - }) - if err != nil || step != 3 { - t.Fatalf("step=%d: %v", step, err) - } - if !legacy && negotiations.Load() != 0 { - t.Fatal("small deltas incurred redundant cache negotiations") +func TestPushWatchStructuralEdits(t *testing.T) { + var 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/diff") { + negotiations.Add(1) } + next.ServeHTTP(w, r) }) + }) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + step := 0 + err := client.WatchPush(ctx, client.PushOptions{PeerURL: peer, Workspace: ws.Name, Root: root, Apply: true}, func(event client.PushWatchEvent) error { + if event.State != "watching" { + return nil + } + switch step { + case 0: + if err := os.MkdirAll(filepath.Join(root, "new", "nested"), 0700); err != nil { + return err + } + if err := os.WriteFile(filepath.Join(root, "new", "nested", "file"), []byte("created\n"), 0600); err != nil { + return err + } + case 1: + if body, err := os.ReadFile(filepath.Join(destination, "new", "nested", "file")); err != nil || string(body) != "created\n" { + return fmt.Errorf("new directory not delivered: %q %v", body, err) + } + if err := os.Rename(filepath.Join(root, "new"), filepath.Join(root, "renamed")); err != nil { + return err + } + if err := os.Remove(filepath.Join(root, "value")); err != nil { + return err + } + if err := os.Symlink("renamed/nested/file", filepath.Join(root, "value")); err != nil { + return err + } + case 2: + if _, err := os.Stat(filepath.Join(destination, "new")); !os.IsNotExist(err) { + return fmt.Errorf("renamed directory remains: %v", err) + } + if link, err := os.Readlink(filepath.Join(destination, "value")); err != nil || link != "renamed/nested/file" { + return fmt.Errorf("file replacement not delivered: %q %v", link, err) + } + cancel() + } + step++ + return nil + }) + if err != nil || step != 3 { + t.Fatalf("step=%d: %v", step, err) + } + if negotiations.Load() != 0 { + t.Fatal("small deltas incurred redundant cache negotiations") } } @@ -84,7 +76,7 @@ func TestPushWatchCancelDrainsUnchangedSelectedPath(t *testing.T) { entered, release := make(chan struct{}, 1), make(chan struct{}) root, _, 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") || strings.HasSuffix(r.URL.Path, "/push/delta-v1") { + if strings.HasSuffix(r.URL.Path, "/push") { entered <- struct{}{} <-release } diff --git a/cmd/errand/transfer_reporting_test.go b/cmd/errand/transfer_reporting_test.go index c6fb2d9..965bf89 100644 --- a/cmd/errand/transfer_reporting_test.go +++ b/cmd/errand/transfer_reporting_test.go @@ -28,7 +28,7 @@ func TestTransferReportingRoundTrip(t *testing.T) { t.Cleanup(func() { _ = d.Close() }) var uploaded, downloaded atomic.Int64 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if strings.HasSuffix(r.URL.Path, "/push") || strings.HasSuffix(r.URL.Path, "/push/delta-v1") { + if strings.HasSuffix(r.URL.Path, "/push") { body, err := io.ReadAll(r.Body) if err != nil { t.Error(err) diff --git a/docs/DESIGN.md b/docs/DESIGN.md index c2a6468..e9480fb 100644 --- a/docs/DESIGN.md +++ b/docs/DESIGN.md @@ -916,12 +916,13 @@ details used by `errand status`; SSE with event IDs powers `GET /v0/jobs//logs?from=`; the signal and kill routes control owned jobs and return `204 No Content` on success; `POST /v0/snapshot/diff` negotiates missing snapshot blobs; `POST /v0/workspaces//push/diff` negotiates the same cache for an owned -workspace and establishes support for partial push archives. Push reconstructs +workspace. Push reconstructs and verifies the complete source before staging its delta; a missing cached body returns `snapshot_cache_miss` before staging, allowing a full-upload retry. -Runners without this endpoint, or with snapshot caching disabled, return 404 -and receive full archives. Both submission and push share the negotiation and -single full-upload fallback policy. Negotiation deduplicates content hashes; +All snapshot negotiation endpoints are required. A disabled cache returns +`200 OK` with all requested hashes missing, so the client sends every body. +Both submission and push share the negotiation and single full-upload retry +policy. Negotiation deduplicates content hashes; its response limit scales with the requested hashes so a valid cold manifest does not exceed an unrelated fixed response cap. Cache corruption is a miss even if deleting the bad cache entry fails; destination failures and cancellation @@ -951,9 +952,9 @@ the client can retry the same job ID with a complete snapshot. Curl-debuggable; the route prefix is the request-protocol version; receipt and change-bundle versions apply only to their persisted formats. Executable versions are diagnostic information. Peers, doctor, and setup report version differences -without blocking commands. Matching versions are recommended when diagnosing -unexpected behavior. There is no version negotiation or alternate behavior for -older formats. Updating the executable requires restarting the daemon. +without blocking commands. The CLI and daemon must run the same version. +Mixed versions are unsupported; there is no version negotiation or alternate +protocol behavior. Updating the executable requires restarting the daemon. Development builds can use an explicit version label to identify them (`go build -ldflags "-X main.version=LABEL" ./cmd/errand`). @@ -1051,21 +1052,19 @@ Cancellation drains the active transaction. A recovered apply is completed with its original identity before reconciling current files. Conflicts and policy or workspace identity changes stop the session; transient transport failures back off. -`GET /v0/workspaces//push/base?client=` negotiates incremental source -uploads and returns that owner's retained directional source checkpoint. Both -ordinary push and watch can send a `PushRequest.delta` with `source_root` instead -of the full manifest through the versioned `/push/delta-v1` endpoint. +`GET /v0/workspaces//push/base?client=` returns that owner's retained +directional source checkpoint. Ordinary push and watch send `PushRequest.delta` +with `source_root` to `/v0/workspaces//push`. The daemon reconstructs metadata from the retained checkpoint, verifies the full root and canonical delta, and checks the complete source against workspace limits. Only changed source bodies are frozen and uploaded. Small deltas (up to 64 KiB of file bodies) skip blob negotiation to save a round trip. The upload gate is released during network I/O; staging rechecks the checkpoint under the apply gate and rejects a stale baseline. Live runner contents never supply unchanged source -entries or merge bases. Local retry state retains the full manifest and frozen -delta bodies; publication and apply retain the existing durable receipts. -Runners without the base endpoint receive ordinary full-manifest push requests. -A runner that loses delta support after negotiation rejects the versioned upload; -the client never sends that partial source to the legacy full-snapshot endpoint. +entries or merge bases. Local retry state retains the delta, complete source +identity and frozen delta bodies; publication and apply retain the existing durable receipts. +Every push requires the source checkpoint endpoint and sends a delta to the +single push endpoint. Missing endpoints fail the operation. A typed checkpoint rejection before staging permits one fresh-source retry. Uncertain apply replies retain the original request identity. @@ -1074,9 +1073,9 @@ reuses them only after recovery and a fresh checkpoint check, retaining normal body verification, source quotas and durable publication. Persistent fetch uses the same plan to materialize only changed remote bodies. Workspace descriptor requests may omit the creation manifest; full reads keep their existing contract. -Creation uploads negotiate the shared verified body cache through separate -snapshot routes, with a full-upload fallback for older runners or explicit cache -misses. Long creation and job uploads use the admission timeout budget. +Creation uploads negotiate the shared verified body cache and use one creation +endpoint for both partial and complete archives. Only an explicit cache miss +permits a complete-archive retry. Long creation and job uploads use the admission timeout budget. An internal receiver-side apply receipt records each application. It binds one application to its owner, destination directory identity, immutable diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index adf9e4b..cc67b84 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -304,8 +304,8 @@ earlier errand, follows the `--older-than` boundary: no command can fetch, apply or recover through it, so keeping it protects nothing. Each record GC cannot collect is named with the reason, and the exit status is nonzero. Cache previews report the selected runner and separate snapshot/named-cache -expiry and byte budgets supplied by that runner. Older runners without policy -reporting are labeled explicitly; the client does not guess their settings. +expiry and byte budgets supplied by that runner. Policy reporting is required; +the client rejects incomplete responses. Cache collection can remove idle caches across owners; leased named caches remain protected. diff --git a/docs/SERVICE_UPGRADES.md b/docs/SERVICE_UPGRADES.md index b3cbc8f..14acaf1 100644 --- a/docs/SERVICE_UPGRADES.md +++ b/docs/SERVICE_UPGRADES.md @@ -66,7 +66,8 @@ Run, attach, and fetch emit one advisory warning to stderr when a bounded info probe reports a different daemon version. Run and attach perform their local environment/forwarding checks first. The probe adds a round trip, bounded to two seconds; failed probes do not gate the requested operation. There is no -version ordering or compatibility negotiation. +version ordering or compatibility negotiation. The CLI and daemon must run the +same version; the diagnostic warning does not provide mixed-version support. The warning identifies the invoking CLI and running daemon separately. It asks the operator to check `errand version` on the runner before choosing setup. diff --git a/docs/TRANSFER_PERFORMANCE.md b/docs/TRANSFER_PERFORMANCE.md index d007e9d..41c8df7 100644 --- a/docs/TRANSFER_PERFORMANCE.md +++ b/docs/TRANSFER_PERFORMANCE.md @@ -57,11 +57,10 @@ receipts, and watches stopped cleanly. Each three-second idle probe reported ## What changed -- Ordinary push freezes and sends delta bodies using a versioned endpoint. - Rollbacks fail closed; only definite checkpoint rejection permits one new - request. Lost apply replies recover the original request first. +- Ordinary push freezes and sends delta bodies using the single push endpoint. + Only definite checkpoint rejection permits one new request. Lost apply replies recover the original request first. - Push resolves a compact workspace descriptor instead of downloading its full - creation manifest. Older servers may still return the full descriptor. + creation manifest. - Watch refreshes hinted files under explicit `.errandignore` policies, with live policy and native directory evidence. Git-driven selection, atomic saves, structural changes and overflow use full reconciliation. Frozen bodies retain diff --git a/docs/USAGE.md b/docs/USAGE.md index c360fb1..036f910 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -299,10 +299,9 @@ workspace name cannot redirect an existing watch. Concurrent commands remain caller-managed writers; watch does not fetch remote edits or resolve application build dependencies automatically. -On a runner supporting incremental push, watch sends changed source entries and -compact metadata against the retained source checkpoint. Unchanged running files -are never used to reconstruct source. Older runners use normal full-snapshot -staging and can be slower. Local native notifications use kqueue on macOS and +Watch sends changed source entries and compact metadata against the retained +source checkpoint. Unchanged running files are never used to reconstruct source. +The CLI and daemon must run the same version. Local native notifications use kqueue on macOS and inotify on Linux. Watch setup reports exhausted OS watch/file-descriptor limits as errors; it does not silently fall back to continuous full-tree polling. @@ -334,8 +333,7 @@ Push supports `--on`, `--url`, `--profile`, `--workspace-root`, `--include-all`, `--json`, and an optional complete changed `PATH` to limit application. Staging describes the complete selected snapshot, but reuses verified file contents in the runner's snapshot cache and uploads only missing bodies. Workspace creation -and successful uploads populate that cache. A cold or disabled cache, or an -older runner, can require a full upload; eviction or corrupt cached content +and successful uploads populate that cache. A cold or disabled cache requires all changed bodies; eviction or corrupt cached content during transfer safely retries the same frozen snapshot. Cache cleanup failures do not prevent that retry. The optional `PATH` limits application, not snapshot selection. Fetch downloads the retained change bundle. Push uses the normal snapshot diff --git a/internal/client/changes_test.go b/internal/client/changes_test.go index 12dbc8b..cb5a756 100644 --- a/internal/client/changes_test.go +++ b/internal/client/changes_test.go @@ -766,7 +766,7 @@ func TestDefiniteSubmitRejectionRemovesChangeState(t *testing.T) { } server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.Method == http.MethodPost && r.URL.Path == "/v0/snapshot/diff" { - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) return } if r.Method == http.MethodPut { @@ -811,7 +811,7 @@ func TestSelectionChangeBeforeSubmitDoesNotClaimPossibleAdmission(t *testing.T) if err := os.WriteFile(filepath.Join(root, "late"), []byte("changed"), 0o600); err != nil { t.Error(err) } - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) return } if r.Method == http.MethodPut { @@ -860,7 +860,7 @@ func TestSubmitRetryConflictRetainsChangeStateAndHandle(t *testing.T) { var puts int server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.Method == http.MethodPost && r.URL.Path == "/v0/snapshot/diff" { - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) return } if r.Method == http.MethodPut { @@ -1068,7 +1068,7 @@ func TestInterruptedSnapshotNegotiationRemovesUnsubmittedChangeState(t *testing. release := make(chan struct{}) server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost || r.URL.Path != "/v0/snapshot/diff" { - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) return } close(entered) diff --git a/internal/client/client.go b/internal/client/client.go index b94ee98..066f8e5 100644 --- a/internal/client/client.go +++ b/internal/client/client.go @@ -338,8 +338,8 @@ func runPrepared(opts RunOptions, prep snapshotPreparation, env, envSources map[ } plan, negErr := negotiation.plan, negotiation.err if negErr != nil { - errf("snapshot negotiation failed (%v); shipping everything", negErr) - plan = shipPlan{} + errf("snapshot negotiation: %v", negErr) + return ExitTransaction, false } if plan.partial { shipFiles, shipBytes := 0, int64(0) diff --git a/internal/client/client_test.go b/internal/client/client_test.go index f824d54..757f8c2 100644 --- a/internal/client/client_test.go +++ b/internal/client/client_test.go @@ -103,7 +103,7 @@ func TestOutputlessRunSendsCanonicalLimits(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { case r.Method == http.MethodPost && r.URL.Path == "/v0/snapshot/diff": - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: mr, err := r.MultipartReader() if err != nil { @@ -188,7 +188,7 @@ func TestAdmissionBookkeepingFailureDoesNotAbandonAdmittedJob(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { case r.Method == http.MethodPost && r.URL.Path == "/v0/snapshot/diff": - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: _, _ = io.Copy(io.Discard, r.Body) _ = json.NewEncoder(w).Encode(proto.JobStatus{State: proto.StateRunning}) @@ -238,7 +238,7 @@ func TestAttachedApplyStartsCompletionWorkerAtAdmission(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { case r.Method == http.MethodPost && r.URL.Path == "/v0/snapshot/diff": - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: _, _ = io.Copy(io.Discard, r.Body) _ = json.NewEncoder(w).Encode(proto.JobStatus{State: proto.StateRunning}) @@ -404,6 +404,8 @@ func TestInterruptIsForwardedBeforeSubmitResponse(t *testing.T) { var controlAttempts atomic.Int32 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { + case r.URL.Path == "/v0/snapshot/diff": + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: if _, err := io.Copy(io.Discard, r.Body); err != nil { t.Errorf("reading submit body: %v", err) @@ -479,6 +481,8 @@ func TestDetachedInterruptWaitsForForwardingAndDoesNotReportSuccess(t *testing.T var controlAttempts atomic.Int32 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { + case r.URL.Path == "/v0/snapshot/diff": + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: if _, err := io.Copy(io.Discard, r.Body); err != nil { t.Errorf("reading submit body: %v", err) @@ -582,6 +586,8 @@ func TestInteractiveDetachRequestedBeforeAdmissionSkipsLogFollowing(t *testing.T var logRequests atomic.Int32 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { + case r.URL.Path == "/v0/snapshot/diff": + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: _, _ = io.Copy(io.Discard, r.Body) json.NewEncoder(w).Encode(proto.JobStatus{State: proto.StateRunning}) @@ -634,7 +640,7 @@ func TestDetachedApplyReportsWorkerLaunchFailure(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/snapshot/diff"): - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPut: _, _ = io.Copy(io.Discard, r.Body) _ = json.NewEncoder(w).Encode(proto.JobStatus{State: proto.StateRunning}) @@ -1138,7 +1144,7 @@ func TestInteractiveDetachCancelsOnlyTheLocalLogFollow(t *testing.T) { <-r.Context().Done() close(logCanceled) case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/snapshot/diff"): - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) case r.Method == http.MethodPost: controlRequests.Add(1) w.WriteHeader(http.StatusOK) diff --git a/internal/client/placement_test.go b/internal/client/placement_test.go index 0825d2d..f4d47dd 100644 --- a/internal/client/placement_test.go +++ b/internal/client/placement_test.go @@ -108,6 +108,10 @@ func TestWhereFallbackReusesSnapshotAndEnvironment(t *testing.T) { envs := make(chan string, 2) server := func(status int) *httptest.Server { return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/diff") { + replyMissingSnapshotBlobs(t, w, r) + return + } if !strings.HasPrefix(r.Header.Get("Content-Type"), "multipart/form-data") { http.NotFound(w, r) return diff --git a/internal/client/push.go b/internal/client/push.go index 8300f26..1d64843 100644 --- a/internal/client/push.go +++ b/internal/client/push.go @@ -32,12 +32,11 @@ type PushOptions struct { } type pushWatchState struct { - manifest string - base *manifeststate.Snapshot - baseChecked bool - generation string - builder snapshot.Builder - watcher *snapshot.Watch + manifest string + base *manifeststate.Snapshot + generation string + builder snapshot.Builder + watcher *snapshot.Watch } var errWatchUnchanged = errors.New("watched source is unchanged") @@ -127,6 +126,9 @@ func pushChangesLocked(opts PushOptions, ws proto.Workspace, origin workspaceOri } else if !os.IsNotExist(err) { return err } + if pending.Request.ID != "" && pending.Request.Delta == nil { + return fmt.Errorf("pending push is missing its source delta") + } if pending.Applying { // Complete the exact interrupted request before accepting another one. The // remote receipt makes a lost response safe to retry despite later job edits. @@ -135,7 +137,6 @@ func pushChangesLocked(opts PushOptions, ws proto.Workspace, origin workspaceOri if opts.watchState != nil { opts.watchState.manifest = "" opts.watchState.base = nil - opts.watchState.baseChecked = false } return err } @@ -174,7 +175,7 @@ func pushChangesLocked(opts PushOptions, ws proto.Workspace, origin workspaceOri if opts.watchState != nil && opts.watchState.manifest == manifestHash { return errWatchUnchanged } - if opts.refreshSource || pending.Request.ID == "" || pushManifestRoot(pending.Request) != manifestHash { + if opts.refreshSource || pending.Request.ID == "" || pending.Request.SourceRoot != manifestHash { clientID, err := localChangeClientID() if err != nil { return err @@ -187,43 +188,31 @@ func pushChangesLocked(opts PushOptions, ws proto.Workspace, origin workspaceOri if state == nil { state = &pushWatchState{} } - if !state.baseChecked { + if state.base == nil { base, err := pushBase(opts.PeerURL, ws.ID, clientID) if err != nil { return err } - if base != nil { - state.base, err = manifeststate.New(context.Background(), *base) - if err != nil { + accepted, err := manifeststate.New(context.Background(), base) + if err != nil { + return err + } + if opts.watchState != nil { + if err := accepted.PrepareUpdates(context.Background()); err != nil { return err } - if opts.watchState != nil { - if err := state.base.PrepareUpdates(context.Background()); err != nil { - return err - } - } } - state.baseChecked = true + state.base = accepted } base := state.base - if base != nil { - plan, err := changeops.PrepareSnapshotDelta(context.Background(), base, sourceState, proto.DefaultLimits().MaxChangeBytes) - if err != nil { - return err - } - preparedDelta = plan - delta := plan.Bundle() - pending.Request.Delta = &delta - pending.Request.SourceRoot = manifestHash - } - } - // Delta recovery uses its frozen changed bodies and SourceRoot. Keeping - // the full inventory here would encode and fsync it on every state write. - if pending.Request.Delta == nil { - pending.Request.Manifest, err = sourceState.Manifest(context.Background()) + plan, err := changeops.PrepareSnapshotDelta(context.Background(), base, sourceState, proto.DefaultLimits().MaxChangeBytes) if err != nil { return err } + preparedDelta = plan + delta := plan.Bundle() + pending.Request.Delta = &delta + pending.Request.SourceRoot = manifestHash } parent := filepath.Join(dir, "push-sources") if err := ensurePrivateLocalDirectory(parent); err != nil { @@ -259,7 +248,6 @@ func pushChangesLocked(opts PushOptions, ws proto.Workspace, origin workspaceOri opts.retryCheckpoint, opts.refreshSource = true, true if opts.watchState != nil { opts.watchState.base = nil - opts.watchState.baseChecked = false opts.watchState.manifest = "" } return pushChangesLocked(opts, ws, origin, dir, result) @@ -293,7 +281,7 @@ func pushChangesLocked(opts PushOptions, ws proto.Workspace, origin workspaceOri err = finishPush(opts.PeerURL, ws.ID, dir, pending, result, opts.meter) if err == nil && opts.watchState != nil { opts.watchState.generation = pending.Request.ID - if opts.watchState.base != nil && pending.Request.Delta != nil { + if opts.watchState.base != nil { var base *manifeststate.Snapshot var err error if preparedDelta != nil { @@ -326,9 +314,12 @@ func prunePushSources(parent, keep string) error { return nil } func uploadPush(peer, workspace, source string, request proto.PushRequest, meter *transferMeter) (proto.PushResult, error) { + if request.Delta == nil { + return proto.PushResult{}, fmt.Errorf("push is missing its source delta") + } // A small delta costs less to send directly than another network round trip // to discover whether its bodies are cached. Larger deltas still negotiate. - if request.Delta != nil && smallPushDelta(request.Delta.RemoteManifest) { + if smallPushDelta(request.Delta.RemoteManifest) { return uploadPushOnce(peer, workspace, source, request, meter, shipPlan{}) } endpoint := strings.TrimSuffix(peer, "/") + "/v0/workspaces/" + workspace + "/push/diff" @@ -347,6 +338,9 @@ func uploadPush(peer, workspace, source string, request proto.PushRequest, meter func uploadPushOnce(peer, workspace, source string, request proto.PushRequest, meter *transferMeter, plan shipPlan) (proto.PushResult, error) { var result proto.PushResult + if request.Delta == nil { + return result, fmt.Errorf("push is missing its source delta") + } pr, pw := io.Pipe() defer pr.Close() mw := multipart.NewWriter(pw) @@ -358,11 +352,7 @@ func uploadPushOnce(peer, workspace, source string, request proto.PushRequest, m if err != nil { return err } - wire := request - if wire.Delta != nil { - wire.Manifest = proto.Manifest{} - } - if err := json.NewEncoder(part).Encode(wire); err != nil { + if err := json.NewEncoder(part).Encode(request); err != nil { return err } part, err = mw.CreateFormFile("workspace", "workspace.tar") @@ -377,11 +367,6 @@ func uploadPushOnce(peer, workspace, source string, request proto.PushRequest, m pw.CloseWithError(err) }() endpoint := strings.TrimSuffix(peer, "/") + "/v0/workspaces/" + workspace + "/push" - if request.Delta != nil { - // Old decoders ignore unknown JSON fields, so capability negotiation - // alone cannot protect a long-running client from a daemon rollback. - endpoint += "/delta-v1" - } req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, meter.readCloser(pr)) if err != nil { return result, err diff --git a/internal/client/push_delta.go b/internal/client/push_delta.go index 0666a7b..f72faa8 100644 --- a/internal/client/push_delta.go +++ b/internal/client/push_delta.go @@ -2,37 +2,28 @@ package client import ( "context" - "errors" - "net/http" "strings" "github.com/lydakis/errand/internal/archive" "github.com/lydakis/errand/internal/proto" ) -func pushBase(peer, workspace, client string) (*proto.Manifest, error) { +func pushBase(peer, workspace, client string) (proto.Manifest, error) { ctx, cancel := context.WithTimeout(context.Background(), controlRequestTimeout) defer cancel() var base proto.Manifest err := getJSONContext(ctx, strings.TrimSuffix(peer, "/")+"/v0/workspaces/"+workspace+"/push/base?client="+client, maxWorkspaceResponseBytes, "push checkpoint", &base) - var response *controlHTTPError - if errors.As(err, &response) && response.statusCode == http.StatusNotFound { - return nil, nil - } // old daemon if err != nil { - return nil, err + return base, err } if err := archive.Validate(base); err != nil { - return nil, err + return base, err } - return &base, nil + return base, nil } func pushSourceManifest(request proto.PushRequest) proto.Manifest { - if request.Delta != nil { - return request.Delta.RemoteManifest - } - return request.Manifest + return request.Delta.RemoteManifest } func smallPushDelta(manifest proto.Manifest) bool { @@ -47,11 +38,3 @@ func smallPushDelta(manifest proto.Manifest) bool { } return true } - -// SourceRoot is bound to the complete manifest when a delta is prepared. -func pushManifestRoot(request proto.PushRequest) string { - if request.Delta != nil { - return request.SourceRoot - } - return request.Manifest.RootHash() -} diff --git a/internal/client/push_generation.go b/internal/client/push_generation.go index b73ceaa..627d054 100644 --- a/internal/client/push_generation.go +++ b/internal/client/push_generation.go @@ -12,7 +12,7 @@ func (s *pushWatchState) observeGeneration(dir string) error { return err } if s.generation != generation { - s.manifest, s.base, s.baseChecked = "", nil, false + s.manifest, s.base = "", nil s.generation = generation } return nil diff --git a/internal/client/snapshot_transfer.go b/internal/client/snapshot_transfer.go index c781892..a0bd8a8 100644 --- a/internal/client/snapshot_transfer.go +++ b/internal/client/snapshot_transfer.go @@ -68,8 +68,6 @@ func negotiateSnapshotAt(ctx context.Context, endpoint string, manifest proto.Ma } switch resp.StatusCode { case http.StatusOK: - case http.StatusNotFound: // Snapshot caching is disabled on this runner. - return shipPlan{}, nil default: return shipPlan{}, &controlHTTPError{statusCode: resp.StatusCode, err: fmt.Errorf("snapshot negotiation: %s: %s", resp.Status, apiError(raw))} } @@ -77,6 +75,17 @@ func negotiateSnapshotAt(ctx context.Context, endpoint string, manifest proto.Ma if err := json.Unmarshal(raw, &diff); err != nil { return shipPlan{}, err } + for _, h := range diff.Missing { + if !seen[h] { + return shipPlan{}, fmt.Errorf("snapshot negotiation returned an unrequested or duplicate hash") + } + delete(seen, h) + } + // Cold and disabled caches need every body. Avoid retaining another hash + // map and checking it for each entry while packing a complete archive. + if len(diff.Missing) == len(refs) { + return shipPlan{}, nil + } ship := make(map[string]bool, len(diff.Missing)) for _, h := range diff.Missing { ship[h] = true diff --git a/internal/client/snapshot_transfer_test.go b/internal/client/snapshot_transfer_test.go index 75f5a70..a9644d8 100644 --- a/internal/client/snapshot_transfer_test.go +++ b/internal/client/snapshot_transfer_test.go @@ -67,7 +67,7 @@ func TestPushAllowsStagingBeyondControlDeadline(t *testing.T) { json.NewEncoder(w).Encode(proto.PushResult{ID: "transfer", WorkspaceID: "workspace"}) })) defer server.Close() - if _, err := uploadPushOnce(server.URL, "workspace", t.TempDir(), proto.PushRequest{ID: "transfer"}, nil, shipPlan{}); err != nil { + if _, err := uploadPushOnce(server.URL, "workspace", t.TempDir(), proto.PushRequest{ID: "transfer", Delta: &proto.ChangeBundle{}}, nil, shipPlan{}); err != nil { t.Fatalf("valid slow staging was cut off by the control timeout: %v", err) } } @@ -108,7 +108,7 @@ func TestPushFallbackRequiresExplicitCacheMissAndStopsAfterFullUpload(t *testing } { t.Run(tc.name, func(t *testing.T) { root := t.TempDir() - if err := os.WriteFile(filepath.Join(root, "file"), []byte("frozen content"), 0600); err != nil { + if err := os.WriteFile(filepath.Join(root, "file"), []byte(strings.Repeat("x", 128<<10)), 0600); err != nil { t.Fatal(err) } manifest, err := snapshot.Build(root, []string{"file"}) @@ -127,7 +127,7 @@ func TestPushFallbackRequiresExplicitCacheMissAndStopsAfterFullUpload(t *testing json.NewEncoder(w).Encode(proto.APIError{Code: tc.code, Error: "rejected"}) })) defer server.Close() - if _, err := uploadPush(server.URL, "workspace", root, proto.PushRequest{ID: "transfer", Manifest: manifest}, nil); err == nil { + if _, err := uploadPush(server.URL, "workspace", root, proto.PushRequest{ID: "transfer", Delta: &proto.ChangeBundle{RemoteManifest: manifest}, SourceRoot: manifest.RootHash()}, nil); err == nil { t.Fatal("rejected push succeeded") } if uploads != tc.uploads { @@ -136,3 +136,50 @@ func TestPushFallbackRequiresExplicitCacheMissAndStopsAfterFullUpload(t *testing }) } } + +func replyMissingSnapshotBlobs(t *testing.T, w http.ResponseWriter, r *http.Request) { + t.Helper() + var req proto.SnapshotDiffRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + t.Error(err) + http.Error(w, "invalid diff", 400) + return + } + missing := make([]string, 0, len(req.Blobs)) + for _, blob := range req.Blobs { + missing = append(missing, blob.SHA256) + } + json.NewEncoder(w).Encode(proto.SnapshotDiffResponse{Missing: missing}) +} + +func TestRequiredSnapshotEndpointDoesNotPermitUploadFallback(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, "input"), []byte("data"), 0600); err != nil { + t.Fatal(err) + } + uploads := 0 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasPrefix(r.Header.Get("Content-Type"), "multipart/") { + uploads++ + } + http.NotFound(w, r) + })) + defer server.Close() + opts := RunOptions{PeerURL: server.URL, Root: root, Argv: []string{"true"}, Stdout: io.Discard, Stderr: io.Discard} + if code := Run(opts); code != ExitTransaction { + t.Fatalf("missing snapshot endpoint: exit=%d", code) + } + if _, err := CreateWorkspace(opts, "dev"); !IsNotFound(err) { + t.Fatalf("missing creation negotiation endpoint: %v", err) + } + if _, err := pushBase(server.URL, proto.NewULID(), "0123456789abcdef0123456789abcdef"); !IsNotFound(err) { + t.Fatalf("missing push checkpoint endpoint: %v", err) + } + if uploads != 0 { + t.Fatalf("sent %d uploads after required endpoints were missing", uploads) + } +} diff --git a/internal/client/upload_deadline_test.go b/internal/client/upload_deadline_test.go index 6e9956a..2a63e06 100644 --- a/internal/client/upload_deadline_test.go +++ b/internal/client/upload_deadline_test.go @@ -49,7 +49,7 @@ func TestUploadsUseAdmissionDeadlineAndPreserveOrigins(t *testing.T) { canceled := 0 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if strings.HasSuffix(r.URL.Path, "/diff") { - http.NotFound(w, r) + replyMissingSnapshotBlobs(t, w, r) return } if _, err := io.Copy(io.Discard, r.Body); err != nil { diff --git a/internal/client/workspace_descriptor_test.go b/internal/client/workspace_descriptor_test.go index 7d7d8cf..0ee9138 100644 --- a/internal/client/workspace_descriptor_test.go +++ b/internal/client/workspace_descriptor_test.go @@ -10,8 +10,8 @@ import ( "github.com/lydakis/errand/internal/proto" ) -func TestWorkspaceDescriptorSupportsLegacyServersAndSharedValidation(t *testing.T) { - for _, scenario := range []string{"compact", "legacy", "invalid-id", "invalid-name"} { +func TestWorkspaceDescriptorUsesCompactResponseAndSharedValidation(t *testing.T) { + for _, scenario := range []string{"compact", "invalid-id", "invalid-name", "ignored-omit"} { t.Run(scenario, func(t *testing.T) { want := proto.Workspace{ID: proto.NewULID(), Name: "descriptor", Selection: proto.SelectionPolicy{Artifacts: []string{"output"}}, Manifest: proto.Manifest{Entries: []proto.ManifestEntry{{Path: "dir", Type: proto.EntryDir, Mode: 0755}}}} if scenario == "invalid-id" { @@ -26,7 +26,7 @@ func TestWorkspaceDescriptorSupportsLegacyServersAndSharedValidation(t *testing. if query != "omit" { t.Errorf("unexpected manifest query: %q", query) } - if scenario != "legacy" { + if scenario != "ignored-omit" { response.Manifest = proto.Manifest{} } } @@ -35,6 +35,12 @@ func TestWorkspaceDescriptorSupportsLegacyServersAndSharedValidation(t *testing. defer server.Close() descriptor, descriptorErr := getWorkspaceDescriptor(server.URL, "descriptor") full, fullErr := GetWorkspace(server.URL, "descriptor") + if scenario == "ignored-omit" { + if descriptorErr == nil || fullErr != nil { + t.Fatalf("ignored descriptor query: descriptor=%v full=%v", descriptorErr, fullErr) + } + return + } if scenario == "invalid-id" || scenario == "invalid-name" { if descriptorErr == nil || fullErr == nil { t.Fatalf("invalid identity accepted: descriptor=%v full=%v", descriptorErr, fullErr) diff --git a/internal/client/workspace_transfer_failure_test.go b/internal/client/workspace_transfer_failure_test.go index 86f07d0..79b5edf 100644 --- a/internal/client/workspace_transfer_failure_test.go +++ b/internal/client/workspace_transfer_failure_test.go @@ -85,7 +85,13 @@ func TestRejectedWorkspaceCreationReclaimsOrigin(t *testing.T) { if err := os.WriteFile(filepath.Join(root, "data"), []byte(strings.Repeat("x", 32768)), 0600); err != nil { t.Fatal(err) } - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { http.Error(w, "creation failed", status) })) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/diff") { + replyMissingSnapshotBlobs(t, w, r) + return + } + http.Error(w, "creation failed", status) + })) defer server.Close() if _, err := CreateWorkspace(RunOptions{PeerURL: server.URL, Root: root, IncludeAll: true}, "duplicate"); err == nil { t.Fatal("expected failure") diff --git a/internal/client/workspaces.go b/internal/client/workspaces.go index 98f81a7..a5d6627 100644 --- a/internal/client/workspaces.go +++ b/internal/client/workspaces.go @@ -16,7 +16,7 @@ import ( "github.com/lydakis/errand/internal/snapshot" ) -// Allow the runner's 64 MiB manifest plus creation metadata in a descriptor. +// Allow the runner's 64 MiB manifest plus creation metadata in a full read. const maxWorkspaceResponseBytes = 65 << 20 // Names resolve to the current workspace identity. IDs also work after removal, @@ -42,7 +42,7 @@ func GetWorkspace(peerURL, name string) (proto.Workspace, error) { } // getWorkspaceDescriptor requests identity and selection metadata without the -// creation manifest. Older runners may ignore the query and return the full row. +// creation manifest. func getWorkspaceDescriptor(peerURL, name string) (proto.Workspace, error) { return getWorkspace(peerURL, name, true) } @@ -52,13 +52,18 @@ func getWorkspace(peerURL, name string, omitManifest bool) (proto.Workspace, err ctx, cancel := context.WithTimeout(context.Background(), controlRequestTimeout) defer cancel() endpoint := strings.TrimSuffix(peerURL, "/") + "/v0/workspaces/" + url.PathEscape(name) + responseLimit := int64(maxWorkspaceResponseBytes) if omitManifest { endpoint += "?manifest=omit" + responseLimit = 2 << 20 // selection metadata plus descriptor framing } - err := getJSONContext(ctx, endpoint, maxWorkspaceResponseBytes, "workspace", &result) + err := getJSONContext(ctx, endpoint, responseLimit, "workspace", &result) if err == nil && (!proto.ValidULID(result.ID) || proto.ValidateWorkspaceName(result.Name) != nil) { err = fmt.Errorf("runner returned an invalid workspace") } + if err == nil && omitManifest && len(result.Manifest.Entries) != 0 { + err = fmt.Errorf("runner did not omit the creation manifest") + } return result, err } @@ -139,9 +144,7 @@ func createPreparedWorkspace(opts RunOptions, prep snapshotPreparation, request endpoint := strings.TrimSuffix(opts.PeerURL, "/") + "/v0/workspaces/" + request.ID + "/snapshot/diff" plan, err := negotiateSnapshotAt(context.Background(), endpoint, prep.manifest) if err != nil { - // Like job submission, cache negotiation is optional. A complete - // upload uses the original endpoint and is safe without capability. - plan = shipPlan{} + return proto.Workspace{}, err } if err := prep.guard.Verify(); err != nil { return proto.Workspace{}, err @@ -199,9 +202,6 @@ func createPreparedWorkspaceOnce(opts RunOptions, prep snapshotPreparation, requ ctx, cancel := context.WithTimeout(context.Background(), submitRequestTimeout) defer cancel() endpoint := strings.TrimSuffix(opts.PeerURL, "/") + "/v0/workspaces/" + request.ID - if plan.partial { - endpoint += "/snapshot" - } req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, pr) if err != nil { return result, err diff --git a/internal/daemon/milestone3_test.go b/internal/daemon/milestone3_test.go index 48d0fcd..8040f6f 100644 --- a/internal/daemon/milestone3_test.go +++ b/internal/daemon/milestone3_test.go @@ -436,8 +436,8 @@ func TestCacheDisabledUsesFullSnapshotFlow(t *testing.T) { t.Fatal(err) } resp.Body.Close() - if resp.StatusCode != http.StatusNotFound { - t.Fatalf("diff with disabled cache = %s, want 404", resp.Status) + if resp.StatusCode != http.StatusOK { + t.Fatalf("diff with disabled cache = %s, want 200", resp.Status) } root := workspaceWith(t, map[string]string{"f.txt": "x"}) diff --git a/internal/daemon/placement_test.go b/internal/daemon/placement_test.go index 058e248..e33a0e7 100644 --- a/internal/daemon/placement_test.go +++ b/internal/daemon/placement_test.go @@ -125,7 +125,7 @@ func TestWhereChecksEffectiveExecutePermission(t *testing.T) { func TestWhereInfoAndWorkspaceCreation(t *testing.T) { _, ts := testDaemon(t) info, err := client.ProbeWhereInfo(context.Background(), ts.URL, "os="+runtime.GOOS, time.Second) - if err != nil || !info.Placement || info.Facts.OS != runtime.GOOS { + if err != nil || info.Facts.OS != runtime.GOOS { t.Fatalf("info=%+v %v", info, err) } wrongOS := "linux" diff --git a/internal/daemon/server.go b/internal/daemon/server.go index e4915b8..df93929 100644 --- a/internal/daemon/server.go +++ b/internal/daemon/server.go @@ -733,14 +733,12 @@ func (d *Daemon) runQueue() { func (d *Daemon) Handler() http.Handler { mux := http.NewServeMux() mux.HandleFunc("POST /v0/gc/changes", d.auth(proto.ActionGCJobs, d.handleTransferGC)) - mux.HandleFunc("POST /v0/workspaces/{id}/push/delta-v1", d.auth(proto.ActionSubmit, d.handleWorkspacePush)) mux.HandleFunc("POST /v0/workspaces/{id}/push", d.auth(proto.ActionSubmit, d.handleWorkspacePush)) mux.HandleFunc("POST /v0/workspaces/{id}/push/diff", d.auth(proto.ActionSubmit, d.handleWorkspacePushDiff)) mux.HandleFunc("GET /v0/workspaces/{id}/push/base", d.auth(proto.ActionSubmit, d.handleWorkspacePushBase)) mux.HandleFunc("POST /v0/workspaces/{id}/push/{transfer}/apply", d.auth(proto.ActionSubmit, d.handleWorkspacePushApply)) mux.HandleFunc("POST /v0/workspaces/{id}", d.auth(proto.ActionSubmit, d.handleWorkspaceCreate)) mux.HandleFunc("POST /v0/workspaces/{id}/snapshot/diff", d.auth(proto.ActionSubmit, d.handleWorkspaceCreateDiff)) - mux.HandleFunc("POST /v0/workspaces/{id}/snapshot", d.auth(proto.ActionSubmit, d.handleWorkspaceCreateSnapshot)) mux.HandleFunc("GET /v0/workspaces", d.auth(proto.ActionReadOwn, d.handleWorkspaceList)) mux.HandleFunc("GET /v0/workspaces/{id}", d.auth(proto.ActionReadOwn, d.handleWorkspaceGet)) mux.HandleFunc("DELETE /v0/workspaces/{id}", d.auth(proto.ActionGCJobs, d.handleWorkspaceRemove)) @@ -886,7 +884,6 @@ func (d *Daemon) handleInfo(w http.ResponseWriter, r *http.Request, id Identity) busy := d.capacityFullLocked() || d.setupQuiesceToken != "" && time.Now().Before(d.setupQuiesceUntil) d.mu.Unlock() writeJSON(w, http.StatusOK, proto.Info{ - Placement: true, SSHDisabled: d.cfg.DisableSSH || d.cfg.LocalOnly, LocalOnly: d.cfg.LocalOnly, Proto: proto.ProtoVersion, @@ -906,15 +903,19 @@ func (d *Daemon) handleInfo(w http.ResponseWriter, r *http.Request, id Identity) // Negotiation is advisory; extraction detects blobs evicted before submission. func (d *Daemon) handleSnapshotDiff(w http.ResponseWriter, r *http.Request, _ Identity) { - if d.cache == nil { - httpError(w, http.StatusNotFound, "snapshot cache is disabled on this runner") - return - } var req proto.SnapshotDiffRequest if err := json.NewDecoder(io.LimitReader(r.Body, maxManifestBytes)).Decode(&req); err != nil { httpError(w, http.StatusBadRequest, err.Error()) return } + if d.cache == nil { + missing := make([]string, 0, len(req.Blobs)) + for _, blob := range req.Blobs { + missing = append(missing, blob.SHA256) + } + writeJSON(w, http.StatusOK, proto.SnapshotDiffResponse{Missing: missing}) + return + } missing, err := d.cache.MissingContext(r.Context(), req.Blobs) if err != nil { if r.Context().Err() != nil { diff --git a/internal/daemon/workspace_create_snapshot_test.go b/internal/daemon/workspace_create_snapshot_test.go index 5219f67..a517efc 100644 --- a/internal/daemon/workspace_create_snapshot_test.go +++ b/internal/daemon/workspace_create_snapshot_test.go @@ -14,7 +14,7 @@ import ( ) func TestWorkspaceCreationSnapshotReuseAndFallback(t *testing.T) { - for _, scenario := range []string{"warm", "old-runner", "disabled", "evicted", "downgraded", "quota", "server-error"} { + for _, scenario := range []string{"warm", "disabled", "evicted", "quota", "server-error"} { t.Run(scenario, func(t *testing.T) { d, err := New(Config{StateDir: t.TempDir(), InsecureNoAuth: true, CacheDisabled: scenario == "disabled"}) if err != nil { @@ -27,10 +27,6 @@ func TestWorkspaceCreationSnapshotReuseAndFallback(t *testing.T) { var recording bool var invalidate func() server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if strings.HasSuffix(r.URL.Path, "/snapshot/diff") && scenario == "old-runner" { - http.NotFound(w, r) - return - } create := r.Method == http.MethodPost && strings.HasPrefix(r.URL.Path, "/v0/workspaces/") && !strings.HasSuffix(r.URL.Path, "/diff") if recording && create { raw, err := io.ReadAll(r.Body) @@ -50,10 +46,6 @@ func TestWorkspaceCreationSnapshotReuseAndFallback(t *testing.T) { httpErrorCode(w, http.StatusInternalServerError, "snapshot_cache_miss", "server failure is not a retry authorization") return } - if scenario == "downgraded" && strings.HasSuffix(r.URL.Path, "/snapshot") { - http.NotFound(w, r) - return - } } handler.ServeHTTP(w, r) })) @@ -78,9 +70,9 @@ func TestWorkspaceCreationSnapshotReuseAndFallback(t *testing.T) { } recording = true second, err := client.CreateWorkspace(client.RunOptions{PeerURL: server.URL, Root: root}, "second") - if scenario == "downgraded" || scenario == "quota" || scenario == "server-error" { + if scenario == "quota" || scenario == "server-error" { if err == nil || len(uploads) != 1 { - t.Fatalf("downgrade retried or succeeded: uploads=%v err=%v", uploads, err) + t.Fatalf("rejected creation retried or succeeded: uploads=%v err=%v", uploads, err) } return } @@ -96,15 +88,15 @@ func TestWorkspaceCreationSnapshotReuseAndFallback(t *testing.T) { if len(uploads) != 1 || uploads[0] > 16<<10 { t.Fatalf("warm creation retransmitted bodies: %v", uploads) } - case "old-runner", "disabled": - if len(uploads) != 1 || uploads[0] < 1<<20 || strings.HasSuffix(uploadPaths[0], "/snapshot") { - t.Fatalf("old runner received partial creation: %v %v", uploads, uploadPaths) + case "disabled": + if len(uploads) != 1 || uploads[0] < 1<<20 { + t.Fatalf("disabled cache received partial creation: %v %v", uploads, uploadPaths) } case "evicted": if len(uploads) != 2 || uploads[0] > 16<<10 || uploads[1] < 1<<20 { t.Fatalf("wrong fallback sequence: %v", uploads) } - if strings.TrimSuffix(uploadPaths[0], "/snapshot") != uploadPaths[1] { + if uploadPaths[0] != uploadPaths[1] { t.Fatalf("fallback changed creation identity: %v", uploadPaths) } } diff --git a/internal/daemon/workspace_push.go b/internal/daemon/workspace_push.go index 6eed9e4..ea89e06 100644 --- a/internal/daemon/workspace_push.go +++ b/internal/daemon/workspace_push.go @@ -94,8 +94,12 @@ func (d *Daemon) handleWorkspacePush(w http.ResponseWriter, r *http.Request, id httpError(w, 400, "invalid push identity") return } + if request.Delta == nil { + httpError(w, 400, "push is missing its source delta") + return + } var prepared changeops.PreparedTransferSource - if request.Delta != nil { + { unlock, err := d.workspaces.lockWorkspaceContext(r.Context(), row.ID) if err != nil { return @@ -106,9 +110,6 @@ func (d *Daemon) handleWorkspacePush(w http.ResponseWriter, r *http.Request, id // outside the apply gate; StagePrepared rechecks the baseline under it. if err == nil { prepared, err = changeops.ExpandTransferSourceBase(r.Context(), base, *request.Delta, request.SourceRoot, d.cfg.MaxLimits.MaxChangeBytes) - if err == nil { - request.Manifest = prepared.Manifest() - } } if err != nil { if errors.Is(err, changeops.ErrCheckpointChanged) { @@ -119,14 +120,9 @@ func (d *Daemon) handleWorkspacePush(w http.ResponseWriter, r *http.Request, id return } } - if request.Delta == nil { - if err := archive.Validate(request.Manifest); err != nil { - httpError(w, 400, err.Error()) - return - } - } + fullManifest := prepared.Manifest() var total int64 - for _, e := range request.Manifest.Entries { + for _, e := range fullManifest.Entries { if e.Type == proto.EntryFile { if e.Size > d.cfg.MaxLimits.MaxWorkspaceBytes-total { httpError(w, 400, "workspace source exceeds byte limit") @@ -139,17 +135,14 @@ func (d *Daemon) handleWorkspacePush(w http.ResponseWriter, r *http.Request, id return } } - sourceManifest := request.Manifest - if request.Delta != nil { - sourceManifest = request.Delta.RemoteManifest - } + sourceManifest := request.Delta.RemoteManifest part, err := nextPart(mr, "workspace") if err != nil { httpError(w, 400, err.Error()) return } extractOpts, restored := d.snapshotExtractOptions(r.Context(), nil) - extractOpts.SymlinkManifest = &request.Manifest + extractOpts.SymlinkManifest = &fullManifest if err := archive.ExtractWith(&contextReader{ctx: r.Context(), r: part}, source, sourceManifest, d.cfg.MaxLimits.MaxWorkspaceBytes, extractOpts); err != nil { if errors.Is(err, archive.ErrCacheMiss) { httpErrorCode(w, http.StatusConflict, proto.ErrorCodeSnapshotCacheMiss, err.Error()) @@ -194,12 +187,7 @@ func (d *Daemon) handleWorkspacePush(w http.ResponseWriter, r *http.Request, id httpError(w, 500, err.Error()) return } - var bundle proto.ChangeBundle - if request.Delta != nil { - _, bundle, err = session.StagePrepared(r.Context(), request.ID, source, prepared) - } else { - _, bundle, err = session.Stage(r.Context(), request.ID, source, request.Manifest) - } + _, bundle, err := session.StagePrepared(r.Context(), request.ID, source, prepared) if err != nil { if errors.Is(err, changeops.ErrCheckpointChanged) { httpErrorCode(w, http.StatusConflict, proto.ErrorCodePushCheckpointChanged, err.Error()) diff --git a/internal/daemon/workspace_push_delta.go b/internal/daemon/workspace_push_delta.go index d46c9a5..a4c1442 100644 --- a/internal/daemon/workspace_push_delta.go +++ b/internal/daemon/workspace_push_delta.go @@ -19,7 +19,7 @@ func (d *Daemon) pushBase(row workspaceRecord, client string) (*changeops.Source } // The retained source checkpoint is independent of the running application's -// mutable files. Advertising it also negotiates delta-source upload support. +// mutable files. Every push uses this checkpoint to prepare its source delta. func (d *Daemon) handleWorkspacePushBase(w http.ResponseWriter, r *http.Request, id Identity) { client := r.URL.Query().Get("client") if !proto.ValidChangeClientID(client) { diff --git a/internal/daemon/workspace_push_delta_test.go b/internal/daemon/workspace_push_delta_test.go index 14adce6..d71cf8e 100644 --- a/internal/daemon/workspace_push_delta_test.go +++ b/internal/daemon/workspace_push_delta_test.go @@ -36,13 +36,15 @@ func TestPushDeltaValidatesRetainedBaseAndFullSourceLimit(t *testing.T) { if err != nil { t.Fatal(err) } - for _, scenario := range []string{"stale", "forged", "quota", "valid"} { + for _, scenario := range []string{"stale", "forged", "quota", "missing-delta", "valid"} { t.Run(scenario, func(t *testing.T) { bundle := delta req := proto.PushRequest{ID: proto.NewULID(), ClientID: "0123456789abcdef0123456789abcdef", Delta: &bundle, SourceRoot: current.RootHash()} limit := d.cfg.MaxLimits.MaxWorkspaceBytes defer func() { d.cfg.MaxLimits.MaxWorkspaceBytes = limit }() switch scenario { + case "missing-delta": + req.Delta = nil case "stale": bundle.BaselineRoot = strings.Repeat("0", 64) case "forged": @@ -75,7 +77,7 @@ func TestPushDeltaValidatesRetainedBaseAndFullSourceLimit(t *testing.T) { w := httptest.NewRecorder() d.handleWorkspacePush(w, r, Identity{}) want := 409 - if scenario == "quota" { + if scenario == "quota" || scenario == "missing-delta" { want = 400 } else if scenario == "valid" { want = 201 diff --git a/internal/daemon/workspace_push_snapshot_test.go b/internal/daemon/workspace_push_snapshot_test.go index 33837da..8c1fa74 100644 --- a/internal/daemon/workspace_push_snapshot_test.go +++ b/internal/daemon/workspace_push_snapshot_test.go @@ -18,7 +18,8 @@ import ( func TestPushReusesSnapshotBodies(t *testing.T) { d, ts := testDaemon(t) - root := workspaceWith(t, map[string]string{"edit": "before\n", "unchanged": strings.Repeat("x", 1<<20)}) + updated := strings.Repeat("u", 128<<10) + root := workspaceWith(t, map[string]string{"edit": "before\n", "cached": updated, "unchanged": strings.Repeat("x", 1<<20)}) ws, err := client.CreateWorkspace(client.RunOptions{PeerURL: ts.URL, Root: root}, "small-push") if err != nil { t.Fatal(err) @@ -28,7 +29,7 @@ func TestPushReusesSnapshotBodies(t *testing.T) { if err := os.WriteFile(filepath.Join(remote, "unchanged"), []byte("runner edit\n"), 0600); err != nil { t.Fatal(err) } - if err := os.WriteFile(filepath.Join(root, "edit"), []byte("after!\n"), 0600); err != nil { + if err := os.WriteFile(filepath.Join(root, "edit"), []byte(updated), 0600); err != nil { t.Fatal(err) } var stats client.TransferStats @@ -39,7 +40,7 @@ func TestPushReusesSnapshotBodies(t *testing.T) { if stats.ChangedPaths != 1 || stats.TransferredBytes > 16<<10 { t.Fatalf("one-line push retransmitted unchanged content: %+v", stats) } - if body, err := os.ReadFile(filepath.Join(remote, "edit")); err != nil || string(body) != "after!\n" { + if body, err := os.ReadFile(filepath.Join(remote, "edit")); err != nil || string(body) != updated { t.Fatalf("edit not applied: %q %v", body, err) } opts.Path = "" @@ -75,7 +76,7 @@ func TestSnapshotIngestionContinuesAfterSourceFailure(t *testing.T) { } func TestPushSnapshotFallback(t *testing.T) { - for _, failure := range []string{"old-runner", "disabled", "cold", "evicted", "corrupt", "corrupt-unremovable"} { + for _, failure := range []string{"disabled", "cold", "evicted", "corrupt", "corrupt-unremovable"} { t.Run(failure, func(t *testing.T) { if failure == "corrupt-unremovable" && os.Geteuid() == 0 { t.Skip("root bypasses directory permissions") @@ -88,15 +89,6 @@ func TestPushSnapshotFallback(t *testing.T) { var bodies []int var invalidate func() server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - // Keep coverage for full-snapshot cache fallback against legacy peers. - if strings.HasSuffix(r.URL.Path, "/push/base") { - http.NotFound(w, r) - return - } - if strings.HasSuffix(r.URL.Path, "/push/diff") && failure == "old-runner" { - http.NotFound(w, r) - return - } if strings.HasSuffix(r.URL.Path, "/push") { body, err := io.ReadAll(r.Body) if err != nil { @@ -114,13 +106,14 @@ func TestPushSnapshotFallback(t *testing.T) { d.Handler().ServeHTTP(w, r) })) defer server.Close() - root := workspaceWith(t, map[string]string{"edit": "before\n", "unchanged": strings.Repeat("x", 1<<20)}) + updated := strings.Repeat("u", 128<<10) + root := workspaceWith(t, map[string]string{"edit": "before\n", "cached": updated, "unchanged": strings.Repeat("x", 1<<20)}) ws, err := client.CreateWorkspace(client.RunOptions{PeerURL: server.URL, Root: root}, "fallback") if err != nil { t.Fatal(err) } for _, e := range ws.Manifest.Entries { - if e.Path != "unchanged" || d.cache == nil { + if e.Path != "cached" || d.cache == nil { continue } invalidate = func() { @@ -146,7 +139,7 @@ func TestPushSnapshotFallback(t *testing.T) { } else if failure != "evicted" && !strings.HasPrefix(failure, "corrupt") { invalidate = nil } - if err := os.WriteFile(filepath.Join(root, "edit"), []byte("after!\n"), 0600); err != nil { + if err := os.WriteFile(filepath.Join(root, "edit"), []byte(updated), 0600); err != nil { t.Fatal(err) } var stats client.TransferStats @@ -164,10 +157,10 @@ func TestPushSnapshotFallback(t *testing.T) { for _, size := range bodies { total += int64(size) } - if len(bodies) != wantUploads || total != stats.TransferredBytes || total < 1<<20 { + if len(bodies) != wantUploads || total != stats.TransferredBytes || total < 128<<10 || total > 256<<10 { t.Fatalf("fallback uploads=%v stats=%+v", bodies, stats) } - if body, err := os.ReadFile(filepath.Join(d.workspaces.dir, ws.ID, "data", "edit")); err != nil || string(body) != "after!\n" { + if body, err := os.ReadFile(filepath.Join(d.workspaces.dir, ws.ID, "data", "edit")); err != nil || string(body) != updated { t.Fatalf("fallback did not apply: %q %v", body, err) } }) diff --git a/internal/daemon/workspace_snapshot_cache.go b/internal/daemon/workspace_snapshot_cache.go index e4a169b..965a6d3 100644 --- a/internal/daemon/workspace_snapshot_cache.go +++ b/internal/daemon/workspace_snapshot_cache.go @@ -11,8 +11,7 @@ import ( "github.com/lydakis/errand/internal/proto" ) -// This endpoint also establishes that the runner can reconstruct partial push -// archives. Older runners return 404 and clients keep using complete uploads. +// Push uploads negotiate verified content against the runner's snapshot cache. func (d *Daemon) handleWorkspacePushDiff(w http.ResponseWriter, r *http.Request, id Identity) { if _, err := d.pushWorkspace(r, id); err != nil { workspaceHTTPError(w, err) @@ -21,8 +20,7 @@ func (d *Daemon) handleWorkspacePushDiff(w http.ResponseWriter, r *http.Request, d.handleSnapshotDiff(w, r, id) } -// Creation has a separate capability and upload endpoint: an older receiver -// that supports job caching must never be sent a partial creation body. +// Creation negotiates content against the same snapshot cache before upload. func (d *Daemon) handleWorkspaceCreateDiff(w http.ResponseWriter, r *http.Request, id Identity) { if !proto.ValidULID(r.PathValue("id")) { httpError(w, 400, "invalid workspace id") diff --git a/internal/daemon/workspace_transfer_test.go b/internal/daemon/workspace_transfer_test.go index e9acfef..c04777c 100644 --- a/internal/daemon/workspace_transfer_test.go +++ b/internal/daemon/workspace_transfer_test.go @@ -26,7 +26,7 @@ func TestPushRecoversCollectedStageFromFrozenSource(t *testing.T) { var uploads atomic.Int32 var interrupted atomic.Bool server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if r.Method == "POST" && (strings.HasSuffix(r.URL.Path, "/push") || strings.HasSuffix(r.URL.Path, "/push/delta-v1")) { + if r.Method == "POST" && strings.HasSuffix(r.URL.Path, "/push") { uploads.Add(1) } if strings.HasSuffix(r.URL.Path, "/apply") && !interrupted.Swap(true) { diff --git a/internal/daemon/workspace_upload_test.go b/internal/daemon/workspace_upload_test.go index ca98d65..5200a92 100644 --- a/internal/daemon/workspace_upload_test.go +++ b/internal/daemon/workspace_upload_test.go @@ -51,13 +51,17 @@ func TestWorkspaceUploadDoesNotBlockCommandsOrRetargetReplacement(t *testing.T) if err != nil { t.Fatal(err) } + delta, err := changeops.PrepareSourceDelta(t.Context(), ws.Manifest, ws.Manifest, proto.DefaultLimits().MaxChangeBytes) + if err != nil { + t.Fatal(err) + } var payload bytes.Buffer mw := multipart.NewWriter(&payload) part, err := mw.CreateFormField("metadata") if err != nil { t.Fatal(err) } - if err := json.NewEncoder(part).Encode(proto.PushRequest{ID: proto.NewULID(), ClientID: "0123456789abcdef0123456789abcdef", Manifest: ws.Manifest}); err != nil { + if err := json.NewEncoder(part).Encode(proto.PushRequest{ID: proto.NewULID(), ClientID: "0123456789abcdef0123456789abcdef", Delta: &delta, SourceRoot: ws.Manifest.RootHash()}); err != nil { t.Fatal(err) } part, err = mw.CreateFormFile("workspace", "workspace.tar") @@ -65,7 +69,7 @@ func TestWorkspaceUploadDoesNotBlockCommandsOrRetargetReplacement(t *testing.T) t.Fatal(err) } split := payload.Len() - if err := snapshot.PackPartial(part, root, ws.Manifest, nil); err != nil { + if err := snapshot.PackPartial(part, root, delta.RemoteManifest, nil); err != nil { t.Fatal(err) } if err := mw.Close(); err != nil { diff --git a/internal/daemon/workspace_upload_unix_test.go b/internal/daemon/workspace_upload_unix_test.go index 79d6881..466fd38 100644 --- a/internal/daemon/workspace_upload_unix_test.go +++ b/internal/daemon/workspace_upload_unix_test.go @@ -16,6 +16,7 @@ import ( "syscall" "testing" + changeops "github.com/lydakis/errand/internal/changes" "github.com/lydakis/errand/internal/client" "github.com/lydakis/errand/internal/proto" "github.com/lydakis/errand/internal/snapshot" @@ -29,13 +30,17 @@ func TestWaitingWorkspacePushesHoldNoDataHandle(t *testing.T) { t.Fatal(err) } body := func() (*bytes.Buffer, string, int) { + delta, err := changeops.PrepareSourceDelta(t.Context(), ws.Manifest, ws.Manifest, proto.DefaultLimits().MaxChangeBytes) + if err != nil { + t.Fatal(err) + } var payload bytes.Buffer mw := multipart.NewWriter(&payload) part, err := mw.CreateFormField("metadata") if err != nil { t.Fatal(err) } - if err := json.NewEncoder(part).Encode(proto.PushRequest{ID: proto.NewULID(), ClientID: "0123456789abcdef0123456789abcdef", Manifest: ws.Manifest}); err != nil { + if err := json.NewEncoder(part).Encode(proto.PushRequest{ID: proto.NewULID(), ClientID: "0123456789abcdef0123456789abcdef", Delta: &delta, SourceRoot: ws.Manifest.RootHash()}); err != nil { t.Fatal(err) } part, err = mw.CreateFormFile("workspace", "workspace.tar") @@ -43,7 +48,7 @@ func TestWaitingWorkspacePushesHoldNoDataHandle(t *testing.T) { t.Fatal(err) } split := payload.Len() - if err := snapshot.PackPartial(part, root, ws.Manifest, nil); err != nil { + if err := snapshot.PackPartial(part, root, delta.RemoteManifest, nil); err != nil { t.Fatal(err) } if err := mw.Close(); err != nil { diff --git a/internal/daemon/workspaces_http.go b/internal/daemon/workspaces_http.go index d93a133..add8a26 100644 --- a/internal/daemon/workspaces_http.go +++ b/internal/daemon/workspaces_http.go @@ -18,14 +18,6 @@ import ( ) func (d *Daemon) handleWorkspaceCreate(w http.ResponseWriter, r *http.Request, id Identity) { - d.createWorkspaceUpload(w, r, id, false) -} - -func (d *Daemon) handleWorkspaceCreateSnapshot(w http.ResponseWriter, r *http.Request, id Identity) { - d.createWorkspaceUpload(w, r, id, true) -} - -func (d *Daemon) createWorkspaceUpload(w http.ResponseWriter, r *http.Request, id Identity, partial bool) { var request proto.Workspace key := r.PathValue("id") if !proto.ValidULID(key) { @@ -117,13 +109,9 @@ func (d *Daemon) createWorkspaceUpload(w http.ResponseWriter, r *http.Request, i httpError(w, 500, err.Error()) return } - var extractOptions archive.ExtractOptions - var restored map[string]bool - if partial { - extractOptions, restored = d.snapshotExtractOptions(r.Context(), nil) - } + extractOptions, restored := d.snapshotExtractOptions(r.Context(), nil) if err := archive.ExtractWith(&contextReader{ctx: r.Context(), r: input}, data, request.Manifest, d.cfg.MaxLimits.MaxWorkspaceBytes, extractOptions); err != nil { - if partial && errors.Is(err, archive.ErrCacheMiss) { + if errors.Is(err, archive.ErrCacheMiss) { httpErrorCode(w, http.StatusConflict, proto.ErrorCodeSnapshotCacheMiss, err.Error()) return } diff --git a/internal/proto/proto.go b/internal/proto/proto.go index 83dc5a5..aeb9b46 100644 --- a/internal/proto/proto.go +++ b/internal/proto/proto.go @@ -490,7 +490,6 @@ type GPU struct { } type Info struct { - Placement bool `json:"placement"` // admission-time requirement validation LocalOnly bool `json:"local_only,omitempty"` // confirms network requests and SSH bridging are disabled SSHDisabled bool `json:"ssh_disabled"` Proto int `json:"proto"` diff --git a/internal/proto/transfer.go b/internal/proto/transfer.go index c65d974..efe620c 100644 --- a/internal/proto/transfer.go +++ b/internal/proto/transfer.go @@ -7,13 +7,11 @@ const ErrorCodePushStageMissing = "push_stage_missing" // PushRequest identifies an immutable source snapshot. Replaying its ID retries // the same transfer; applying a staged transfer uses that ID and no new upload. type PushRequest struct { - ID string `json:"id"` - ClientID string `json:"client_id"` - Manifest Manifest `json:"manifest"` + ID string `json:"id"` + ClientID string `json:"client_id"` // Delta carries changed source entries against a retained checkpoint. - // When present, Manifest is omitted on the wire and reconstructed remotely. - Delta *ChangeBundle `json:"delta,omitempty"` - SourceRoot string `json:"source_root,omitempty"` + Delta *ChangeBundle `json:"delta"` + SourceRoot string `json:"source_root"` } type PushApplyRequest struct { Path string `json:"path,omitempty"`