diff --git a/docs/GRPC.md b/docs/GRPC.md index ca557ac..7064c9f 100644 --- a/docs/GRPC.md +++ b/docs/GRPC.md @@ -281,13 +281,98 @@ protoc --go_out=. --go-grpc_out=. protoapi.proto --- -## 7. 한 걸음 더 나아가기 (다음 단계) +## 7. 대용량 데이터 전송을 위한 스트리밍(Streaming) 구현 + +일반적인 단발성 요청/응답(Unary) 통신은 전송할 전체 데이터를 단일 메모리에 전부 올려 적재한 상태에서 동작하므로, 펌웨어나 대형 이미지 같은 대용량 데이터를 다룰 때 메모리 고갈(OOM)이나 네트워크 대역폭 병목을 초래하기 쉽습니다. gRPC는 HTTP/2 프로토콜의 스트림(Stream) 채널을 기본 가용하므로, 데이터를 일정 크기(Chunk) 단위로 쪼개 연속적으로 전송할 수 있는 강력한 **스트리밍(Streaming)** 기법을 지원합니다. 본 예제에서는 클라이언트가 가상의 10KB 파일을 조각내어 연속 송출하는 **클라이언트 스트리밍(Client-side Streaming)** 기법을 구현했습니다. + +### 7.1 스키마 설계 (`protoapi.proto`) +데이터를 전송할 때 파일 식별 메타와 조각난 데이터 조각(`bytes`)을 순차적으로 실어 보낼 수 있도록 인터페이스 명세를 아래와 같이 선언해 둡니다: +```proto +message FileChunk { + string FileName = 1; + bytes Content = 2; // 쪼개진 실제 데이터 알맹이 +} + +message UploadStatus { + bool Success = 1; + string Message = 2; + int64 BytesUploaded = 3; // 서버가 조립 완료한 누적 수신량 +} +``` +* **설계 포인트**: 원격 호출 함수의 입력 매개변수 앞에 `stream` 키워드를 정의함으로써 단발성이 아닌 스트림 채널을 활성화하도록 컴파일러에 제약합니다. + +### 7.2 서버 수신 핸들러 구현 ([server.go](../examples/grpcentity/server.go)) +```go +func (IoTServer) UploadFile(stream protoapi.IoTService_UploadFileServer) error { + var totalBytes int64 + var fileName string + + for { + chunk, err := stream.Recv() + if err == io.EOF { + // io.EOF는 클라이언트가 데이터 전송 완료 후 채널을 닫았음을 의미함 + fmt.Printf("File upload completed. Received %d bytes for file '%s'\n", totalBytes, fileName) + return stream.SendAndClose(&protoapi.UploadStatus{ + Success: true, + Message: fmt.Sprintf("File '%s' uploaded successfully.", fileName), + BytesUploaded: totalBytes, + }) + } + if err != nil { + return err + } + + if fileName == "" { + fileName = chunk.GetFileName() + } + totalBytes += int64(len(chunk.GetContent())) + } +} +``` +* **쉬운 설명**: 서버는 통신 채널에 조각 상자가 도달할 때까지 루프를 통해 대기(`stream.Recv()`)합니다. 더 들어올 데이터가 없어 마감 표시(`io.EOF`)가 수신되면 루프를 빠져나와 총 수령 바이트 크기와 정합성 지표를 적재한 성공 영수증을 클라이언트에게 돌려주며 최종 처리 채널을 마감합니다. +* **상세 설명**: `stream.Recv()`는 수신 파이프라인 버퍼에서 이벤트를 논블로킹 대기(Blocking Wait) 형태로 반환받는 동작 방식입니다. 클라이언트가 송출 완료 신호(Half-close)를 보내면 `io.EOF` 로 캡처되며, 이후 수집된 누적 메타 정보를 바탕으로 `SendAndClose()`를 기동해 단일 응답 객체를 전송하고 통신 리소스를 해제합니다. + +### 7.3 클라이언트 송신 기동 ([client.go](../examples/grpcentity/client.go)) +```go +func AskUploadFile(ctx context.Context, m protoapi.IoTServiceClient, fileName string, fileData []byte) (*protoapi.UploadStatus, error) { + stream, err := m.UploadFile(ctx) + if err != nil { + return nil, err + } + + chunkSize := 1024 // 1KB 크기 단위로 조각 설정 + totalBytes := len(fileData) + + for i := 0; i < totalBytes; i += chunkSize { + end := i + chunkSize + if end > totalBytes { + end = totalBytes + } + + err := stream.Send(&protoapi.FileChunk{ + FileName: fileName, + Content: fileData[i:end], + }) + if err != nil { + return nil, err + } + } + + return stream.CloseAndRecv() +} +``` +* **쉬운 설명**: 클라이언트는 전송할 데이터를 준비한 뒤 이를 1KB 단위의 작은 상자들로 조각조각 잘라냅니다. 그 후 연속적인 루프를 돌면서 채널을 타고 상자들을 차례대로 전송(`stream.Send`)합니다. 전송이 모두 끝나면 통신선을 닫고 서버가 수집 집계를 마치고 영수증을 돌려줄 때까지 대기(`CloseAndRecv`)합니다. +* **상세 설명**: gRPC 컴파일러가 도출한 스트림 인스턴스의 `Send()` 메서드를 사용하여 순차적인 패킷 세그먼트를 윈도우 슬라이싱 크기만큼 나누어 송신합니다. 데이터 분할 루프가 마감되면 `CloseAndRecv()`를 트리거하여 스트림의 송신 플래그를 차단(Half-close)하고 서버로부터 결과 영수증이 도달할 때까지 동기식으로 블로킹 대기합니다. + +--- + +## 8. 한 걸음 더 나아가기 (다음 단계) 본 기초 실습을 끝마치셨다면, 아래 과제를 해결해보세요: **약속 스펙 확장해 보기**: [protoapi.proto](../examples/grpcentity/protoapi.proto) 파일에 새로운 환경 데이터(예: 미세먼지 수치 `double Dust = 4;`)를 슬쩍 얹어본 뒤, 직접 번역기를 새로 돌리고 Go 소스코드를 고치며 확장해 봅니다. --- -## 8. 참고 자료 +## 9. 참고 자료 * [gRPC와 REST의 차이점 (AWS)](https://aws.amazon.com/ko/compare/the-difference-between-grpc-and-rest/): 두 방식의 특징과 언제 어떤 기술을 선택해야 하는지 친절하게 정리된 공식 블로그 자료입니다. diff --git a/examples/grpcentity/client.go b/examples/grpcentity/client.go index e315556..42b772f 100644 --- a/examples/grpcentity/client.go +++ b/examples/grpcentity/client.go @@ -37,6 +37,33 @@ func AskUpdateSensingData(ctx context.Context, m protoapi.IoTServiceClient, devi return m.UpdateSensingData(ctx, request) } +func AskUploadFile(ctx context.Context, m protoapi.IoTServiceClient, fileName string, fileData []byte) (*protoapi.UploadStatus, error) { + stream, err := m.UploadFile(ctx) + if err != nil { + return nil, err + } + + chunkSize := 1024 // 1KB 단위 청크 + totalBytes := len(fileData) + + for i := 0; i < totalBytes; i += chunkSize { + end := i + chunkSize + if end > totalBytes { + end = totalBytes + } + + err := stream.Send(&protoapi.FileChunk{ + FileName: fileName, + Content: fileData[i:end], + }) + if err != nil { + return nil, err + } + } + + return stream.CloseAndRecv() +} + func ClientRun(addr string) { conn, err := grpc.NewClient(addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { @@ -67,4 +94,19 @@ func ClientRun(addr string) { } fmt.Println("Sensing Update Success:", res.Success) fmt.Println("Sensing Update Message:", res.Message) + + // 4단계: 파일 업로드 스트리밍 실행 예제 + dummyData := make([]byte, 10240) // 10KB 가상 더미 데이터 + for i := range dummyData { + dummyData[i] = byte(rand.Intn(256)) + } + fmt.Println("Uploading dummy file (10KB) via Client Streaming...") + status, err := AskUploadFile(context.Background(), client, "firmware.bin", dummyData) + if err != nil { + fmt.Println("File upload failed:", err) + return + } + fmt.Println("Upload Success:", status.Success) + fmt.Println("Upload Message:", status.Message) + fmt.Printf("Uploaded Bytes: %d bytes\n", status.BytesUploaded) } diff --git a/examples/grpcentity/protoapi.proto b/examples/grpcentity/protoapi.proto index 8410e7c..bff4a60 100644 --- a/examples/grpcentity/protoapi.proto +++ b/examples/grpcentity/protoapi.proto @@ -6,6 +6,18 @@ service IoTService { rpc GetDate (RequestDateTime) returns (DateTime); rpc UpdateSensingData (SensingData) returns (SensingResponse); rpc GetRandomPass (RequestPass) returns (RandomPass); + rpc UploadFile (stream FileChunk) returns (UploadStatus); +} + +message FileChunk { + string FileName = 1; + bytes Content = 2; +} + +message UploadStatus { + bool Success = 1; + string Message = 2; + int64 BytesUploaded = 3; } message SensingData { diff --git a/examples/grpcentity/protoapi/protoapi.pb.go b/examples/grpcentity/protoapi/protoapi.pb.go index a01ae6d..e0d8809 100644 --- a/examples/grpcentity/protoapi/protoapi.pb.go +++ b/examples/grpcentity/protoapi/protoapi.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go. DO NOT EDIT. // versions: // protoc-gen-go v1.36.11 -// protoc v3.21.12 +// protoc v5.27.2 // source: protoapi.proto package protoapi @@ -21,6 +21,118 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) +type FileChunk struct { + state protoimpl.MessageState `protogen:"open.v1"` + FileName string `protobuf:"bytes,1,opt,name=FileName,proto3" json:"FileName,omitempty"` + Content []byte `protobuf:"bytes,2,opt,name=Content,proto3" json:"Content,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *FileChunk) Reset() { + *x = FileChunk{} + mi := &file_protoapi_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *FileChunk) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*FileChunk) ProtoMessage() {} + +func (x *FileChunk) ProtoReflect() protoreflect.Message { + mi := &file_protoapi_proto_msgTypes[0] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use FileChunk.ProtoReflect.Descriptor instead. +func (*FileChunk) Descriptor() ([]byte, []int) { + return file_protoapi_proto_rawDescGZIP(), []int{0} +} + +func (x *FileChunk) GetFileName() string { + if x != nil { + return x.FileName + } + return "" +} + +func (x *FileChunk) GetContent() []byte { + if x != nil { + return x.Content + } + return nil +} + +type UploadStatus struct { + state protoimpl.MessageState `protogen:"open.v1"` + Success bool `protobuf:"varint,1,opt,name=Success,proto3" json:"Success,omitempty"` + Message string `protobuf:"bytes,2,opt,name=Message,proto3" json:"Message,omitempty"` + BytesUploaded int64 `protobuf:"varint,3,opt,name=BytesUploaded,proto3" json:"BytesUploaded,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *UploadStatus) Reset() { + *x = UploadStatus{} + mi := &file_protoapi_proto_msgTypes[1] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *UploadStatus) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*UploadStatus) ProtoMessage() {} + +func (x *UploadStatus) ProtoReflect() protoreflect.Message { + mi := &file_protoapi_proto_msgTypes[1] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use UploadStatus.ProtoReflect.Descriptor instead. +func (*UploadStatus) Descriptor() ([]byte, []int) { + return file_protoapi_proto_rawDescGZIP(), []int{1} +} + +func (x *UploadStatus) GetSuccess() bool { + if x != nil { + return x.Success + } + return false +} + +func (x *UploadStatus) GetMessage() string { + if x != nil { + return x.Message + } + return "" +} + +func (x *UploadStatus) GetBytesUploaded() int64 { + if x != nil { + return x.BytesUploaded + } + return 0 +} + type SensingData struct { state protoimpl.MessageState `protogen:"open.v1"` DeviceId string `protobuf:"bytes,1,opt,name=DeviceId,proto3" json:"DeviceId,omitempty"` @@ -32,7 +144,7 @@ type SensingData struct { func (x *SensingData) Reset() { *x = SensingData{} - mi := &file_protoapi_proto_msgTypes[0] + mi := &file_protoapi_proto_msgTypes[2] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -44,7 +156,7 @@ func (x *SensingData) String() string { func (*SensingData) ProtoMessage() {} func (x *SensingData) ProtoReflect() protoreflect.Message { - mi := &file_protoapi_proto_msgTypes[0] + mi := &file_protoapi_proto_msgTypes[2] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -57,7 +169,7 @@ func (x *SensingData) ProtoReflect() protoreflect.Message { // Deprecated: Use SensingData.ProtoReflect.Descriptor instead. func (*SensingData) Descriptor() ([]byte, []int) { - return file_protoapi_proto_rawDescGZIP(), []int{0} + return file_protoapi_proto_rawDescGZIP(), []int{2} } func (x *SensingData) GetDeviceId() string { @@ -91,7 +203,7 @@ type SensingResponse struct { func (x *SensingResponse) Reset() { *x = SensingResponse{} - mi := &file_protoapi_proto_msgTypes[1] + mi := &file_protoapi_proto_msgTypes[3] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -103,7 +215,7 @@ func (x *SensingResponse) String() string { func (*SensingResponse) ProtoMessage() {} func (x *SensingResponse) ProtoReflect() protoreflect.Message { - mi := &file_protoapi_proto_msgTypes[1] + mi := &file_protoapi_proto_msgTypes[3] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -116,7 +228,7 @@ func (x *SensingResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use SensingResponse.ProtoReflect.Descriptor instead. func (*SensingResponse) Descriptor() ([]byte, []int) { - return file_protoapi_proto_rawDescGZIP(), []int{1} + return file_protoapi_proto_rawDescGZIP(), []int{3} } func (x *SensingResponse) GetSuccess() bool { @@ -142,7 +254,7 @@ type DateTime struct { func (x *DateTime) Reset() { *x = DateTime{} - mi := &file_protoapi_proto_msgTypes[2] + mi := &file_protoapi_proto_msgTypes[4] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -154,7 +266,7 @@ func (x *DateTime) String() string { func (*DateTime) ProtoMessage() {} func (x *DateTime) ProtoReflect() protoreflect.Message { - mi := &file_protoapi_proto_msgTypes[2] + mi := &file_protoapi_proto_msgTypes[4] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -167,7 +279,7 @@ func (x *DateTime) ProtoReflect() protoreflect.Message { // Deprecated: Use DateTime.ProtoReflect.Descriptor instead. func (*DateTime) Descriptor() ([]byte, []int) { - return file_protoapi_proto_rawDescGZIP(), []int{2} + return file_protoapi_proto_rawDescGZIP(), []int{4} } func (x *DateTime) GetValue() string { @@ -186,7 +298,7 @@ type RequestDateTime struct { func (x *RequestDateTime) Reset() { *x = RequestDateTime{} - mi := &file_protoapi_proto_msgTypes[3] + mi := &file_protoapi_proto_msgTypes[5] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -198,7 +310,7 @@ func (x *RequestDateTime) String() string { func (*RequestDateTime) ProtoMessage() {} func (x *RequestDateTime) ProtoReflect() protoreflect.Message { - mi := &file_protoapi_proto_msgTypes[3] + mi := &file_protoapi_proto_msgTypes[5] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -211,7 +323,7 @@ func (x *RequestDateTime) ProtoReflect() protoreflect.Message { // Deprecated: Use RequestDateTime.ProtoReflect.Descriptor instead. func (*RequestDateTime) Descriptor() ([]byte, []int) { - return file_protoapi_proto_rawDescGZIP(), []int{3} + return file_protoapi_proto_rawDescGZIP(), []int{5} } func (x *RequestDateTime) GetValue() string { @@ -231,7 +343,7 @@ type RequestPass struct { func (x *RequestPass) Reset() { *x = RequestPass{} - mi := &file_protoapi_proto_msgTypes[4] + mi := &file_protoapi_proto_msgTypes[6] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -243,7 +355,7 @@ func (x *RequestPass) String() string { func (*RequestPass) ProtoMessage() {} func (x *RequestPass) ProtoReflect() protoreflect.Message { - mi := &file_protoapi_proto_msgTypes[4] + mi := &file_protoapi_proto_msgTypes[6] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -256,7 +368,7 @@ func (x *RequestPass) ProtoReflect() protoreflect.Message { // Deprecated: Use RequestPass.ProtoReflect.Descriptor instead. func (*RequestPass) Descriptor() ([]byte, []int) { - return file_protoapi_proto_rawDescGZIP(), []int{4} + return file_protoapi_proto_rawDescGZIP(), []int{6} } func (x *RequestPass) GetSeed() int64 { @@ -282,7 +394,7 @@ type RandomPass struct { func (x *RandomPass) Reset() { *x = RandomPass{} - mi := &file_protoapi_proto_msgTypes[5] + mi := &file_protoapi_proto_msgTypes[7] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -294,7 +406,7 @@ func (x *RandomPass) String() string { func (*RandomPass) ProtoMessage() {} func (x *RandomPass) ProtoReflect() protoreflect.Message { - mi := &file_protoapi_proto_msgTypes[5] + mi := &file_protoapi_proto_msgTypes[7] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -307,7 +419,7 @@ func (x *RandomPass) ProtoReflect() protoreflect.Message { // Deprecated: Use RandomPass.ProtoReflect.Descriptor instead. func (*RandomPass) Descriptor() ([]byte, []int) { - return file_protoapi_proto_rawDescGZIP(), []int{5} + return file_protoapi_proto_rawDescGZIP(), []int{7} } func (x *RandomPass) GetPassword() string { @@ -321,7 +433,14 @@ var File_protoapi_proto protoreflect.FileDescriptor const file_protoapi_proto_rawDesc = "" + "\n" + - "\x0eprotoapi.proto\"g\n" + + "\x0eprotoapi.proto\"A\n" + + "\tFileChunk\x12\x1a\n" + + "\bFileName\x18\x01 \x01(\tR\bFileName\x12\x18\n" + + "\aContent\x18\x02 \x01(\fR\aContent\"h\n" + + "\fUploadStatus\x12\x18\n" + + "\aSuccess\x18\x01 \x01(\bR\aSuccess\x12\x18\n" + + "\aMessage\x18\x02 \x01(\tR\aMessage\x12$\n" + + "\rBytesUploaded\x18\x03 \x01(\x03R\rBytesUploaded\"g\n" + "\vSensingData\x12\x1a\n" + "\bDeviceId\x18\x01 \x01(\tR\bDeviceId\x12 \n" + "\vTemperature\x18\x02 \x01(\x01R\vTemperature\x12\x1a\n" + @@ -338,12 +457,15 @@ const file_protoapi_proto_rawDesc = "" + "\x06Length\x18\b \x01(\x03R\x06Length\"(\n" + "\n" + "RandomPass\x12\x1a\n" + - "\bPassword\x18\x01 \x01(\tR\bPassword2\x95\x01\n" + + "\bPassword\x18\x01 \x01(\tR\bPassword2\xc0\x01\n" + "\n" + "IoTService\x12&\n" + "\aGetDate\x12\x10.RequestDateTime\x1a\t.DateTime\x123\n" + "\x11UpdateSensingData\x12\f.SensingData\x1a\x10.SensingResponse\x12*\n" + - "\rGetRandomPass\x12\f.RequestPass\x1a\v.RandomPassB\x16Z\x14./protoapi/;protoapib\x06proto3" + "\rGetRandomPass\x12\f.RequestPass\x1a\v.RandomPass\x12)\n" + + "\n" + + "UploadFile\x12\n" + + ".FileChunk\x1a\r.UploadStatus(\x01B\x16Z\x14./protoapi/;protoapib\x06proto3" var ( file_protoapi_proto_rawDescOnce sync.Once @@ -357,24 +479,28 @@ func file_protoapi_proto_rawDescGZIP() []byte { return file_protoapi_proto_rawDescData } -var file_protoapi_proto_msgTypes = make([]protoimpl.MessageInfo, 6) +var file_protoapi_proto_msgTypes = make([]protoimpl.MessageInfo, 8) var file_protoapi_proto_goTypes = []any{ - (*SensingData)(nil), // 0: SensingData - (*SensingResponse)(nil), // 1: SensingResponse - (*DateTime)(nil), // 2: DateTime - (*RequestDateTime)(nil), // 3: RequestDateTime - (*RequestPass)(nil), // 4: RequestPass - (*RandomPass)(nil), // 5: RandomPass + (*FileChunk)(nil), // 0: FileChunk + (*UploadStatus)(nil), // 1: UploadStatus + (*SensingData)(nil), // 2: SensingData + (*SensingResponse)(nil), // 3: SensingResponse + (*DateTime)(nil), // 4: DateTime + (*RequestDateTime)(nil), // 5: RequestDateTime + (*RequestPass)(nil), // 6: RequestPass + (*RandomPass)(nil), // 7: RandomPass } var file_protoapi_proto_depIdxs = []int32{ - 3, // 0: IoTService.GetDate:input_type -> RequestDateTime - 0, // 1: IoTService.UpdateSensingData:input_type -> SensingData - 4, // 2: IoTService.GetRandomPass:input_type -> RequestPass - 2, // 3: IoTService.GetDate:output_type -> DateTime - 1, // 4: IoTService.UpdateSensingData:output_type -> SensingResponse - 5, // 5: IoTService.GetRandomPass:output_type -> RandomPass - 3, // [3:6] is the sub-list for method output_type - 0, // [0:3] is the sub-list for method input_type + 5, // 0: IoTService.GetDate:input_type -> RequestDateTime + 2, // 1: IoTService.UpdateSensingData:input_type -> SensingData + 6, // 2: IoTService.GetRandomPass:input_type -> RequestPass + 0, // 3: IoTService.UploadFile:input_type -> FileChunk + 4, // 4: IoTService.GetDate:output_type -> DateTime + 3, // 5: IoTService.UpdateSensingData:output_type -> SensingResponse + 7, // 6: IoTService.GetRandomPass:output_type -> RandomPass + 1, // 7: IoTService.UploadFile:output_type -> UploadStatus + 4, // [4:8] is the sub-list for method output_type + 0, // [0:4] 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 extendee 0, // [0:0] is the sub-list for field type_name @@ -391,7 +517,7 @@ func file_protoapi_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_protoapi_proto_rawDesc), len(file_protoapi_proto_rawDesc)), NumEnums: 0, - NumMessages: 6, + NumMessages: 8, NumExtensions: 0, NumServices: 1, }, diff --git a/examples/grpcentity/protoapi/protoapi_grpc.pb.go b/examples/grpcentity/protoapi/protoapi_grpc.pb.go index 8db0d84..058d39a 100644 --- a/examples/grpcentity/protoapi/protoapi_grpc.pb.go +++ b/examples/grpcentity/protoapi/protoapi_grpc.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go-grpc. DO NOT EDIT. // versions: // - protoc-gen-go-grpc v1.6.2 -// - protoc v3.21.12 +// - protoc v5.27.2 // source: protoapi.proto package protoapi @@ -22,6 +22,7 @@ const ( IoTService_GetDate_FullMethodName = "/IoTService/GetDate" IoTService_UpdateSensingData_FullMethodName = "/IoTService/UpdateSensingData" IoTService_GetRandomPass_FullMethodName = "/IoTService/GetRandomPass" + IoTService_UploadFile_FullMethodName = "/IoTService/UploadFile" ) // IoTServiceClient is the client API for IoTService service. @@ -31,6 +32,7 @@ type IoTServiceClient interface { GetDate(ctx context.Context, in *RequestDateTime, opts ...grpc.CallOption) (*DateTime, error) UpdateSensingData(ctx context.Context, in *SensingData, opts ...grpc.CallOption) (*SensingResponse, error) GetRandomPass(ctx context.Context, in *RequestPass, opts ...grpc.CallOption) (*RandomPass, error) + UploadFile(ctx context.Context, opts ...grpc.CallOption) (grpc.ClientStreamingClient[FileChunk, UploadStatus], error) } type ioTServiceClient struct { @@ -71,6 +73,19 @@ func (c *ioTServiceClient) GetRandomPass(ctx context.Context, in *RequestPass, o return out, nil } +func (c *ioTServiceClient) UploadFile(ctx context.Context, opts ...grpc.CallOption) (grpc.ClientStreamingClient[FileChunk, UploadStatus], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &IoTService_ServiceDesc.Streams[0], IoTService_UploadFile_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[FileChunk, UploadStatus]{ClientStream: stream} + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type IoTService_UploadFileClient = grpc.ClientStreamingClient[FileChunk, UploadStatus] + // IoTServiceServer is the server API for IoTService service. // All implementations must embed UnimplementedIoTServiceServer // for forward compatibility. @@ -78,6 +93,7 @@ type IoTServiceServer interface { GetDate(context.Context, *RequestDateTime) (*DateTime, error) UpdateSensingData(context.Context, *SensingData) (*SensingResponse, error) GetRandomPass(context.Context, *RequestPass) (*RandomPass, error) + UploadFile(grpc.ClientStreamingServer[FileChunk, UploadStatus]) error mustEmbedUnimplementedIoTServiceServer() } @@ -97,6 +113,9 @@ func (UnimplementedIoTServiceServer) UpdateSensingData(context.Context, *Sensing func (UnimplementedIoTServiceServer) GetRandomPass(context.Context, *RequestPass) (*RandomPass, error) { return nil, status.Error(codes.Unimplemented, "method GetRandomPass not implemented") } +func (UnimplementedIoTServiceServer) UploadFile(grpc.ClientStreamingServer[FileChunk, UploadStatus]) error { + return status.Error(codes.Unimplemented, "method UploadFile not implemented") +} func (UnimplementedIoTServiceServer) mustEmbedUnimplementedIoTServiceServer() {} func (UnimplementedIoTServiceServer) testEmbeddedByValue() {} @@ -172,6 +191,13 @@ func _IoTService_GetRandomPass_Handler(srv interface{}, ctx context.Context, dec return interceptor(ctx, in, info, handler) } +func _IoTService_UploadFile_Handler(srv interface{}, stream grpc.ServerStream) error { + return srv.(IoTServiceServer).UploadFile(&grpc.GenericServerStream[FileChunk, UploadStatus]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type IoTService_UploadFileServer = grpc.ClientStreamingServer[FileChunk, UploadStatus] + // IoTService_ServiceDesc is the grpc.ServiceDesc for IoTService service. // It's only intended for direct use with grpc.RegisterService, // and not to be introspected or modified (even as a copy) @@ -192,6 +218,12 @@ var IoTService_ServiceDesc = grpc.ServiceDesc{ Handler: _IoTService_GetRandomPass_Handler, }, }, - Streams: []grpc.StreamDesc{}, + Streams: []grpc.StreamDesc{ + { + StreamName: "UploadFile", + Handler: _IoTService_UploadFile_Handler, + ClientStreams: true, + }, + }, Metadata: "protoapi.proto", } diff --git a/examples/grpcentity/server.go b/examples/grpcentity/server.go index dbf0f2d..612836c 100644 --- a/examples/grpcentity/server.go +++ b/examples/grpcentity/server.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "grpccanary/examples/grpcentity/protoapi" + "io" "math/rand" "net" "time" @@ -85,6 +86,32 @@ func (IoTServer) GetRandomPass(ctx context.Context, r *protoapi.RequestPass) (*p return response, nil } +func (IoTServer) UploadFile(stream protoapi.IoTService_UploadFileServer) error { + var totalBytes int64 + var fileName string + + for { + chunk, err := stream.Recv() + if err == io.EOF { + fmt.Printf("File upload completed. Received %d bytes for file '%s'\n", totalBytes, fileName) + return stream.SendAndClose(&protoapi.UploadStatus{ + Success: true, + Message: fmt.Sprintf("File '%s' uploaded successfully.", fileName), + BytesUploaded: totalBytes, + }) + } + if err != nil { + fmt.Println("File upload error:", err) + return err + } + + if fileName == "" { + fileName = chunk.GetFileName() + } + totalBytes += int64(len(chunk.GetContent())) + } +} + func ServerRun(addr string) { server := grpc.NewServer() var iotServer IoTServer @@ -92,7 +119,7 @@ func ServerRun(addr string) { reflection.Register(server) - listen, err := net.Listen("tcp", port) + listen, err := net.Listen("tcp", addr) if err != nil { fmt.Println(err) return