Bläddra i källkod

refactoring and error checking fix

Varun 6 år sedan
förälder
incheckning
fcf2bf22c7
1 ändrade filer med 72 tillägg och 38 borttagningar
  1. 72 38
      broker.go

+ 72 - 38
broker.go

@@ -18,19 +18,18 @@ import (
 
 
 // Broker represents a single Kafka broker connection. All operations on this object are entirely concurrency-safe.
 // Broker represents a single Kafka broker connection. All operations on this object are entirely concurrency-safe.
 type Broker struct {
 type Broker struct {
-	id   int32
-	addr string
+	conf *Config
 	rack *string
 	rack *string
 
 
-	conf          *Config
+	id            int32
+	addr          string
 	correlationID int32
 	correlationID int32
 	conn          net.Conn
 	conn          net.Conn
 	connErr       error
 	connErr       error
 	lock          sync.Mutex
 	lock          sync.Mutex
 	opened        int32
 	opened        int32
-
-	responses chan responsePromise
-	done      chan bool
+	responses     chan responsePromise
+	done          chan bool
 
 
 	incomingByteRate       metrics.Meter
 	incomingByteRate       metrics.Meter
 	requestRate            metrics.Meter
 	requestRate            metrics.Meter
@@ -228,6 +227,7 @@ func (b *Broker) Connected() (bool, error) {
 	return b.conn != nil, b.connErr
 	return b.conn != nil, b.connErr
 }
 }
 
 
+//Close closes the broker resources
 func (b *Broker) Close() error {
 func (b *Broker) Close() error {
 	b.lock.Lock()
 	b.lock.Lock()
 	defer b.lock.Unlock()
 	defer b.lock.Unlock()
@@ -285,6 +285,7 @@ func (b *Broker) Rack() string {
 	return *b.rack
 	return *b.rack
 }
 }
 
 
+//GetMetadata send a metadata request and returns a metadata response or error
 func (b *Broker) GetMetadata(request *MetadataRequest) (*MetadataResponse, error) {
 func (b *Broker) GetMetadata(request *MetadataRequest) (*MetadataResponse, error) {
 	response := new(MetadataResponse)
 	response := new(MetadataResponse)
 
 
@@ -297,6 +298,7 @@ func (b *Broker) GetMetadata(request *MetadataRequest) (*MetadataResponse, error
 	return response, nil
 	return response, nil
 }
 }
 
 
+//GetConsumerMetadata send a consumer metadata request and returns a consumer metadata response or error
 func (b *Broker) GetConsumerMetadata(request *ConsumerMetadataRequest) (*ConsumerMetadataResponse, error) {
 func (b *Broker) GetConsumerMetadata(request *ConsumerMetadataRequest) (*ConsumerMetadataResponse, error) {
 	response := new(ConsumerMetadataResponse)
 	response := new(ConsumerMetadataResponse)
 
 
@@ -309,6 +311,7 @@ func (b *Broker) GetConsumerMetadata(request *ConsumerMetadataRequest) (*Consume
 	return response, nil
 	return response, nil
 }
 }
 
 
+//FindCoordinator sends a find coordinate request and returns a response or error
 func (b *Broker) FindCoordinator(request *FindCoordinatorRequest) (*FindCoordinatorResponse, error) {
 func (b *Broker) FindCoordinator(request *FindCoordinatorRequest) (*FindCoordinatorResponse, error) {
 	response := new(FindCoordinatorResponse)
 	response := new(FindCoordinatorResponse)
 
 
@@ -321,6 +324,7 @@ func (b *Broker) FindCoordinator(request *FindCoordinatorRequest) (*FindCoordina
 	return response, nil
 	return response, nil
 }
 }
 
 
+//GetAvailableOffsets return an offset response or error
 func (b *Broker) GetAvailableOffsets(request *OffsetRequest) (*OffsetResponse, error) {
 func (b *Broker) GetAvailableOffsets(request *OffsetRequest) (*OffsetResponse, error) {
 	response := new(OffsetResponse)
 	response := new(OffsetResponse)
 
 
@@ -333,9 +337,12 @@ func (b *Broker) GetAvailableOffsets(request *OffsetRequest) (*OffsetResponse, e
 	return response, nil
 	return response, nil
 }
 }
 
 
+//Produce returns a produce response or error
 func (b *Broker) Produce(request *ProduceRequest) (*ProduceResponse, error) {
 func (b *Broker) Produce(request *ProduceRequest) (*ProduceResponse, error) {
-	var response *ProduceResponse
-	var err error
+	var (
+		response *ProduceResponse
+		err      error
+	)
 
 
 	if request.RequiredAcks == NoResponse {
 	if request.RequiredAcks == NoResponse {
 		err = b.sendAndReceive(request, nil)
 		err = b.sendAndReceive(request, nil)
@@ -351,11 +358,11 @@ func (b *Broker) Produce(request *ProduceRequest) (*ProduceResponse, error) {
 	return response, nil
 	return response, nil
 }
 }
 
 
+//Fetch returns a FetchResponse or error
 func (b *Broker) Fetch(request *FetchRequest) (*FetchResponse, error) {
 func (b *Broker) Fetch(request *FetchRequest) (*FetchResponse, error) {
 	response := new(FetchResponse)
 	response := new(FetchResponse)
 
 
 	err := b.sendAndReceive(request, response)
 	err := b.sendAndReceive(request, response)
-
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
@@ -363,11 +370,11 @@ func (b *Broker) Fetch(request *FetchRequest) (*FetchResponse, error) {
 	return response, nil
 	return response, nil
 }
 }
 
 
+//CommitOffset return an Offset commit reponse or error
 func (b *Broker) CommitOffset(request *OffsetCommitRequest) (*OffsetCommitResponse, error) {
 func (b *Broker) CommitOffset(request *OffsetCommitRequest) (*OffsetCommitResponse, error) {
 	response := new(OffsetCommitResponse)
 	response := new(OffsetCommitResponse)
 
 
 	err := b.sendAndReceive(request, response)
 	err := b.sendAndReceive(request, response)
-
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
@@ -375,11 +382,11 @@ func (b *Broker) CommitOffset(request *OffsetCommitRequest) (*OffsetCommitRespon
 	return response, nil
 	return response, nil
 }
 }
 
 
+//FetchOffset returns an offset fetch response or error
 func (b *Broker) FetchOffset(request *OffsetFetchRequest) (*OffsetFetchResponse, error) {
 func (b *Broker) FetchOffset(request *OffsetFetchRequest) (*OffsetFetchResponse, error) {
 	response := new(OffsetFetchResponse)
 	response := new(OffsetFetchResponse)
 
 
 	err := b.sendAndReceive(request, response)
 	err := b.sendAndReceive(request, response)
-
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
@@ -387,6 +394,7 @@ func (b *Broker) FetchOffset(request *OffsetFetchRequest) (*OffsetFetchResponse,
 	return response, nil
 	return response, nil
 }
 }
 
 
+//JoinGroup returns a join group response or error
 func (b *Broker) JoinGroup(request *JoinGroupRequest) (*JoinGroupResponse, error) {
 func (b *Broker) JoinGroup(request *JoinGroupRequest) (*JoinGroupResponse, error) {
 	response := new(JoinGroupResponse)
 	response := new(JoinGroupResponse)
 
 
@@ -398,6 +406,7 @@ func (b *Broker) JoinGroup(request *JoinGroupRequest) (*JoinGroupResponse, error
 	return response, nil
 	return response, nil
 }
 }
 
 
+//SyncGroup returns a sync group response or error
 func (b *Broker) SyncGroup(request *SyncGroupRequest) (*SyncGroupResponse, error) {
 func (b *Broker) SyncGroup(request *SyncGroupRequest) (*SyncGroupResponse, error) {
 	response := new(SyncGroupResponse)
 	response := new(SyncGroupResponse)
 
 
@@ -409,6 +418,7 @@ func (b *Broker) SyncGroup(request *SyncGroupRequest) (*SyncGroupResponse, error
 	return response, nil
 	return response, nil
 }
 }
 
 
+//LeaveGroup return a leave group response or error
 func (b *Broker) LeaveGroup(request *LeaveGroupRequest) (*LeaveGroupResponse, error) {
 func (b *Broker) LeaveGroup(request *LeaveGroupRequest) (*LeaveGroupResponse, error) {
 	response := new(LeaveGroupResponse)
 	response := new(LeaveGroupResponse)
 
 
@@ -420,6 +430,7 @@ func (b *Broker) LeaveGroup(request *LeaveGroupRequest) (*LeaveGroupResponse, er
 	return response, nil
 	return response, nil
 }
 }
 
 
+//Heartbeat returns a heartbeat response or error
 func (b *Broker) Heartbeat(request *HeartbeatRequest) (*HeartbeatResponse, error) {
 func (b *Broker) Heartbeat(request *HeartbeatRequest) (*HeartbeatResponse, error) {
 	response := new(HeartbeatResponse)
 	response := new(HeartbeatResponse)
 
 
@@ -431,6 +442,7 @@ func (b *Broker) Heartbeat(request *HeartbeatRequest) (*HeartbeatResponse, error
 	return response, nil
 	return response, nil
 }
 }
 
 
+//ListGroups return a list group response or error
 func (b *Broker) ListGroups(request *ListGroupsRequest) (*ListGroupsResponse, error) {
 func (b *Broker) ListGroups(request *ListGroupsRequest) (*ListGroupsResponse, error) {
 	response := new(ListGroupsResponse)
 	response := new(ListGroupsResponse)
 
 
@@ -442,6 +454,7 @@ func (b *Broker) ListGroups(request *ListGroupsRequest) (*ListGroupsResponse, er
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DescribeGroups return describe group response or error
 func (b *Broker) DescribeGroups(request *DescribeGroupsRequest) (*DescribeGroupsResponse, error) {
 func (b *Broker) DescribeGroups(request *DescribeGroupsRequest) (*DescribeGroupsResponse, error) {
 	response := new(DescribeGroupsResponse)
 	response := new(DescribeGroupsResponse)
 
 
@@ -453,6 +466,7 @@ func (b *Broker) DescribeGroups(request *DescribeGroupsRequest) (*DescribeGroups
 	return response, nil
 	return response, nil
 }
 }
 
 
+//ApiVersions return api version response or error
 func (b *Broker) ApiVersions(request *ApiVersionsRequest) (*ApiVersionsResponse, error) {
 func (b *Broker) ApiVersions(request *ApiVersionsRequest) (*ApiVersionsResponse, error) {
 	response := new(ApiVersionsResponse)
 	response := new(ApiVersionsResponse)
 
 
@@ -464,6 +478,7 @@ func (b *Broker) ApiVersions(request *ApiVersionsRequest) (*ApiVersionsResponse,
 	return response, nil
 	return response, nil
 }
 }
 
 
+//CreateTopics send a create topic request and returns create topic response
 func (b *Broker) CreateTopics(request *CreateTopicsRequest) (*CreateTopicsResponse, error) {
 func (b *Broker) CreateTopics(request *CreateTopicsRequest) (*CreateTopicsResponse, error) {
 	response := new(CreateTopicsResponse)
 	response := new(CreateTopicsResponse)
 
 
@@ -475,6 +490,7 @@ func (b *Broker) CreateTopics(request *CreateTopicsRequest) (*CreateTopicsRespon
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DeleteTopics sends a delete topic request and returns delete topic response
 func (b *Broker) DeleteTopics(request *DeleteTopicsRequest) (*DeleteTopicsResponse, error) {
 func (b *Broker) DeleteTopics(request *DeleteTopicsRequest) (*DeleteTopicsResponse, error) {
 	response := new(DeleteTopicsResponse)
 	response := new(DeleteTopicsResponse)
 
 
@@ -486,6 +502,8 @@ func (b *Broker) DeleteTopics(request *DeleteTopicsRequest) (*DeleteTopicsRespon
 	return response, nil
 	return response, nil
 }
 }
 
 
+//CreatePartitions sends a create partition request and returns create
+//partitions response or error
 func (b *Broker) CreatePartitions(request *CreatePartitionsRequest) (*CreatePartitionsResponse, error) {
 func (b *Broker) CreatePartitions(request *CreatePartitionsRequest) (*CreatePartitionsResponse, error) {
 	response := new(CreatePartitionsResponse)
 	response := new(CreatePartitionsResponse)
 
 
@@ -497,6 +515,8 @@ func (b *Broker) CreatePartitions(request *CreatePartitionsRequest) (*CreatePart
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DeleteRecords send a request to delete records and return delete record
+//response or error
 func (b *Broker) DeleteRecords(request *DeleteRecordsRequest) (*DeleteRecordsResponse, error) {
 func (b *Broker) DeleteRecords(request *DeleteRecordsRequest) (*DeleteRecordsResponse, error) {
 	response := new(DeleteRecordsResponse)
 	response := new(DeleteRecordsResponse)
 
 
@@ -508,6 +528,7 @@ func (b *Broker) DeleteRecords(request *DeleteRecordsRequest) (*DeleteRecordsRes
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DescribeAcls sends a describe acl request and returns a response or error
 func (b *Broker) DescribeAcls(request *DescribeAclsRequest) (*DescribeAclsResponse, error) {
 func (b *Broker) DescribeAcls(request *DescribeAclsRequest) (*DescribeAclsResponse, error) {
 	response := new(DescribeAclsResponse)
 	response := new(DescribeAclsResponse)
 
 
@@ -519,6 +540,7 @@ func (b *Broker) DescribeAcls(request *DescribeAclsRequest) (*DescribeAclsRespon
 	return response, nil
 	return response, nil
 }
 }
 
 
+//CreateAcls sends a create acl request and returns a response or error
 func (b *Broker) CreateAcls(request *CreateAclsRequest) (*CreateAclsResponse, error) {
 func (b *Broker) CreateAcls(request *CreateAclsRequest) (*CreateAclsResponse, error) {
 	response := new(CreateAclsResponse)
 	response := new(CreateAclsResponse)
 
 
@@ -530,6 +552,7 @@ func (b *Broker) CreateAcls(request *CreateAclsRequest) (*CreateAclsResponse, er
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DeleteAcls sends a delete acl request and returns a response or error
 func (b *Broker) DeleteAcls(request *DeleteAclsRequest) (*DeleteAclsResponse, error) {
 func (b *Broker) DeleteAcls(request *DeleteAclsRequest) (*DeleteAclsResponse, error) {
 	response := new(DeleteAclsResponse)
 	response := new(DeleteAclsResponse)
 
 
@@ -541,6 +564,7 @@ func (b *Broker) DeleteAcls(request *DeleteAclsRequest) (*DeleteAclsResponse, er
 	return response, nil
 	return response, nil
 }
 }
 
 
+//InitProducerID sends an init producer request and returns a response or error
 func (b *Broker) InitProducerID(request *InitProducerIDRequest) (*InitProducerIDResponse, error) {
 func (b *Broker) InitProducerID(request *InitProducerIDRequest) (*InitProducerIDResponse, error) {
 	response := new(InitProducerIDResponse)
 	response := new(InitProducerIDResponse)
 
 
@@ -552,6 +576,8 @@ func (b *Broker) InitProducerID(request *InitProducerIDRequest) (*InitProducerID
 	return response, nil
 	return response, nil
 }
 }
 
 
+//AddPartitionsToTxn send a request to add partition to txn and returns
+//a response or error
 func (b *Broker) AddPartitionsToTxn(request *AddPartitionsToTxnRequest) (*AddPartitionsToTxnResponse, error) {
 func (b *Broker) AddPartitionsToTxn(request *AddPartitionsToTxnRequest) (*AddPartitionsToTxnResponse, error) {
 	response := new(AddPartitionsToTxnResponse)
 	response := new(AddPartitionsToTxnResponse)
 
 
@@ -563,6 +589,8 @@ func (b *Broker) AddPartitionsToTxn(request *AddPartitionsToTxnRequest) (*AddPar
 	return response, nil
 	return response, nil
 }
 }
 
 
+//AddOffsetsToTxn sends a request to add offsets to txn and returns a response
+//or error
 func (b *Broker) AddOffsetsToTxn(request *AddOffsetsToTxnRequest) (*AddOffsetsToTxnResponse, error) {
 func (b *Broker) AddOffsetsToTxn(request *AddOffsetsToTxnRequest) (*AddOffsetsToTxnResponse, error) {
 	response := new(AddOffsetsToTxnResponse)
 	response := new(AddOffsetsToTxnResponse)
 
 
@@ -574,6 +602,7 @@ func (b *Broker) AddOffsetsToTxn(request *AddOffsetsToTxnRequest) (*AddOffsetsTo
 	return response, nil
 	return response, nil
 }
 }
 
 
+//EndTxn sends a request to end txn and returns a response or error
 func (b *Broker) EndTxn(request *EndTxnRequest) (*EndTxnResponse, error) {
 func (b *Broker) EndTxn(request *EndTxnRequest) (*EndTxnResponse, error) {
 	response := new(EndTxnResponse)
 	response := new(EndTxnResponse)
 
 
@@ -585,6 +614,8 @@ func (b *Broker) EndTxn(request *EndTxnRequest) (*EndTxnResponse, error) {
 	return response, nil
 	return response, nil
 }
 }
 
 
+//TxnOffsetCommit sends a request to commit transaction offsets and returns
+//a response or error
 func (b *Broker) TxnOffsetCommit(request *TxnOffsetCommitRequest) (*TxnOffsetCommitResponse, error) {
 func (b *Broker) TxnOffsetCommit(request *TxnOffsetCommitRequest) (*TxnOffsetCommitResponse, error) {
 	response := new(TxnOffsetCommitResponse)
 	response := new(TxnOffsetCommitResponse)
 
 
@@ -596,6 +627,8 @@ func (b *Broker) TxnOffsetCommit(request *TxnOffsetCommitRequest) (*TxnOffsetCom
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DescribeConfigs sends a request to describe config and returns a response or
+//error
 func (b *Broker) DescribeConfigs(request *DescribeConfigsRequest) (*DescribeConfigsResponse, error) {
 func (b *Broker) DescribeConfigs(request *DescribeConfigsRequest) (*DescribeConfigsResponse, error) {
 	response := new(DescribeConfigsResponse)
 	response := new(DescribeConfigsResponse)
 
 
@@ -607,6 +640,7 @@ func (b *Broker) DescribeConfigs(request *DescribeConfigsRequest) (*DescribeConf
 	return response, nil
 	return response, nil
 }
 }
 
 
+//AlterConfigs sends a request to alter config and return a response or error
 func (b *Broker) AlterConfigs(request *AlterConfigsRequest) (*AlterConfigsResponse, error) {
 func (b *Broker) AlterConfigs(request *AlterConfigsRequest) (*AlterConfigsResponse, error) {
 	response := new(AlterConfigsResponse)
 	response := new(AlterConfigsResponse)
 
 
@@ -618,6 +652,7 @@ func (b *Broker) AlterConfigs(request *AlterConfigsRequest) (*AlterConfigsRespon
 	return response, nil
 	return response, nil
 }
 }
 
 
+//DeleteGroups sends a request to delete groups and returns a response or error
 func (b *Broker) DeleteGroups(request *DeleteGroupsRequest) (*DeleteGroupsResponse, error) {
 func (b *Broker) DeleteGroups(request *DeleteGroupsRequest) (*DeleteGroupsResponse, error) {
 	response := new(DeleteGroupsResponse)
 	response := new(DeleteGroupsResponse)
 
 
@@ -656,7 +691,7 @@ func (b *Broker) send(rb protocolBody, promiseResponse bool) (*responsePromise,
 
 
 	requestTime := time.Now()
 	requestTime := time.Now()
 	bytes, err := b.conn.Write(buf)
 	bytes, err := b.conn.Write(buf)
-	b.updateOutgoingCommunicationMetrics(bytes)
+	b.updateOutgoingCommunicationMetrics(bytes) //TODO: should it be after error check
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
@@ -676,7 +711,6 @@ func (b *Broker) send(rb protocolBody, promiseResponse bool) (*responsePromise,
 
 
 func (b *Broker) sendAndReceive(req protocolBody, res versionedDecoder) error {
 func (b *Broker) sendAndReceive(req protocolBody, res versionedDecoder) error {
 	promise, err := b.send(req, res != nil)
 	promise, err := b.send(req, res != nil)
-
 	if err != nil {
 	if err != nil {
 		return err
 		return err
 	}
 	}
@@ -725,11 +759,11 @@ func (b *Broker) decode(pd packetDecoder, version int16) (err error) {
 }
 }
 
 
 func (b *Broker) encode(pe packetEncoder, version int16) (err error) {
 func (b *Broker) encode(pe packetEncoder, version int16) (err error) {
-
 	host, portstr, err := net.SplitHostPort(b.addr)
 	host, portstr, err := net.SplitHostPort(b.addr)
 	if err != nil {
 	if err != nil {
 		return err
 		return err
 	}
 	}
+
 	port, err := strconv.Atoi(portstr)
 	port, err := strconv.Atoi(portstr)
 	if err != nil {
 	if err != nil {
 		return err
 		return err
@@ -757,6 +791,7 @@ func (b *Broker) encode(pe packetEncoder, version int16) (err error) {
 func (b *Broker) responseReceiver() {
 func (b *Broker) responseReceiver() {
 	var dead error
 	var dead error
 	header := make([]byte, 8)
 	header := make([]byte, 8)
+
 	for response := range b.responses {
 	for response := range b.responses {
 		if dead != nil {
 		if dead != nil {
 			response.errors <- dead
 			response.errors <- dead
@@ -819,7 +854,6 @@ func (b *Broker) authenticateViaSASL() error {
 	default:
 	default:
 		return b.sendAndReceiveSASLPlainAuth()
 		return b.sendAndReceiveSASLPlainAuth()
 	}
 	}
-
 }
 }
 
 
 func (b *Broker) sendAndReceiveSASLHandshake(saslType SASLMechanism, version int16) error {
 func (b *Broker) sendAndReceiveSASLHandshake(saslType SASLMechanism, version int16) error {
@@ -851,6 +885,7 @@ func (b *Broker) sendAndReceiveSASLHandshake(saslType SASLMechanism, version int
 		Logger.Printf("Failed to read SASL handshake header : %s\n", err.Error())
 		Logger.Printf("Failed to read SASL handshake header : %s\n", err.Error())
 		return err
 		return err
 	}
 	}
+
 	length := binary.BigEndian.Uint32(header[:4])
 	length := binary.BigEndian.Uint32(header[:4])
 	payload := make([]byte, length-4)
 	payload := make([]byte, length-4)
 	n, err := io.ReadFull(b.conn, payload)
 	n, err := io.ReadFull(b.conn, payload)
@@ -858,17 +893,21 @@ func (b *Broker) sendAndReceiveSASLHandshake(saslType SASLMechanism, version int
 		Logger.Printf("Failed to read SASL handshake payload : %s\n", err.Error())
 		Logger.Printf("Failed to read SASL handshake payload : %s\n", err.Error())
 		return err
 		return err
 	}
 	}
+
 	b.updateIncomingCommunicationMetrics(n+8, time.Since(requestTime))
 	b.updateIncomingCommunicationMetrics(n+8, time.Since(requestTime))
 	res := &SaslHandshakeResponse{}
 	res := &SaslHandshakeResponse{}
+
 	err = versionedDecode(payload, res, 0)
 	err = versionedDecode(payload, res, 0)
 	if err != nil {
 	if err != nil {
 		Logger.Printf("Failed to parse SASL handshake : %s\n", err.Error())
 		Logger.Printf("Failed to parse SASL handshake : %s\n", err.Error())
 		return err
 		return err
 	}
 	}
+
 	if res.Err != ErrNoError {
 	if res.Err != ErrNoError {
 		Logger.Printf("Invalid SASL Mechanism : %s\n", res.Err.Error())
 		Logger.Printf("Invalid SASL Mechanism : %s\n", res.Err.Error())
 		return res.Err
 		return res.Err
 	}
 	}
+
 	Logger.Print("Successful SASL handshake. Available mechanisms: ", res.EnabledMechanisms)
 	Logger.Print("Successful SASL handshake. Available mechanisms: ", res.EnabledMechanisms)
 	return nil
 	return nil
 }
 }
@@ -899,6 +938,7 @@ func (b *Broker) sendAndReceiveSASLPlainAuth() error {
 			return handshakeErr
 			return handshakeErr
 		}
 		}
 	}
 	}
+
 	length := 1 + len(b.conf.Net.SASL.User) + 1 + len(b.conf.Net.SASL.Password)
 	length := 1 + len(b.conf.Net.SASL.User) + 1 + len(b.conf.Net.SASL.Password)
 	authBytes := make([]byte, length+4) //4 byte length header + auth data
 	authBytes := make([]byte, length+4) //4 byte length header + auth data
 	binary.BigEndian.PutUint32(authBytes, uint32(length))
 	binary.BigEndian.PutUint32(authBytes, uint32(length))
@@ -935,33 +975,27 @@ func (b *Broker) sendAndReceiveSASLPlainAuth() error {
 // sendAndReceiveSASLOAuth performs the authentication flow as described by KIP-255
 // sendAndReceiveSASLOAuth performs the authentication flow as described by KIP-255
 // https://cwiki.apache.org/confluence/pages/viewpage.action?pageId=75968876
 // https://cwiki.apache.org/confluence/pages/viewpage.action?pageId=75968876
 func (b *Broker) sendAndReceiveSASLOAuth(provider AccessTokenProvider) error {
 func (b *Broker) sendAndReceiveSASLOAuth(provider AccessTokenProvider) error {
-
 	if err := b.sendAndReceiveSASLHandshake(SASLTypeOAuth, SASLHandshakeV1); err != nil {
 	if err := b.sendAndReceiveSASLHandshake(SASLTypeOAuth, SASLHandshakeV1); err != nil {
 		return err
 		return err
 	}
 	}
 
 
 	token, err := provider.Token()
 	token, err := provider.Token()
-
 	if err != nil {
 	if err != nil {
 		return err
 		return err
 	}
 	}
 
 
 	requestTime := time.Now()
 	requestTime := time.Now()
-
 	correlationID := b.correlationID
 	correlationID := b.correlationID
 
 
 	bytesWritten, err := b.sendSASLOAuthBearerClientResponse(token, correlationID)
 	bytesWritten, err := b.sendSASLOAuthBearerClientResponse(token, correlationID)
-
 	if err != nil {
 	if err != nil {
 		return err
 		return err
 	}
 	}
 
 
 	b.updateOutgoingCommunicationMetrics(bytesWritten)
 	b.updateOutgoingCommunicationMetrics(bytesWritten)
-
 	b.correlationID++
 	b.correlationID++
 
 
 	bytesRead, err := b.receiveSASLOAuthBearerServerResponse(correlationID)
 	bytesRead, err := b.receiveSASLOAuthBearerServerResponse(correlationID)
-
 	if err != nil {
 	if err != nil {
 		return err
 		return err
 	}
 	}
@@ -1012,6 +1046,7 @@ func (b *Broker) sendAndReceiveSASLSCRAMv1() error {
 			return err
 			return err
 		}
 		}
 	}
 	}
+
 	Logger.Println("SASL authentication succeeded")
 	Logger.Println("SASL authentication succeeded")
 	return nil
 	return nil
 }
 }
@@ -1023,9 +1058,11 @@ func (b *Broker) sendSaslAuthenticateRequest(correlationID int32, msg []byte) (i
 	if err != nil {
 	if err != nil {
 		return 0, err
 		return 0, err
 	}
 	}
+
 	if err := b.conn.SetWriteDeadline(time.Now().Add(b.conf.Net.WriteTimeout)); err != nil {
 	if err := b.conn.SetWriteDeadline(time.Now().Add(b.conf.Net.WriteTimeout)); err != nil {
 		return 0, err
 		return 0, err
 	}
 	}
+
 	return b.conn.Write(buf)
 	return b.conn.Write(buf)
 }
 }
 
 
@@ -1035,37 +1072,39 @@ func (b *Broker) receiveSaslAuthenticateResponse(correlationID int32) ([]byte, e
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
+
 	header := responseHeader{}
 	header := responseHeader{}
 	err = decode(buf, &header)
 	err = decode(buf, &header)
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
+
 	if header.correlationID != correlationID {
 	if header.correlationID != correlationID {
 		return nil, fmt.Errorf("correlation ID didn't match, wanted %d, got %d", b.correlationID, header.correlationID)
 		return nil, fmt.Errorf("correlation ID didn't match, wanted %d, got %d", b.correlationID, header.correlationID)
 	}
 	}
+
 	buf = make([]byte, header.length-correlationIDSize)
 	buf = make([]byte, header.length-correlationIDSize)
 	c, err := io.ReadFull(b.conn, buf)
 	c, err := io.ReadFull(b.conn, buf)
 	bytesRead += c
 	bytesRead += c
 	if err != nil {
 	if err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
+
 	res := &SaslAuthenticateResponse{}
 	res := &SaslAuthenticateResponse{}
 	if err := versionedDecode(buf, res, 0); err != nil {
 	if err := versionedDecode(buf, res, 0); err != nil {
 		return nil, err
 		return nil, err
 	}
 	}
-	if err != nil {
-		return nil, err
-	}
+
 	if res.Err != ErrNoError {
 	if res.Err != ErrNoError {
 		return nil, res.Err
 		return nil, res.Err
 	}
 	}
+
 	return res.SaslAuthBytes, nil
 	return res.SaslAuthBytes, nil
 }
 }
 
 
 // Build SASL/OAUTHBEARER initial client response as described by RFC-7628
 // Build SASL/OAUTHBEARER initial client response as described by RFC-7628
 // https://tools.ietf.org/html/rfc7628
 // https://tools.ietf.org/html/rfc7628
 func buildClientInitialResponse(token *AccessToken) ([]byte, error) {
 func buildClientInitialResponse(token *AccessToken) ([]byte, error) {
-
 	var ext string
 	var ext string
 
 
 	if token.Extensions != nil && len(token.Extensions) > 0 {
 	if token.Extensions != nil && len(token.Extensions) > 0 {
@@ -1083,7 +1122,6 @@ func buildClientInitialResponse(token *AccessToken) ([]byte, error) {
 // mapToString returns a list of key-value pairs ordered by key.
 // mapToString returns a list of key-value pairs ordered by key.
 // keyValSep separates the key from the value. elemSep separates each pair.
 // keyValSep separates the key from the value. elemSep separates each pair.
 func mapToString(extensions map[string]string, keyValSep string, elemSep string) string {
 func mapToString(extensions map[string]string, keyValSep string, elemSep string) string {
-
 	buf := make([]string, 0, len(extensions))
 	buf := make([]string, 0, len(extensions))
 
 
 	for k, v := range extensions {
 	for k, v := range extensions {
@@ -1096,9 +1134,7 @@ func mapToString(extensions map[string]string, keyValSep string, elemSep string)
 }
 }
 
 
 func (b *Broker) sendSASLOAuthBearerClientResponse(token *AccessToken, correlationID int32) (int, error) {
 func (b *Broker) sendSASLOAuthBearerClientResponse(token *AccessToken, correlationID int32) (int, error) {
-
 	initialResp, err := buildClientInitialResponse(token)
 	initialResp, err := buildClientInitialResponse(token)
-
 	if err != nil {
 	if err != nil {
 		return 0, err
 		return 0, err
 	}
 	}
@@ -1108,7 +1144,6 @@ func (b *Broker) sendSASLOAuthBearerClientResponse(token *AccessToken, correlati
 	req := &request{correlationID: correlationID, clientID: b.conf.ClientID, body: rb}
 	req := &request{correlationID: correlationID, clientID: b.conf.ClientID, body: rb}
 
 
 	buf, err := encode(req, b.conf.MetricRegistry)
 	buf, err := encode(req, b.conf.MetricRegistry)
-
 	if err != nil {
 	if err != nil {
 		return 0, err
 		return 0, err
 	}
 	}
@@ -1121,11 +1156,9 @@ func (b *Broker) sendSASLOAuthBearerClientResponse(token *AccessToken, correlati
 }
 }
 
 
 func (b *Broker) receiveSASLOAuthBearerServerResponse(correlationID int32) (int, error) {
 func (b *Broker) receiveSASLOAuthBearerServerResponse(correlationID int32) (int, error) {
-
 	buf := make([]byte, 8)
 	buf := make([]byte, 8)
 
 
 	bytesRead, err := io.ReadFull(b.conn, buf)
 	bytesRead, err := io.ReadFull(b.conn, buf)
-
 	if err != nil {
 	if err != nil {
 		return bytesRead, err
 		return bytesRead, err
 	}
 	}
@@ -1133,7 +1166,6 @@ func (b *Broker) receiveSASLOAuthBearerServerResponse(correlationID int32) (int,
 	header := responseHeader{}
 	header := responseHeader{}
 
 
 	err = decode(buf, &header)
 	err = decode(buf, &header)
-
 	if err != nil {
 	if err != nil {
 		return bytesRead, err
 		return bytesRead, err
 	}
 	}
@@ -1145,9 +1177,7 @@ func (b *Broker) receiveSASLOAuthBearerServerResponse(correlationID int32) (int,
 	buf = make([]byte, header.length-4)
 	buf = make([]byte, header.length-4)
 
 
 	c, err := io.ReadFull(b.conn, buf)
 	c, err := io.ReadFull(b.conn, buf)
-
 	bytesRead += c
 	bytesRead += c
-
 	if err != nil {
 	if err != nil {
 		return bytesRead, err
 		return bytesRead, err
 	}
 	}
@@ -1158,10 +1188,6 @@ func (b *Broker) receiveSASLOAuthBearerServerResponse(correlationID int32) (int,
 		return bytesRead, err
 		return bytesRead, err
 	}
 	}
 
 
-	if err != nil {
-		return bytesRead, err
-	}
-
 	if res.Err != ErrNoError {
 	if res.Err != ErrNoError {
 		return bytesRead, res.Err
 		return bytesRead, res.Err
 	}
 	}
@@ -1176,14 +1202,17 @@ func (b *Broker) receiveSASLOAuthBearerServerResponse(correlationID int32) (int,
 func (b *Broker) updateIncomingCommunicationMetrics(bytes int, requestLatency time.Duration) {
 func (b *Broker) updateIncomingCommunicationMetrics(bytes int, requestLatency time.Duration) {
 	b.updateRequestLatencyMetrics(requestLatency)
 	b.updateRequestLatencyMetrics(requestLatency)
 	b.responseRate.Mark(1)
 	b.responseRate.Mark(1)
+
 	if b.brokerResponseRate != nil {
 	if b.brokerResponseRate != nil {
 		b.brokerResponseRate.Mark(1)
 		b.brokerResponseRate.Mark(1)
 	}
 	}
+
 	responseSize := int64(bytes)
 	responseSize := int64(bytes)
 	b.incomingByteRate.Mark(responseSize)
 	b.incomingByteRate.Mark(responseSize)
 	if b.brokerIncomingByteRate != nil {
 	if b.brokerIncomingByteRate != nil {
 		b.brokerIncomingByteRate.Mark(responseSize)
 		b.brokerIncomingByteRate.Mark(responseSize)
 	}
 	}
+
 	b.responseSize.Update(responseSize)
 	b.responseSize.Update(responseSize)
 	if b.brokerResponseSize != nil {
 	if b.brokerResponseSize != nil {
 		b.brokerResponseSize.Update(responseSize)
 		b.brokerResponseSize.Update(responseSize)
@@ -1193,9 +1222,11 @@ func (b *Broker) updateIncomingCommunicationMetrics(bytes int, requestLatency ti
 func (b *Broker) updateRequestLatencyMetrics(requestLatency time.Duration) {
 func (b *Broker) updateRequestLatencyMetrics(requestLatency time.Duration) {
 	requestLatencyInMs := int64(requestLatency / time.Millisecond)
 	requestLatencyInMs := int64(requestLatency / time.Millisecond)
 	b.requestLatency.Update(requestLatencyInMs)
 	b.requestLatency.Update(requestLatencyInMs)
+
 	if b.brokerRequestLatency != nil {
 	if b.brokerRequestLatency != nil {
 		b.brokerRequestLatency.Update(requestLatencyInMs)
 		b.brokerRequestLatency.Update(requestLatencyInMs)
 	}
 	}
+
 }
 }
 
 
 func (b *Broker) updateOutgoingCommunicationMetrics(bytes int) {
 func (b *Broker) updateOutgoingCommunicationMetrics(bytes int) {
@@ -1203,13 +1234,16 @@ func (b *Broker) updateOutgoingCommunicationMetrics(bytes int) {
 	if b.brokerRequestRate != nil {
 	if b.brokerRequestRate != nil {
 		b.brokerRequestRate.Mark(1)
 		b.brokerRequestRate.Mark(1)
 	}
 	}
+
 	requestSize := int64(bytes)
 	requestSize := int64(bytes)
 	b.outgoingByteRate.Mark(requestSize)
 	b.outgoingByteRate.Mark(requestSize)
 	if b.brokerOutgoingByteRate != nil {
 	if b.brokerOutgoingByteRate != nil {
 		b.brokerOutgoingByteRate.Mark(requestSize)
 		b.brokerOutgoingByteRate.Mark(requestSize)
 	}
 	}
+
 	b.requestSize.Update(requestSize)
 	b.requestSize.Update(requestSize)
 	if b.brokerRequestSize != nil {
 	if b.brokerRequestSize != nil {
 		b.brokerRequestSize.Update(requestSize)
 		b.brokerRequestSize.Update(requestSize)
 	}
 	}
+
 }
 }