* fix: use keyed fields in struct literals - Replace unsafe reflect.StringHeader/SliceHeader with safe unsafe.String/Slice (weed/query/sqltypes/unsafe.go) - Add field names to Type_ScalarType struct literals (weed/mq/schema/schema_builder.go) - Add Duration field name to FlexibleDuration struct literals across test files - Add field names to bson.D struct literals (weed/filer/mongodb/mongodb_store_kv.go) Fixes go vet warnings about unkeyed struct literals. * fix: remove unreachable code - Remove unreachable return statements after infinite for loops - Remove unreachable code after if/else blocks where all paths return - Simplify recursive logic by removing unnecessary for loop (inode_to_path.go) - Fix Type_ScalarType literal to use enum value directly (schema_builder.go) - Call onCompletionFn on stream error (subscribe_session.go) Files fixed: - weed/query/sqltypes/unsafe.go - weed/mq/schema/schema_builder.go - weed/mq/client/sub_client/connect_to_sub_coordinator.go - weed/filer/redis3/ItemList.go - weed/mq/client/agent_client/subscribe_session.go - weed/mq/broker/broker_grpc_pub_balancer.go - weed/mount/inode_to_path.go - weed/util/skiplist/name_list.go * fix: avoid copying lock values in protobuf messages - Use proto.Merge() instead of direct assignment to avoid copying sync.Mutex in S3ApiConfiguration (iamapi_server.go) - Add explicit comments noting that channel-received values are already copies before taking addresses (volume_grpc_client_to_master.go) The protobuf messages contain sync.Mutex fields from the message state, which should not be copied. Using proto.Merge() properly merges messages without copying the embedded mutex. * fix: correct byte array size for uint32 bit shift operations The generateAccountId() function only needs 4 bytes to create a uint32 value. Changed from allocating 8 bytes to 4 bytes to match the actual usage. This fixes go vet warning about shifting 8-bit values (bytes) by more than 8 bits. * fix: ensure context cancellation on all error paths In broker_client_subscribe.go, ensure subscriberCancel() is called on all error return paths: - When stream creation fails - When partition assignment fails - When sending initialization message fails This prevents context leaks when an error occurs during subscriber creation. * fix: ensure subscriberCancel called for CreateFreshSubscriber stream.Send error Ensure subscriberCancel() is called when stream.Send fails in CreateFreshSubscriber. * ci: add go vet step to prevent future lint regressions - Add go vet step to GitHub Actions workflow - Filter known protobuf lock warnings (MessageState sync.Mutex) These are expected in generated protobuf code and are safe - Prevents accumulation of go vet errors in future PRs - Step runs before build to catch issues early * fix: resolve remaining syntax and logic errors in vet fixes - Fixed syntax errors in filer_sync.go caused by missing closing braces - Added missing closing brace for if block and function - Synchronized fixes to match previous commits on branch * fix: add missing return statements to daemon functions - Add 'return false' after infinite loops in filer_backup.go and filer_meta_backup.go - Satisfies declared bool return type signatures - Maintains consistency with other daemon functions (runMaster, runFilerSynchronize, runWorker) - While unreachable, explicitly declares the return satisfies function signature contract * fix: add nil check for onCompletionFn in SubscribeMessageRecord - Check if onCompletionFn is not nil before calling it - Prevents potential panic if nil function is passed - Matches pattern used in other callback functions * docs: clarify unreachable return statements in daemon functions - Add comments documenting that return statements satisfy function signature - Explains that these returns follow infinite loops and are unreachable - Improves code clarity for future maintainers
88 lines
2.4 KiB
Go
88 lines
2.4 KiB
Go
package agent_client
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/mq/topic"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/mq_agent_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/schema_pb"
|
|
"google.golang.org/grpc"
|
|
"google.golang.org/grpc/credentials/insecure"
|
|
)
|
|
|
|
type SubscribeOption struct {
|
|
ConsumerGroup string
|
|
ConsumerGroupInstanceId string
|
|
Topic topic.Topic
|
|
OffsetType schema_pb.OffsetType
|
|
OffsetTsNs int64
|
|
Filter string
|
|
MaxSubscribedPartitions int32
|
|
SlidingWindowSize int32
|
|
}
|
|
|
|
type SubscribeSession struct {
|
|
Option *SubscribeOption
|
|
stream grpc.BidiStreamingClient[mq_agent_pb.SubscribeRecordRequest, mq_agent_pb.SubscribeRecordResponse]
|
|
}
|
|
|
|
func NewSubscribeSession(agentAddress string, option *SubscribeOption) (*SubscribeSession, error) {
|
|
// call local agent grpc server to create a new session
|
|
clientConn, err := grpc.NewClient(agentAddress, grpc.WithTransportCredentials(insecure.NewCredentials()))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("dial agent server %s: %v", agentAddress, err)
|
|
}
|
|
agentClient := mq_agent_pb.NewSeaweedMessagingAgentClient(clientConn)
|
|
|
|
initRequest := &mq_agent_pb.SubscribeRecordRequest_InitSubscribeRecordRequest{
|
|
ConsumerGroup: option.ConsumerGroup,
|
|
ConsumerGroupInstanceId: option.ConsumerGroupInstanceId,
|
|
Topic: &schema_pb.Topic{
|
|
Namespace: option.Topic.Namespace,
|
|
Name: option.Topic.Name,
|
|
},
|
|
OffsetType: option.OffsetType,
|
|
OffsetTsNs: option.OffsetTsNs,
|
|
MaxSubscribedPartitions: option.MaxSubscribedPartitions,
|
|
Filter: option.Filter,
|
|
SlidingWindowSize: option.SlidingWindowSize,
|
|
}
|
|
|
|
stream, err := agentClient.SubscribeRecord(context.Background())
|
|
if err != nil {
|
|
return nil, fmt.Errorf("subscribe record: %w", err)
|
|
}
|
|
|
|
if err = stream.Send(&mq_agent_pb.SubscribeRecordRequest{
|
|
Init: initRequest,
|
|
}); err != nil {
|
|
return nil, fmt.Errorf("send session id: %w", err)
|
|
}
|
|
|
|
return &SubscribeSession{
|
|
Option: option,
|
|
stream: stream,
|
|
}, nil
|
|
}
|
|
|
|
func (s *SubscribeSession) CloseSession() error {
|
|
err := s.stream.CloseSend()
|
|
return err
|
|
}
|
|
|
|
func (a *SubscribeSession) SubscribeMessageRecord(
|
|
onEachMessageFn func(key []byte, record *schema_pb.RecordValue),
|
|
onCompletionFn func()) error {
|
|
for {
|
|
resp, err := a.stream.Recv()
|
|
if err != nil {
|
|
if onCompletionFn != nil {
|
|
onCompletionFn()
|
|
}
|
|
return err
|
|
}
|
|
onEachMessageFn(resp.Key, resp.Value)
|
|
}
|
|
}
|