Skip to content
Draft
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
1 change: 1 addition & 0 deletions MODULE.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ use_repo(
"com_github_aws_aws_sdk_go_v2_service_sts",
"com_github_bazelbuild_buildtools",
"com_github_bazelbuild_remote_apis",
"com_github_buildbarn_go_cdc",
"com_github_buildbarn_go_sha256tree",
"com_github_fxtlabs_primes",
"com_github_go_jose_go_jose_v3",
Expand Down
54 changes: 22 additions & 32 deletions cmd/bb_copy/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,52 +41,42 @@ func main() {
return util.StatusWrapf(err, "Failed to read configuration from %s", os.Args[1])
}

if configuration.TraversalConcurrency <= 0 {
return status.Errorf(codes.InvalidArgument, "traversal_concurrency must be > 0, got %d", configuration.TraversalConcurrency)
}

grpcClientFactory := grpc.NewBaseClientFactory(grpc.BaseClientDialer, nil, nil, nil)

blobAccessCreator := blobstore_configuration.NewCASBlobAccessCreator(
grpcClientFactory,
int(configuration.MaximumMessageSizeBytes),
bb_zstd.NewPoolFromConfiguration(nil),
)
source, err := blobstore_configuration.NewBlobAccessFromConfiguration(
dependenciesGroup,
configuration.Source,
blobAccessCreator,
)
zstdPool := bb_zstd.NewPoolFromConfiguration(nil)

source, _, _, _, _, err := blobstore_configuration.NewCASFromConfiguration(dependenciesGroup, configuration.Source, grpcClientFactory, int(configuration.MaximumMessageSizeBytes), zstdPool)
if err != nil {
return util.StatusWrap(err, "Failed to create source")
}
sink, err := blobstore_configuration.NewBlobAccessFromConfiguration(
dependenciesGroup,
configuration.Sink,
blobAccessCreator,
)

sink, _, _, _, _, err := blobstore_configuration.NewCASFromConfiguration(dependenciesGroup, configuration.Sink, grpcClientFactory, int(configuration.MaximumMessageSizeBytes), zstdPool)
if err != nil {
return util.StatusWrap(err, "Failed to create sink")
}
blobReplicator, err := blobstore_configuration.NewBlobReplicatorFromConfiguration(
dependenciesGroup,
configuration.Replicator,
source.BlobAccess,
sink,
blobstore_configuration.NewCASBlobReplicatorCreator(grpcClientFactory),
)

instanceName, err := digest.NewInstanceName(configuration.InstanceName)
if err != nil {
return util.StatusWrap(err, "Invalid instance name")
}

replicator := cas.NewReplicator(source, sink, instanceName)
if err != nil {
return util.StatusWrap(err, "Failed to create blob replicator")
}
nestedReplicator := cas.NewNestedBlobReplicator(
cas.NewBlobAccessReplicator(blobReplicator),
replicator,
int(configuration.MaximumMessageSizeBytes),
cas.NewBlobAccessMessageReader[remoteexecution.Action](source.BlobAccess, int(configuration.MaximumMessageSizeBytes)),
cas.NewBlobAccessMessageReader[remoteexecution.Directory](source.BlobAccess, int(configuration.MaximumMessageSizeBytes)),
cas.NewBlobAccessStreamReader(source.BlobAccess),
sink.DigestKeyFormat,
cas.NewMessageReader[remoteexecution.Action](source, int(configuration.MaximumMessageSizeBytes)),
cas.NewMessageReader[remoteexecution.Directory](source, int(configuration.MaximumMessageSizeBytes)),
cas.NewStreamReader(source),
sink.GetDigestKeyFormat(),
)

instanceName, err := digest.NewInstanceName(configuration.InstanceName)
if err != nil {
return util.StatusWrap(err, "Invalid instance name")
}
digestFunction, err := instanceName.GetDigestFunction(configuration.DigestFunction, 0)
if err != nil {
return util.StatusWrap(err, "Invalid digest function")
Expand All @@ -105,7 +95,7 @@ func main() {
if err != nil {
return util.StatusWrapf(err, "Invalid blob digest at index %d", i)
}
if err := blobReplicator.ReplicateMultiple(ctx, blobDigest.ToSingletonSet()); err != nil {
if err := replicator.Replicate(ctx, blobDigest.ToSingletonSet()); err != nil {
return util.StatusWrapf(err, "Failed to schedule replication of blob with digest %#v", blobDigest.String())
}
}
Expand Down
4 changes: 2 additions & 2 deletions cmd/bb_replicator/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ func main() {
return util.StatusWrap(err, "Failed to apply global configuration options")
}

blobAccessCreator := blobstore_configuration.NewCASBlobAccessCreator(
blobAccessCreator := blobstore_configuration.NewCSBlobAccessCreator(
grpcClientFactory,
int(configuration.MaximumMessageSizeBytes),
bb_zstd.NewPoolFromConfiguration(nil),
Expand All @@ -59,7 +59,7 @@ func main() {
configuration.Replicator,
source.BlobAccess,
sink,
blobstore_configuration.NewCASBlobReplicatorCreator(grpcClientFactory),
blobstore_configuration.NewCSBlobReplicatorCreator(grpcClientFactory),
)
if err != nil {
return util.StatusWrap(err, "Failed to create replicator")
Expand Down
2 changes: 2 additions & 0 deletions cmd/bb_storage/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,12 @@ go_library(
"//pkg/auth",
"//pkg/auth/configuration",
"//pkg/blobstore",
"//pkg/blobstore/cdc",
"//pkg/blobstore/configuration",
"//pkg/blobstore/grpcservers",
"//pkg/builder",
"//pkg/capabilities",
"//pkg/cas",
"//pkg/global",
"//pkg/grpc",
"//pkg/program",
Expand Down
70 changes: 48 additions & 22 deletions cmd/bb_storage/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,12 @@ import (
"github.com/buildbarn/bb-storage/pkg/auth"
auth_configuration "github.com/buildbarn/bb-storage/pkg/auth/configuration"
"github.com/buildbarn/bb-storage/pkg/blobstore"
"github.com/buildbarn/bb-storage/pkg/blobstore/cdc"
blobstore_configuration "github.com/buildbarn/bb-storage/pkg/blobstore/configuration"
"github.com/buildbarn/bb-storage/pkg/blobstore/grpcservers"
"github.com/buildbarn/bb-storage/pkg/builder"
"github.com/buildbarn/bb-storage/pkg/capabilities"
"github.com/buildbarn/bb-storage/pkg/cas"
"github.com/buildbarn/bb-storage/pkg/global"
bb_grpc "github.com/buildbarn/bb-storage/pkg/grpc"
"github.com/buildbarn/bb-storage/pkg/program"
Expand Down Expand Up @@ -54,34 +56,58 @@ func main() {
var cacheCapabilitiesAuthorizers []auth.Authorizer

// Content Addressable Storage (CAS).
var contentAddressableStorageInfo *blobstore_configuration.BlobAccessInfo
var contentAddressableStorage blobstore.BlobAccess
if configuration.ContentAddressableStorage != nil {
info, authorizedBackend, allAuthorizers, err := newScannableBlobAccess(
var contentAddressableStorage cas.ContentAddressableStorage
var authorizedContentAddressableStorage cas.ContentAddressableStorage
if configuration.ContentAddressableStorageServer != nil {
var chunkStorage, chunkListStorage blobstore.BlobAccess
var cdcParametersFetcher cdc.ParametersFetcher
contentAddressableStorage, chunkStorage, chunkListStorage, _, cdcParametersFetcher, err = blobstore_configuration.NewCASFromConfiguration(
dependenciesGroup,
configuration.ContentAddressableStorage,
blobstore_configuration.NewCASBlobAccessCreator(
grpcClientFactory,
int(configuration.MaximumMessageSizeBytes),
zstdPool,
),
configuration.ContentAddressableStorageServer.ContentAddressableStorage,
grpcClientFactory,
int(configuration.MaximumMessageSizeBytes),
zstdPool,
)
if err != nil {
return util.StatusWrap(err, "Failed to create Content Addressable Storage")
}

// Create authorizers.
getAuthorizer, err := auth_configuration.DefaultAuthorizerFactory.NewAuthorizerFromConfiguration(configuration.ContentAddressableStorageServer.GetAuthorizer, dependenciesGroup, grpcClientFactory)
if err != nil {
return util.StatusWrap(err, "Failed to create Get() authorizer for Content Addressable Storage")
}
putAuthorizer, err := auth_configuration.DefaultAuthorizerFactory.NewAuthorizerFromConfiguration(configuration.ContentAddressableStorageServer.PutAuthorizer, dependenciesGroup, grpcClientFactory)
if err != nil {
return util.StatusWrap(err, "Failed to create Put() authorizer for Content Addressable Storage")
}
findMissingAuthorizer, err := auth_configuration.DefaultAuthorizerFactory.NewAuthorizerFromConfiguration(configuration.ContentAddressableStorageServer.FindMissingAuthorizer, dependenciesGroup, grpcClientFactory)
if err != nil {
return util.StatusWrap(err, "Failed to create FindMissing() authorizer for Content Addressable Storage")
}

// Create authorized versions of the backends.
authorizedChunkStorage := blobstore.NewAuthorizingBlobAccess(chunkStorage, getAuthorizer, putAuthorizer, findMissingAuthorizer)
authorizedChunkListStorage := blobstore.NewAuthorizingBlobAccess(chunkListStorage, getAuthorizer, putAuthorizer, findMissingAuthorizer)
authorizedChunkListFetcher := blobstore.NewBlobAccessChunkListFetcher(authorizedChunkListStorage, int(configuration.MaximumMessageSizeBytes))
authorizedContentAddressableStorage = cas.NewBlobAccessContentAddressableStorage(
authorizedChunkStorage,
authorizedChunkListStorage,
authorizedChunkListFetcher,
cdcParametersFetcher,
contentAddressableStorage.GetDigestKeyFormat(),
)
// Create the Chunk Storage (CS).
cacheCapabilitiesProviders = append(
cacheCapabilitiesProviders,
info.BlobAccess,
chunkStorage,
capabilities.NewStaticProvider(&remoteexecution.ServerCapabilities{
CacheCapabilities: &remoteexecution.CacheCapabilities{
SupportedCompressors: configuration.SupportedCompressors,
},
}),
)
cacheCapabilitiesAuthorizers = append(cacheCapabilitiesAuthorizers, allAuthorizers...)
contentAddressableStorageInfo = &info
contentAddressableStorage = authorizedBackend
cacheCapabilitiesAuthorizers = append(cacheCapabilitiesAuthorizers, getAuthorizer, putAuthorizer, findMissingAuthorizer)
}

// Action Cache (AC).
Expand All @@ -91,7 +117,7 @@ func main() {
dependenciesGroup,
configuration.ActionCache,
blobstore_configuration.NewACBlobAccessCreator(
contentAddressableStorageInfo,
contentAddressableStorage,
grpcClientFactory,
int(configuration.MaximumMessageSizeBytes),
),
Expand Down Expand Up @@ -192,19 +218,19 @@ func main() {
if err := bb_grpc.NewServersFromConfigurationAndServe(
configuration.GrpcServers,
func(s grpc.ServiceRegistrar) {
if contentAddressableStorage != nil {
if authorizedContentAddressableStorage != nil {
contentAddressableStorageServer := grpcservers.NewContentAddressableStorageServer(
authorizedContentAddressableStorage,
configuration.MaximumMessageSizeBytes,
)
remoteexecution.RegisterContentAddressableStorageServer(
s,
grpcservers.NewContentAddressableStorageServer(
contentAddressableStorage,
configuration.MaximumMessageSizeBytes,
),
contentAddressableStorageServer,
)
bytestream.RegisterByteStreamServer(
s,
grpcservers.NewByteStreamServer(
contentAddressableStorage,
1<<16,
authorizedContentAddressableStorage,
zstdPool,
),
)
Expand Down
1 change: 1 addition & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ require (
github.com/bazelbuild/buildtools v0.0.0-20260527135131-3b47c424ecf5
github.com/bazelbuild/remote-apis v0.0.0-20260331222004-becdd8f9ff81
github.com/bazelbuild/rules_go v0.62.0
github.com/buildbarn/go-cdc v0.0.9
github.com/buildbarn/go-sha256tree v0.0.0-20250310211320-0f70f20e855b
github.com/fxtlabs/primes v0.0.0-20150821004651-dad82d10a449
github.com/go-jose/go-jose/v3 v3.0.5
Expand Down
2 changes: 2 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,8 @@ github.com/bazelbuild/rules_go v0.62.0/go.mod h1:6YghDRf6l3FSiAwncK+Ww9jE1naoQzN
github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/buildbarn/go-cdc v0.0.9 h1:bWfgn92ed8Oo2zZKJdMAfB0APGz7Q8zvnqUn3hPuihM=
github.com/buildbarn/go-cdc v0.0.9/go.mod h1:KUMqSMvoRlby3uak9aKIvgz3KgNqwm2CMUoVX1EDr8k=
github.com/buildbarn/go-sha256tree v0.0.0-20250310211320-0f70f20e855b h1:IKUxixGBm9UxobU7c248z0BF0ojG19uoSLz8MFZM/KA=
github.com/buildbarn/go-sha256tree v0.0.0-20250310211320-0f70f20e855b/go.mod h1:e7g3/yWApcg+PpDqd4eQEEV8pexQmfCgK3frP+1Wuvk=
github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU=
Expand Down
19 changes: 18 additions & 1 deletion internal/mock/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -156,8 +156,9 @@ gomock(
name = "cas",
out = "cas.go",
interfaces = [
"StreamReader",
"ContentAddressableStorage",
"Replicator",
"StreamReader",
],
library = "//pkg/cas",
mockgen_model_library = "@org_uber_go_mock//mockgen/model",
Expand All @@ -178,6 +179,18 @@ gomock(
source = "//pkg/cas:message_reader.go",
)

gomock(
name = "cdc",
out = "cdc.go",
interfaces = [
"ParametersFetcher",
],
library = "//pkg/blobstore/cdc",
mockgen_model_library = "@org_uber_go_mock//mockgen/model",
mockgen_tool = "@org_uber_go_mock//mockgen",
package = "mock",
)

gomock(
name = "clock",
out = "clock.go",
Expand Down Expand Up @@ -395,6 +408,7 @@ go_library(
"capabilities.go",
"cas.go",
"cas_message_reader.go",
"cdc.go",
"clock.go",
"cloud_aws.go",
"cloud_gcp.go",
Expand All @@ -419,10 +433,13 @@ go_library(
"//pkg/auth",
"//pkg/blobstore",
"//pkg/blobstore/buffer",
"//pkg/blobstore/cdc",
"//pkg/blobstore/chunklist",
"//pkg/blobstore/local",
"//pkg/blobstore/sharding",
"//pkg/blobstore/slicing",
"//pkg/builder",
"//pkg/cas",
"//pkg/clock",
"//pkg/cloud/gcp",
"//pkg/digest",
Expand Down
17 changes: 4 additions & 13 deletions pkg/blobstore/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ go_library(
"authorizing_blob_access.go",
"blob_access.go",
"cas_read_buffer_factory.go",
"chunk_list_fetcher.go",
"cls_read_buffer_factory.go",
"deadline_enforcing_blob_access.go",
"demultiplexing_blob_access.go",
"empty_blob_injecting_blob_access.go",
Expand All @@ -21,7 +23,6 @@ go_library(
"metrics_blob_access.go",
"read_buffer_factory.go",
"read_canarying_blob_access.go",
"reference_expanding_blob_access.go",
"validation_caching_read_buffer_factory.go",
"visit_topologically_sorted_tree.go",
"zip_reading_blob_access.go",
Expand All @@ -32,21 +33,18 @@ go_library(
deps = [
"//pkg/auth",
"//pkg/blobstore/buffer",
"//pkg/blobstore/chunklist",
"//pkg/blobstore/slicing",
"//pkg/capabilities",
"//pkg/clock",
"//pkg/cloud/aws",
"//pkg/cloud/gcp",
"//pkg/digest",
"//pkg/eviction",
"//pkg/proto/blobstore/chunklist",
"//pkg/proto/fsac",
"//pkg/proto/icas",
"//pkg/proto/iscc",
"//pkg/util",
"//pkg/zstd",
"@bazel_remote_apis//build/bazel/remote/execution/v2:remote_execution_go_proto",
"@com_github_aws_aws_sdk_go_v2//aws",
"@com_github_aws_aws_sdk_go_v2_service_s3//:s3",
"@com_github_prometheus_client_golang//prometheus",
"@org_golang_google_grpc//codes",
"@org_golang_google_grpc//status",
Expand All @@ -67,7 +65,6 @@ go_test(
"existence_caching_blob_access_test.go",
"hierarchical_instance_names_blob_access_test.go",
"read_canarying_blob_access_test.go",
"reference_expanding_blob_access_test.go",
"validation_caching_read_buffer_factory_test.go",
"visit_topologically_sorted_tree_test.go",
"zip_reading_blob_access_test.go",
Expand All @@ -79,16 +76,10 @@ go_test(
"//pkg/blobstore/buffer",
"//pkg/digest",
"//pkg/eviction",
"//pkg/proto/icas",
"//pkg/testutil",
"//pkg/util",
"//pkg/zstd",
"@bazel_remote_apis//build/bazel/remote/execution/v2:remote_execution_go_proto",
"@bazel_remote_apis//build/bazel/semver:semver_go_proto",
"@com_github_aws_aws_sdk_go_v2//aws",
"@com_github_aws_aws_sdk_go_v2_service_s3//:s3",
"@com_github_aws_aws_sdk_go_v2_service_s3//types",
"@com_github_klauspost_compress//zstd",
"@com_github_stretchr_testify//require",
"@org_golang_google_grpc//codes",
"@org_golang_google_grpc//status",
Expand Down
Loading