Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CMakeLists.txt

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions build_autogenerated.yaml

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

20 changes: 13 additions & 7 deletions src/core/lib/debug/trace_flags.cc

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 0 additions & 1 deletion src/python/grpcio_tests/tests/interop/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -318,4 +318,3 @@ def test_interoperability(args):

if __name__ == "__main__":
app.run(test_interoperability, flags_parser=parse_interop_client_args)

Original file line number Diff line number Diff line change
Expand Up @@ -159,7 +159,6 @@ def shutdown_tracer_provider():
time.sleep(0.5)



def pack_grpc_trace_bin(
trace_id_int: int, span_id_int: int, is_sampled: bool = True
) -> bytes:
Expand Down
5 changes: 2 additions & 3 deletions test/cpp/interop/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -364,7 +364,6 @@ py_test(
tags = ["no_windows"],
)


grpc_cc_library(
name = "xds_stats_watcher",
srcs = ["xds_stats_watcher.cc"],
Expand Down Expand Up @@ -847,9 +846,9 @@ grpc_cc_binary(
"//:gpr",
"//:grpc++",
"//src/core:instrument",
"@opentelemetry_proto//:trace_service_grpc_cc",
"@opentelemetry_proto//:trace_service_proto_cc",
"@opentelemetry_proto//:metrics_service_grpc_cc",
"@opentelemetry_proto//:metrics_service_proto_cc",
"@opentelemetry_proto//:trace_service_grpc_cc",
"@opentelemetry_proto//:trace_service_proto_cc",
],
)
4 changes: 2 additions & 2 deletions test/cpp/interop/otel_helper.cc
Original file line number Diff line number Diff line change
Expand Up @@ -140,8 +140,8 @@ void MaybeRegisterOpenTelemetry() {
LOG(ERROR) << "Failed to register gRPC OpenTelemetry Plugin: "
<< status.ToString();
} else {
LOG(INFO)
<< "Successfully registered gRPC OpenTelemetry Plugin for tracing and metrics.";
LOG(INFO) << "Successfully registered gRPC OpenTelemetry Plugin for "
"tracing and metrics.";
}
});
#endif
Expand Down
15 changes: 6 additions & 9 deletions test/cpp/interop/otlp_collector.cc
Original file line number Diff line number Diff line change
Expand Up @@ -114,12 +114,11 @@ class MetricsServiceServiceImpl final
explicit MetricsServiceServiceImpl(std::string file_path)
: file_path_(std::move(file_path)) {}

