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"); } }; }