(pe packetEncoder)
| 70 | } |
| 71 | |
| 72 | func (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 |
nothing calls this directly
no test coverage detected