From c82f2409e36e4c1edec1f02779543c2074ad19c6 Mon Sep 17 00:00:00 2001 From: gunli Date: Fri, 7 Mar 2025 10:42:30 +0800 Subject: [PATCH 1/3] opt: delete redundant channel --- pulsar/internal/connection.go | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/pulsar/internal/connection.go b/pulsar/internal/connection.go index 1faccc918e..bf474148f9 100644 --- a/pulsar/internal/connection.go +++ b/pulsar/internal/connection.go @@ -160,7 +160,6 @@ type connection struct { incomingRequestsWG sync.WaitGroup incomingRequestsCh chan *request - incomingCmdCh chan *incomingCmd closeCh chan struct{} readyCh chan struct{} writeRequestsCh chan Buffer @@ -209,7 +208,6 @@ func newConnection(opts connectionOptions) *connection { closeCh: make(chan struct{}), readyCh: make(chan struct{}), incomingRequestsCh: make(chan *request, 10), - incomingCmdCh: make(chan *incomingCmd, 10), // This channel is used to pass data from producers to the connection // go routine. It can become contended or blocking if we have multiple @@ -438,9 +436,6 @@ func (c *connection) run() { select { case <-c.closeCh: return - - case cmd := <-c.incomingCmdCh: - c.internalReceivedCommand(cmd.cmd, cmd.headersAndPayload) case data := <-c.writeRequestsCh: if data == nil { return @@ -534,7 +529,7 @@ func (c *connection) writeCommand(cmd *pb.BaseCommand) { } func (c *connection) receivedCommand(cmd *pb.BaseCommand, headersAndPayload Buffer) { - c.incomingCmdCh <- &incomingCmd{cmd, headersAndPayload} + c.internalReceivedCommand(cmd, headersAndPayload) } func (c *connection) internalReceivedCommand(cmd *pb.BaseCommand, headersAndPayload Buffer) { From 4b0853aedae96bdb329d539d9e457e50585d1607 Mon Sep 17 00:00:00 2001 From: gunli Date: Fri, 7 Mar 2025 11:04:21 +0800 Subject: [PATCH 2/3] fix lint error --- pulsar/internal/connection.go | 5 ----- 1 file changed, 5 deletions(-) diff --git a/pulsar/internal/connection.go b/pulsar/internal/connection.go index bf474148f9..58d5b9f406 100644 --- a/pulsar/internal/connection.go +++ b/pulsar/internal/connection.go @@ -129,11 +129,6 @@ type request struct { callback func(command *pb.BaseCommand, err error) } -type incomingCmd struct { - cmd *pb.BaseCommand - headersAndPayload Buffer -} - type connection struct { started int32 connectionTimeout time.Duration From e668eb14c2263520e906ae90e59fc4da7748682b Mon Sep 17 00:00:00 2001 From: gunli Date: Fri, 7 Mar 2025 11:23:16 +0800 Subject: [PATCH 3/3] merge internalReceivedCommand and receivedCommand --- pulsar/internal/connection.go | 4 ---- 1 file changed, 4 deletions(-) diff --git a/pulsar/internal/connection.go b/pulsar/internal/connection.go index 58d5b9f406..a04fa2a270 100644 --- a/pulsar/internal/connection.go +++ b/pulsar/internal/connection.go @@ -524,10 +524,6 @@ func (c *connection) writeCommand(cmd *pb.BaseCommand) { } func (c *connection) receivedCommand(cmd *pb.BaseCommand, headersAndPayload Buffer) { - c.internalReceivedCommand(cmd, headersAndPayload) -} - -func (c *connection) internalReceivedCommand(cmd *pb.BaseCommand, headersAndPayload Buffer) { c.log.Debugf("Received command: %s -- payload: %v", cmd, headersAndPayload) c.setLastDataReceived(time.Now())