From e32916507d6403458a37558138359a42dda5031d Mon Sep 17 00:00:00 2001 From: lixiachen Date: Thu, 25 Sep 2025 15:28:12 -0400 Subject: [PATCH 1/2] textproxy: Allow testproxy to build its own proto registry Change-Id: Ie930064363d92d61daad498e8380dd87ab6722dd --- .../bigtable/testproxy/CbtTestProxy.java | 3 +- .../testproxy/ResultSetSerializer.java | 178 +++++++++++++++++- test-proxy/src/main/proto/test_proxy.proto | 4 + 3 files changed, 177 insertions(+), 8 deletions(-) diff --git a/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/CbtTestProxy.java b/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/CbtTestProxy.java index da205c3d3dba..122bc514394c 100644 --- a/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/CbtTestProxy.java +++ b/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/CbtTestProxy.java @@ -706,7 +706,8 @@ public void executeQuery( .dataClient() .executeQuery( BoundStatementDeserializer.toBoundStatement(preparedStatement, request)); - responseObserver.onNext(ResultSetSerializer.toExecuteQueryResult(resultSet)); + responseObserver.onNext( + new ResultSetSerializer(request.getProtoDescriptors()).toExecuteQueryResult(resultSet)); } catch (InterruptedException e) { responseObserver.onError(e); return; diff --git a/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java b/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java index 1868108efc70..f6fa72b31314 100644 --- a/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java +++ b/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java @@ -35,13 +35,155 @@ import com.google.cloud.bigtable.data.v2.models.sql.ResultSet; import com.google.cloud.bigtable.data.v2.models.sql.SqlType; import com.google.cloud.bigtable.data.v2.models.sql.StructReader; +import com.google.protobuf.AbstractMessage; import com.google.protobuf.ByteString; +import com.google.protobuf.DescriptorProtos.FileDescriptorProto; +import com.google.protobuf.DescriptorProtos.FileDescriptorSet; +import com.google.protobuf.Descriptors.Descriptor; +import com.google.protobuf.Descriptors.DescriptorValidationException; +import com.google.protobuf.Descriptors.EnumDescriptor; +import com.google.protobuf.Descriptors.FileDescriptor; +import com.google.protobuf.DynamicMessage; +import com.google.protobuf.ProtocolMessageEnum; import java.time.Instant; +import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import java.util.concurrent.ExecutionException; public class ResultSetSerializer { - public static ExecuteQueryResult toExecuteQueryResult(ResultSet resultSet) + + // This is a helper enum to satisfy the type constraints of {@link StructReader#getProtoEnum}. + private static class DummyEnum implements ProtocolMessageEnum { + private final int value; + private final EnumDescriptor descriptor; + + private DummyEnum(int value, EnumDescriptor descriptor) { + this.value = value; + this.descriptor = descriptor; + } + + @Override + public int getNumber() { + return value; + } + + @Override + public com.google.protobuf.Descriptors.EnumValueDescriptor getValueDescriptor() { + return descriptor.findValueByNumber(value); + } + + @Override + public com.google.protobuf.Descriptors.EnumDescriptor getDescriptorForType() { + return descriptor; + } + } + + /** + * A map of all known message descriptors, keyed by their fully qualified name (e.g., + * "my.package.MyMessage"). + */ + private final java.util.Map messageDescriptorMap; + + /** + * A map of all known enum descriptors, keyed by their fully qualified name (e.g., + * "my.package.MyEnum"). + */ + private final java.util.Map enumDescriptorMap; + + /** + * Helper function to recursively adds a message descriptor and all its nested types to the map. + */ + private void populateDescriptorMapsRecursively(Descriptor descriptor) { + messageDescriptorMap.put(descriptor.getFullName(), descriptor); + + for (EnumDescriptor nestedEnum : descriptor.getEnumTypes()) { + enumDescriptorMap.put(nestedEnum.getFullName(), nestedEnum); + } + for (Descriptor nestedMessage : descriptor.getNestedTypes()) { + populateDescriptorMapsRecursively(nestedMessage); + } + } + + /** + * Creates a serializer with a descriptor cache built from the provided FileDescriptorSet. This is + * useful for handling PROTO or ENUM types that require schema lookup. + * + * @param descriptorSet A set containing one or more .proto file definitions and all their + * non-standard dependencies. + * @throws IllegalArgumentException if the descriptorSet contains unresolvable dependencies. + */ + public ResultSetSerializer(FileDescriptorSet descriptorSet) throws IllegalArgumentException { + this.messageDescriptorMap = new HashMap<>(); + this.enumDescriptorMap = new HashMap<>(); + + // Java's dynamic descriptor loading requires manual dependency resolution. + // We must build a DAG and build files only after their dependencies are built. + java.util.Map pendingProtos = new HashMap<>(); + java.util.Map builtDescriptors = new HashMap<>(); + for (FileDescriptorProto fileDescriptorProto : descriptorSet.getFileList()) { + pendingProtos.put(fileDescriptorProto.getName(), fileDescriptorProto); + } + + // Keep looping as long as we are making progress. + int filesBuiltThisPass; + do { + filesBuiltThisPass = 0; + // Use a copy of the values, to avoid modifying the pendingProtos map while iterating it. + List protosToBuild = new ArrayList<>(pendingProtos.values()); + + for (FileDescriptorProto proto : protosToBuild) { + boolean allDependenciesMet = true; + List dependencies = new ArrayList<>(); + + for (String dependencyName : proto.getDependencyList()) { + FileDescriptor dependency = builtDescriptors.get(dependencyName); + if (dependency != null) { + // Dependency is already built, add it. + dependencies.add(dependency); + } else if (pendingProtos.containsKey(dependencyName)) { + // Dependency is not yet built, we must wait. + allDependenciesMet = false; + break; + } else { + // Dependency is not in our set. We assume it's a well-known type (e.g., + // google/protobuf/timestamp.proto) that buildFrom() can find and link automatically. + } + } + + if (allDependenciesMet) { + try { + // All dependencies are met, we can build this file. + FileDescriptor fileDescriptor = + FileDescriptor.buildFrom(proto, dependencies.toArray(new FileDescriptor[0])); + + builtDescriptors.put(fileDescriptor.getName(), fileDescriptor); + pendingProtos.remove(proto.getName()); + filesBuiltThisPass++; + + // Now, populate both message and enum maps with all messages/enums in this file. + for (EnumDescriptor enumDescriptor : fileDescriptor.getEnumTypes()) { + enumDescriptorMap.put(enumDescriptor.getFullName(), enumDescriptor); + } + for (Descriptor messageDescriptor : fileDescriptor.getMessageTypes()) { + populateDescriptorMapsRecursively(messageDescriptor); + } + } catch (DescriptorValidationException e) { + throw new IllegalArgumentException( + "Failed to build descriptor for " + proto.getName(), e); + } + } + } + } while (filesBuiltThisPass > 0); + + // Throw if we finished looping but still have pending protos. + if (!pendingProtos.isEmpty()) { + throw new IllegalArgumentException( + "Unresolvable or circular dependencies found for: " + pendingProtos.keySet()); + } + } + + public ExecuteQueryResult toExecuteQueryResult(ResultSet resultSet) throws ExecutionException, InterruptedException { ExecuteQueryResult.Builder resultBuilder = ExecuteQueryResult.newBuilder(); for (ColumnMetadata columnMetadata : resultSet.getMetadata().getColumns()) { @@ -64,7 +206,7 @@ public static ExecuteQueryResult toExecuteQueryResult(ResultSet resultSet) return resultBuilder.build(); } - private static Value toProtoValue(Object value, SqlType type) { + private Value toProtoValue(Object value, SqlType type) { if (value == null) { return Value.getDefaultInstance(); } @@ -72,16 +214,20 @@ private static Value toProtoValue(Object value, SqlType type) { Value.Builder valueBuilder = Value.newBuilder(); switch (type.getCode()) { case BYTES: - case PROTO: valueBuilder.setBytesValue((ByteString) value); break; + case PROTO: + valueBuilder.setBytesValue(((AbstractMessage) value).toByteString()); + break; case STRING: valueBuilder.setStringValue((String) value); break; case INT64: - case ENUM: valueBuilder.setIntValue((Long) value); break; + case ENUM: + valueBuilder.setIntValue(((ProtocolMessageEnum) value).getNumber()); + break; case FLOAT32: valueBuilder.setFloatValue((Float) value); break; @@ -151,7 +297,7 @@ private static Value toProtoValue(Object value, SqlType type) { return valueBuilder.build(); } - private static Object getColumn(StructReader struct, int fieldIndex, SqlType fieldType) { + private Object getColumn(StructReader struct, int fieldIndex, SqlType fieldType) { if (struct.isNull(fieldIndex)) { return null; } @@ -162,8 +308,15 @@ private static Object getColumn(StructReader struct, int fieldIndex, SqlType case BOOL: return struct.getBoolean(fieldIndex); case BYTES: - case PROTO: return struct.getBytes(fieldIndex); + case PROTO: + SchemalessProto protoType = (SchemalessProto) fieldType; + Descriptor descriptor = messageDescriptorMap.get(protoType.getMessageName()); + if (descriptor == null) { + throw new IllegalArgumentException( + "Descriptor for message " + protoType.getMessageName() + " could not be found"); + } + return struct.getProtoMessage(fieldIndex, DynamicMessage.getDefaultInstance(descriptor)); case DATE: return struct.getDate(fieldIndex); case FLOAT32: @@ -171,8 +324,19 @@ private static Object getColumn(StructReader struct, int fieldIndex, SqlType case FLOAT64: return struct.getDouble(fieldIndex); case INT64: - case ENUM: return struct.getLong(fieldIndex); + case ENUM: + SchemalessEnum enumType = (SchemalessEnum) fieldType; + EnumDescriptor enumDescriptor = enumDescriptorMap.get(enumType.getEnumName()); + if (enumDescriptor == null) { + throw new IllegalArgumentException( + "Descriptor for enum " + enumType.getEnumName() + " could not be found"); + } + // We need to extract the integer value of the enum. `getProtoEnum` is the only + // available method, but it is designed for static enum types. To work around this, + // we can pass a lambda that constructs our DummyEnum with the captured integer value + // and the descriptor from the outer scope. + return struct.getProtoEnum(fieldIndex, number -> new DummyEnum(number, enumDescriptor)); case MAP: return struct.getMap(fieldIndex, (SqlType.Map) fieldType); case STRING: diff --git a/test-proxy/src/main/proto/test_proxy.proto b/test-proxy/src/main/proto/test_proxy.proto index b82354b08e24..34cf534425c2 100644 --- a/test-proxy/src/main/proto/test_proxy.proto +++ b/test-proxy/src/main/proto/test_proxy.proto @@ -20,6 +20,7 @@ import "google/api/client.proto"; import "google/bigtable/v2/bigtable.proto"; import "google/bigtable/v2/data.proto"; import "google/protobuf/duration.proto"; +import "google/protobuf/descriptor.proto"; import "google/rpc/status.proto"; option go_package = "./testproxypb"; @@ -256,6 +257,9 @@ message ExecuteQueryRequest { // The raw request to the Bigtable server. google.bigtable.v2.ExecuteQueryRequest request = 2; + + // The file descriptor set for the query. + google.protobuf.FileDescriptorSet proto_descriptors = 3; } // Response from test proxy service for ExecuteQueryRequest. From a61c4b5f575a751615f4c31fd5d0a90a233eb8d8 Mon Sep 17 00:00:00 2001 From: lixiachen Date: Wed, 1 Oct 2025 13:51:26 -0400 Subject: [PATCH 2/2] testproxy: Remove manual dependency resolution from ResultSetSerializer We now require the clients to provide FileDescriptorSet with files in dependency order Change-Id: I240cfe8f4d499ba7053dfab1f119972de1810fce --- .../testproxy/ResultSetSerializer.java | 87 ++++++------------- 1 file changed, 28 insertions(+), 59 deletions(-) diff --git a/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java b/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java index f6fa72b31314..2467ce529967 100644 --- a/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java +++ b/test-proxy/src/main/java/com/google/cloud/bigtable/testproxy/ResultSetSerializer.java @@ -55,6 +55,7 @@ public class ResultSetSerializer { // This is a helper enum to satisfy the type constraints of {@link StructReader#getProtoEnum}. private static class DummyEnum implements ProtocolMessageEnum { + private final int value; private final EnumDescriptor descriptor; @@ -110,76 +111,44 @@ private void populateDescriptorMapsRecursively(Descriptor descriptor) { * useful for handling PROTO or ENUM types that require schema lookup. * * @param descriptorSet A set containing one or more .proto file definitions and all their - * non-standard dependencies. + * non-standard dependencies. All .proto file must be provided in dependency order. * @throws IllegalArgumentException if the descriptorSet contains unresolvable dependencies. */ public ResultSetSerializer(FileDescriptorSet descriptorSet) throws IllegalArgumentException { this.messageDescriptorMap = new HashMap<>(); this.enumDescriptorMap = new HashMap<>(); - - // Java's dynamic descriptor loading requires manual dependency resolution. - // We must build a DAG and build files only after their dependencies are built. - java.util.Map pendingProtos = new HashMap<>(); java.util.Map builtDescriptors = new HashMap<>(); - for (FileDescriptorProto fileDescriptorProto : descriptorSet.getFileList()) { - pendingProtos.put(fileDescriptorProto.getName(), fileDescriptorProto); - } - - // Keep looping as long as we are making progress. - int filesBuiltThisPass; - do { - filesBuiltThisPass = 0; - // Use a copy of the values, to avoid modifying the pendingProtos map while iterating it. - List protosToBuild = new ArrayList<>(pendingProtos.values()); - - for (FileDescriptorProto proto : protosToBuild) { - boolean allDependenciesMet = true; - List dependencies = new ArrayList<>(); - for (String dependencyName : proto.getDependencyList()) { - FileDescriptor dependency = builtDescriptors.get(dependencyName); - if (dependency != null) { - // Dependency is already built, add it. - dependencies.add(dependency); - } else if (pendingProtos.containsKey(dependencyName)) { - // Dependency is not yet built, we must wait. - allDependenciesMet = false; - break; - } else { - // Dependency is not in our set. We assume it's a well-known type (e.g., - // google/protobuf/timestamp.proto) that buildFrom() can find and link automatically. - } + for (FileDescriptorProto fileDescriptorProto : descriptorSet.getFileList()) { + // Collect dependencies. This code require files inside the descriptor set to be sorted + // according to the dependency order. + List dependencies = new ArrayList<>(); + for (String dependencyName : fileDescriptorProto.getDependencyList()) { + FileDescriptor dependency = builtDescriptors.get(dependencyName); + if (dependency != null) { + // Dependency is already built, add it. + dependencies.add(dependency); } + // Dependency is not in our set. We assume it's a well-known type (e.g., + // google/protobuf/timestamp.proto) that buildFrom() can find and link automatically. + } - if (allDependenciesMet) { - try { - // All dependencies are met, we can build this file. - FileDescriptor fileDescriptor = - FileDescriptor.buildFrom(proto, dependencies.toArray(new FileDescriptor[0])); - - builtDescriptors.put(fileDescriptor.getName(), fileDescriptor); - pendingProtos.remove(proto.getName()); - filesBuiltThisPass++; - - // Now, populate both message and enum maps with all messages/enums in this file. - for (EnumDescriptor enumDescriptor : fileDescriptor.getEnumTypes()) { - enumDescriptorMap.put(enumDescriptor.getFullName(), enumDescriptor); - } - for (Descriptor messageDescriptor : fileDescriptor.getMessageTypes()) { - populateDescriptorMapsRecursively(messageDescriptor); - } - } catch (DescriptorValidationException e) { - throw new IllegalArgumentException( - "Failed to build descriptor for " + proto.getName(), e); - } + try { + FileDescriptor fileDescriptor = + FileDescriptor.buildFrom( + fileDescriptorProto, dependencies.toArray(new FileDescriptor[0])); + builtDescriptors.put(fileDescriptor.getName(), fileDescriptor); + // Now, populate both message and enum maps with all messages/enums in this file. + for (EnumDescriptor enumDescriptor : fileDescriptor.getEnumTypes()) { + enumDescriptorMap.put(enumDescriptor.getFullName(), enumDescriptor); + } + for (Descriptor messageDescriptor : fileDescriptor.getMessageTypes()) { + populateDescriptorMapsRecursively(messageDescriptor); } + } catch (DescriptorValidationException e) { + throw new IllegalArgumentException( + "Failed to build descriptor for " + fileDescriptorProto.getName(), e); } - } while (filesBuiltThisPass > 0); - - // Throw if we finished looping but still have pending protos. - if (!pendingProtos.isEmpty()) { - throw new IllegalArgumentException( - "Unresolvable or circular dependencies found for: " + pendingProtos.keySet()); } }