From b4356fa2c1c322ef163a97f7477db02284361013 Mon Sep 17 00:00:00 2001 From: George Lydakis Date: Thu, 1 Oct 2026 03:26:58 +0000 Subject: [PATCH 1/2] Use one current client and daemon protocol --- cmd/errand/placement.go | 2 - cmd/errand/placement_integration_test.go | 12 +- cmd/errand/placement_test.go | 12 +- cmd/errand/push_batch_test.go | 26 ++-- cmd/errand/push_benchmark_test.go | 2 +- cmd/errand/push_record_test.go | 8 +- cmd/errand/push_watch_delta_test.go | 116 ++++++++---------- cmd/errand/transfer_reporting_test.go | 2 +- docs/DESIGN.md | 37 +++--- docs/OPERATIONS.md | 4 +- docs/SERVICE_UPGRADES.md | 3 +- docs/TRANSFER_PERFORMANCE.md | 7 +- docs/USAGE.md | 10 +- internal/client/changes_test.go | 8 +- internal/client/client.go | 4 +- internal/client/client_test.go | 16 ++- internal/client/placement_test.go | 4 + internal/client/push.go | 77 +++++------- internal/client/push_delta.go | 27 +--- internal/client/push_generation.go | 2 +- internal/client/snapshot_transfer.go | 13 +- internal/client/snapshot_transfer_test.go | 53 +++++++- internal/client/upload_deadline_test.go | 2 +- internal/client/workspace_descriptor_test.go | 12 +- .../client/workspace_transfer_failure_test.go | 8 +- internal/client/workspaces.go | 18 +-- internal/daemon/milestone3_test.go | 4 +- internal/daemon/placement_test.go | 2 +- internal/daemon/server.go | 15 +-- .../daemon/workspace_create_snapshot_test.go | 22 ++-- internal/daemon/workspace_push.go | 29 ++--- internal/daemon/workspace_push_delta.go | 2 +- internal/daemon/workspace_push_delta_test.go | 6 +- .../daemon/workspace_push_snapshot_test.go | 29 ++--- internal/daemon/workspace_snapshot_cache.go | 6 +- internal/daemon/workspace_transfer_test.go | 2 +- internal/daemon/workspace_upload_test.go | 16 ++- internal/daemon/workspaces_http.go | 16 +-- internal/proto/proto.go | 1 - internal/proto/transfer.go | 10 +- 40 files changed, 323 insertions(+), 322 deletions(-) diff --git a/cmd/errand/placement.go b/cmd/errand/placement.go index 148de9ff..253b65c0 100644 --- a/cmd/errand/placement.go +++ b/cmd/errand/placement.go @@ -68,8 +68,6 @@ func chooseRunners(ctx context.Context, e config.EffectiveRun, probe placementPr return } switch { - case !info.Placement: - reasons[i] = "runner does not support requirement validation; upgrade it" case info.Busy: reasons[i] = "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 ca2d1751..48efd76c 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 9922a248..5aa752dd 100644 --- a/cmd/errand/placement_test.go +++ b/cmd/errand/placement_test.go @@ -15,18 +15,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": @@ -43,7 +41,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) } @@ -112,7 +110,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 @@ -130,7 +128,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 diff --git a/cmd/errand/push_batch_test.go b/cmd/errand/push_batch_test.go index 30994f97..3c1cbdd5 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 5645e731..a0d6b0f2 100644 --- a/cmd/errand/push_benchmark_test.go +++ b/cmd/errand/push_benchmark_test.go @@ -108,7 +108,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 f2f982b9..852f57b7 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 97d7a616..7a8a954c 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 c6fb2d90..965bf890 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 aea07d1a..2d539873 100644 --- a/docs/DESIGN.md +++ b/docs/DESIGN.md @@ -889,12 +889,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 @@ -918,9 +919,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`). @@ -1018,21 +1019,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. @@ -1041,9 +1040,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 ad835956..63c0a684 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -279,8 +279,8 @@ the requested `--older-than` boundary. Unresolved submitted jobs remain protected for 30 days, after which an explicit local GC may retire abandoned state. 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 b3cbc8f6..14acaf1b 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 d007e9d2..41c8df7f 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 d70b38bf..ca5382df 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -287,10 +287,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. @@ -322,8 +321,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 81516e89..c0931a5f 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 8c2b47a9..36e30625 100644 --- a/internal/client/client.go +++ b/internal/client/client.go @@ -302,8 +302,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 f824d541..757f8c2b 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 0825d2d1..f4d47dd9 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 f9f79e70..354ccc2f 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") @@ -122,6 +121,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. @@ -130,7 +132,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 } @@ -169,7 +170,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 @@ -182,43 +183,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 { @@ -254,7 +243,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) @@ -288,7 +276,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 { @@ -321,9 +309,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" @@ -342,6 +333,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) @@ -353,11 +347,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") @@ -372,11 +362,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 0666a7bc..f72faa87 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 b73ceaaa..627d0541 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 c7818928..a0bd8a85 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 75f5a704..a9644d8f 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 6e9956a6..2a63e060 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 7d7d8cfe..0ee9138a 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 86f07d0d..79b5edf8 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 7eaabc94..585aeb08 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 } @@ -132,9 +137,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 @@ -192,9 +195,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 f03d0262..f7723d51 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 1cb726de..c81e00e1 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 2732a519..184ab861 100644 --- a/internal/daemon/server.go +++ b/internal/daemon/server.go @@ -683,14 +683,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)) @@ -817,7 +815,6 @@ func (d *Daemon) handleInfo(w http.ResponseWriter, r *http.Request, _ 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, @@ -835,15 +832,19 @@ func (d *Daemon) handleInfo(w http.ResponseWriter, r *http.Request, _ 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 5219f671..a517efc4 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 8d2d551c..fa295db3 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.ExpandTransferSource(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,8 @@ 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 - } - } var total int64 - for _, e := range request.Manifest.Entries { + for _, e := range prepared.Manifest().Entries { if e.Type == proto.EntryFile { if e.Size > d.cfg.MaxLimits.MaxWorkspaceBytes-total { httpError(w, 400, "workspace source exceeds byte limit") @@ -139,10 +134,7 @@ 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()) @@ -193,12 +185,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 77d6d798..a34807b1 100644 --- a/internal/daemon/workspace_push_delta.go +++ b/internal/daemon/workspace_push_delta.go @@ -16,7 +16,7 @@ func (d *Daemon) pushBase(row workspaceRecord, client string) (proto.Manifest, e } // 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 6627540d..88206db7 100644 --- a/internal/daemon/workspace_push_delta_test.go +++ b/internal/daemon/workspace_push_delta_test.go @@ -35,13 +35,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": @@ -74,7 +76,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 262abad8..4690b5ed 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 054d581e..059086e0 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 2e87a1c1..e5a0bb4a 100644 --- a/internal/daemon/workspace_transfer_test.go +++ b/internal/daemon/workspace_transfer_test.go @@ -25,7 +25,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 1cb39deb..b362552f 100644 --- a/internal/daemon/workspace_upload_test.go +++ b/internal/daemon/workspace_upload_test.go @@ -53,13 +53,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") @@ -67,7 +71,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 { @@ -257,13 +261,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") @@ -271,7 +279,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 a80857f3..caae619a 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()) - } + extractOptions, restored := d.snapshotExtractOptions(r.Context()) 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 ce60be53..9c16adc2 100644 --- a/internal/proto/proto.go +++ b/internal/proto/proto.go @@ -467,7 +467,6 @@ type Facts 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 c65d974c..efe620c7 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"` From 6d3e841f123f8963d64ab10aa271df757b08499c Mon Sep 17 00:00:00 2001 From: George Lydakis Date: Thu, 1 Oct 2026 03:35:07 +0000 Subject: [PATCH 2/2] Use the current snapshot response in the GC test fixture --- cmd/errand/gc_test.go | 21 +++++++++++++++++++-- 1 file changed, 19 insertions(+), 2 deletions(-) diff --git a/cmd/errand/gc_test.go b/cmd/errand/gc_test.go index a089d355..cccfddfd 100644 --- a/cmd/errand/gc_test.go +++ b/cmd/errand/gc_test.go @@ -172,20 +172,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 {