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

gnodet pushed a commit to branch camel-4.18.x
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/camel-4.18.x by this push:
     new da0dfd76a8ce [backport camel-4.18.x] CAMEL-24814: 
camel-elasticsearch/camel-opensearch - BulkRequestAggregationStrategy must 
return the new exchange (#26724)
da0dfd76a8ce is described below

commit da0dfd76a8ce94145afeca719f39b93b0860b557
Author: Guillaume Nodet <[email protected]>
AuthorDate: Tue Sep 22 18:01:28 2026 +0200

    [backport camel-4.18.x] CAMEL-24814: camel-elasticsearch/camel-opensearch - 
BulkRequestAggregationStrategy must return the new exchange (#26724)
---
 ...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 eb7d2e91b65a..fa56af0bb4bd 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