This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new b4aab1b04bd8 CAMEL-24814: camel-elasticsearch/camel-opensearch - 
BulkRequestAggregationStrategy must return the new exchange
b4aab1b04bd8 is described below

commit b4aab1b04bd86b7d37aac285f6878144f0938f87
Author: Andrea Cosentino <[email protected]>
AuthorDate: Tue Sep 22 09:12:28 2026 +0200

    CAMEL-24814: camel-elasticsearch/camel-opensearch - 
BulkRequestAggregationStrategy must return the new exchange
    
    ElasticsearchBulkRequestAggregationStrategy and 
OpensearchBulkRequestAggregationStrategy
    merge BulkOperation[] messages into a single BulkRequest for use with the 
Aggregate EIP.
    Both stored the merged request on newExchange but returned oldExchange, 
which is null on
    the first call of every aggregation group. AggregateProcessor rejects a 
null return, so
    the first exchange of every group failed and later returns never carried 
the merged
    request, leaving the strategy unusable. Neither component had a test for it.
    
    Return newExchange from both strategies, add the already-aggregated 
operations before
    the new ones so the merged request preserves insertion order, and add unit 
tests
    covering the first call and a 3-step aggregation chain in both components.
    
    Closes #26588
    
    Co-authored-by: Claude <[email protected]>
---
 ...lasticsearchBulkRequestAggregationStrategy.java |  7 +-
 ...icsearchBulkRequestAggregationStrategyTest.java | 89 ++++++++++++++++++++++
 .../OpensearchBulkRequestAggregationStrategy.java  |  7 +-
 ...ensearchBulkRequestAggregationStrategyTest.java | 89 ++++++++++++++++++++++
 4 files changed, 188 insertions(+), 4 deletions(-)

diff --git 
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
 
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
index cd5ca8f91dcb..f1b7bca87c44 100644
--- 
a/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
+++ 
b/components/camel-elasticsearch/src/main/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategy.java
@@ -42,12 +42,15 @@ public class ElasticsearchBulkRequestAggregationStrategy 
implements AggregationS
 
         BulkOperation[] newBody = (BulkOperation[]) objBody;
         BulkRequest.Builder builder = new BulkRequest.Builder();
-        builder.operations(List.of(newBody));
         if (oldExchange != null) {
+            // add the already-aggregated operations first so the merged 
request keeps insertion order
             BulkRequest request = 
oldExchange.getIn().getBody(BulkRequest.class);
             builder.operations(request.operations());
         }
+        builder.operations(List.of(newBody));
+        // the merged BulkRequest is stored on the new exchange, so the new 
exchange must be the one returned
+        // (returning oldExchange would return null on the first aggregation 
call, which is not allowed)
         newExchange.getIn().setBody(builder.build());
-        return oldExchange;
+        return newExchange;
     }
 }
diff --git 
a/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategyTest.java
 
b/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategyTest.java
new file mode 100644
index 000000000000..cbf4ea2e387c
--- /dev/null
+++ 
b/components/camel-elasticsearch/src/test/java/org/apache/camel/component/es/aggregation/ElasticsearchBulkRequestAggregationStrategyTest.java
@@ -0,0 +1,89 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.es.aggregation;
+
+import java.util.List;
+import java.util.Map;
+
+import co.elastic.clients.elasticsearch.core.BulkRequest;
+import co.elastic.clients.elasticsearch.core.bulk.BulkOperation;
+import co.elastic.clients.elasticsearch.core.bulk.IndexOperation;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+public class ElasticsearchBulkRequestAggregationStrategyTest {
+
+    private final ElasticsearchBulkRequestAggregationStrategy strategy = new 
ElasticsearchBulkRequestAggregationStrategy();
+    private CamelContext context;
+
+    @BeforeEach
+    void setUp() {
+        context = new DefaultCamelContext();
+    }
+
+    @AfterEach
+    void tearDown() {
+        context.stop();
+    }
+
+    private Exchange exchangeWith(String id) {
+        Exchange exchange = new DefaultExchange(context);
+        BulkOperation operation = new BulkOperation.Builder()
+                .index(new 
IndexOperation.Builder<>().index("idx").id(id).document(Map.of("id", 
id)).build())
+                .build();
+        exchange.getIn().setBody(new BulkOperation[] { operation });
+        return exchange;
+    }
+
+    private static List<String> ids(BulkRequest request) {
+        return request.operations().stream().map(op -> 
op.index().id()).toList();
+    }
+
+    @Test
+    void firstAggregationMustReturnNewExchangeCarryingTheRequest() {
+        Exchange newExchange = exchangeWith("1");
+
+        Exchange result = strategy.aggregate(null, newExchange);
+
+        // On the first call oldExchange is null; returning it (the previous 
bug) would make the
+        // AggregationStrategy return null, which AggregateProcessor rejects.
+        assertSame(newExchange, result);
+        BulkRequest request = result.getIn().getBody(BulkRequest.class);
+        assertNotNull(request);
+        assertEquals(List.of("1"), ids(request));
+    }
+
+    @Test
+    void subsequentAggregationMergesAllOperationsInInsertionOrder() {
+        Exchange step1 = strategy.aggregate(null, exchangeWith("1"));
+        Exchange step2 = strategy.aggregate(step1, exchangeWith("2"));
+        Exchange step3 = strategy.aggregate(step2, exchangeWith("3"));
+
+        // a 3-step chain proves the merged request preserves insertion order 
cumulatively
+        assertEquals(List.of("1", "2"), 
ids(step2.getIn().getBody(BulkRequest.class)));
+        assertEquals(List.of("1", "2", "3"), 
ids(step3.getIn().getBody(BulkRequest.class)));
+    }
+}
diff --git 
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
 
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
index a4e15fd73296..914c0251027e 100644
--- 
a/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
+++ 
b/components/camel-opensearch/src/main/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategy.java
@@ -46,12 +46,15 @@ public class OpensearchBulkRequestAggregationStrategy 
implements AggregationStra
 
         BulkOperation[] newBody = (BulkOperation[]) objBody;
         BulkRequest.Builder builder = new BulkRequest.Builder();
