From 891d47e1c0a826b2fe47e166aacee77daa20a7f3 Mon Sep 17 00:00:00 2001 From: MarcoDN Date: Wed, 7 Oct 2026 15:43:48 +0100 Subject: [PATCH] Add GCE and GKE integration tests for the default OTel config The terraform/gcp/gce and terraform/gcp/gke modules provision a GCE instance and a zonal GKE cluster, federate a per-run IAM role to AWS via web identity, install the agent with the default OTel config, generate OTLP load, and run the Go suites validating metrics, logs, and traces in CloudWatch. Adds the GCE and GKE compute types and a README covering the one-time GCP project and AWS account setup. # Conflicts: # environment/metadata.go # test/azure/aks/aks_test.go # test/azure/vm/azurevm_test.go --- environment/computetype/compute_type.go | 7 + environment/metadata.go | 12 +- terraform/azure/aks/main.tf | 3 +- terraform/gcp/README.md | 119 +++++ terraform/gcp/gce/iam.tf | 93 ++++ terraform/gcp/gce/main.tf | 126 +++++ terraform/gcp/gce/outputs.tf | 22 + terraform/gcp/gce/providers.tf | 26 ++ terraform/gcp/gce/variables.tf | 97 ++++ terraform/gcp/gke/.gitignore | 1 + terraform/gcp/gke/main.tf | 429 ++++++++++++++++++ terraform/gcp/gke/providers.tf | 53 +++ terraform/gcp/gke/variables.tf | 81 ++++ .../{azure/aks => }/otlp_load_generator.sh | 8 +- test/gcp/gce/gce_test.go | 246 ++++++++++ test/gcp/gke/gke_test.go | 165 +++++++ test/otel_collect/otlpvalidation/payloads.go | 210 +++++++++ 17 files changed, 1691 insertions(+), 7 deletions(-) create mode 100644 terraform/gcp/README.md create mode 100644 terraform/gcp/gce/iam.tf create mode 100644 terraform/gcp/gce/main.tf create mode 100644 terraform/gcp/gce/outputs.tf create mode 100644 terraform/gcp/gce/providers.tf create mode 100644 terraform/gcp/gce/variables.tf create mode 100644 terraform/gcp/gke/.gitignore create mode 100644 terraform/gcp/gke/main.tf create mode 100644 terraform/gcp/gke/providers.tf create mode 100644 terraform/gcp/gke/variables.tf rename terraform/{azure/aks => }/otlp_load_generator.sh (72%) create mode 100644 test/gcp/gce/gce_test.go create mode 100644 test/gcp/gke/gke_test.go create mode 100644 test/otel_collect/otlpvalidation/payloads.go diff --git a/environment/computetype/compute_type.go b/environment/computetype/compute_type.go index 4ea79253e..a2c96e2cf 100644 --- a/environment/computetype/compute_type.go +++ b/environment/computetype/compute_type.go @@ -16,6 +16,11 @@ const ( // AKS is an Azure Kubernetes Service cluster authenticating to AWS via the projected // service-account web-identity credential chain. AKS ComputeType = "AKS" + // GCE is a non-AWS host authenticating to AWS via the GCP web-identity credential chain. + GCE ComputeType = "GCE" + // GKE is a Google Kubernetes Engine cluster authenticating to AWS via the projected + // service-account web-identity credential chain. + GKE ComputeType = "GKE" ) var ( @@ -25,6 +30,8 @@ var ( "EKS": EKS, "AZUREVM": AzureVM, "AKS": AKS, + "GCE": GCE, + "GKE": GKE, } ) diff --git a/environment/metadata.go b/environment/metadata.go index dbb79ff42..00edccda1 100644 --- a/environment/metadata.go +++ b/environment/metadata.go @@ -44,6 +44,7 @@ type MetaData struct { AzureVMSize string AzureResourceGroup string AzureLocation string + GKEClusterName string ProxyUrl string AssumeRoleArn string InstanceArn string @@ -96,6 +97,7 @@ type MetaDataStrings struct { AzureVMSize string AzureResourceGroup string AzureLocation string + GKEClusterName string ProxyUrl string AssumeRoleArn string InstanceArn string @@ -129,7 +131,7 @@ type MetaDataStrings struct { } func registerComputeType(dataString *MetaDataStrings) { - flag.StringVar(&(dataString.ComputeType), "computeType", "", "EC2/ECS/EKS/AZUREVM/AKS") + flag.StringVar(&(dataString.ComputeType), "computeType", "", "EC2/ECS/EKS/AZUREVM/AKS/GCE/GKE") } func registerBucket(dataString *MetaDataStrings) { flag.StringVar(&(dataString.Bucket), "bucket", "", "s3 bucket ex cloudwatch-agent-integration-bucket") @@ -163,6 +165,10 @@ func registerAzureVMData(d *MetaDataStrings) { flag.StringVar(&(d.AzureLocation), "azureLocation", "", "expected cloud.region resource attribute (Azure location, e.g. eastus)") } +func registerGKEData(d *MetaDataStrings) { + flag.StringVar(&(d.GKEClusterName), "gkeClusterName", "", "GKE cluster name") +} + func registerEKSData(d *MetaDataStrings) { flag.StringVar(&(d.EKSClusterName), "eksClusterName", "", "EKS cluster name") flag.StringVar(&(d.EksDeploymentStrategy), "eksDeploymentStrategy", "", "Daemon/Replica/Sidecar") @@ -208,7 +214,7 @@ func registerProxyUrl(dataString *MetaDataStrings) { func fillComputeType(e *MetaData, data *MetaDataStrings) { computeType, ok := computetype.FromString(data.ComputeType) if !ok { - log.Panic("Invalid compute type. Needs to be EC2/ECS/EKS/AZUREVM/AKS. Compute Type is a required flag. :" + data.ComputeType) + log.Panic("Invalid compute type. Needs to be EC2/ECS/EKS/AZUREVM/AKS/GCE/GKE. Compute Type is a required flag. :" + data.ComputeType) } e.ComputeType = computeType } @@ -344,6 +350,7 @@ func RegisterEnvironmentMetaDataFlags() *MetaDataStrings { registerEKSData(registeredMetaDataStrings) registerAKSData(registeredMetaDataStrings) registerAzureVMData(registeredMetaDataStrings) + registerGKEData(registeredMetaDataStrings) registerEKSE2ETestData(registeredMetaDataStrings) registerBucket(registeredMetaDataStrings) registerS3Key(registeredMetaDataStrings) @@ -387,6 +394,7 @@ func GetEnvironmentMetaData() *MetaData { metaDataStorage.AzureVMSize = registeredMetaDataStrings.AzureVMSize metaDataStorage.AzureResourceGroup = registeredMetaDataStrings.AzureResourceGroup metaDataStorage.AzureLocation = registeredMetaDataStrings.AzureLocation + metaDataStorage.GKEClusterName = registeredMetaDataStrings.GKEClusterName metaDataStorage.InstancePlatform = registeredMetaDataStrings.InstancePlatform metaDataStorage.AgentStartCommand = registeredMetaDataStrings.AgentStartCommand metaDataStorage.EksGpuType = registeredMetaDataStrings.EksGpuType diff --git a/terraform/azure/aks/main.tf b/terraform/azure/aks/main.tf index 4a7e18208..a093255d7 100644 --- a/terraform/azure/aks/main.tf +++ b/terraform/azure/aks/main.tf @@ -435,7 +435,8 @@ resource "kubernetes_job_v1" "otlp_load" { name = "load-gen" image = "curlimages/curl:8.8.0" command = ["/bin/sh", "-c"] - args = [templatefile("${path.module}/otlp_load_generator.sh", { + args = [templatefile("${path.module}/../../otlp_load_generator.sh", { + prefix = "aks" service_name = local.load_gen_service_name instance_id = azurerm_kubernetes_cluster.cwagent.name endpoint = "http://127.0.0.1:4318" diff --git a/terraform/gcp/README.md b/terraform/gcp/README.md new file mode 100644 index 000000000..52162d457 --- /dev/null +++ b/terraform/gcp/README.md @@ -0,0 +1,119 @@ +# GCP integration tests — environment setup + +One-time setup the GCP project and AWS account need before the `terraform/gcp` +modules can run. Everything else is created and destroyed per run by the +modules themselves. + +The commands below use these placeholders: + +```sh +PROJECT_ID= +CI_SA=otel-collector-integ-tests@$PROJECT_ID.iam.gserviceaccount.com +``` + +## GCP project + +Enable the APIs: + +```sh +gcloud services enable compute.googleapis.com iam.googleapis.com \ + container.googleapis.com --project "$PROJECT_ID" +``` + +An existing VPC network and subnetwork must be passed as `gcp_network_name` / +`gcp_subnetwork_name`. Prefer a dedicated network over the auto-created +`default` one: the default network comes with pre-populated firewall rules +(`default-allow-ssh` among them) that admit any source IP, which defeats the +modules' runner-IP-scoped SSH rule. A dedicated network with no extra rules +leaves the modules' own firewall rules as the only ingress: + +```sh +gcloud compute networks create cwagent-integ --subnet-mode auto --project "$PROJECT_ID" +``` + +Create the CI service account and grant it — and any human running the suites +locally — these project roles: + +| Role | Needed for | +|---|---| +| `roles/compute.admin` | VM and firewall lifecycle (`gce`) | +| `roles/iam.serviceAccountAdmin` | creating the per-run service account (`gce`) | +| `roles/iam.serviceAccountUser` | attaching service accounts to VMs and GKE nodes | +| `roles/container.admin` | cluster lifecycle and in-cluster RBAC objects (`gke`; project editor alone is not enough) | + +```sh +gcloud iam service-accounts create otel-collector-integ-tests --project "$PROJECT_ID" + +for role in roles/compute.admin roles/iam.serviceAccountAdmin \ + roles/iam.serviceAccountUser roles/container.admin; do + gcloud projects add-iam-policy-binding "$PROJECT_ID" \ + --member "serviceAccount:$CI_SA" --role "$role" +done +``` + +For a human, the same bindings with `--member "user:"`. + +## CI authentication + +**Keyless (Workload Identity Federation)** — required when org policy disables +service-account key creation. The provider resource path and the +service-account email go in the workflow's auth step; the path embeds the +project number, so a project move means updating the workflow file as well. + +```sh +gcloud iam workload-identity-pools create github-actions \ + --project "$PROJECT_ID" --location global --display-name "GitHub Actions" + +gcloud iam workload-identity-pools providers create-oidc github \ + --project "$PROJECT_ID" --location global \ + --workload-identity-pool github-actions --display-name "GitHub" \ + --issuer-uri "https://token.actions.githubusercontent.com" \ + --attribute-mapping "google.subject=assertion.sub,attribute.repository=assertion.repository" \ + --attribute-condition "assertion.repository == 'aws/amazon-cloudwatch-agent'" + +PROJECT_NUMBER=$(gcloud projects describe "$PROJECT_ID" --format "value(projectNumber)") + +gcloud iam service-accounts add-iam-policy-binding "$CI_SA" \ + --project "$PROJECT_ID" --role roles/iam.workloadIdentityUser \ + --member "principalSet://iam.googleapis.com/projects/$PROJECT_NUMBER/locations/global/workloadIdentityPools/github-actions/attribute.repository/aws/amazon-cloudwatch-agent" +``` + +**Static key** — where key creation is allowed, the simpler alternative: + +```sh +gcloud iam service-accounts keys create key.json --iam-account "$CI_SA" +gh secret set GCP_CREDENTIALS < key.json +``` + +## GitHub repository configuration + +| Setting | Value | +|---|---| +| `vars.GCP_PROJECT` | project ID (not the display name or number) | +| `vars.GCP_NETWORK_NAME` | VPC network name | +| `secrets.GCP_CREDENTIALS` | service-account key JSON (static-key model only) | + +```sh +gh variable set GCP_PROJECT --body "$PROJECT_ID" +gh variable set GCP_NETWORK_NAME --body default +``` + +## AWS account + +Enable Transaction Search in the suite region (us-east-2): trace validation +reads the `aws/spans` log group, which only exists where the account's X-Ray +trace destination is CloudWatch Logs. See the `region` variable comment in +`gce/variables.tf`. + +## Local runs + +- Google credentials come from application-default credentials + (`gcloud auth application-default login`; on Workspace-managed accounts that + reject it with `admin_policy_enforced`, use `gcloud auth login --update-adc`). +- `gce`: a locally built agent `.deb` (`agent_deb_path`) and a test-repo clone + URL/branch reachable from the VM (`github_test_repo`, + `github_test_repo_branch`). +- `gke`: `kubectl` on PATH, and an agent container image in an ECR repository + your AWS credentials can access (`cwagent_image_repo`, `cwagent_image_tag`, + `ecr_region`). The agent repo's `make docker-build-amd64 IMAGE=...` target + builds a suitable image from local binaries. diff --git a/terraform/gcp/gce/iam.tf b/terraform/gcp/gce/iam.tf new file mode 100644 index 000000000..d2d74e394 --- /dev/null +++ b/terraform/gcp/gce/iam.tf @@ -0,0 +1,93 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +# CWAGENT_ROLE: assumed via web identity (GCP service-account identity token). Carries the agent's +# default:otel CloudWatch writes plus the reads the on-VM test needs to assert delivery. + +module "common" { + source = "../../common" +} + +# Per-run service account: the identity the VM mints tokens as. It holds no GCP permissions: it +# exists only so the metadata server mints identity tokens for it, and its unique ID is what the +# role trust pins. account_id caps at 30 chars, hence the short prefix. +resource "google_service_account" "cwagent" { + account_id = "cwa-gce-${module.common.testing_id}" + display_name = "cwa-gce-integ-${module.common.testing_id}" +} + +data "aws_iam_policy_document" "cwagent_assume_role" { + statement { + effect = "Allow" + actions = ["sts:AssumeRoleWithWebIdentity"] + + # Google is a built-in web-identity provider, so no IAM OIDC provider resource is involved: + # the principal is accounts.google.com itself. + principals { + type = "Federated" + identifiers = ["accounts.google.com"] + } + + # Pin all three Google condition keys -- the recommended trust policy for Google-issued tokens + # and the same shape the onboarding scripts create. :sub is the service account's unique ID, + # :aud reads the token's azp claim (the unique ID again on service-account tokens), and :oaud + # is the audience the token was requested with, rejecting tokens minted for other services. + # https://aws.amazon.com/blogs/security/access-aws-using-a-google-cloud-platform-native-workload-identity/ + condition { + test = "StringEquals" + variable = "accounts.google.com:aud" + values = [google_service_account.cwagent.unique_id] + } + + condition { + test = "StringEquals" + variable = "accounts.google.com:sub" + values = [google_service_account.cwagent.unique_id] + } + + condition { + test = "StringEquals" + variable = "accounts.google.com:oaud" + values = [var.gcp_token_audience] + } + } +} + +resource "aws_iam_role" "cwagent" { + name = "cwa-gce-integ-role-${module.common.testing_id}" + assume_role_policy = data.aws_iam_policy_document.cwagent_assume_role.json +} + +# The agent's own writes come from the same AWS-managed policy customers are told to use, so a green run +# also proves that documented policy is sufficient over the GCP web-identity path. +resource "aws_iam_role_policy_attachment" "cwagent_server_policy" { + role = aws_iam_role.cwagent.name + policy_arn = "arn:aws:iam::aws:policy/CloudWatchAgentServerPolicy" +} + +# Validation reads only -- the agent's own writes are fully covered by CloudWatchAgentServerPolicy. The +# test binary runs on the VM under this same role, so these have to live here. +data "aws_iam_policy_document" "cwagent_permissions" { + statement { + effect = "Allow" + actions = [ + "cloudwatch:ListMetrics", + "cloudwatch:GetMetricData", + "logs:GetLogEvents", + # Cleanup: the test deletes its own stream from the shared /aws/cwagent/otlp group when it finishes. + "logs:DeleteLogStream", + # StartQuery/GetQueryResults validate OTLP trace delivery via the aws/spans log group. That group + # is only populated where the X-Ray trace segment destination is set to CloudWatchLogs, which is a + # per-region setting -- hence the region default in variables.tf. + "logs:StartQuery", + "logs:GetQueryResults", + ] + resources = ["*"] + } +} + +resource "aws_iam_role_policy" "cwagent" { + name = "cwa-gce-integ-policy-${module.common.testing_id}" + role = aws_iam_role.cwagent.id + policy = data.aws_iam_policy_document.cwagent_permissions.json +} diff --git a/terraform/gcp/gce/main.tf b/terraform/gcp/gce/main.tf new file mode 100644 index 000000000..76df52878 --- /dev/null +++ b/terraform/gcp/gce/main.tf @@ -0,0 +1,126 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +##################################################################### +# SSH key for connecting to the GCE VM +##################################################################### +resource "tls_private_key" "ssh_key" { + algorithm = "RSA" + rsa_bits = 4096 +} + +locals { + # Subnetworks are regional; derive the region from the zone input (us-east1-b -> us-east1). + gcp_region = join("-", slice(split("-", var.gcp_zone), 0, 2)) +} + +# Attach to an existing network/subnetwork in the project so CI needs no networking-create perms. +data "google_compute_network" "selected" { + name = var.gcp_network_name +} + +data "google_compute_subnetwork" "selected" { + name = var.gcp_subnetwork_name + region = local.gcp_region +} + +# Allow inbound SSH from the runner only, scoped to this VM via its target tag. Explicit so the +# module does not depend on the network's pre-existing firewall rules. +resource "google_compute_firewall" "cwagent" { + name = "cwa-gce-integ-fw-${module.common.testing_id}" + network = data.google_compute_network.selected.self_link + + allow { + protocol = "tcp" + ports = ["22"] + } + + source_ranges = [var.runner_ip] + target_tags = ["cwa-gce-integ-${module.common.testing_id}"] +} + +# GCE Linux VM; its attached service account is what oidctoken exchanges for an AWS session. +resource "google_compute_instance" "cwagent" { + name = "cwa-gce-integ-${module.common.testing_id}" + machine_type = var.gcp_machine_type + zone = var.gcp_zone + tags = ["cwa-gce-integ-${module.common.testing_id}"] + + boot_disk { + initialize_params { + image = var.gcp_image + } + } + + network_interface { + subnetwork = data.google_compute_subnetwork.selected.self_link + + # Ephemeral public IP for the runner's SSH connection. + access_config {} + } + + # Attached per-run service account: the token source for cross-cloud AssumeRoleWithWebIdentity. + # cloud-platform is the standard access scope; the account holds no IAM roles, so it grants nothing. + service_account { + email = google_service_account.cwagent.email + scopes = ["cloud-platform"] + } + + metadata = { + ssh-keys = "${var.admin_username}:${tls_private_key.ssh_key.public_key_openssh}" + } +} + +##################################################################### +# Install the agent, start it with default:otel, and run the test. +##################################################################### +resource "null_resource" "integration_test" { + connection { + type = "ssh" + user = var.admin_username + private_key = tls_private_key.ssh_key.private_key_pem + host = google_compute_instance.cwagent.network_interface[0].access_config[0].nat_ip + } + + # Upload the runner-built .deb straight over the SSH connection (no S3 or public URL). + provisioner "file" { + source = var.agent_deb_path + destination = "/home/${var.admin_username}/amazon-cloudwatch-agent.deb" + } + + # Install Go, clone the test repo, and install the uploaded agent .deb. + provisioner "remote-exec" { + inline = [ + "cloud-init status --wait", + "echo sha ${var.cwa_github_sha}", + "sudo apt-get update -y && sudo apt-get install -y golang-go git", + "git clone --branch ${var.github_test_repo_branch} ${var.github_test_repo} -q", + "sudo dpkg -i -E amazon-cloudwatch-agent.deb", + ] + } + + # Persist env vars with the ctl set-env action: the agent loads env-config.json at startup, making them + # available to OTel expandconverter which resolves ${AWS_REGION} and ${CWAGENT_ROLE_ARN} in the + # translated YAML. set-env runs before fetch-config so both are set on the first agent start. + provisioner "remote-exec" { + inline = [ + "sudo /opt/aws/amazon-cloudwatch-agent/bin/amazon-cloudwatch-agent-ctl -a set-env -e 'AWS_REGION=${var.region}'", + "sudo /opt/aws/amazon-cloudwatch-agent/bin/amazon-cloudwatch-agent-ctl -a set-env -e 'CWAGENT_ROLE_ARN=${aws_iam_role.cwagent.arn}'", + "sudo /opt/aws/amazon-cloudwatch-agent/bin/amazon-cloudwatch-agent-ctl -a fetch-config -m auto -s -c default:otel", + # The test binary validates delivery via AWS reads; give it the same web-identity chain the agent + # uses. The token is a bearer credential, so create it 0600 (umask before redirect, since the + # shell creates the file) and remove it once the test finishes. + "umask 077 && curl -s -H 'Metadata-Flavor: Google' 'http://169.254.169.254/computeMetadata/v1/instance/service-accounts/default/identity?audience=${var.gcp_token_audience}' > /tmp/gcp-identity-token", + "export AWS_WEB_IDENTITY_TOKEN_FILE=/tmp/gcp-identity-token AWS_ROLE_ARN=${aws_iam_role.cwagent.arn} AWS_REGION=${var.region}", + "cd amazon-cloudwatch-agent-test", + "go test -tags integration ${var.test_dir} -p 1 -timeout 30m -computeType=GCE -region=${var.region} -cwaCommitSha=${var.cwa_github_sha} -instanceId=${google_compute_instance.cwagent.instance_id} -assumeRoleArn=${aws_iam_role.cwagent.arn} -v; test_rc=$?; rm -f /tmp/gcp-identity-token; exit $test_rc", + ] + } + + depends_on = [ + google_compute_instance.cwagent, + google_compute_firewall.cwagent, + aws_iam_role_policy.cwagent, + aws_iam_role_policy_attachment.cwagent_server_policy, + ] +} diff --git a/terraform/gcp/gce/outputs.tf b/terraform/gcp/gce/outputs.tf new file mode 100644 index 000000000..49e339ec6 --- /dev/null +++ b/terraform/gcp/gce/outputs.tf @@ -0,0 +1,22 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +output "cwagent_public_ip" { + value = google_compute_instance.cwagent.network_interface[0].access_config[0].nat_ip +} + +output "cwagent_instance_id" { + value = google_compute_instance.cwagent.instance_id +} + +output "cwagent_role_arn" { + value = aws_iam_role.cwagent.arn +} + +output "cwagent_sa_unique_id" { + value = google_service_account.cwagent.unique_id +} + +output "testing_id" { + value = module.common.testing_id +} diff --git a/terraform/gcp/gce/providers.tf b/terraform/gcp/gce/providers.tf new file mode 100644 index 000000000..28d755a5c --- /dev/null +++ b/terraform/gcp/gce/providers.tf @@ -0,0 +1,26 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +terraform { + required_providers { + aws = { + source = "hashicorp/aws" + version = "!= 6.22.0" + } + google = { + source = "hashicorp/google" + version = "~> 6.0" + } + } +} + +provider "aws" { + region = var.region +} + +# Credentials come from GOOGLE_APPLICATION_CREDENTIALS or gcloud application-default +# credentials; project and zone are explicit module inputs. +provider "google" { + project = var.gcp_project + zone = var.gcp_zone +} diff --git a/terraform/gcp/gce/variables.tf b/terraform/gcp/gce/variables.tf new file mode 100644 index 000000000..5a0444581 --- /dev/null +++ b/terraform/gcp/gce/variables.tf @@ -0,0 +1,97 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +##################################################################### +# AWS side +##################################################################### + +# Trace validation reads the aws/spans log group, which only exists where the X-Ray trace segment +# destination is CloudWatchLogs. That destination is a per-region setting, so this suite cannot use the +# repo-wide us-west-2 default: us-west-2 is deliberately left on the legacy XRay destination because the +# App Signals trace suite there validates through the X-Ray query APIs, which Transaction Search would +# break. +variable "region" { + type = string + default = "us-east-2" +} + +variable "test_dir" { + type = string + default = "./test/gcp/gce" +} + +variable "cwa_github_sha" { + type = string + default = "" +} + +variable "github_test_repo" { + type = string + default = "https://github.com/aws/amazon-cloudwatch-agent-test.git" +} + +variable "github_test_repo_branch" { + type = string + default = "main" +} + +# Local path to the agent .deb (built on the runner); uploaded to the VM over SSH, so no S3/public URL needed. +variable "agent_deb_path" { + type = string + default = "" +} + +##################################################################### +# GCP side +##################################################################### + +# Existing project the VM is created in (input so CI needs no project create/delete perms). +variable "gcp_project" { + type = string + default = "" +} + +variable "gcp_zone" { + type = string + default = "us-east1-b" +} + +variable "gcp_machine_type" { + type = string + default = "e2-standard-2" +} + +# Existing VPC network + subnetwork the VM's NIC attaches to (must exist; no safe default). +variable "gcp_network_name" { + type = string + default = "" +} + +variable "gcp_subnetwork_name" { + type = string + default = "default" +} + +variable "admin_username" { + type = string + default = "cwagent" +} + +# CIDR allowed inbound SSH to the VM (the CI runner's public IP, e.g. "1.2.3.4/32"). +# Required, not defaulted: this is the only thing scoping the firewall rule's Allow-22 reach. +variable "runner_ip" { + type = string +} + +# Ubuntu image the VM boots; matches the Debian-package install path main.tf uses. +variable "gcp_image" { + type = string + default = "ubuntu-os-cloud/ubuntu-2404-lts-amd64" +} + +# Token audience (aud) requested from the metadata identity endpoint; must match the role trust +# condition. sts.amazonaws.com is the audience the agent's own token provider requests. +variable "gcp_token_audience" { + type = string + default = "sts.amazonaws.com" +} diff --git a/terraform/gcp/gke/.gitignore b/terraform/gcp/gke/.gitignore new file mode 100644 index 000000000..f940e0b55 --- /dev/null +++ b/terraform/gcp/gke/.gitignore @@ -0,0 +1 @@ +kubeconfig diff --git a/terraform/gcp/gke/main.tf b/terraform/gcp/gke/main.tf new file mode 100644 index 000000000..26756754b --- /dev/null +++ b/terraform/gcp/gke/main.tf @@ -0,0 +1,429 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +module "common" { + source = "../../common" +} + +##################################################################### +# GKE cluster (its OIDC issuer is what AWS trusts for cross-cloud web-identity) +##################################################################### +# GKE always serves the cluster's OIDC discovery document and projects service-account +# tokens; there is no issuer opt-in flag to set. +resource "google_container_cluster" "cwagent" { + name = "cwa-gke-integ-${module.common.testing_id}" + location = var.gcp_zone + min_master_version = var.kubernetes_version + network = var.gcp_network_name + subnetwork = var.gcp_subnetwork_name + + initial_node_count = var.gke_node_count + + # The google provider defaults this to true, which would make terraform destroy fail. + deletion_protection = false + + # Terraform drives the cluster over the public API server, so restrict it to the runner that created it. + # runner_ip is required, so there is no path where this silently ends up open to all. + master_authorized_networks_config { + cidr_blocks { + cidr_block = var.runner_ip + } + } + + node_config { + machine_type = var.gke_node_machine_type + disk_size_gb = 50 + oauth_scopes = ["https://www.googleapis.com/auth/cloud-platform"] + } +} + +##################################################################### +# AWS IAM: trust GKE OIDC issuer for cross-cloud federation +##################################################################### +locals { + # GKE serves the issuer at this deterministic URL; the cluster resource does not export it + # as an attribute. Referencing the resource's name/location makes everything derived from + # this URL wait for the cluster, so the discovery endpoint is live before it is read. + gke_oidc_issuer_url = "https://container.googleapis.com/v1/projects/${var.gcp_project}/locations/${google_container_cluster.cwagent.location}/clusters/${google_container_cluster.cwagent.name}" +} + +data "tls_certificate" "gke_oidc" { + url = local.gke_oidc_issuer_url +} + +resource "aws_iam_openid_connect_provider" "gke" { + url = local.gke_oidc_issuer_url + client_id_list = ["sts.amazonaws.com"] + thumbprint_list = [data.tls_certificate.gke_oidc.certificates[0].sha1_fingerprint] +} + +locals { + gke_oidc_issuer_host = replace(local.gke_oidc_issuer_url, "https://", "") + namespace = "amazon-cloudwatch" + service_account_name = "cloudwatch-agent" + cwagent_role_name = "cwa-gke-integ-role-${module.common.testing_id}" + + # Must match serviceName in test/gcp/gke/gke_test.go -- the test derives the expected log stream + # and the trace query filter from it. + load_gen_service_name = "gke-otlp-test-service" + load_gen_duration_seconds = 180 +} + +data "aws_iam_policy_document" "cwagent_assume_role" { + statement { + effect = "Allow" + actions = ["sts:AssumeRoleWithWebIdentity"] + + principals { + type = "Federated" + identifiers = [aws_iam_openid_connect_provider.gke.arn] + } + + condition { + test = "StringEquals" + variable = "${local.gke_oidc_issuer_host}:sub" + values = ["system:serviceaccount:${local.namespace}:${local.service_account_name}"] + } + + condition { + test = "StringEquals" + variable = "${local.gke_oidc_issuer_host}:aud" + values = ["sts.amazonaws.com"] + } + } + +} + +resource "aws_iam_role" "cwagent" { + name = local.cwagent_role_name + assume_role_policy = data.aws_iam_policy_document.cwagent_assume_role.json +} + +# The agent's own writes come from the same AWS-managed policy customers are told to use, so a green run +# also proves that documented policy is sufficient over the GKE projected-token path. +resource "aws_iam_role_policy_attachment" "cwagent_server_policy" { + role = aws_iam_role.cwagent.name + policy_arn = "arn:aws:iam::aws:policy/CloudWatchAgentServerPolicy" +} + +# No inline policy: the GKE test binary runs on the runner under its own credentials, so this role needs +# agent writes only -- and CloudWatchAgentServerPolicy alone covers them, OTLP traces included. + +##################################################################### +# Kubernetes resources: deploy CWA DaemonSet from ECR image +##################################################################### +resource "kubernetes_namespace" "cwagent" { + metadata { + name = local.namespace + } +} + +resource "kubernetes_service_account" "cwagent" { + metadata { + name = local.service_account_name + namespace = kubernetes_namespace.cwagent.metadata[0].name + } +} + +resource "kubernetes_cluster_role" "cwagent" { + metadata { + name = "cwa-gke-integ-${module.common.testing_id}" + } + + rule { + api_groups = [""] + resources = ["pods", "nodes", "endpoints", "services", "namespaces"] + verbs = ["list", "watch", "get"] + } + rule { + api_groups = ["apps"] + resources = ["replicasets", "daemonsets", "deployments"] + verbs = ["list", "watch", "get"] + } + rule { + api_groups = ["batch"] + resources = ["jobs"] + verbs = ["list", "watch", "get"] + } +} + +resource "kubernetes_cluster_role_binding" "cwagent" { + metadata { + name = "cwa-gke-integ-${module.common.testing_id}" + } + role_ref { + api_group = "rbac.authorization.k8s.io" + kind = "ClusterRole" + name = kubernetes_cluster_role.cwagent.metadata[0].name + } + subject { + kind = "ServiceAccount" + name = kubernetes_service_account.cwagent.metadata[0].name + namespace = kubernetes_namespace.cwagent.metadata[0].name + } +} + +# ECR pull secret so GKE nodes can pull the CWA image from AWS ECR. +# The 12h auth token is fetched here with the runner's AWS credentials rather +# than passed in as a variable, which cannot survive the workflow's shell quoting. +# The integration-test image is published to us-west-2 only, while the job's +# CloudWatch region may differ -- pin the registry host to the ECR region. +locals { + cwagent_image_repo = replace(var.cwagent_image_repo, "/\\.ecr\\.[a-z0-9-]+\\./", ".ecr.${var.ecr_region}.") +} + +data "aws_ecr_authorization_token" "ecr" { + provider = aws.ecr +} + +resource "kubernetes_secret" "ecr_pull" { + metadata { + name = "ecr-pull-secret" + namespace = kubernetes_namespace.cwagent.metadata[0].name + } + type = "kubernetes.io/dockerconfigjson" + data = { + ".dockerconfigjson" = jsonencode({ + auths = { + (split("/", local.cwagent_image_repo)[0]) = { + auth = data.aws_ecr_authorization_token.ecr.authorization_token + } + } + }) + } +} + +resource "kubernetes_daemon_set_v1" "cwagent" { + metadata { + name = "cloudwatch-agent" + namespace = kubernetes_namespace.cwagent.metadata[0].name + } + + spec { + selector { + match_labels = { app = "cloudwatch-agent" } + } + + template { + metadata { + labels = { app = "cloudwatch-agent" } + } + + spec { + service_account_name = kubernetes_service_account.cwagent.metadata[0].name + host_network = true + dns_policy = "ClusterFirstWithHostNet" + + image_pull_secrets { + name = kubernetes_secret.ecr_pull.metadata[0].name + } + + container { + name = "cloudwatch-agent" + image = "${local.cwagent_image_repo}:${var.cwagent_image_tag}" + image_pull_policy = "Always" + + env { + name = "AWS_REGION" + value = var.region + } + env { + name = "AWS_WEB_IDENTITY_TOKEN_FILE" + value = "/var/run/secrets/aws/token" + } + env { + name = "AWS_ROLE_ARN" + value = aws_iam_role.cwagent.arn + } + # CWAGENT_ROLE_ARN is deliberately unset. It only feeds sigv4auth's role_arn, and leaving that + # empty makes the extension fall through to the default credential chain, which picks up the + # projected token via AWS_ROLE_ARN + AWS_WEB_IDENTITY_TOKEN_FILE. Setting it would layer a + # redundant sts:AssumeRole of the same role on top of the session we already have. + env { + name = "RUN_IN_CONTAINER" + value = "True" + } + # Explicit GKE signal so mode detection selects the GCP credential/region + # path without depending on a metadata-server probe from the pod. + env { + name = "RUN_IN_GKE" + value = "True" + } + env { + name = "USE_DEFAULT_CONFIG" + value = "otel" + } + env { + name = "K8S_NODE_NAME" + value_from { + field_ref { + field_path = "spec.nodeName" + } + } + } + env { + name = "HOST_IP" + value_from { + field_ref { + field_path = "status.hostIP" + } + } + } + + volume_mount { + name = "aws-token" + mount_path = "/var/run/secrets/aws" + read_only = true + } + volume_mount { + name = "rootfs" + mount_path = "/rootfs" + read_only = true + } + } + + volume { + name = "aws-token" + projected { + sources { + service_account_token { + audience = "sts.amazonaws.com" + expiration_seconds = 86400 + path = "token" + } + } + } + } + volume { + name = "rootfs" + host_path { + path = "/" + } + } + } + } + } + + depends_on = [ + kubernetes_cluster_role_binding.cwagent, + aws_iam_role_policy_attachment.cwagent_server_policy, + ] +} + +##################################################################### +# Load generator: pushes OTLP to localhost:4318 for 3 min via hostNetwork +##################################################################### +# See otlp_load_generator.sh for what the payloads carry and why. +resource "kubernetes_job_v1" "otlp_load" { + metadata { + name = "otlp-load-generator" + namespace = kubernetes_namespace.cwagent.metadata[0].name + } + + spec { + backoff_limit = 0 + + template { + metadata { + labels = { app = "otlp-load" } + } + + spec { + host_network = true + dns_policy = "ClusterFirstWithHostNet" + restart_policy = "Never" + + container { + name = "load-gen" + image = "curlimages/curl:8.8.0" + command = ["/bin/sh", "-c"] + args = [templatefile("${path.module}/../../otlp_load_generator.sh", { + prefix = "gke" + service_name = local.load_gen_service_name + instance_id = google_container_cluster.cwagent.name + endpoint = "http://127.0.0.1:4318" + duration_seconds = local.load_gen_duration_seconds + })] + } + } + } + } + + wait_for_completion = true + timeouts { + create = "10m" + } + + depends_on = [kubernetes_daemon_set_v1.cwagent] +} + +##################################################################### +# Diagnostics: surface agent pod state and logs in the job output so +# delivery failures are debuggable after the cluster is destroyed. +##################################################################### +# GKE has no ready-made kubeconfig attribute, so build a static one from the cluster +# endpoint and the caller's ADC bearer token (valid ~1h, longer than a run) -- kubectl +# then needs no gcloud auth plugin on the machine running terraform. +resource "local_sensitive_file" "kubeconfig" { + content = yamlencode({ + apiVersion = "v1" + kind = "Config" + clusters = [{ + name = google_container_cluster.cwagent.name + cluster = { + server = "https://${google_container_cluster.cwagent.endpoint}" + "certificate-authority-data" = google_container_cluster.cwagent.master_auth[0].cluster_ca_certificate + } + }] + users = [{ + name = "terraform" + user = { + token = data.google_client_config.current.access_token + } + }] + contexts = [{ + name = google_container_cluster.cwagent.name + context = { + cluster = google_container_cluster.cwagent.name + user = "terraform" + } + }] + "current-context" = google_container_cluster.cwagent.name + }) + filename = "${path.module}/kubeconfig" + file_permission = "0600" +} + +resource "null_resource" "agent_diagnostics" { + provisioner "local-exec" { + command = <<-EOT + kubectl --kubeconfig='${local_sensitive_file.kubeconfig.filename}' get pods -n amazon-cloudwatch -o wide || true + kubectl --kubeconfig='${local_sensitive_file.kubeconfig.filename}' logs -n amazon-cloudwatch -l app=cloudwatch-agent --tail=200 --prefix || true + EOT + } + + depends_on = [kubernetes_job_v1.otlp_load] +} + +##################################################################### +# Run Go integration test from the runner (validates CloudWatch) +##################################################################### +resource "null_resource" "integration_test" { + provisioner "local-exec" { + working_dir = "${path.module}/../../../" + command = <<-EOT + go test -tags integration ${var.test_dir} -p 1 -timeout 30m \ + -computeType=GKE \ + -region=${var.region} \ + -cwaCommitSha=${var.cwa_github_sha} \ + -gkeClusterName=${google_container_cluster.cwagent.name} \ + -v + EOT + + environment = { + AWS_REGION = var.region + } + } + + depends_on = [kubernetes_job_v1.otlp_load, null_resource.agent_diagnostics] +} diff --git a/terraform/gcp/gke/providers.tf b/terraform/gcp/gke/providers.tf new file mode 100644 index 000000000..accf81303 --- /dev/null +++ b/terraform/gcp/gke/providers.tf @@ -0,0 +1,53 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +terraform { + required_providers { + aws = { + source = "hashicorp/aws" + version = "!= 6.22.0" + } + google = { + source = "hashicorp/google" + version = "~> 6.0" + } + kubernetes = { + source = "hashicorp/kubernetes" + version = "~> 2.0" + } + tls = { + source = "hashicorp/tls" + version = "~> 4.0" + } + local = { + source = "hashicorp/local" + version = "~> 2.4" + } + } +} + +provider "aws" { + region = var.region +} + +provider "aws" { + alias = "ecr" + region = var.ecr_region +} + +# Credentials come from GOOGLE_APPLICATION_CREDENTIALS or gcloud application-default +# credentials; project and zone are explicit module inputs. +provider "google" { + project = var.gcp_project + zone = var.gcp_zone +} + +# The same ADC identity's bearer token authenticates directly against the cluster +# endpoint, so the machine running terraform needs no exec-based kubectl auth plugin. +data "google_client_config" "current" {} + +provider "kubernetes" { + host = "https://${google_container_cluster.cwagent.endpoint}" + token = data.google_client_config.current.access_token + cluster_ca_certificate = base64decode(google_container_cluster.cwagent.master_auth[0].cluster_ca_certificate) +} diff --git a/terraform/gcp/gke/variables.tf b/terraform/gcp/gke/variables.tf new file mode 100644 index 000000000..7440d418d --- /dev/null +++ b/terraform/gcp/gke/variables.tf @@ -0,0 +1,81 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +# us-east-2 for the Transaction Search reason documented in terraform/gcp/gce/variables.tf: +# trace validation reads the aws/spans log group, which only exists where the X-Ray trace +# segment destination is CloudWatchLogs. +variable "region" { + type = string + default = "us-east-2" +} + +variable "test_dir" { + type = string + default = "./test/gcp/gke" +} + +variable "cwa_github_sha" { + type = string + default = "" +} + +# Existing project the cluster is created in (input so CI needs no project create/delete perms). +variable "gcp_project" { + type = string + default = "" +} + +# GKE node capacity in us-east1-b has a history of stockouts. +variable "gcp_zone" { + type = string + default = "us-east1-c" +} + +# Existing VPC network + subnetwork the cluster attaches to (must exist; no safe default). +variable "gcp_network_name" { + type = string + default = "" +} + +variable "gcp_subnetwork_name" { + type = string + default = "default" +} + +# Required, not defaulted: the API server is public and this is the only thing scoping it to the runner. +variable "runner_ip" { + type = string + description = "Runner public IP CIDR (e.g. \"1.2.3.4/32\") allowed to reach the GKE API server." +} + +variable "gke_node_machine_type" { + type = string + default = "e2-standard-4" +} + +variable "gke_node_count" { + type = number + default = 1 +} + +variable "kubernetes_version" { + type = string + description = "GKE Kubernetes version. null lets GKE pick the default channel's version." + default = null +} + +variable "cwagent_image_repo" { + type = string + description = "ECR repository URI for the pre-built CWA container image." +} + +variable "ecr_region" { + type = string + description = "Region of the integration-test ECR repository (the build publishes to us-west-2 only)." + default = "us-west-2" +} + +variable "cwagent_image_tag" { + type = string + description = "Image tag (build_id / commit SHA)." +} diff --git a/terraform/azure/aks/otlp_load_generator.sh b/terraform/otlp_load_generator.sh similarity index 72% rename from terraform/azure/aks/otlp_load_generator.sh rename to terraform/otlp_load_generator.sh index f294292a9..0ed32dc0f 100644 --- a/terraform/azure/aks/otlp_load_generator.sh +++ b/terraform/otlp_load_generator.sh @@ -2,7 +2,7 @@ # Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. # SPDX-License-Identifier: MIT # -# OTLP load generator for the AKS integration test. Runs as a Kubernetes Job on the node's network +# OTLP load generator for the Kubernetes integration tests (AKS, GKE). Runs as a Kubernetes Job on the node's network # namespace so 127.0.0.1:4318 reaches the agent's OTLP receiver on the same node. # # Rendered by terraform via templatefile(), so a single-dollar brace expansion is a template @@ -31,13 +31,13 @@ while [ $(date +%s) -lt $END ]; do SPAN_ID=$(printf '%016x' "$NOW_S$SEQ") curl -sf -X POST "$ENDPOINT/v1/metrics" -H "Content-Type: application/json" \ - -d "{\"resourceMetrics\":[{\"resource\":{\"attributes\":[{\"key\":\"service.name\",\"value\":{\"stringValue\":\"$SERVICE_NAME\"}},{\"key\":\"host.id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.cluster.name\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.namespace.name\",\"value\":{\"stringValue\":\"amazon-cloudwatch\"}}]},\"scopeMetrics\":[{\"scope\":{\"name\":\"aks-otlp-test\"},\"metrics\":[{\"name\":\"aks_otlp_counter\",\"unit\":\"1\",\"sum\":{\"aggregationTemporality\":2,\"isMonotonic\":true,\"dataPoints\":[{\"asInt\":\"$SEQ\",\"startTimeUnixNano\":\"$${START}000000000\",\"timeUnixNano\":\"$NOW_NS\",\"attributes\":[{\"key\":\"test_id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}}]}]}}]}]}]}" || true + -d "{\"resourceMetrics\":[{\"resource\":{\"attributes\":[{\"key\":\"service.name\",\"value\":{\"stringValue\":\"$SERVICE_NAME\"}},{\"key\":\"host.id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.cluster.name\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.namespace.name\",\"value\":{\"stringValue\":\"amazon-cloudwatch\"}}]},\"scopeMetrics\":[{\"scope\":{\"name\":\"${prefix}-otlp-test\"},\"metrics\":[{\"name\":\"${prefix}_otlp_counter\",\"unit\":\"1\",\"sum\":{\"aggregationTemporality\":2,\"isMonotonic\":true,\"dataPoints\":[{\"asInt\":\"$SEQ\",\"startTimeUnixNano\":\"$${START}000000000\",\"timeUnixNano\":\"$NOW_NS\",\"attributes\":[{\"key\":\"test_id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}}]}]}}]}]}]}" || true curl -sf -X POST "$ENDPOINT/v1/logs" -H "Content-Type: application/json" \ - -d "{\"resourceLogs\":[{\"resource\":{\"attributes\":[{\"key\":\"service.name\",\"value\":{\"stringValue\":\"$SERVICE_NAME\"}},{\"key\":\"host.id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.cluster.name\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.namespace.name\",\"value\":{\"stringValue\":\"amazon-cloudwatch\"}}]},\"scopeLogs\":[{\"scope\":{\"name\":\"aks-otlp-test\"},\"logRecords\":[{\"timeUnixNano\":\"$NOW_NS\",\"severityText\":\"INFO\",\"body\":{\"stringValue\":\"aks_otlp_log_$INSTANCE_ID\"},\"attributes\":[{\"key\":\"test_id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}}]}]}]}]}" || true + -d "{\"resourceLogs\":[{\"resource\":{\"attributes\":[{\"key\":\"service.name\",\"value\":{\"stringValue\":\"$SERVICE_NAME\"}},{\"key\":\"host.id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.cluster.name\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.namespace.name\",\"value\":{\"stringValue\":\"amazon-cloudwatch\"}}]},\"scopeLogs\":[{\"scope\":{\"name\":\"${prefix}-otlp-test\"},\"logRecords\":[{\"timeUnixNano\":\"$NOW_NS\",\"severityText\":\"INFO\",\"body\":{\"stringValue\":\"${prefix}_otlp_log_$INSTANCE_ID\"},\"attributes\":[{\"key\":\"test_id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}}]}]}]}]}" || true curl -sf -X POST "$ENDPOINT/v1/traces" -H "Content-Type: application/json" \ - -d "{\"resourceSpans\":[{\"resource\":{\"attributes\":[{\"key\":\"service.name\",\"value\":{\"stringValue\":\"$SERVICE_NAME\"}},{\"key\":\"host.id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.cluster.name\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.namespace.name\",\"value\":{\"stringValue\":\"amazon-cloudwatch\"}}]},\"scopeSpans\":[{\"scope\":{\"name\":\"aks-otlp-test\"},\"spans\":[{\"traceId\":\"$TRACE_ID\",\"spanId\":\"$SPAN_ID\",\"name\":\"aks-otlp-test-span\",\"kind\":2,\"startTimeUnixNano\":\"$START_NS\",\"endTimeUnixNano\":\"$NOW_NS\",\"attributes\":[{\"key\":\"test_id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}}]}]}]}]}" || true + -d "{\"resourceSpans\":[{\"resource\":{\"attributes\":[{\"key\":\"service.name\",\"value\":{\"stringValue\":\"$SERVICE_NAME\"}},{\"key\":\"host.id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.cluster.name\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}},{\"key\":\"k8s.namespace.name\",\"value\":{\"stringValue\":\"amazon-cloudwatch\"}}]},\"scopeSpans\":[{\"scope\":{\"name\":\"${prefix}-otlp-test\"},\"spans\":[{\"traceId\":\"$TRACE_ID\",\"spanId\":\"$SPAN_ID\",\"name\":\"${prefix}-otlp-test-span\",\"kind\":2,\"startTimeUnixNano\":\"$START_NS\",\"endTimeUnixNano\":\"$NOW_NS\",\"attributes\":[{\"key\":\"test_id\",\"value\":{\"stringValue\":\"$INSTANCE_ID\"}}]}]}]}]}" || true sleep 10 done diff --git a/test/gcp/gce/gce_test.go b/test/gcp/gce/gce_test.go new file mode 100644 index 000000000..c486a2c07 --- /dev/null +++ b/test/gcp/gce/gce_test.go @@ -0,0 +1,246 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +//go:build integration + +// Package gce validates the agent on a real GCE VM running default:otel: it pushes OTLP to the +// pre-provisioned collector and verifies metrics/logs/traces reach CloudWatch via the GCP web-identity chain. +// Uses the TestMain/pre-provisioned pattern (not test_runner.TestRunner, which would restart the agent). +package gce + +import ( + "flag" + "fmt" + "log" + "os" + "strings" + "testing" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/stretchr/testify/require" + + "github.com/aws/amazon-cloudwatch-agent-test/environment" + "github.com/aws/amazon-cloudwatch-agent-test/test/otel_collect/otlpvalidation" + "github.com/aws/amazon-cloudwatch-agent-test/test/status" + "github.com/aws/amazon-cloudwatch-agent-test/util/awsservice" + "github.com/aws/amazon-cloudwatch-agent-test/util/common" +) + +const ( + // loadWindow is how long OTLP telemetry is pushed before validation; delivery + CloudWatch ingestion + // need headroom beyond the push window. + loadWindow = 3 * time.Minute + otlpEndpoint = "http://127.0.0.1:4318" + // otlpLogGroup is where default:otel routes OTLP logs: "/aws/cwagent" + "/" + aws.log.source ("otlp"). + otlpLogGroup = "/aws/cwagent/otlp" + // agentLogFile lets us confirm the collector booted the GCP web-identity pipeline before asserting delivery. + agentLogFile = "/opt/aws/amazon-cloudwatch-agent/logs/amazon-cloudwatch-agent.log" + // serviceName tags emitted telemetry so validation can isolate this test's records from other traffic. + serviceName = "gce-otlp-test-service" + // spansLogGroup is where Transaction Search stores 100% of spans ingested via the X-Ray OTLP endpoint. + spansLogGroup = "aws/spans" +) + +var ( + env *environment.MetaData + // payloadPrefix namespaces the emitted telemetry (see otlpvalidation.PayloadConfig); it is the + // lowercase compute type, so generation here matches what the validations below grep for. + payloadPrefix string + payloadCfg otlpvalidation.PayloadConfig + traces otlpvalidation.TraceRecorder +) + +func TestMain(m *testing.M) { + environment.RegisterEnvironmentMetaDataFlags() + flag.Parse() + env = environment.GetEnvironmentMetaData() + if env.InstanceId == "" { + fmt.Fprintln(os.Stderr, "instanceId flag is required (GCE numeric instance ID) to scope telemetry") + os.Exit(1) + } + payloadPrefix = strings.ToLower(string(env.ComputeType)) + payloadCfg = otlpvalidation.PayloadConfig{ + Prefix: payloadPrefix, + ServiceName: serviceName, + InstanceID: env.InstanceId, + } + os.Exit(m.Run()) +} + +// TestGCE confirms the pre-provisioned default:otel agent detected GCE, then pushes OTLP and validates +// that all three signals reach CloudWatch via the GCP web-identity chain. +func TestGCE(t *testing.T) { + // The agent must already be running default:otel and have detected GCE before we generate load. + agentLog := common.ReadAgentLogfile(agentLogFile) + require.Contains(t, agentLog, "gcp", + "agent log has no \"gcp\" marker; the default:otel GCE detection path was not exercised") + + // Push OTLP for the load window, then validate. + stop := make(chan struct{}) + senderDone := make(chan struct{}) + go func() { + defer close(senderDone) + sendTelemetry(stop) + }() + time.Sleep(loadWindow) + close(stop) + // Join the sender before reading what it recorded, rather than assuming the settle + // sleep is long enough for its final iteration to finish. + <-senderDone + // Allow final export + CloudWatch ingestion to settle before querying. + time.Sleep(30 * time.Second) + + // Snapshot the accepted trace IDs. The sender has exited, and the recorder's lock still + // gives a clean happens-before with its last append. + traceIDsCopy := traces.Snapshot() + + t.Run("Metrics", func(t *testing.T) { + group := otlpvalidation.ValidateOtlpMetricsWithLabels( + "GCEDefaultOtel", env.Region, otlpvalidation.MeasuredMetricNames(payloadPrefix), + map[string]string{ + "@resource.host.id": env.InstanceId, + "@resource.cloud.provider": "gcp", + }, + ) + for _, r := range group.TestResults { + require.Equal(t, status.SUCCESSFUL, r.Status, "metric %s: %v", r.Name, r.Reason) + } + }) + + t.Run("Logs", func(t *testing.T) { + require.Equal(t, status.SUCCESSFUL, validateLogs().Status) + }) + + t.Run("Traces", func(t *testing.T) { + // Dump agent log errors/warnings from the load window to diagnose trace export issues. + postLoadLog := common.ReadAgentLogfile(agentLogFile) + for _, line := range otlpvalidation.FilterLogLines(postLoadLog, "error", "warn", "xray", "traces", "401", "403", "500") { + t.Logf("agent: %s", line) + } + r := validateTraces(traceIDsCopy) + require.Equal(t, status.SUCCESSFUL, r.Status, "trace validation failed: %v", r.Reason) + }) +} + +// validateLogs confirms the OTLP log record landed in the default:otel log group on the stream the +// agent's log routing is expected to derive for this host. +func validateLogs() status.TestResult { + testResult := status.TestResult{Name: "GCE_Logs", Status: status.FAILED} + + // The agent routes OTLP logs to {host.id}/{service.name}, so assert that exact stream: it makes the + // check prove log routing rather than just delivery, and keeps cost flat as the shared group + // accumulates a stream per VM. Retries because the stream and events both lag. + logStream := fmt.Sprintf("%s/%s", env.InstanceId, serviceName) + // Clean up only on success: the group is shared by every VM run, so drop this run's stream but never + // the group. On failure the stream is left in place as evidence for whoever debugs the run. + defer func() { + if testResult.Status == status.SUCCESSFUL { + awsservice.DeleteLogStream(otlpLogGroup, logStream) + } + }() + marker := otlpvalidation.LogMarker(payloadPrefix, env.InstanceId) + const maxRetries = 4 + const retryInterval = 30 * time.Second + for attempt := 1; attempt <= maxRetries; attempt++ { + since := time.Now().Add(-loadWindow - time.Minute) + until := time.Now() + log.Printf("[GCE_Logs] attempt %d: checking %s/%s", attempt, otlpLogGroup, logStream) + err := awsservice.ValidateLogs( + otlpLogGroup, logStream, &since, &until, + awsservice.AssertLogsNotEmpty(), + awsservice.AssertPerLog(awsservice.AssertLogContainsSubstring(marker)), + ) + if err == nil { + testResult.Status = status.SUCCESSFUL + return testResult + } + testResult.Reason = err + if attempt < maxRetries { + log.Printf("[GCE_Logs] %v — retrying in %v", testResult.Reason, retryInterval) + time.Sleep(retryInterval) + } + } + return testResult +} + +// validateTraces confirms every OTLP span emitted during the load window reached AWS through the +// X-Ray OTLP endpoint. That endpoint requires Transaction Search (trace segment destination = +// CloudWatchLogs), which stores 100% of ingested spans in the aws/spans log group; the X-Ray query +// APIs (GetTraceSummaries/BatchGetTraces) only see the indexed subset (1% by default), so aws/spans +// is the authoritative surface for OTLP trace delivery. Ingestion lags a few minutes, hence retries. +func validateTraces(traceIDs []string) status.TestResult { + testResult := status.TestResult{Name: "GCE_Traces", Status: status.FAILED} + + if len(traceIDs) == 0 { + testResult.Reason = fmt.Errorf("no trace IDs were generated during the load window") + return testResult + } + + quoted := make([]string, len(traceIDs)) + for i, id := range traceIDs { + quoted[i] = fmt.Sprintf("%q", id) + } + query := fmt.Sprintf("fields traceId | filter traceId in [%s] | dedup traceId", strings.Join(quoted, ", ")) + log.Printf("[GCE_Traces] expecting %d trace IDs in %s (sample: %s)", len(traceIDs), spansLogGroup, traceIDs[0]) + + const maxRetries = 5 + const retryInterval = 60 * time.Second + for attempt := 1; attempt <= maxRetries; attempt++ { + since := time.Now().Add(-loadWindow - 10*time.Minute) + rows, err := awsservice.GetLogQueryResults(spansLogGroup, since.Unix(), time.Now().Unix(), query) + if err != nil { + testResult.Reason = fmt.Errorf("attempt %d: %s query failed (is Transaction Search enabled in the account?): %w", + attempt, spansLogGroup, err) + } else { + found := make(map[string]bool, len(rows)) + for _, row := range rows { + for _, field := range row { + if aws.ToString(field.Field) == "traceId" { + found[aws.ToString(field.Value)] = true + } + } + } + var missing []string + for _, id := range traceIDs { + if !found[id] { + missing = append(missing, id) + } + } + if len(missing) == 0 { + log.Printf("[GCE_Traces] attempt %d: all %d traces found in %s", attempt, len(traceIDs), spansLogGroup) + testResult.Status = status.SUCCESSFUL + return testResult + } + testResult.Reason = fmt.Errorf("attempt %d: %d/%d traces missing from %s (first missing: %s)", + attempt, len(missing), len(traceIDs), spansLogGroup, missing[0]) + } + if attempt < maxRetries { + log.Printf("[GCE_Traces] %v — retrying in %v", testResult.Reason, retryInterval) + time.Sleep(retryInterval) + } + } + return testResult +} + +// sendTelemetry pushes OTLP metrics, logs, and traces to the local collector until stop is closed. +func sendTelemetry(stop <-chan struct{}) { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + for { + select { + case <-stop: + return + case <-ticker.C: + otlpvalidation.PostOTLP(otlpEndpoint, "/v1/metrics", otlpvalidation.BuildMetricsPayload(payloadCfg)) + otlpvalidation.PostOTLP(otlpEndpoint, "/v1/logs", otlpvalidation.BuildLogsPayload(payloadCfg)) + // Only record the trace ID once the collector has accepted the span. Recording it + // unconditionally would make a single transient POST failure guarantee a validation + // failure for a trace that was never actually sent. + payload, traceID := otlpvalidation.BuildTracesPayload(payloadCfg) + if otlpvalidation.PostOTLP(otlpEndpoint, "/v1/traces", payload) { + traces.Record(traceID) + } + } + } +} diff --git a/test/gcp/gke/gke_test.go b/test/gcp/gke/gke_test.go new file mode 100644 index 000000000..ea7a9d0d3 --- /dev/null +++ b/test/gcp/gke/gke_test.go @@ -0,0 +1,165 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +//go:build integration + +// Package gke validates the agent on a real GKE cluster running default:otel: a load-generator Job +// pushes OTLP to the DaemonSet agent via hostNetwork, and this test validates metrics/logs/traces +// reach CloudWatch via the GKE projected-token → AWS STS web-identity federation chain. +package gke + +import ( + "flag" + "fmt" + "log" + "os" + "strings" + "testing" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/stretchr/testify/require" + + "github.com/aws/amazon-cloudwatch-agent-test/environment" + "github.com/aws/amazon-cloudwatch-agent-test/test/otel_collect/otlpvalidation" + "github.com/aws/amazon-cloudwatch-agent-test/test/status" + "github.com/aws/amazon-cloudwatch-agent-test/util/awsservice" +) + +const ( + spansLogGroup = "aws/spans" + serviceName = "gke-otlp-test-service" + // The load generator runs for 3 minutes; allow extra ingestion time. + validationWindow = 10 * time.Minute +) + +var ( + env *environment.MetaData + // payloadPrefix is the lowercase compute type: the prefix the load generator template was + // rendered with, so validation greps for exactly what generation emitted. + payloadPrefix string +) + +func TestMain(m *testing.M) { + environment.RegisterEnvironmentMetaDataFlags() + flag.Parse() + env = environment.GetEnvironmentMetaData() + if env.GKEClusterName == "" { + fmt.Fprintln(os.Stderr, "gkeClusterName flag is required to scope telemetry to this cluster") + os.Exit(1) + } + payloadPrefix = strings.ToLower(string(env.ComputeType)) + os.Exit(m.Run()) +} + +func TestGKE(t *testing.T) { + t.Run("Metrics", func(t *testing.T) { + // test_id is a datapoint attribute, the one surface no resource processor rewrites, so it + // isolates this run. cloud.platform=gcp_kubernetes_engine comes only from the gcp detector: + // proves detection ran. + group := otlpvalidation.ValidateOtlpMetricsWithLabels( + "GKEDefaultOtel", env.Region, []string{payloadPrefix + "_otlp_counter"}, + map[string]string{ + "test_id": env.GKEClusterName, + "@resource.cloud.platform": "gcp_kubernetes_engine", + "@resource.cloud.provider": "gcp", + }, + ) + for _, r := range group.TestResults { + require.Equal(t, status.SUCCESSFUL, r.Status, "metric %s: %v", r.Name, r.Reason) + } + }) + + t.Run("Logs", func(t *testing.T) { + r := validateLogs() + require.Equal(t, status.SUCCESSFUL, r.Status, "log validation failed: %v", r.Reason) + }) + + t.Run("Traces", func(t *testing.T) { + r := validateTraces() + require.Equal(t, status.SUCCESSFUL, r.Status, "trace validation failed: %v", r.Reason) + }) +} + +func validateLogs() status.TestResult { + testResult := status.TestResult{Name: "GKE_Logs", Status: status.FAILED} + + // The agent's k8s logs routing derives the destination from the k8s.cluster.name and + // k8s.namespace.name resource attributes the load generator sends, so it is unique to + // this cluster. The stream is {k8s.namespace.name}/{service.namespace}/{service.name}, + // where the agent's identity transform fills service.namespace from k8s.namespace.name. + // AssertLogsNotEmpty guards against a vacuous pass on an empty window. + logGroup := fmt.Sprintf("/aws/cwagent/%s/otlp", env.GKEClusterName) + // Clean up only on success: the group name carries this run's cluster so the whole group is + // disposable, but on failure it is left in place as evidence for whoever debugs the run. + defer func() { + if testResult.Status == status.SUCCESSFUL { + awsservice.DeleteLogGroup(logGroup) + } + }() + logStream := fmt.Sprintf("amazon-cloudwatch/amazon-cloudwatch/%s", serviceName) + marker := otlpvalidation.LogMarker(payloadPrefix, env.GKEClusterName) + const maxRetries = 4 + const retryInterval = 30 * time.Second + for attempt := 1; attempt <= maxRetries; attempt++ { + since := time.Now().Add(-validationWindow) + until := time.Now() + log.Printf("[GKE_Logs] attempt %d: checking %s/%s", attempt, logGroup, logStream) + err := awsservice.ValidateLogs( + logGroup, logStream, &since, &until, + awsservice.AssertLogsNotEmpty(), + awsservice.AssertPerLog(awsservice.AssertLogContainsSubstring(marker)), + ) + if err == nil { + testResult.Status = status.SUCCESSFUL + return testResult + } + testResult.Reason = err + if attempt < maxRetries { + log.Printf("[GKE_Logs] %v — retrying in %v", testResult.Reason, retryInterval) + time.Sleep(retryInterval) + } + } + return testResult +} + +// validateTraces queries aws/spans (Transaction Search) for spans with our cluster's service name. +func validateTraces() status.TestResult { + testResult := status.TestResult{Name: "GKE_Traces", Status: status.FAILED} + + query := fmt.Sprintf( + `fields traceId | filter @message like "%s" and @message like "%s" | dedup traceId | limit 5`, + serviceName, env.GKEClusterName, + ) + log.Printf("[GKE_Traces] querying %s for spans from service=%s instance=%s", spansLogGroup, serviceName, env.GKEClusterName) + + const maxRetries = 5 + const retryInterval = 60 * time.Second + for attempt := 1; attempt <= maxRetries; attempt++ { + since := time.Now().Add(-validationWindow) + rows, err := awsservice.GetLogQueryResults(spansLogGroup, since.Unix(), time.Now().Unix(), query) + if err != nil { + testResult.Reason = fmt.Errorf("attempt %d: %s query failed: %w", attempt, spansLogGroup, err) + } else { + found := 0 + for _, row := range rows { + for _, field := range row { + if aws.ToString(field.Field) == "traceId" && aws.ToString(field.Value) != "" { + found++ + } + } + } + if found > 0 { + log.Printf("[GKE_Traces] attempt %d: found %d traces in %s", attempt, found, spansLogGroup) + testResult.Status = status.SUCCESSFUL + return testResult + } + testResult.Reason = fmt.Errorf("attempt %d: 0 traces found in %s for service=%s", attempt, spansLogGroup, serviceName) + } + if attempt < maxRetries { + log.Printf("[GKE_Traces] %v — retrying in %v", testResult.Reason, retryInterval) + time.Sleep(retryInterval) + } + } + return testResult +} diff --git a/test/otel_collect/otlpvalidation/payloads.go b/test/otel_collect/otlpvalidation/payloads.go new file mode 100644 index 000000000..50a060d1f --- /dev/null +++ b/test/otel_collect/otlpvalidation/payloads.go @@ -0,0 +1,210 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: MIT + +package otlpvalidation + +import ( + "bytes" + "fmt" + "io" + "log" + "net/http" + "strings" + "sync" + "time" +) + +// PayloadConfig parameterizes the OTLP payloads shared by the VM integration +// tests. Prefix namespaces every metric, log marker, scope, and span name +// (e.g. "gce" produces gce_otlp_counter), so each platform's telemetry is +// distinguishable in the shared destinations. +type PayloadConfig struct { + Prefix string + ServiceName string + InstanceID string +} + +// startTimeNano is captured once so counter data points share a stable start time across the run. +var startTimeNano = time.Now().UnixNano() + +// MeasuredMetricNames returns the metric names BuildMetricsPayload emits for a prefix. +func MeasuredMetricNames(prefix string) []string { + return []string{prefix + "_otlp_counter", prefix + "_otlp_gauge"} +} + +// LogMarker returns the log body BuildLogsPayload emits: a per-host marker the +// log validation greps for. +func LogMarker(prefix, id string) string { + return fmt.Sprintf("%s_otlp_log_%s", prefix, id) +} + +// BuildMetricsPayload emits a monotonic counter and a gauge tagged with this host's instance id so the +// CloudWatch OTLP PromQL query can isolate them. service.name lets the collector derive a stable stream/scope. +func BuildMetricsPayload(cfg PayloadConfig) []byte { + now := time.Now().UnixNano() + return []byte(fmt.Sprintf(`{ + "resourceMetrics": [{ + "resource": {"attributes": [ + {"key": "service.name", "value": {"stringValue": "%s"}}, + {"key": "host.id", "value": {"stringValue": "%s"}} + ]}, + "scopeMetrics": [{ + "scope": {"name": "%s-otlp-test-metrics", "version": "1.0.0"}, + "metrics": [ + { + "name": "%s_otlp_counter", + "unit": "1", + "sum": { + "aggregationTemporality": 2, + "isMonotonic": true, + "dataPoints": [{"asInt": "1", "startTimeUnixNano": "%d", "timeUnixNano": "%d", "attributes": [{"key": "InstanceId", "value": {"stringValue": "%s"}}]}] + } + }, + { + "name": "%s_otlp_gauge", + "unit": "1", + "gauge": { + "dataPoints": [{"asDouble": 42.0, "timeUnixNano": "%d", "attributes": [{"key": "InstanceId", "value": {"stringValue": "%s"}}]}] + } + } + ] + }] + }] +}`, cfg.ServiceName, cfg.InstanceID, cfg.Prefix, cfg.Prefix, startTimeNano, now, cfg.InstanceID, cfg.Prefix, now, cfg.InstanceID)) +} + +// BuildLogsPayload emits an OTLP log whose body carries a per-host marker (see LogMarker) +// for CloudWatch Logs validation. +func BuildLogsPayload(cfg PayloadConfig) []byte { + now := time.Now().UnixNano() + return []byte(fmt.Sprintf(`{ + "resourceLogs": [{ + "resource": {"attributes": [ + {"key": "service.name", "value": {"stringValue": "%s"}}, + {"key": "host.id", "value": {"stringValue": "%s"}} + ]}, + "scopeLogs": [{ + "scope": {"name": "%s-otlp-test-logs"}, + "logRecords": [{ + "timeUnixNano": "%d", + "severityText": "INFO", + "body": {"stringValue": "%s"}, + "attributes": [{"key": "InstanceId", "value": {"stringValue": "%s"}}] + }] + }] + }] +}`, cfg.ServiceName, cfg.InstanceID, cfg.Prefix, now, LogMarker(cfg.Prefix, cfg.InstanceID), cfg.InstanceID)) +} + +// traceSeq is an incrementing counter ensuring unique trace/span IDs across calls. +var traceSeq uint64 + +// traceSeqMu protects traceSeq, which is incremented from the sender goroutine. +var traceSeqMu sync.Mutex + +// BuildTracesPayload emits an OTLP span whose trace ID follows the X-Ray format: the first 4 bytes hold +// the Unix epoch in seconds, matching what the X-Ray propagator's own ID generator does and what every +// other trace producer in this repo does. X-Ray rejects IDs whose embedded date is too far in the past +// with InvalidTraceId, and randomly generated IDs were silently dropped during bring-up. +// +// Transaction Search stores spans with W3C trace IDs, so it is possible this prefix is no longer +// required on the ingest path -- that has not been re-verified. Keeping it is valid either way. +func BuildTracesPayload(cfg PayloadConfig) ([]byte, string) { + traceSeqMu.Lock() + traceSeq++ + now := time.Now() + nowNano := now.UnixNano() + startNano := nowNano - int64(time.Second) + // First 4 bytes: unix seconds (X-Ray requirement). Remaining 12 bytes: sequence + padding for uniqueness. + traceID := fmt.Sprintf("%08x0000000000000000%08x", now.Unix(), traceSeq) + spanID := fmt.Sprintf("%016x", nowNano) + traceSeqMu.Unlock() + return []byte(fmt.Sprintf(`{ + "resourceSpans": [{ + "resource": {"attributes": [ + {"key": "service.name", "value": {"stringValue": "%s"}}, + {"key": "host.id", "value": {"stringValue": "%s"}} + ]}, + "scopeSpans": [{ + "scope": {"name": "%s-otlp-test-traces"}, + "spans": [{ + "traceId": "%s", + "spanId": "%s", + "name": "%s-otlp-test-span", + "kind": 2, + "startTimeUnixNano": "%d", + "endTimeUnixNano": "%d", + "attributes": [{"key": "instance_id", "value": {"stringValue": "%s"}}] + }] + }] + }] +}`, cfg.ServiceName, cfg.InstanceID, cfg.Prefix, traceID, spanID, cfg.Prefix, startNano, nowNano, cfg.InstanceID)), traceID +} + +// TraceRecorder collects the trace IDs the collector accepted during the load +// window. It is written by the sender goroutine and read by the test goroutine +// after the window closes, so access is mutex-guarded. +type TraceRecorder struct { + mu sync.Mutex + ids []string +} + +// Record marks a trace ID as successfully delivered, so trace validation expects to find it. +func (r *TraceRecorder) Record(traceID string) { + r.mu.Lock() + defer r.mu.Unlock() + r.ids = append(r.ids, traceID) +} + +// Snapshot returns a copy of the recorded trace IDs. +func (r *TraceRecorder) Snapshot() []string { + r.mu.Lock() + defer r.mu.Unlock() + out := make([]string, len(r.ids)) + copy(out, r.ids) + return out +} + +// PostOTLP sends an OTLP payload and reports whether the collector accepted it. +func PostOTLP(endpoint, path string, payload []byte) bool { + req, err := http.NewRequest("POST", endpoint+path, bytes.NewReader(payload)) + if err != nil { + log.Printf("failed to build OTLP request for %s: %v", path, err) + return false + } + req.Header.Set("Content-Type", "application/json") + resp, err := http.DefaultClient.Do(req) + if err != nil { + log.Printf("failed to POST OTLP to %s: %v", path, err) + return false + } + // Drain before closing so the connection can be reused. + defer func() { + _, _ = io.Copy(io.Discard, resp.Body) + resp.Body.Close() + }() + if resp.StatusCode < 200 || resp.StatusCode > 299 { + log.Printf("OTLP POST to %s returned %s", path, resp.Status) + return false + } + return true +} + +// FilterLogLines returns lines from a multi-line string that contain any of the given +// substrings (case-insensitive), capped to the trailing 50 matches. +func FilterLogLines(text string, substrs ...string) []string { + var result []string + for _, line := range strings.Split(text, "\n") { + lower := strings.ToLower(line) + for _, s := range substrs { + if strings.Contains(lower, strings.ToLower(s)) { + result = append(result, line) + break + } + } + } + if len(result) > 50 { + result = result[len(result)-50:] + } + return result +}