diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3d91b734b..ad29226f5 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -80,10 +80,10 @@ jobs: mise exec -- uv run --frozen --project . --no-sync ruff check --select E,F,I,UP,B tools/ci mise exec -- uv run --frozen --project . --no-sync ruff check --select E9,F63,F7,F82 tools/mise_tasks mise exec -- uv run --frozen --project . --no-sync pytest -q tools/ci/tests --ignore tools/ci/tests/test_required_ci.py --ignore tools/ci/tests/test_public_tree.py --ignore tools/ci/tests/test_helm_render.py - - name: Pinned collector trace privacy + - name: Pinned collector privacy run: | mise run helm -- dependencies - mise exec -- uv run --frozen --project . --no-sync pytest -q -m docker tools/ci/tests/test_collector_trace_privacy.py + mise exec -- uv run --frozen --project . --no-sync pytest -q -m docker tools/ci/tests/test_collector_*privacy.py typescript: name: CI / TypeScript diff --git a/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl b/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl index 26206f607..8a0da8a69 100644 --- a/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl +++ b/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl @@ -295,6 +295,105 @@ processors: - 'resource.attributes["service.name"] == "sie-worker" and not IsMatch(name, "^(sie[.]worker[.]queue[.]duration|sie[.]worker[.]queue[.]depth|sie[.]worker[.]batch[.]size|sie[.]worker[.]batch[.]cost|sie[.]worker[.]batch[.]fill_ratio|sie[.]worker[.]queue[.]pending_at_dispatch|sie[.]worker[.]scheduler[.]adaptive[.]wait|sie[.]worker[.]scheduler[.]adaptive[.]cost|sie[.]worker[.]scheduler[.]adaptive[.]p50|sie[.]worker[.]scheduler[.]starvation[.]resets|sie[.]worker[.]runtime[.]batch[.]size|sie[.]worker[.]runtime[.]batch[.]subgroups|sie[.]worker[.]runtime[.]subgroup[.]size|sie[.]worker[.]requests|sie[.]worker[.]request[.]duration|sie[.]worker[.]inference[.]duration|sie[.]worker[.]units|sie[.]worker[.]model[.]loaded|sie[.]worker[.]model[.]load[.]duration|sie[.]worker[.]model[.]memory|sie[.]worker[.]oom[.]recoveries|sie[.]worker[.]model[.]evictions|sie[.]worker[.]generation[.]worker_wait|sie[.]worker[.]generation[.]ttft|sie[.]worker[.]generation[.]tpot|sie[.]worker[.]generation[.]tokens|sie[.]worker[.]generation[.]inflight|sie[.]worker[.]generation[.]kv[.]reserved|sie[.]worker[.]generation[.]kv[.]budget|sie[.]worker[.]generation[.]admission[.]decisions|sie[.]worker[.]generation[.]duplicate_prevented|sie[.]worker[.]generation[.]grammar[.]compile[.]duration|sie[.]worker[.]generation[.]grammar[.]cache[.]lookups|sie[.]worker[.]generation[.]grammar[.]requests|sie[.]worker[.]runtime[.]forward[.]duration|sie[.]worker[.]runtime[.]forward[.]permit[.]wait|sie[.]worker[.]runtime[.]forward[.]concurrent|sie[.]worker[.]runtime[.]forward[.]limit)$")' {{- end }} {{- if $betterStack.enabled }} + # Scalar lookups select one value; keep_keys alone retains duplicate OTLP + # keys. Rebuild only on remote export after per-metric field pruning. + transform/remote_metric_scalars: + error_mode: propagate + metric_statements: + - context: resource + statements: + - keep_keys(cache, []) + - 'set(cache["service.name"], attributes["service.name"]) where IsString(attributes["service.name"])' + - 'set(cache["service.instance.id"], attributes["service.instance.id"]) where IsString(attributes["service.instance.id"])' + - 'set(cache["deployment.environment"], attributes["deployment.environment"]) where IsString(attributes["deployment.environment"])' + - 'set(cache["cloud.region"], attributes["cloud.region"]) where IsString(attributes["cloud.region"])' + - 'set(cache["service.version"], attributes["service.version"]) where IsString(attributes["service.version"])' + - keep_keys(attributes, []) + - set(attributes["service.name"], cache["service.name"]) + - set(attributes["service.instance.id"], cache["service.instance.id"]) + - set(attributes["deployment.environment"], cache["deployment.environment"]) + - set(attributes["cloud.region"], cache["cloud.region"]) + - set(attributes["service.version"], cache["service.version"]) + - context: datapoint + statements: + - keep_keys(cache, []) + - 'set(cache["backend"], attributes["backend"]) where IsString(attributes["backend"]) or IsInt(attributes["backend"]) or IsDouble(attributes["backend"]) or IsBool(attributes["backend"])' + - 'set(cache["bundle"], attributes["bundle"]) where IsString(attributes["bundle"]) or IsInt(attributes["bundle"]) or IsDouble(attributes["bundle"]) or IsBool(attributes["bundle"])' + - 'set(cache["dispatch.path"], attributes["dispatch.path"]) where IsString(attributes["dispatch.path"]) or IsInt(attributes["dispatch.path"]) or IsDouble(attributes["dispatch.path"]) or IsBool(attributes["dispatch.path"])' + - 'set(cache["event"], attributes["event"]) where IsString(attributes["event"]) or IsInt(attributes["event"]) or IsDouble(attributes["event"]) or IsBool(attributes["event"])' + - 'set(cache["fallback.reason"], attributes["fallback.reason"]) where IsString(attributes["fallback.reason"]) or IsInt(attributes["fallback.reason"]) or IsDouble(attributes["fallback.reason"]) or IsBool(attributes["fallback.reason"])' + - 'set(cache["flush.reason"], attributes["flush.reason"]) where IsString(attributes["flush.reason"]) or IsInt(attributes["flush.reason"]) or IsDouble(attributes["flush.reason"]) or IsBool(attributes["flush.reason"])' + - 'set(cache["gpu_class"], attributes["gpu_class"]) where IsString(attributes["gpu_class"]) or IsInt(attributes["gpu_class"]) or IsDouble(attributes["gpu_class"]) or IsBool(attributes["gpu_class"])' + - 'set(cache["grammar"], attributes["grammar"]) where IsString(attributes["grammar"]) or IsInt(attributes["grammar"]) or IsDouble(attributes["grammar"]) or IsBool(attributes["grammar"])' + - 'set(cache["grammar.backend"], attributes["grammar.backend"]) where IsString(attributes["grammar.backend"]) or IsInt(attributes["grammar.backend"]) or IsDouble(attributes["grammar.backend"]) or IsBool(attributes["grammar.backend"])' + - 'set(cache["http.method"], attributes["http.method"]) where IsString(attributes["http.method"]) or IsInt(attributes["http.method"]) or IsDouble(attributes["http.method"]) or IsBool(attributes["http.method"])' + - 'set(cache["http.route"], attributes["http.route"]) where IsString(attributes["http.route"]) or IsInt(attributes["http.route"]) or IsDouble(attributes["http.route"]) or IsBool(attributes["http.route"])' + - 'set(cache["http.status_code"], attributes["http.status_code"]) where IsString(attributes["http.status_code"]) or IsInt(attributes["http.status_code"]) or IsDouble(attributes["http.status_code"]) or IsBool(attributes["http.status_code"])' + - 'set(cache["input.source"], attributes["input.source"]) where IsString(attributes["input.source"]) or IsInt(attributes["input.source"]) or IsDouble(attributes["input.source"]) or IsBool(attributes["input.source"])' + - 'set(cache["kind"], attributes["kind"]) where IsString(attributes["kind"]) or IsInt(attributes["kind"]) or IsDouble(attributes["kind"]) or IsBool(attributes["kind"])' + - 'set(cache["lane"], attributes["lane"]) where IsString(attributes["lane"]) or IsInt(attributes["lane"]) or IsDouble(attributes["lane"]) or IsBool(attributes["lane"])' + - 'set(cache["machine_profile"], attributes["machine_profile"]) where IsString(attributes["machine_profile"]) or IsInt(attributes["machine_profile"]) or IsDouble(attributes["machine_profile"]) or IsBool(attributes["machine_profile"])' + - 'set(cache["method"], attributes["method"]) where IsString(attributes["method"]) or IsInt(attributes["method"]) or IsDouble(attributes["method"]) or IsBool(attributes["method"])' + - 'set(cache["mode"], attributes["mode"]) where IsString(attributes["mode"]) or IsInt(attributes["mode"]) or IsDouble(attributes["mode"]) or IsBool(attributes["mode"])' + - 'set(cache["model"], attributes["model"]) where IsString(attributes["model"]) or IsInt(attributes["model"]) or IsDouble(attributes["model"]) or IsBool(attributes["model"])' + - 'set(cache["operation"], attributes["operation"]) where IsString(attributes["operation"]) or IsInt(attributes["operation"]) or IsDouble(attributes["operation"]) or IsBool(attributes["operation"])' + - 'set(cache["outcome"], attributes["outcome"]) where IsString(attributes["outcome"]) or IsInt(attributes["outcome"]) or IsDouble(attributes["outcome"]) or IsBool(attributes["outcome"])' + - 'set(cache["output.path"], attributes["output.path"]) where IsString(attributes["output.path"]) or IsInt(attributes["output.path"]) or IsDouble(attributes["output.path"]) or IsBool(attributes["output.path"])' + - 'set(cache["phase"], attributes["phase"]) where IsString(attributes["phase"]) or IsInt(attributes["phase"]) or IsDouble(attributes["phase"]) or IsBool(attributes["phase"])' + - 'set(cache["pool"], attributes["pool"]) where IsString(attributes["pool"]) or IsInt(attributes["pool"]) or IsDouble(attributes["pool"]) or IsBool(attributes["pool"])' + - 'set(cache["profile"], attributes["profile"]) where IsString(attributes["profile"]) or IsInt(attributes["profile"]) or IsDouble(attributes["profile"]) or IsBool(attributes["profile"])' + - 'set(cache["reason"], attributes["reason"]) where IsString(attributes["reason"]) or IsInt(attributes["reason"]) or IsDouble(attributes["reason"]) or IsBool(attributes["reason"])' + - 'set(cache["redelivered"], attributes["redelivered"]) where IsString(attributes["redelivered"]) or IsInt(attributes["redelivered"]) or IsDouble(attributes["redelivered"]) or IsBool(attributes["redelivered"])' + - 'set(cache["result"], attributes["result"]) where IsString(attributes["result"]) or IsInt(attributes["result"]) or IsDouble(attributes["result"]) or IsBool(attributes["result"])' + - 'set(cache["scaling_action"], attributes["scaling_action"]) where IsString(attributes["scaling_action"]) or IsInt(attributes["scaling_action"]) or IsDouble(attributes["scaling_action"]) or IsBool(attributes["scaling_action"])' + - 'set(cache["source"], attributes["source"]) where IsString(attributes["source"]) or IsInt(attributes["source"]) or IsDouble(attributes["source"]) or IsBool(attributes["source"])' + - 'set(cache["stage"], attributes["stage"]) where IsString(attributes["stage"]) or IsInt(attributes["stage"]) or IsDouble(attributes["stage"]) or IsBool(attributes["stage"])' + - 'set(cache["state"], attributes["state"]) where IsString(attributes["state"]) or IsInt(attributes["state"]) or IsDouble(attributes["state"]) or IsBool(attributes["state"])' + - 'set(cache["strategy"], attributes["strategy"]) where IsString(attributes["strategy"]) or IsInt(attributes["strategy"]) or IsDouble(attributes["strategy"]) or IsBool(attributes["strategy"])' + - 'set(cache["surface"], attributes["surface"]) where IsString(attributes["surface"]) or IsInt(attributes["surface"]) or IsDouble(attributes["surface"]) or IsBool(attributes["surface"])' + - 'set(cache["token.kind"], attributes["token.kind"]) where IsString(attributes["token.kind"]) or IsInt(attributes["token.kind"]) or IsDouble(attributes["token.kind"]) or IsBool(attributes["token.kind"])' + - 'set(cache["token.type"], attributes["token.type"]) where IsString(attributes["token.type"]) or IsInt(attributes["token.type"]) or IsDouble(attributes["token.type"]) or IsBool(attributes["token.type"])' + - 'set(cache["transport"], attributes["transport"]) where IsString(attributes["transport"]) or IsInt(attributes["transport"]) or IsDouble(attributes["transport"]) or IsBool(attributes["transport"])' + - 'set(cache["unit.type"], attributes["unit.type"]) where IsString(attributes["unit.type"]) or IsInt(attributes["unit.type"]) or IsDouble(attributes["unit.type"]) or IsBool(attributes["unit.type"])' + - keep_keys(attributes, []) + - set(attributes["backend"], cache["backend"]) + - set(attributes["bundle"], cache["bundle"]) + - set(attributes["dispatch.path"], cache["dispatch.path"]) + - set(attributes["event"], cache["event"]) + - set(attributes["fallback.reason"], cache["fallback.reason"]) + - set(attributes["flush.reason"], cache["flush.reason"]) + - set(attributes["gpu_class"], cache["gpu_class"]) + - set(attributes["grammar"], cache["grammar"]) + - set(attributes["grammar.backend"], cache["grammar.backend"]) + - set(attributes["http.method"], cache["http.method"]) + - set(attributes["http.route"], cache["http.route"]) + - set(attributes["http.status_code"], cache["http.status_code"]) + - set(attributes["input.source"], cache["input.source"]) + - set(attributes["kind"], cache["kind"]) + - set(attributes["lane"], cache["lane"]) + - set(attributes["machine_profile"], cache["machine_profile"]) + - set(attributes["method"], cache["method"]) + - set(attributes["mode"], cache["mode"]) + - set(attributes["model"], cache["model"]) + - set(attributes["operation"], cache["operation"]) + - set(attributes["outcome"], cache["outcome"]) + - set(attributes["output.path"], cache["output.path"]) + - set(attributes["phase"], cache["phase"]) + - set(attributes["pool"], cache["pool"]) + - set(attributes["profile"], cache["profile"]) + - set(attributes["reason"], cache["reason"]) + - set(attributes["redelivered"], cache["redelivered"]) + - set(attributes["result"], cache["result"]) + - set(attributes["scaling_action"], cache["scaling_action"]) + - set(attributes["source"], cache["source"]) + - set(attributes["stage"], cache["stage"]) + - set(attributes["state"], cache["state"]) + - set(attributes["strategy"], cache["strategy"]) + - set(attributes["surface"], cache["surface"]) + - set(attributes["token.kind"], cache["token.kind"]) + - set(attributes["token.type"], cache["token.type"]) + - set(attributes["transport"], cache["transport"]) + - set(attributes["unit.type"], cache["unit.type"]) # Collector implementation health is isolated from the application # contract and reduced to nine stable families before remote export. filter/collector_self_contract: @@ -560,11 +659,11 @@ service: {{- if $betterStack.enabled }} metrics/betterstack/gateway: receivers: [otlp/gateway] - processors: [memory_limiter, filter/remote_gateway_contract, resource/gateway_identity, transform/contract_metrics, transform/remote_queue_identity, batch] + processors: [memory_limiter, filter/remote_gateway_contract, resource/gateway_identity, transform/contract_metrics, transform/remote_metric_scalars, transform/remote_queue_identity, batch] exporters: [otlphttp/betterstack] metrics/betterstack/application: receivers: [otlp/application] - processors: [memory_limiter, filter/remote_application_contract, resource/application_identity, transform/contract_metrics, batch] + processors: [memory_limiter, filter/remote_application_contract, resource/application_identity, transform/contract_metrics, transform/remote_metric_scalars, batch] exporters: [otlphttp/betterstack] {{- end }} {{- end }} diff --git a/telemetry/README.md b/telemetry/README.md index 8fa4a271d..0292f3b1a 100644 --- a/telemetry/README.md +++ b/telemetry/README.md @@ -634,3 +634,14 @@ CI should reject a change unless it proves all of the following: 10. an end-to-end KEDA signal respects the declared five-second OTLP export and Prometheus scrape budgets, and collector, producer-export, scrape, or query failure activates the worker lane's declared safe fallback. + +Remote metric maps are reconstructed after the existing per-metric attribute +allowlist. Each retained resource key has at most one string value; each retained +point key has at most one string, integer, double or boolean value. The first +lookup value wins, matching the value used by preceding filters. Missing or +non-scalar values are omitted, and scratch state is cleared for every resource +and point. This removes duplicate protobuf keys without changing metric values, +histograms, timestamps, temporality or producer domain policies. Local Prometheus +processing remains unchanged. The pinned collector regression sends raw duplicate +keys through both rendered receiver branches, including reversed service claims, +nested values and missing optional fields on successive points. diff --git a/telemetry/contract.yaml b/telemetry/contract.yaml index 31bb6ed1b..f07aadfb4 100644 --- a/telemetry/contract.yaml +++ b/telemetry/contract.yaml @@ -596,6 +596,14 @@ metric_processing: datapoint_attributes: policy: exact_metric_attribute_set unknown: drop + remote_scalar_reconstruction: + ordering: after_exact_metric_attribute_set_before_queue_identity + resource_types: [string] + datapoint_types: [string, int, double, bool] + duplicates: first_lookup_value_only + non_scalar_or_missing: omit + cache: clear_for_each_resource_and_datapoint + local_prometheus_processing: unchanged exemplars: policy: forbidden producer_filter: always_off diff --git a/tools/ci/tests/test_collector_metric_privacy.py b/tools/ci/tests/test_collector_metric_privacy.py new file mode 100644 index 000000000..88a350c70 --- /dev/null +++ b/tools/ci/tests/test_collector_metric_privacy.py @@ -0,0 +1,222 @@ +"""Raw OTLP duplicate-key regression against the rendered pinned collector.""" + +from __future__ import annotations + +import json +import urllib.request +from pathlib import Path + +import pytest +import yaml +from opentelemetry.proto.collector.metrics.v1.metrics_service_pb2 import ExportMetricsServiceRequest +from opentelemetry.proto.metrics.v1.metrics_pb2 import AGGREGATION_TEMPORALITY_DELTA + +from tools.ci.tests.test_collector_trace_privacy import SENTINEL, run, wait_for + +pytestmark = pytest.mark.docker + + +def rendered_config(): + values = { + "payloadStore.enabled": "false", + "observability.otel.metrics.enabled": "true", + "observability.otel.collector.install": "true", + "observability.otel.collector.betterStack.enabled": "true", + "observability.otel.collector.betterStack.endpoint": "https://example.invalid", + "observability.otel.collector.betterStack.existingSecret": "synthetic", + "observability.otel.resource.deploymentEnvironment": "dev", + "observability.otel.resource.cloudRegion": "test-region", + } + documents = list( + yaml.safe_load_all( + run( + "mise", + "exec", + "--", + "helm", + "template", + "sie", + "deploy/helm/sie-cluster", + "--namespace", + "sie", + "--show-only", + "templates/otel-collector.yaml", + *(arg for key, value in values.items() for arg in ["--set", f"{key}={value}"]), + ) + ) + ) + config = next(d for d in documents if d and d.get("kind") == "ConfigMap") + deployment = next(d for d in documents if d and d.get("kind") == "Deployment") + return yaml.safe_load(config["data"]["collector.yaml"]), deployment["spec"]["template"]["spec"]["containers"][0][ + "image" + ] + + +def add(attributes, key, value): + attr = attributes.add(key=key) + if isinstance(value, str): + attr.value.string_value = value + elif isinstance(value, int): + attr.value.int_value = value + else: + attr.value.kvlist_value.values.add(key="payload").value.string_value = SENTINEL + + +def payload(receiver): + wire = ExportMetricsServiceRequest() + service = "sie-gateway" if receiver == "gateway" else "sie-worker" + prefix = "sie.gateway" if receiver == "gateway" else "sie.worker" + for reverse_identity in [False, True]: + rs = wire.resource_metrics.add() + for key, value in { + "service.name": service, + "service.instance.id": "instance", + "service.version": "version", + "deployment.environment": "producer-env", + "cloud.region": "producer-region", + }.items(): + if key == "service.name" and reverse_identity: + add(rs.resource.attributes, key, SENTINEL) + add(rs.resource.attributes, key, value) + add(rs.resource.attributes, key, SENTINEL) + add(rs.resource.attributes, key, None) + scope = rs.scope_metrics.add() + for suffix, kind in [("requests", "sum"), ("request.duration", "histogram")]: + metric = scope.metrics.add(name=f"{prefix}.{suffix}") + data = getattr(metric, kind) + data.aggregation_temporality = AGGREGATION_TEMPORALITY_DELTA + if kind == "sum": + data.is_monotonic = True + for case in range(3): + point = data.data_points.add(start_time_unix_nano=1_000, time_unix_nano=2_000 + case) + if kind == "sum": + point.as_int = 10 + case + else: + point.count = 2 + point.sum = 0.5 + case + point.explicit_bounds.extend([1.0]) + point.bucket_counts.extend([1, 1]) + add(point.attributes, "outcome", "success") + add(point.attributes, "outcome", None) + if case == 0: + add(point.attributes, "operation", "encode") + add(point.attributes, "operation", SENTINEL) + add(point.attributes, "operation", None) + elif case == 2: + add(point.attributes, "operation", None) + add(point.attributes, "operation", "encode") + if receiver == "gateway": + add(point.attributes, "http.status_code", 200) + add(point.attributes, "http.status_code", None) + add(point.attributes, "cloud.region", SENTINEL) + return wire.SerializeToString() + + +def batches(path: Path): + return [json.loads(line) for line in path.read_text().splitlines()] + + +@pytest.mark.parametrize("receiver", ["gateway", "application"]) +def test_remote_metric_maps_have_one_scalar_per_retained_key(tmp_path, receiver): + config, image = rendered_config() + assert image.endswith(":0.119.0") + processors = config["processors"] + shared = processors["transform/contract_metrics"]["metric_statements"] + shared_points = next(group["statements"] for group in shared if group["context"] == "datapoint") + allowed = { + key + for statement in shared_points + for key in json.loads(statement.split("keep_keys(attributes, ", 1)[1].split(") where", 1)[0]) + } + remote_groups = processors["transform/remote_metric_scalars"]["metric_statements"] + remote_points = next(group["statements"] for group in remote_groups if group["context"] == "datapoint") + rebuilt = { + statement.split('set(attributes["', 1)[1].split('"]', 1)[0] + for statement in remote_points + if statement.startswith('set(attributes["') + } + assert rebuilt == allowed + for name, pipeline in config["service"]["pipelines"].items(): + if name.startswith("metrics/prometheus/"): + assert "transform/remote_metric_scalars" not in pipeline["processors"] + config["exporters"] = { + f"file/{name}": {"path": f"/out/{name}.json", "flush_interval": "100ms"} for name in ["local", "remote", "self"] + } + for name, pipeline in config["service"]["pipelines"].items(): + sink = "self" if name == "metrics/self" else "remote" if "betterstack" in name else "local" + pipeline["exporters"] = [f"file/{sink}"] + config["receivers"][f"otlp/{receiver}"]["protocols"]["http"] = {"endpoint": "0.0.0.0:4338"} + config["processors"]["batch"]["timeout"] = "100ms" + (tmp_path / "collector.yaml").write_text(yaml.safe_dump(config)) + container = run( + "docker", + "run", + "-d", + "--user", + "0", + "-p", + "127.0.0.1::4338", + "-p", + "127.0.0.1::13133", + "-v", + f"{tmp_path}:/out", + image, + "--config=/out/collector.yaml", + ) + try: + ports = {p: run("docker", "port", container, f"{p}/tcp").split(":")[-1] for p in [4338, 13133]} + wait_for(lambda: urllib.request.urlopen(f"http://127.0.0.1:{ports[13133]}", timeout=1).status == 200) + request = urllib.request.Request( + f"http://127.0.0.1:{ports[4338]}/v1/metrics", + data=payload(receiver), + headers={"Content-Type": "application/x-protobuf"}, + ) + with urllib.request.urlopen(request, timeout=5) as response: # noqa: S310 - fixed loopback HTTP + assert response.status == 200 + remote = wait_for(lambda: batches(tmp_path / "remote.json")) + wait_for(lambda: batches(tmp_path / "local.json")) + assert SENTINEL in (tmp_path / "local.json").read_text() + assert SENTINEL not in (tmp_path / "remote.json").read_text() + seen = [] + for batch in remote: + for rs in batch["resourceMetrics"]: + attrs = rs["resource"]["attributes"] + assert len(attrs) == 5 + assert {a["key"]: a["value"] for a in attrs} == { + "service.name": {"stringValue": "sie-gateway" if receiver == "gateway" else "sie-worker"}, + "service.instance.id": {"stringValue": "instance"}, + "service.version": {"stringValue": "version"}, + "deployment.environment": {"stringValue": "dev"}, + "cloud.region": {"stringValue": "test-region"}, + } + for scope in rs["scopeMetrics"]: + for metric in scope["metrics"]: + kind = "sum" if "sum" in metric else "histogram" + data = metric[kind] + assert data["aggregationTemporality"] == 1 + for point in data["dataPoints"]: + seen.append(point) + case = int(point["timeUnixNano"]) - 2_000 + expected = {"outcome": {"stringValue": "success"}} + if case == 0: + expected["operation"] = {"stringValue": "encode"} + if receiver == "gateway": + expected["http.status_code"] = {"intValue": "200"} + assert len(point["attributes"]) == len(expected) + assert {a["key"]: a["value"] for a in point["attributes"]} == expected + assert int(point["startTimeUnixNano"]) == 1_000 + if kind == "sum": + assert int(point["asInt"]) == 10 + case + else: + assert int(point["count"]) == 2 + assert point["sum"] == 0.5 + case + assert point["bucketCounts"] == ["1", "1"] + assert point["explicitBounds"] == [1] + assert len(seen) == 6 + finally: + try: + run("docker", "stop", "--time", "10", container) + logs = run("docker", "logs", container, include_stderr=True) + assert "Error: " not in logs, logs + finally: + run("docker", "rm", "-f", container)