Skip to content

Commit d3f8ea0

Browse files
committed
update
1 parent c05ccda commit d3f8ea0

24 files changed

Lines changed: 420 additions & 441 deletions

.idea/runConfigurations/Debug_OpenSearch.xml

Lines changed: 0 additions & 11 deletions
This file was deleted.

.idea/vcs.xml

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

modules/transport-grpc/spi/src/main/java/org/opensearch/transport/grpc/spi/AggregateProtoConverter.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,9 @@
1717
* SPI interface for converting OpenSearch InternalAggregation objects to Protocol Buffer Aggregate messages.
1818
* Follows the same pattern as {@link AggregationBuilderProtoConverter} for request-side conversions.
1919
*
20-
* <p>The registry handles metadata centrally. Converters should delegate to existing
21-
* {@code *AggregateProtoUtils} classes for the actual conversion logic.
20+
* <p>Implementations should build the typed inner message (e.g., {@code SingleMetricAggregateBase},
21+
* {@code LongTermsAggregate}) with metadata included, then wrap it in {@code Aggregate.Builder}
22+
* using the appropriate setter (e.g., {@code setMax()}, {@code setLterms()}).
2223
*/
2324
public interface AggregateProtoConverter {
2425

modules/transport-grpc/src/main/java/org/opensearch/transport/grpc/proto/response/search/aggregation/bucket/terms/DoubleTermsAggregateConverter.java

Lines changed: 36 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -2,20 +2,21 @@
22

33
import org.opensearch.core.xcontent.XContentBuilder;
44
import org.opensearch.protobufs.Aggregate;
5-
import org.opensearch.protobufs.ObjectMap;
5+
import org.opensearch.protobufs.DoubleTermsAggregate;
6+
import org.opensearch.protobufs.DoubleTermsBucket;
67
import org.opensearch.search.DocValueFormat;
78
import org.opensearch.search.aggregations.Aggregation;
89
import org.opensearch.search.aggregations.InternalAggregation;
910
import org.opensearch.search.aggregations.bucket.terms.DoubleTerms;
11+
import org.opensearch.transport.grpc.proto.response.search.aggregation.AggregateProtoUtils;
12+
import org.opensearch.transport.grpc.spi.AggregateProtoConverter;
1013

1114
import java.io.IOException;
1215

13-
import static org.opensearch.transport.grpc.proto.response.search.aggregation.AggregateProtoUtils.newValue;
14-
1516
/**
1617
* Proto converter for {@link DoubleTerms}
1718
*/
18-
public class DoubleTermsAggregateConverter extends TermsAggregatesProtoConverter<DoubleTerms.Bucket> {
19+
public class DoubleTermsAggregateConverter implements AggregateProtoConverter {
1920
@Override
2021
public Class<? extends InternalAggregation> getHandledAggregationType() {
2122
return DoubleTerms.class;
@@ -24,21 +25,43 @@ public Class<? extends InternalAggregation> getHandledAggregationType() {
2425
@Override
2526
public Aggregate.Builder toProto(InternalAggregation aggregation) throws IOException {
2627
DoubleTerms doubleTerms = (DoubleTerms) aggregation;
27-
Aggregate.Builder protoBuilder = Aggregate.newBuilder();
28-
convertCommon(protoBuilder, doubleTerms.getDocCountError(), doubleTerms.getSumOfOtherDocCounts(), doubleTerms.getBuckets());
29-
return protoBuilder;
28+
DoubleTermsAggregate.Builder termsBuilder = DoubleTermsAggregate.newBuilder();
29+
30+
termsBuilder.setDocCountErrorUpperBound(doubleTerms.getDocCountError());
31+
termsBuilder.setSumOtherDocCount(doubleTerms.getSumOfOtherDocCounts());
32+
33+
for (DoubleTerms.Bucket bucket : doubleTerms.getBuckets()) {
34+
termsBuilder.addBuckets(convertBucket(bucket));
35+
}
36+
37+
TermsAggregateProtoUtils.applyMetadata(termsBuilder::setMeta, doubleTerms);
38+
39+
return Aggregate.newBuilder().setDterms(termsBuilder);
3040
}
3141

3242
/**
3343
* Mirroring {@link DoubleTerms.Bucket#keyToXContent(XContentBuilder)}
34-
*
35-
* {@inheritDoc}
3644
*/
37-
@Override
38-
void convertBucketKey(ObjectMap.Builder builder, DoubleTerms.Bucket bucket) {
39-
builder.putFields(Aggregation.CommonFields.KEY.getPreferredName(), newValue((double) bucket.getKey()));
45+
private DoubleTermsBucket convertBucket(DoubleTerms.Bucket bucket) throws IOException {
46+
DoubleTermsBucket.Builder builder = DoubleTermsBucket.newBuilder();
47+
48+
builder.setKey((double) bucket.getKey());
49+
4050
if (bucket.getFormat() != DocValueFormat.RAW) {
41-
builder.putFields(Aggregation.CommonFields.KEY_AS_STRING.getPreferredName(), newValue(bucket.getKeyAsString()));
51+
builder.setKeyAsString(bucket.getKeyAsString());
52+
}
53+
54+
builder.setDocCount(bucket.getDocCount());
55+
if (bucket.showDocCountError()) {
56+
builder.setDocCountErrorUpperBound(bucket.getDocCountError());
4257
}
58+
59+
for (Aggregation subAgg : bucket.getAggregations()) {
60+
if (subAgg instanceof InternalAggregation internalAgg) {
61+
builder.getMutableAggregate().put(subAgg.getName(), AggregateProtoUtils.toProto(internalAgg));
62+
}
63+
}
64+
65+
return builder.build();
4366
}
4467
}

modules/transport-grpc/src/main/java/org/opensearch/transport/grpc/proto/response/search/aggregation/bucket/terms/LongTermsAggregateConverter.java

Lines changed: 40 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -2,20 +2,25 @@
22

33
import org.opensearch.core.xcontent.XContentBuilder;
44
import org.opensearch.protobufs.Aggregate;
5+
import org.opensearch.protobufs.LongTermsAggregate;
6+
import org.opensearch.protobufs.LongTermsBucket;
7+
import org.opensearch.protobufs.LongTermsBucketKey;
58
import org.opensearch.protobufs.ObjectMap;
69
import org.opensearch.search.DocValueFormat;
710
import org.opensearch.search.aggregations.Aggregation;
811
import org.opensearch.search.aggregations.InternalAggregation;
12+
import org.opensearch.search.aggregations.bucket.terms.InternalTerms;
913
import org.opensearch.search.aggregations.bucket.terms.LongTerms;
14+
import org.opensearch.transport.grpc.proto.response.common.ObjectMapProtoUtils;
15+
import org.opensearch.transport.grpc.proto.response.search.aggregation.AggregateProtoUtils;
16+
import org.opensearch.transport.grpc.spi.AggregateProtoConverter;
1017

1118
import java.io.IOException;
1219

13-
import static org.opensearch.transport.grpc.proto.response.search.aggregation.AggregateProtoUtils.newValue;
14-
1520
/**
1621
* Proto converter for {@link LongTerms}
1722
*/
18-
public class LongTermsAggregateConverter extends TermsAggregatesProtoConverter<LongTerms.Bucket> {
23+
public class LongTermsAggregateConverter implements AggregateProtoConverter {
1924
@Override
2025
public Class<? extends InternalAggregation> getHandledAggregationType() {
2126
return LongTerms.class;
@@ -24,28 +29,48 @@ public Class<? extends InternalAggregation> getHandledAggregationType() {
2429
@Override
2530
public Aggregate.Builder toProto(InternalAggregation aggregation) throws IOException {
2631
LongTerms longTerms = (LongTerms) aggregation;
27-
Aggregate.Builder protoBuilder = Aggregate.newBuilder();
28-
convertCommon(protoBuilder, longTerms.getDocCountError(), longTerms.getSumOfOtherDocCounts(), longTerms.getBuckets());
29-
return protoBuilder;
32+
LongTermsAggregate.Builder termsBuilder = LongTermsAggregate.newBuilder();
33+
34+
termsBuilder.setDocCountErrorUpperBound(longTerms.getDocCountError());
35+
termsBuilder.setSumOtherDocCount(longTerms.getSumOfOtherDocCounts());
36+
37+
for (LongTerms.Bucket bucket : longTerms.getBuckets()) {
38+
termsBuilder.addBuckets(convertBucket(bucket));
39+
}
40+
41+
TermsAggregateProtoUtils.applyMetadata(termsBuilder::setMeta, longTerms);
42+
43+
return Aggregate.newBuilder().setLterms(termsBuilder);
3044
}
3145

3246
/**
3347
* Mirroring {@link LongTerms.Bucket#keyToXContent(XContentBuilder)}
34-
*
35-
* {@inheritDoc}
3648
*/
37-
@Override
38-
void convertBucketKey(ObjectMap.Builder builder, LongTerms.Bucket bucket) {
49+
private LongTermsBucket convertBucket(LongTerms.Bucket bucket) throws IOException {
50+
LongTermsBucket.Builder builder = LongTermsBucket.newBuilder();
51+
3952
Object key = bucket.getKey();
40-
// the key could be a long or a BigInteger produced by UNSIGNED_LONG_SHIFTED
4153
if (key instanceof Long) {
42-
builder.putFields(Aggregation.CommonFields.KEY.getPreferredName(), newValue((long) key));
54+
builder.setKey(LongTermsBucketKey.newBuilder().setSigned((long) key));
4355
} else {
44-
// BigInteger's case, use string to represent
45-
builder.putFields(Aggregation.CommonFields.KEY.getPreferredName(), newValue(key.toString()));
56+
builder.setKey(LongTermsBucketKey.newBuilder().setUnsigned(key.toString()));
4657
}
58+
4759
if (bucket.getFormat() != DocValueFormat.RAW && bucket.getFormat() != DocValueFormat.UNSIGNED_LONG_SHIFTED) {
48-
builder.putFields(Aggregation.CommonFields.KEY_AS_STRING.getPreferredName(), newValue(bucket.getKeyAsString()));
60+
builder.setKeyAsString(bucket.getKeyAsString());
61+
}
62+
63+
builder.setDocCount(bucket.getDocCount());
64+
if (bucket.showDocCountError()) {
65+
builder.setDocCountErrorUpperBound(bucket.getDocCountError());
4966
}
67+
68+
for (Aggregation subAgg : bucket.getAggregations()) {
69+
if (subAgg instanceof InternalAggregation internalAgg) {
70+
builder.getMutableAggregate().put(subAgg.getName(), AggregateProtoUtils.toProto(internalAgg));
71+
}
72+
}
73+
74+
return builder.build();
5075
}
5176
}

modules/transport-grpc/src/main/java/org/opensearch/transport/grpc/proto/response/search/aggregation/bucket/terms/StringTermsAggregateConverter.java

Lines changed: 34 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -2,19 +2,20 @@
22

33
import org.opensearch.core.xcontent.XContentBuilder;
44
import org.opensearch.protobufs.Aggregate;
5-
import org.opensearch.protobufs.ObjectMap;
5+
import org.opensearch.protobufs.StringTermsAggregate;
6+
import org.opensearch.protobufs.StringTermsBucket;
67
import org.opensearch.search.aggregations.Aggregation;
78
import org.opensearch.search.aggregations.InternalAggregation;
89
import org.opensearch.search.aggregations.bucket.terms.StringTerms;
10+
import org.opensearch.transport.grpc.proto.response.search.aggregation.AggregateProtoUtils;
11+
import org.opensearch.transport.grpc.spi.AggregateProtoConverter;
912

1013
import java.io.IOException;
1114

12-
import static org.opensearch.transport.grpc.proto.response.search.aggregation.AggregateProtoUtils.newValue;
13-
1415
/**
1516
* Proto converter for {@link StringTerms}
1617
*/
17-
public class StringTermsAggregateConverter extends TermsAggregatesProtoConverter<StringTerms.Bucket> {
18+
public class StringTermsAggregateConverter implements AggregateProtoConverter {
1819
@Override
1920
public Class<? extends InternalAggregation> getHandledAggregationType() {
2021
return StringTerms.class;
@@ -23,18 +24,39 @@ public Class<? extends InternalAggregation> getHandledAggregationType() {
2324
@Override
2425
public Aggregate.Builder toProto(InternalAggregation aggregation) throws IOException {
2526
StringTerms stringTerms = (StringTerms) aggregation;
26-
Aggregate.Builder protoBuilder = Aggregate.newBuilder();
27-
convertCommon(protoBuilder, stringTerms.getDocCountError(), stringTerms.getSumOfOtherDocCounts(), stringTerms.getBuckets());
28-
return protoBuilder;
27+
StringTermsAggregate.Builder termsBuilder = StringTermsAggregate.newBuilder();
28+
29+
termsBuilder.setDocCountErrorUpperBound(stringTerms.getDocCountError());
30+
termsBuilder.setSumOtherDocCount(stringTerms.getSumOfOtherDocCounts());
31+
32+
for (StringTerms.Bucket bucket : stringTerms.getBuckets()) {
33+
termsBuilder.addBuckets(convertBucket(bucket));
34+
}
35+
36+
TermsAggregateProtoUtils.applyMetadata(termsBuilder::setMeta, stringTerms);
37+
38+
return Aggregate.newBuilder().setSterms(termsBuilder);
2939
}
3040

3141
/**
3242
* Mirroring {@link StringTerms.Bucket#keyToXContent(XContentBuilder)}
33-
*
34-
* {@inheritDoc}
3543
*/
36-
@Override
37-
void convertBucketKey(ObjectMap.Builder builder, StringTerms.Bucket bucket) {
38-
builder.putFields(Aggregation.CommonFields.KEY.getPreferredName(), newValue(bucket.getKeyAsString()));
44+
private StringTermsBucket convertBucket(StringTerms.Bucket bucket) throws IOException {
45+
StringTermsBucket.Builder builder = StringTermsBucket.newBuilder();
46+
47+
builder.setKey(bucket.getKeyAsString());
48+
49+
builder.setDocCount(bucket.getDocCount());
50+
if (bucket.showDocCountError()) {
51+
builder.setDocCountErrorUpperBound(bucket.getDocCountError());
52+
}
53+
54+
for (Aggregation subAgg : bucket.getAggregations()) {
55+
if (subAgg instanceof InternalAggregation internalAgg) {
56+
builder.getMutableAggregate().put(subAgg.getName(), AggregateProtoUtils.toProto(internalAgg));
57+
}
58+
}
59+
60+
return builder.build();
3961
}
4062
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
/*
2+
* SPDX-License-Identifier: Apache-2.0
3+
*
4+
* The OpenSearch Contributors require contributions made to
5+
* this file be licensed under the Apache-2.0 license or a
6+
* compatible open source license.
7+
*/
8+
package org.opensearch.transport.grpc.proto.response.search.aggregation.bucket.terms;
9+
10+
import org.opensearch.protobufs.ObjectMap;
11+
import org.opensearch.search.aggregations.InternalAggregation;
12+
import org.opensearch.transport.grpc.proto.response.common.ObjectMapProtoUtils;
13+
14+
import java.util.function.Consumer;
15+
16+
/**
17+
* Shared utility methods for terms aggregation proto converters.
18+
*/
19+
class TermsAggregateProtoUtils {
20+
21+
private TermsAggregateProtoUtils() {}
22+
23+
/**
24+
* Applies metadata from an InternalAggregation to a typed terms aggregate builder
25+
* via the provided setter function.
26+
*
27+
* @param metaSetter the setMeta method reference on the typed builder
28+
* @param aggregation the source aggregation
29+
*/
30+
static void applyMetadata(Consumer<ObjectMap> metaSetter, InternalAggregation aggregation) {
31+
if (aggregation.getMetadata() != null && !aggregation.getMetadata().isEmpty()) {
32+
ObjectMap.Value metaValue = ObjectMapProtoUtils.toProto(aggregation.getMetadata());
33+
if (metaValue.hasObjectMap()) {
34+
metaSetter.accept(metaValue.getObjectMap());
35+
}
36+
}
37+
}
38+
}

0 commit comments

Comments
 (0)