grpc::Status Export(
grpc::ServerContext* /*context*/,
const opentelemetry::proto::collector::metrics::v1::
ExportMetricsServiceRequest* request,
opentelemetry::proto::collector::metrics::v1::
ExportMetricsServiceResponse* /*response*/) override {
grpc::Status Export(grpc::ServerContext* /*context*/,
const opentelemetry::proto::collector::metrics::v1::
ExportMetricsServiceRequest* request,
opentelemetry::proto::collector::metrics::v1::
ExportMetricsServiceResponse* /*response*/) override {
std::string json_string;
google::protobuf::json::PrintOptions options;
options.add_whitespace = true;
Expand Down Expand Up @@ -201,8 +200,7 @@ int main(int argc, char** argv) {

std::unique_ptr<TraceServiceServiceImpl> trace_service;
if (!trace_file_path.empty()) {
trace_service =
std::make_unique<TraceServiceServiceImpl>(trace_file_path);
trace_service = std::make_unique<TraceServiceServiceImpl>(trace_file_path);
builder.RegisterService(trace_service.get());
}

Expand All @@ -226,4 +224,3 @@ int main(int argc, char** argv) {

return 0;
}

66 changes: 50 additions & 16 deletions test/cpp/interop/run_otel_interop_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ def resolve_binary(path):
if os.path.exists(clean_path):
return clean_path
if clean_path.startswith("bazel-bin/"):
alt_path = clean_path[len("bazel-bin/"):]
alt_path = clean_path[len("bazel-bin/") :]
if os.path.exists(alt_path):
return alt_path
return path
Expand Down Expand Up @@ -300,23 +300,33 @@ def verify_metrics(metrics_file, server_lang="c++"):
collected_metrics[metric_name] = []
collected_metrics[metric_name].append(metric)

if all(any(k in m for m in collected_metrics) for k in expected_keywords):
if all(
any(k in m for m in collected_metrics) for k in expected_keywords
):
break
time.sleep(0.5)

print(f"Collected {len(collected_metrics)} metric names: {set(collected_metrics.keys())}")
print(
f"Collected {len(collected_metrics)} metric names: {set(collected_metrics.keys())}"
)
if not collected_metrics:
if server_lang != "c++":
print(f"No metrics collected for non-C++ server ({server_lang}). TCP metrics only required for C++ servers. Metrics verification passed.")
print(
f"No metrics collected for non-C++ server ({server_lang}). TCP metrics only required for C++ servers. Metrics verification passed."
)
return True
print("Assertion Failed: No metrics collected.")
return False

# 1. Assert domain presence
if server_lang == "c++":
found = all(any(k in m for m in collected_metrics) for k in expected_keywords)
found = all(
any(k in m for m in collected_metrics) for k in expected_keywords
)
if not found:
print(f"Assertion Failed: Not all of {expected_keywords} found in collected metrics.")
print(
f"Assertion Failed: Not all of {expected_keywords} found in collected metrics."
)
return False

expected_tcp_units = {
Expand All @@ -333,13 +343,17 @@ def verify_metrics(metrics_file, server_lang="c++"):
for mname, expected_unit in expected_tcp_units.items():
if mname not in collected_metrics:
if mname in required_tcp_metrics:
print(f"Assertion Failed: Required TCP metric '{mname}' not found.")
print(
f"Assertion Failed: Required TCP metric '{mname}' not found."
)
return False
continue
metrics_with_name = collected_metrics[mname]
actual_unit = metrics_with_name[0].get("unit", "")
if actual_unit != expected_unit:
print(f"Assertion Failed: Metric '{mname}' unit '{actual_unit}' != expected '{expected_unit}'")
print(
f"Assertion Failed: Metric '{mname}' unit '{actual_unit}' != expected '{expected_unit}'"
)
return False

# 3. Assert TCP metric label keys
Expand All @@ -354,7 +368,12 @@ def verify_metrics(metrics_file, server_lang="c++"):
for mname, metrics_list in collected_metrics.items():
if mname.startswith("grpc.tcp."):
for metric in metrics_list:
for dp_type in ("gauge", "sum", "histogram", "exponential_histogram"):
for dp_type in (
"gauge",
"sum",
"histogram",
"exponential_histogram",
):
if dp_type in metric:
data_points = metric[dp_type].get("data_points", [])
for dp in data_points:
Expand All @@ -364,7 +383,9 @@ def verify_metrics(metrics_file, server_lang="c++"):

missing_keys = required_label_keys - found_tcp_label_keys
if missing_keys:
print(f"Assertion Failed: Missing TCP label keys: {missing_keys}. Found keys: {found_tcp_label_keys}")
print(
f"Assertion Failed: Missing TCP label keys: {missing_keys}. Found keys: {found_tcp_label_keys}"
)
return False
else:
# For non-C++ servers, check if any grpc metric was recorded
Expand Down Expand Up @@ -472,8 +493,12 @@ def main():

collector_port = get_free_port()
server_port = get_free_port()
spans_file = os.path.abspath(f"captured_spans_{args.client}_{args.server}.json")
metrics_file = os.path.abspath(f"captured_metrics_{args.client}_{args.server}.json")
spans_file = os.path.abspath(
f"captured_spans_{args.client}_{args.server}.json"
)
metrics_file = os.path.abspath(
f"captured_metrics_{args.client}_{args.server}.json"
)
if os.path.exists(spans_file):
os.remove(spans_file)
if os.path.exists(metrics_file):
Expand Down Expand Up @@ -513,7 +538,9 @@ def main():
if args.server == "c++":
server_proc = start_proc(
[
resolve_binary("./bazel-bin/test/cpp/interop/interop_server"),
resolve_binary(
"./bazel-bin/test/cpp/interop/interop_server"
),
f"--port={server_port}",
"--enable_opentelemetry=true",
"--enable_tcp_metrics=true",
Expand All @@ -535,7 +562,9 @@ def main():
elif args.server == "python":
server_proc = start_proc(
[
resolve_binary("./bazel-bin/src/python/grpcio_tests/tests/interop/server_bin"),
resolve_binary(
"./bazel-bin/src/python/grpcio_tests/tests/interop/server_bin"
),
f"--port={server_port}",
"--use_tls=false",
"--enable_opentelemetry=true",
Expand Down Expand Up @@ -566,7 +595,9 @@ def main():
if args.client == "c++":
client_res = run_cmd(
[
resolve_binary("./bazel-bin/test/cpp/interop/interop_client"),
resolve_binary(
"./bazel-bin/test/cpp/interop/interop_client"
),
"--server_host=localhost",
f"--server_port={server_port}",
"--test_case=empty_unary",
Expand All @@ -592,7 +623,9 @@ def main():
elif args.client == "python":
client_res = run_cmd(
[
resolve_binary("./bazel-bin/src/python/grpcio_tests/tests/interop/client"),
resolve_binary(
"./bazel-bin/src/python/grpcio_tests/tests/interop/client"
),
"--server_host=localhost",
f"--server_port={server_port}",
"--test_case=empty_unary",
Expand Down Expand Up @@ -627,6 +660,7 @@ def main():
print("Terminating server...")
try:
import signal

server_proc.send_signal(signal.SIGINT)
server_proc.wait(timeout=4)
except Exception:
Expand Down