-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathsync_group.go
58 lines (47 loc) · 1.55 KB
/
sync_group.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
package client
// SyncGroupRequest is used to synchronize state for all members of a group.
type SyncGroupRequest struct {
GroupID string
GenerationID int32
MemberID string
GroupAssignment map[string][]byte
}
// Key returns the Kafka API key for SyncGroupRequest.
func (*SyncGroupRequest) Key() int16 {
return 14
}
// Version returns the Kafka request version for backwards compatibility.
func (*SyncGroupRequest) Version() int16 {
return 0
}
func (sgr *SyncGroupRequest) Write(encoder Encoder) {
encoder.WriteString(sgr.GroupID)
encoder.WriteInt32(sgr.GenerationID)
encoder.WriteString(sgr.MemberID)
encoder.WriteInt32(int32(len(sgr.GroupAssignment)))
for memberID, memberAssignment := range sgr.GroupAssignment {
encoder.WriteString(memberID)
encoder.WriteBytes(memberAssignment)
}
}
// SyncGroupResponse contains information about partition distribution within a group.
type SyncGroupResponse struct {
Error error
MemberAssignment []byte
}
func (sgr *SyncGroupResponse) Read(decoder Decoder) *DecodingError {
errCode, err := decoder.GetInt16()
if err != nil {
return NewDecodingError(err, reasonInvalidSyncGroupResponseErrorCode)
}
sgr.Error = BrokerErrors[errCode]
sgr.MemberAssignment, err = decoder.GetBytes()
if err != nil {
return NewDecodingError(err, reasonInvalidSyncGroupResponseMemberAssignment)
}
return nil
}
var (
reasonInvalidSyncGroupResponseErrorCode = "Invalid error code in SyncGroupResponse"
reasonInvalidSyncGroupResponseMemberAssignment = "Invalid SyncGroupResponse MemberAssignment"
)