From 7c2f020059b92267349bac16fdc2f1729226b7f1 Mon Sep 17 00:00:00 2001 From: Andrea Cosentino Date: Mon, 27 Jul 2026 22:46:48 +0200 Subject: [PATCH] CAMEL-24277: camel-aws2-msk - throw when pojoRequest=true and the body is the wrong type Child of CAMEL-24261. MSK2Producer's listClusters, createCluster, deleteCluster and describeCluster branches only acted when the body was the matching request type under pojoRequest=true; any other body silently fell through with no AWS call and no error. Add the missing else that throws IllegalArgumentException naming the required type, consistent with CAMEL-23462. Tests send a wrong-typed body to pojoRequest=true routes and assert the message; each was verified to fail (silent no-op) before the fix. The shared 4.22 upgrade-guide entry was added with CAMEL-24263. Co-authored-by: Claude Opus 4.8 Signed-off-by: Andrea Cosentino --- .../component/aws2/msk/MSK2Producer.java | 12 +++++++ .../component/aws2/msk/MSKProducerTest.java | 34 +++++++++++++++++++ 2 files changed, 46 insertions(+) diff --git a/components/camel-aws/camel-aws2-msk/src/main/java/org/apache/camel/component/aws2/msk/MSK2Producer.java b/components/camel-aws/camel-aws2-msk/src/main/java/org/apache/camel/component/aws2/msk/MSK2Producer.java index c100cbe5184d6..4d87c38854c81 100644 --- a/components/camel-aws/camel-aws2-msk/src/main/java/org/apache/camel/component/aws2/msk/MSK2Producer.java +++ b/components/camel-aws/camel-aws2-msk/src/main/java/org/apache/camel/component/aws2/msk/MSK2Producer.java @@ -115,6 +115,9 @@ private void listClusters(KafkaClient mskClient, Exchange exchange) throws Inval message.setBody(result); message.setHeader(MSK2Constants.NEXT_TOKEN, result.nextToken()); message.setHeader(MSK2Constants.IS_TRUNCATED, ObjectHelper.isNotEmpty(result.nextToken())); + } else { + throw new IllegalArgumentException( + "listClusters operation requires ListClustersRequest in POJO mode"); } } else { ListClustersRequest.Builder builder = ListClustersRequest.builder(); @@ -157,6 +160,9 @@ private void createCluster(KafkaClient mskClient, Exchange exchange) throws Inva } Message message = getMessageForResponse(exchange); message.setBody(response); + } else { + throw new IllegalArgumentException( + "createCluster operation requires CreateClusterRequest in POJO mode"); } } else { CreateClusterRequest.Builder builder = CreateClusterRequest.builder(); @@ -212,6 +218,9 @@ private void deleteCluster(KafkaClient mskClient, Exchange exchange) throws Inva } Message message = getMessageForResponse(exchange); message.setBody(result); + } else { + throw new IllegalArgumentException( + "deleteCluster operation requires DeleteClusterRequest in POJO mode"); } } else { DeleteClusterRequest.Builder builder = DeleteClusterRequest.builder(); @@ -246,6 +255,9 @@ private void describeCluster(KafkaClient mskClient, Exchange exchange) throws In } Message message = getMessageForResponse(exchange); message.setBody(result); + } else { + throw new IllegalArgumentException( + "describeCluster operation requires DescribeClusterRequest in POJO mode"); } } else { DescribeClusterRequest.Builder builder = DescribeClusterRequest.builder(); diff --git a/components/camel-aws/camel-aws2-msk/src/test/java/org/apache/camel/component/aws2/msk/MSKProducerTest.java b/components/camel-aws/camel-aws2-msk/src/test/java/org/apache/camel/component/aws2/msk/MSKProducerTest.java index 7c1e1b1b437fd..676a16f9fa614 100644 --- a/components/camel-aws/camel-aws2-msk/src/test/java/org/apache/camel/component/aws2/msk/MSKProducerTest.java +++ b/components/camel-aws/camel-aws2-msk/src/test/java/org/apache/camel/component/aws2/msk/MSKProducerTest.java @@ -149,6 +149,34 @@ public void process(Exchange exchange) { assertEquals(ClusterState.ACTIVE.name(), resultGet.clusterInfo().state().toString()); } + @Test + void listClustersWithPojoRequestAndWrongBodyTypeThrows() { + assertThatThrownBy(() -> template.requestBody("direct:listClustersPojo", "not a ListClustersRequest")) + .hasRootCauseInstanceOf(IllegalArgumentException.class) + .hasRootCauseMessage("listClusters operation requires ListClustersRequest in POJO mode"); + } + + @Test + void createClusterWithPojoRequestAndWrongBodyTypeThrows() { + assertThatThrownBy(() -> template.requestBody("direct:createClusterPojo", "not a CreateClusterRequest")) + .hasRootCauseInstanceOf(IllegalArgumentException.class) + .hasRootCauseMessage("createCluster operation requires CreateClusterRequest in POJO mode"); + } + + @Test + void deleteClusterWithPojoRequestAndWrongBodyTypeThrows() { + assertThatThrownBy(() -> template.requestBody("direct:deleteClusterPojo", "not a DeleteClusterRequest")) + .hasRootCauseInstanceOf(IllegalArgumentException.class) + .hasRootCauseMessage("deleteCluster operation requires DeleteClusterRequest in POJO mode"); + } + + @Test + void describeClusterWithPojoRequestAndWrongBodyTypeThrows() { + assertThatThrownBy(() -> template.requestBody("direct:describeClusterPojo", "not a DescribeClusterRequest")) + .hasRootCauseInstanceOf(IllegalArgumentException.class) + .hasRootCauseMessage("describeCluster operation requires DescribeClusterRequest in POJO mode"); + } + @Override protected RouteBuilder createRouteBuilder() { return new RouteBuilder() { @@ -165,6 +193,12 @@ public void configure() { .to("mock:result"); from("direct:describeCluster").to("aws2-msk://test?mskClient=#amazonMskClient&operation=describeCluster") .to("mock:result"); + from("direct:createClusterPojo") + .to("aws2-msk://test?mskClient=#amazonMskClient&operation=createCluster&pojoRequest=true"); + from("direct:deleteClusterPojo") + .to("aws2-msk://test?mskClient=#amazonMskClient&operation=deleteCluster&pojoRequest=true"); + from("direct:describeClusterPojo") + .to("aws2-msk://test?mskClient=#amazonMskClient&operation=describeCluster&pojoRequest=true"); } }; }