|
4 | 4 | import static org.mockito.ArgumentMatchers.eq; |
5 | 5 | import static org.mockito.Mockito.doReturn; |
6 | 6 | import static org.mockito.Mockito.mock; |
| 7 | +import static org.mockito.Mockito.reset; |
7 | 8 | import static org.mockito.Mockito.spy; |
8 | 9 | import static org.mockito.Mockito.times; |
9 | 10 | import static org.mockito.Mockito.verify; |
@@ -597,52 +598,32 @@ public void testGetAllTopicRetentionsThrowsException() throws ExecutionException |
597 | 598 | } |
598 | 599 |
|
599 | 600 | @Test |
600 | | - public void testSetTopicConfigWithUncleanLeaderElectionDisabled() { |
601 | | - PubSubTopicConfiguration topicConfiguration = |
602 | | - new PubSubTopicConfiguration(Optional.of(1111L), true, Optional.of(222), 333L, Optional.empty()); |
603 | | - topicConfiguration.setUncleanLeaderElectionEnable(Optional.of(false)); |
604 | | - |
605 | | - AlterConfigsResult alterConfigsResultMock = mock(AlterConfigsResult.class); |
606 | | - KafkaFuture<Void> alterConfigsKafkaFutureMock = mock(KafkaFuture.class); |
607 | | - |
608 | | - Map<ConfigResource, Config> resourceConfigMap = new HashMap<>(); |
609 | | - when(internalKafkaAdminClientMock.alterConfigs(any())).thenAnswer(invocation -> { |
610 | | - resourceConfigMap.putAll(invocation.getArgument(0)); |
611 | | - return alterConfigsResultMock; |
612 | | - }); |
613 | | - when(alterConfigsResultMock.all()).thenReturn(alterConfigsKafkaFutureMock); |
614 | | - |
615 | | - kafkaAdminAdapter.setTopicConfig(testPubSubTopic, topicConfiguration); |
616 | | - |
617 | | - assertEquals(resourceConfigMap.size(), 1); |
618 | | - for (Map.Entry<ConfigResource, Config> entry: resourceConfigMap.entrySet()) { |
619 | | - Config config = entry.getValue(); |
620 | | - assertEquals(config.get(TopicConfig.UNCLEAN_LEADER_ELECTION_ENABLE_CONFIG).value(), "false"); |
621 | | - } |
622 | | - } |
623 | | - |
624 | | - @Test |
625 | | - public void testSetTopicConfigWithUncleanLeaderElectionEnabled() { |
626 | | - PubSubTopicConfiguration topicConfiguration = |
627 | | - new PubSubTopicConfiguration(Optional.of(1111L), true, Optional.of(222), 333L, Optional.empty()); |
628 | | - topicConfiguration.setUncleanLeaderElectionEnable(Optional.of(true)); |
629 | | - |
630 | | - AlterConfigsResult alterConfigsResultMock = mock(AlterConfigsResult.class); |
631 | | - KafkaFuture<Void> alterConfigsKafkaFutureMock = mock(KafkaFuture.class); |
632 | | - |
633 | | - Map<ConfigResource, Config> resourceConfigMap = new HashMap<>(); |
634 | | - when(internalKafkaAdminClientMock.alterConfigs(any())).thenAnswer(invocation -> { |
635 | | - resourceConfigMap.putAll(invocation.getArgument(0)); |
636 | | - return alterConfigsResultMock; |
637 | | - }); |
638 | | - when(alterConfigsResultMock.all()).thenReturn(alterConfigsKafkaFutureMock); |
639 | | - |
640 | | - kafkaAdminAdapter.setTopicConfig(testPubSubTopic, topicConfiguration); |
641 | | - |
642 | | - assertEquals(resourceConfigMap.size(), 1); |
643 | | - for (Map.Entry<ConfigResource, Config> entry: resourceConfigMap.entrySet()) { |
644 | | - Config config = entry.getValue(); |
645 | | - assertEquals(config.get(TopicConfig.UNCLEAN_LEADER_ELECTION_ENABLE_CONFIG).value(), "true"); |
| 601 | + public void testSetTopicConfigWithUncleanLeaderElectionSet() { |
| 602 | + for (boolean enabled: new boolean[] { false, true }) { |
| 603 | + // Reset mock state between iterations |
| 604 | + reset(internalKafkaAdminClientMock); |
| 605 | + |
| 606 | + PubSubTopicConfiguration topicConfiguration = |
| 607 | + new PubSubTopicConfiguration(Optional.of(1111L), true, Optional.of(222), 333L, Optional.empty()); |
| 608 | + topicConfiguration.setUncleanLeaderElectionEnable(Optional.of(enabled)); |
| 609 | + |
| 610 | + AlterConfigsResult alterConfigsResultMock = mock(AlterConfigsResult.class); |
| 611 | + KafkaFuture<Void> alterConfigsKafkaFutureMock = mock(KafkaFuture.class); |
| 612 | + |
| 613 | + Map<ConfigResource, Config> resourceConfigMap = new HashMap<>(); |
| 614 | + when(internalKafkaAdminClientMock.alterConfigs(any())).thenAnswer(invocation -> { |
| 615 | + resourceConfigMap.putAll(invocation.getArgument(0)); |
| 616 | + return alterConfigsResultMock; |
| 617 | + }); |
| 618 | + when(alterConfigsResultMock.all()).thenReturn(alterConfigsKafkaFutureMock); |
| 619 | + |
| 620 | + kafkaAdminAdapter.setTopicConfig(testPubSubTopic, topicConfiguration); |
| 621 | + |
| 622 | + assertEquals(resourceConfigMap.size(), 1); |
| 623 | + for (Map.Entry<ConfigResource, Config> entry: resourceConfigMap.entrySet()) { |
| 624 | + Config config = entry.getValue(); |
| 625 | + assertEquals(config.get(TopicConfig.UNCLEAN_LEADER_ELECTION_ENABLE_CONFIG).value(), Boolean.toString(enabled)); |
| 626 | + } |
646 | 627 | } |
647 | 628 | } |
648 | 629 |
|
|
0 commit comments