diff --git a/go.mod b/go.mod index 18f2cfcbab..b8f5bcaec9 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module github.com/buildkite/agent/v4 go 1.26.5 require ( + buf.build/gen/go/namespace/cloud/protocolbuffers/go v1.36.12-20260820164744-a1973bcf4d87.1 cloud.google.com/go/compute/metadata v0.9.0 cloud.google.com/go/kms v1.33.0 connectrpc.com/connect v1.20.0 @@ -70,11 +71,14 @@ require ( golang.org/x/sys v0.47.0 golang.org/x/term v0.45.0 google.golang.org/api v0.292.0 + google.golang.org/grpc v1.83.0 google.golang.org/protobuf v1.36.12 gopkg.in/yaml.v3 v3.0.1 + namespacelabs.dev/integrations v0.0.11-0.20260508113815-6a8135624a35 // Pinned for CreateNewVersionIfExists support. ) require ( + buf.build/gen/go/namespace/cloud/grpc/go v1.6.2-20260820164744-a1973bcf4d87.1 // indirect cloud.google.com/go v0.123.0 // indirect cloud.google.com/go/auth v0.22.0 // indirect cloud.google.com/go/auth/oauth2adapt v0.2.8 // indirect @@ -114,6 +118,7 @@ require ( github.com/go-logr/logr v1.4.4 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/goccy/go-json v0.10.6 // indirect + github.com/golang-jwt/jwt/v4 v4.5.2 // indirect github.com/golang-jwt/jwt/v5 v5.3.1 // indirect github.com/google/s2a-go v0.1.9 // indirect github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510 // indirect @@ -157,7 +162,6 @@ require ( google.golang.org/genproto v0.0.0-20260406210006-6f92a3bedf2d // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260803160001-6ac0973c030d // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260803160001-6ac0973c030d // indirect - google.golang.org/grpc v1.83.0 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect gotest.tools/gotestsum v1.13.0 // indirect mvdan.cc/gofumpt v0.9.2 // indirect diff --git a/go.sum b/go.sum index c9186d8d25..eaa99a28da 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,7 @@ +buf.build/gen/go/namespace/cloud/grpc/go v1.6.2-20260820164744-a1973bcf4d87.1 h1:Szd8+tmae26rv8IG14gItvykV9vhFzvrBcUg47ZAVHM= +buf.build/gen/go/namespace/cloud/grpc/go v1.6.2-20260820164744-a1973bcf4d87.1/go.mod h1:v2Fn3cyk8eT581ReENee/l/lFl1JNYMU/WspzgSjXiU= +buf.build/gen/go/namespace/cloud/protocolbuffers/go v1.36.12-20260820164744-a1973bcf4d87.1 h1:rha4+KA9AGYPQktS1LYeHjpXDafnpoPLziO0FGhd1iw= +buf.build/gen/go/namespace/cloud/protocolbuffers/go v1.36.12-20260820164744-a1973bcf4d87.1/go.mod h1:qZ8DAMLNpt2zvGyC5UsmblfkkKyUWbVHEkLodBFoMKg= cloud.google.com/go v0.123.0 h1:2NAUJwPR47q+E35uaJeYoNhuNEM9kM8SjgRgdeOJUSE= cloud.google.com/go v0.123.0/go.mod h1:xBoMV08QcqUGuPW65Qfm1o9Y4zKZBpGS+7bImXLTAZU= cloud.google.com/go/auth v0.22.0 h1:Xp9wAKkLoeaYb5pYZZoQGz4E9sdPxIbzS3gywZE3ciQ= @@ -168,6 +172,8 @@ github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU= github.com/goccy/go-json v0.10.6/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= github.com/gofrs/flock v0.13.0 h1:95JolYOvGMqeH31+FC7D2+uULf6mG61mEZ/A8dRYMzw= github.com/gofrs/flock v0.13.0/go.mod h1:jxeyy9R1auM5S6JYDBhDt+E2TCo7DkratH4Pgi8P+Z0= +github.com/golang-jwt/jwt/v4 v4.5.2 h1:YtQM7lnr8iZ+j5q71MGKkNw9Mn7AjHM68uc9g5fXeUI= +github.com/golang-jwt/jwt/v4 v4.5.2/go.mod h1:m21LjoU+eqJr34lmDMbreY2eSTRJ1cv77w39/MY0Ch0= github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY= github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= @@ -188,8 +194,8 @@ github.com/googleapis/enterprise-certificate-proxy v0.3.19 h1:mMOE7DN2+p76/EdIrm github.com/googleapis/enterprise-certificate-proxy v0.3.19/go.mod h1:rSEsBUemEBZEexP2y6jPp16LUmUbjmSbcPMQizR0o4k= github.com/googleapis/gax-go/v2 v2.23.0 h1:Tchl7qkvE7Ip3y+ztvNufYFvkfqTe7NfLTYGIdJRLuE= github.com/googleapis/gax-go/v2 v2.23.0/go.mod h1:rBQKOVJCdb8IFEzg+FCwlt1LP/xMDGuqUXhUG+XMXEg= -github.com/gorilla/websocket v1.5.0 h1:PPwGk2jz7EePpoHN/+ClbZu8SPxiqlu12wZP/3sWmnc= -github.com/gorilla/websocket v1.5.0/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +github.com/gorilla/websocket v1.5.1 h1:gmztn0JnHVt9JZquRuzLw3g4wouNVzKL15iLr/zn/QY= +github.com/gorilla/websocket v1.5.1/go.mod h1:x3kM2JMyaluk02fnUJpQuwD2dCS5NDG2ZHL0uE0tcaY= github.com/gowebpki/jcs v1.0.1 h1:Qjzg8EOkrOTuWP7DqQ1FbYtcpEbeTzUoTN9bptp8FOU= github.com/gowebpki/jcs v1.0.1/go.mod h1:CID1cNZ+sHp1CCpAR8mPf6QRtagFBgPJE0FCUQ6+BrI= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk= @@ -389,3 +395,5 @@ gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= gotest.tools/v3 v3.5.2/go.mod h1:LtdLGcnqToBH83WByAAi/wiwSFCArdFIUV/xxN4pcjA= mvdan.cc/gofumpt v0.9.2 h1:zsEMWL8SVKGHNztrx6uZrXdp7AX8r421Vvp23sz7ik4= mvdan.cc/gofumpt v0.9.2/go.mod h1:iB7Hn+ai8lPvofHd9ZFGVg2GOr8sBUw1QUWjNbmIL/s= +namespacelabs.dev/integrations v0.0.11-0.20260508113815-6a8135624a35 h1:khQ002qSHfL0QocNYvjbEJHCHZCNrlNbfZJA/hPCqPA= +namespacelabs.dev/integrations v0.0.11-0.20260508113815-6a8135624a35/go.mod h1:E8wND7kuqalgSfiaGuxkfYd6xRq918FGX8vpnUmwOu8= diff --git a/internal/cache/cache.go b/internal/cache/cache.go index 4e1b95e89c..3ecb429890 100644 --- a/internal/cache/cache.go +++ b/internal/cache/cache.go @@ -34,6 +34,19 @@ type cacheOps interface { ListCaches() []configuration.Cache } +type cacheClientCloser interface { + close() error +} + +func withClientCleanup(l logger.Logger, c cacheClientCloser, run func() error) error { + defer func() { + if closeErr := c.close(); closeErr != nil { + l.Warnf("Failed to close Namespace storage client: %v", closeErr) + } + }() + return run() +} + // RunSave saves caches based on the provided configuration and logs results as // each cache is processed. func RunSave(ctx context.Context, l logger.Logger, apiClient *api.Client, cfg Config) error { @@ -45,7 +58,9 @@ func RunSave(ctx context.Context, l logger.Logger, apiClient *api.Client, cfg Co l.Infof("No caches defined in the cache configuration file, nothing to save") return nil } - return saveWithClient(ctx, l, c, cacheIDs, cfg.Concurrency) + return withClientCleanup(l, c, func() error { + return saveWithClient(ctx, l, c, cacheIDs, cfg.Concurrency) + }) } // RunRestore restores caches based on the provided configuration and logs results @@ -59,7 +74,9 @@ func RunRestore(ctx context.Context, l logger.Logger, apiClient *api.Client, cfg l.Infof("No caches defined in the cache configuration file, nothing to restore") return nil } - return restoreWithClient(ctx, l, c, cacheIDs, cfg.Concurrency) + return withClientCleanup(l, c, func() error { + return restoreWithClient(ctx, l, c, cacheIDs, cfg.Concurrency) + }) } // ListCaches returns all cache definitions configured on the client. diff --git a/internal/cache/cache_test.go b/internal/cache/cache_test.go index ba92806cd6..60573d1b9c 100644 --- a/internal/cache/cache_test.go +++ b/internal/cache/cache_test.go @@ -43,6 +43,63 @@ func (m *mockCacheClient) ListCaches() []configuration.Cache { return nil } +type fakeCacheClientCloser struct { + closeFunc func() error +} + +func (c fakeCacheClientCloser) close() error { + return c.closeFunc() +} + +func TestWithClientCleanup(t *testing.T) { + t.Run("closes after work completes", func(t *testing.T) { + workComplete := false + closed := false + closer := fakeCacheClientCloser{closeFunc: func() error { + if !workComplete { + t.Error("client closed before work completed") + } + closed = true + return nil + }} + + err := withClientCleanup(logger.Discard, closer, func() error { + workComplete = true + return nil + }) + if err != nil { + t.Fatalf("withClientCleanup error = %v, want nil", err) + } + if !closed { + t.Error("client was not closed") + } + }) + + t.Run("close failure is logged without masking work error", func(t *testing.T) { + workErr := errors.New("worker failed") + closeErr := errors.New("close failed") + log := logger.NewBuffer() + closer := fakeCacheClientCloser{closeFunc: func() error { return closeErr }} + + err := withClientCleanup(log, closer, func() error { return workErr }) + if !errors.Is(err, workErr) { + t.Fatalf("withClientCleanup error = %v, want %v", err, workErr) + } + if got := strings.Join(log.Messages, "\n"); !strings.Contains(got, closeErr.Error()) { + t.Errorf("log messages = %q, want close error", got) + } + }) + + t.Run("close failure does not turn success into failure", func(t *testing.T) { + closeErr := errors.New("close failed") + closer := fakeCacheClientCloser{closeFunc: func() error { return closeErr }} + + if err := withClientCleanup(logger.Discard, closer, func() error { return nil }); err != nil { + t.Fatalf("withClientCleanup error = %v, want nil", err) + } + }) +} + // Test helpers func createTempCacheConfig(t *testing.T, content string) string { @@ -396,6 +453,12 @@ func TestNewClient_ValidCacheIDs(t *testing.T) { if got := client; got == nil { t.Fatalf("newClient(logger.Discard, nil, cfg) = %v, want non-nil value", got) } + if client.nscClient == nil { + t.Fatal("newClient did not create a Namespace client holder") + } + if err := client.close(); err != nil { + t.Fatalf("close unused client: %v", err) + } if diff := cmp.Diff(cacheIDs, []string{"cache1", "cache2"}); diff != "" { t.Fatalf("newClient(logger.Discard, nil, cfg) diff (-got +want):\n%s", diff) } diff --git a/internal/cache/client.go b/internal/cache/client.go index eede5562c4..9b7647dab9 100644 --- a/internal/cache/client.go +++ b/internal/cache/client.go @@ -15,6 +15,7 @@ import ( "github.com/buildkite/agent/v4/api" "github.com/buildkite/agent/v4/internal/cache/configuration" + "github.com/buildkite/agent/v4/internal/cache/store" "github.com/buildkite/agent/v4/logger" ) @@ -47,9 +48,9 @@ var ( // client is a configured handle for cache Save and Restore operations. // -// It is not a network connection; it just bundles the API client, storage -// bucket and the expanded, validated cache definitions used by every call. -// Safe for concurrent use; honours context cancellation. +// It bundles the API client, storage bucket, expanded cache definitions, and a +// lazily initialized Namespace connection. Safe for concurrent use; honours +// context cancellation. type client struct { api cacheAPI bucketURL string @@ -57,6 +58,7 @@ type client struct { platform string registry string caches []configuration.Cache + nscClient *store.NscClient onProgress ProgressCallback } @@ -100,6 +102,7 @@ func newClient(l logger.Logger, apiClient cacheAPI, cfg Config) (*client, []stri platform: fmt.Sprintf("%s/%s", runtime.GOOS, runtime.GOARCH), registry: registry, caches: expanded, + nscClient: store.NewNscClient(), onProgress: func(cacheID, stage, message string, _, _ int) { l.WithFields( logger.StringField("cache_id", cacheID), @@ -116,6 +119,10 @@ func newClient(l logger.Logger, apiClient cacheAPI, cfg Config) (*client, []stri return c, names, nil } +func (c *client) close() error { + return c.nscClient.Close() +} + // resolveCacheNames returns requested if non-empty (after validating every name // exists), otherwise returns every cache name configured on the client. func (c *client) resolveCacheNames(requested []string) ([]string, error) { diff --git a/internal/cache/restore.go b/internal/cache/restore.go index 7919d55c51..b3b95c8987 100644 --- a/internal/cache/restore.go +++ b/internal/cache/restore.go @@ -407,7 +407,7 @@ func (c *client) downloadCache(ctx context.Context, retrieveResp api.CacheEntryR ) // Create blob store - blobStore, err := store.NewBlobStore(ctx, retrieveResp.Store, bucketURL) + blobStore, err := store.NewBlobStore(ctx, retrieveResp.Store, bucketURL, c.nscClient) if err != nil { span.RecordError(err) span.SetStatus(codes.Error, "failed to create blob store") @@ -422,7 +422,9 @@ func (c *client) downloadCache(ctx context.Context, retrieveResp api.CacheEntryR return "", "", nil, fmt.Errorf("failed to create temp directory: %w", err) } - archiveFile = filepath.Join(tmpDir, storeObjectName) + // The object name comes from the cache API. Keep it as the remote lookup key, + // but never use it as a local path component. + archiveFile = filepath.Join(tmpDir, "archive") // Download archive transferInfo, err = blobStore.Download(ctx, storeObjectName, archiveFile) diff --git a/internal/cache/save.go b/internal/cache/save.go index 8409989e90..006152202e 100644 --- a/internal/cache/save.go +++ b/internal/cache/save.go @@ -277,7 +277,7 @@ func (c *client) Save(ctx context.Context, cacheID string) (SaveResult, error) { c.callProgress(cacheID, "uploading", "Uploading cache archive", 0, int(archiveInfo.Size)) // Upload archive - blobStore, err := store.NewBlobStore(ctx, registryResp.Store, c.bucketURL) + blobStore, err := store.NewBlobStore(ctx, registryResp.Store, c.bucketURL, c.nscClient) if err != nil { span.RecordError(err) span.SetStatus(codes.Error, "failed to create blob store") diff --git a/internal/cache/store/blob.go b/internal/cache/store/blob.go index 46c3f77d37..d898ebe62d 100644 --- a/internal/cache/store/blob.go +++ b/internal/cache/store/blob.go @@ -20,13 +20,13 @@ type Blob interface { Download(ctx context.Context, key, destPath string) (*TransferInfo, error) } -func NewBlobStore(ctx context.Context, store, bucketURL string) (Blob, error) { +func NewBlobStore(ctx context.Context, store, bucketURL string, nscClient *NscClient) (Blob, error) { switch store { case AgentManaged: scheme, _, _ := strings.Cut(bucketURL, "://") switch scheme { case nscScheme: - return NewNscStore(bucketURL) + return NewNscStore(bucketURL, nscClient) case "file": // Supported only for local testing, kept consistent with validateCacheStore. return NewLocalFileBlob(ctx, bucketURL) diff --git a/internal/cache/store/file.go b/internal/cache/store/file.go index 6bcc085ec2..0e3c8e4d1f 100644 --- a/internal/cache/store/file.go +++ b/internal/cache/store/file.go @@ -53,6 +53,27 @@ type FileMetadata struct { Version int `json:"version"` // Metadata schema version } +func validateFilePath(filePath string) error { + if filePath == "" { + return fmt.Errorf("file path cannot be empty") + } + + cleanPath := filepath.Clean(filePath) + dangerousChars := []string{";", "&", "|", "`", "$", "(", ")", "{", "}", "[", "]", "<", ">", "\"", "'"} + if runtime.GOOS != "windows" { + dangerousChars = append(dangerousChars, "\\") + } + for _, char := range dangerousChars { + if strings.Contains(cleanPath, char) { + return fmt.Errorf("file path contains potentially dangerous character: %s", char) + } + } + if strings.Contains(cleanPath, "..") { + return fmt.Errorf("file path contains path traversal sequence") + } + return nil +} + // NewLocalFileBlob creates a new local file storage backend from a file:// URL. // // Supported URL formats: diff --git a/internal/cache/store/file_test.go b/internal/cache/store/file_test.go index fc79588c67..04a0662114 100644 --- a/internal/cache/store/file_test.go +++ b/internal/cache/store/file_test.go @@ -10,6 +10,41 @@ import ( "testing" ) +func TestValidateFilePath(t *testing.T) { + tests := []struct { + name string + filePath string + wantError string + }{ + {name: "valid simple path", filePath: "test.txt"}, + {name: "valid relative path", filePath: "dir/subdir/file.txt"}, + {name: "valid absolute path", filePath: "/tmp/test.txt"}, + {name: "empty path", wantError: "file path cannot be empty"}, + {name: "path with semicolon", filePath: "file;rm -rf /", wantError: "file path contains potentially dangerous character: ;"}, + {name: "path with ampersand", filePath: "file&malicious", wantError: "file path contains potentially dangerous character: &"}, + {name: "path with pipe", filePath: "file|cat /etc/passwd", wantError: "file path contains potentially dangerous character: |"}, + {name: "path with backtick", filePath: "file`whoami`", wantError: "file path contains potentially dangerous character: `"}, + {name: "path with dollar sign", filePath: "file$(whoami)", wantError: "file path contains potentially dangerous character: $"}, + {name: "path traversal attempt", filePath: "../../../etc/passwd", wantError: "file path contains path traversal sequence"}, + {name: "path with quotes", filePath: `file"test"`, wantError: `file path contains potentially dangerous character: "`}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + err := validateFilePath(tt.filePath) + if tt.wantError == "" { + if err != nil { + t.Fatalf("validateFilePath: %v", err) + } + return + } + if err == nil || !strings.Contains(err.Error(), tt.wantError) { + t.Fatalf("validateFilePath error = %v, want error containing %q", err, tt.wantError) + } + }) + } +} + func TestNewLocalFileBlob(t *testing.T) { ctx := t.Context() @@ -560,7 +595,7 @@ func TestNewBlobStoreLocalFile(t *testing.T) { ctx := t.Context() tmpDir := t.TempDir() - blob, err := NewBlobStore(ctx, LocalFileStore, fileURL(tmpDir)) + blob, err := NewBlobStore(ctx, LocalFileStore, fileURL(tmpDir), nil) if err != nil { t.Fatalf("NewBlobStore: %v", err) } @@ -576,7 +611,7 @@ func TestNewBlobStoreLocalFile(t *testing.T) { // file:// is accepted for agent_managed (local testing) and must reach the local // file store rather than the S3 store, matching validateCacheStore. func TestNewBlobStoreAgentManagedFile(t *testing.T) { - blob, err := NewBlobStore(t.Context(), AgentManaged, fileURL(t.TempDir())) + blob, err := NewBlobStore(t.Context(), AgentManaged, fileURL(t.TempDir()), nil) if err != nil { t.Fatalf("NewBlobStore: %v", err) } diff --git a/internal/cache/store/nsc.go b/internal/cache/store/nsc.go index b9f0fbf7b6..f9a08c7b4a 100644 --- a/internal/cache/store/nsc.go +++ b/internal/cache/store/nsc.go @@ -1,166 +1,189 @@ package store import ( - "bytes" "context" + "errors" "fmt" + "io" "log/slog" "net/url" "os" - "os/exec" - "path/filepath" - "regexp" - "runtime" + "strconv" "strings" + "sync" "time" + storagev1beta "buf.build/gen/go/namespace/cloud/protocolbuffers/go/proto/namespace/cloud/storage/v1beta" + "github.com/buildkite/agent/v4/api" "github.com/buildkite/agent/v4/internal/cache/internal/trace" + "github.com/buildkite/roko" "go.opentelemetry.io/otel/attribute" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/types/known/durationpb" + "namespacelabs.dev/integrations/api/storage" + "namespacelabs.dev/integrations/auth" ) // nscScheme is the URL scheme that routes an agent-managed cache store to NSC. const nscScheme = "nsc" -// nscDefaultExpiry is the artifact lifetime used both when uploading a cache -// entry (--expires_in) and when refreshing it on access (--ensure_minimum). -// Cache entries are content-addressed and short-lived, so we cap storage growth -// rather than relying on NSC's no-expiry default. Every restore pushes the -// expiry back out to this duration from now, keeping hot caches alive while -// letting cold ones expire. -const nscDefaultExpiry = "24h" - -// commandRunner executes an external command. It is a seam so tests can assert -// the arguments passed to the nsc CLI without invoking the real binary. -type commandRunner func(ctx context.Context, workingDir string, args ...string) (*CommandResult, error) - -// NscStore implements the Blob interface for NSC artifact storage which uses the nsc CLI tool -// https://namespace.so/docs/reference/cli/artifact-download -// https://namespace.so/docs/reference/cli/artifact-upload -type NscStore struct { - // namespace is taken from the nsc:// cache store URL and passed - // as --namespace to the nsc CLI. It is always non-empty. - namespace string - run commandRunner +const ( + // nscDefaultExpiry is the artifact lifetime used both when uploading a cache + // entry and when refreshing it on access. + nscDefaultExpiry = 24 * time.Hour + + // Whole-file retries are bounded because, unlike the NSC CLI's ranged + // downloader, every attempt starts the download again from the beginning. + nscDownloadAttempts = 3 +) + +type nscAPIClient interface { + UploadArtifact(context.Context, string, string, io.Reader, storage.UploadOpts) error + ResolveArtifactStream(context.Context, string, string) (io.ReadCloser, error) + ExtendArtifact(context.Context, *storagev1beta.ExtendArtifactRequest) error + Close() error } -// NewNscStore creates a store backed by the nsc CLI. The namespace is parsed -// from the nsc:// cache store URL and is required. -func NewNscStore(bucketURL string) (*NscStore, error) { - namespace, err := parseNscNamespace(bucketURL) - if err != nil { - return nil, err - } - return &NscStore{namespace: namespace, run: runCommand}, nil +// Namespace Storage API documentation: +// https://buf.build/namespace/cloud/docs/main:namespace.cloud.storage.v1beta +type nscStorageClient struct { + client storage.Client } -// parseNscNamespace extracts the namespace from an nsc:// cache store -// URL. It errors on a non-nsc URL or a missing namespace. -func parseNscNamespace(bucketURL string) (string, error) { - u, err := url.Parse(bucketURL) - if err != nil { - return "", fmt.Errorf("invalid cache store URL %q: %w", bucketURL, err) - } - if u.Scheme != nscScheme { - return "", fmt.Errorf("expected %s:// cache store URL, got %q", nscScheme, bucketURL) - } - if u.Host == "" { - return "", fmt.Errorf("nsc:// URL must include a namespace, e.g. nsc://my-namespace") - } - return u.Host, nil +func (c *nscStorageClient) UploadArtifact(ctx context.Context, nsc, path string, r io.Reader, opts storage.UploadOpts) error { + _, err := storage.UploadArtifactWithOpts(ctx, c.client, nsc, path, r, opts) + return err } -// artifactArgs builds the nsc CLI argument list, including --namespace. -func (n *NscStore) artifactArgs(args ...string) []string { - cmd := append([]string{"nsc", "artifact"}, args...) - return append(cmd, "--namespace", n.namespace) +func (c *nscStorageClient) ResolveArtifactStream(ctx context.Context, nsc, path string) (io.ReadCloser, error) { + return storage.ResolveArtifactStream(ctx, c.client, nsc, path) } -// validateFilePath validates that a file path is safe for use in commands -func validateFilePath(filePath string) error { - if filePath == "" { - return fmt.Errorf("file path cannot be empty") - } +func (c *nscStorageClient) ExtendArtifact(ctx context.Context, req *storagev1beta.ExtendArtifactRequest) error { + _, err := c.client.Artifacts.ExtendArtifact(ctx, req) + return err +} - // Clean the path to normalize it - cleanPath := filepath.Clean(filePath) +func (c *nscStorageClient) Close() error { + return c.client.Close() +} - // Check for potentially dangerous characters that could be used for command injection. - // Backslash is the path separator on Windows so it must be allowed there. - dangerousChars := []string{";", "&", "|", "`", "$", "(", ")", "{", "}", "[", "]", "<", ">", "\"", "'"} - if runtime.GOOS != "windows" { - dangerousChars = append(dangerousChars, "\\") - } - for _, char := range dangerousChars { - if strings.Contains(cleanPath, char) { - return fmt.Errorf("file path contains potentially dangerous character: %s", char) - } - } +type nscClientFactory func(context.Context) (nscAPIClient, error) + +// NscClient lazily initializes one Namespace Storage API client for use by all +// NSC cache transfers in a cache command. +type NscClient struct { + initialize nscClientFactory + initOnce sync.Once + client nscAPIClient + initErr error + closeOnce sync.Once + closeErr error +} - // Check for path traversal attempts - if strings.Contains(cleanPath, "..") { - return fmt.Errorf("file path contains path traversal sequence") - } +// NewNscClient creates a lazy Namespace Storage API client holder. +func NewNscClient() *NscClient { + return newNscClient(newNscAPIClient) +} - return nil +func newNscClient(initialize nscClientFactory) *NscClient { + return &NscClient{initialize: initialize} } -// validateKey validates that an artifact key is safe for use in commands -func validateKey(key string) error { - if key == "" { - return fmt.Errorf("key cannot be empty") +func newNscAPIClient(ctx context.Context) (nscAPIClient, error) { + token, err := auth.LoadDefaults() + if err != nil { + return nil, fmt.Errorf("load Namespace authentication: %w", err) } - // Check length - reasonable limit for artifact keys - if len(key) > 256 { - return fmt.Errorf("key too long (max 256 characters)") + client, err := storage.NewClient(ctx, token) + if err != nil { + return nil, fmt.Errorf("create Namespace Storage API client: %w", err) } + return &nscStorageClient{client: client}, nil +} - // NSC artifact keys should be alphanumeric with some safe special characters - // Allow: alphanumeric, hyphens, underscores, dots, forward slashes - validKeyPattern := regexp.MustCompile(`^[a-zA-Z0-9._/-]+$`) - if !validKeyPattern.MatchString(key) { - return fmt.Errorf("key contains invalid characters (only alphanumeric, ., _, /, - are allowed)") +func (c *NscClient) get(ctx context.Context) (nscAPIClient, error) { + c.initOnce.Do(func() { + c.client, c.initErr = c.initialize(ctx) + }) + return c.client, c.initErr +} + +// Close closes the initialized Namespace Storage API client. It is a no-op if +// no NSC cache transfer initialized the client. +func (c *NscClient) Close() error { + if c.client == nil { + return nil } + c.closeOnce.Do(func() { + c.closeErr = c.client.Close() + }) + return c.closeErr +} - // Check for potentially dangerous patterns - dangerousPatterns := []string{"../", "./", "//", "&&", "||", ";", "`", "$"} - for _, pattern := range dangerousPatterns { - if strings.Contains(key, pattern) { - return fmt.Errorf("key contains potentially dangerous pattern: %s", pattern) - } +// NscStore implements Blob using the Namespace Storage API. +type NscStore struct { + nsc string + client *NscClient +} + +// NewNscStore creates a Namespace Storage API-backed store. +func NewNscStore(bucketURL string, client *NscClient) (*NscStore, error) { + if client == nil { + return nil, fmt.Errorf("namespace client is required for nsc:// cache stores") + } + nsc, err := parseNscNamespace(bucketURL) + if err != nil { + return nil, err } + return &NscStore{nsc: nsc, client: client}, nil +} - return nil +// parseNscNamespace extracts the namespace from an nsc:// cache store +// URL. It errors on a non-nsc URL or a missing namespace. +func parseNscNamespace(bucketURL string) (string, error) { + u, err := url.Parse(bucketURL) + if err != nil { + return "", fmt.Errorf("invalid cache store URL %q: %w", bucketURL, err) + } + if u.Scheme != nscScheme { + return "", fmt.Errorf("expected %s:// cache store URL, got %q", nscScheme, bucketURL) + } + if u.Host == "" { + return "", fmt.Errorf("nsc:// URL must include a namespace, e.g. nsc://my-namespace") + } + return u.Host, nil } func (n *NscStore) Upload(ctx context.Context, filePath, key string) (*TransferInfo, error) { _, span := trace.Start(ctx, "NscStore.Upload") defer span.End() - // Validate input parameters to prevent command injection - if err := validateFilePath(filePath); err != nil { - return nil, fmt.Errorf("invalid file path: %w", err) - } - if err := validateKey(key); err != nil { - return nil, fmt.Errorf("invalid key: %w", err) + file, err := os.Open(filePath) + if err != nil { + return nil, fmt.Errorf("open upload file: %w", err) } + defer func() { _ = file.Close() }() - start := time.Now() - - // Execute nsc artifact upload command - result, err := n.run(ctx, "", n.artifactArgs("upload", filePath, key, "--expires_in", nscDefaultExpiry)...) + fileInfo, err := file.Stat() if err != nil { - return nil, fmt.Errorf("failed to execute nsc upload command: %w", err) + return nil, fmt.Errorf("get upload file info: %w", err) } - if result.ExitCode != 0 { - return nil, fmt.Errorf("nsc upload failed with exit code %d: %s", result.ExitCode, result.Stderr) + client, err := n.client.get(ctx) + if err != nil { + return nil, fmt.Errorf("initialize Namespace client for upload: %w", err) } - // Get file size for transfer info - fileInfo, err := os.Stat(filePath) - if err != nil { - return nil, fmt.Errorf("failed to get file info: %w", err) + start := time.Now() + expiresAt := start.Add(nscDefaultExpiry) + if err := client.UploadArtifact(ctx, n.nsc, key, file, storage.UploadOpts{ + ExpiresAt: &expiresAt, + Length: fileInfo.Size(), + }); err != nil { + return nil, fmt.Errorf("upload Namespace artifact %q: %w", key, err) } duration := time.Since(start) @@ -176,7 +199,7 @@ func (n *NscStore) Upload(ctx context.Context, filePath, key string) (*TransferI return &TransferInfo{ BytesTransferred: bytesTransferred, TransferSpeed: averageSpeed, - RequestID: "", // NSC doesn't expose request IDs + RequestID: "", // The Namespace Storage API does not expose request IDs. Duration: duration, }, nil } @@ -185,41 +208,35 @@ func (n *NscStore) Download(ctx context.Context, key, filePath string) (*Transfe _, span := trace.Start(ctx, "NscStore.Download") defer span.End() - // Validate input parameters to prevent command injection - if err := validateKey(key); err != nil { - return nil, fmt.Errorf("invalid key: %w", err) - } - if err := validateFilePath(filePath); err != nil { - return nil, fmt.Errorf("invalid file path: %w", err) - } - - start := time.Now() - - // Execute nsc artifact download command - result, err := n.run(ctx, "", n.artifactArgs("download", key, filePath)...) + client, err := n.client.get(ctx) if err != nil { - return nil, fmt.Errorf("failed to execute nsc download command: %w", err) + return nil, fmt.Errorf("initialize Namespace client for download: %w", err) } - if result.ExitCode != 0 { - // Map a missing or expired artifact to ErrBlobNotFound so the restore - // treats it as a cache miss (invalidating the stale entry and continuing) - // rather than failing the build. - if strings.Contains(result.Stderr, "not found") || - strings.Contains(result.Stderr, "has expired") { - return nil, fmt.Errorf("%w: nsc key %s: %s", ErrBlobNotFound, key, strings.TrimSpace(result.Stderr)) + start := time.Now() + var bytesTransferred int64 + err = roko.NewRetrier( + roko.WithMaxAttempts(nscDownloadAttempts), + roko.WithStrategy(roko.ExponentialSubsecond(200*time.Millisecond)), + roko.WithJitterRange(0, 250*time.Millisecond), + ).DoWithContext(ctx, func(r *roko.Retrier) error { + var err error + bytesTransferred, err = n.downloadOnce(ctx, client, key, filePath) + if err == nil { + return nil } - return nil, fmt.Errorf("nsc download failed with exit code %d: %s", result.ExitCode, result.Stderr) - } - - // Get file size for transfer info - fileInfo, err := os.Stat(filePath) + if errors.Is(err, ErrBlobNotFound) || !isRetryableNscDownloadError(err) { + r.Break() + return err + } + slog.Warn("NSC cache download failed, retrying", "key", key, "error", err, "retrier", r.String()) + return err + }) if err != nil { - return nil, fmt.Errorf("failed to get downloaded file info: %w", err) + return nil, err } duration := time.Since(start) - bytesTransferred := fileInfo.Size() averageSpeed := calculateTransferSpeedMBps(bytesTransferred, duration) span.SetAttributes( @@ -228,103 +245,93 @@ func (n *NscStore) Download(ctx context.Context, key, filePath string) (*Transfe attribute.String("nsc_key", key), ) - // Refresh the artifact's TTL on access so hot caches stay alive, mirroring - // the S3 store's self-CopyObject refresh. Unlike CopyObject this is cheap, - // so we refresh on every restore rather than gating it behind a minimum - // interval. - n.refreshExpiry(ctx, key) + n.refreshExpiry(ctx, client, key) return &TransferInfo{ BytesTransferred: bytesTransferred, TransferSpeed: averageSpeed, - RequestID: "", // NSC doesn't expose request IDs + RequestID: "", // The Namespace Storage API does not expose request IDs. Duration: duration, }, nil } -// refreshExpiry pushes the artifact's expiry out to at least nscDefaultExpiry -// from now via `nsc artifact extend --ensure_minimum`. Using --ensure_minimum -// (rather than the additive --by) makes the refresh idempotent, so calling it -// on every restore keeps a hot cache alive without growing its expiry unbounded. -// -// This is best-effort: any failure is logged and swallowed so a restore never -// fails because its TTL could not be refreshed. -func (n *NscStore) refreshExpiry(ctx context.Context, key string) { - if !n.extendSupported(ctx) { - // `nsc artifact extend` is still being rolled out by Namespace. Update - // the nsc CLI as a contingency so the command becomes available. - slog.Debug("nsc artifact extend unavailable, updating nsc CLI") - if _, err := n.run(ctx, "", "nsc", "version", "update"); err != nil { - slog.Warn("failed to update nsc CLI, skipping cache TTL refresh (non-fatal)", - "key", key, "error", err) - return - } +func (n *NscStore) downloadOnce(ctx context.Context, client nscAPIClient, key, filePath string) (int64, error) { + body, err := client.ResolveArtifactStream(ctx, n.nsc, key) + httpStatus, hasHTTPStatus := nscDownloadHTTPStatus(err) + if status.Code(err) == codes.NotFound || hasHTTPStatus && httpStatus == 404 { + return 0, fmt.Errorf("%w: nsc key %s: %w", ErrBlobNotFound, key, err) + } + if err != nil { + return 0, fmt.Errorf("resolve Namespace artifact %q: %w", key, err) } - result, err := n.run(ctx, "", n.artifactArgs("extend", key, "--ensure_minimum", nscDefaultExpiry)...) + dest, err := os.Create(filePath) + if err != nil { + _ = body.Close() + return 0, fmt.Errorf("create download file: %w", err) + } + + bytesTransferred, copyErr := io.Copy(dest, body) + destCloseErr := dest.Close() + bodyCloseErr := body.Close() switch { - case err != nil: - slog.Warn("failed to refresh cache TTL, continuing (non-fatal)", "key", key, "error", err) - case result.ExitCode != 0: - slog.Warn("failed to refresh cache TTL, continuing (non-fatal)", - "key", key, "exit_code", result.ExitCode, "stderr", result.Stderr) + case status.Code(copyErr) == codes.NotFound: + return 0, fmt.Errorf("%w: nsc key %s: %w", ErrBlobNotFound, key, copyErr) + case copyErr != nil: + return 0, fmt.Errorf("download Namespace artifact %q: %w", key, copyErr) + case destCloseErr != nil: + return 0, fmt.Errorf("close download file: %w", destCloseErr) + case bodyCloseErr != nil: + return 0, fmt.Errorf("close Namespace artifact stream %q: %w", key, bodyCloseErr) default: - slog.Debug("refreshed cache TTL", "key", key) + return bytesTransferred, nil } } -// extendSupported reports whether the installed nsc CLI supports -// `nsc artifact extend`. During Namespace's rollout of this command, older CLIs -// exit non-zero for its --help. -func (n *NscStore) extendSupported(ctx context.Context) bool { - result, err := n.run(ctx, "", "nsc", "artifact", "extend", "--help") - return err == nil && result.ExitCode == 0 && - strings.Contains(result.Stdout, "--ensure_minimum") -} - -type CommandResult struct { - Stdout string - Stderr string - ExitCode int +func isRetryableNscDownloadError(err error) bool { + // net/http represents HTTP/2 stream resets with an unexported error type, + // so its stable Error prefix is the only available classification seam. + if strings.Contains(err.Error(), "stream error: stream ID ") { + return true + } + if statusCode, ok := nscDownloadHTTPStatus(err); ok { + return statusCode == 429 || statusCode >= 500 + } + switch status.Code(err) { + case codes.Aborted, codes.DeadlineExceeded, codes.Internal, codes.ResourceExhausted, codes.Unavailable: + return true + default: + return api.IsRetryableError(err) + } } -func runCommand(ctx context.Context, workingDir string, args ...string) (*CommandResult, error) { - _, span := trace.Start(ctx, "runCommand") - defer span.End() - - // Validate that args is not empty to prevent panic - if len(args) == 0 { - return nil, fmt.Errorf("no command provided") +func nscDownloadHTTPStatus(err error) (int, bool) { + if err == nil { + return 0, false } - - span.SetAttributes(attribute.StringSlice("command", args)) - - cr := &CommandResult{} - - // #nosec G204 - args are validated by callers (validateFilePath, validateKey) - // and this function is internal to the store package with controlled usage - cmd := exec.CommandContext(ctx, args[0], args[1:]...) - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - cmd.Env = os.Environ() // inherit the environment - - if workingDir != "" { - cmd.Dir = workingDir + _, statusText, ok := strings.Cut(err.Error(), "failed to download file: status ") + if !ok { + return 0, false } + fields := strings.Fields(statusText) + if len(fields) == 0 { + return 0, false + } + statusCode, err := strconv.Atoi(fields[0]) + return statusCode, err == nil +} - err := cmd.Run() +// refreshExpiry best-effort ensures the artifact has at least 24 hours until +// expiry. A refresh failure never turns a successful download into a failure. +func (n *NscStore) refreshExpiry(ctx context.Context, client nscAPIClient, key string) { + err := client.ExtendArtifact(ctx, &storagev1beta.ExtendArtifactRequest{ + Path: key, + Namespace: n.nsc, + EnsureMinimum: durationpb.New(nscDefaultExpiry), + }) if err != nil { - span.RecordError(err) - if exitError, ok := err.(*exec.ExitError); ok { - cr.ExitCode = exitError.ExitCode() - } else { - return nil, err - } + slog.Warn("failed to refresh cache TTL, continuing (non-fatal)", "key", key, "error", err) + return } - - cr.Stdout = stdout.String() - cr.Stderr = stderr.String() - - return cr, nil + slog.Debug("refreshed cache TTL", "key", key) } diff --git a/internal/cache/store/nsc_test.go b/internal/cache/store/nsc_test.go index dce9d0fc9c..66fdf47015 100644 --- a/internal/cache/store/nsc_test.go +++ b/internal/cache/store/nsc_test.go @@ -1,228 +1,193 @@ package store import ( + "bytes" "context" "errors" + "fmt" + "io" "os" "path/filepath" "strings" + "sync" + "sync/atomic" "testing" + "time" - "github.com/google/go-cmp/cmp" + storagev1beta "buf.build/gen/go/namespace/cloud/protocolbuffers/go/proto/namespace/cloud/storage/v1beta" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "namespacelabs.dev/integrations/api/storage" ) func TestNscStore_Interface(t *testing.T) { - // This test ensures that NscStore properly implements the Blob interface var _ Blob = (*NscStore)(nil) } -func TestValidateFilePath(t *testing.T) { - tests := []struct { - name string - filePath string - expectError bool - errorMsg string - }{ - { - name: "valid simple path", - filePath: "test.txt", - expectError: false, - }, - { - name: "valid relative path", - filePath: "dir/subdir/file.txt", - expectError: false, - }, - { - name: "valid absolute path", - filePath: "/tmp/test.txt", - expectError: false, - }, - { - name: "empty path", - filePath: "", - expectError: true, - errorMsg: "file path cannot be empty", - }, - { - name: "path with semicolon", - filePath: "file;rm -rf /", - expectError: true, - errorMsg: "file path contains potentially dangerous character: ;", - }, - { - name: "path with ampersand", - filePath: "file&malicious", - expectError: true, - errorMsg: "file path contains potentially dangerous character: &", - }, - { - name: "path with pipe", - filePath: "file|cat /etc/passwd", - expectError: true, - errorMsg: "file path contains potentially dangerous character: |", - }, - { - name: "path with backtick", - filePath: "file`whoami`", - expectError: true, - errorMsg: "file path contains potentially dangerous character: `", - }, - { - name: "path with dollar sign", - filePath: "file$(whoami)", - expectError: true, - errorMsg: "file path contains potentially dangerous character: $", - }, - { - name: "path traversal attempt", - filePath: "../../../etc/passwd", - expectError: true, - errorMsg: "file path contains path traversal sequence", - }, - { - name: "path with quotes", - filePath: `file"test"`, - expectError: true, - errorMsg: `file path contains potentially dangerous character: "`, - }, +type fakeNscAPIClient struct { + upload func(context.Context, string, string, io.Reader, storage.UploadOpts) error + resolve func(context.Context, string, string) (io.ReadCloser, error) + extend func(context.Context, *storagev1beta.ExtendArtifactRequest) error + close func() error + closeCalls atomic.Int32 +} + +func (c *fakeNscAPIClient) UploadArtifact(ctx context.Context, nsc, path string, r io.Reader, opts storage.UploadOpts) error { + if c.upload == nil { + return nil } + return c.upload(ctx, nsc, path, r, opts) +} - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - err := validateFilePath(tt.filePath) - - if tt.expectError { - if err == nil { - t.Fatal("expected error, got nil") - } - if tt.errorMsg != "" && !strings.Contains(err.Error(), tt.errorMsg) { - t.Errorf("error %q does not contain %q", err.Error(), tt.errorMsg) - } - } else { - if err != nil { - t.Fatalf("validateFilePath: %v", err) - } - } - }) +func (c *fakeNscAPIClient) ResolveArtifactStream(ctx context.Context, nsc, path string) (io.ReadCloser, error) { + if c.resolve == nil { + return io.NopCloser(strings.NewReader("")), nil } + return c.resolve(ctx, nsc, path) } -func TestValidateKey(t *testing.T) { - tests := []struct { - name string - key string - expectError bool - errorMsg string - }{ - { - name: "valid simple key", - key: "mykey", - expectError: false, - }, - { - name: "valid key with path", - key: "builds/123/artifacts/report.txt", - expectError: false, - }, - { - name: "valid key with dots and underscores", - key: "test_file.tar.gz", - expectError: false, - }, - { - name: "valid key with hyphens", - key: "build-artifact-v1.0.0", - expectError: false, - }, - { - name: "empty key", - key: "", - expectError: true, - errorMsg: "key cannot be empty", - }, - { - name: "key too long", - key: string(make([]byte, 257)), // 257 characters - expectError: true, - errorMsg: "key too long (max 256 characters)", - }, - { - name: "key with invalid characters", - key: "key with spaces", - expectError: true, - errorMsg: "key contains invalid characters", - }, - { - name: "key with special characters", - key: "key@domain.com", - expectError: true, - errorMsg: "key contains invalid characters", - }, - { - name: "key with path traversal", - key: "../secret", - expectError: true, - errorMsg: "key contains potentially dangerous pattern: ../", - }, - { - name: "key with command injection attempt", - key: "file&&rm-rf", - expectError: true, - errorMsg: "key contains invalid characters", // regex catches it first - }, - { - name: "key with backtick", - key: "file`whoami`", - expectError: true, - errorMsg: "key contains invalid characters", // regex catches it first - }, +func (c *fakeNscAPIClient) ExtendArtifact(ctx context.Context, req *storagev1beta.ExtendArtifactRequest) error { + if c.extend == nil { + return nil } + return c.extend(ctx, req) +} - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - err := validateKey(tt.key) - - if tt.expectError { - if err == nil { - t.Fatal("expected error, got nil") - } - if tt.errorMsg != "" && !strings.Contains(err.Error(), tt.errorMsg) { - t.Errorf("error %q does not contain %q", err.Error(), tt.errorMsg) - } - } else { - if err != nil { - t.Fatalf("validateKey: %v", err) - } - } - }) +func (c *fakeNscAPIClient) Close() error { + c.closeCalls.Add(1) + if c.close == nil { + return nil + } + return c.close() +} + +func TestNscClient_IsLazy(t *testing.T) { + var initializeCalls atomic.Int32 + client := newNscClient(func(context.Context) (nscAPIClient, error) { + initializeCalls.Add(1) + return &fakeNscAPIClient{}, nil + }) + + if got := initializeCalls.Load(); got != 0 { + t.Fatalf("initialization calls after construction = %d, want 0", got) + } + if err := client.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + if got := initializeCalls.Load(); got != 0 { + t.Fatalf("initialization calls after closing unused client = %d, want 0", got) + } +} + +func TestNscClient_ConcurrentInitialization(t *testing.T) { + var initializeCalls atomic.Int32 + want := &fakeNscAPIClient{} + client := newNscClient(func(context.Context) (nscAPIClient, error) { + initializeCalls.Add(1) + return want, nil + }) + + const callers = 32 + results := make(chan nscAPIClient, callers) + errs := make(chan error, callers) + var wg sync.WaitGroup + for range callers { + wg.Add(1) + go func() { + defer wg.Done() + got, err := client.get(t.Context()) + results <- got + errs <- err + }() + } + wg.Wait() + close(results) + close(errs) + + for err := range errs { + if err != nil { + t.Errorf("get: %v", err) + } + } + for got := range results { + if got != want { + t.Errorf("get returned %p, want %p", got, want) + } + } + if got := initializeCalls.Load(); got != 1 { + t.Errorf("initialization calls = %d, want 1", got) + } +} + +func TestNscClient_InitializationError(t *testing.T) { + var initializeCalls atomic.Int32 + wantErr := errors.New("initialization failed") + client := newNscClient(func(context.Context) (nscAPIClient, error) { + initializeCalls.Add(1) + return nil, wantErr + }) + + for range 2 { + got, err := client.get(t.Context()) + if got != nil { + t.Errorf("get returned client %v, want nil", got) + } + if !errors.Is(err, wantErr) { + t.Errorf("get error = %v, want %v", err, wantErr) + } + } + if got := initializeCalls.Load(); got != 1 { + t.Errorf("initialization calls = %d, want 1", got) } } -func TestRunCommandValidation(t *testing.T) { - ctx := t.Context() +func TestNscClient_CloseOnce(t *testing.T) { + want := &fakeNscAPIClient{} + client := newNscClient(func(context.Context) (nscAPIClient, error) { + return want, nil + }) + if _, err := client.get(t.Context()); err != nil { + t.Fatalf("get: %v", err) + } + + for range 2 { + if err := client.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + } + if got := want.closeCalls.Load(); got != 1 { + t.Errorf("close calls = %d, want 1", got) + } +} - // Test empty args - result, err := runCommand(ctx, "" /* no args */) - if err == nil { - t.Fatal("expected error, got nil") +func TestNscClient_CloseError(t *testing.T) { + wantErr := errors.New("close failed") + want := &fakeNscAPIClient{close: func() error { return wantErr }} + client := newNscClient(func(context.Context) (nscAPIClient, error) { + return want, nil + }) + if _, err := client.get(t.Context()); err != nil { + t.Fatalf("get: %v", err) } - if !strings.Contains(err.Error(), "no command provided") { - t.Errorf("error %q does not contain %q", err.Error(), "no command provided") + + for range 2 { + if err := client.Close(); !errors.Is(err, wantErr) { + t.Fatalf("Close error = %v, want %v", err, wantErr) + } } - if result != nil { - t.Errorf("expected nil result, got %v", result) + if got := want.closeCalls.Load(); got != 1 { + t.Errorf("close calls = %d, want 1", got) } } func TestParseNscNamespace(t *testing.T) { tests := []struct { - name string - url string - namespace string - wantErr bool + name string + url string + nsc string + wantErr bool }{ - {name: "nsc with namespace", url: "nsc://my-namespace", namespace: "my-namespace"}, + {name: "nsc with namespace", url: "nsc://my-namespace", nsc: "my-namespace"}, {name: "not nsc", url: "s3://my-bucket", wantErr: true}, {name: "nsc without namespace", url: "nsc://", wantErr: true}, {name: "invalid url", url: "nsc://host:notaport", wantErr: true}, @@ -230,227 +195,438 @@ func TestParseNscNamespace(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - namespace, err := parseNscNamespace(tt.url) + nsc, err := parseNscNamespace(tt.url) if (err != nil) != tt.wantErr { t.Fatalf("parseNscNamespace(%q) err = %v, wantErr %v", tt.url, err, tt.wantErr) } - if namespace != tt.namespace { - t.Errorf("parseNscNamespace(%q) namespace = %q, want %q", tt.url, namespace, tt.namespace) + if nsc != tt.nsc { + t.Errorf("parseNscNamespace(%q) namespace = %q, want %q", tt.url, nsc, tt.nsc) } }) } } -// fakeRunner records the args of the last command and returns a successful result. -func fakeRunner(captured *[]string) commandRunner { - return func(_ context.Context, _ string, args ...string) (*CommandResult, error) { - *captured = args - return &CommandResult{}, nil +func newTestNscStore(t *testing.T, api nscAPIClient) *NscStore { + t.Helper() + holder := newNscClient(func(context.Context) (nscAPIClient, error) { + return api, nil + }) + store, err := NewNscStore("nsc://test-namespace", holder) + if err != nil { + t.Fatalf("NewNscStore: %v", err) } + return store } -func TestNscStore_PassesNamespace(t *testing.T) { - ctx := t.Context() - - tmpDir := t.TempDir() - testFile := filepath.Join(tmpDir, "test.txt") - if err := os.WriteFile(testFile, []byte("test content"), 0o600); err != nil { +func TestNscStore_Upload(t *testing.T) { + content := []byte("cache content") + filePath := filepath.Join(t.TempDir(), "cache archive") + if err := os.WriteFile(filePath, content, 0o600); err != nil { t.Fatalf("WriteFile: %v", err) } - var captured []string - store := &NscStore{namespace: "my-namespace", run: fakeRunner(&captured)} + var gotNsc, gotPath string + var gotContent []byte + var gotOpts storage.UploadOpts + api := &fakeNscAPIClient{upload: func(_ context.Context, nsc, path string, r io.Reader, opts storage.UploadOpts) error { + gotNsc = nsc + gotPath = path + gotOpts = opts + var err error + gotContent, err = io.ReadAll(r) + return err + }} + store := newTestNscStore(t, api) - if _, err := store.Upload(ctx, testFile, "key"); err != nil { + before := time.Now() + info, err := store.Upload(t.Context(), filePath, "artifact-key") + after := time.Now() + if err != nil { t.Fatalf("Upload: %v", err) } + if gotNsc != "test-namespace" { + t.Errorf("namespace = %q, want test-namespace", gotNsc) + } + if gotPath != "artifact-key" { + t.Errorf("path = %q, want artifact-key", gotPath) + } + if !bytes.Equal(gotContent, content) { + t.Errorf("content = %q, want %q", gotContent, content) + } + if gotOpts.Length != int64(len(content)) { + t.Errorf("Length = %d, want %d", gotOpts.Length, len(content)) + } + if gotOpts.ExpiresAt == nil || gotOpts.ExpiresAt.Before(before.Add(nscDefaultExpiry)) || gotOpts.ExpiresAt.After(after.Add(nscDefaultExpiry)) { + t.Errorf("ExpiresAt = %v, want between %v and %v", gotOpts.ExpiresAt, before.Add(nscDefaultExpiry), after.Add(nscDefaultExpiry)) + } + if info.BytesTransferred != int64(len(content)) { + t.Errorf("BytesTransferred = %d, want %d", info.BytesTransferred, len(content)) + } +} - wantArgs := []string{"nsc", "artifact", "upload", testFile, "key", "--expires_in", "24h", "--namespace", "my-namespace"} - if diff := cmp.Diff(wantArgs, captured); diff != "" { - t.Errorf("upload args mismatch (-want +got):\n%s", diff) +func TestNscStore_UploadFailure(t *testing.T) { + wantErr := errors.New("upload failed") + api := &fakeNscAPIClient{upload: func(context.Context, string, string, io.Reader, storage.UploadOpts) error { + return wantErr + }} + store := newTestNscStore(t, api) + filePath := filepath.Join(t.TempDir(), "cache") + if err := os.WriteFile(filePath, []byte("data"), 0o600); err != nil { + t.Fatalf("WriteFile: %v", err) + } + + if _, err := store.Upload(t.Context(), filePath, "key"); !errors.Is(err, wantErr) { + t.Fatalf("Upload error = %v, want %v", err, wantErr) } } -// isCommand reports whether args starts with the given nsc subcommand tokens. -func isCommand(args []string, tokens ...string) bool { - if len(args) < len(tokens) { - return false +func TestNscStore_UploadOpensFileBeforeInitializingClient(t *testing.T) { + var initializeCalls atomic.Int32 + holder := newNscClient(func(context.Context) (nscAPIClient, error) { + initializeCalls.Add(1) + return &fakeNscAPIClient{}, nil + }) + store, err := NewNscStore("nsc://ns", holder) + if err != nil { + t.Fatalf("NewNscStore: %v", err) + } + + if _, err := store.Upload(t.Context(), filepath.Join(t.TempDir(), "missing"), "key"); err == nil { + t.Fatal("Upload: expected error") } - for i, tok := range tokens { - if args[i] != tok { - return false - } + if got := initializeCalls.Load(); got != 0 { + t.Errorf("initialization calls = %d, want 0", got) } - return true } -// recordingRunner records every command invocation and delegates the result to -// respond. When respond is nil, every command succeeds. Download commands write -// their destination file so the store's os.Stat succeeds. -func recordingRunner(calls *[][]string, respond func(args []string) (*CommandResult, error)) commandRunner { - return func(_ context.Context, _ string, args ...string) (*CommandResult, error) { - *calls = append(*calls, args) +type trackingReadCloser struct { + io.Reader + closed bool + closeErr error +} - if isCommand(args, "nsc", "artifact", "download") { - // download args: nsc artifact download ... - if err := os.WriteFile(args[4], []byte("data"), 0o600); err != nil { - return nil, err - } - } +func (r *trackingReadCloser) Close() error { + r.closed = true + return r.closeErr +} - if respond != nil { - return respond(args) - } - return &CommandResult{}, nil - } +type errorReader struct { + err error } -func TestNscStore_RefreshesTTLOnDownload(t *testing.T) { - ctx := t.Context() - dest := filepath.Join(t.TempDir(), "out.txt") +func (r errorReader) Read([]byte) (int, error) { + return 0, r.err +} - var calls [][]string - store := &NscStore{namespace: "my-namespace", run: recordingRunner(&calls, nil)} +type partialErrorReader struct { + sent bool +} - if _, err := store.Download(ctx, "key", dest); err != nil { - t.Fatalf("Download: %v", err) +func (r *partialErrorReader) Read(p []byte) (int, error) { + if r.sent { + return 0, io.ErrUnexpectedEOF + } + r.sent = true + return copy(p, "partial"), nil +} + +func TestNscStore_Download(t *testing.T) { + body := &trackingReadCloser{Reader: strings.NewReader("downloaded content")} + var gotNsc, gotPath string + var extendReq *storagev1beta.ExtendArtifactRequest + api := &fakeNscAPIClient{ + resolve: func(_ context.Context, nsc, path string) (io.ReadCloser, error) { + gotNsc = nsc + gotPath = path + return body, nil + }, + extend: func(_ context.Context, req *storagev1beta.ExtendArtifactRequest) error { + extendReq = req + return nil + }, } + store := newTestNscStore(t, api) + dest := filepath.Join(t.TempDir(), "downloaded") - wantExtend := []string{"nsc", "artifact", "extend", "key", "--ensure_minimum", "24h", "--namespace", "my-namespace"} - var gotExtend []string - for _, c := range calls { - if isCommand(c, "nsc", "artifact", "extend") && !isCommand(c, "nsc", "artifact", "extend", "--help") { - gotExtend = c - } + info, err := store.Download(t.Context(), "artifact-key", dest) + if err != nil { + t.Fatalf("Download: %v", err) } - if diff := cmp.Diff(wantExtend, gotExtend); diff != "" { - t.Errorf("extend args mismatch (-want +got):\n%s", diff) + if gotNsc != "test-namespace" || gotPath != "artifact-key" { + t.Errorf("ResolveArtifactStream(%q, %q), want (%q, %q)", gotNsc, gotPath, "test-namespace", "artifact-key") + } + if !body.closed { + t.Error("download body was not closed") + } + got, err := os.ReadFile(dest) + if err != nil { + t.Fatalf("ReadFile: %v", err) + } + if string(got) != "downloaded content" { + t.Errorf("downloaded content = %q", got) + } + if info.BytesTransferred != int64(len(got)) { + t.Errorf("BytesTransferred = %d, want %d", info.BytesTransferred, len(got)) + } + if extendReq == nil { + t.Fatal("ExtendArtifact was not called") + } + if extendReq.GetNamespace() != "test-namespace" || extendReq.GetPath() != "artifact-key" { + t.Errorf("ExtendArtifact request namespace/path = %q/%q", extendReq.GetNamespace(), extendReq.GetPath()) + } + if got := extendReq.GetEnsureMinimum().AsDuration(); got != nscDefaultExpiry { + t.Errorf("EnsureMinimum = %v, want %v", got, nscDefaultExpiry) } } -func TestNscStore_UpdatesCLIWhenExtendUnsupported(t *testing.T) { - ctx := t.Context() - dest := filepath.Join(t.TempDir(), "out.txt") +func TestNscStore_DownloadNotFound(t *testing.T) { + resolveCalls := 0 + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + resolveCalls++ + return nil, status.Error(codes.NotFound, "artifact expired") + }} + store := newTestNscStore(t, api) - var calls [][]string - respond := func(args []string) (*CommandResult, error) { - // Simulate an old CLI: `nsc artifact extend --help` exits non-zero. - if isCommand(args, "nsc", "artifact", "extend", "--help") { - return &CommandResult{ExitCode: 1}, nil - } - return &CommandResult{}, nil + _, err := store.Download(t.Context(), "missing", filepath.Join(t.TempDir(), "dest")) + if !errors.Is(err, ErrBlobNotFound) { + t.Fatalf("Download error = %v, want ErrBlobNotFound", err) } - store := &NscStore{namespace: "ns", run: recordingRunner(&calls, respond)} - - if _, err := store.Download(ctx, "key", dest); err != nil { - t.Fatalf("Download: %v", err) + if resolveCalls != 1 { + t.Errorf("resolve calls = %d, want 1", resolveCalls) } +} - var updated, extended bool - for _, c := range calls { - if isCommand(c, "nsc", "version", "update") { - updated = true - } - if isCommand(c, "nsc", "artifact", "extend", "key") { - extended = true - } +func TestNscStore_DownloadStreamNotFound(t *testing.T) { + resolveCalls := 0 + body := &trackingReadCloser{Reader: errorReader{err: status.Error(codes.NotFound, "artifact expired")}} + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + resolveCalls++ + return body, nil + }} + store := newTestNscStore(t, api) + + _, err := store.Download(t.Context(), "missing", filepath.Join(t.TempDir(), "dest")) + if !errors.Is(err, ErrBlobNotFound) { + t.Fatalf("Download error = %v, want ErrBlobNotFound", err) } - if !updated { - t.Error("expected nsc version update to run when extend is unsupported") + if resolveCalls != 1 { + t.Errorf("resolve calls = %d, want 1", resolveCalls) } - if !extended { - t.Error("expected nsc artifact extend to run after updating the CLI") + if !body.closed { + t.Error("download body was not closed") } } -func TestNscStore_UpdatesCLIWhenExtendHelpLacksEnsureMinimum(t *testing.T) { - ctx := t.Context() - dest := filepath.Join(t.TempDir(), "out.txt") +func TestNscStore_DownloadHTTPNotFound(t *testing.T) { + resolveCalls := 0 + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + resolveCalls++ + return nil, errors.New("failed to download file: status 404") + }} + store := newTestNscStore(t, api) - var calls [][]string - respond := func(args []string) (*CommandResult, error) { - if isCommand(args, "nsc", "artifact", "extend", "--help") { - return &CommandResult{ - ExitCode: 0, - Stdout: "Artifact-related activities.\n" + - "Usage:\n nsc artifact [command]\n" + - "Available Commands:\n" + - " download Download an artifact.\n" + - " upload Upload an artifact.\n", - }, nil - } - return &CommandResult{}, nil + _, err := store.Download(t.Context(), "missing", filepath.Join(t.TempDir(), "dest")) + if !errors.Is(err, ErrBlobNotFound) { + t.Fatalf("Download error = %v, want ErrBlobNotFound", err) + } + if resolveCalls != 1 { + t.Errorf("resolve calls = %d, want 1", resolveCalls) + } +} + +func TestNscStore_DownloadResolveFailure(t *testing.T) { + wantErr := status.Error(codes.PermissionDenied, "permission denied") + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + return nil, wantErr + }} + store := newTestNscStore(t, api) + + _, err := store.Download(t.Context(), "key", filepath.Join(t.TempDir(), "dest")) + if !errors.Is(err, wantErr) { + t.Fatalf("Download error = %v, want %v", err, wantErr) + } + if errors.Is(err, ErrBlobNotFound) { + t.Fatalf("Download error = %v, should not be ErrBlobNotFound", err) } - store := &NscStore{namespace: "ns", run: recordingRunner(&calls, respond)} +} - if _, err := store.Download(ctx, "key", dest); err != nil { +func TestNscStore_DownloadRetriesTransientResolveFailure(t *testing.T) { + resolveCalls := 0 + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + resolveCalls++ + if resolveCalls == 1 { + return nil, status.Error(codes.Unavailable, "unavailable") + } + return io.NopCloser(strings.NewReader("complete")), nil + }} + store := newTestNscStore(t, api) + dest := filepath.Join(t.TempDir(), "dest") + + if _, err := store.Download(t.Context(), "key", dest); err != nil { t.Fatalf("Download: %v", err) } + if resolveCalls != 2 { + t.Errorf("resolve calls = %d, want 2", resolveCalls) + } +} - var updated, extended bool - for _, c := range calls { - if isCommand(c, "nsc", "version", "update") { - updated = true +func TestNscStore_DownloadRetriesTransientHTTPStatus(t *testing.T) { + resolveCalls := 0 + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + resolveCalls++ + if resolveCalls == 1 { + return nil, errors.New("failed to download file: status 503") } - if isCommand(c, "nsc", "artifact", "extend", "key") { - extended = true + return io.NopCloser(strings.NewReader("complete")), nil + }} + store := newTestNscStore(t, api) + dest := filepath.Join(t.TempDir(), "dest") + + if _, err := store.Download(t.Context(), "key", dest); err != nil { + t.Fatalf("Download: %v", err) + } + if resolveCalls != 2 { + t.Errorf("resolve calls = %d, want 2", resolveCalls) + } +} + +func TestIsRetryableNscDownloadError_HTTPStatus(t *testing.T) { + tests := []struct { + status int + retryable bool + }{ + {status: 429, retryable: true}, + {status: 500, retryable: true}, + {status: 503, retryable: true}, + {status: 404, retryable: false}, + } + for _, test := range tests { + err := fmt.Errorf("resolve artifact: failed to download file: status %d", test.status) + if got := isRetryableNscDownloadError(err); got != test.retryable { + t.Errorf("isRetryableNscDownloadError(status %d) = %t, want %t", test.status, got, test.retryable) } } - if !updated { - t.Error("expected nsc version update to run when extend --help exits 0 but omits --ensure_minimum") +} + +func TestNscStore_DownloadClosesBodyWhenDestinationCreationFails(t *testing.T) { + body := &trackingReadCloser{Reader: strings.NewReader("data")} + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + return body, nil + }} + store := newTestNscStore(t, api) + + if _, err := store.Download(t.Context(), "key", t.TempDir()); err == nil { + t.Fatal("Download: expected error") } - if !extended { - t.Error("expected nsc artifact extend to run after updating the CLI") + if !body.closed { + t.Error("download body was not closed") } } -func TestNscStore_DownloadSucceedsWhenRefreshFails(t *testing.T) { - ctx := t.Context() - dest := filepath.Join(t.TempDir(), "out.txt") +func TestNscStore_DownloadClosesBodyWhenCopyFails(t *testing.T) { + wantErr := errors.New("read failed") + body := &trackingReadCloser{Reader: errorReader{err: wantErr}} + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + return body, nil + }} + store := newTestNscStore(t, api) - respond := func(args []string) (*CommandResult, error) { - if isCommand(args, "nsc", "artifact", "extend", "key") { - return &CommandResult{ExitCode: 1, Stderr: "boom"}, nil - } - return &CommandResult{}, nil + _, err := store.Download(t.Context(), "key", filepath.Join(t.TempDir(), "dest")) + if !errors.Is(err, wantErr) { + t.Fatalf("Download error = %v, want %v", err, wantErr) + } + if !body.closed { + t.Error("download body was not closed") } - var calls [][]string - store := &NscStore{namespace: "ns", run: recordingRunner(&calls, respond)} +} + +func TestNscStore_DownloadRetriesCopyFailureAndTruncatesDestination(t *testing.T) { + firstBody := &trackingReadCloser{Reader: &partialErrorReader{}} + secondBody := &trackingReadCloser{Reader: strings.NewReader("complete")} + resolveCalls := 0 + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + resolveCalls++ + if resolveCalls == 1 { + return firstBody, nil + } + return secondBody, nil + }} + store := newTestNscStore(t, api) + dest := filepath.Join(t.TempDir(), "dest") - if _, err := store.Download(ctx, "key", dest); err != nil { - t.Fatalf("Download should succeed despite a failed TTL refresh: %v", err) + if _, err := store.Download(t.Context(), "key", dest); err != nil { + t.Fatalf("Download: %v", err) + } + if resolveCalls != 2 { + t.Errorf("resolve calls = %d, want 2", resolveCalls) + } + if !firstBody.closed || !secondBody.closed { + t.Errorf("body closure = first:%t second:%t, want both closed", firstBody.closed, secondBody.closed) + } + content, err := os.ReadFile(dest) + if err != nil { + t.Fatalf("ReadFile: %v", err) + } + if got := string(content); got != "complete" { + t.Errorf("downloaded content = %q, want complete", got) } } -func TestNewNscStore_RequiresNamespace(t *testing.T) { - if _, err := NewNscStore("nsc://"); err == nil { - t.Error(`NewNscStore("nsc://"): expected error, got nil`) +func TestIsRetryableNscDownloadError_HTTP2StreamError(t *testing.T) { + err := fmt.Errorf("read response body: %w", errors.New("stream error: stream ID 1; CANCEL; received from peer")) + if !isRetryableNscDownloadError(err) { + t.Errorf("isRetryableNscDownloadError(%v) = false, want true", err) } } -// TestNscStore_ValidationShortCircuits ensures unsafe inputs are rejected before -// the CLI is ever invoked. -func TestNscStore_ValidationShortCircuits(t *testing.T) { - ctx := t.Context() - ran := false - store := &NscStore{namespace: "ns", run: func(context.Context, string, ...string) (*CommandResult, error) { - ran = true - return &CommandResult{}, nil +func TestNscStore_DownloadFailsWhenBodyCloseFails(t *testing.T) { + wantErr := errors.New("close failed") + body := &trackingReadCloser{Reader: strings.NewReader("data"), closeErr: wantErr} + api := &fakeNscAPIClient{resolve: func(context.Context, string, string) (io.ReadCloser, error) { + return body, nil }} + store := newTestNscStore(t, api) + + _, err := store.Download(t.Context(), "key", filepath.Join(t.TempDir(), "dest")) + if !errors.Is(err, wantErr) { + t.Fatalf("Download error = %v, want %v", err, wantErr) + } +} + +func TestNscStore_DownloadSucceedsWhenRefreshFails(t *testing.T) { + wantErr := errors.New("extend failed") + api := &fakeNscAPIClient{ + resolve: func(context.Context, string, string) (io.ReadCloser, error) { + return io.NopCloser(strings.NewReader("data")), nil + }, + extend: func(context.Context, *storagev1beta.ExtendArtifactRequest) error { + return wantErr + }, + } + store := newTestNscStore(t, api) - if _, err := store.Upload(ctx, "invalid;path", "valid-key"); err == nil { - t.Error("Upload with unsafe path: expected error, got nil") + if _, err := store.Download(t.Context(), "key", filepath.Join(t.TempDir(), "dest")); err != nil { + t.Fatalf("Download should succeed despite failed expiry refresh: %v", err) } - if _, err := store.Download(ctx, "invalid key with spaces", "dest.txt"); err == nil { - t.Error("Download with unsafe key: expected error, got nil") +} + +func TestNewNscStore_RequiresNamespaceAndClient(t *testing.T) { + holder := newNscClient(func(context.Context) (nscAPIClient, error) { + return &fakeNscAPIClient{}, nil + }) + if _, err := NewNscStore("nsc://", holder); err == nil { + t.Error(`NewNscStore("nsc://"): expected error`) } - if ran { - t.Error("nsc CLI should not run when input validation fails") + if _, err := NewNscStore("nsc://namespace", nil); err == nil { + t.Error("NewNscStore with nil client: expected error") } } func TestNewBlobStore_NscScheme(t *testing.T) { - blob, err := NewBlobStore(t.Context(), AgentManaged, "nsc://my-namespace") + holder := newNscClient(func(context.Context) (nscAPIClient, error) { + return &fakeNscAPIClient{}, nil + }) + blob, err := NewBlobStore(t.Context(), AgentManaged, "nsc://my-namespace", holder) if err != nil { t.Fatalf("NewBlobStore: %v", err) } @@ -458,123 +634,51 @@ func TestNewBlobStore_NscScheme(t *testing.T) { if !ok { t.Fatalf("NewBlobStore returned %T, want *NscStore", blob) } - if nsc.namespace != "my-namespace" { - t.Errorf("namespace = %q, want %q", nsc.namespace, "my-namespace") + if nsc.nsc != "my-namespace" { + t.Errorf("namespace = %q, want my-namespace", nsc.nsc) + } + + if _, err := NewBlobStore(t.Context(), AgentManaged, "nsc://my-namespace", nil); err == nil { + t.Error("NewBlobStore with nil Namespace client: expected error") } } -// TestNscStore_Integration runs integration tests if NSC CLI is available -// This test can be skipped if NSC is not installed or configured func TestNscStore_Integration(t *testing.T) { - // Skip this test if NSC_INTEGRATION_TEST environment variable is not set if os.Getenv("NSC_INTEGRATION_TEST") == "" { t.Skip("Skipping NSC integration test (set NSC_INTEGRATION_TEST=1 to run)") } - // "main" is the nsc CLI's default namespace. - store, err := NewNscStore("nsc://main") + holder := NewNscClient() + t.Cleanup(func() { + if err := holder.Close(); err != nil { + t.Errorf("close Namespace client: %v", err) + } + }) + store, err := NewNscStore("nsc://main", holder) if err != nil { t.Fatalf("NewNscStore: %v", err) } - ctx := t.Context() - - // Create temporary directories and files - tmpDir, err := os.MkdirTemp("", "nsc-integration-test") - if err != nil { - t.Fatalf("MkdirTemp: %v", err) - } - defer func() { _ = os.RemoveAll(tmpDir) }() - - // Create a test file - testFile := filepath.Join(tmpDir, "test-upload.txt") + testFile := filepath.Join(t.TempDir(), "test-upload.txt") testContent := "Hello from NSC integration test!" - err = os.WriteFile(testFile, []byte(testContent), 0o600) - if err != nil { + if err := os.WriteFile(testFile, []byte(testContent), 0o600); err != nil { t.Fatalf("WriteFile: %v", err) } - // Test upload key := "integration-test/test-file.txt" - transferInfo, err := store.Upload(ctx, testFile, key) - if err != nil { - t.Fatalf("Upload should succeed with valid NSC setup: %v", err) - } - - if transferInfo.BytesTransferred <= 0 { - t.Errorf("expected BytesTransferred > 0, got %d", transferInfo.BytesTransferred) - } - if transferInfo.TransferSpeed <= 0.0 { - t.Errorf("expected TransferSpeed > 0, got %f", transferInfo.TransferSpeed) - } - if transferInfo.Duration <= 0 { - t.Errorf("expected Duration > 0, got %v", transferInfo.Duration) - } - - // Test download - downloadFile := filepath.Join(tmpDir, "test-download.txt") - transferInfo, err = store.Download(ctx, key, downloadFile) - if err != nil { - t.Fatalf("Download should succeed: %v", err) + if _, err := store.Upload(t.Context(), testFile, key); err != nil { + t.Fatalf("Upload: %v", err) } - if transferInfo.BytesTransferred <= 0 { - t.Errorf("expected BytesTransferred > 0, got %d", transferInfo.BytesTransferred) + downloadFile := filepath.Join(t.TempDir(), "test-download.txt") + if _, err := store.Download(t.Context(), key, downloadFile); err != nil { + t.Fatalf("Download: %v", err) } - - // Verify downloaded content downloadedContent, err := os.ReadFile(downloadFile) if err != nil { t.Fatalf("ReadFile: %v", err) } if string(downloadedContent) != testContent { - t.Errorf("downloaded content mismatch: got %q, want %q", string(downloadedContent), testContent) + t.Errorf("downloaded content = %q, want %q", downloadedContent, testContent) } - - t.Logf("Upload: %d bytes at %.2f MB/s in %v", - transferInfo.BytesTransferred, - transferInfo.TransferSpeed, - transferInfo.Duration) -} - -// TestNscStore_DownloadNotFound checks the store-specific not-found mapping -func TestNscStore_DownloadNotFound(t *testing.T) { - ctx := t.Context() - dest := filepath.Join(t.TempDir(), "dest") - - t.Run("stderr not found maps to ErrBlobNotFound", func(t *testing.T) { - store := &NscStore{namespace: "ns", run: func(context.Context, string, ...string) (*CommandResult, error) { - return &CommandResult{ExitCode: 1, Stderr: "Error: artifact not found"}, nil - }} - _, err := store.Download(ctx, "valid-key", dest) - if !errors.Is(err, ErrBlobNotFound) { - t.Fatalf("Download err = %v, want ErrBlobNotFound", err) - } - }) - - t.Run("stderr expired maps to ErrBlobNotFound", func(t *testing.T) { - store := &NscStore{namespace: "ns", run: func(context.Context, string, ...string) (*CommandResult, error) { - return &CommandResult{ - ExitCode: 1, - Stderr: "Failed: the artifact has expired at 2026-07-14T02:06:24Z (request id: cbdorlqepas5e10ldvfvdpog40).", - }, nil - }} - _, err := store.Download(ctx, "valid-key", dest) - if !errors.Is(err, ErrBlobNotFound) { - t.Fatalf("Download err = %v, want ErrBlobNotFound", err) - } - }) - - t.Run("other failures are not ErrBlobNotFound", func(t *testing.T) { - store := &NscStore{namespace: "ns", run: func(context.Context, string, ...string) (*CommandResult, error) { - return &CommandResult{ExitCode: 1, Stderr: "connection refused"}, nil - }} - _, err := store.Download(ctx, "valid-key", dest) - if err == nil { - t.Fatal("Download: expected error, got nil") - } - if errors.Is(err, ErrBlobNotFound) { - t.Errorf("Download err = %v, should not be ErrBlobNotFound", err) - } - }) } diff --git a/scripts/generate-acknowledgements.sh b/scripts/generate-acknowledgements.sh index 6565d61174..6518d179df 100755 --- a/scripts/generate-acknowledgements.sh +++ b/scripts/generate-acknowledgements.sh @@ -32,7 +32,13 @@ export TEMPDIR="$(mktemp -d /tmp/generate-acknowledgements.XXXXXX)" export TEMPFILE="$(mktemp /tmp/acknowledgements.XXXXXX)" trap "rm -fr ${TEMPDIR} ${TEMPFILE}" EXIT -"${GO_LICENSES}" save . --save_path="${TEMPDIR}" --force +# Namespace's Buf-generated Go modules do not currently include license files. +# Track their request to publish them at: +# https://buildkite-corp.slack.com/archives/C0BBGCKJ2HJ/p1787538095689349 +"${GO_LICENSES}" save . \ + --ignore buf.build/gen/go/namespace/cloud \ + --save_path="${TEMPDIR}" \ + --force # Build acknowledgements file cat > "${TEMPFILE}" <