feat(tvix/store): make blobstore stream chunks
This changes the RPC methods to return/consume a stream of chunks, instead of a very big message containing the whole blob, to keep message sizes in manageable sizes (less than 4MiB). Change-Id: I2a3a50f07b059d8a2f5196860254adff98c8a352 Reviewed-on: https://cl.tvl.fyi/c/depot/+/7651 Reviewed-by: tazjin <tazjin@tvl.su> Tested-by: BuildkiteCI
This commit is contained in:
parent
f879993cc4
commit
1c15154b83
3 changed files with 200 additions and 201 deletions
|
@ -71,100 +71,6 @@ func (x *GetBlobRequest) GetDigest() []byte {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
type GetBlobResponse struct {
|
|
||||||
state protoimpl.MessageState
|
|
||||||
sizeCache protoimpl.SizeCache
|
|
||||||
unknownFields protoimpl.UnknownFields
|
|
||||||
|
|
||||||
Data []byte `protobuf:"bytes,1,opt,name=data,proto3" json:"data,omitempty"`
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *GetBlobResponse) Reset() {
|
|
||||||
*x = GetBlobResponse{}
|
|
||||||
if protoimpl.UnsafeEnabled {
|
|
||||||
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[1]
|
|
||||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
|
||||||
ms.StoreMessageInfo(mi)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *GetBlobResponse) String() string {
|
|
||||||
return protoimpl.X.MessageStringOf(x)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (*GetBlobResponse) ProtoMessage() {}
|
|
||||||
|
|
||||||
func (x *GetBlobResponse) ProtoReflect() protoreflect.Message {
|
|
||||||
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[1]
|
|
||||||
if protoimpl.UnsafeEnabled && x != nil {
|
|
||||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
|
||||||
if ms.LoadMessageInfo() == nil {
|
|
||||||
ms.StoreMessageInfo(mi)
|
|
||||||
}
|
|
||||||
return ms
|
|
||||||
}
|
|
||||||
return mi.MessageOf(x)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Deprecated: Use GetBlobResponse.ProtoReflect.Descriptor instead.
|
|
||||||
func (*GetBlobResponse) Descriptor() ([]byte, []int) {
|
|
||||||
return file_tvix_store_protos_rpc_blobstore_proto_rawDescGZIP(), []int{1}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *GetBlobResponse) GetData() []byte {
|
|
||||||
if x != nil {
|
|
||||||
return x.Data
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
type PutBlobRequest struct {
|
|
||||||
state protoimpl.MessageState
|
|
||||||
sizeCache protoimpl.SizeCache
|
|
||||||
unknownFields protoimpl.UnknownFields
|
|
||||||
|
|
||||||
Data []byte `protobuf:"bytes,1,opt,name=data,proto3" json:"data,omitempty"`
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *PutBlobRequest) Reset() {
|
|
||||||
*x = PutBlobRequest{}
|
|
||||||
if protoimpl.UnsafeEnabled {
|
|
||||||
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[2]
|
|
||||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
|
||||||
ms.StoreMessageInfo(mi)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *PutBlobRequest) String() string {
|
|
||||||
return protoimpl.X.MessageStringOf(x)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (*PutBlobRequest) ProtoMessage() {}
|
|
||||||
|
|
||||||
func (x *PutBlobRequest) ProtoReflect() protoreflect.Message {
|
|
||||||
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[2]
|
|
||||||
if protoimpl.UnsafeEnabled && x != nil {
|
|
||||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
|
||||||
if ms.LoadMessageInfo() == nil {
|
|
||||||
ms.StoreMessageInfo(mi)
|
|
||||||
}
|
|
||||||
return ms
|
|
||||||
}
|
|
||||||
return mi.MessageOf(x)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Deprecated: Use PutBlobRequest.ProtoReflect.Descriptor instead.
|
|
||||||
func (*PutBlobRequest) Descriptor() ([]byte, []int) {
|
|
||||||
return file_tvix_store_protos_rpc_blobstore_proto_rawDescGZIP(), []int{2}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *PutBlobRequest) GetData() []byte {
|
|
||||||
if x != nil {
|
|
||||||
return x.Data
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
type PutBlobResponse struct {
|
type PutBlobResponse struct {
|
||||||
state protoimpl.MessageState
|
state protoimpl.MessageState
|
||||||
sizeCache protoimpl.SizeCache
|
sizeCache protoimpl.SizeCache
|
||||||
|
@ -177,7 +83,7 @@ type PutBlobResponse struct {
|
||||||
func (x *PutBlobResponse) Reset() {
|
func (x *PutBlobResponse) Reset() {
|
||||||
*x = PutBlobResponse{}
|
*x = PutBlobResponse{}
|
||||||
if protoimpl.UnsafeEnabled {
|
if protoimpl.UnsafeEnabled {
|
||||||
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[3]
|
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[1]
|
||||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||||
ms.StoreMessageInfo(mi)
|
ms.StoreMessageInfo(mi)
|
||||||
}
|
}
|
||||||
|
@ -190,7 +96,7 @@ func (x *PutBlobResponse) String() string {
|
||||||
func (*PutBlobResponse) ProtoMessage() {}
|
func (*PutBlobResponse) ProtoMessage() {}
|
||||||
|
|
||||||
func (x *PutBlobResponse) ProtoReflect() protoreflect.Message {
|
func (x *PutBlobResponse) ProtoReflect() protoreflect.Message {
|
||||||
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[3]
|
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[1]
|
||||||
if protoimpl.UnsafeEnabled && x != nil {
|
if protoimpl.UnsafeEnabled && x != nil {
|
||||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||||
if ms.LoadMessageInfo() == nil {
|
if ms.LoadMessageInfo() == nil {
|
||||||
|
@ -203,7 +109,7 @@ func (x *PutBlobResponse) ProtoReflect() protoreflect.Message {
|
||||||
|
|
||||||
// Deprecated: Use PutBlobResponse.ProtoReflect.Descriptor instead.
|
// Deprecated: Use PutBlobResponse.ProtoReflect.Descriptor instead.
|
||||||
func (*PutBlobResponse) Descriptor() ([]byte, []int) {
|
func (*PutBlobResponse) Descriptor() ([]byte, []int) {
|
||||||
return file_tvix_store_protos_rpc_blobstore_proto_rawDescGZIP(), []int{3}
|
return file_tvix_store_protos_rpc_blobstore_proto_rawDescGZIP(), []int{1}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (x *PutBlobResponse) GetDigest() []byte {
|
func (x *PutBlobResponse) GetDigest() []byte {
|
||||||
|
@ -213,6 +119,55 @@ func (x *PutBlobResponse) GetDigest() []byte {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// This represents a part of a chunk.
|
||||||
|
// Blobs are sent in smaller chunks to keep message sizes manageable.
|
||||||
|
type BlobChunk struct {
|
||||||
|
state protoimpl.MessageState
|
||||||
|
sizeCache protoimpl.SizeCache
|
||||||
|
unknownFields protoimpl.UnknownFields
|
||||||
|
|
||||||
|
Data []byte `protobuf:"bytes,1,opt,name=data,proto3" json:"data,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *BlobChunk) Reset() {
|
||||||
|
*x = BlobChunk{}
|
||||||
|
if protoimpl.UnsafeEnabled {
|
||||||
|
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[2]
|
||||||
|
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||||
|
ms.StoreMessageInfo(mi)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *BlobChunk) String() string {
|
||||||
|
return protoimpl.X.MessageStringOf(x)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (*BlobChunk) ProtoMessage() {}
|
||||||
|
|
||||||
|
func (x *BlobChunk) ProtoReflect() protoreflect.Message {
|
||||||
|
mi := &file_tvix_store_protos_rpc_blobstore_proto_msgTypes[2]
|
||||||
|
if protoimpl.UnsafeEnabled && x != nil {
|
||||||
|
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||||
|
if ms.LoadMessageInfo() == nil {
|
||||||
|
ms.StoreMessageInfo(mi)
|
||||||
|
}
|
||||||
|
return ms
|
||||||
|
}
|
||||||
|
return mi.MessageOf(x)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Deprecated: Use BlobChunk.ProtoReflect.Descriptor instead.
|
||||||
|
func (*BlobChunk) Descriptor() ([]byte, []int) {
|
||||||
|
return file_tvix_store_protos_rpc_blobstore_proto_rawDescGZIP(), []int{2}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *BlobChunk) GetData() []byte {
|
||||||
|
if x != nil {
|
||||||
|
return x.Data
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
var File_tvix_store_protos_rpc_blobstore_proto protoreflect.FileDescriptor
|
var File_tvix_store_protos_rpc_blobstore_proto protoreflect.FileDescriptor
|
||||||
|
|
||||||
var file_tvix_store_protos_rpc_blobstore_proto_rawDesc = []byte{
|
var file_tvix_store_protos_rpc_blobstore_proto_rawDesc = []byte{
|
||||||
|
@ -222,27 +177,24 @@ var file_tvix_store_protos_rpc_blobstore_proto_rawDesc = []byte{
|
||||||
0x6f, 0x72, 0x65, 0x2e, 0x76, 0x31, 0x22, 0x28, 0x0a, 0x0e, 0x47, 0x65, 0x74, 0x42, 0x6c, 0x6f,
|
0x6f, 0x72, 0x65, 0x2e, 0x76, 0x31, 0x22, 0x28, 0x0a, 0x0e, 0x47, 0x65, 0x74, 0x42, 0x6c, 0x6f,
|
||||||
0x62, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x12, 0x16, 0x0a, 0x06, 0x64, 0x69, 0x67, 0x65,
|
0x62, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x12, 0x16, 0x0a, 0x06, 0x64, 0x69, 0x67, 0x65,
|
||||||
0x73, 0x74, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x06, 0x64, 0x69, 0x67, 0x65, 0x73, 0x74,
|
0x73, 0x74, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x06, 0x64, 0x69, 0x67, 0x65, 0x73, 0x74,
|
||||||
0x22, 0x25, 0x0a, 0x0f, 0x47, 0x65, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x73, 0x70, 0x6f,
|
0x22, 0x29, 0x0a, 0x0f, 0x50, 0x75, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x73, 0x70, 0x6f,
|
||||||
0x6e, 0x73, 0x65, 0x12, 0x12, 0x0a, 0x04, 0x64, 0x61, 0x74, 0x61, 0x18, 0x01, 0x20, 0x01, 0x28,
|
0x6e, 0x73, 0x65, 0x12, 0x16, 0x0a, 0x06, 0x64, 0x69, 0x67, 0x65, 0x73, 0x74, 0x18, 0x01, 0x20,
|
||||||
0x0c, 0x52, 0x04, 0x64, 0x61, 0x74, 0x61, 0x22, 0x24, 0x0a, 0x0e, 0x50, 0x75, 0x74, 0x42, 0x6c,
|
0x01, 0x28, 0x0c, 0x52, 0x06, 0x64, 0x69, 0x67, 0x65, 0x73, 0x74, 0x22, 0x1f, 0x0a, 0x09, 0x42,
|
||||||
0x6f, 0x62, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x12, 0x12, 0x0a, 0x04, 0x64, 0x61, 0x74,
|
0x6c, 0x6f, 0x62, 0x43, 0x68, 0x75, 0x6e, 0x6b, 0x12, 0x12, 0x0a, 0x04, 0x64, 0x61, 0x74, 0x61,
|
||||||
0x61, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x04, 0x64, 0x61, 0x74, 0x61, 0x22, 0x29, 0x0a,
|
0x18, 0x01, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x04, 0x64, 0x61, 0x74, 0x61, 0x32, 0x92, 0x01, 0x0a,
|
||||||
0x0f, 0x50, 0x75, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65,
|
0x0b, 0x42, 0x6c, 0x6f, 0x62, 0x53, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x12, 0x40, 0x0a, 0x03,
|
||||||
0x12, 0x16, 0x0a, 0x06, 0x64, 0x69, 0x67, 0x65, 0x73, 0x74, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c,
|
0x47, 0x65, 0x74, 0x12, 0x1d, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f, 0x72, 0x65,
|
||||||
0x52, 0x06, 0x64, 0x69, 0x67, 0x65, 0x73, 0x74, 0x32, 0x99, 0x01, 0x0a, 0x0b, 0x42, 0x6c, 0x6f,
|
0x2e, 0x76, 0x31, 0x2e, 0x47, 0x65, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x71, 0x75, 0x65,
|
||||||
0x62, 0x53, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x12, 0x44, 0x0a, 0x03, 0x47, 0x65, 0x74, 0x12,
|
0x73, 0x74, 0x1a, 0x18, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x2e,
|
||||||
0x1d, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x2e, 0x76, 0x31, 0x2e,
|
0x76, 0x31, 0x2e, 0x42, 0x6c, 0x6f, 0x62, 0x43, 0x68, 0x75, 0x6e, 0x6b, 0x30, 0x01, 0x12, 0x41,
|
||||||
0x47, 0x65, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x1e,
|
0x0a, 0x03, 0x50, 0x75, 0x74, 0x12, 0x18, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f,
|
||||||
0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x2e, 0x76, 0x31, 0x2e, 0x47,
|
0x72, 0x65, 0x2e, 0x76, 0x31, 0x2e, 0x42, 0x6c, 0x6f, 0x62, 0x43, 0x68, 0x75, 0x6e, 0x6b, 0x1a,
|
||||||
0x65, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x44,
|
0x1e, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x2e, 0x76, 0x31, 0x2e,
|
||||||
0x0a, 0x03, 0x50, 0x75, 0x74, 0x12, 0x1d, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f,
|
0x50, 0x75, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x28,
|
||||||
0x72, 0x65, 0x2e, 0x76, 0x31, 0x2e, 0x50, 0x75, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x71,
|
0x01, 0x42, 0x28, 0x5a, 0x26, 0x63, 0x6f, 0x64, 0x65, 0x2e, 0x74, 0x76, 0x6c, 0x2e, 0x66, 0x79,
|
||||||
0x75, 0x65, 0x73, 0x74, 0x1a, 0x1e, 0x2e, 0x74, 0x76, 0x69, 0x78, 0x2e, 0x73, 0x74, 0x6f, 0x72,
|
0x69, 0x2f, 0x74, 0x76, 0x69, 0x78, 0x2f, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x2f, 0x70, 0x72, 0x6f,
|
||||||
0x65, 0x2e, 0x76, 0x31, 0x2e, 0x50, 0x75, 0x74, 0x42, 0x6c, 0x6f, 0x62, 0x52, 0x65, 0x73, 0x70,
|
0x74, 0x6f, 0x73, 0x3b, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x76, 0x31, 0x62, 0x06, 0x70, 0x72, 0x6f,
|
||||||
0x6f, 0x6e, 0x73, 0x65, 0x42, 0x28, 0x5a, 0x26, 0x63, 0x6f, 0x64, 0x65, 0x2e, 0x74, 0x76, 0x6c,
|
0x74, 0x6f, 0x33,
|
||||||
0x2e, 0x66, 0x79, 0x69, 0x2f, 0x74, 0x76, 0x69, 0x78, 0x2f, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x2f,
|
|
||||||
0x70, 0x72, 0x6f, 0x74, 0x6f, 0x73, 0x3b, 0x73, 0x74, 0x6f, 0x72, 0x65, 0x76, 0x31, 0x62, 0x06,
|
|
||||||
0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
|
@ -257,18 +209,17 @@ func file_tvix_store_protos_rpc_blobstore_proto_rawDescGZIP() []byte {
|
||||||
return file_tvix_store_protos_rpc_blobstore_proto_rawDescData
|
return file_tvix_store_protos_rpc_blobstore_proto_rawDescData
|
||||||
}
|
}
|
||||||
|
|
||||||
var file_tvix_store_protos_rpc_blobstore_proto_msgTypes = make([]protoimpl.MessageInfo, 4)
|
var file_tvix_store_protos_rpc_blobstore_proto_msgTypes = make([]protoimpl.MessageInfo, 3)
|
||||||
var file_tvix_store_protos_rpc_blobstore_proto_goTypes = []interface{}{
|
var file_tvix_store_protos_rpc_blobstore_proto_goTypes = []interface{}{
|
||||||
(*GetBlobRequest)(nil), // 0: tvix.store.v1.GetBlobRequest
|
(*GetBlobRequest)(nil), // 0: tvix.store.v1.GetBlobRequest
|
||||||
(*GetBlobResponse)(nil), // 1: tvix.store.v1.GetBlobResponse
|
(*PutBlobResponse)(nil), // 1: tvix.store.v1.PutBlobResponse
|
||||||
(*PutBlobRequest)(nil), // 2: tvix.store.v1.PutBlobRequest
|
(*BlobChunk)(nil), // 2: tvix.store.v1.BlobChunk
|
||||||
(*PutBlobResponse)(nil), // 3: tvix.store.v1.PutBlobResponse
|
|
||||||
}
|
}
|
||||||
var file_tvix_store_protos_rpc_blobstore_proto_depIdxs = []int32{
|
var file_tvix_store_protos_rpc_blobstore_proto_depIdxs = []int32{
|
||||||
0, // 0: tvix.store.v1.BlobService.Get:input_type -> tvix.store.v1.GetBlobRequest
|
0, // 0: tvix.store.v1.BlobService.Get:input_type -> tvix.store.v1.GetBlobRequest
|
||||||
2, // 1: tvix.store.v1.BlobService.Put:input_type -> tvix.store.v1.PutBlobRequest
|
2, // 1: tvix.store.v1.BlobService.Put:input_type -> tvix.store.v1.BlobChunk
|
||||||
1, // 2: tvix.store.v1.BlobService.Get:output_type -> tvix.store.v1.GetBlobResponse
|
2, // 2: tvix.store.v1.BlobService.Get:output_type -> tvix.store.v1.BlobChunk
|
||||||
3, // 3: tvix.store.v1.BlobService.Put:output_type -> tvix.store.v1.PutBlobResponse
|
1, // 3: tvix.store.v1.BlobService.Put:output_type -> tvix.store.v1.PutBlobResponse
|
||||||
2, // [2:4] is the sub-list for method output_type
|
2, // [2:4] is the sub-list for method output_type
|
||||||
0, // [0:2] is the sub-list for method input_type
|
0, // [0:2] is the sub-list for method input_type
|
||||||
0, // [0:0] is the sub-list for extension type_name
|
0, // [0:0] is the sub-list for extension type_name
|
||||||
|
@ -295,7 +246,7 @@ func file_tvix_store_protos_rpc_blobstore_proto_init() {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
file_tvix_store_protos_rpc_blobstore_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} {
|
file_tvix_store_protos_rpc_blobstore_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} {
|
||||||
switch v := v.(*GetBlobResponse); i {
|
switch v := v.(*PutBlobResponse); i {
|
||||||
case 0:
|
case 0:
|
||||||
return &v.state
|
return &v.state
|
||||||
case 1:
|
case 1:
|
||||||
|
@ -307,19 +258,7 @@ func file_tvix_store_protos_rpc_blobstore_proto_init() {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
file_tvix_store_protos_rpc_blobstore_proto_msgTypes[2].Exporter = func(v interface{}, i int) interface{} {
|
file_tvix_store_protos_rpc_blobstore_proto_msgTypes[2].Exporter = func(v interface{}, i int) interface{} {
|
||||||
switch v := v.(*PutBlobRequest); i {
|
switch v := v.(*BlobChunk); i {
|
||||||
case 0:
|
|
||||||
return &v.state
|
|
||||||
case 1:
|
|
||||||
return &v.sizeCache
|
|
||||||
case 2:
|
|
||||||
return &v.unknownFields
|
|
||||||
default:
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
}
|
|
||||||
file_tvix_store_protos_rpc_blobstore_proto_msgTypes[3].Exporter = func(v interface{}, i int) interface{} {
|
|
||||||
switch v := v.(*PutBlobResponse); i {
|
|
||||||
case 0:
|
case 0:
|
||||||
return &v.state
|
return &v.state
|
||||||
case 1:
|
case 1:
|
||||||
|
@ -337,7 +276,7 @@ func file_tvix_store_protos_rpc_blobstore_proto_init() {
|
||||||
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
|
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
|
||||||
RawDescriptor: file_tvix_store_protos_rpc_blobstore_proto_rawDesc,
|
RawDescriptor: file_tvix_store_protos_rpc_blobstore_proto_rawDesc,
|
||||||
NumEnums: 0,
|
NumEnums: 0,
|
||||||
NumMessages: 4,
|
NumMessages: 3,
|
||||||
NumExtensions: 0,
|
NumExtensions: 0,
|
||||||
NumServices: 1,
|
NumServices: 1,
|
||||||
},
|
},
|
||||||
|
|
|
@ -7,8 +7,9 @@ package tvix.store.v1;
|
||||||
option go_package = "code.tvl.fyi/tvix/store/protos;storev1";
|
option go_package = "code.tvl.fyi/tvix/store/protos;storev1";
|
||||||
|
|
||||||
service BlobService {
|
service BlobService {
|
||||||
rpc Get(GetBlobRequest) returns (GetBlobResponse);
|
rpc Get(GetBlobRequest) returns (stream BlobChunk);
|
||||||
rpc Put(PutBlobRequest) returns (PutBlobResponse);
|
|
||||||
|
rpc Put(stream BlobChunk) returns (PutBlobResponse);
|
||||||
|
|
||||||
// TODO(flokli): We can get fancy here, and add methods to retrieve
|
// TODO(flokli): We can get fancy here, and add methods to retrieve
|
||||||
// [Bao](https://github.com/oconnor663/bao/blob/master/docs/spec.md), and
|
// [Bao](https://github.com/oconnor663/bao/blob/master/docs/spec.md), and
|
||||||
|
@ -20,15 +21,13 @@ message GetBlobRequest {
|
||||||
bytes digest = 1;
|
bytes digest = 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
message GetBlobResponse {
|
|
||||||
bytes data = 1;
|
|
||||||
}
|
|
||||||
|
|
||||||
message PutBlobRequest {
|
|
||||||
bytes data = 1;
|
|
||||||
}
|
|
||||||
|
|
||||||
message PutBlobResponse {
|
message PutBlobResponse {
|
||||||
// The blake3 digest of the data that was sent.
|
// The blake3 digest of the data that was sent.
|
||||||
bytes digest = 1;
|
bytes digest = 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// This represents a part of a chunk.
|
||||||
|
// Blobs are sent in smaller chunks to keep message sizes manageable.
|
||||||
|
message BlobChunk {
|
||||||
|
bytes data = 1;
|
||||||
|
}
|
||||||
|
|
|
@ -22,8 +22,8 @@ const _ = grpc.SupportPackageIsVersion7
|
||||||
//
|
//
|
||||||
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
|
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
|
||||||
type BlobServiceClient interface {
|
type BlobServiceClient interface {
|
||||||
Get(ctx context.Context, in *GetBlobRequest, opts ...grpc.CallOption) (*GetBlobResponse, error)
|
Get(ctx context.Context, in *GetBlobRequest, opts ...grpc.CallOption) (BlobService_GetClient, error)
|
||||||
Put(ctx context.Context, in *PutBlobRequest, opts ...grpc.CallOption) (*PutBlobResponse, error)
|
Put(ctx context.Context, opts ...grpc.CallOption) (BlobService_PutClient, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
type blobServiceClient struct {
|
type blobServiceClient struct {
|
||||||
|
@ -34,30 +34,78 @@ func NewBlobServiceClient(cc grpc.ClientConnInterface) BlobServiceClient {
|
||||||
return &blobServiceClient{cc}
|
return &blobServiceClient{cc}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *blobServiceClient) Get(ctx context.Context, in *GetBlobRequest, opts ...grpc.CallOption) (*GetBlobResponse, error) {
|
func (c *blobServiceClient) Get(ctx context.Context, in *GetBlobRequest, opts ...grpc.CallOption) (BlobService_GetClient, error) {
|
||||||
out := new(GetBlobResponse)
|
stream, err := c.cc.NewStream(ctx, &BlobService_ServiceDesc.Streams[0], "/tvix.store.v1.BlobService/Get", opts...)
|
||||||
err := c.cc.Invoke(ctx, "/tvix.store.v1.BlobService/Get", in, out, opts...)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return out, nil
|
x := &blobServiceGetClient{stream}
|
||||||
|
if err := x.ClientStream.SendMsg(in); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if err := x.ClientStream.CloseSend(); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return x, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *blobServiceClient) Put(ctx context.Context, in *PutBlobRequest, opts ...grpc.CallOption) (*PutBlobResponse, error) {
|
type BlobService_GetClient interface {
|
||||||
out := new(PutBlobResponse)
|
Recv() (*BlobChunk, error)
|
||||||
err := c.cc.Invoke(ctx, "/tvix.store.v1.BlobService/Put", in, out, opts...)
|
grpc.ClientStream
|
||||||
|
}
|
||||||
|
|
||||||
|
type blobServiceGetClient struct {
|
||||||
|
grpc.ClientStream
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *blobServiceGetClient) Recv() (*BlobChunk, error) {
|
||||||
|
m := new(BlobChunk)
|
||||||
|
if err := x.ClientStream.RecvMsg(m); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return m, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *blobServiceClient) Put(ctx context.Context, opts ...grpc.CallOption) (BlobService_PutClient, error) {
|
||||||
|
stream, err := c.cc.NewStream(ctx, &BlobService_ServiceDesc.Streams[1], "/tvix.store.v1.BlobService/Put", opts...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return out, nil
|
x := &blobServicePutClient{stream}
|
||||||
|
return x, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type BlobService_PutClient interface {
|
||||||
|
Send(*BlobChunk) error
|
||||||
|
CloseAndRecv() (*PutBlobResponse, error)
|
||||||
|
grpc.ClientStream
|
||||||
|
}
|
||||||
|
|
||||||
|
type blobServicePutClient struct {
|
||||||
|
grpc.ClientStream
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *blobServicePutClient) Send(m *BlobChunk) error {
|
||||||
|
return x.ClientStream.SendMsg(m)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *blobServicePutClient) CloseAndRecv() (*PutBlobResponse, error) {
|
||||||
|
if err := x.ClientStream.CloseSend(); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
m := new(PutBlobResponse)
|
||||||
|
if err := x.ClientStream.RecvMsg(m); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return m, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// BlobServiceServer is the server API for BlobService service.
|
// BlobServiceServer is the server API for BlobService service.
|
||||||
// All implementations must embed UnimplementedBlobServiceServer
|
// All implementations must embed UnimplementedBlobServiceServer
|
||||||
// for forward compatibility
|
// for forward compatibility
|
||||||
type BlobServiceServer interface {
|
type BlobServiceServer interface {
|
||||||
Get(context.Context, *GetBlobRequest) (*GetBlobResponse, error)
|
Get(*GetBlobRequest, BlobService_GetServer) error
|
||||||
Put(context.Context, *PutBlobRequest) (*PutBlobResponse, error)
|
Put(BlobService_PutServer) error
|
||||||
mustEmbedUnimplementedBlobServiceServer()
|
mustEmbedUnimplementedBlobServiceServer()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -65,11 +113,11 @@ type BlobServiceServer interface {
|
||||||
type UnimplementedBlobServiceServer struct {
|
type UnimplementedBlobServiceServer struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (UnimplementedBlobServiceServer) Get(context.Context, *GetBlobRequest) (*GetBlobResponse, error) {
|
func (UnimplementedBlobServiceServer) Get(*GetBlobRequest, BlobService_GetServer) error {
|
||||||
return nil, status.Errorf(codes.Unimplemented, "method Get not implemented")
|
return status.Errorf(codes.Unimplemented, "method Get not implemented")
|
||||||
}
|
}
|
||||||
func (UnimplementedBlobServiceServer) Put(context.Context, *PutBlobRequest) (*PutBlobResponse, error) {
|
func (UnimplementedBlobServiceServer) Put(BlobService_PutServer) error {
|
||||||
return nil, status.Errorf(codes.Unimplemented, "method Put not implemented")
|
return status.Errorf(codes.Unimplemented, "method Put not implemented")
|
||||||
}
|
}
|
||||||
func (UnimplementedBlobServiceServer) mustEmbedUnimplementedBlobServiceServer() {}
|
func (UnimplementedBlobServiceServer) mustEmbedUnimplementedBlobServiceServer() {}
|
||||||
|
|
||||||
|
@ -84,40 +132,51 @@ func RegisterBlobServiceServer(s grpc.ServiceRegistrar, srv BlobServiceServer) {
|
||||||
s.RegisterService(&BlobService_ServiceDesc, srv)
|
s.RegisterService(&BlobService_ServiceDesc, srv)
|
||||||
}
|
}
|
||||||
|
|
||||||
func _BlobService_Get_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
|
func _BlobService_Get_Handler(srv interface{}, stream grpc.ServerStream) error {
|
||||||
in := new(GetBlobRequest)
|
m := new(GetBlobRequest)
|
||||||
if err := dec(in); err != nil {
|
if err := stream.RecvMsg(m); err != nil {
|
||||||
return nil, err
|
return err
|
||||||
}
|
}
|
||||||
if interceptor == nil {
|
return srv.(BlobServiceServer).Get(m, &blobServiceGetServer{stream})
|
||||||
return srv.(BlobServiceServer).Get(ctx, in)
|
|
||||||
}
|
|
||||||
info := &grpc.UnaryServerInfo{
|
|
||||||
Server: srv,
|
|
||||||
FullMethod: "/tvix.store.v1.BlobService/Get",
|
|
||||||
}
|
|
||||||
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
|
|
||||||
return srv.(BlobServiceServer).Get(ctx, req.(*GetBlobRequest))
|
|
||||||
}
|
|
||||||
return interceptor(ctx, in, info, handler)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func _BlobService_Put_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
|
type BlobService_GetServer interface {
|
||||||
in := new(PutBlobRequest)
|
Send(*BlobChunk) error
|
||||||
if err := dec(in); err != nil {
|
grpc.ServerStream
|
||||||
|
}
|
||||||
|
|
||||||
|
type blobServiceGetServer struct {
|
||||||
|
grpc.ServerStream
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *blobServiceGetServer) Send(m *BlobChunk) error {
|
||||||
|
return x.ServerStream.SendMsg(m)
|
||||||
|
}
|
||||||
|
|
||||||
|
func _BlobService_Put_Handler(srv interface{}, stream grpc.ServerStream) error {
|
||||||
|
return srv.(BlobServiceServer).Put(&blobServicePutServer{stream})
|
||||||
|
}
|
||||||
|
|
||||||
|
type BlobService_PutServer interface {
|
||||||
|
SendAndClose(*PutBlobResponse) error
|
||||||
|
Recv() (*BlobChunk, error)
|
||||||
|
grpc.ServerStream
|
||||||
|
}
|
||||||
|
|
||||||
|
type blobServicePutServer struct {
|
||||||
|
grpc.ServerStream
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *blobServicePutServer) SendAndClose(m *PutBlobResponse) error {
|
||||||
|
return x.ServerStream.SendMsg(m)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *blobServicePutServer) Recv() (*BlobChunk, error) {
|
||||||
|
m := new(BlobChunk)
|
||||||
|
if err := x.ServerStream.RecvMsg(m); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
if interceptor == nil {
|
return m, nil
|
||||||
return srv.(BlobServiceServer).Put(ctx, in)
|
|
||||||
}
|
|
||||||
info := &grpc.UnaryServerInfo{
|
|
||||||
Server: srv,
|
|
||||||
FullMethod: "/tvix.store.v1.BlobService/Put",
|
|
||||||
}
|
|
||||||
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
|
|
||||||
return srv.(BlobServiceServer).Put(ctx, req.(*PutBlobRequest))
|
|
||||||
}
|
|
||||||
return interceptor(ctx, in, info, handler)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// BlobService_ServiceDesc is the grpc.ServiceDesc for BlobService service.
|
// BlobService_ServiceDesc is the grpc.ServiceDesc for BlobService service.
|
||||||
|
@ -126,16 +185,18 @@ func _BlobService_Put_Handler(srv interface{}, ctx context.Context, dec func(int
|
||||||
var BlobService_ServiceDesc = grpc.ServiceDesc{
|
var BlobService_ServiceDesc = grpc.ServiceDesc{
|
||||||
ServiceName: "tvix.store.v1.BlobService",
|
ServiceName: "tvix.store.v1.BlobService",
|
||||||
HandlerType: (*BlobServiceServer)(nil),
|
HandlerType: (*BlobServiceServer)(nil),
|
||||||
Methods: []grpc.MethodDesc{
|
Methods: []grpc.MethodDesc{},
|
||||||
|
Streams: []grpc.StreamDesc{
|
||||||
{
|
{
|
||||||
MethodName: "Get",
|
StreamName: "Get",
|
||||||
Handler: _BlobService_Get_Handler,
|
Handler: _BlobService_Get_Handler,
|
||||||
|
ServerStreams: true,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
MethodName: "Put",
|
StreamName: "Put",
|
||||||
Handler: _BlobService_Put_Handler,
|
Handler: _BlobService_Put_Handler,
|
||||||
|
ClientStreams: true,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
Streams: []grpc.StreamDesc{},
|
|
||||||
Metadata: "tvix/store/protos/rpc_blobstore.proto",
|
Metadata: "tvix/store/protos/rpc_blobstore.proto",
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in a new issue