Skip to content

Commit 50003b7

Browse files
3AceShowHanddnwe
andauthored
fix: correct ApiVersionsResponse handling of ErrUnsupportedVersion (IBM#3337) (#20)
Improve the api_versions_response_test.go to be more representative of the different good and bad responses it might receive. Subsequently fix the ApiVersionsResponse decoding so that it correctly downgrades the decoder from flexible to non-flexible after reading the ErrorCode of UnsupportedVersion Signed-off-by: Dominic Evans <dominic.evans@uk.ibm.com> Co-authored-by: Dominic Evans <8060970+dnwe@users.noreply.github.com>
1 parent 9ac62a1 commit 50003b7

3 files changed

Lines changed: 161 additions & 41 deletions

File tree

api_versions_response.go

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
package sarama
22

3-
import "time"
3+
import (
4+
"time"
5+
)
46

57
// ApiVersionsResponseKey contains the APIs supported by the broker.
68
type ApiVersionsResponseKey struct {
@@ -91,6 +93,15 @@ func (r *ApiVersionsResponse) decode(pd packetDecoder, version int16) (err error
9193
return err
9294
}
9395

96+
// KIP-511: if broker didn't understand the ApiVersionsRequest version then
97+
// it replies with a V0 non-flexible ApiVersionResponse where its supported
98+
// ApiVersionsRequest version is available in ApiKeys
99+
if r.ErrorCode == int16(ErrUnsupportedVersion) {
100+
// drop version to 0 and to revert packageDecoder to non-flexible for remaining decoding
101+
r.Version = 0
102+
pd = downgradeFlexibleDecoder(pd)
103+
}
104+
94105
numApiKeys, err := pd.getArrayLength()
95106
if err != nil {
96107
return err

api_versions_response_test.go

Lines changed: 142 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -2,60 +2,162 @@
22

33
package sarama
44

5-
import "testing"
5+
import (
6+
"testing"
7+
8+
assert "github.com/stretchr/testify/require"
9+
)
610

711
var (
8-
apiVersionResponse = []byte{
9-
0x00, 0x00,
10-
0x00, 0x00, 0x00, 0x01,
11-
0x00, 0x03,
12-
0x00, 0x02,
13-
0x00, 0x01,
12+
apiVersionResponseV0 = []byte{
13+
0x00, 0x00, // no error
14+
0x00, 0x00, 0x00, 0x04, // array length 4 (APIs)
15+
0x00, 0x00, 0x00, 0x00, 0x00, 0x02, // API Version Produce (v0-2)
16+
0x00, 0x01, 0x00, 0x00, 0x00, 0x03, // API Version Fetch (v0-3)
17+
0x00, 0x02, 0x00, 0x00, 0x00, 0x01, // API Version Offsets (v0-1)
18+
0x00, 0x03, 0x00, 0x00, 0x00, 0x02, // API Version Metadata (v0-2)
19+
}
20+
21+
apiVersionResponseV1V2 = []byte{
22+
0x00, 0x00, // no error
23+
0x00, 0x00, 0x00, 0x05, // array length 5 (APIs)
24+
0x00, 0x00, 0x00, 0x00, 0x00, 0x07, // API Version Produce (v0-7)
25+
0x00, 0x01, 0x00, 0x00, 0x00, 0x0b, // API Version Fetch (v0-11)
26+
0x00, 0x02, 0x00, 0x00, 0x00, 0x05, // API Version Offsets (v0-5)
27+
0x00, 0x03, 0x00, 0x00, 0x00, 0x08, // API Version Metadata (v0-8)
28+
0x00, 0x04, 0x00, 0x00, 0x00, 0x02, // API Version LeaderAndIsr (v0-2)
29+
0x00, 0x00, 0x00, 0x40, // throttle time (64ms)
1430
}
1531

1632
apiVersionResponseV3 = []byte{
1733
0x00, 0x00, // no error
18-
0x02, // compact array length 1
19-
0x00, 0x03,
20-
0x00, 0x02,
21-
0x00, 0x01,
22-
0x00, // tagged fields
23-
0x00, 0x00, 0x00, 0x00, // throttle time
24-
0x01, 0x01, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // tagged fields (empty SupportedFeatures)
34+
0x07, // compact array length 6 (APIs)
35+
0x00, 0x00, 0x00, 0x00, 0x00, 0x08, // API Version Produce (v0-8)
36+
0x00, // empty tagged fields
37+
0x00, 0x01, 0x00, 0x00, 0x00, 0x0b, // API Version Fetch (v0-11)
38+
0x00, // empty tagged fields
39+
0x00, 0x02, 0x00, 0x00, 0x00, 0x05, // API Version Offsets (v0-5)
40+
0x00, // empty tagged fields
41+
0x00, 0x03, 0x00, 0x00, 0x00, 0x09, // API Version Metadata (v0-9)
42+
0x00, // empty tagged fields
43+
0x00, 0x04, 0x00, 0x00, 0x00, 0x04, // API Version LeaderAndIsr (v0-4)
44+
0x00, // empty tagged fields
45+
0x00, 0x05, 0x00, 0x00, 0x00, 0x02, // API Version StopReplica (v0-2)
46+
0x00, // empty tagged fields
47+
0x00, 0x00, 0x00, 0x80, // throttle time (128ms)
48+
0x00, // empty tagged fields
2549
}
26-
)
2750

28-
func TestApiVersionsResponse(t *testing.T) {
29-
response := new(ApiVersionsResponse)
30-
testVersionDecodable(t, "no error", response, apiVersionResponse, 0)
31-
if response.ErrorCode != int16(ErrNoError) {
32-
t.Error("Decoding error failed: no error expected but found", response.ErrorCode)
51+
// unsupported version from kafka 0.10.2.1
52+
apiVersionsResponseUnsupportedVersionV0 = []byte{
53+
0x00, 0x23, // unsupported version error
54+
0x00, 0x00, 0x00, 0x00, // array length 0
3355
}
34-
if response.ApiKeys[0].ApiKey != 0x03 {
35-
t.Error("Decoding error: expected 0x03 but got", response.ApiKeys[0].ApiKey)
56+
57+
// unsupported version from kafka 2.3.0
58+
apiVersionsResponseUnsupportedVersionV1V2 = []byte{
59+
0x00, 0x23, // unsupported version error
60+
0x00, 0x00, 0x00, 0x00, // array length 0
61+
}
62+
63+
// unsupported version from kafka 2.4.0
64+
apiVersionsResponseUnsupportedVersionV3 = []byte{
65+
0x00, 0x23, // unsupported version error
66+
0x00, 0x00, 0x00, 0x01, // array length 1
67+
0x00, 0x12, 0x00, 0x00, 0x00, 0x03, // API Version ApiVersions (v0-3)
3668
}
37-
if response.ApiKeys[0].MinVersion != 0x02 {
38-
t.Error("Decoding error: expected 0x02 but got", response.ApiKeys[0].MinVersion)
69+
70+
// unsupported version from kafka 4.1.0
71+
apiVersionsResponseUnsupportedVersionV4 = []byte{
72+
0x00, 0x23, // unsupported version error
73+
0x00, 0x00, 0x00, 0x01, // array length 1
74+
0x00, 0x12, 0x00, 0x00, 0x00, 0x04, // API Version ApiVersions (v0-4)
3975
}
40-
if response.ApiKeys[0].MaxVersion != 0x01 {
41-
t.Error("Decoding error: expected 0x01 but got", response.ApiKeys[0].MaxVersion)
76+
)
77+
78+
func TestApiVersionsResponseV0(t *testing.T) {
79+
const v = 0
80+
response := new(ApiVersionsResponse)
81+
testVersionDecodable(t, "no error V0", response, apiVersionResponseV0, v)
82+
83+
assert.Equal(t, int16(ErrNoError), response.ErrorCode)
84+
assert.Equal(t, []ApiVersionsResponseKey{
85+
{v, 0, 0, 2}, // API Version Produce (v0-2)
86+
{v, 1, 0, 3}, // API Version Fetch (v0-3)
87+
{v, 2, 0, 1}, // API Version Offsets (v0-1)
88+
{v, 3, 0, 2}, // API Version Metadata (v0-2)
89+
}, response.ApiKeys)
90+
}
91+
92+
func TestApiVersionsResponseV1V2(t *testing.T) {
93+
response := new(ApiVersionsResponse)
94+
95+
for _, v := range []int16{1, 2} {
96+
testVersionDecodable(t, "no error V1V2", response, apiVersionResponseV1V2, v)
97+
98+
assert.Equal(t, int16(ErrNoError), response.ErrorCode)
99+
assert.Equal(t, []ApiVersionsResponseKey{
100+
{v, 0, 0, 7}, // API Version Produce (v0-7)
101+
{v, 1, 0, 11}, // API Version Fetch (v0-11)
102+
{v, 2, 0, 5}, // API Version Offsets (v0-5)
103+
{v, 3, 0, 8}, // API Version Metadata (v0-8)
104+
{v, 4, 0, 2}, // API Version LeaderAndIsr (v0-2)
105+
}, response.ApiKeys)
106+
assert.Equal(t, int32(64), response.ThrottleTimeMs)
42107
}
43108
}
44109

45110
func TestApiVersionsResponseV3(t *testing.T) {
111+
const v = 3
46112
response := new(ApiVersionsResponse)
47-
response.Version = 3
48-
testVersionDecodable(t, "no error", response, apiVersionResponseV3, 3)
49-
if response.ErrorCode != int16(ErrNoError) {
50-
t.Error("Decoding error failed: no error expected but found", response.ErrorCode)
51-
}
52-
if response.ApiKeys[0].ApiKey != 0x03 {
53-
t.Error("Decoding error: expected 0x03 but got", response.ApiKeys[0].ApiKey)
54-
}
55-
if response.ApiKeys[0].MinVersion != 0x02 {
56-
t.Error("Decoding error: expected 0x02 but got", response.ApiKeys[0].MinVersion)
57-
}
58-
if response.ApiKeys[0].MaxVersion != 0x01 {
59-
t.Error("Decoding error: expected 0x01 but got", response.ApiKeys[0].MaxVersion)
60-
}
113+
response.Version = v
114+
testVersionDecodable(t, "no error V3", response, apiVersionResponseV3, v)
115+
assert.Equal(t, int16(ErrNoError), response.ErrorCode)
116+
assert.Equal(t, []ApiVersionsResponseKey{
117+
{v, 0, 0, 8}, // API Version Produce (v0-8)
118+
{v, 1, 0, 11}, // API Version Fetch (v0-11)
119+
{v, 2, 0, 5}, // API Version Offsets (v0-5)
120+
{v, 3, 0, 9}, // API Version Metadata (v0-9)
121+
{v, 4, 0, 4}, // API Version LeaderAndIsr (v0-4)
122+
{v, 5, 0, 2}, // API Version StopReplica (v0-2)
123+
}, response.ApiKeys)
124+
assert.Equal(t, int32(128), response.ThrottleTimeMs)
125+
}
126+
127+
func TestApiVersionsResponseUnsupportedVersion(t *testing.T) {
128+
t.Run("V0", func(t *testing.T) {
129+
response := new(ApiVersionsResponse)
130+
response.Version = 3
131+
testVersionDecodable(t, "unsupported", response, apiVersionsResponseUnsupportedVersionV0, 3)
132+
assert.Equal(t, int16(ErrUnsupportedVersion), response.ErrorCode)
133+
assert.Empty(t, response.ApiKeys)
134+
})
135+
136+
t.Run("V1V2", func(t *testing.T) {
137+
response := new(ApiVersionsResponse)
138+
response.Version = 3
139+
testVersionDecodable(t, "unsupported", response, apiVersionsResponseUnsupportedVersionV1V2, 3)
140+
assert.Equal(t, int16(ErrUnsupportedVersion), response.ErrorCode)
141+
assert.Empty(t, response.ApiKeys)
142+
})
143+
144+
t.Run("V3", func(t *testing.T) {
145+
response := new(ApiVersionsResponse)
146+
response.Version = 3
147+
testVersionDecodable(t, "unsupported", response, apiVersionsResponseUnsupportedVersionV3, 3)
148+
assert.Equal(t, int16(ErrUnsupportedVersion), response.ErrorCode)
149+
assert.Equal(t, []ApiVersionsResponseKey{
150+
{0, 18, 0, 3}, // API Version ApiVersions (v0-3)
151+
}, response.ApiKeys)
152+
})
153+
154+
t.Run("V4", func(t *testing.T) {
155+
response := new(ApiVersionsResponse)
156+
response.Version = 4
157+
testVersionDecodable(t, "unsupported", response, apiVersionsResponseUnsupportedVersionV4, 4)
158+
assert.Equal(t, int16(ErrUnsupportedVersion), response.ErrorCode)
159+
assert.Equal(t, []ApiVersionsResponseKey{
160+
{0, 18, 0, 4}, // API Version ApiVersions (v0-4)
161+
}, response.ApiKeys)
162+
})
61163
}

encoder_decoder.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,3 +125,10 @@ func prepareFlexibleEncoder(pe packetEncoder, req encoder) packetEncoder {
125125
}
126126
return pe
127127
}
128+
129+
func downgradeFlexibleDecoder(pd packetDecoder) packetDecoder {
130+
if f, ok := pd.(*realFlexibleDecoder); ok {
131+
return f.realDecoder
132+
}
133+
return pd
134+
}

0 commit comments

Comments
 (0)