MCPcopy Create free account
hub / github.com/IBM/sarama / encode

Method encode

offset_fetch_request.go:72–165  ·  view source on GitHub ↗
(pe packetEncoder)

Source from the content-addressed store, hash-verified

70}
71
72func (r *OffsetFetchRequest) encode(pe packetEncoder) (err error) {
73 if r.Version < 0 || r.Version > 8 {
74 return PacketEncodingError{"invalid or unsupported OffsetFetchRequest version field"}
75 }
76
77 if r.RequireStable && r.Version < 7 {
78 return PacketEncodingError{"requireStable is not supported. use version 7 or later"}
79 }
80
81 if r.Version >= 8 {
82 if len(r.Groups) == 0 {
83 return PacketEncodingError{"version 8 or later requires Groups to be populated"}
84 }
85 if err := pe.putArrayLength(len(r.Groups)); err != nil {
86 return err
87 }
88 for _, g := range r.Groups {
89 if err := pe.putString(g.GroupId); err != nil {
90 return err
91 }
92
93 // nil Partitions encodes as null topics array (fetch all)
94 topicCount := len(g.Partitions)
95 if g.Partitions == nil {
96 topicCount = -1
97 }
98 if err := pe.putArrayLength(topicCount); err != nil {
99 return err
100 }
101 for topic, partitions := range g.Partitions {
102 if err := pe.putString(topic); err != nil {
103 return err
104 }
105 if err := pe.putInt32Array(partitions); err != nil {
106 return err
107 }
108 pe.putEmptyTaggedFieldArray()
109 }
110
111 pe.putEmptyTaggedFieldArray()
112 }
113
114 pe.putBool(r.RequireStable)
115 pe.putEmptyTaggedFieldArray()
116 return nil
117 }
118
119 if len(r.Groups) > 1 {
120 return PacketEncodingError{"multiple groups require version 8 or later"}
121 }
122
123 consumerGroup := r.ConsumerGroup
124 partitions := r.partitions
125 if consumerGroup == "" && len(r.Groups) == 1 {
126 consumerGroup = r.Groups[0].GroupId
127 partitions = r.Groups[0].Partitions
128 }
129

Callers

nothing calls this directly

Calls 5

putArrayLengthMethod · 0.65
putStringMethod · 0.65
putInt32ArrayMethod · 0.65
putBoolMethod · 0.65

Tested by

no test coverage detected