Test fetchInitialOffset retry on ErrOffsetsLoadInProgress
(t *testing.T)
| 524 | |
| 525 | // Test fetchInitialOffset retry on ErrOffsetsLoadInProgress |
| 526 | func TestOffsetManagerFetchInitialLoadInProgress(t *testing.T) { |
| 527 | var retryCount atomic.Int32 |
| 528 | backoff := func(retries, maxRetries int) time.Duration { |
| 529 | retryCount.Add(1) |
| 530 | return 0 |
| 531 | } |
| 532 | om, testClient, broker, coordinator := initOffsetManagerWithBackoffFunc(t, 0, backoff, NewTestConfig()) |
| 533 | defer broker.Close() |
| 534 | defer coordinator.Close() |
| 535 | |
| 536 | // Error on first fetchInitialOffset call |
| 537 | responseBlock := OffsetFetchResponseBlock{ |
| 538 | Err: ErrOffsetsLoadInProgress, |
| 539 | Offset: 5, |
| 540 | Metadata: "test_meta", |
| 541 | } |
| 542 | |
| 543 | fetchResponse := new(OffsetFetchResponse) |
| 544 | fetchResponse.AddBlock("my_topic", 0, &responseBlock) |
| 545 | coordinator.Returns(fetchResponse) |
| 546 | |
| 547 | // Second fetchInitialOffset call is fine |
| 548 | fetchResponse2 := new(OffsetFetchResponse) |
| 549 | responseBlock2 := responseBlock |
| 550 | responseBlock2.Err = ErrNoError |
| 551 | fetchResponse2.AddBlock("my_topic", 0, &responseBlock2) |
| 552 | coordinator.Returns(fetchResponse2) |
| 553 | |
| 554 | pom, err := om.ManagePartition("my_topic", 0) |
| 555 | if err != nil { |
| 556 | t.Error(err) |
| 557 | } |
| 558 | |
| 559 | safeClose(t, pom) |
| 560 | safeClose(t, om) |
| 561 | safeClose(t, testClient) |
| 562 | |
| 563 | if retryCount.Load() == 0 { |
| 564 | t.Fatal("Expected at least one retry") |
| 565 | } |
| 566 | } |
| 567 | |
| 568 | // fetchInitialOffset must retry when OffsetFetchResponse v2+ surfaces a |
| 569 | // retriable coordinator error at the top level with no per-partition blocks |
nothing calls this directly
no test coverage detected