|
4 | 4 | "bytes"
|
5 | 5 | "testing"
|
6 | 6 |
|
| 7 | + kcronsumer "github.com/Trendyol/kafka-cronsumer/pkg/kafka" |
| 8 | + |
7 | 9 | "github.com/segmentio/kafka-go"
|
8 | 10 | )
|
9 | 11 |
|
@@ -113,3 +115,135 @@ func TestMessage_RemoveHeader(t *testing.T) {
|
113 | 115 | t.Fatalf("Header length must be equal to 0")
|
114 | 116 | }
|
115 | 117 | }
|
| 118 | + |
| 119 | +func TestMessage_toRetryableMessage(t *testing.T) { |
| 120 | + t.Run("When_error_description_exist", func(t *testing.T) { |
| 121 | + // Given |
| 122 | + message := Message{ |
| 123 | + Key: []byte("key"), |
| 124 | + Value: []byte("value"), |
| 125 | + Headers: []Header{ |
| 126 | + { |
| 127 | + Key: "x-custom-client-header", |
| 128 | + Value: []byte("bar"), |
| 129 | + }, |
| 130 | + }, |
| 131 | + ErrDescription: "some error description", |
| 132 | + } |
| 133 | + expected := kcronsumer.Message{ |
| 134 | + Topic: "retry-topic", |
| 135 | + Key: []byte("key"), |
| 136 | + Value: []byte("value"), |
| 137 | + Headers: []kcronsumer.Header{ |
| 138 | + { |
| 139 | + Key: "x-custom-client-header", |
| 140 | + Value: []byte("bar"), |
| 141 | + }, |
| 142 | + { |
| 143 | + Key: "x-error-message", |
| 144 | + Value: []byte("some error description"), |
| 145 | + }, |
| 146 | + }, |
| 147 | + } |
| 148 | + |
| 149 | + // When |
| 150 | + actual := message.toRetryableMessage("retry-topic", "consumeFn error") |
| 151 | + |
| 152 | + // Then |
| 153 | + if actual.Topic != expected.Topic { |
| 154 | + t.Errorf("topic must be %q", expected.Topic) |
| 155 | + } |
| 156 | + |
| 157 | + if !bytes.Equal(actual.Key, expected.Key) { |
| 158 | + t.Errorf("Key must be equal to %q", string(expected.Key)) |
| 159 | + } |
| 160 | + |
| 161 | + if !bytes.Equal(actual.Value, expected.Value) { |
| 162 | + t.Errorf("Value must be equal to %q", string(expected.Value)) |
| 163 | + } |
| 164 | + |
| 165 | + if len(actual.Headers) != 2 { |
| 166 | + t.Error("Header length must be equal to 2") |
| 167 | + } |
| 168 | + |
| 169 | + if actual.Headers[0].Key != expected.Headers[0].Key { |
| 170 | + t.Errorf("First Header key must be equal to %q", expected.Headers[0].Key) |
| 171 | + } |
| 172 | + |
| 173 | + if !bytes.Equal(actual.Headers[0].Value, expected.Headers[0].Value) { |
| 174 | + t.Errorf("First Header value must be equal to %q", expected.Headers[0].Value) |
| 175 | + } |
| 176 | + |
| 177 | + if actual.Headers[1].Key != expected.Headers[1].Key { |
| 178 | + t.Errorf("Second Header key must be equal to %q", expected.Headers[1].Key) |
| 179 | + } |
| 180 | + |
| 181 | + if !bytes.Equal(actual.Headers[1].Value, expected.Headers[1].Value) { |
| 182 | + t.Errorf("Second Header value must be equal to %q", expected.Headers[1].Value) |
| 183 | + } |
| 184 | + }) |
| 185 | + t.Run("When_error_description_does_not_exist", func(t *testing.T) { |
| 186 | + // Given |
| 187 | + message := Message{ |
| 188 | + Key: []byte("key"), |
| 189 | + Value: []byte("value"), |
| 190 | + Headers: []Header{ |
| 191 | + { |
| 192 | + Key: "x-custom-client-header", |
| 193 | + Value: []byte("bar"), |
| 194 | + }, |
| 195 | + }, |
| 196 | + } |
| 197 | + expected := kcronsumer.Message{ |
| 198 | + Topic: "retry-topic", |
| 199 | + Key: []byte("key"), |
| 200 | + Value: []byte("value"), |
| 201 | + Headers: []kcronsumer.Header{ |
| 202 | + { |
| 203 | + Key: "x-custom-client-header", |
| 204 | + Value: []byte("bar"), |
| 205 | + }, |
| 206 | + { |
| 207 | + Key: "x-error-message", |
| 208 | + Value: []byte("consumeFn error"), |
| 209 | + }, |
| 210 | + }, |
| 211 | + } |
| 212 | + |
| 213 | + // When |
| 214 | + actual := message.toRetryableMessage("retry-topic", "consumeFn error") |
| 215 | + |
| 216 | + // Then |
| 217 | + if actual.Topic != expected.Topic { |
| 218 | + t.Errorf("topic must be %q", expected.Topic) |
| 219 | + } |
| 220 | + |
| 221 | + if !bytes.Equal(actual.Key, expected.Key) { |
| 222 | + t.Errorf("Key must be equal to %q", string(expected.Key)) |
| 223 | + } |
| 224 | + |
| 225 | + if !bytes.Equal(actual.Value, expected.Value) { |
| 226 | + t.Errorf("Value must be equal to %q", string(expected.Value)) |
| 227 | + } |
| 228 | + |
| 229 | + if len(actual.Headers) != 2 { |
| 230 | + t.Error("Header length must be equal to 2") |
| 231 | + } |
| 232 | + |
| 233 | + if actual.Headers[0].Key != expected.Headers[0].Key { |
| 234 | + t.Errorf("First Header key must be equal to %q", expected.Headers[0].Key) |
| 235 | + } |
| 236 | + |
| 237 | + if !bytes.Equal(actual.Headers[0].Value, expected.Headers[0].Value) { |
| 238 | + t.Errorf("First Header value must be equal to %q", expected.Headers[0].Value) |
| 239 | + } |
| 240 | + |
| 241 | + if actual.Headers[1].Key != expected.Headers[1].Key { |
| 242 | + t.Errorf("Second Header key must be equal to %q", expected.Headers[1].Key) |
| 243 | + } |
| 244 | + |
| 245 | + if !bytes.Equal(actual.Headers[1].Value, expected.Headers[1].Value) { |
| 246 | + t.Errorf("Second Header value must be equal to %q", expected.Headers[1].Value) |
| 247 | + } |
| 248 | + }) |
| 249 | +} |
0 commit comments