From 8b15171a222e90963124b02fedb922baf0f7fb98 Mon Sep 17 00:00:00 2001 From: lczyk Date: Tue, 14 Jul 2026 18:20:16 +0200 Subject: [PATCH 1/8] feat(source): resumable chunked pull for remote images Layer blobs download through ranged requests with digest-named partials in the cache tmp area, so an interrupted pull resumes from the last received byte instead of restarting. Blobs verify against manifest digests before landing in the layout; progress renders from real byte counts. Replaces the layer-stream progress wrappers, which could not survive process death. --- internal/source/copy.go | 65 +++-- internal/source/copy_e2e_test.go | 186 ++++++++++++++ internal/source/fetch.go | 412 +++++++++++++++++++++++++++++++ internal/source/fetch_test.go | 203 +++++++++++++++ internal/source/progress.go | 239 ------------------ internal/source/progress_test.go | 78 ------ 6 files changed, 844 insertions(+), 339 deletions(-) create mode 100644 internal/source/copy_e2e_test.go create mode 100644 internal/source/fetch.go create mode 100644 internal/source/fetch_test.go delete mode 100644 internal/source/progress.go delete mode 100644 internal/source/progress_test.go diff --git a/internal/source/copy.go b/internal/source/copy.go index 1812a60..f02ed5e 100644 --- a/internal/source/copy.go +++ b/internal/source/copy.go @@ -1,6 +1,7 @@ package source import ( + "bytes" "context" "encoding/base64" "encoding/json" @@ -125,23 +126,26 @@ func writeRemoteLayout(ctx context.Context, dir, sourceRef string, platform Plat if err != nil { return err } - if platform.All { - desc, err := remote.Get(ref, options...) + desc, err := remote.Get(ref, options...) + if err != nil { + return err + } + + plan := &pullPlan{} + var rootDesc v1.Descriptor + if platform.All && desc.MediaType.IsIndex() { + idx, err := desc.ImageIndex() if err != nil { return err } - if desc.MediaType.IsIndex() { - idx, err := desc.ImageIndex() - if err != nil { - return err - } - if progress != nil { - fmt.Fprintln(progress, "olav: writing OCI image index...") - } - wrapped, finish := withIndexProgress(idx, progress) - defer finish() - return path.AppendIndex(wrapped, layout.WithAnnotations(map[string]string{"org.opencontainers.image.ref.name": "olav"})) + if progress != nil { + fmt.Fprintln(progress, "olav: writing OCI image index...") } + rootDesc, err = plan.addIndex(idx) + if err != nil { + return err + } + } else { img, err := desc.Image() if err != nil { return err @@ -149,20 +153,37 @@ func writeRemoteLayout(ctx context.Context, dir, sourceRef string, platform Plat if progress != nil { fmt.Fprintln(progress, "olav: writing OCI image...") } - wrapped, finish := withImageProgress(img, progress) - defer finish() - return path.AppendImage(wrapped, layout.WithAnnotations(map[string]string{"org.opencontainers.image.ref.name": "olav"})) + rootDesc, err = plan.addImage(img) + if err != nil { + return err + } } - img, err := remote.Image(ref, options...) + + for _, blob := range plan.raw { + if err := path.WriteBlob(blob.digest, io.NopCloser(bytes.NewReader(blob.data))); err != nil { + return err + } + } + + counter := newProgressCounter(plan.fetchTotal(), progress) + fetcher, err := newBlobFetcher(ctx, ref, counter, progress) if err != nil { return err } - if progress != nil { - fmt.Fprintln(progress, "olav: writing OCI image...") + defer counter.finish() + for _, blob := range plan.fetch { + dst := filepath.Join(dir, "blobs", blob.digest.Algorithm, blob.digest.Hex) + if err := fetcher.fetchBlob(ctx, blob.digest, blob.size, dst); err != nil { + return err + } + } + + rootDesc.Annotations = map[string]string{"org.opencontainers.image.ref.name": "olav"} + if err := path.AppendDescriptor(rootDesc); err != nil { + return err } - wrapped, finish := withImageProgress(img, progress) - defer finish() - return path.AppendImage(wrapped, layout.WithAnnotations(map[string]string{"org.opencontainers.image.ref.name": "olav"})) + fetcher.cleanupPartials() + return nil } func writeDaemonLayout(ctx context.Context, dir, sourceRef string, progress io.Writer) error { diff --git a/internal/source/copy_e2e_test.go b/internal/source/copy_e2e_test.go new file mode 100644 index 0000000..fe49753 --- /dev/null +++ b/internal/source/copy_e2e_test.go @@ -0,0 +1,186 @@ +package source + +import ( + "context" + "io" + "log" + "net/http/httptest" + "net/url" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/google/go-containerregistry/pkg/name" + "github.com/google/go-containerregistry/pkg/registry" + "github.com/google/go-containerregistry/pkg/v1/random" + "github.com/google/go-containerregistry/pkg/v1/remote" +) + +func testRegistry(t *testing.T) string { + t.Helper() + server := httptest.NewServer(registry.New(registry.Logger(log.New(io.Discard, "", 0)))) + t.Cleanup(server.Close) + u, err := url.Parse(server.URL) + if err != nil { + t.Fatal(err) + } + return u.Host +} + +func TestResolveRemotePullsCompleteLayout(t *testing.T) { + t.Setenv("XDG_CACHE_HOME", t.TempDir()) + host := testRegistry(t) + + img, err := random.Image(4096, 3) + if err != nil { + t.Fatal(err) + } + ref, err := name.ParseReference(host + "/test/img:latest") + if err != nil { + t.Fatal(err) + } + if err := remote.Write(ref, img); err != nil { + t.Fatal(err) + } + + var progress strings.Builder + resolved, err := Resolve(context.Background(), Options{Input: "docker://" + host + "/test/img:latest", Progress: &progress}) + if err != nil { + t.Fatal(err) + } + + manifest, err := img.Manifest() + if err != nil { + t.Fatal(err) + } + digest, err := img.Digest() + if err != nil { + t.Fatal(err) + } + blobs := []string{digest.Hex, manifest.Config.Digest.Hex} + for _, layer := range manifest.Layers { + blobs = append(blobs, layer.Digest.Hex) + } + for _, hex := range blobs { + if _, err := os.Stat(filepath.Join(resolved.LocalPath, "blobs", "sha256", hex)); err != nil { + t.Fatalf("missing blob %s: %v", hex, err) + } + } + if !strings.Contains(progress.String(), "100%") { + t.Fatalf("progress output missing 100%%: %q", progress.String()) + } + + // Second resolve must hit the cache without touching blobs again. + progress.Reset() + cached, err := Resolve(context.Background(), Options{Input: "docker://" + host + "/test/img:latest", Progress: &progress}) + if err != nil { + t.Fatal(err) + } + if cached.LocalPath != resolved.LocalPath { + t.Fatalf("cache miss on second resolve: %q vs %q", cached.LocalPath, resolved.LocalPath) + } + if !strings.Contains(progress.String(), "using cached") { + t.Fatalf("expected cached resolve, got %q", progress.String()) + } +} + +func TestResolveRemotePlatformAllPullsIndex(t *testing.T) { + t.Setenv("XDG_CACHE_HOME", t.TempDir()) + host := testRegistry(t) + + idx, err := random.Index(2048, 2, 2) + if err != nil { + t.Fatal(err) + } + ref, err := name.ParseReference(host + "/test/multi:latest") + if err != nil { + t.Fatal(err) + } + if err := remote.WriteIndex(ref, idx); err != nil { + t.Fatal(err) + } + + resolved, err := Resolve(context.Background(), Options{Input: "docker://" + host + "/test/multi:latest", Platform: "all"}) + if err != nil { + t.Fatal(err) + } + + manifest, err := idx.IndexManifest() + if err != nil { + t.Fatal(err) + } + for _, child := range manifest.Manifests { + childImg, err := idx.Image(child.Digest) + if err != nil { + t.Fatal(err) + } + childManifest, err := childImg.Manifest() + if err != nil { + t.Fatal(err) + } + for _, layer := range childManifest.Layers { + if _, err := os.Stat(filepath.Join(resolved.LocalPath, "blobs", "sha256", layer.Digest.Hex)); err != nil { + t.Fatalf("missing child layer %s: %v", layer.Digest.Hex, err) + } + } + } +} + +func TestResolveRemoteResumesInterruptedPull(t *testing.T) { + cacheDir := t.TempDir() + t.Setenv("XDG_CACHE_HOME", cacheDir) + host := testRegistry(t) + + img, err := random.Image(8192, 1) + if err != nil { + t.Fatal(err) + } + ref, err := name.ParseReference(host + "/test/resume:latest") + if err != nil { + t.Fatal(err) + } + if err := remote.Write(ref, img); err != nil { + t.Fatal(err) + } + + // Simulate an interrupted earlier run: half the layer already sits in the + // partials dir under its digest name. + manifest, err := img.Manifest() + if err != nil { + t.Fatal(err) + } + layerDigest := manifest.Layers[0].Digest + layer, err := img.LayerByDigest(layerDigest) + if err != nil { + t.Fatal(err) + } + rc, err := layer.Compressed() + if err != nil { + t.Fatal(err) + } + full := make([]byte, manifest.Layers[0].Size) + if _, err := io.ReadFull(rc, full); err != nil { + t.Fatal(err) + } + partialDir := filepath.Join(cacheDir, "olav", "tmp", "partials") + if err := os.MkdirAll(partialDir, 0o755); err != nil { + t.Fatal(err) + } + partial := filepath.Join(partialDir, "sha256-"+layerDigest.Hex+".partial") + if err := os.WriteFile(partial, full[:len(full)/2], 0o600); err != nil { + t.Fatal(err) + } + + var progress strings.Builder + resolved, err := Resolve(context.Background(), Options{Input: "docker://" + host + "/test/resume:latest", Progress: &progress}) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(progress.String(), "resuming blob") { + t.Fatalf("expected resume message, got %q", progress.String()) + } + if _, err := os.Stat(filepath.Join(resolved.LocalPath, "blobs", "sha256", layerDigest.Hex)); err != nil { + t.Fatalf("missing resumed layer: %v", err) + } +} diff --git a/internal/source/fetch.go b/internal/source/fetch.go new file mode 100644 index 0000000..6d30f54 --- /dev/null +++ b/internal/source/fetch.go @@ -0,0 +1,412 @@ +package source + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "strings" + "sync/atomic" + + "github.com/google/go-containerregistry/pkg/name" + v1 "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/remote/transport" +) + +const fetchAttempts = 4 + +// blobFetcher downloads registry blobs with HTTP range requests. Partial +// downloads persist in the cache tmp area named by digest, so an interrupted +// pull resumes from the last received byte on the next run. +type blobFetcher struct { + blobURL func(v1.Hash) string + client *http.Client + partialDir string + counter *progressCounter + progress io.Writer + consumed []string +} + +func newBlobFetcher(ctx context.Context, ref name.Reference, counter *progressCounter, progress io.Writer) (*blobFetcher, error) { + repo := ref.Context() + auth, err := authKeychain().Resolve(repo.Registry) + if err != nil { + return nil, err + } + rt, err := transport.NewWithContext(ctx, repo.Registry, auth, http.DefaultTransport, []string{repo.Scope(transport.PullScope)}) + if err != nil { + return nil, err + } + root, err := cacheRoot() + if err != nil { + return nil, err + } + partialDir := filepath.Join(root, "tmp", "partials") + if err := os.MkdirAll(partialDir, 0o755); err != nil { + return nil, err + } + base := fmt.Sprintf("%s://%s/v2/%s/blobs/", repo.Registry.Scheme(), repo.RegistryStr(), repo.RepositoryStr()) + return &blobFetcher{ + blobURL: func(h v1.Hash) string { return base + h.String() }, + client: &http.Client{Transport: rt}, + partialDir: partialDir, + counter: counter, + progress: progress, + }, nil +} + +func (f *blobFetcher) fetchBlob(ctx context.Context, digest v1.Hash, size int64, dst string) error { + if digest.Algorithm != "sha256" { + return fmt.Errorf("unsupported digest algorithm %q", digest.Algorithm) + } + if info, err := os.Stat(dst); err == nil && info.Size() == size { + f.counter.add(size) + return nil + } + partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") + if info, err := os.Stat(partial); err == nil && info.Size() > 0 { + if f.progress != nil { + fmt.Fprintf(f.progress, "olav: resuming blob %s at %s\n", shortHex(digest), formatBytes(info.Size())) + } + f.counter.add(info.Size()) + } + var lastErr error + for attempt := 0; attempt < fetchAttempts; attempt++ { + if err := ctx.Err(); err != nil { + return err + } + lastErr = f.fetchOnce(ctx, digest, size, partial) + if lastErr == nil { + break + } + if !isRetryableFetchError(lastErr) { + return lastErr + } + } + if lastErr != nil { + return lastErr + } + if err := verifyBlobDigest(partial, digest); err != nil { + _ = os.Remove(partial) + return err + } + if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil { + return err + } + // Hardlink instead of rename: the partial keeps the bytes reachable for + // resume until the whole layout lands in the cache, in case a later blob + // of the same pull is interrupted. + if err := os.Link(partial, dst); err != nil { + return err + } + f.consumed = append(f.consumed, partial) + return nil +} + +func (f *blobFetcher) fetchOnce(ctx context.Context, digest v1.Hash, size int64, partial string) error { + offset := int64(0) + if info, err := os.Stat(partial); err == nil { + offset = info.Size() + } + if size > 0 && offset > size { + if err := os.Truncate(partial, 0); err != nil { + return err + } + offset = 0 + } + if size > 0 && offset == size { + return nil + } + req, err := http.NewRequestWithContext(ctx, http.MethodGet, f.blobURL(digest), nil) + if err != nil { + return err + } + if offset > 0 { + req.Header.Set("Range", fmt.Sprintf("bytes=%d-", offset)) + } + resp, err := f.client.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + + flags := os.O_CREATE | os.O_WRONLY | os.O_APPEND + switch { + case offset > 0 && resp.StatusCode == http.StatusPartialContent: + case resp.StatusCode == http.StatusOK: + // Server ignored the range request; start over with the full body. + flags = os.O_CREATE | os.O_WRONLY | os.O_TRUNC + case offset > 0 && resp.StatusCode == http.StatusRequestedRangeNotSatisfiable: + // Server refuses ranges outright; drop the partial and fetch in full. + if err := os.Truncate(partial, 0); err != nil { + return err + } + return f.fetchOnce(ctx, digest, size, partial) + default: + return &fetchStatusError{status: resp.StatusCode, url: req.URL.String()} + } + file, err := os.OpenFile(partial, flags, 0o600) + if err != nil { + return err + } + defer file.Close() + _, err = io.Copy(file, &countingReader{r: resp.Body, counter: f.counter}) + return err +} + +// cleanupPartials removes partial files whose bytes are fully linked into the +// written layout. Call only once the layout is complete. +func (f *blobFetcher) cleanupPartials() { + for _, partial := range f.consumed { + _ = os.Remove(partial) + } + f.consumed = nil +} + +func verifyBlobDigest(path string, digest v1.Hash) error { + file, err := os.Open(path) + if err != nil { + return err + } + defer file.Close() + h := sha256.New() + if _, err := io.Copy(h, file); err != nil { + return err + } + if got := hex.EncodeToString(h.Sum(nil)); got != digest.Hex { + return fmt.Errorf("digest mismatch for blob %s: downloaded sha256:%s", digest, got) + } + return nil +} + +type fetchStatusError struct { + status int + url string +} + +func (e *fetchStatusError) Error() string { + return fmt.Sprintf("fetch %s: unexpected status %d", e.url, e.status) +} + +func isRetryableFetchError(err error) bool { + var statusErr *fetchStatusError + if errors.As(err, &statusErr) { + return statusErr.status >= 500 || statusErr.status == http.StatusTooManyRequests + } + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + return false + } + // Network-level errors: the partial keeps what already arrived, retry + // resumes from there. + return true +} + +func shortHex(digest v1.Hash) string { + if len(digest.Hex) > 12 { + return digest.Hex[:12] + } + return digest.Hex +} + +type countingReader struct { + r io.Reader + counter *progressCounter +} + +func (c *countingReader) Read(p []byte) (int, error) { + n, err := c.r.Read(p) + c.counter.add(int64(n)) + return n, err +} + +type progressCounter struct { + total int64 + w io.Writer + complete atomic.Int64 + lastPct atomic.Int64 +} + +func newProgressCounter(total int64, w io.Writer) *progressCounter { + c := &progressCounter{total: total, w: w} + c.lastPct.Store(-1) + return c +} + +func (c *progressCounter) add(n int64) { + if c == nil || c.w == nil || n <= 0 { + return + } + complete := c.complete.Add(n) + if c.total <= 0 { + return + } + pct := complete * 100 / c.total + if pct > 100 { + pct = 100 + } + if c.lastPct.Swap(pct) == pct { + return + } + renderProgressLine(c.w, complete, c.total) +} + +func (c *progressCounter) finish() { + if c == nil || c.w == nil || c.lastPct.Load() < 0 { + return + } + renderProgressLine(c.w, c.complete.Load(), c.total) + fmt.Fprintln(c.w) +} + +func renderProgressLine(w io.Writer, complete, total int64) { + const width = 24 + if total <= 0 { + fmt.Fprintf(w, "\rolav: copying image blobs...") + return + } + if complete > total { + complete = total + } + filled := int(float64(complete) / float64(total) * width) + if filled > width { + filled = width + } + bar := strings.Repeat("=", filled) + strings.Repeat(" ", width-filled) + percent := int(float64(complete) / float64(total) * 100) + fmt.Fprintf(w, "\rolav: [%s] %3d%% %s / %s", bar, percent, formatBytes(complete), formatBytes(total)) +} + +func formatBytes(n int64) string { + const mib = 1 << 20 + if n < mib { + return fmt.Sprintf("%d B", n) + } + return fmt.Sprintf("%.1f MiB", float64(n)/mib) +} + +type rawBlob struct { + digest v1.Hash + data []byte +} + +type fetchBlobRef struct { + digest v1.Hash + size int64 +} + +// pullPlan splits an image or index into blobs already in memory (manifests, +// configs) and blobs worth fetching with the resumable fetcher (layers and +// other large artifacts). +type pullPlan struct { + raw []rawBlob + fetch []fetchBlobRef + seen map[string]bool +} + +func (p *pullPlan) mark(digest v1.Hash) bool { + if p.seen == nil { + p.seen = map[string]bool{} + } + if p.seen[digest.String()] { + return false + } + p.seen[digest.String()] = true + return true +} + +func (p *pullPlan) addRaw(digest v1.Hash, data []byte) { + if p.mark(digest) { + p.raw = append(p.raw, rawBlob{digest: digest, data: data}) + } +} + +func (p *pullPlan) addFetch(digest v1.Hash, size int64) { + if p.mark(digest) { + p.fetch = append(p.fetch, fetchBlobRef{digest: digest, size: size}) + } +} + +func (p *pullPlan) fetchTotal() int64 { + var total int64 + for _, blob := range p.fetch { + total += blob.size + } + return total +} + +func (p *pullPlan) addImage(img v1.Image) (v1.Descriptor, error) { + rawManifest, err := img.RawManifest() + if err != nil { + return v1.Descriptor{}, err + } + digest, err := img.Digest() + if err != nil { + return v1.Descriptor{}, err + } + mediaType, err := img.MediaType() + if err != nil { + return v1.Descriptor{}, err + } + manifest, err := img.Manifest() + if err != nil { + return v1.Descriptor{}, err + } + rawConfig, err := img.RawConfigFile() + if err != nil { + return v1.Descriptor{}, err + } + p.addRaw(digest, rawManifest) + p.addRaw(manifest.Config.Digest, rawConfig) + for _, layer := range manifest.Layers { + p.addFetch(layer.Digest, layer.Size) + } + return v1.Descriptor{MediaType: mediaType, Digest: digest, Size: int64(len(rawManifest))}, nil +} + +func (p *pullPlan) addIndex(idx v1.ImageIndex) (v1.Descriptor, error) { + rawManifest, err := idx.RawManifest() + if err != nil { + return v1.Descriptor{}, err + } + digest, err := idx.Digest() + if err != nil { + return v1.Descriptor{}, err + } + mediaType, err := idx.MediaType() + if err != nil { + return v1.Descriptor{}, err + } + p.addRaw(digest, rawManifest) + manifest, err := idx.IndexManifest() + if err != nil { + return v1.Descriptor{}, err + } + for _, child := range manifest.Manifests { + switch { + case child.MediaType.IsIndex(): + childIdx, err := idx.ImageIndex(child.Digest) + if err != nil { + return v1.Descriptor{}, err + } + if _, err := p.addIndex(childIdx); err != nil { + return v1.Descriptor{}, err + } + case child.MediaType.IsImage(): + childImg, err := idx.Image(child.Digest) + if err != nil { + return v1.Descriptor{}, err + } + if _, err := p.addImage(childImg); err != nil { + return v1.Descriptor{}, err + } + default: + p.addFetch(child.Digest, child.Size) + } + } + return v1.Descriptor{MediaType: mediaType, Digest: digest, Size: int64(len(rawManifest))}, nil +} diff --git a/internal/source/fetch_test.go b/internal/source/fetch_test.go new file mode 100644 index 0000000..72416dd --- /dev/null +++ b/internal/source/fetch_test.go @@ -0,0 +1,203 @@ +package source + +import ( + "bytes" + "context" + "crypto/rand" + "crypto/sha256" + "encoding/hex" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strconv" + "strings" + "testing" + + v1 "github.com/google/go-containerregistry/pkg/v1" +) + +type blobServer struct { + data []byte + ignoreRange bool + failures int + ranges []string +} + +func (s *blobServer) ServeHTTP(w http.ResponseWriter, r *http.Request) { + if s.failures > 0 { + s.failures-- + w.WriteHeader(http.StatusInternalServerError) + return + } + rangeHeader := r.Header.Get("Range") + s.ranges = append(s.ranges, rangeHeader) + if rangeHeader == "" || s.ignoreRange { + w.WriteHeader(http.StatusOK) + _, _ = w.Write(s.data) + return + } + offsetStr := strings.TrimSuffix(strings.TrimPrefix(rangeHeader, "bytes="), "-") + offset, err := strconv.ParseInt(offsetStr, 10, 64) + if err != nil || offset < 0 || offset >= int64(len(s.data)) { + w.WriteHeader(http.StatusRequestedRangeNotSatisfiable) + return + } + w.Header().Set("Content-Range", fmt.Sprintf("bytes %d-%d/%d", offset, len(s.data)-1, len(s.data))) + w.WriteHeader(http.StatusPartialContent) + _, _ = w.Write(s.data[offset:]) +} + +func testBlob(t *testing.T, size int) ([]byte, v1.Hash) { + t.Helper() + data := make([]byte, size) + if _, err := rand.Read(data); err != nil { + t.Fatal(err) + } + sum := sha256.Sum256(data) + return data, v1.Hash{Algorithm: "sha256", Hex: hex.EncodeToString(sum[:])} +} + +func testFetcher(t *testing.T, server *httptest.Server) *blobFetcher { + t.Helper() + return &blobFetcher{ + blobURL: func(h v1.Hash) string { return server.URL + "/blobs/" + h.String() }, + client: server.Client(), + partialDir: t.TempDir(), + counter: newProgressCounter(0, nil), + } +} + +func TestFetchBlobFresh(t *testing.T) { + data, digest := testBlob(t, 4096) + backend := &blobServer{data: data} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + dst := filepath.Join(t.TempDir(), "blob") + if err := f.fetchBlob(context.Background(), digest, int64(len(data)), dst); err != nil { + t.Fatal(err) + } + got, err := os.ReadFile(dst) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, data) { + t.Fatal("downloaded blob differs from source") + } + f.cleanupPartials() + entries, err := os.ReadDir(f.partialDir) + if err != nil { + t.Fatal(err) + } + if len(entries) != 0 { + t.Fatalf("cleanupPartials left %d files", len(entries)) + } +} + +func TestFetchBlobResumesFromPartial(t *testing.T) { + data, digest := testBlob(t, 4096) + backend := &blobServer{data: data} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + partial := filepath.Join(f.partialDir, "sha256-"+digest.Hex+".partial") + if err := os.WriteFile(partial, data[:1000], 0o600); err != nil { + t.Fatal(err) + } + dst := filepath.Join(t.TempDir(), "blob") + if err := f.fetchBlob(context.Background(), digest, int64(len(data)), dst); err != nil { + t.Fatal(err) + } + if len(backend.ranges) != 1 || backend.ranges[0] != "bytes=1000-" { + t.Fatalf("expected one resume request from byte 1000, got %v", backend.ranges) + } + got, err := os.ReadFile(dst) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, data) { + t.Fatal("resumed blob differs from source") + } +} + +func TestFetchBlobRestartsWhenRangeIgnored(t *testing.T) { + data, digest := testBlob(t, 4096) + backend := &blobServer{data: data, ignoreRange: true} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + partial := filepath.Join(f.partialDir, "sha256-"+digest.Hex+".partial") + // Poison the partial: server ignoring ranges must overwrite, not append. + if err := os.WriteFile(partial, []byte("garbage"), 0o600); err != nil { + t.Fatal(err) + } + dst := filepath.Join(t.TempDir(), "blob") + if err := f.fetchBlob(context.Background(), digest, int64(len(data)), dst); err != nil { + t.Fatal(err) + } + got, err := os.ReadFile(dst) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, data) { + t.Fatal("restarted blob differs from source") + } +} + +func TestFetchBlobDigestMismatch(t *testing.T) { + data, _ := testBlob(t, 1024) + _, wrongDigest := testBlob(t, 8) + backend := &blobServer{data: data} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + dst := filepath.Join(t.TempDir(), "blob") + err := f.fetchBlob(context.Background(), wrongDigest, int64(len(data)), dst) + if err == nil || !strings.Contains(err.Error(), "digest mismatch") { + t.Fatalf("expected digest mismatch error, got %v", err) + } + partial := filepath.Join(f.partialDir, "sha256-"+wrongDigest.Hex+".partial") + if _, statErr := os.Stat(partial); !os.IsNotExist(statErr) { + t.Fatal("mismatched partial should be removed") + } +} + +func TestFetchBlobRetriesServerErrors(t *testing.T) { + data, digest := testBlob(t, 1024) + backend := &blobServer{data: data, failures: 2} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + dst := filepath.Join(t.TempDir(), "blob") + if err := f.fetchBlob(context.Background(), digest, int64(len(data)), dst); err != nil { + t.Fatal(err) + } + got, err := os.ReadFile(dst) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, data) { + t.Fatal("blob differs after retries") + } +} + +func TestFetchBlobGivesUpAfterAttempts(t *testing.T) { + data, digest := testBlob(t, 1024) + backend := &blobServer{data: data, failures: fetchAttempts} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + dst := filepath.Join(t.TempDir(), "blob") + err := f.fetchBlob(context.Background(), digest, int64(len(data)), dst) + if err == nil || !strings.Contains(err.Error(), "unexpected status") { + t.Fatalf("expected status error after exhausted retries, got %v", err) + } +} diff --git a/internal/source/progress.go b/internal/source/progress.go deleted file mode 100644 index 5a8aa8d..0000000 --- a/internal/source/progress.go +++ /dev/null @@ -1,239 +0,0 @@ -package source - -import ( - "fmt" - "io" - "strings" - "sync" - "sync/atomic" - - v1 "github.com/google/go-containerregistry/pkg/v1" - "github.com/google/go-containerregistry/pkg/v1/types" -) - -type progressCounter struct { - total int64 - w io.Writer - complete atomic.Int64 - mu sync.Mutex - lastPct int - rendered bool -} - -func newProgressCounter(total int64, w io.Writer) *progressCounter { - return &progressCounter{total: total, w: w, lastPct: -1} -} - -func (c *progressCounter) add(n int64) { - if n <= 0 { - return - } - complete := c.complete.Add(n) - c.mu.Lock() - defer c.mu.Unlock() - if c.total <= 0 { - if !c.rendered { - c.rendered = true - renderProgressLine(c.w, complete, c.total) - } - return - } - if complete > c.total { - complete = c.total - } - pct := int(float64(complete) / float64(c.total) * 100) - if pct == c.lastPct { - return - } - c.lastPct = pct - c.rendered = true - renderProgressLine(c.w, complete, c.total) -} - -func (c *progressCounter) finish() { - c.mu.Lock() - defer c.mu.Unlock() - if !c.rendered { - return - } - renderProgressLine(c.w, c.complete.Load(), c.total) - fmt.Fprintln(c.w) -} - -func renderProgressLine(w io.Writer, complete, total int64) { - const width = 24 - if total <= 0 { - fmt.Fprintf(w, "\rolav: copying image blobs...") - return - } - if complete > total { - complete = total - } - filled := int(float64(complete) / float64(total) * width) - if filled > width { - filled = width - } - bar := strings.Repeat("=", filled) + strings.Repeat(" ", width-filled) - percent := int(float64(complete) / float64(total) * 100) - fmt.Fprintf(w, "\rolav: [%s] %3d%% %s / %s", bar, percent, formatBytes(complete), formatBytes(total)) -} - -func formatBytes(n int64) string { - const mib = 1 << 20 - if n < mib { - return fmt.Sprintf("%d B", n) - } - return fmt.Sprintf("%.1f MiB", float64(n)/mib) -} - -type progressReader struct { - io.ReadCloser - c *progressCounter -} - -func (r *progressReader) Read(p []byte) (int, error) { - n, err := r.ReadCloser.Read(p) - r.c.add(int64(n)) - return n, err -} - -type progressLayer struct { - v1.Layer - c *progressCounter -} - -func (l *progressLayer) Compressed() (io.ReadCloser, error) { - rc, err := l.Layer.Compressed() - if err != nil { - return nil, err - } - return &progressReader{ReadCloser: rc, c: l.c}, nil -} - -type progressImage struct { - v1.Image - c *progressCounter -} - -func (i *progressImage) Layers() ([]v1.Layer, error) { - layers, err := i.Image.Layers() - if err != nil { - return nil, err - } - wrapped := make([]v1.Layer, len(layers)) - for j, l := range layers { - wrapped[j] = &progressLayer{Layer: l, c: i.c} - } - return wrapped, nil -} - -type progressIndex struct { - idx v1.ImageIndex - c *progressCounter -} - -func (ix *progressIndex) MediaType() (types.MediaType, error) { return ix.idx.MediaType() } -func (ix *progressIndex) Digest() (v1.Hash, error) { return ix.idx.Digest() } -func (ix *progressIndex) Size() (int64, error) { return ix.idx.Size() } -func (ix *progressIndex) IndexManifest() (*v1.IndexManifest, error) { - return ix.idx.IndexManifest() -} -func (ix *progressIndex) RawManifest() ([]byte, error) { return ix.idx.RawManifest() } - -func (ix *progressIndex) Image(h v1.Hash) (v1.Image, error) { - img, err := ix.idx.Image(h) - if err != nil { - return nil, err - } - return &progressImage{Image: img, c: ix.c}, nil -} - -func (ix *progressIndex) ImageIndex(h v1.Hash) (v1.ImageIndex, error) { - child, err := ix.idx.ImageIndex(h) - if err != nil { - return nil, err - } - return &progressIndex{idx: child, c: ix.c}, nil -} - -// Layer mirrors the optional method go-containerregistry's layout writer uses -// for index children that are neither manifests nor indexes. -func (ix *progressIndex) Layer(h v1.Hash) (v1.Layer, error) { - wl, ok := ix.idx.(interface { - Layer(v1.Hash) (v1.Layer, error) - }) - if !ok { - return nil, fmt.Errorf("index does not expose blob %s", h) - } - l, err := wl.Layer(h) - if err != nil { - return nil, err - } - return &progressLayer{Layer: l, c: ix.c}, nil -} - -// imageBlobTotal sums only the blobs that stream through the wrapped layers; -// config and manifest blobs are written directly by the layout writer and -// never pass the progress reader, so they are left out of the total. -func imageBlobTotal(img v1.Image) int64 { - m, err := img.Manifest() - if err != nil { - return 0 - } - var total int64 - for _, l := range m.Layers { - total += l.Size - } - return total -} - -func indexBlobTotal(idx v1.ImageIndex) int64 { - m, err := idx.IndexManifest() - if err != nil { - return 0 - } - var ( - total int64 - failed bool - ) - for _, desc := range m.Manifests { - switch { - case desc.MediaType.IsIndex(): - child, err := idx.ImageIndex(desc.Digest) - if err != nil { - failed = true - continue - } - total += indexBlobTotal(child) - case desc.MediaType.IsImage(): - child, err := idx.Image(desc.Digest) - if err != nil { - failed = true - continue - } - total += imageBlobTotal(child) - default: - total += desc.Size - } - } - if failed { - return 0 - } - return total -} - -func withImageProgress(img v1.Image, progress io.Writer) (v1.Image, func()) { - if progress == nil { - return img, func() {} - } - c := newProgressCounter(imageBlobTotal(img), progress) - return &progressImage{Image: img, c: c}, c.finish -} - -func withIndexProgress(idx v1.ImageIndex, progress io.Writer) (v1.ImageIndex, func()) { - if progress == nil { - return idx, func() {} - } - c := newProgressCounter(indexBlobTotal(idx), progress) - return &progressIndex{idx: idx, c: c}, c.finish -} diff --git a/internal/source/progress_test.go b/internal/source/progress_test.go deleted file mode 100644 index dfe6cc9..0000000 --- a/internal/source/progress_test.go +++ /dev/null @@ -1,78 +0,0 @@ -package source - -import ( - "strings" - "testing" - - "github.com/google/go-containerregistry/pkg/v1/empty" - "github.com/google/go-containerregistry/pkg/v1/layout" - "github.com/google/go-containerregistry/pkg/v1/random" -) - -func TestImageProgressReachesTotal(t *testing.T) { - img, err := random.Image(2048, 3) - if err != nil { - t.Fatal(err) - } - total := imageBlobTotal(img) - if total <= 0 { - t.Fatalf("imageBlobTotal = %d, want > 0", total) - } - - var out strings.Builder - wrapped, finish := withImageProgress(img, &out) - path, err := layout.Write(t.TempDir(), empty.Index) - if err != nil { - t.Fatal(err) - } - if err := path.AppendImage(wrapped); err != nil { - t.Fatal(err) - } - finish() - - got := out.String() - if !strings.Contains(got, "100%") { - t.Fatalf("progress output missing 100%%: %q", got) - } - if !strings.HasSuffix(got, "\n") { - t.Fatalf("progress output does not end with newline: %q", got) - } -} - -func TestIndexProgressReachesTotal(t *testing.T) { - idx, err := random.Index(1024, 2, 2) - if err != nil { - t.Fatal(err) - } - total := indexBlobTotal(idx) - if total <= 0 { - t.Fatalf("indexBlobTotal = %d, want > 0", total) - } - - var out strings.Builder - wrapped, finish := withIndexProgress(idx, &out) - path, err := layout.Write(t.TempDir(), empty.Index) - if err != nil { - t.Fatal(err) - } - if err := path.AppendIndex(wrapped); err != nil { - t.Fatal(err) - } - finish() - - if got := out.String(); !strings.Contains(got, "100%") { - t.Fatalf("progress output missing 100%%: %q", got) - } -} - -func TestNilProgressWriterIsNoop(t *testing.T) { - img, err := random.Image(256, 1) - if err != nil { - t.Fatal(err) - } - wrapped, finish := withImageProgress(img, nil) - if wrapped != img { - t.Fatal("nil progress writer should return the image unwrapped") - } - finish() -} From 451b5c93289eeea2dda8914393f40263aca3c4d7 Mon Sep 17 00:00:00 2001 From: Marcin Konowalczyk Date: Wed, 15 Jul 2026 12:57:17 +0200 Subject: [PATCH 2/8] feat: apply suggestions from code review Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Signed-off-by: Marcin Konowalczyk --- internal/source/copy_e2e_test.go | 1 + internal/source/fetch.go | 23 ++++++++++++++++++++++- 2 files changed, 23 insertions(+), 1 deletion(-) diff --git a/internal/source/copy_e2e_test.go b/internal/source/copy_e2e_test.go index fe49753..9c9cad0 100644 --- a/internal/source/copy_e2e_test.go +++ b/internal/source/copy_e2e_test.go @@ -159,6 +159,7 @@ func TestResolveRemoteResumesInterruptedPull(t *testing.T) { if err != nil { t.Fatal(err) } + defer rc.Close() full := make([]byte, manifest.Layers[0].Size) if _, err := io.ReadFull(rc, full); err != nil { t.Fatal(err) diff --git a/internal/source/fetch.go b/internal/source/fetch.go index 6d30f54..2fb3057 100644 --- a/internal/source/fetch.go +++ b/internal/source/fetch.go @@ -102,7 +102,24 @@ func (f *blobFetcher) fetchBlob(ctx context.Context, digest v1.Hash, size int64, // resume until the whole layout lands in the cache, in case a later blob // of the same pull is interrupted. if err := os.Link(partial, dst); err != nil { - return err + // Fall back to copying if hardlinks aren't supported. + src, openErr := os.Open(partial) + if openErr != nil { + return openErr + } + defer src.Close() + + out, createErr := os.OpenFile(dst, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644) + if createErr != nil { + return createErr + } + if _, copyErr := io.Copy(out, src); copyErr != nil { + _ = out.Close() + return copyErr + } + if closeErr := out.Close(); closeErr != nil { + return closeErr + } } f.consumed = append(f.consumed, partial) return nil @@ -243,6 +260,10 @@ func (c *progressCounter) add(n int64) { } complete := c.complete.Add(n) if c.total <= 0 { + if c.lastPct.Swap(0) == 0 { + return + } + renderProgressLine(c.w, complete, c.total) return } pct := complete * 100 / c.total From a6702d98e5e908af6305626f9f963b6c4184fae7 Mon Sep 17 00:00:00 2001 From: lczyk Date: Wed, 15 Jul 2026 13:02:18 +0200 Subject: [PATCH 3/8] fix(source): double-counted progress on resume restart --- internal/source/fetch.go | 14 ++++++++++++++ internal/source/fetch_test.go | 24 ++++++++++++++++++++++++ 2 files changed, 38 insertions(+) diff --git a/internal/source/fetch.go b/internal/source/fetch.go index 2fb3057..ed471e6 100644 --- a/internal/source/fetch.go +++ b/internal/source/fetch.go @@ -131,6 +131,7 @@ func (f *blobFetcher) fetchOnce(ctx context.Context, digest v1.Hash, size int64, offset = info.Size() } if size > 0 && offset > size { + f.counter.sub(offset) if err := os.Truncate(partial, 0); err != nil { return err } @@ -157,9 +158,11 @@ func (f *blobFetcher) fetchOnce(ctx context.Context, digest v1.Hash, size int64, case offset > 0 && resp.StatusCode == http.StatusPartialContent: case resp.StatusCode == http.StatusOK: // Server ignored the range request; start over with the full body. + f.counter.sub(offset) flags = os.O_CREATE | os.O_WRONLY | os.O_TRUNC case offset > 0 && resp.StatusCode == http.StatusRequestedRangeNotSatisfiable: // Server refuses ranges outright; drop the partial and fetch in full. + f.counter.sub(offset) if err := os.Truncate(partial, 0); err != nil { return err } @@ -276,6 +279,17 @@ func (c *progressCounter) add(n int64) { renderProgressLine(c.w, complete, c.total) } +// sub rolls back bytes that were counted but then discarded, e.g. a resumed +// partial the server declines to honour. The next add re-renders at the +// corrected, lower percentage. +func (c *progressCounter) sub(n int64) { + if c == nil || c.w == nil || n <= 0 { + return + } + c.complete.Add(-n) + c.lastPct.Store(-1) +} + func (c *progressCounter) finish() { if c == nil || c.w == nil || c.lastPct.Load() < 0 { return diff --git a/internal/source/fetch_test.go b/internal/source/fetch_test.go index 72416dd..00f17db 100644 --- a/internal/source/fetch_test.go +++ b/internal/source/fetch_test.go @@ -7,6 +7,7 @@ import ( "crypto/sha256" "encoding/hex" "fmt" + "io" "net/http" "net/http/httptest" "os" @@ -149,6 +150,29 @@ func TestFetchBlobRestartsWhenRangeIgnored(t *testing.T) { } } +func TestFetchBlobRestartDoesNotDoubleCountProgress(t *testing.T) { + data, digest := testBlob(t, 4096) + backend := &blobServer{data: data, ignoreRange: true} + server := httptest.NewServer(backend) + defer server.Close() + + f := testFetcher(t, server) + f.counter = newProgressCounter(int64(len(data)), io.Discard) + partial := filepath.Join(f.partialDir, "sha256-"+digest.Hex+".partial") + // Pre-existing partial gets counted, then discarded when the server ignores + // the range: those bytes must be rolled back, not double-counted. + if err := os.WriteFile(partial, []byte("garbage"), 0o600); err != nil { + t.Fatal(err) + } + dst := filepath.Join(t.TempDir(), "blob") + if err := f.fetchBlob(context.Background(), digest, int64(len(data)), dst); err != nil { + t.Fatal(err) + } + if got := f.counter.complete.Load(); got != int64(len(data)) { + t.Fatalf("counter = %d, want %d (resumed bytes double-counted)", got, len(data)) + } +} + func TestFetchBlobDigestMismatch(t *testing.T) { data, _ := testBlob(t, 1024) _, wrongDigest := testBlob(t, 8) From 05dacf1498969ece90fa9526fba3846c18654d95 Mon Sep 17 00:00:00 2001 From: lczyk Date: Wed, 15 Jul 2026 13:23:58 +0200 Subject: [PATCH 4/8] chore!: readd progress --- internal/source/progress.go | 239 +++++++++++++++++++++++++++++++ internal/source/progress_test.go | 78 ++++++++++ 2 files changed, 317 insertions(+) create mode 100644 internal/source/progress.go create mode 100644 internal/source/progress_test.go diff --git a/internal/source/progress.go b/internal/source/progress.go new file mode 100644 index 0000000..5a8aa8d --- /dev/null +++ b/internal/source/progress.go @@ -0,0 +1,239 @@ +package source + +import ( + "fmt" + "io" + "strings" + "sync" + "sync/atomic" + + v1 "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/types" +) + +type progressCounter struct { + total int64 + w io.Writer + complete atomic.Int64 + mu sync.Mutex + lastPct int + rendered bool +} + +func newProgressCounter(total int64, w io.Writer) *progressCounter { + return &progressCounter{total: total, w: w, lastPct: -1} +} + +func (c *progressCounter) add(n int64) { + if n <= 0 { + return + } + complete := c.complete.Add(n) + c.mu.Lock() + defer c.mu.Unlock() + if c.total <= 0 { + if !c.rendered { + c.rendered = true + renderProgressLine(c.w, complete, c.total) + } + return + } + if complete > c.total { + complete = c.total + } + pct := int(float64(complete) / float64(c.total) * 100) + if pct == c.lastPct { + return + } + c.lastPct = pct + c.rendered = true + renderProgressLine(c.w, complete, c.total) +} + +func (c *progressCounter) finish() { + c.mu.Lock() + defer c.mu.Unlock() + if !c.rendered { + return + } + renderProgressLine(c.w, c.complete.Load(), c.total) + fmt.Fprintln(c.w) +} + +func renderProgressLine(w io.Writer, complete, total int64) { + const width = 24 + if total <= 0 { + fmt.Fprintf(w, "\rolav: copying image blobs...") + return + } + if complete > total { + complete = total + } + filled := int(float64(complete) / float64(total) * width) + if filled > width { + filled = width + } + bar := strings.Repeat("=", filled) + strings.Repeat(" ", width-filled) + percent := int(float64(complete) / float64(total) * 100) + fmt.Fprintf(w, "\rolav: [%s] %3d%% %s / %s", bar, percent, formatBytes(complete), formatBytes(total)) +} + +func formatBytes(n int64) string { + const mib = 1 << 20 + if n < mib { + return fmt.Sprintf("%d B", n) + } + return fmt.Sprintf("%.1f MiB", float64(n)/mib) +} + +type progressReader struct { + io.ReadCloser + c *progressCounter +} + +func (r *progressReader) Read(p []byte) (int, error) { + n, err := r.ReadCloser.Read(p) + r.c.add(int64(n)) + return n, err +} + +type progressLayer struct { + v1.Layer + c *progressCounter +} + +func (l *progressLayer) Compressed() (io.ReadCloser, error) { + rc, err := l.Layer.Compressed() + if err != nil { + return nil, err + } + return &progressReader{ReadCloser: rc, c: l.c}, nil +} + +type progressImage struct { + v1.Image + c *progressCounter +} + +func (i *progressImage) Layers() ([]v1.Layer, error) { + layers, err := i.Image.Layers() + if err != nil { + return nil, err + } + wrapped := make([]v1.Layer, len(layers)) + for j, l := range layers { + wrapped[j] = &progressLayer{Layer: l, c: i.c} + } + return wrapped, nil +} + +type progressIndex struct { + idx v1.ImageIndex + c *progressCounter +} + +func (ix *progressIndex) MediaType() (types.MediaType, error) { return ix.idx.MediaType() } +func (ix *progressIndex) Digest() (v1.Hash, error) { return ix.idx.Digest() } +func (ix *progressIndex) Size() (int64, error) { return ix.idx.Size() } +func (ix *progressIndex) IndexManifest() (*v1.IndexManifest, error) { + return ix.idx.IndexManifest() +} +func (ix *progressIndex) RawManifest() ([]byte, error) { return ix.idx.RawManifest() } + +func (ix *progressIndex) Image(h v1.Hash) (v1.Image, error) { + img, err := ix.idx.Image(h) + if err != nil { + return nil, err + } + return &progressImage{Image: img, c: ix.c}, nil +} + +func (ix *progressIndex) ImageIndex(h v1.Hash) (v1.ImageIndex, error) { + child, err := ix.idx.ImageIndex(h) + if err != nil { + return nil, err + } + return &progressIndex{idx: child, c: ix.c}, nil +} + +// Layer mirrors the optional method go-containerregistry's layout writer uses +// for index children that are neither manifests nor indexes. +func (ix *progressIndex) Layer(h v1.Hash) (v1.Layer, error) { + wl, ok := ix.idx.(interface { + Layer(v1.Hash) (v1.Layer, error) + }) + if !ok { + return nil, fmt.Errorf("index does not expose blob %s", h) + } + l, err := wl.Layer(h) + if err != nil { + return nil, err + } + return &progressLayer{Layer: l, c: ix.c}, nil +} + +// imageBlobTotal sums only the blobs that stream through the wrapped layers; +// config and manifest blobs are written directly by the layout writer and +// never pass the progress reader, so they are left out of the total. +func imageBlobTotal(img v1.Image) int64 { + m, err := img.Manifest() + if err != nil { + return 0 + } + var total int64 + for _, l := range m.Layers { + total += l.Size + } + return total +} + +func indexBlobTotal(idx v1.ImageIndex) int64 { + m, err := idx.IndexManifest() + if err != nil { + return 0 + } + var ( + total int64 + failed bool + ) + for _, desc := range m.Manifests { + switch { + case desc.MediaType.IsIndex(): + child, err := idx.ImageIndex(desc.Digest) + if err != nil { + failed = true + continue + } + total += indexBlobTotal(child) + case desc.MediaType.IsImage(): + child, err := idx.Image(desc.Digest) + if err != nil { + failed = true + continue + } + total += imageBlobTotal(child) + default: + total += desc.Size + } + } + if failed { + return 0 + } + return total +} + +func withImageProgress(img v1.Image, progress io.Writer) (v1.Image, func()) { + if progress == nil { + return img, func() {} + } + c := newProgressCounter(imageBlobTotal(img), progress) + return &progressImage{Image: img, c: c}, c.finish +} + +func withIndexProgress(idx v1.ImageIndex, progress io.Writer) (v1.ImageIndex, func()) { + if progress == nil { + return idx, func() {} + } + c := newProgressCounter(indexBlobTotal(idx), progress) + return &progressIndex{idx: idx, c: c}, c.finish +} diff --git a/internal/source/progress_test.go b/internal/source/progress_test.go new file mode 100644 index 0000000..dfe6cc9 --- /dev/null +++ b/internal/source/progress_test.go @@ -0,0 +1,78 @@ +package source + +import ( + "strings" + "testing" + + "github.com/google/go-containerregistry/pkg/v1/empty" + "github.com/google/go-containerregistry/pkg/v1/layout" + "github.com/google/go-containerregistry/pkg/v1/random" +) + +func TestImageProgressReachesTotal(t *testing.T) { + img, err := random.Image(2048, 3) + if err != nil { + t.Fatal(err) + } + total := imageBlobTotal(img) + if total <= 0 { + t.Fatalf("imageBlobTotal = %d, want > 0", total) + } + + var out strings.Builder + wrapped, finish := withImageProgress(img, &out) + path, err := layout.Write(t.TempDir(), empty.Index) + if err != nil { + t.Fatal(err) + } + if err := path.AppendImage(wrapped); err != nil { + t.Fatal(err) + } + finish() + + got := out.String() + if !strings.Contains(got, "100%") { + t.Fatalf("progress output missing 100%%: %q", got) + } + if !strings.HasSuffix(got, "\n") { + t.Fatalf("progress output does not end with newline: %q", got) + } +} + +func TestIndexProgressReachesTotal(t *testing.T) { + idx, err := random.Index(1024, 2, 2) + if err != nil { + t.Fatal(err) + } + total := indexBlobTotal(idx) + if total <= 0 { + t.Fatalf("indexBlobTotal = %d, want > 0", total) + } + + var out strings.Builder + wrapped, finish := withIndexProgress(idx, &out) + path, err := layout.Write(t.TempDir(), empty.Index) + if err != nil { + t.Fatal(err) + } + if err := path.AppendIndex(wrapped); err != nil { + t.Fatal(err) + } + finish() + + if got := out.String(); !strings.Contains(got, "100%") { + t.Fatalf("progress output missing 100%%: %q", got) + } +} + +func TestNilProgressWriterIsNoop(t *testing.T) { + img, err := random.Image(256, 1) + if err != nil { + t.Fatal(err) + } + wrapped, finish := withImageProgress(img, nil) + if wrapped != img { + t.Fatal("nil progress writer should return the image unwrapped") + } + finish() +} From 692b0fd0a78f0482d74138ca913be61fb602c3de Mon Sep 17 00:00:00 2001 From: lczyk Date: Wed, 15 Jul 2026 13:28:10 +0200 Subject: [PATCH 5/8] refactor(source): reuse progress counter, drop layer wrappers --- internal/source/fetch.go | 93 ----------------- internal/source/progress.go | 174 ++++--------------------------- internal/source/progress_test.go | 106 +++++++++---------- 3 files changed, 75 insertions(+), 298 deletions(-) diff --git a/internal/source/fetch.go b/internal/source/fetch.go index ed471e6..8290e7b 100644 --- a/internal/source/fetch.go +++ b/internal/source/fetch.go @@ -10,8 +10,6 @@ import ( "net/http" "os" "path/filepath" - "strings" - "sync/atomic" "github.com/google/go-containerregistry/pkg/name" v1 "github.com/google/go-containerregistry/pkg/v1" @@ -233,97 +231,6 @@ func shortHex(digest v1.Hash) string { return digest.Hex } -type countingReader struct { - r io.Reader - counter *progressCounter -} - -func (c *countingReader) Read(p []byte) (int, error) { - n, err := c.r.Read(p) - c.counter.add(int64(n)) - return n, err -} - -type progressCounter struct { - total int64 - w io.Writer - complete atomic.Int64 - lastPct atomic.Int64 -} - -func newProgressCounter(total int64, w io.Writer) *progressCounter { - c := &progressCounter{total: total, w: w} - c.lastPct.Store(-1) - return c -} - -func (c *progressCounter) add(n int64) { - if c == nil || c.w == nil || n <= 0 { - return - } - complete := c.complete.Add(n) - if c.total <= 0 { - if c.lastPct.Swap(0) == 0 { - return - } - renderProgressLine(c.w, complete, c.total) - return - } - pct := complete * 100 / c.total - if pct > 100 { - pct = 100 - } - if c.lastPct.Swap(pct) == pct { - return - } - renderProgressLine(c.w, complete, c.total) -} - -// sub rolls back bytes that were counted but then discarded, e.g. a resumed -// partial the server declines to honour. The next add re-renders at the -// corrected, lower percentage. -func (c *progressCounter) sub(n int64) { - if c == nil || c.w == nil || n <= 0 { - return - } - c.complete.Add(-n) - c.lastPct.Store(-1) -} - -func (c *progressCounter) finish() { - if c == nil || c.w == nil || c.lastPct.Load() < 0 { - return - } - renderProgressLine(c.w, c.complete.Load(), c.total) - fmt.Fprintln(c.w) -} - -func renderProgressLine(w io.Writer, complete, total int64) { - const width = 24 - if total <= 0 { - fmt.Fprintf(w, "\rolav: copying image blobs...") - return - } - if complete > total { - complete = total - } - filled := int(float64(complete) / float64(total) * width) - if filled > width { - filled = width - } - bar := strings.Repeat("=", filled) + strings.Repeat(" ", width-filled) - percent := int(float64(complete) / float64(total) * 100) - fmt.Fprintf(w, "\rolav: [%s] %3d%% %s / %s", bar, percent, formatBytes(complete), formatBytes(total)) -} - -func formatBytes(n int64) string { - const mib = 1 << 20 - if n < mib { - return fmt.Sprintf("%d B", n) - } - return fmt.Sprintf("%.1f MiB", float64(n)/mib) -} - type rawBlob struct { digest v1.Hash data []byte diff --git a/internal/source/progress.go b/internal/source/progress.go index 5a8aa8d..2bef31d 100644 --- a/internal/source/progress.go +++ b/internal/source/progress.go @@ -6,9 +6,6 @@ import ( "strings" "sync" "sync/atomic" - - v1 "github.com/google/go-containerregistry/pkg/v1" - "github.com/google/go-containerregistry/pkg/v1/types" ) type progressCounter struct { @@ -25,7 +22,7 @@ func newProgressCounter(total int64, w io.Writer) *progressCounter { } func (c *progressCounter) add(n int64) { - if n <= 0 { + if c == nil || c.w == nil || n <= 0 { return } complete := c.complete.Add(n) @@ -50,7 +47,23 @@ func (c *progressCounter) add(n int64) { renderProgressLine(c.w, complete, c.total) } +// sub rolls back bytes that were counted but then discarded, e.g. a resumed +// partial the server declines to honour. The next add re-renders at the +// corrected, lower percentage. +func (c *progressCounter) sub(n int64) { + if c == nil || c.w == nil || n <= 0 { + return + } + c.complete.Add(-n) + c.mu.Lock() + c.lastPct = -1 + c.mu.Unlock() +} + func (c *progressCounter) finish() { + if c == nil || c.w == nil { + return + } c.mu.Lock() defer c.mu.Unlock() if !c.rendered { @@ -86,154 +99,13 @@ func formatBytes(n int64) string { return fmt.Sprintf("%.1f MiB", float64(n)/mib) } -type progressReader struct { - io.ReadCloser - c *progressCounter +type countingReader struct { + r io.Reader + counter *progressCounter } -func (r *progressReader) Read(p []byte) (int, error) { - n, err := r.ReadCloser.Read(p) - r.c.add(int64(n)) +func (c *countingReader) Read(p []byte) (int, error) { + n, err := c.r.Read(p) + c.counter.add(int64(n)) return n, err } - -type progressLayer struct { - v1.Layer - c *progressCounter -} - -func (l *progressLayer) Compressed() (io.ReadCloser, error) { - rc, err := l.Layer.Compressed() - if err != nil { - return nil, err - } - return &progressReader{ReadCloser: rc, c: l.c}, nil -} - -type progressImage struct { - v1.Image - c *progressCounter -} - -func (i *progressImage) Layers() ([]v1.Layer, error) { - layers, err := i.Image.Layers() - if err != nil { - return nil, err - } - wrapped := make([]v1.Layer, len(layers)) - for j, l := range layers { - wrapped[j] = &progressLayer{Layer: l, c: i.c} - } - return wrapped, nil -} - -type progressIndex struct { - idx v1.ImageIndex - c *progressCounter -} - -func (ix *progressIndex) MediaType() (types.MediaType, error) { return ix.idx.MediaType() } -func (ix *progressIndex) Digest() (v1.Hash, error) { return ix.idx.Digest() } -func (ix *progressIndex) Size() (int64, error) { return ix.idx.Size() } -func (ix *progressIndex) IndexManifest() (*v1.IndexManifest, error) { - return ix.idx.IndexManifest() -} -func (ix *progressIndex) RawManifest() ([]byte, error) { return ix.idx.RawManifest() } - -func (ix *progressIndex) Image(h v1.Hash) (v1.Image, error) { - img, err := ix.idx.Image(h) - if err != nil { - return nil, err - } - return &progressImage{Image: img, c: ix.c}, nil -} - -func (ix *progressIndex) ImageIndex(h v1.Hash) (v1.ImageIndex, error) { - child, err := ix.idx.ImageIndex(h) - if err != nil { - return nil, err - } - return &progressIndex{idx: child, c: ix.c}, nil -} - -// Layer mirrors the optional method go-containerregistry's layout writer uses -// for index children that are neither manifests nor indexes. -func (ix *progressIndex) Layer(h v1.Hash) (v1.Layer, error) { - wl, ok := ix.idx.(interface { - Layer(v1.Hash) (v1.Layer, error) - }) - if !ok { - return nil, fmt.Errorf("index does not expose blob %s", h) - } - l, err := wl.Layer(h) - if err != nil { - return nil, err - } - return &progressLayer{Layer: l, c: ix.c}, nil -} - -// imageBlobTotal sums only the blobs that stream through the wrapped layers; -// config and manifest blobs are written directly by the layout writer and -// never pass the progress reader, so they are left out of the total. -func imageBlobTotal(img v1.Image) int64 { - m, err := img.Manifest() - if err != nil { - return 0 - } - var total int64 - for _, l := range m.Layers { - total += l.Size - } - return total -} - -func indexBlobTotal(idx v1.ImageIndex) int64 { - m, err := idx.IndexManifest() - if err != nil { - return 0 - } - var ( - total int64 - failed bool - ) - for _, desc := range m.Manifests { - switch { - case desc.MediaType.IsIndex(): - child, err := idx.ImageIndex(desc.Digest) - if err != nil { - failed = true - continue - } - total += indexBlobTotal(child) - case desc.MediaType.IsImage(): - child, err := idx.Image(desc.Digest) - if err != nil { - failed = true - continue - } - total += imageBlobTotal(child) - default: - total += desc.Size - } - } - if failed { - return 0 - } - return total -} - -func withImageProgress(img v1.Image, progress io.Writer) (v1.Image, func()) { - if progress == nil { - return img, func() {} - } - c := newProgressCounter(imageBlobTotal(img), progress) - return &progressImage{Image: img, c: c}, c.finish -} - -func withIndexProgress(idx v1.ImageIndex, progress io.Writer) (v1.ImageIndex, func()) { - if progress == nil { - return idx, func() {} - } - c := newProgressCounter(indexBlobTotal(idx), progress) - return &progressIndex{idx: idx, c: c}, c.finish -} diff --git a/internal/source/progress_test.go b/internal/source/progress_test.go index dfe6cc9..d359c0d 100644 --- a/internal/source/progress_test.go +++ b/internal/source/progress_test.go @@ -1,78 +1,76 @@ package source import ( + "bytes" + "io" "strings" "testing" - - "github.com/google/go-containerregistry/pkg/v1/empty" - "github.com/google/go-containerregistry/pkg/v1/layout" - "github.com/google/go-containerregistry/pkg/v1/random" ) -func TestImageProgressReachesTotal(t *testing.T) { - img, err := random.Image(2048, 3) - if err != nil { - t.Fatal(err) - } - total := imageBlobTotal(img) - if total <= 0 { - t.Fatalf("imageBlobTotal = %d, want > 0", total) - } - - var out strings.Builder - wrapped, finish := withImageProgress(img, &out) - path, err := layout.Write(t.TempDir(), empty.Index) - if err != nil { - t.Fatal(err) +func TestProgressCounterReachesTotal(t *testing.T) { + var out bytes.Buffer + c := newProgressCounter(100, &out) + c.add(40) + c.add(60) + c.finish() + if !strings.Contains(out.String(), "100%") { + t.Fatalf("counter never reached 100%%: %q", out.String()) } - if err := path.AppendImage(wrapped); err != nil { - t.Fatal(err) + if !strings.HasSuffix(out.String(), "\n") { + t.Fatalf("finish should end with a newline: %q", out.String()) } - finish() +} - got := out.String() - if !strings.Contains(got, "100%") { - t.Fatalf("progress output missing 100%%: %q", got) +func TestProgressCounterClampsOvershoot(t *testing.T) { + var out bytes.Buffer + c := newProgressCounter(100, &out) + c.add(250) + // The rendered line is clamped to the total; the raw count is left intact. + if strings.Contains(out.String(), "250") { + t.Fatalf("render should clamp to total, got %q", out.String()) } - if !strings.HasSuffix(got, "\n") { - t.Fatalf("progress output does not end with newline: %q", got) + if !strings.Contains(out.String(), "100%") { + t.Fatalf("render should show 100%%, got %q", out.String()) } } -func TestIndexProgressReachesTotal(t *testing.T) { - idx, err := random.Index(1024, 2, 2) - if err != nil { - t.Fatal(err) - } - total := indexBlobTotal(idx) - if total <= 0 { - t.Fatalf("indexBlobTotal = %d, want > 0", total) +func TestProgressCounterSubRollsBack(t *testing.T) { + var out bytes.Buffer + c := newProgressCounter(100, &out) + c.add(30) + c.sub(30) + c.add(100) + if got := c.complete.Load(); got != 100 { + t.Fatalf("complete = %d, want 100 after rollback", got) } +} - var out strings.Builder - wrapped, finish := withIndexProgress(idx, &out) - path, err := layout.Write(t.TempDir(), empty.Index) - if err != nil { - t.Fatal(err) +func TestProgressCounterUnknownTotal(t *testing.T) { + var out bytes.Buffer + c := newProgressCounter(0, &out) + c.add(10) + c.add(10) + if !strings.Contains(out.String(), "copying image blobs") { + t.Fatalf("unknown total should render the fallback line, got %q", out.String()) } - if err := path.AppendIndex(wrapped); err != nil { - t.Fatal(err) - } - finish() +} - if got := out.String(); !strings.Contains(got, "100%") { - t.Fatalf("progress output missing 100%%: %q", got) - } +func TestProgressCounterNilWriterIsNoop(t *testing.T) { + c := newProgressCounter(100, nil) + c.add(50) + c.sub(10) + c.finish() // must not panic writing to a nil writer } -func TestNilProgressWriterIsNoop(t *testing.T) { - img, err := random.Image(256, 1) - if err != nil { +func TestCountingReaderCountsBytes(t *testing.T) { + var out bytes.Buffer + c := newProgressCounter(4, &out) + r := &countingReader{r: strings.NewReader("data"), counter: c} + buf := make([]byte, 4) + if _, err := io.ReadFull(r, buf); err != nil { t.Fatal(err) } - wrapped, finish := withImageProgress(img, nil) - if wrapped != img { - t.Fatal("nil progress writer should return the image unwrapped") + if got := c.complete.Load(); got != 4 { + t.Fatalf("counter = %d, want 4", got) } - finish() } From f6b080e14fedf2eb0afd9396174f2a01e05167f6 Mon Sep 17 00:00:00 2001 From: lczyk Date: Wed, 15 Jul 2026 14:32:07 +0200 Subject: [PATCH 6/8] refactor(source): compact blob fetch and pull plan --- internal/source/fetch.go | 63 +++++++++++++++++----------------------- 1 file changed, 26 insertions(+), 37 deletions(-) diff --git a/internal/source/fetch.go b/internal/source/fetch.go index 8290e7b..bcf004c 100644 --- a/internal/source/fetch.go +++ b/internal/source/fetch.go @@ -13,6 +13,7 @@ import ( "github.com/google/go-containerregistry/pkg/name" v1 "github.com/google/go-containerregistry/pkg/v1" + "github.com/google/go-containerregistry/pkg/v1/partial" "github.com/google/go-containerregistry/pkg/v1/remote/transport" ) @@ -69,7 +70,7 @@ func (f *blobFetcher) fetchBlob(ctx context.Context, digest v1.Hash, size int64, partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") if info, err := os.Stat(partial); err == nil && info.Size() > 0 { if f.progress != nil { - fmt.Fprintf(f.progress, "olav: resuming blob %s at %s\n", shortHex(digest), formatBytes(info.Size())) + fmt.Fprintf(f.progress, "olav: resuming blob %s at %s\n", digest.Hex[:12], formatBytes(info.Size())) } f.counter.add(info.Size()) } @@ -101,22 +102,8 @@ func (f *blobFetcher) fetchBlob(ctx context.Context, digest v1.Hash, size int64, // of the same pull is interrupted. if err := os.Link(partial, dst); err != nil { // Fall back to copying if hardlinks aren't supported. - src, openErr := os.Open(partial) - if openErr != nil { - return openErr - } - defer src.Close() - - out, createErr := os.OpenFile(dst, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644) - if createErr != nil { - return createErr - } - if _, copyErr := io.Copy(out, src); copyErr != nil { - _ = out.Close() - return copyErr - } - if closeErr := out.Close(); closeErr != nil { - return closeErr + if err := copyFile(partial, dst); err != nil { + return err } } f.consumed = append(f.consumed, partial) @@ -224,11 +211,21 @@ func isRetryableFetchError(err error) bool { return true } -func shortHex(digest v1.Hash) string { - if len(digest.Hex) > 12 { - return digest.Hex[:12] +func copyFile(src, dst string) error { + in, err := os.Open(src) + if err != nil { + return err + } + defer in.Close() + out, err := os.Create(dst) + if err != nil { + return err + } + if _, err := io.Copy(out, in); err != nil { + _ = out.Close() + return err } - return digest.Hex + return out.Close() } type rawBlob struct { @@ -282,15 +279,11 @@ func (p *pullPlan) fetchTotal() int64 { } func (p *pullPlan) addImage(img v1.Image) (v1.Descriptor, error) { - rawManifest, err := img.RawManifest() - if err != nil { - return v1.Descriptor{}, err - } - digest, err := img.Digest() + desc, err := partial.Descriptor(img) if err != nil { return v1.Descriptor{}, err } - mediaType, err := img.MediaType() + rawManifest, err := img.RawManifest() if err != nil { return v1.Descriptor{}, err } @@ -302,28 +295,24 @@ func (p *pullPlan) addImage(img v1.Image) (v1.Descriptor, error) { if err != nil { return v1.Descriptor{}, err } - p.addRaw(digest, rawManifest) + p.addRaw(desc.Digest, rawManifest) p.addRaw(manifest.Config.Digest, rawConfig) for _, layer := range manifest.Layers { p.addFetch(layer.Digest, layer.Size) } - return v1.Descriptor{MediaType: mediaType, Digest: digest, Size: int64(len(rawManifest))}, nil + return *desc, nil } func (p *pullPlan) addIndex(idx v1.ImageIndex) (v1.Descriptor, error) { - rawManifest, err := idx.RawManifest() + desc, err := partial.Descriptor(idx) if err != nil { return v1.Descriptor{}, err } - digest, err := idx.Digest() - if err != nil { - return v1.Descriptor{}, err - } - mediaType, err := idx.MediaType() + rawManifest, err := idx.RawManifest() if err != nil { return v1.Descriptor{}, err } - p.addRaw(digest, rawManifest) + p.addRaw(desc.Digest, rawManifest) manifest, err := idx.IndexManifest() if err != nil { return v1.Descriptor{}, err @@ -350,5 +339,5 @@ func (p *pullPlan) addIndex(idx v1.ImageIndex) (v1.Descriptor, error) { p.addFetch(child.Digest, child.Size) } } - return v1.Descriptor{MediaType: mediaType, Digest: digest, Size: int64(len(rawManifest))}, nil + return *desc, nil } From e58545f256a050b435ce6dadbe81c0ecda6ba7c3 Mon Sep 17 00:00:00 2001 From: Marcin Konowalczyk Date: Wed, 15 Jul 2026 15:24:36 +0200 Subject: [PATCH 7/8] feat: apply suggestions from code review Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Signed-off-by: Marcin Konowalczyk --- internal/source/copy.go | 5 ++++- internal/source/copy_e2e_test.go | 2 +- internal/source/fetch.go | 11 ++++++++--- 3 files changed, 13 insertions(+), 5 deletions(-) diff --git a/internal/source/copy.go b/internal/source/copy.go index f02ed5e..292f5b2 100644 --- a/internal/source/copy.go +++ b/internal/source/copy.go @@ -178,7 +178,10 @@ func writeRemoteLayout(ctx context.Context, dir, sourceRef string, platform Plat } } - rootDesc.Annotations = map[string]string{"org.opencontainers.image.ref.name": "olav"} + if rootDesc.Annotations == nil { + rootDesc.Annotations = map[string]string{} + } + rootDesc.Annotations["org.opencontainers.image.ref.name"] = "olav" if err := path.AppendDescriptor(rootDesc); err != nil { return err } diff --git a/internal/source/copy_e2e_test.go b/internal/source/copy_e2e_test.go index 9c9cad0..dbae685 100644 --- a/internal/source/copy_e2e_test.go +++ b/internal/source/copy_e2e_test.go @@ -160,7 +160,7 @@ func TestResolveRemoteResumesInterruptedPull(t *testing.T) { t.Fatal(err) } defer rc.Close() - full := make([]byte, manifest.Layers[0].Size) +full := make([]byte, int(manifest.Layers[0].Size)) if _, err := io.ReadFull(rc, full); err != nil { t.Fatal(err) } diff --git a/internal/source/fetch.go b/internal/source/fetch.go index bcf004c..0ee15c7 100644 --- a/internal/source/fetch.go +++ b/internal/source/fetch.go @@ -63,9 +63,14 @@ func (f *blobFetcher) fetchBlob(ctx context.Context, digest v1.Hash, size int64, if digest.Algorithm != "sha256" { return fmt.Errorf("unsupported digest algorithm %q", digest.Algorithm) } - if info, err := os.Stat(dst); err == nil && info.Size() == size { - f.counter.add(size) - return nil + if size > 0 { + if info, err := os.Stat(dst); err == nil && info.Size() == size { + // Best-effort cleanup: the blob is already present in the layout. + partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") + _ = os.Remove(partial) + f.counter.add(size) + return nil + } } partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") if info, err := os.Stat(partial); err == nil && info.Size() > 0 { From daba56fc8287d75b414b1a4ae0bcb0bc064836cb Mon Sep 17 00:00:00 2001 From: lczyk Date: Wed, 15 Jul 2026 15:28:39 +0200 Subject: [PATCH 8/8] chore: appease gofmt --- internal/source/copy_e2e_test.go | 2 +- internal/source/fetch.go | 3 +-- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/internal/source/copy_e2e_test.go b/internal/source/copy_e2e_test.go index dbae685..8b7385c 100644 --- a/internal/source/copy_e2e_test.go +++ b/internal/source/copy_e2e_test.go @@ -160,7 +160,7 @@ func TestResolveRemoteResumesInterruptedPull(t *testing.T) { t.Fatal(err) } defer rc.Close() -full := make([]byte, int(manifest.Layers[0].Size)) + full := make([]byte, int(manifest.Layers[0].Size)) if _, err := io.ReadFull(rc, full); err != nil { t.Fatal(err) } diff --git a/internal/source/fetch.go b/internal/source/fetch.go index 0ee15c7..ec81092 100644 --- a/internal/source/fetch.go +++ b/internal/source/fetch.go @@ -63,16 +63,15 @@ func (f *blobFetcher) fetchBlob(ctx context.Context, digest v1.Hash, size int64, if digest.Algorithm != "sha256" { return fmt.Errorf("unsupported digest algorithm %q", digest.Algorithm) } + partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") if size > 0 { if info, err := os.Stat(dst); err == nil && info.Size() == size { // Best-effort cleanup: the blob is already present in the layout. - partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") _ = os.Remove(partial) f.counter.add(size) return nil } } - partial := filepath.Join(f.partialDir, digest.Algorithm+"-"+digest.Hex+".partial") if info, err := os.Stat(partial); err == nil && info.Size() > 0 { if f.progress != nil { fmt.Fprintf(f.progress, "olav: resuming blob %s at %s\n", digest.Hex[:12], formatBytes(info.Size()))