-        builder.operations(List.of(newBody));
         if (ObjectHelper.isNotEmpty(oldExchange)) {
+            // add the already-aggregated operations first so the merged 
request keeps insertion order
             BulkRequest request = 
oldExchange.getIn().getBody(BulkRequest.class);
             builder.operations(request.operations());
         }
+        builder.operations(List.of(newBody));
+        // the merged BulkRequest is stored on the new exchange, so the new 
exchange must be the one returned
+        // (returning oldExchange would return null on the first aggregation 
call, which is not allowed)
         newExchange.getIn().setBody(builder.build());
-        return oldExchange;
+        return newExchange;
     }
 }
diff --git 
a/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategyTest.java
 
b/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategyTest.java
new file mode 100644
index 000000000000..9c4b6d247c18
--- /dev/null
+++ 
b/components/camel-opensearch/src/test/java/org/apache/camel/component/opensearch/aggregation/OpensearchBulkRequestAggregationStrategyTest.java
@@ -0,0 +1,89 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.opensearch.aggregation;
+
+import java.util.List;
+import java.util.Map;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.opensearch.client.opensearch.core.BulkRequest;
+import org.opensearch.client.opensearch.core.bulk.BulkOperation;
+import org.opensearch.client.opensearch.core.bulk.IndexOperation;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+
+public class OpensearchBulkRequestAggregationStrategyTest {
+
+    private final OpensearchBulkRequestAggregationStrategy strategy = new 
OpensearchBulkRequestAggregationStrategy();
+    private CamelContext context;
+
+    @BeforeEach
+    void setUp() {
+        context = new DefaultCamelContext();
+    }
+
+    @AfterEach
+    void tearDown() {
+        context.stop();
+    }
+
+    private Exchange exchangeWith(String id) {
+        Exchange exchange = new DefaultExchange(context);
+        BulkOperation operation = new BulkOperation.Builder()
+                .index(new 
IndexOperation.Builder<>().index("idx").id(id).document(Map.of("id", 
id)).build())
+                .build();
+        exchange.getIn().setBody(new BulkOperation[] { operation });
+        return exchange;
+    }
+
+    private static List<String> ids(BulkRequest request) {
+        return request.operations().stream().map(op -> 
op.index().id()).toList();
+    }
+
+    @Test
+    void firstAggregationMustReturnNewExchangeCarryingTheRequest() {
+        Exchange newExchange = exchangeWith("1");
+
+        Exchange result = strategy.aggregate(null, newExchange);
+
+        // On the first call oldExchange is null; returning it (the previous 
bug) would make the
+        // AggregationStrategy return null, which AggregateProcessor rejects.
+        assertSame(newExchange, result);
+        BulkRequest request = result.getIn().getBody(BulkRequest.class);
+        assertNotNull(request);
+        assertEquals(List.of("1"), ids(request));
+    }
+
+    @Test
+    void subsequentAggregationMergesAllOperationsInInsertionOrder() {
+        Exchange step1 = strategy.aggregate(null, exchangeWith("1"));
+        Exchange step2 = strategy.aggregate(step1, exchangeWith("2"));
+        Exchange step3 = strategy.aggregate(step2, exchangeWith("3"));
+
+        // a 3-step chain proves the merged request preserves insertion order 
cumulatively
+        assertEquals(List.of("1", "2"), 
ids(step2.getIn().getBody(BulkRequest.class)));
+        assertEquals(List.of("1", "2", "3"), 
ids(step3.getIn().getBody(BulkRequest.class)));
+    }
+}

Reply via email to