diff --git a/.github/scripts/deploy_observability_collector.sh b/.github/scripts/deploy_observability_collector.sh new file mode 100755 index 000000000..e2ec1f45d --- /dev/null +++ b/.github/scripts/deploy_observability_collector.sh @@ -0,0 +1,148 @@ +#!/usr/bin/env bash + +set -euo pipefail + +required=( + OBSERVABILITY_PROJECT_ID + OBSERVABILITY_COLLECTOR_REGION + OBSERVABILITY_COLLECTOR_SERVICE + OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY + OBSERVABILITY_COLLECTOR_IMAGE_NAME + OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT + OBSERVABILITY_METRIC_LOCATION + GITHUB_SHA +) + +missing=() +for name in "${required[@]}"; do + if [[ -z "${!name:-}" ]]; then + missing+=("${name}") + fi +done +if (( ${#missing[@]} > 0 )); then + echo "Missing collector deployment configuration: ${missing[*]}" >&2 + exit 1 +fi + +[[ "${OBSERVABILITY_PROJECT_ID}" =~ ^[a-z][a-z0-9-]{4,28}[a-z0-9]$ ]] \ + || { echo "OBSERVABILITY_PROJECT_ID is invalid." >&2; exit 1; } +[[ "${OBSERVABILITY_COLLECTOR_REGION}" =~ ^[a-z]+-[a-z]+[0-9]+$ ]] \ + || { echo "OBSERVABILITY_COLLECTOR_REGION is invalid." >&2; exit 1; } +[[ "${OBSERVABILITY_COLLECTOR_SERVICE}" =~ ^[a-z]([a-z0-9-]{0,61}[a-z0-9])?$ ]] \ + || { echo "OBSERVABILITY_COLLECTOR_SERVICE is invalid." >&2; exit 1; } +[[ "${OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY}" =~ ^[a-z][a-z0-9._-]{0,62}$ ]] \ + || { echo "OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY is invalid." >&2; exit 1; } +[[ "${OBSERVABILITY_COLLECTOR_IMAGE_NAME}" =~ ^[a-z0-9][a-z0-9._-]{0,127}$ ]] \ + || { echo "OBSERVABILITY_COLLECTOR_IMAGE_NAME is invalid." >&2; exit 1; } +[[ "${OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT}" =~ ^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.gserviceaccount\.com$ ]] \ + || { echo "OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT is invalid." >&2; exit 1; } +[[ "${OBSERVABILITY_METRIC_LOCATION}" =~ ^[a-z]+-[a-z]+[0-9]+$ ]] \ + || { echo "OBSERVABILITY_METRIC_LOCATION is invalid." >&2; exit 1; } +[[ "${GITHUB_SHA}" =~ ^[0-9a-f]{40}$ ]] \ + || { echo "GITHUB_SHA must be a 40-character lowercase commit SHA." >&2; exit 1; } + +registry="${OBSERVABILITY_COLLECTOR_REGION}-docker.pkg.dev" +image_uri="${registry}/${OBSERVABILITY_PROJECT_ID}/${OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY}/${OBSERVABILITY_COLLECTOR_IMAGE_NAME}:${GITHUB_SHA}" + +gcloud artifacts repositories describe \ + "${OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --location "${OBSERVABILITY_COLLECTOR_REGION}" >/dev/null +gcloud auth configure-docker "${registry}" --quiet +docker build \ + --platform linux/amd64 \ + --file gcp/observability/collector/Dockerfile \ + --tag "${image_uri}" \ + gcp/observability/collector +docker push "${image_uri}" + +gcloud run deploy "${OBSERVABILITY_COLLECTOR_SERVICE}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --region "${OBSERVABILITY_COLLECTOR_REGION}" \ + --platform managed \ + --image "${image_uri}" \ + --service-account "${OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT}" \ + --invoker-iam-check \ + --execution-environment gen2 \ + --port 8080 \ + --use-http2 \ + --cpu 1 \ + --no-cpu-throttling \ + --cpu-boost \ + --memory 512Mi \ + --timeout 30 \ + --min 1 \ + --max 10 \ + --concurrency 100 \ + --startup-probe \ + 'httpGet.path=/,httpGet.port=13133,periodSeconds=2,failureThreshold=30,timeoutSeconds=1' \ + --liveness-probe \ + 'httpGet.path=/,httpGet.port=13133,periodSeconds=30,failureThreshold=3,timeoutSeconds=2' \ + --set-env-vars \ + "OBSERVABILITY_PROJECT_ID=${OBSERVABILITY_PROJECT_ID},OBSERVABILITY_METRIC_LOCATION=${OBSERVABILITY_METRIC_LOCATION}" \ + --quiet + +service_json="$( + gcloud run services describe "${OBSERVABILITY_COLLECTOR_SERVICE}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --region "${OBSERVABILITY_COLLECTOR_REGION}" \ + --format=json +)" +revision="$(jq -r '.status.latestReadyRevisionName // ""' <<< "${service_json}")" +service_url="$(jq -r '.status.url // ""' <<< "${service_json}")" +ready="$( + jq -r \ + '[.status.conditions[]? | select(.type == "Ready") | .status][0] // ""' \ + <<< "${service_json}" +)" + +for public_member in allUsers allAuthenticatedUsers; do + if gcloud run services get-iam-policy \ + "${OBSERVABILITY_COLLECTOR_SERVICE}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --region "${OBSERVABILITY_COLLECTOR_REGION}" \ + --flatten='bindings[].members' \ + --filter="bindings.role=roles/run.invoker AND bindings.members=${public_member}" \ + --format='value(bindings.members)' | grep -qx "${public_member}"; then + gcloud run services remove-iam-policy-binding \ + "${OBSERVABILITY_COLLECTOR_SERVICE}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --region "${OBSERVABILITY_COLLECTOR_REGION}" \ + --member="${public_member}" \ + --role=roles/run.invoker \ + --quiet + fi +done + +public_members="$( + gcloud run services get-iam-policy \ + "${OBSERVABILITY_COLLECTOR_SERVICE}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --region "${OBSERVABILITY_COLLECTOR_REGION}" \ + --flatten='bindings[].members' \ + --filter='bindings.members:(allUsers OR allAuthenticatedUsers)' \ + --format='value(bindings.members)' +)" +invoker_iam_disabled="$( + gcloud run services describe "${OBSERVABILITY_COLLECTOR_SERVICE}" \ + --project "${OBSERVABILITY_PROJECT_ID}" \ + --region "${OBSERVABILITY_COLLECTOR_REGION}" \ + --format="value(metadata.annotations.'run.googleapis.com/invoker-iam-disabled')" +)" + +if [[ -z "${revision}" || -z "${service_url}" || "${ready}" != "True" ]]; then + echo "Collector deployment did not produce a healthy ready revision and URL." >&2 + exit 1 +fi +if [[ -n "${public_members}" || "${invoker_iam_disabled}" == "true" ]]; then + echo "Collector Cloud Run service permits unauthenticated invocation." >&2 + exit 1 +fi + +if [[ -n "${GITHUB_OUTPUT:-}" ]]; then + echo "revision=${revision}" >> "${GITHUB_OUTPUT}" + echo "service_url=${service_url}" >> "${GITHUB_OUTPUT}" + echo "image_uri=${image_uri}" >> "${GITHUB_OUTPUT}" +fi + +printf 'Deployed collector revision %s at %s\n' "${revision}" "${service_url}" diff --git a/.github/scripts/verify_observability_collector.py b/.github/scripts/verify_observability_collector.py new file mode 100644 index 000000000..8156a70bc --- /dev/null +++ b/.github/scripts/verify_observability_collector.py @@ -0,0 +1,392 @@ +"""Verify authenticated collector ingestion in Google Cloud.""" + +from __future__ import annotations + +import argparse +import json +import os +import time +import uuid +from datetime import UTC, datetime +from typing import NamedTuple +from urllib.error import HTTPError +from urllib.parse import quote, urlencode, urlparse +from urllib.request import Request, urlopen + +import grpc +from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ( + ExportLogsServiceRequest, +) +from opentelemetry.proto.collector.logs.v1.logs_service_pb2_grpc import ( + LogsServiceStub, +) +from opentelemetry.proto.collector.metrics.v1.metrics_service_pb2 import ( + ExportMetricsServiceRequest, +) +from opentelemetry.proto.collector.metrics.v1.metrics_service_pb2_grpc import ( + MetricsServiceStub, +) +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ( + ExportTraceServiceRequest, +) +from opentelemetry.proto.collector.trace.v1.trace_service_pb2_grpc import ( + TraceServiceStub, +) +from opentelemetry.proto.common.v1.common_pb2 import ( + AnyValue, + InstrumentationScope, + KeyValue, +) +from opentelemetry.proto.logs.v1.logs_pb2 import ( + LogRecord, + ResourceLogs, + ScopeLogs, +) +from opentelemetry.proto.metrics.v1.metrics_pb2 import ( + Gauge, + Metric, + NumberDataPoint, + ResourceMetrics, + ScopeMetrics, +) +from opentelemetry.proto.resource.v1.resource_pb2 import Resource +from opentelemetry.proto.trace.v1.trace_pb2 import ( + ResourceSpans, + ScopeSpans, + Span, +) + +METRIC_NAME = "policyengine.collector.verification" +METRIC_TYPE = f"prometheus.googleapis.com/{METRIC_NAME}/gauge" +SPAN_NAME = "policyengine.collector.verification" +SERVICE_NAME = "policyengine-observability-verifier" +SERVICE_NAMESPACE = "policyengine.api-v1" + + +class VerificationProbe(NamedTuple): + verification_id: str + trace_id: str + started_at: float + + +def _authority(endpoint: str) -> str: + parsed = urlparse(endpoint) + if parsed.scheme != "https" or not parsed.hostname: + raise ValueError("collector endpoint must be an HTTPS URL") + return f"{parsed.hostname}:{parsed.port or 443}" + + +def _attribute(key: str, value: str) -> KeyValue: + return KeyValue(key=key, value=AnyValue(string_value=value)) + + +def _resource(verification_id: str) -> Resource: + return Resource( + attributes=[ + _attribute("service.name", SERVICE_NAME), + _attribute("service.namespace", SERVICE_NAMESPACE), + _attribute("service.instance.id", verification_id), + _attribute("deployment.environment.name", "production"), + ] + ) + + +def _probe_requests( + verification_id: str, + timestamp_ns: int, +) -> tuple[ + ExportTraceServiceRequest, + ExportMetricsServiceRequest, + ExportLogsServiceRequest, + str, +]: + resource = _resource(verification_id) + scope = InstrumentationScope(name="policyengine.collector.verifier") + trace_id_bytes = uuid.uuid4().bytes + span = Span( + trace_id=trace_id_bytes, + span_id=uuid.uuid4().bytes[:8], + name=SPAN_NAME, + kind=Span.SPAN_KIND_INTERNAL, + start_time_unix_nano=timestamp_ns - 1_000_000, + end_time_unix_nano=timestamp_ns, + attributes=[_attribute("verification.id", verification_id)], + ) + trace_request = ExportTraceServiceRequest( + resource_spans=[ + ResourceSpans( + resource=resource, + scope_spans=[ScopeSpans(scope=scope, spans=[span])], + ) + ] + ) + metric_request = ExportMetricsServiceRequest( + resource_metrics=[ + ResourceMetrics( + resource=resource, + scope_metrics=[ + ScopeMetrics( + scope=scope, + metrics=[ + Metric( + name=METRIC_NAME, + description=("Collector deployment verification value"), + unit="1", + gauge=Gauge( + data_points=[ + NumberDataPoint( + time_unix_nano=timestamp_ns, + as_int=1, + attributes=[ + _attribute( + "verification.id", + verification_id, + ) + ], + ) + ] + ), + ) + ], + ) + ], + ) + ] + ) + log_request = ExportLogsServiceRequest( + resource_logs=[ + ResourceLogs( + resource=resource, + scope_logs=[ + ScopeLogs( + scope=scope, + log_records=[ + LogRecord( + time_unix_nano=timestamp_ns, + body=AnyValue( + string_value=( + "collector log rejection verification" + ) + ), + attributes=[ + _attribute("verification.id", verification_id) + ], + ) + ], + ) + ], + ) + ] + ) + return trace_request, metric_request, log_request, trace_id_bytes.hex() + + +def _raise_for_partial_rejection( + response: object, + *, + signal: str, + rejected_field: str, +) -> None: + partial = getattr(response, "partial_success", None) + rejected = getattr(partial, rejected_field, 0) + if rejected: + message = getattr(partial, "error_message", "") + raise RuntimeError(f"collector rejected {rejected} {signal}: {message}") + + +def _get_json(url: str, access_token: str) -> dict[str, object] | None: + request = Request( + url, + headers={"Authorization": f"Bearer {access_token}"}, + ) + try: + with urlopen(request, timeout=15) as response: + return json.load(response) + except HTTPError as error: + if error.code == 404: + return None + detail = error.read().decode("utf-8", errors="replace") + raise RuntimeError( + f"Google Cloud verification request failed ({error.code}): {detail}" + ) from error + + +def _trace_available( + project_id: str, + trace_id: str, + access_token: str, +) -> bool: + project = quote(project_id, safe="") + trace = quote(trace_id, safe="") + url = f"https://cloudtrace.googleapis.com/v1/projects/{project}/traces/{trace}" + payload = _get_json(url, access_token) + return payload is not None and payload.get("traceId") == trace_id + + +def _metric_available( + project_id: str, + verification_id: str, + metric_location: str, + access_token: str, + started_at: float, +) -> bool: + project = quote(project_id, safe="") + monitoring_filter = " AND ".join( + ( + f'metric.type = "{METRIC_TYPE}"', + 'resource.type = "prometheus_target"', + f'resource.labels.location = "{metric_location}"', + f'resource.labels.instance = "{verification_id}"', + ) + ) + start = datetime.fromtimestamp(started_at - 60, UTC).isoformat() + end = datetime.now(UTC).isoformat() + parameters = { + "filter": monitoring_filter, + "interval.startTime": start, + "interval.endTime": end, + "view": "FULL", + "pageSize": "1000", + } + base_url = f"https://monitoring.googleapis.com/v3/projects/{project}/timeSeries" + page_token: str | None = None + while True: + if page_token: + parameters["pageToken"] = page_token + url = f"{base_url}?{urlencode(parameters)}" + payload = _get_json(url, access_token) or {} + if payload.get("timeSeries"): + return True + raw_page_token = payload.get("nextPageToken") + page_token = raw_page_token if isinstance(raw_page_token, str) else None + if not page_token: + return False + + +def _wait_for_delivery( + *, + project_id: str, + trace_id: str, + verification_id: str, + metric_location: str, + access_token: str, + started_at: float, + timeout_seconds: float, + poll_seconds: float = 5.0, +) -> None: + deadline = time.monotonic() + timeout_seconds + trace_found = False + metric_found = False + while True: + if not trace_found: + trace_found = _trace_available(project_id, trace_id, access_token) + if not metric_found: + metric_found = _metric_available( + project_id, + verification_id, + metric_location, + access_token, + started_at, + ) + if trace_found and metric_found: + return + remaining = deadline - time.monotonic() + if remaining <= 0: + missing = ", ".join( + signal + for signal, found in ( + ("trace", trace_found), + ("metric", metric_found), + ) + if not found + ) + raise RuntimeError(f"collector verification data was not stored: {missing}") + time.sleep(min(poll_seconds, remaining)) + + +def verify( + endpoint: str, + token: str, + *, + project_id: str, + metric_location: str, + access_token: str, + timeout_seconds: float = 15.0, + delivery_timeout_seconds: float = 600.0, +) -> VerificationProbe: + """Export real telemetry, verify storage, and require log rejection.""" + + started_at = time.time() + verification_id = f"collector-verify-{uuid.uuid4().hex}" + trace_request, metric_request, log_request, trace_id = _probe_requests( + verification_id, + time.time_ns(), + ) + metadata = (("authorization", f"Bearer {token}"),) + with grpc.secure_channel( + _authority(endpoint), grpc.ssl_channel_credentials() + ) as channel: + trace_response = TraceServiceStub(channel).Export( + trace_request, + metadata=metadata, + timeout=timeout_seconds, + ) + _raise_for_partial_rejection( + trace_response, + signal="spans", + rejected_field="rejected_spans", + ) + metric_response = MetricsServiceStub(channel).Export( + metric_request, + metadata=metadata, + timeout=timeout_seconds, + ) + _raise_for_partial_rejection( + metric_response, + signal="metric points", + rejected_field="rejected_data_points", + ) + try: + LogsServiceStub(channel).Export( + log_request, + metadata=metadata, + timeout=timeout_seconds, + ) + except grpc.RpcError as error: + if error.code() is not grpc.StatusCode.UNIMPLEMENTED: + raise + else: + raise RuntimeError("collector unexpectedly accepted application logs") + _wait_for_delivery( + project_id=project_id, + trace_id=trace_id, + verification_id=verification_id, + metric_location=metric_location, + access_token=access_token, + started_at=started_at, + timeout_seconds=delivery_timeout_seconds, + ) + return VerificationProbe(verification_id, trace_id, started_at) + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--endpoint", required=True) + parser.add_argument("--project-id", required=True) + parser.add_argument("--metric-location", required=True) + args = parser.parse_args() + probe = verify( + args.endpoint, + os.environ["COLLECTOR_ID_TOKEN"], + project_id=args.project_id, + metric_location=args.metric_location, + access_token=os.environ["GOOGLE_OAUTH_ACCESS_TOKEN"], + ) + print( + "Verified collector ingestion " + f"(verification_id={probe.verification_id}, trace_id={probe.trace_id})" + ) + + +if __name__ == "__main__": + main() diff --git a/.github/scripts/verify_observability_collector.sh b/.github/scripts/verify_observability_collector.sh new file mode 100755 index 000000000..0b887d121 --- /dev/null +++ b/.github/scripts/verify_observability_collector.sh @@ -0,0 +1,21 @@ +#!/usr/bin/env bash + +set -euo pipefail + +endpoint="${1:-}" +if [[ -z "${endpoint}" ]]; then + echo "Usage: $0 COLLECTOR_ENDPOINT" >&2 + exit 1 +fi + +export COLLECTOR_ID_TOKEN +COLLECTOR_ID_TOKEN="$(gcloud auth print-identity-token --audiences="${endpoint}")" +export GOOGLE_OAUTH_ACCESS_TOKEN +GOOGLE_OAUTH_ACCESS_TOKEN="$(gcloud auth print-access-token)" +echo "::add-mask::${COLLECTOR_ID_TOKEN}" +echo "::add-mask::${GOOGLE_OAUTH_ACCESS_TOKEN}" +uv run --no-project --with grpcio --with opentelemetry-proto \ + python .github/scripts/verify_observability_collector.py \ + --endpoint "${endpoint}" \ + --project-id "${OBSERVABILITY_PROJECT_ID}" \ + --metric-location "${OBSERVABILITY_METRIC_LOCATION}" diff --git a/.github/workflows/deploy-observability-collector.yml b/.github/workflows/deploy-observability-collector.yml new file mode 100644 index 000000000..bc4dde1be --- /dev/null +++ b/.github/workflows/deploy-observability-collector.yml @@ -0,0 +1,59 @@ +name: Deploy observability collector + +on: + push: + branches: + - master + paths: + - "gcp/observability/collector/**" + - ".github/scripts/deploy_observability_collector.sh" + - ".github/scripts/verify_observability_collector.sh" + - ".github/scripts/verify_observability_collector.py" + - ".github/workflows/deploy-observability-collector.yml" + workflow_dispatch: + +concurrency: + group: observability-collector-production + cancel-in-progress: false + +jobs: + deploy: + name: Build, deploy, and verify collector + if: github.repository == 'PolicyEngine/policyengine-api' + runs-on: ubuntu-latest + environment: production + permissions: + contents: read + id-token: write + env: + OBSERVABILITY_PROJECT_ID: ${{ vars.OBSERVABILITY_PROJECT_ID }} + OBSERVABILITY_COLLECTOR_REGION: ${{ vars.OBSERVABILITY_COLLECTOR_REGION }} + OBSERVABILITY_COLLECTOR_SERVICE: ${{ vars.OBSERVABILITY_COLLECTOR_SERVICE }} + OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY: ${{ vars.OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY }} + OBSERVABILITY_COLLECTOR_IMAGE_NAME: ${{ vars.OBSERVABILITY_COLLECTOR_IMAGE_NAME }} + OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT: ${{ vars.OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT }} + OBSERVABILITY_METRIC_LOCATION: ${{ vars.OBSERVABILITY_METRIC_LOCATION }} + steps: + - name: Require master branch + if: github.ref != 'refs/heads/master' + run: exit 1 + - name: Checkout repository + uses: actions/checkout@v4 + - name: Validate collector configuration + run: docker run --rm --env OBSERVABILITY_PROJECT_ID="${OBSERVABILITY_PROJECT_ID}" --env OBSERVABILITY_METRIC_LOCATION="${OBSERVABILITY_METRIC_LOCATION}" --volume "${PWD}/gcp/observability/collector/config.yaml:/etc/otelcol-google/config.yaml:ro" us-docker.pkg.dev/cloud-ops-agents-artifacts/google-cloud-opentelemetry-collector/otelcol-google:0.160.0 validate --config=/etc/otelcol-google/config.yaml + - name: Authenticate to Google Cloud + uses: google-github-actions/auth@v2 + with: + workload_identity_provider: ${{ secrets.GCP_WORKLOAD_IDENTITY_PROVIDER }} + service_account: ${{ secrets.GCP_DEPLOY_SERVICE_ACCOUNT }} + - name: Set up Google Cloud CLI + uses: google-github-actions/setup-gcloud@v2 + - name: Set up uv + uses: astral-sh/setup-uv@v6 + with: + version: "0.12.1" + - name: Build and deploy collector + id: collector + run: bash .github/scripts/deploy_observability_collector.sh + - name: Verify Google Cloud ingestion and log rejection + run: bash .github/scripts/verify_observability_collector.sh '${{ steps.collector.outputs.service_url }}' diff --git a/changelog.d/3860.fixed.md b/changelog.d/3860.fixed.md new file mode 100644 index 000000000..914c4ed67 --- /dev/null +++ b/changelog.d/3860.fixed.md @@ -0,0 +1 @@ +Give every API process a unique metric resource identity and add repeatable deployment of the central OpenTelemetry Collector with a valid Google Monitoring location. diff --git a/docs/engineering/skills/observability.md b/docs/engineering/skills/observability.md index b44ba8242..d52f959f0 100644 --- a/docs/engineering/skills/observability.md +++ b/docs/engineering/skills/observability.md @@ -78,6 +78,20 @@ When adding a calculation configuration or stage: 3. Import that plan in runtime code. 4. Add focused tests for the stage and identifier lifecycle. +## Metric resource identity + +`service.instance.id` identifies one telemetry-producing process. Cloud Run +revision names identify deployed code and are shared by multiple containers and +Gunicorn workers, so they must not be used alone as the instance identifier. +Construct the value with `policyengine_observability.process_instance_id` after +the worker process starts. + +The central collector adds the configured Google Monitoring `location` only to +metrics. Logs and traces retain the workload's actual `cloud.region`. Deploy +collector configuration changes with the `Deploy observability collector` +workflow; committing `gcp/observability/collector/config.yaml` alone does not +change the live Cloud Run service. + ## Trace boundaries HTTP instrumentation carries W3C trace context across synchronous calls. The diff --git a/gcp/observability/README.md b/gcp/observability/README.md index e40a3acda..a7d384dd3 100644 --- a/gcp/observability/README.md +++ b/gcp/observability/README.md @@ -8,9 +8,9 @@ Collector: | `collector/config.yaml` | OTLP gRPC receiver, bounded processing, trace sampling, and Google Telemetry API export configuration | | `collector/Dockerfile` | Collector container image built with that configuration | -The Google Cloud IAM, logging sinks, collector service, dashboard, and alert -policies were provisioned separately. This repository does not manage or apply -those resources. +The Google Cloud IAM, logging sinks, dashboard, and alert policies were +provisioned separately. This repository owns the collector image, +configuration, and deployment workflow. ## Collector behavior @@ -25,10 +25,38 @@ retains errors, operations lasting at least 30 seconds, and currently 100% of all remaining traces. Participating SDKs therefore use 100% head sampling so the collector can evaluate complete traces. -Changing `collector/config.yaml` does not update the live service. The image -must be rebuilt and the existing `policyengine-api-v1-otel-collector` Cloud Run -service must be updated through a separately managed deployment process. No -collector deployment workflow exists in this repository. +The metrics pipeline assigns `location` from +`OBSERVABILITY_METRIC_LOCATION`. This is the Google Monitoring resource +location for centrally collected metrics; workload execution regions remain +available on logs and traces. + +Run the `Deploy observability collector` workflow from `master` after changing +the collector image or configuration. It validates the configuration, builds +an image tagged with the source commit, updates the authenticated +`policyengine-api-v1-otel-collector` Cloud Run service, waits for its health +probes, sends authenticated trace and metric data, and reads both records back +from Cloud Trace and Cloud Monitoring. It also sends a log record and confirms +that the OTLP logs RPC rejects it. The deployment restores the Cloud Run +invoker IAM check, removes public invoker bindings, and verifies the resulting +access policy before reporting success. + +The production GitHub environment supplies: + +- `OBSERVABILITY_PROJECT_ID` +- `OBSERVABILITY_COLLECTOR_REGION` +- `OBSERVABILITY_COLLECTOR_SERVICE` +- `OBSERVABILITY_COLLECTOR_ARTIFACT_REPOSITORY` +- `OBSERVABILITY_COLLECTOR_IMAGE_NAME` +- `OBSERVABILITY_COLLECTOR_SERVICE_ACCOUNT` +- `OBSERVABILITY_METRIC_LOCATION` + +The GitHub deployment service account has the custom +`collectorDeploymentIam` role on only the collector service. That role contains +`run.services.getIamPolicy` and `run.services.setIamPolicy`, which let the +workflow remove public invoker bindings. The project-level custom +`collectorDeploymentVerifier` role contains only `cloudtrace.traces.get`, +`monitoring.timeSeries.list`, and `resourcemanager.projects.get`, which let the +workflow read back its uniquely identified verification data. ## Participating workloads @@ -85,10 +113,24 @@ The infrastructure was applied and verified on 2026-09-22: - The alert policies have no notification channels, so they record incidents without sending email, Slack, or paging notifications. -The consumer services require `policyengine-observability` 3.0.1. Record the +The consumer services require `policyengine-observability` 3.0.2. Record the deployed consumer revisions and a representative cost and volume observation interval after the API v1 rollout. +## Metric delivery verification + +After deploying the collector and both consumers, query the collector's Cloud +Run logs from the deployment timestamp forward. A successful rollout has no +new metric export errors containing `InvalidArgument`, `out-of-order`, +`Duplicate TimeSeries`, `frequency`, or a missing `location`. Then run one API +v1 society report and confirm that Cloud Monitoring receives new +`policyengine.*` metric points from the API, simulation entry service, Modal +gateway, coordinator, and both Stage 12 simulation workers. + +Compare the deployed interval with the recorded pre-fix baseline: 46 collector +metric export errors comprising 863 rejected points. Trace export had no +corresponding rejected points in that interval. + ## Rollback 1. Remove the OTel endpoint from participating service configuration. diff --git a/gcp/observability/collector/config.yaml b/gcp/observability/collector/config.yaml index 2fa3ced55..72a17b37b 100644 --- a/gcp/observability/collector/config.yaml +++ b/gcp/observability/collector/config.yaml @@ -14,6 +14,11 @@ processors: - key: gcp.project_id value: ${env:OBSERVABILITY_PROJECT_ID} action: upsert + resource/metric_location: + attributes: + - key: location + value: ${env:OBSERVABILITY_METRIC_LOCATION} + action: upsert tail_sampling: decision_wait: 30s num_traces: 50000 @@ -64,7 +69,8 @@ service: exporters: [otlp_grpc] metrics: receivers: [otlp] - processors: [memory_limiter, resource/destination, batch] + processors: + [memory_limiter, resource/destination, resource/metric_location, batch] exporters: [otlp_grpc] telemetry: logs: diff --git a/policyengine_api/observability/runtime.py b/policyengine_api/observability/runtime.py index 857cada38..eefd45185 100644 --- a/policyengine_api/observability/runtime.py +++ b/policyengine_api/observability/runtime.py @@ -15,6 +15,7 @@ ServiceIdentity, StdoutLogDestination, configure, + process_instance_id, ) @@ -46,7 +47,10 @@ def _build_runtime() -> ObservabilityRuntime: environment=environment, platform="google_cloud_run", region=os.getenv("CLOUD_RUN_REGION") or "us-central1", - instance_id=os.getenv("K_REVISION"), + instance_id=process_instance_id( + "policyengine-api", + os.getenv("K_REVISION"), + ), ), logging=LoggingConfig( destinations=(StdoutLogDestination(formatter=formatter),), diff --git a/pyproject.toml b/pyproject.toml index a9d8a7cd3..024edd114 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -39,7 +39,7 @@ dependencies = [ "microdf_python>=1.0.0", "openai", "packaging>=24,<27", - "policyengine-observability[flask,google,httpx,otlp-grpc]>=3.0.1,<4", + "policyengine-observability[flask,google,httpx,otlp-grpc]>=3.0.2,<4", "policyengine_canada==0.96.3", "policyengine-ng==0.5.1", "policyengine-il==0.1.0", diff --git a/tests/unit/test_observability_collector_config.py b/tests/unit/test_observability_collector_config.py index e2bc0e1a7..264527cee 100644 --- a/tests/unit/test_observability_collector_config.py +++ b/tests/unit/test_observability_collector_config.py @@ -2,6 +2,8 @@ from pathlib import Path +import yaml + ROOT = Path(__file__).parents[2] DEPLOY = ROOT / "gcp" / "observability" @@ -14,3 +16,78 @@ def test_collector_accepts_only_traces_and_metrics() -> None: assert " traces:" in config assert " metrics:" in config assert " logs:\n receivers:" not in config + + +def test_collector_assigns_metric_location_without_overwriting_trace_region() -> None: + config = yaml.safe_load( + (DEPLOY / "collector" / "config.yaml").read_text(encoding="utf-8") + ) + + location = config["processors"]["resource/metric_location"]["attributes"] + assert location == [ + { + "key": "location", + "value": "${env:OBSERVABILITY_METRIC_LOCATION}", + "action": "upsert", + } + ] + assert ( + "resource/metric_location" + in config["service"]["pipelines"]["metrics"]["processors"] + ) + assert ( + "resource/metric_location" + not in config["service"]["pipelines"]["traces"]["processors"] + ) + + +def test_collector_has_repeatable_authenticated_deployment() -> None: + workflow = ( + ROOT / ".github" / "workflows" / "deploy-observability-collector.yml" + ).read_text(encoding="utf-8") + deploy_script = ( + ROOT / ".github" / "scripts" / "deploy_observability_collector.sh" + ).read_text(encoding="utf-8") + verify_script = ( + ROOT / ".github" / "scripts" / "verify_observability_collector.sh" + ).read_text(encoding="utf-8") + verifier = ( + ROOT / ".github" / "scripts" / "verify_observability_collector.py" + ).read_text(encoding="utf-8") + + assert "workflow_dispatch:" in workflow + assert "push:" in workflow + assert "- master" in workflow + assert '"gcp/observability/collector/**"' in workflow + assert '".github/scripts/deploy_observability_collector.sh"' in workflow + assert '".github/scripts/verify_observability_collector.sh"' in workflow + assert '".github/scripts/verify_observability_collector.py"' in workflow + assert '".github/workflows/deploy-observability-collector.yml"' in workflow + assert "environment: production" in workflow + assert "GCP_WORKLOAD_IDENTITY_PROVIDER" in workflow + assert "GCP_DEPLOY_SERVICE_ACCOUNT" in workflow + assert "validate --config=" in workflow + assert "bash .github/scripts/deploy_observability_collector.sh" in workflow + assert "bash .github/scripts/verify_observability_collector.sh" in workflow + assert ":${GITHUB_SHA}" in deploy_script + assert "--invoker-iam-check" in deploy_script + assert "remove-iam-policy-binding" in deploy_script + assert "allUsers allAuthenticatedUsers" in deploy_script + assert "invoker-iam-disabled" in deploy_script + assert "--use-http2" in deploy_script + assert "OBSERVABILITY_METRIC_LOCATION" in deploy_script + assert "httpGet.port=13133" in deploy_script + assert "TraceServiceStub" in verifier + assert "MetricsServiceStub" in verifier + assert "LogsServiceStub" in verifier + assert "StatusCode.UNIMPLEMENTED" in verifier + assert "print-identity-token" in verify_script + assert "print-access-token" in verify_script + assert "uv run --no-project" in verify_script + assert "cloudtrace.googleapis.com" in verifier + assert "monitoring.googleapis.com" in verifier + assert "--format=json" in deploy_script + assert 'select(.type == "Ready")' in deploy_script + assert 'status.latestReadyRevisionName // ""' in deploy_script + assert 'status.url // ""' in deploy_script + assert "delivery_timeout_seconds: float = 600.0" in verifier diff --git a/tests/unit/test_observability_collector_verifier.py b/tests/unit/test_observability_collector_verifier.py new file mode 100644 index 000000000..1c3ea1601 --- /dev/null +++ b/tests/unit/test_observability_collector_verifier.py @@ -0,0 +1,216 @@ +from __future__ import annotations + +import importlib.util +from contextlib import nullcontext +from pathlib import Path +from types import SimpleNamespace +from urllib.parse import parse_qs, urlparse + +import grpc +import pytest + +ROOT = Path(__file__).parents[2] +VERIFIER_PATH = ROOT / ".github" / "scripts" / "verify_observability_collector.py" + + +def _load_verifier(): + spec = importlib.util.spec_from_file_location("collector_verifier", VERIFIER_PATH) + assert spec is not None + assert spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +class _Unimplemented(grpc.RpcError): + def code(self): + return grpc.StatusCode.UNIMPLEMENTED + + +def _accepted_response(): + return SimpleNamespace( + partial_success=SimpleNamespace( + rejected_spans=0, + rejected_data_points=0, + error_message="", + ) + ) + + +def _attributes(items) -> dict[str, str]: + return {item.key: item.value.string_value for item in items} + + +def test_verify_exports_real_telemetry_reads_it_back_and_rejects_logs( + monkeypatch: pytest.MonkeyPatch, +) -> None: + verifier = _load_verifier() + exports: dict[str, object] = {} + delivery: dict[str, object] = {} + + def stub(signal: str, *, reject: bool = False): + class Stub: + def __init__(self, _channel): + pass + + def Export(self, request, *, metadata, timeout): + exports[signal] = (request, metadata, timeout) + if reject: + raise _Unimplemented() + return _accepted_response() + + return Stub + + monkeypatch.setattr( + verifier.grpc, + "secure_channel", + lambda *_args, **_kwargs: nullcontext(object()), + ) + monkeypatch.setattr(verifier, "TraceServiceStub", stub("traces")) + monkeypatch.setattr(verifier, "MetricsServiceStub", stub("metrics")) + monkeypatch.setattr(verifier, "LogsServiceStub", stub("logs", reject=True)) + monkeypatch.setattr( + verifier, + "_wait_for_delivery", + lambda **kwargs: delivery.update(kwargs), + ) + + probe = verifier.verify( + "https://collector.example", + "identity-token", + project_id="observability-project", + metric_location="us-central1", + access_token="access-token", + timeout_seconds=3.0, + delivery_timeout_seconds=9.0, + ) + + trace_request, metadata, timeout = exports["traces"] + trace_resource = trace_request.resource_spans[0] + span = trace_resource.scope_spans[0].spans[0] + assert span.name == verifier.SPAN_NAME + assert len(span.trace_id) == 16 + assert ( + _attributes(trace_resource.resource.attributes)["service.instance.id"] + == probe.verification_id + ) + assert metadata == (("authorization", "Bearer identity-token"),) + assert timeout == 3.0 + + metric_request, _, _ = exports["metrics"] + metric_resource = metric_request.resource_metrics[0] + metric = metric_resource.scope_metrics[0].metrics[0] + assert metric.name == verifier.METRIC_NAME + assert metric.gauge.data_points[0].as_int == 1 + assert ( + _attributes(metric_resource.resource.attributes)["service.instance.id"] + == probe.verification_id + ) + + log_request, _, _ = exports["logs"] + assert ( + log_request.resource_logs[0].scope_logs[0].log_records[0].body.string_value + == "collector log rejection verification" + ) + assert delivery == { + "project_id": "observability-project", + "trace_id": probe.trace_id, + "verification_id": probe.verification_id, + "metric_location": "us-central1", + "access_token": "access-token", + "started_at": probe.started_at, + "timeout_seconds": 9.0, + } + + +def test_verify_fails_when_the_collector_accepts_logs( + monkeypatch: pytest.MonkeyPatch, +) -> None: + verifier = _load_verifier() + + class AcceptingStub: + def __init__(self, _channel): + pass + + def Export(self, *_args, **_kwargs): + return _accepted_response() + + monkeypatch.setattr( + verifier.grpc, + "secure_channel", + lambda *_args, **_kwargs: nullcontext(object()), + ) + monkeypatch.setattr(verifier, "TraceServiceStub", AcceptingStub) + monkeypatch.setattr(verifier, "MetricsServiceStub", AcceptingStub) + monkeypatch.setattr(verifier, "LogsServiceStub", AcceptingStub) + + with pytest.raises(RuntimeError, match="unexpectedly accepted application logs"): + verifier.verify( + "https://collector.example", + "identity-token", + project_id="observability-project", + metric_location="us-central1", + access_token="access-token", + ) + + +def test_metric_lookup_follows_all_pages( + monkeypatch: pytest.MonkeyPatch, +) -> None: + verifier = _load_verifier() + urls: list[str] = [] + + def get_json(url: str, _access_token: str): + urls.append(url) + query = parse_qs(urlparse(url).query) + if "pageToken" not in query: + return {"nextPageToken": "second-page"} + return {"timeSeries": [{"points": [{"value": {"int64Value": "1"}}]}]} + + monkeypatch.setattr(verifier, "_get_json", get_json) + + assert verifier._metric_available( + "observability-project", + "collector-verify-123", + "us-central1", + "access-token", + 1_700_000_000.0, + ) + assert len(urls) == 2 + assert parse_qs(urlparse(urls[1]).query)["pageToken"] == ["second-page"] + monitoring_filter = parse_qs(urlparse(urls[0]).query)["filter"][0] + assert ( + 'metric.type = "prometheus.googleapis.com/' + 'policyengine.collector.verification/gauge"' in monitoring_filter + ) + assert 'resource.type = "prometheus_target"' in monitoring_filter + assert 'resource.labels.location = "us-central1"' in monitoring_filter + assert 'resource.labels.instance = "collector-verify-123"' in monitoring_filter + + +def test_partial_metric_rejection_fails_verification() -> None: + verifier = _load_verifier() + response = SimpleNamespace( + partial_success=SimpleNamespace( + rejected_data_points=1, + error_message="invalid location", + ) + ) + + with pytest.raises(RuntimeError, match="invalid location"): + verifier._raise_for_partial_rejection( + response, + signal="metric points", + rejected_field="rejected_data_points", + ) + + +@pytest.mark.parametrize( + "endpoint", + ["http://collector.example", "collector.example", "https:///missing-host"], +) +def test_authority_rejects_non_https_or_hostless_endpoints(endpoint: str) -> None: + verifier = _load_verifier() + + with pytest.raises(ValueError, match="HTTPS URL"): + verifier._authority(endpoint) diff --git a/tests/unit/test_observability_runtime.py b/tests/unit/test_observability_runtime.py index 8eba73d22..70740dfc9 100644 --- a/tests/unit/test_observability_runtime.py +++ b/tests/unit/test_observability_runtime.py @@ -1,3 +1,5 @@ +import re + from policyengine_observability import ( GoogleCloudLogFormatter, StdoutLogDestination, @@ -15,6 +17,7 @@ def test_runtime_uses_consumer_owned_identity_and_stdout(monkeypatch): monkeypatch.setenv("OTEL_TRACES_SAMPLER_ARG", "0.01") monkeypatch.setenv("OBSERVABILITY_SERVICE_NAMESPACE", "example.stack") monkeypatch.setenv("OBSERVABILITY_TRACE_PROJECT_ID", "trace-project") + monkeypatch.setenv("K_REVISION", "policyengine-api-00123-test") runtime = _build_runtime() try: @@ -22,6 +25,10 @@ def test_runtime_uses_consumer_owned_identity_and_stdout(monkeypatch): assert runtime.config.otel.sampling_ratio == 1.0 assert runtime.config.application_attribute_keys is None assert runtime.config.dispatch_attribute_keys == frozenset({"observability_id"}) + assert re.fullmatch( + r"policyengine-api-00123-test:\d+:[0-9a-f]{32}", + runtime.config.deployment.instance_id or "", + ) assert len(runtime.config.logging.destinations) == 1 destination = runtime.config.logging.destinations[0] assert isinstance(destination, StdoutLogDestination) diff --git a/uv.lock b/uv.lock index 0f7818484..7ea6142d3 100644 --- a/uv.lock +++ b/uv.lock @@ -2848,7 +2848,7 @@ requires-dist = [ { name = "policyengine-canada", specifier = "==0.96.3" }, { name = "policyengine-il", specifier = "==0.1.0" }, { name = "policyengine-ng", specifier = "==0.5.1" }, - { name = "policyengine-observability", extras = ["flask", "google", "httpx", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["flask", "google", "httpx", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "psycopg", extras = ["binary"], specifier = ">=3.3,<4" }, { name = "pydantic" }, { name = "pymysql" }, @@ -2951,11 +2951,11 @@ wheels = [ [[package]] name = "policyengine-observability" -version = "3.0.1" +version = "3.0.2" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/00/a6/f5d3e49523bf3d3231e9d2297fb03a4ae0704ae246d86153356dcb7eee7c/policyengine_observability-3.0.1.tar.gz", hash = "sha256:35e934b31843545a13570d6ba2e6a4b0cb002488f2c7bb1840351f6f90f30104", size = 124370, upload-time = "2026-09-28T16:48:19.916Z" } +sdist = { url = "https://files.pythonhosted.org/packages/74/2a/0c261c8ca693bb3a13d7fe837d3f89f7fdda245496d4b3e92eefed3c972e/policyengine_observability-3.0.2.tar.gz", hash = "sha256:554f907e43d8eeb274c1298fbe990b1626ac45f194c49995ddde091ddcdc2a64", size = 125718, upload-time = "2026-09-30T21:36:02.028Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/27/6c/879ecc01a618885e4bc741d8abb35e601d638f878e3145cda677b75c46b7/policyengine_observability-3.0.1-py3-none-any.whl", hash = "sha256:b7630d4d2d91c183b733d2e3470d3c0cbad67fd59a8c5e2f2bcc75872eea3c02", size = 40859, upload-time = "2026-09-28T16:48:18.661Z" }, + { url = "https://files.pythonhosted.org/packages/a5/9c/b48fb64c1bebf83bd5fdeebe914188806c3513cf2de0fc95ee1d95082fc4/policyengine_observability-3.0.2-py3-none-any.whl", hash = "sha256:d84c7564922f6a3294268e89549065924edde932c31adf8b7c714558e107cb31", size = 42029, upload-time = "2026-09-30T21:36:00.841Z" }, ] [package.optional-dependencies]