You can not select more than 25 topics
			Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
		
		
		
		
		
			
		
			
				
					
					
						
							228 lines
						
					
					
						
							11 KiB
						
					
					
				
			
		
		
		
			
			
			
		
		
	
	
							228 lines
						
					
					
						
							11 KiB
						
					
					
				
								// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
							 | 
						|
								// versions:
							 | 
						|
								// - protoc-gen-go-grpc v1.5.1
							 | 
						|
								// - protoc             v5.29.3
							 | 
						|
								// source: mq_agent.proto
							 | 
						|
								
							 | 
						|
								package mq_agent_pb
							 | 
						|
								
							 | 
						|
								import (
							 | 
						|
									context "context"
							 | 
						|
									grpc "google.golang.org/grpc"
							 | 
						|
									codes "google.golang.org/grpc/codes"
							 | 
						|
									status "google.golang.org/grpc/status"
							 | 
						|
								)
							 | 
						|
								
							 | 
						|
								// This is a compile-time assertion to ensure that this generated file
							 | 
						|
								// is compatible with the grpc package it is being compiled against.
							 | 
						|
								// Requires gRPC-Go v1.64.0 or later.
							 | 
						|
								const _ = grpc.SupportPackageIsVersion9
							 | 
						|
								
							 | 
						|
								const (
							 | 
						|
									SeaweedMessagingAgent_StartPublishSession_FullMethodName = "/messaging_pb.SeaweedMessagingAgent/StartPublishSession"
							 | 
						|
									SeaweedMessagingAgent_ClosePublishSession_FullMethodName = "/messaging_pb.SeaweedMessagingAgent/ClosePublishSession"
							 | 
						|
									SeaweedMessagingAgent_PublishRecord_FullMethodName       = "/messaging_pb.SeaweedMessagingAgent/PublishRecord"
							 | 
						|
									SeaweedMessagingAgent_SubscribeRecord_FullMethodName     = "/messaging_pb.SeaweedMessagingAgent/SubscribeRecord"
							 | 
						|
								)
							 | 
						|
								
							 | 
						|
								// SeaweedMessagingAgentClient is the client API for SeaweedMessagingAgent service.
							 | 
						|
								//
							 | 
						|
								// 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 SeaweedMessagingAgentClient interface {
							 | 
						|
									// Publishing
							 | 
						|
									StartPublishSession(ctx context.Context, in *StartPublishSessionRequest, opts ...grpc.CallOption) (*StartPublishSessionResponse, error)
							 | 
						|
									ClosePublishSession(ctx context.Context, in *ClosePublishSessionRequest, opts ...grpc.CallOption) (*ClosePublishSessionResponse, error)
							 | 
						|
									PublishRecord(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[PublishRecordRequest, PublishRecordResponse], error)
							 | 
						|
									// Subscribing
							 | 
						|
									SubscribeRecord(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[SubscribeRecordRequest, SubscribeRecordResponse], error)
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								type seaweedMessagingAgentClient struct {
							 | 
						|
									cc grpc.ClientConnInterface
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func NewSeaweedMessagingAgentClient(cc grpc.ClientConnInterface) SeaweedMessagingAgentClient {
							 | 
						|
									return &seaweedMessagingAgentClient{cc}
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func (c *seaweedMessagingAgentClient) StartPublishSession(ctx context.Context, in *StartPublishSessionRequest, opts ...grpc.CallOption) (*StartPublishSessionResponse, error) {
							 | 
						|
									cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
							 | 
						|
									out := new(StartPublishSessionResponse)
							 | 
						|
									err := c.cc.Invoke(ctx, SeaweedMessagingAgent_StartPublishSession_FullMethodName, in, out, cOpts...)
							 | 
						|
									if err != nil {
							 | 
						|
										return nil, err
							 | 
						|
									}
							 | 
						|
									return out, nil
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func (c *seaweedMessagingAgentClient) ClosePublishSession(ctx context.Context, in *ClosePublishSessionRequest, opts ...grpc.CallOption) (*ClosePublishSessionResponse, error) {
							 | 
						|
									cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
							 | 
						|
									out := new(ClosePublishSessionResponse)
							 | 
						|
									err := c.cc.Invoke(ctx, SeaweedMessagingAgent_ClosePublishSession_FullMethodName, in, out, cOpts...)
							 | 
						|
									if err != nil {
							 | 
						|
										return nil, err
							 | 
						|
									}
							 | 
						|
									return out, nil
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func (c *seaweedMessagingAgentClient) PublishRecord(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[PublishRecordRequest, PublishRecordResponse], error) {
							 | 
						|
									cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
							 | 
						|
									stream, err := c.cc.NewStream(ctx, &SeaweedMessagingAgent_ServiceDesc.Streams[0], SeaweedMessagingAgent_PublishRecord_FullMethodName, cOpts...)
							 | 
						|
									if err != nil {
							 | 
						|
										return nil, err
							 | 
						|
									}
							 | 
						|
									x := &grpc.GenericClientStream[PublishRecordRequest, PublishRecordResponse]{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 SeaweedMessagingAgent_PublishRecordClient = grpc.BidiStreamingClient[PublishRecordRequest, PublishRecordResponse]
							 | 
						|
								
							 | 
						|
								func (c *seaweedMessagingAgentClient) SubscribeRecord(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[SubscribeRecordRequest, SubscribeRecordResponse], error) {
							 | 
						|
									cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
							 | 
						|
									stream, err := c.cc.NewStream(ctx, &SeaweedMessagingAgent_ServiceDesc.Streams[1], SeaweedMessagingAgent_SubscribeRecord_FullMethodName, cOpts...)
							 | 
						|
									if err != nil {
							 | 
						|
										return nil, err
							 | 
						|
									}
							 | 
						|
									x := &grpc.GenericClientStream[SubscribeRecordRequest, SubscribeRecordResponse]{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 SeaweedMessagingAgent_SubscribeRecordClient = grpc.BidiStreamingClient[SubscribeRecordRequest, SubscribeRecordResponse]
							 | 
						|
								
							 | 
						|
								// SeaweedMessagingAgentServer is the server API for SeaweedMessagingAgent service.
							 | 
						|
								// All implementations must embed UnimplementedSeaweedMessagingAgentServer
							 | 
						|
								// for forward compatibility.
							 | 
						|
								type SeaweedMessagingAgentServer interface {
							 | 
						|
									// Publishing
							 | 
						|
									StartPublishSession(context.Context, *StartPublishSessionRequest) (*StartPublishSessionResponse, error)
							 | 
						|
									ClosePublishSession(context.Context, *ClosePublishSessionRequest) (*ClosePublishSessionResponse, error)
							 | 
						|
									PublishRecord(grpc.BidiStreamingServer[PublishRecordRequest, PublishRecordResponse]) error
							 | 
						|
									// Subscribing
							 | 
						|
									SubscribeRecord(grpc.BidiStreamingServer[SubscribeRecordRequest, SubscribeRecordResponse]) error
							 | 
						|
									mustEmbedUnimplementedSeaweedMessagingAgentServer()
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								// UnimplementedSeaweedMessagingAgentServer must be embedded to have
							 | 
						|
								// forward compatible implementations.
							 | 
						|
								//
							 | 
						|
								// NOTE: this should be embedded by value instead of pointer to avoid a nil
							 | 
						|
								// pointer dereference when methods are called.
							 | 
						|
								type UnimplementedSeaweedMessagingAgentServer struct{}
							 | 
						|
								
							 | 
						|
								func (UnimplementedSeaweedMessagingAgentServer) StartPublishSession(context.Context, *StartPublishSessionRequest) (*StartPublishSessionResponse, error) {
							 | 
						|
									return nil, status.Errorf(codes.Unimplemented, "method StartPublishSession not implemented")
							 | 
						|
								}
							 | 
						|
								func (UnimplementedSeaweedMessagingAgentServer) ClosePublishSession(context.Context, *ClosePublishSessionRequest) (*ClosePublishSessionResponse, error) {
							 | 
						|
									return nil, status.Errorf(codes.Unimplemented, "method ClosePublishSession not implemented")
							 | 
						|
								}
							 | 
						|
								func (UnimplementedSeaweedMessagingAgentServer) PublishRecord(grpc.BidiStreamingServer[PublishRecordRequest, PublishRecordResponse]) error {
							 | 
						|
									return status.Errorf(codes.Unimplemented, "method PublishRecord not implemented")
							 | 
						|
								}
							 | 
						|
								func (UnimplementedSeaweedMessagingAgentServer) SubscribeRecord(grpc.BidiStreamingServer[SubscribeRecordRequest, SubscribeRecordResponse]) error {
							 | 
						|
									return status.Errorf(codes.Unimplemented, "method SubscribeRecord not implemented")
							 | 
						|
								}
							 | 
						|
								func (UnimplementedSeaweedMessagingAgentServer) mustEmbedUnimplementedSeaweedMessagingAgentServer() {}
							 | 
						|
								func (UnimplementedSeaweedMessagingAgentServer) testEmbeddedByValue()                               {}
							 | 
						|
								
							 | 
						|
								// UnsafeSeaweedMessagingAgentServer may be embedded to opt out of forward compatibility for this service.
							 | 
						|
								// Use of this interface is not recommended, as added methods to SeaweedMessagingAgentServer will
							 | 
						|
								// result in compilation errors.
							 | 
						|
								type UnsafeSeaweedMessagingAgentServer interface {
							 | 
						|
									mustEmbedUnimplementedSeaweedMessagingAgentServer()
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func RegisterSeaweedMessagingAgentServer(s grpc.ServiceRegistrar, srv SeaweedMessagingAgentServer) {
							 | 
						|
									// If the following call pancis, it indicates UnimplementedSeaweedMessagingAgentServer was
							 | 
						|
									// embedded by pointer and is nil.  This will cause panics if an
							 | 
						|
									// unimplemented method is ever invoked, so we test this at initialization
							 | 
						|
									// time to prevent it from happening at runtime later due to I/O.
							 | 
						|
									if t, ok := srv.(interface{ testEmbeddedByValue() }); ok {
							 | 
						|
										t.testEmbeddedByValue()
							 | 
						|
									}
							 | 
						|
									s.RegisterService(&SeaweedMessagingAgent_ServiceDesc, srv)
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func _SeaweedMessagingAgent_StartPublishSession_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
							 | 
						|
									in := new(StartPublishSessionRequest)
							 | 
						|
									if err := dec(in); err != nil {
							 | 
						|
										return nil, err
							 | 
						|
									}
							 | 
						|
									if interceptor == nil {
							 | 
						|
										return srv.(SeaweedMessagingAgentServer).StartPublishSession(ctx, in)
							 | 
						|
									}
							 | 
						|
									info := &grpc.UnaryServerInfo{
							 | 
						|
										Server:     srv,
							 | 
						|
										FullMethod: SeaweedMessagingAgent_StartPublishSession_FullMethodName,
							 | 
						|
									}
							 | 
						|
									handler := func(ctx context.Context, req interface{}) (interface{}, error) {
							 | 
						|
										return srv.(SeaweedMessagingAgentServer).StartPublishSession(ctx, req.(*StartPublishSessionRequest))
							 | 
						|
									}
							 | 
						|
									return interceptor(ctx, in, info, handler)
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func _SeaweedMessagingAgent_ClosePublishSession_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
							 | 
						|
									in := new(ClosePublishSessionRequest)
							 | 
						|
									if err := dec(in); err != nil {
							 | 
						|
										return nil, err
							 | 
						|
									}
							 | 
						|
									if interceptor == nil {
							 | 
						|
										return srv.(SeaweedMessagingAgentServer).ClosePublishSession(ctx, in)
							 | 
						|
									}
							 | 
						|
									info := &grpc.UnaryServerInfo{
							 | 
						|
										Server:     srv,
							 | 
						|
										FullMethod: SeaweedMessagingAgent_ClosePublishSession_FullMethodName,
							 | 
						|
									}
							 | 
						|
									handler := func(ctx context.Context, req interface{}) (interface{}, error) {
							 | 
						|
										return srv.(SeaweedMessagingAgentServer).ClosePublishSession(ctx, req.(*ClosePublishSessionRequest))
							 | 
						|
									}
							 | 
						|
									return interceptor(ctx, in, info, handler)
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								func _SeaweedMessagingAgent_PublishRecord_Handler(srv interface{}, stream grpc.ServerStream) error {
							 | 
						|
									return srv.(SeaweedMessagingAgentServer).PublishRecord(&grpc.GenericServerStream[PublishRecordRequest, PublishRecordResponse]{ServerStream: stream})
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
							 | 
						|
								type SeaweedMessagingAgent_PublishRecordServer = grpc.BidiStreamingServer[PublishRecordRequest, PublishRecordResponse]
							 | 
						|
								
							 | 
						|
								func _SeaweedMessagingAgent_SubscribeRecord_Handler(srv interface{}, stream grpc.ServerStream) error {
							 | 
						|
									return srv.(SeaweedMessagingAgentServer).SubscribeRecord(&grpc.GenericServerStream[SubscribeRecordRequest, SubscribeRecordResponse]{ServerStream: stream})
							 | 
						|
								}
							 | 
						|
								
							 | 
						|
								// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
							 | 
						|
								type SeaweedMessagingAgent_SubscribeRecordServer = grpc.BidiStreamingServer[SubscribeRecordRequest, SubscribeRecordResponse]
							 | 
						|
								
							 | 
						|
								// SeaweedMessagingAgent_ServiceDesc is the grpc.ServiceDesc for SeaweedMessagingAgent service.
							 | 
						|
								// It's only intended for direct use with grpc.RegisterService,
							 | 
						|
								// and not to be introspected or modified (even as a copy)
							 | 
						|
								var SeaweedMessagingAgent_ServiceDesc = grpc.ServiceDesc{
							 | 
						|
									ServiceName: "messaging_pb.SeaweedMessagingAgent",
							 | 
						|
									HandlerType: (*SeaweedMessagingAgentServer)(nil),
							 | 
						|
									Methods: []grpc.MethodDesc{
							 | 
						|
										{
							 | 
						|
											MethodName: "StartPublishSession",
							 | 
						|
											Handler:    _SeaweedMessagingAgent_StartPublishSession_Handler,
							 | 
						|
										},
							 | 
						|
										{
							 | 
						|
											MethodName: "ClosePublishSession",
							 | 
						|
											Handler:    _SeaweedMessagingAgent_ClosePublishSession_Handler,
							 | 
						|
										},
							 | 
						|
									},
							 | 
						|
									Streams: []grpc.StreamDesc{
							 | 
						|
										{
							 | 
						|
											StreamName:    "PublishRecord",
							 | 
						|
											Handler:       _SeaweedMessagingAgent_PublishRecord_Handler,
							 | 
						|
											ServerStreams: true,
							 | 
						|
											ClientStreams: true,
							 | 
						|
										},
							 | 
						|
										{
							 | 
						|
											StreamName:    "SubscribeRecord",
							 | 
						|
											Handler:       _SeaweedMessagingAgent_SubscribeRecord_Handler,
							 | 
						|
											ServerStreams: true,
							 | 
						|
											ClientStreams: true,
							 | 
						|
										},
							 | 
						|
									},
							 | 
						|
									Metadata: "mq_agent.proto",
							 | 
						|
								}
							 |