Skip to content

Commit ed3a308

Browse files
committed
Set the resource-name header as part of Bytestream requests
Even though I don't think we should do anything on the Buildbarn side to attempt to interpret these headers, let's at least emit it as part of outgoing Bytestream requests. That's likely going to make life easier for people who want to do load balancing this way. More details: bazelbuild/remote-apis#303
1 parent f24d5ce commit ed3a308

3 files changed

Lines changed: 20 additions & 9 deletions

File tree

pkg/blobstore/grpcclients/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ go_library(
2525
"@org_golang_google_genproto_googleapis_bytestream//:bytestream",
2626
"@org_golang_google_grpc//:grpc",
2727
"@org_golang_google_grpc//codes",
28+
"@org_golang_google_grpc//metadata",
2829
"@org_golang_google_grpc//status",
2930
],
3031
)

pkg/blobstore/grpcclients/cas_blob_access.go

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import (
1414

1515
"google.golang.org/genproto/googleapis/bytestream"
1616
"google.golang.org/grpc"
17+
"google.golang.org/grpc/metadata"
1718
)
1819

1920
type casBlobAccess struct {
@@ -61,11 +62,17 @@ func (r *byteStreamChunkReader) Close() {
6162
}
6263
}
6364

65+
const resourceNameHeader = "build.bazel.remote.execution.v2.resource-name"
66+
6467
func (ba *casBlobAccess) Get(ctx context.Context, digest digest.Digest) buffer.Buffer {
6568
ctxWithCancel, cancel := context.WithCancel(ctx)
66-
client, err := ba.byteStreamClient.Read(ctxWithCancel, &bytestream.ReadRequest{
67-
ResourceName: digest.GetByteStreamReadPath(remoteexecution.Compressor_IDENTITY),
68-
})
69+
resourceName := digest.GetByteStreamReadPath(remoteexecution.Compressor_IDENTITY)
70+
client, err := ba.byteStreamClient.Read(
71+
metadata.AppendToOutgoingContext(ctxWithCancel, resourceNameHeader, resourceName),
72+
&bytestream.ReadRequest{
73+
ResourceName: resourceName,
74+
},
75+
)
6976
if err != nil {
7077
cancel()
7178
return buffer.NewBufferFromError(err)
@@ -86,13 +93,15 @@ func (ba *casBlobAccess) Put(ctx context.Context, digest digest.Digest, b buffer
8693
defer r.Close()
8794

8895
ctxWithCancel, cancel := context.WithCancel(ctx)
89-
client, err := ba.byteStreamClient.Write(ctxWithCancel)
96+
resourceName := digest.GetByteStreamWritePath(uuid.Must(ba.uuidGenerator()), remoteexecution.Compressor_IDENTITY)
97+
client, err := ba.byteStreamClient.Write(
98+
metadata.AppendToOutgoingContext(ctxWithCancel, resourceNameHeader, resourceName),
99+
)
90100
if err != nil {
91101
cancel()
92102
return err
93103
}
94104

95-
resourceName := digest.GetByteStreamWritePath(uuid.Must(ba.uuidGenerator()), remoteexecution.Compressor_IDENTITY)
96105
writeOffset := int64(0)
97106
for {
98107
if data, err := r.Read(); err == nil {

pkg/blobstore/grpcclients/cas_blob_access_test.go

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ func TestCASBlobAccessPut(t *testing.T) {
3636

3737
t.Run("InitialFailure", func(t *testing.T) {
3838
// Failure to create the outgoing connection.
39+
uuidGenerator.EXPECT().Call().Return(uuid, nil)
3940
client.EXPECT().NewStream(gomock.Any(), gomock.Any(), "/google.bytestream.ByteStream/Write").
4041
Return(nil, status.Error(codes.Internal, "Failed to create outgoing connection"))
4142
r := mock.NewMockFileReader(ctrl)
@@ -52,12 +53,12 @@ func TestCASBlobAccessPut(t *testing.T) {
5253
// should be returned.
5354
clientStream := mock.NewMockClientStream(ctrl)
5455
var savedCtx context.Context
56+
uuidGenerator.EXPECT().Call().Return(uuid, nil)
5557
client.EXPECT().NewStream(gomock.Any(), gomock.Any(), "/google.bytestream.ByteStream/Write").
5658
DoAndReturn(func(ctx context.Context, desc *grpc.StreamDesc, method string, opts ...grpc.CallOption) (grpc.ClientStream, error) {
5759
savedCtx = ctx
5860
return clientStream, nil
5961
})
60-
uuidGenerator.EXPECT().Call().Return(uuid, nil)
6162
r := mock.NewMockFileReader(ctrl)
6263
r.EXPECT().ReadAt(gomock.Len(5), int64(0)).Return(0, status.Error(codes.Internal, "Disk on fire"))
6364
clientStream.EXPECT().CloseSend().DoAndReturn(func() error {
@@ -79,9 +80,9 @@ func TestCASBlobAccessPut(t *testing.T) {
7980
// error message that is returned by
8081
// ClientStream.CloseSend().
8182
clientStream := mock.NewMockClientStream(ctrl)
83+
uuidGenerator.EXPECT().Call().Return(uuid, nil)
8284
client.EXPECT().NewStream(gomock.Any(), gomock.Any(), "/google.bytestream.ByteStream/Write").
8385
Return(clientStream, nil)
84-
uuidGenerator.EXPECT().Call().Return(uuid, nil)
8586
r := mock.NewMockFileReader(ctrl)
8687
r.EXPECT().ReadAt(gomock.Len(5), int64(0)).DoAndReturn(func(p []byte, off int64) (int, error) {
8788
copy(p, "Hello")
@@ -104,9 +105,9 @@ func TestCASBlobAccessPut(t *testing.T) {
104105
// Similar to the previous test, ClientStream.SendMsg()
105106
// may fail with io.EOF for the final call.
106107
clientStream := mock.NewMockClientStream(ctrl)
108+
uuidGenerator.EXPECT().Call().Return(uuid, nil)
107109
client.EXPECT().NewStream(gomock.Any(), gomock.Any(), "/google.bytestream.ByteStream/Write").
108110
Return(clientStream, nil)
109-
uuidGenerator.EXPECT().Call().Return(uuid, nil)
110111
r := mock.NewMockFileReader(ctrl)
111112
r.EXPECT().ReadAt(gomock.Len(5), int64(0)).DoAndReturn(func(p []byte, off int64) (int, error) {
112113
copy(p, "Hello")
@@ -135,9 +136,9 @@ func TestCASBlobAccessPut(t *testing.T) {
135136
// ClientStream.CloseSend() still fails. The error must
136137
// still be propagated.
137138
clientStream := mock.NewMockClientStream(ctrl)
139+
uuidGenerator.EXPECT().Call().Return(uuid, nil)
138140
client.EXPECT().NewStream(gomock.Any(), gomock.Any(), "/google.bytestream.ByteStream/Write").
139141
Return(clientStream, nil)
140-
uuidGenerator.EXPECT().Call().Return(uuid, nil)
141142
r := mock.NewMockFileReader(ctrl)
142143
r.EXPECT().ReadAt(gomock.Len(5), int64(0)).DoAndReturn(func(p []byte, off int64) (int, error) {
143144
copy(p, "Hello")

0 commit comments

Comments
 (0)