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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 51 additions & 6 deletions pkg/xpkg/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ const cacheContentExt = ".gz"
// A PackageCache caches package content.
type PackageCache interface {
Has(id string) bool
// Get returns cached package content. The caller must close the returned reader.
Get(id string) (io.ReadCloser, error)
Store(id string, content io.ReadCloser) error
Delete(id string) error
Expand All @@ -49,6 +50,23 @@ type FsPackageCache struct {
mu sync.RWMutex
}

type unlockingReadCloser struct {
io.ReadCloser

once sync.Once
err error
unlock func()
}

func (r *unlockingReadCloser) Close() error {
r.once.Do(func() {
r.err = r.ReadCloser.Close()
r.unlock()
})

return r.err
}

// NewFsPackageCache creates a new FsPackageCache.
func NewFsPackageCache(dir string, fs afero.Fs) *FsPackageCache {
return &FsPackageCache{
Expand All @@ -59,53 +77,80 @@ func NewFsPackageCache(dir string, fs afero.Fs) *FsPackageCache {

// Has indicates whether an item with the given id is in the cache.
func (c *FsPackageCache) Has(id string) bool {
c.mu.RLock()
defer c.mu.RUnlock()

if fi, err := c.fs.Stat(BuildPath(c.dir, id, cacheContentExt)); err == nil && !fi.IsDir() {
return true
}

return false
}

// Get retrieves package contents from the cache.
// Get retrieves package contents from the cache. It holds a read lock until the
// returned reader is closed.
func (c *FsPackageCache) Get(id string) (io.ReadCloser, error) {
c.mu.RLock()
defer c.mu.RUnlock()

f, err := c.fs.Open(BuildPath(c.dir, id, cacheContentExt))
if err != nil {
c.mu.RUnlock()
return nil, err
}

return GzipReadCloser(f)
r, err := GzipReadCloser(f)
if err != nil {
_ = f.Close()
c.mu.RUnlock()
return nil, err
}

return &unlockingReadCloser{
ReadCloser: r,
unlock: c.mu.RUnlock,
}, nil
}
Comment on lines 92 to 112

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Description: Inspect the Get doc comment and locate all callers of PackageCache.Get.
set -euo pipefail

echo "=== cache.go Get doc comment ==="
sed -n '85,112p' pkg/xpkg/cache.go

echo "=== PackageCache interface definition ==="
ast-grep run --pattern 'type PackageCache interface { $$$ }' --lang go pkg/xpkg/cache.go

echo "=== callers of .cache.Get / cache.Get across repo ==="
rg -nP --type=go -C4 '\bcache\.Get\s*\(' -g '!**/vendor/**'

echo "=== FsPackageCache construction sites ==="
rg -nP --type=go -C3 '\bNewFsPackageCache\s*\('

Repository: crossplane/crossplane-runtime

Length of output: 1087


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "=== File locations ==="
git ls-files | rg '(^|/)(cache\.go|reader\.go|client\.go)$|_test\.go$' | sort | sed -n '1,120p'

echo "=== cache.go relevant outline ==="
ast-grep outline pkg/xpkg/cache.go --view expanded || true

echo "=== cache.go full relevant section ==="
sed -n '1,150p' pkg/xpkg/cache.go

echo "=== reader.go relevant GzipReadCloser ==="
sed -n '1,180p' pkg/xpkg/reader.go

echo "=== calls to Get methods in Go files ==="
rg -nP --type=go -C3 '\.Get\s*\(|\.get\s*\(' | head -200 || true

echo "=== FsPackageCache construction ==="
rg -nP --type=go -C3 '\bNewFsPackageCache\s*\(' | head -200 || true

echo "=== repo package metadata / breaking label references ==="
git ls-files | rg '(\.github|breaking|CHANGELOG|README|.*\.md)$' | xargs -r rg -n "breaking-change|breaking" | head -200 || true

Repository: crossplane/crossplane-runtime

Length of output: 9250


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "=== xpkg client cache usage ==="
sed -n '1,260p' pkg/xpkg/client.go

echo "=== cache tests around Get ==="
sed -n '1,260p' pkg/xpkg/cache_test.go

echo "=== all PackageCache references (read-only) ==="
rg -nP --type=go -C3 '\bPackageCache\b|\bNewFsPackageCache\s*\(' | head -200

echo "=== exported xpkg symbols in cache/client modules ==="
python3 - <<'PY'
from pathlib import Path
for p in ['pkg/xpkg/cache.go','pkg/xpkg/client.go']:
    print(f'--- {p} ---')
    data=Path(p).read_text()
    for i,line in enumerate(data.splitlines(),1):
        if line.startswith(('type ','Func','func ')) and i>=60 and i<=130:
            print(f'{i}: {line}')
PY

Repository: crossplane/crossplane-runtime

Length of output: 14137


Document the cache read-lock lifecycle on PackageCache.Get / FsPackageCache.Get.

Get now returns an io.ReadCloser that releases the cache RUnlock only when closed. Add a doc requirement such as: the caller must close the returned reader to release the cache read lock; leaking it can block later Store and Delete calls. Also note whether this counts as a breaking-change because the lock lifetime changed from the previous exported API behavior.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/xpkg/cache.go` around lines 90 - 110, Document the read-lock lifecycle on
the exported PackageCache.Get and FsPackageCache.Get APIs: callers must close
the returned io.ReadCloser to release the cache read lock, and failing to do so
can block subsequent Store and Delete calls. Explicitly mark this as a breaking
change if the project’s API documentation supports such annotations, since the
lock now remains held until the reader is closed.

Source: Coding guidelines


// Store saves the package contents to the cache.
func (c *FsPackageCache) Store(id string, content io.ReadCloser) error {
c.mu.Lock()
defer c.mu.Unlock()

cf, err := c.fs.Create(BuildPath(c.dir, id, cacheContentExt))
path := BuildPath(c.dir, id, cacheContentExt)
cf, err := c.fs.Create(path)
if err != nil {
return err
}
defer cf.Close() //nolint:errcheck // Error is checked in the happy path.
cleanup := func() {
_ = cf.Close()
_ = c.fs.Remove(path)
}

w, err := gzip.NewWriterLevel(cf, gzip.BestSpeed)
if err != nil {
cleanup()
return err
}

_, err = io.Copy(w, content)
if err != nil {
_ = w.Close()
cleanup()
return err
}
// NOTE(hasheddan): gzip writer must be closed to ensure all data is flushed
// to file.
if err := w.Close(); err != nil {
cleanup()
return err
}

return cf.Close()
if err := cf.Close(); err != nil {
cleanup()
return err
}

return nil
}

// Delete removes package contents from the cache.
Expand Down
177 changes: 177 additions & 0 deletions pkg/xpkg/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,14 @@ package xpkg
import (
"bytes"
"compress/gzip"
"errors"
"io"
"os"
"strings"
"syscall"
"testing"
"testing/iotest"
"time"

"github.com/google/go-cmp/cmp"
"github.com/spf13/afero"
Expand All @@ -32,6 +36,44 @@ import (

var _ PackageCache = &FsPackageCache{}

type errorFile struct {
afero.File
writeErr error
closeErr error
}

func (f *errorFile) Write(p []byte) (int, error) {
if f.writeErr != nil {
return 0, f.writeErr
}

return f.File.Write(p)
}

func (f *errorFile) Close() error {
err := f.File.Close()
if f.closeErr != nil {
return f.closeErr
}

return err
}

type errorFs struct {
afero.Fs
writeErr error
closeErr error
}

func (f *errorFs) Create(name string) (afero.File, error) {
file, err := f.Fs.Create(name)
if err != nil {
return nil, err
}

return &errorFile{File: file, writeErr: f.writeErr, closeErr: f.closeErr}, nil
}

func TestHas(t *testing.T) {
fs := afero.NewMemMapFs()
cf, _ := fs.Create("/cache/exists.gz")
Expand Down Expand Up @@ -181,6 +223,141 @@ func TestStore(t *testing.T) {
}
}

func TestStoreRoundTrip(t *testing.T) {
cache := NewFsPackageCache("/cache", afero.NewMemMapFs())
want := "package content"

if err := cache.Store("package", io.NopCloser(strings.NewReader(want))); err != nil {
t.Fatalf("Store(...): unexpected error: %v", err)
}

r, err := cache.Get("package")
if err != nil {
t.Fatalf("Get(...): unexpected error: %v", err)
}
defer r.Close()

got, err := io.ReadAll(r)
if err != nil {
t.Fatalf("Read(...): unexpected error: %v", err)
}
if diff := cmp.Diff(want, string(got)); diff != "" {
t.Errorf("Store(...): -want content, +got content:\n%s", diff)
}
}

func TestStoreRemovesFailedWrites(t *testing.T) {
errWrite := errors.New("write failed")
cases := map[string]struct {
fs afero.Fs
content io.ReadCloser
}{
"ContentReadError": {
fs: afero.NewMemMapFs(),
content: io.NopCloser(io.MultiReader(strings.NewReader("partial"), iotest.ErrReader(errWrite))),
},
"GzipCloseError": {
fs: &errorFs{Fs: afero.NewMemMapFs(), writeErr: errWrite},
content: io.NopCloser(strings.NewReader("")),
},
"FileCloseError": {
fs: &errorFs{Fs: afero.NewMemMapFs(), closeErr: errWrite},
content: io.NopCloser(strings.NewReader("content")),
},
}

for name, tc := range cases {
t.Run(name, func(t *testing.T) {
cache := NewFsPackageCache("/cache", tc.fs)
if err := cache.Store("package", tc.content); err == nil {
t.Fatal("Store(...): expected an error")
}
if cache.Has("package") {
t.Fatal("Store(...): failed write remained in cache")
}
})
}
}

func TestStoreWaitsForReader(t *testing.T) {
cache := NewFsPackageCache("/cache", afero.NewMemMapFs())
if err := cache.Store("package", io.NopCloser(strings.NewReader("old"))); err != nil {
t.Fatalf("Store(...): unexpected error: %v", err)
}

r, err := cache.Get("package")
if err != nil {
t.Fatalf("Get(...): unexpected error: %v", err)
}

started := make(chan struct{})
done := make(chan error, 1)
go func() {
close(started)
done <- cache.Store("package", io.NopCloser(strings.NewReader("new")))
}()
<-started

select {
case err := <-done:
t.Fatalf("Store(...) completed while cache reader was open: %v", err)
case <-time.After(100 * time.Millisecond):
}

got, err := io.ReadAll(r)
if err != nil {
t.Fatalf("Read(...): unexpected error: %v", err)
}
if diff := cmp.Diff("old", string(got)); diff != "" {
t.Errorf("Get(...): -want content, +got content:\n%s", diff)
}
if err := r.Close(); err != nil {
t.Fatalf("Close(...): unexpected error: %v", err)
}

select {
case err := <-done:
if err != nil {
t.Fatalf("Store(...): unexpected error: %v", err)
}
case <-time.After(time.Second):
t.Fatal("Store(...) did not complete after cache reader closed")
}
}

func TestGetErrorReleasesLock(t *testing.T) {
fs := afero.NewMemMapFs()
f, err := fs.Create(BuildPath("/cache", "package", cacheContentExt))
if err != nil {
t.Fatalf("Create(...): unexpected error: %v", err)
}
if _, err := f.WriteString("not gzip"); err != nil {
t.Fatalf("WriteString(...): unexpected error: %v", err)
}
if err := f.Close(); err != nil {
t.Fatalf("Close(...): unexpected error: %v", err)
}

cache := NewFsPackageCache("/cache", fs)
if _, err := cache.Get("package"); err == nil {
t.Fatal("Get(...): expected an error")
}

done := make(chan error, 1)
go func() {
done <- cache.Delete("package")
}()

select {
case err := <-done:
if err != nil {
t.Fatalf("Delete(...): unexpected error: %v", err)
}
case <-time.After(time.Second):
t.Fatal("Delete(...) blocked after Get(...) failed")
}
}

func TestDelete(t *testing.T) {
fs := afero.NewMemMapFs()
_, _ = fs.Create("/cache/exists.xpkg")
Expand Down
Loading