From caed144a59f49288e7f35e5f889a673a2ce5db95 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Thu, 26 Jun 2025 18:55:00 +0100 Subject: [PATCH 1/5] KAFKA-19440: Handle top-level errors in AlterShareGroupOffsets RPC --- .../admin/AlterShareGroupOffsetsResult.java | 29 +++---- .../kafka/clients/admin/KafkaAdminClient.java | 2 +- .../AlterShareGroupOffsetsHandler.java | 85 ++++++++----------- .../AlterShareGroupOffsetsRequest.java | 31 ++++--- .../main/scala/kafka/server/KafkaApis.scala | 4 +- .../group/GroupCoordinatorService.java | 45 +++++----- 6 files changed, 89 insertions(+), 107 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java index 7c41852231d90..9cb7b1197ec8e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java @@ -20,12 +20,10 @@ import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; +import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.internals.KafkaFutureImpl; -import org.apache.kafka.common.protocol.Errors; -import java.util.List; import java.util.Map; -import java.util.stream.Collectors; /** * The result of the {@link Admin#alterShareGroupOffsets(String, Map, AlterShareGroupOffsetsOptions)} call. @@ -35,9 +33,9 @@ @InterfaceStability.Evolving public class AlterShareGroupOffsetsResult { - private final KafkaFuture> future; + private final KafkaFuture> future; - AlterShareGroupOffsetsResult(KafkaFuture> future) { + AlterShareGroupOffsetsResult(KafkaFuture> future) { this.future = future; } @@ -54,11 +52,11 @@ public KafkaFuture partitionResult(final TopicPartition partition) { result.completeExceptionally(new IllegalArgumentException( "Alter offset for partition \"" + partition + "\" was not attempted")); } else { - final Errors error = topicPartitions.get(partition); - if (error == Errors.NONE) { + final ApiException exception = topicPartitions.get(partition); + if (exception == null) { result.complete(null); } else { - result.completeExceptionally(error.exception()); + result.completeExceptionally(exception); } } }); @@ -70,20 +68,13 @@ public KafkaFuture partitionResult(final TopicPartition partition) { * Return a future which succeeds if all the alter offsets succeed. */ public KafkaFuture all() { - return this.future.thenApply(topicPartitionErrorsMap -> { - List partitionsFailed = topicPartitionErrorsMap.entrySet() - .stream() - .filter(e -> e.getValue() != Errors.NONE) - .map(Map.Entry::getKey) - .collect(Collectors.toList()); - for (Errors error : topicPartitionErrorsMap.values()) { - if (error != Errors.NONE) { - throw error.exception( - "Failed altering share group offsets for the following partitions: " + partitionsFailed); + return this.future.thenApply(topicPartitionErrorsMap -> { + for (ApiException exception : topicPartitionErrorsMap.values()) { + if (exception != null) { + throw exception; } } return null; }); } - } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index b283d65cbee06..9d16f6e540f6c 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3805,7 +3805,7 @@ public DescribeShareGroupsResult describeShareGroups(final Collection gr @Override public AlterShareGroupOffsetsResult alterShareGroupOffsets(String groupId, Map offsets, AlterShareGroupOffsetsOptions options) { - SimpleAdminApiFuture> future = AlterShareGroupOffsetsHandler.newFuture(groupId); + SimpleAdminApiFuture> future = AlterShareGroupOffsetsHandler.newFuture(groupId); AlterShareGroupOffsetsHandler handler = new AlterShareGroupOffsetsHandler(groupId, offsets, logContext); invokeDriver(handler, future, options.timeoutMs); return new AlterShareGroupOffsetsResult(future.get(CoordinatorKey.byGroupId(groupId))); diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java b/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java index f66f597283630..9bc5dfd2e593e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java @@ -21,8 +21,8 @@ import org.apache.kafka.clients.admin.KafkaAdminClient; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.message.AlterShareGroupOffsetsRequestData; -import org.apache.kafka.common.message.AlterShareGroupOffsetsResponseData; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.requests.AbstractResponse; import org.apache.kafka.common.requests.AlterShareGroupOffsetsRequest; @@ -33,7 +33,6 @@ import org.slf4j.Logger; import java.util.ArrayList; -import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Map; @@ -42,7 +41,7 @@ /** * This class is the handler for {@link KafkaAdminClient#alterShareGroupOffsets(String, Map, AlterShareGroupOffsetsOptions)} call */ -public class AlterShareGroupOffsetsHandler extends AdminApiHandler.Batched> { +public class AlterShareGroupOffsetsHandler extends AdminApiHandler.Batched> { private final CoordinatorKey groupId; @@ -60,8 +59,8 @@ public AlterShareGroupOffsetsHandler(String groupId, Map o this.lookupStrategy = new CoordinatorStrategy(FindCoordinatorRequest.CoordinatorType.GROUP, logContext); } - public static AdminApiFuture.SimpleAdminApiFuture> newFuture(String groupId) { - return AdminApiFuture.forKeys(Collections.singleton(CoordinatorKey.byGroupId(groupId))); + public static AdminApiFuture.SimpleAdminApiFuture> newFuture(String groupId) { + return AdminApiFuture.forKeys(Set.of(CoordinatorKey.byGroupId(groupId))); } @Override @@ -87,57 +86,49 @@ public String apiName() { } @Override - public ApiResult> handleResponse(Node broker, Set keys, AbstractResponse abstractResponse) { + public ApiResult> handleResponse(Node broker, Set keys, AbstractResponse abstractResponse) { AlterShareGroupOffsetsResponse response = (AlterShareGroupOffsetsResponse) abstractResponse; - final Map partitionResults = new HashMap<>(); - final Set groupsToUnmap = new HashSet<>(); - final Set groupsToRetry = new HashSet<>(); - - for (AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopic topic : response.data().responses()) { - for (AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponsePartition partition : topic.partitions()) { - TopicPartition topicPartition = new TopicPartition(topic.topicName(), partition.partitionIndex()); - Errors error = Errors.forCode(partition.errorCode()); - - if (error != Errors.NONE) { - handleError( - groupId, - topicPartition, - error, - partitionResults, - groupsToUnmap, - groupsToRetry - ); - } else { - partitionResults.put(topicPartition, error); + + final Errors topLevelError = Errors.forCode(response.data().errorCode()); + final String topLevelErrorMessage = response.data().errorMessage(); + + if (topLevelError != Errors.NONE) { + final Set groupsToUnmap = new HashSet<>(); + final Map groupsFailed = new HashMap<>(); + handleGroupError(groupId, topLevelError, topLevelErrorMessage, groupsToUnmap, groupsFailed); + + return new ApiResult<>(Map.of(), groupsFailed, new ArrayList<>(groupsToUnmap)); + } else { + final Map partitionResults = new HashMap<>(); + response.data().responses().forEach(topic -> topic.partitions().forEach(partition -> { + if (partition.errorCode() != Errors.NONE.code()) { + final Errors partitionError = Errors.forCode(partition.errorCode()); + final String partitionErrorMessage = partition.errorMessage(); + log.debug("AlterShareGroupOffsets request for group id {} and topic-partition {}-{} failed and returned error {}." + partitionErrorMessage, + groupId.idValue, topic.topicName(), partition.partitionIndex(), partitionError); } - } - } + partitionResults.put(new TopicPartition(topic.topicName(), partition.partitionIndex()), Errors.forCode(partition.errorCode()).exception(partition.errorMessage())); + })); - if (groupsToUnmap.isEmpty() && groupsToRetry.isEmpty()) { return ApiResult.completed(groupId, partitionResults); - } else { - return ApiResult.unmapped(new ArrayList<>(groupsToUnmap)); } } - private void handleError( - CoordinatorKey groupId, - TopicPartition topicPartition, - Errors error, - Map partitionResults, - Set groupsToUnmap, - Set groupsToRetry + private void handleGroupError( + CoordinatorKey groupId, + Errors error, + String errorMessage, + Set groupsToUnmap, + Map groupsFailed ) { switch (error) { case COORDINATOR_LOAD_IN_PROGRESS: case REBALANCE_IN_PROGRESS: - log.debug("AlterShareGroupOffsets request for group id {} returned error {}. Will retry.", - groupId.idValue, error); - groupsToRetry.add(groupId); + log.debug("AlterShareGroupOffsets request for group id {} returned error {}. Will retry." + errorMessage, groupId.idValue, error); break; case COORDINATOR_NOT_AVAILABLE: case NOT_COORDINATOR: - log.debug("AlterShareGroupOffsets request for group id {} returned error {}. Will rediscover the coordinator and retry.", + log.debug("AlterShareGroupOffsets request for group id {} returned error {}. Will rediscover the coordinator and retry." + errorMessage, groupId.idValue, error); groupsToUnmap.add(groupId); break; @@ -147,14 +138,12 @@ private void handleError( case UNKNOWN_SERVER_ERROR: case KAFKA_STORAGE_ERROR: case GROUP_AUTHORIZATION_FAILED: - log.debug("AlterShareGroupOffsets request for group id {} and partition {} failed due" + - " to error {}.", groupId.idValue, topicPartition, error); - partitionResults.put(topicPartition, error); + log.debug("AlterShareGroupOffsets request for group id {} failed due to error {}." + errorMessage, groupId.idValue, error); + groupsFailed.put(groupId, error.exception(errorMessage)); break; default: - log.error("AlterShareGroupOffsets request for group id {} and partition {} failed due" + - " to unexpected error {}.", groupId.idValue, topicPartition, error); - partitionResults.put(topicPartition, error); + log.error("AlterShareGroupOffsets request for group id {} failed due to unexpected error {}." + errorMessage, groupId.idValue, error); + groupsFailed.put(groupId, error.exception(errorMessage)); } } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java index 2eb9e37bc507c..fe58ee1a3f9fe 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java @@ -53,26 +53,31 @@ public String toString() { } @Override - public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) { - Errors error = Errors.forException(e); - return new AlterShareGroupOffsetsResponse(getErrorResponse(throttleTimeMs, error)); + public AlterShareGroupOffsetsResponse getErrorResponse(int throttleTimeMs, Throwable e) { + return getErrorResponse(throttleTimeMs, Errors.forException(e)); } - public static AlterShareGroupOffsetsResponseData getErrorResponse(int throttleTimeMs, Errors error) { - return new AlterShareGroupOffsetsResponseData() - .setThrottleTimeMs(throttleTimeMs) - .setErrorCode(error.code()) - .setErrorMessage(error.message()); + public AlterShareGroupOffsetsResponse getErrorResponse(int throttleTimeMs, Errors error) { + return getErrorResponse(throttleTimeMs, error.code(), error.message()); + } + + public AlterShareGroupOffsetsResponse getErrorResponse(int throttleTimeMs, short errorCode, String message) { + return new AlterShareGroupOffsetsResponse( + new AlterShareGroupOffsetsResponseData() + .setThrottleTimeMs(throttleTimeMs) + .setErrorCode(errorCode) + .setErrorMessage(message) + ); } - public static AlterShareGroupOffsetsResponseData getErrorResponse(Errors error) { - return getErrorResponse(error.code(), error.message()); + public static AlterShareGroupOffsetsResponseData getErrorResponseData(Errors error) { + return getErrorResponseData(error.code(), error.message()); } - public static AlterShareGroupOffsetsResponseData getErrorResponse(short errorCode, String errorMessage) { + public static AlterShareGroupOffsetsResponseData getErrorResponseData(short errorCode, String errorMessage) { return new AlterShareGroupOffsetsResponseData() - .setErrorCode(errorCode) - .setErrorMessage(errorMessage); + .setErrorCode(errorCode) + .setErrorMessage(errorMessage); } public static AlterShareGroupOffsetsRequest parse(Readable readable, short version) { diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 5eb249c54d6b0..443eccdf56427 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -3756,7 +3756,7 @@ class KafkaApis(val requestChannel: RequestChannel, val groupId = alterShareGroupOffsetsRequest.data.groupId if (!isShareGroupProtocolEnabled) { - requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(AbstractResponse.DEFAULT_THROTTLE_TIME, Errors.UNSUPPORTED_VERSION.exception)) + requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(Errors.UNSUPPORTED_VERSION.exception)) return CompletableFuture.completedFuture[Unit](()) } else if (!authHelper.authorize(request.context, READ, GROUP, groupId)) { requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(Errors.GROUP_AUTHORIZATION_FAILED.exception)) @@ -3792,6 +3792,8 @@ class KafkaApis(val requestChannel: RequestChannel, ).handle[Unit] { (response, exception) => if (exception != null) { requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(exception)) + } else if (response.errorCode() != Errors.NONE.code) { + requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(AbstractResponse.DEFAULT_THROTTLE_TIME, response.errorCode(), response.errorMessage())) } else { requestHelper.sendMaybeThrottle(request, responseBuilder.merge(response, metadataCache.topicNamesToIds()).build()) } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java index 099201ecabb82..4ca62c77a46d9 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java @@ -1233,13 +1233,17 @@ public CompletableFuture> shareGroupDescribe( * See {@link GroupCoordinator#alterShareGroupOffsets(AuthorizableRequestContext, String, AlterShareGroupOffsetsRequestData)}. */ @Override - public CompletableFuture alterShareGroupOffsets(AuthorizableRequestContext context, String groupId, AlterShareGroupOffsetsRequestData request) { + public CompletableFuture alterShareGroupOffsets( + AuthorizableRequestContext context, + String groupId, + AlterShareGroupOffsetsRequestData request + ) { if (!isActive.get() || metadataImage == null) { - return CompletableFuture.completedFuture(AlterShareGroupOffsetsRequest.getErrorResponse(Errors.COORDINATOR_NOT_AVAILABLE)); + return CompletableFuture.completedFuture(AlterShareGroupOffsetsRequest.getErrorResponseData(Errors.COORDINATOR_NOT_AVAILABLE)); } if (groupId == null || groupId.isEmpty()) { - return CompletableFuture.completedFuture(AlterShareGroupOffsetsRequest.getErrorResponse(Errors.INVALID_GROUP_ID)); + return CompletableFuture.completedFuture(AlterShareGroupOffsetsRequest.getErrorResponseData(Errors.INVALID_GROUP_ID)); } if (request.topics() == null || request.topics().isEmpty()) { @@ -1257,7 +1261,7 @@ public CompletableFuture alterShareGroupOffs "share-group-offsets-alter", request, exception, - (error, message) -> AlterShareGroupOffsetsRequest.getErrorResponse(error), + (error, __) -> AlterShareGroupOffsetsRequest.getErrorResponseData(error), log )); } @@ -1822,26 +1826,18 @@ public CompletableFuture deleteShareGroupOf AuthorizableRequestContext context, DeleteShareGroupOffsetsRequestData requestData ) { - if (!isActive.get()) { - return CompletableFuture.completedFuture( - DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(Errors.COORDINATOR_NOT_AVAILABLE)); - } - - if (metadataImage == null) { - return CompletableFuture.completedFuture( - DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(Errors.COORDINATOR_NOT_AVAILABLE)); + if (!isActive.get() || metadataImage == null) { + return CompletableFuture.completedFuture(DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(Errors.COORDINATOR_NOT_AVAILABLE)); } String groupId = requestData.groupId(); if (!isGroupIdNotEmpty(groupId)) { - return CompletableFuture.completedFuture( - DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(Errors.INVALID_GROUP_ID)); + return CompletableFuture.completedFuture(DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(Errors.INVALID_GROUP_ID)); } if (requestData.topics() == null || requestData.topics().isEmpty()) { - return CompletableFuture.completedFuture( - new DeleteShareGroupOffsetsResponseData() + return CompletableFuture.completedFuture(new DeleteShareGroupOffsetsResponseData() ); } @@ -1850,15 +1846,14 @@ public CompletableFuture deleteShareGroupOf topicPartitionFor(groupId), Duration.ofMillis(config.offsetCommitTimeoutMs()), coordinator -> coordinator.initiateDeleteShareGroupOffsets(groupId, requestData) - ) - .thenCompose(resultHolder -> deleteShareGroupOffsetsState(groupId, resultHolder)) - .exceptionally(exception -> handleOperationException( - "initiate-delete-share-group-offsets", - groupId, - exception, - (error, __) -> DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(error), - log - )); + ).thenCompose(resultHolder -> deleteShareGroupOffsetsState(groupId, resultHolder) + ).exceptionally(exception -> handleOperationException( + "initiate-delete-share-group-offsets", + groupId, + exception, + (error, __) -> DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(error), + log + )); } private CompletableFuture deleteShareGroupOffsetsState( From e13607befd1e21b98f6bf5cd2b23451b61aae271 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Fri, 27 Jun 2025 23:03:40 +0100 Subject: [PATCH 2/5] Improve tests --- .../admin/AlterShareGroupOffsetsResult.java | 14 ++++- .../kafka/clients/admin/KafkaAdminClient.java | 8 ++- .../AlterShareGroupOffsetsHandler.java | 54 +++++++++++++------ .../clients/admin/KafkaAdminClientTest.java | 30 +++++++++-- 4 files changed, 82 insertions(+), 24 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java index 9cb7b1197ec8e..293daaadbb925 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/AlterShareGroupOffsetsResult.java @@ -22,8 +22,11 @@ import org.apache.kafka.common.annotation.InterfaceStability; import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.internals.KafkaFutureImpl; +import org.apache.kafka.common.protocol.Errors; +import java.util.List; import java.util.Map; +import java.util.stream.Collectors; /** * The result of the {@link Admin#alterShareGroupOffsets(String, Map, AlterShareGroupOffsetsOptions)} call. @@ -66,12 +69,19 @@ public KafkaFuture partitionResult(final TopicPartition partition) { /** * Return a future which succeeds if all the alter offsets succeed. + * If not, the first topic error shall be returned. */ public KafkaFuture all() { - return this.future.thenApply(topicPartitionErrorsMap -> { + return this.future.thenApply(topicPartitionErrorsMap -> { + List partitionsFailed = topicPartitionErrorsMap.entrySet() + .stream() + .filter(e -> e.getValue() != null) + .map(Map.Entry::getKey) + .collect(Collectors.toList()); for (ApiException exception : topicPartitionErrorsMap.values()) { if (exception != null) { - throw exception; + throw Errors.forException(exception).exception( + "Failed altering group offsets for the following partitions: " + partitionsFailed); } } return null; diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 9d16f6e540f6c..e1be4304950a3 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3804,7 +3804,9 @@ public DescribeShareGroupsResult describeShareGroups(final Collection gr } @Override - public AlterShareGroupOffsetsResult alterShareGroupOffsets(String groupId, Map offsets, AlterShareGroupOffsetsOptions options) { + public AlterShareGroupOffsetsResult alterShareGroupOffsets(final String groupId, + final Map offsets, + final AlterShareGroupOffsetsOptions options) { SimpleAdminApiFuture> future = AlterShareGroupOffsetsHandler.newFuture(groupId); AlterShareGroupOffsetsHandler handler = new AlterShareGroupOffsetsHandler(groupId, offsets, logContext); invokeDriver(handler, future, options.timeoutMs); @@ -3821,7 +3823,9 @@ public ListShareGroupOffsetsResult listShareGroupOffsets(final Map topics, DeleteShareGroupOffsetsOptions options) { + public DeleteShareGroupOffsetsResult deleteShareGroupOffsets(final String groupId, + final Set topics, + final DeleteShareGroupOffsetsOptions options) { SimpleAdminApiFuture> future = DeleteShareGroupOffsetsHandler.newFuture(groupId); DeleteShareGroupOffsetsHandler handler = new DeleteShareGroupOffsetsHandler(groupId, topics, logContext); invokeDriver(handler, future, options.timeoutMs); diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java b/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java index 9bc5dfd2e593e..ef21be6b6d2b0 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/internals/AlterShareGroupOffsetsHandler.java @@ -51,7 +51,6 @@ public class AlterShareGroupOffsetsHandler extends AdminApiHandler.Batched offsets, LogContext logContext) { this.groupId = CoordinatorKey.byGroupId(groupId); this.offsets = offsets; @@ -63,6 +62,13 @@ public static AdminApiFuture.SimpleAdminApiFuture groupIds) { + if (!groupIds.equals(Set.of(groupId))) { + throw new IllegalArgumentException("Received unexpected group ids " + groupIds + + " (expected only " + Set.of(groupId) + ")"); + } + } + @Override AlterShareGroupOffsetsRequest.Builder buildBatchedRequest(int brokerId, Set groupIds) { var data = new AlterShareGroupOffsetsRequestData().setGroupId(groupId.idValue); @@ -87,19 +93,28 @@ public String apiName() { @Override public ApiResult> handleResponse(Node broker, Set keys, AbstractResponse abstractResponse) { - AlterShareGroupOffsetsResponse response = (AlterShareGroupOffsetsResponse) abstractResponse; + validateKeys(keys); - final Errors topLevelError = Errors.forCode(response.data().errorCode()); - final String topLevelErrorMessage = response.data().errorMessage(); - - if (topLevelError != Errors.NONE) { - final Set groupsToUnmap = new HashSet<>(); - final Map groupsFailed = new HashMap<>(); - handleGroupError(groupId, topLevelError, topLevelErrorMessage, groupsToUnmap, groupsFailed); - - return new ApiResult<>(Map.of(), groupsFailed, new ArrayList<>(groupsToUnmap)); + AlterShareGroupOffsetsResponse response = (AlterShareGroupOffsetsResponse) abstractResponse; + final Set groupsToUnmap = new HashSet<>(); + final Set groupsToRetry = new HashSet<>(); + final Map partitionResults = new HashMap<>(); + + if (response.data().errorCode() != Errors.NONE.code()) { + final Errors topLevelError = Errors.forCode(response.data().errorCode()); + final String topLevelErrorMessage = response.data().errorMessage(); + + offsets.forEach((topicPartition, offset) -> + handleError( + groupId, + topicPartition, + topLevelError, + topLevelErrorMessage, + partitionResults, + groupsToUnmap, + groupsToRetry + )); } else { - final Map partitionResults = new HashMap<>(); response.data().responses().forEach(topic -> topic.partitions().forEach(partition -> { if (partition.errorCode() != Errors.NONE.code()) { final Errors partitionError = Errors.forCode(partition.errorCode()); @@ -109,22 +124,29 @@ public ApiResult> handleRespon } partitionResults.put(new TopicPartition(topic.topicName(), partition.partitionIndex()), Errors.forCode(partition.errorCode()).exception(partition.errorMessage())); })); + } + if (groupsToUnmap.isEmpty() && groupsToRetry.isEmpty()) { return ApiResult.completed(groupId, partitionResults); + } else { + return ApiResult.unmapped(new ArrayList<>(groupsToUnmap)); } } - private void handleGroupError( + private void handleError( CoordinatorKey groupId, + TopicPartition topicPartition, Errors error, String errorMessage, + Map partitionResults, Set groupsToUnmap, - Map groupsFailed + Set groupsToRetry ) { switch (error) { case COORDINATOR_LOAD_IN_PROGRESS: case REBALANCE_IN_PROGRESS: log.debug("AlterShareGroupOffsets request for group id {} returned error {}. Will retry." + errorMessage, groupId.idValue, error); + groupsToRetry.add(groupId); break; case COORDINATOR_NOT_AVAILABLE: case NOT_COORDINATOR: @@ -139,11 +161,11 @@ private void handleGroupError( case KAFKA_STORAGE_ERROR: case GROUP_AUTHORIZATION_FAILED: log.debug("AlterShareGroupOffsets request for group id {} failed due to error {}." + errorMessage, groupId.idValue, error); - groupsFailed.put(groupId, error.exception(errorMessage)); + partitionResults.put(topicPartition, error.exception(errorMessage)); break; default: log.error("AlterShareGroupOffsets request for group id {} failed due to unexpected error {}." + errorMessage, groupId.idValue, error); - groupsFailed.put(groupId, error.exception(errorMessage)); + partitionResults.put(topicPartition, error.exception(errorMessage)); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 1d516cf66483c..1098078b582fc 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -56,7 +56,6 @@ import org.apache.kafka.common.errors.DuplicateVoterException; import org.apache.kafka.common.errors.FencedInstanceIdException; import org.apache.kafka.common.errors.GroupAuthorizationException; -import org.apache.kafka.common.errors.GroupNotEmptyException; import org.apache.kafka.common.errors.GroupSubscribedToTopicException; import org.apache.kafka.common.errors.InvalidConfigurationException; import org.apache.kafka.common.errors.InvalidReplicaAssignmentException; @@ -11351,6 +11350,28 @@ public void testAlterShareGroupOffsets() throws Exception { } } + @Test + public void testAlterShareGroupOffsetsWithTopLevelError() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + + AlterShareGroupOffsetsResponseData data = new AlterShareGroupOffsetsResponseData().setErrorCode(Errors.GROUP_AUTHORIZATION_FAILED.code()).setErrorMessage("Group authorization failed."); + + TopicPartition fooTopicPartition0 = new TopicPartition("foo", 0); + TopicPartition fooTopicPartition1 = new TopicPartition("foo", 1); + TopicPartition barPartition0 = new TopicPartition("bar", 0); + TopicPartition zooTopicPartition0 = new TopicPartition("zoo", 0); + + env.kafkaClient().prepareResponse(new AlterShareGroupOffsetsResponse(data)); + final AlterShareGroupOffsetsResult result = env.adminClient().alterShareGroupOffsets(GROUP_ID, Map.of(fooTopicPartition0, 1L, fooTopicPartition1, 2L, barPartition0, 1L)); + + TestUtils.assertFutureThrows(GroupAuthorizationException.class, result.all()); + TestUtils.assertFutureThrows(GroupAuthorizationException.class, result.partitionResult(fooTopicPartition1)); + TestUtils.assertFutureThrows(IllegalArgumentException.class, result.partitionResult(zooTopicPartition0)); + } + } + @Test public void testAlterShareGroupOffsetsWithErrorInOnePartition() throws Exception { try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { @@ -11359,7 +11380,8 @@ public void testAlterShareGroupOffsetsWithErrorInOnePartition() throws Exception AlterShareGroupOffsetsResponseData data = new AlterShareGroupOffsetsResponseData().setResponses( new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopicCollection(List.of( - new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopic().setTopicName("foo").setPartitions(List.of(new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponsePartition().setPartitionIndex(0), new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponsePartition().setPartitionIndex(1).setErrorCode(Errors.NON_EMPTY_GROUP.code()).setErrorMessage("The group is not empty"))), + new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopic().setTopicName("foo").setPartitions(List.of(new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponsePartition().setPartitionIndex(0), + new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponsePartition().setPartitionIndex(1).setErrorCode(Errors.TOPIC_AUTHORIZATION_FAILED.code()).setErrorMessage("Topic authorization failed."))), new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopic().setTopicName("bar").setPartitions(List.of(new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponsePartition().setPartitionIndex(0))) ).iterator()) ); @@ -11371,9 +11393,9 @@ public void testAlterShareGroupOffsetsWithErrorInOnePartition() throws Exception env.kafkaClient().prepareResponse(new AlterShareGroupOffsetsResponse(data)); final AlterShareGroupOffsetsResult result = env.adminClient().alterShareGroupOffsets(GROUP_ID, Map.of(fooTopicPartition0, 1L, fooTopicPartition1, 2L, barPartition0, 1L)); - TestUtils.assertFutureThrows(GroupNotEmptyException.class, result.all()); + TestUtils.assertFutureThrows(TopicAuthorizationException.class, result.all()); assertNull(result.partitionResult(fooTopicPartition0).get()); - TestUtils.assertFutureThrows(GroupNotEmptyException.class, result.partitionResult(fooTopicPartition1)); + TestUtils.assertFutureThrows(TopicAuthorizationException.class, result.partitionResult(fooTopicPartition1)); assertNull(result.partitionResult(barPartition0).get()); } } From b9316ddc0669679e3a5443f9f355bdef7a30cb87 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Tue, 1 Jul 2025 18:40:19 +0100 Subject: [PATCH 3/5] Tighten up code --- .../AlterShareGroupOffsetsRequest.java | 8 +-- .../DeleteShareGroupOffsetsRequest.java | 10 ++- .../main/scala/kafka/server/KafkaApis.scala | 11 +--- .../unit/kafka/server/KafkaApisTest.scala | 12 +++- .../group/GroupCoordinatorService.java | 12 ++-- .../group/GroupCoordinatorServiceTest.java | 61 ++++++++++++++++--- 6 files changed, 84 insertions(+), 30 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java index fe58ee1a3f9fe..be04568e1a395 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/AlterShareGroupOffsetsRequest.java @@ -71,13 +71,13 @@ public AlterShareGroupOffsetsResponse getErrorResponse(int throttleTimeMs, short } public static AlterShareGroupOffsetsResponseData getErrorResponseData(Errors error) { - return getErrorResponseData(error.code(), error.message()); + return getErrorResponseData(error, null); } - public static AlterShareGroupOffsetsResponseData getErrorResponseData(short errorCode, String errorMessage) { + public static AlterShareGroupOffsetsResponseData getErrorResponseData(Errors error, String errorMessage) { return new AlterShareGroupOffsetsResponseData() - .setErrorCode(errorCode) - .setErrorMessage(errorMessage); + .setErrorCode(error.code()) + .setErrorMessage(errorMessage == null ? error.message() : errorMessage); } public static AlterShareGroupOffsetsRequest parse(Readable readable, short version) { diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java index bec0077b9b3c5..c946f28723a91 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java @@ -80,12 +80,18 @@ public static DeleteShareGroupOffsetsRequest parse(Readable readable, short vers } public static DeleteShareGroupOffsetsResponseData getErrorDeleteResponseData(Errors error) { - return getErrorDeleteResponseData(error.code(), error.message()); + return getErrorDeleteResponseData(error, null); } public static DeleteShareGroupOffsetsResponseData getErrorDeleteResponseData(short errorCode, String errorMessage) { return new DeleteShareGroupOffsetsResponseData() .setErrorCode(errorCode) - .setErrorMessage(errorMessage); + .setErrorMessage(errorMessage == null ? Errors.forCode(errorCode).message(): errorMessage); + } + + public static DeleteShareGroupOffsetsResponseData getErrorDeleteResponseData(Errors error, String errorMessage) { + return new DeleteShareGroupOffsetsResponseData() + .setErrorCode(error.code()) + .setErrorMessage(errorMessage == null ? error.message() : errorMessage); } } \ No newline at end of file diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 443eccdf56427..5574d09d688e6 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -3778,7 +3778,7 @@ class KafkaApis(val requestChannel: RequestChannel, case Some(error) => topic.partitions().forEach(partition => responseBuilder.addPartition(topic.topicName(), partition.partitionIndex(), metadataCache.topicNamesToIds(), error.error)) case None => - authorizedTopicPartitions.add(topic) + authorizedTopicPartitions.add(topic.duplicate) } }) @@ -3833,15 +3833,6 @@ class KafkaApis(val requestChannel: RequestChannel, } } - if (authorizedTopics.isEmpty) { - requestHelper.sendMaybeThrottle( - request, - new DeleteShareGroupOffsetsResponse( - new DeleteShareGroupOffsetsResponseData() - .setResponses(deleteShareGroupOffsetsResponseTopics))) - return CompletableFuture.completedFuture[Unit](()) - } - groupCoordinator.deleteShareGroupOffsets( request.context, new DeleteShareGroupOffsetsRequestData().setGroupId(groupId).setTopics(authorizedTopics) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index a6c2658963515..e3a396ba9d52e 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -12863,10 +12863,18 @@ class KafkaApisTest extends Logging { def testDeleteShareGroupOffsetsRequestEmptyTopicsSuccess(): Unit = { metadataCache = initializeMetadataCacheWithShareGroupsEnabled() - val deleteShareGroupOffsetsRequest = new DeleteShareGroupOffsetsRequestData() + val deleteShareGroupOffsetsRequestData = new DeleteShareGroupOffsetsRequestData() .setGroupId("group") - val requestChannelRequest = buildRequest(new DeleteShareGroupOffsetsRequest.Builder(deleteShareGroupOffsetsRequest).build) + val requestChannelRequest = buildRequest(new DeleteShareGroupOffsetsRequest.Builder(deleteShareGroupOffsetsRequestData).build) + + val groupCoordinatorResponse: DeleteShareGroupOffsetsResponseData = new DeleteShareGroupOffsetsResponseData() + .setErrorCode(Errors.NONE.code()) + + when(groupCoordinator.deleteShareGroupOffsets( + requestChannelRequest.context, + deleteShareGroupOffsetsRequestData + )).thenReturn(CompletableFuture.completedFuture(groupCoordinatorResponse)) val resultFuture = new CompletableFuture[DeleteShareGroupOffsetsResponseData] kafkaApis = createKafkaApis() diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java index 4ca62c77a46d9..812853a263f6c 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java @@ -698,6 +698,10 @@ CompletableFuture persisterInitialize( InitializeShareGroupStateParameters request, AlterShareGroupOffsetsResponseData response ) { + if (request.groupTopicPartitionData().topicsData().isEmpty()) { + return CompletableFuture.completedFuture(response); + } + return persister.initializeState(request) .handle((result, exp) -> { if (exp == null) { @@ -1246,7 +1250,7 @@ public CompletableFuture alterShareGroupOffs return CompletableFuture.completedFuture(AlterShareGroupOffsetsRequest.getErrorResponseData(Errors.INVALID_GROUP_ID)); } - if (request.topics() == null || request.topics().isEmpty()) { + if (request.topics() == null) { return CompletableFuture.completedFuture(new AlterShareGroupOffsetsResponseData()); } @@ -1261,7 +1265,7 @@ public CompletableFuture alterShareGroupOffs "share-group-offsets-alter", request, exception, - (error, __) -> AlterShareGroupOffsetsRequest.getErrorResponseData(error), + (error, message) -> AlterShareGroupOffsetsRequest.getErrorResponseData(error, message), log )); } @@ -1836,7 +1840,7 @@ public CompletableFuture deleteShareGroupOf return CompletableFuture.completedFuture(DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(Errors.INVALID_GROUP_ID)); } - if (requestData.topics() == null || requestData.topics().isEmpty()) { + if (requestData.topics() == null) { return CompletableFuture.completedFuture(new DeleteShareGroupOffsetsResponseData() ); } @@ -1851,7 +1855,7 @@ public CompletableFuture deleteShareGroupOf "initiate-delete-share-group-offsets", groupId, exception, - (error, __) -> DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(error), + (error, message) -> DeleteShareGroupOffsetsRequest.getErrorDeleteResponseData(error, message), log )); } diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java index 01c87696053a2..4517d0cb8f372 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java @@ -4280,10 +4280,26 @@ public void testDeleteShareGroupOffsetsEmptyRequest() throws InterruptedExceptio service.startup(() -> 1); DeleteShareGroupOffsetsRequestData requestData = new DeleteShareGroupOffsetsRequestData() - .setGroupId("share-group-id"); + .setGroupId("share-group-id") + .setTopics(List.of()); DeleteShareGroupOffsetsResponseData responseData = new DeleteShareGroupOffsetsResponseData(); + GroupCoordinatorShard.DeleteShareGroupOffsetsResultHolder deleteShareGroupOffsetsResultHolder = + new GroupCoordinatorShard.DeleteShareGroupOffsetsResultHolder( + Errors.NONE.code(), + null, + null, + null + ); + + when(runtime.scheduleWriteOperation( + ArgumentMatchers.eq("initiate-delete-share-group-offsets"), + ArgumentMatchers.eq(new TopicPartition(Topic.GROUP_METADATA_TOPIC_NAME, 0)), + ArgumentMatchers.eq(Duration.ofMillis(5000)), + ArgumentMatchers.any() + )).thenReturn(CompletableFuture.completedFuture(deleteShareGroupOffsetsResultHolder)); + CompletableFuture future = service.deleteShareGroupOffsets(requestContext(ApiKeys.DELETE_SHARE_GROUP_OFFSETS), requestData); @@ -4291,7 +4307,7 @@ public void testDeleteShareGroupOffsetsEmptyRequest() throws InterruptedExceptio } @Test - public void testDeleteShareGroupOffsetsNullTopicsInRequest() throws InterruptedException, ExecutionException { + public void testDeleteShareGroupOffsetsEmptyTopicsInRequest() throws InterruptedException, ExecutionException { CoordinatorRuntime runtime = mockRuntime(); Persister persister = mock(DefaultStatePersister.class); GroupCoordinatorService service = new GroupCoordinatorServiceBuilder() @@ -4303,10 +4319,25 @@ public void testDeleteShareGroupOffsetsNullTopicsInRequest() throws InterruptedE DeleteShareGroupOffsetsRequestData requestData = new DeleteShareGroupOffsetsRequestData() .setGroupId("share-group-id") - .setTopics(null); + .setTopics(List.of()); DeleteShareGroupOffsetsResponseData responseData = new DeleteShareGroupOffsetsResponseData(); + GroupCoordinatorShard.DeleteShareGroupOffsetsResultHolder deleteShareGroupOffsetsResultHolder = + new GroupCoordinatorShard.DeleteShareGroupOffsetsResultHolder( + Errors.NONE.code(), + null, + null, + null + ); + + when(runtime.scheduleWriteOperation( + ArgumentMatchers.eq("initiate-delete-share-group-offsets"), + ArgumentMatchers.eq(new TopicPartition(Topic.GROUP_METADATA_TOPIC_NAME, 0)), + ArgumentMatchers.eq(Duration.ofMillis(5000)), + ArgumentMatchers.any() + )).thenReturn(CompletableFuture.completedFuture(deleteShareGroupOffsetsResultHolder)); + CompletableFuture future = service.deleteShareGroupOffsets(requestContext(ApiKeys.DELETE_SHARE_GROUP_OFFSETS), requestData); @@ -4400,9 +4431,7 @@ public void testDeleteShareGroupOffsetsRequestReturnsGroupIdNotFound() throws In DeleteShareGroupOffsetsRequestData requestData = new DeleteShareGroupOffsetsRequestData() .setGroupId(groupId) - .setTopics(List.of(new DeleteShareGroupOffsetsRequestData.DeleteShareGroupOffsetsRequestTopic() - .setTopicName(TOPIC_NAME) - )); + .setTopics(List.of()); DeleteShareGroupOffsetsResponseData responseData = new DeleteShareGroupOffsetsResponseData() .setErrorCode(Errors.GROUP_ID_NOT_FOUND.code()) @@ -5376,13 +5405,29 @@ public void testAlterShareGroupOffsetsEmptyRequest() throws ExecutionException, AlterShareGroupOffsetsRequestData request = new AlterShareGroupOffsetsRequestData() .setGroupId(groupId); + AlterShareGroupOffsetsResponseData data = new AlterShareGroupOffsetsResponseData(); + + Map.Entry alterShareGroupOffsetsIntermediate = + Map.entry( + new AlterShareGroupOffsetsResponseData() + .setResponses(new AlterShareGroupOffsetsResponseData.AlterShareGroupOffsetsResponseTopicCollection()), + new InitializeShareGroupStateParameters.Builder() + .setGroupTopicPartitionData(new GroupTopicPartitionData<>("share-group", List.of())) + .build()); + + when(runtime.scheduleWriteOperation( + ArgumentMatchers.eq("share-group-offsets-alter"), + ArgumentMatchers.eq(new TopicPartition(Topic.GROUP_METADATA_TOPIC_NAME, 0)), + ArgumentMatchers.eq(Duration.ofMillis(5000)), + ArgumentMatchers.any() + )).thenReturn(CompletableFuture.completedFuture(alterShareGroupOffsetsIntermediate)); + CompletableFuture future = service.alterShareGroupOffsets( requestContext(ApiKeys.ALTER_SHARE_GROUP_OFFSETS), groupId, request ); - AlterShareGroupOffsetsResponseData data = new AlterShareGroupOffsetsResponseData(); assertEquals(data, future.get()); } @@ -5416,7 +5461,7 @@ public void testAlterShareGroupOffsetsRequestReturnsGroupNotEmpty() throws Execu AlterShareGroupOffsetsResponseData response = new AlterShareGroupOffsetsResponseData() .setErrorCode(Errors.NON_EMPTY_GROUP.code()) - .setErrorMessage(Errors.NON_EMPTY_GROUP.message()); + .setErrorMessage("bad stuff"); when(runtime.scheduleWriteOperation( ArgumentMatchers.eq("share-group-offsets-alter"), From d1f7d0d12fe7234f687ea2bdc4f1afd5b9da3b87 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Tue, 1 Jul 2025 20:04:18 +0100 Subject: [PATCH 4/5] Checkstyle --- .../kafka/common/requests/DeleteShareGroupOffsetsRequest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java index c946f28723a91..1e28115bada87 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DeleteShareGroupOffsetsRequest.java @@ -86,7 +86,7 @@ public static DeleteShareGroupOffsetsResponseData getErrorDeleteResponseData(Err public static DeleteShareGroupOffsetsResponseData getErrorDeleteResponseData(short errorCode, String errorMessage) { return new DeleteShareGroupOffsetsResponseData() .setErrorCode(errorCode) - .setErrorMessage(errorMessage == null ? Errors.forCode(errorCode).message(): errorMessage); + .setErrorMessage(errorMessage == null ? Errors.forCode(errorCode).message() : errorMessage); } public static DeleteShareGroupOffsetsResponseData getErrorDeleteResponseData(Errors error, String errorMessage) { From 8d91d2a01c3fddfeb32a76d3e613bd8035779b41 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Wed, 2 Jul 2025 09:59:46 +0100 Subject: [PATCH 5/5] Fix authorizer test --- .../main/scala/kafka/server/KafkaApis.scala | 20 +++++++++---------- .../kafka/api/AuthorizerIntegrationTest.scala | 14 +++++++++++++ 2 files changed, 24 insertions(+), 10 deletions(-) diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 5574d09d688e6..21edc36c13dbc 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -3766,9 +3766,9 @@ class KafkaApis(val requestChannel: RequestChannel, alterShareGroupOffsetsRequest.data.topics.forEach(topic => { val topicError = { - if (!authHelper.authorize(request.context, READ, TOPIC, topic.topicName())) { + if (!authHelper.authorize(request.context, READ, TOPIC, topic.topicName)) { Some(new ApiError(Errors.TOPIC_AUTHORIZATION_FAILED)) - } else if (!metadataCache.contains(topic.topicName())) { + } else if (!metadataCache.contains(topic.topicName)) { Some(new ApiError(Errors.UNKNOWN_TOPIC_OR_PARTITION)) } else { None @@ -3776,7 +3776,7 @@ class KafkaApis(val requestChannel: RequestChannel, } topicError match { case Some(error) => - topic.partitions().forEach(partition => responseBuilder.addPartition(topic.topicName(), partition.partitionIndex(), metadataCache.topicNamesToIds(), error.error)) + topic.partitions.forEach(partition => responseBuilder.addPartition(topic.topicName, partition.partitionIndex, metadataCache.topicNamesToIds, error.error)) case None => authorizedTopicPartitions.add(topic.duplicate) } @@ -3792,10 +3792,10 @@ class KafkaApis(val requestChannel: RequestChannel, ).handle[Unit] { (response, exception) => if (exception != null) { requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(exception)) - } else if (response.errorCode() != Errors.NONE.code) { - requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(AbstractResponse.DEFAULT_THROTTLE_TIME, response.errorCode(), response.errorMessage())) + } else if (response.errorCode != Errors.NONE.code) { + requestHelper.sendMaybeThrottle(request, alterShareGroupOffsetsRequest.getErrorResponse(AbstractResponse.DEFAULT_THROTTLE_TIME, response.errorCode, response.errorMessage)) } else { - requestHelper.sendMaybeThrottle(request, responseBuilder.merge(response, metadataCache.topicNamesToIds()).build()) + requestHelper.sendMaybeThrottle(request, responseBuilder.merge(response, metadataCache.topicNamesToIds).build()) } } } @@ -3826,7 +3826,7 @@ class KafkaApis(val requestChannel: RequestChannel, new DeleteShareGroupOffsetsResponseData.DeleteShareGroupOffsetsResponseTopic() .setTopicName(topic.topicName) .setErrorCode(Errors.TOPIC_AUTHORIZATION_FAILED.code) - .setErrorMessage(Errors.TOPIC_AUTHORIZATION_FAILED.message()) + .setErrorMessage(Errors.TOPIC_AUTHORIZATION_FAILED.message) ) } else { authorizedTopics.add(topic) @@ -3840,12 +3840,12 @@ class KafkaApis(val requestChannel: RequestChannel, if (exception != null) { requestHelper.sendMaybeThrottle(request, deleteShareGroupOffsetsRequest.getErrorResponse( AbstractResponse.DEFAULT_THROTTLE_TIME, - Errors.forException(exception).code(), - exception.getMessage())) + Errors.forException(exception).code, + exception.getMessage)) } else if (responseData.errorCode() != Errors.NONE.code) { requestHelper.sendMaybeThrottle( request, - deleteShareGroupOffsetsRequest.getErrorResponse(AbstractResponse.DEFAULT_THROTTLE_TIME, responseData.errorCode(), responseData.errorMessage()) + deleteShareGroupOffsetsRequest.getErrorResponse(AbstractResponse.DEFAULT_THROTTLE_TIME, responseData.errorCode, responseData.errorMessage) ) } else { responseData.responses.forEach { topic => { diff --git a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala index 424772275ea0a..f950362354c8c 100644 --- a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala @@ -3254,6 +3254,18 @@ class AuthorizerIntegrationTest extends AbstractAuthorizerIntegrationTest { removeAllClientAcls() } + private def createEmptyShareGroup(): Unit = { + createTopicWithBrokerPrincipal(topic) + addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WILDCARD_HOST, READ, ALLOW)), shareGroupResource) + addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WILDCARD_HOST, READ, ALLOW)), topicResource) + shareConsumerConfig.put(ConsumerConfig.GROUP_ID_CONFIG, shareGroup) + val consumer = createShareConsumer() + consumer.subscribe(util.Set.of(topic)) + consumer.poll(Duration.ofMillis(500L)) + consumer.close() + removeAllClientAcls() + } + @Test def testShareGroupDescribeWithGroupDescribeAndTopicDescribeAcl(): Unit = { createShareGroupToDescribe() @@ -3614,6 +3626,7 @@ class AuthorizerIntegrationTest extends AbstractAuthorizerIntegrationTest { @Test def testDeleteShareGroupOffsetsWithoutTopicReadAcl(): Unit = { + createEmptyShareGroup() addAndVerifyAcls(shareGroupDeleteAcl(shareGroupResource), shareGroupResource) val request = deleteShareGroupOffsetsRequest @@ -3663,6 +3676,7 @@ class AuthorizerIntegrationTest extends AbstractAuthorizerIntegrationTest { @Test def testAlterShareGroupOffsetsWithoutTopicReadAcl(): Unit = { + createEmptyShareGroup() addAndVerifyAcls(shareGroupReadAcl(shareGroupResource), shareGroupResource) val request = alterShareGroupOffsetsRequest