From 14b73a18a2226cb25431a966903a38b0c8c5084c Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 17:03:29 +0200 Subject: [PATCH 1/4] feat(data-warehouse): implement the aws compute optimizer import source Rebase the AWS Compute Optimizer source onto master and preserve the source catalog changes from both branches. --- .../temporal/data_imports/sources/SOURCES.md | 2 +- .../aws_compute_optimizer.py | 233 ++++++++++++ .../canonical_descriptions.py | 75 ++++ .../sources/aws_compute_optimizer/settings.py | 69 ++++ .../sources/aws_compute_optimizer/source.py | 140 ++++++- .../tests/test_aws_compute_optimizer.py | 343 ++++++++++++++++++ .../generated_configs/awscomputeoptimizer.py | 5 +- 7 files changed, 859 insertions(+), 8 deletions(-) create mode 100644 products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/aws_compute_optimizer.py create mode 100644 products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/canonical_descriptions.py create mode 100644 products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/settings.py create mode 100644 products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/tests/test_aws_compute_optimizer.py diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md b/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md index f0134539c223..4bdbeb1d7dcc 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md @@ -104,6 +104,7 @@ the row lists both. | awin | HTTP | requests | ✅ | | aws_budgets | HTTP | requests | ✅ | | aws_cloudtrail | HTTP | requests | ✅ | +| aws_compute_optimizer | HTTP | requests | ✅ | | aws_cost_anomaly_detection | HTTP | requests | ✅ | | aws_cost_explorer | HTTP | requests | ✅ | | aws_glue_data_catalog | HTTP | requests | ✅ | @@ -918,7 +919,6 @@ doesn't conflict with concurrent PRs. - aws_athena - aws_batch - aws_cloudformation -- aws_compute_optimizer - aws_config - aws_connect - aws_cost_and_usage_report diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/aws_compute_optimizer.py b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/aws_compute_optimizer.py new file mode 100644 index 000000000000..c4c302b6228e --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/aws_compute_optimizer.py @@ -0,0 +1,233 @@ +import re +import json +from collections.abc import Iterator +from datetime import UTC, datetime +from typing import Any + +import requests +from botocore.auth import SigV4Auth +from botocore.awsrequest import AWSRequest +from botocore.credentials import Credentials +from tenacity import retry, retry_if_exception_type, stop_after_attempt, wait_exponential_jitter + +from posthog.dataclasses import frozen + +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.settings import ( + API_VERSION, + DEFAULT_REGION, + ENDPOINTS, + ERROR_MESSAGES, + PAGE_SIZE, + TARGET_PREFIXES, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.common.http import make_tracked_session +from products.warehouse_sources.backend.temporal.data_imports.sources.common.http.transport import BoundedRetry +from products.warehouse_sources.backend.temporal.data_imports.sources.common.resumable import ResumableSourceManager +from products.warehouse_sources.backend.temporal.data_imports.sources.common.typings import SourceResponse +from products.warehouse_sources.backend.temporal.data_imports.sources.generated_configs.awscomputeoptimizer import ( + AwsComputeOptimizerSourceConfig, +) + +_CAMEL_BOUNDARY = re.compile(r"(?<=[a-z0-9])(?=[A-Z])|(?<=[A-Z])(?=[A-Z][a-z])") +TRANSPORT_RETRY = BoundedRetry( + total=3, + backoff_factor=1, + status_forcelist=(429, 500, 502, 503, 504), + allowed_methods=frozenset({"POST"}), + raise_on_status=False, +) + + +@frozen +class AwsComputeOptimizerResumeConfig: + next_token: str | None = None + complete: bool = False + + +class AwsComputeOptimizerError(Exception): + def __init__(self, code: str) -> None: + self.code = code + super().__init__(f"AWS Compute Optimizer request failed: {code}") + + +class AwsComputeOptimizerThrottleError(AwsComputeOptimizerError): + pass + + +def error_for_response(response: requests.Response) -> AwsComputeOptimizerError: + try: + body = response.json() + except ValueError: + body = {} + if not isinstance(body, dict): + body = {} + raw_code = response.headers.get("x-amzn-ErrorType") or body.get("__type") or body.get("code") + code = str(raw_code or f"HTTP {response.status_code}").split(":", 1)[0].rsplit("#", 1)[-1] + # The transport retries HTTP statuses. AWS also returns throttling errors with HTTP 400. + if response.status_code == 400 and code in { + "ThrottlingException", + "TooManyRequestsException", + "RequestLimitExceeded", + }: + return AwsComputeOptimizerThrottleError(code) + return AwsComputeOptimizerError(code) + + +class AwsComputeOptimizerClient: + def __init__(self, config: AwsComputeOptimizerSourceConfig, api_version: str | None = None) -> None: + version = api_version or API_VERSION + if version not in TARGET_PREFIXES: + raise ValueError("Unsupported AWS Compute Optimizer API version.") + if not config.aws_access_key_id or not config.aws_secret_access_key: + raise ValueError("Enter both an AWS access key ID and a secret access key.") + self.region = config.region or DEFAULT_REGION + if not re.fullmatch(r"(?:us|eu|ap|sa|ca|me|af|il|mx|cn)-(?:[a-z]+-)+[0-9]+", self.region): + raise ValueError("Enter an AWS region, such as us-east-1.") + suffix = "amazonaws.com.cn" if self.region.startswith("cn-") else "amazonaws.com" + self.url = f"https://compute-optimizer.{self.region}.{suffix}/" + self.target_prefix = TARGET_PREFIXES[version] + self.signer = SigV4Auth( + Credentials(config.aws_access_key_id, config.aws_secret_access_key, config.aws_session_token or None), + "compute-optimizer", + self.region, + ) + self.session = make_tracked_session( + retry=TRANSPORT_RETRY, + redact_values=tuple( + value + for value in (config.aws_access_key_id, config.aws_secret_access_key, config.aws_session_token) + if value + ), + ) + + def close(self) -> None: + self.session.close() + + @retry( + retry=retry_if_exception_type(AwsComputeOptimizerThrottleError), + stop=stop_after_attempt(4), + wait=wait_exponential_jitter(initial=1, max=10), + reraise=True, + ) + def request(self, operation: str, payload: dict[str, Any]) -> dict[str, Any]: + body = json.dumps(payload).encode("utf-8") + request = AWSRequest( + method="POST", + url=self.url, + data=body, + headers={ + "Content-Type": "application/x-amz-json-1.0", + "X-Amz-Target": f"{self.target_prefix}.{operation}", + }, + ) + self.signer.add_auth(request) + response = self.session.post( + self.url, data=body, headers=dict(request.headers), timeout=60, allow_redirects=False + ) + if not 200 <= response.status_code < 300: + raise error_for_response(response) + result = response.json() + if not isinstance(result, dict): + raise ValueError("AWS Compute Optimizer returned an invalid response.") + if result.get("errors"): + error = result["errors"][0] + raise AwsComputeOptimizerError(str(error.get("code") or "PartialResponseError")) + return result + + +def normalize_row(item: dict[str, Any], region: str) -> dict[str, Any]: + def flatten(obj: dict[str, Any], prefix: str = "") -> dict[str, Any]: + row: dict[str, Any] = {} + for key, value in obj.items(): + column = prefix + _CAMEL_BOUNDARY.sub("_", key).lower() + if isinstance(value, dict): + row.update(flatten(value, column + "_")) + else: + row[column] = value + return row + + row = flatten(item) + timestamp = row.get("last_refresh_timestamp") + if isinstance(timestamp, (int, float)) and not isinstance(timestamp, bool): + row["last_refresh_timestamp"] = datetime.fromtimestamp(timestamp, tz=UTC) + row["region"] = region + return row + + +def get_rows( + config: AwsComputeOptimizerSourceConfig, + endpoint: str, + manager: ResumableSourceManager[AwsComputeOptimizerResumeConfig], + api_version: str | None, +) -> Iterator[list[dict[str, Any]]]: + settings = ENDPOINTS[endpoint] + state = manager.load_state() or AwsComputeOptimizerResumeConfig() + if state.complete: + return + token = state.next_token + seen_tokens: set[str] = {token} if token else set() + client = AwsComputeOptimizerClient(config, api_version) + try: + while True: + payload: dict[str, Any] = {"maxResults": PAGE_SIZE} + if token: + payload["nextToken"] = token + result = client.request(settings.operation, payload) + rows = [normalize_row(item, client.region) for item in result.get(settings.result_key, [])] + token = result.get("nextToken") or None + if token is not None: + if token in seen_tokens: + raise ValueError("AWS Compute Optimizer repeated a pagination token. Restart the sync.") + seen_tokens.add(token) + manager.save_state(AwsComputeOptimizerResumeConfig(next_token=token, complete=token is None)) + if rows: + yield rows + manager.safe_point() + if token is None: + return + finally: + client.close() + + +def validate_credentials( + config: AwsComputeOptimizerSourceConfig, schema_name: str | None = None, api_version: str | None = None +) -> tuple[bool, str | None]: + if schema_name is not None and schema_name not in ENDPOINTS: + return False, "Unknown AWS Compute Optimizer table. Select a supported table." + try: + client = AwsComputeOptimizerClient(config, api_version) + except ValueError as error: + return False, str(error) + operation = ENDPOINTS[schema_name].operation if schema_name is not None else "GetEnrollmentStatus" + try: + result = client.request(operation, {"maxResults": 1} if schema_name is not None else {}) + if schema_name is None and result.get("status") != "Active": + return False, "Enable AWS Compute Optimizer for this account and wait until enrollment is active." + except AwsComputeOptimizerError as error: + if error.code in {"AccessDenied", "AccessDeniedException"}: + if schema_name is None: + return True, None + return False, f"AWS denied access. Grant compute-optimizer:{operation} to the IAM user or role." + return False, ERROR_MESSAGES.get(error.code, "Could not read AWS Compute Optimizer. Try again.") + except (requests.RequestException, ValueError): + return False, "Could not reach the AWS Compute Optimizer API. Try again." + finally: + client.close() + return True, None + + +def aws_compute_optimizer_source( + config: AwsComputeOptimizerSourceConfig, + endpoint: str, + manager: ResumableSourceManager[AwsComputeOptimizerResumeConfig], + api_version: str | None = None, +) -> SourceResponse: + if endpoint not in ENDPOINTS: + raise ValueError("Unknown AWS Compute Optimizer table. Select a supported table.") + return SourceResponse( + name=endpoint, + items=lambda: get_rows(config, endpoint, manager, api_version), + primary_keys=list(ENDPOINTS[endpoint].primary_keys), + sort_mode=None, + on_complete=manager.clear_state, + ) diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/canonical_descriptions.py b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/canonical_descriptions.py new file mode 100644 index 000000000000..746836aa5f18 --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/canonical_descriptions.py @@ -0,0 +1,75 @@ +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.settings import ( + API_DOCS_URL, + ENDPOINTS, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.common.canonical_descriptions import ( + CanonicalDescriptions, +) + +CANONICAL_DESCRIPTIONS: CanonicalDescriptions = { + name: { + "description": endpoint.description, + "docs_url": f"{API_DOCS_URL}API_{endpoint.operation}.html", + "columns": { + "account_id": "AWS account that owns the resources.", + "region": "AWS region from which this source reads recommendations.", + }, + } + for name, endpoint in ENDPOINTS.items() +} + +for name in ENDPOINTS: + if name != "recommendation_summaries": + CANONICAL_DESCRIPTIONS[name]["columns"].update( + { + "finding": "Classification of the resource's optimization status.", + "last_refresh_timestamp": "Time when Compute Optimizer last refreshed the recommendation, in UTC.", + ( + "lookback_period_in_days" + if name in {"lambda_function_recommendations", "ecs_service_recommendations"} + else "look_back_period_in_days" + ): "Number of days of utilization data used to generate the recommendation.", + "utilization_metrics": "Utilization measurements used to evaluate the resource.", + } + ) + +CANONICAL_DESCRIPTIONS["ec2_instance_recommendations"]["columns"].update( + { + "instance_arn": "Amazon Resource Name of the EC2 instance.", + "current_instance_type": "EC2 instance type used by the resource.", + "recommendation_options": "Recommended instance types, performance risks, utilization projections, and potential savings.", + } +) +CANONICAL_DESCRIPTIONS["auto_scaling_group_recommendations"]["columns"].update( + { + "auto_scaling_group_arn": "Amazon Resource Name of the Auto Scaling group.", + "recommendation_options": "Recommended configurations for the Auto Scaling group.", + } +) +CANONICAL_DESCRIPTIONS["lambda_function_recommendations"]["columns"].update( + { + "function_arn": "Amazon Resource Name of the Lambda function.", + "function_version": "Version of the Lambda function evaluated by Compute Optimizer.", + "memory_size_recommendation_options": "Recommended memory sizes with projected utilization and potential savings.", + } +) +CANONICAL_DESCRIPTIONS["ecs_service_recommendations"]["columns"].update( + { + "service_arn": "Amazon Resource Name of the ECS service.", + "service_recommendation_options": "Recommended CPU and memory configurations with potential savings.", + } +) +CANONICAL_DESCRIPTIONS["ebs_volume_recommendations"]["columns"].update( + { + "volume_arn": "Amazon Resource Name of the EBS volume.", + "volume_recommendation_options": "Recommended volume configurations with performance risks and potential savings.", + } +) +CANONICAL_DESCRIPTIONS["recommendation_summaries"]["columns"].update( + { + "recommendation_resource_type": "Type of AWS resource covered by the summary.", + "summaries": "Counts of resources grouped by optimization finding.", + "savings_opportunity_estimated_monthly_savings_value": "Estimated monthly savings for this resource type.", + "savings_opportunity_estimated_monthly_savings_currency": "Currency used for the estimated monthly savings.", + } +) diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/settings.py b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/settings.py new file mode 100644 index 000000000000..b6e5106a81cc --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/settings.py @@ -0,0 +1,69 @@ +from posthog.dataclasses import frozen + +API_VERSION = "2019-11-01" +API_DOCS_URL = "https://docs.aws.amazon.com/compute-optimizer/latest/APIReference/" +TARGET_PREFIXES = {API_VERSION: "ComputeOptimizerService"} +DEFAULT_REGION = "us-east-1" +PAGE_SIZE = 100 + + +@frozen +class Endpoint: + operation: str + result_key: str + primary_keys: tuple[str, ...] + description: str + + +ENDPOINTS: dict[str, Endpoint] = { + "ec2_instance_recommendations": Endpoint( + operation="GetEC2InstanceRecommendations", + result_key="instanceRecommendations", + primary_keys=("instance_arn",), + description="EC2 instance recommendations with utilization metrics and estimated savings.", + ), + "auto_scaling_group_recommendations": Endpoint( + operation="GetAutoScalingGroupRecommendations", + result_key="autoScalingGroupRecommendations", + primary_keys=("auto_scaling_group_arn",), + description="Auto Scaling group recommendations with configuration options and estimated savings.", + ), + "lambda_function_recommendations": Endpoint( + operation="GetLambdaFunctionRecommendations", + result_key="lambdaFunctionRecommendations", + primary_keys=("function_arn", "function_version"), + description="Lambda function recommendations with memory options and estimated savings.", + ), + "ecs_service_recommendations": Endpoint( + operation="GetECSServiceRecommendations", + result_key="ecsServiceRecommendations", + primary_keys=("service_arn",), + description="ECS service recommendations with CPU, memory, and estimated savings.", + ), + "ebs_volume_recommendations": Endpoint( + operation="GetEBSVolumeRecommendations", + result_key="volumeRecommendations", + primary_keys=("volume_arn",), + description="EBS volume recommendations with configuration options and estimated savings.", + ), + "recommendation_summaries": Endpoint( + operation="GetRecommendationSummaries", + result_key="recommendationSummaries", + primary_keys=("account_id", "recommendation_resource_type", "region"), + description="Counts of optimization findings and estimated savings by account and resource type.", + ), +} + +ERROR_MESSAGES = { + "AccessDeniedException": "AWS denied access. Grant the compute-optimizer read permission for the selected table.", + "AccessDenied": "AWS denied access. Grant the compute-optimizer read permission for the selected table.", + "UnrecognizedClientException": "AWS rejected the credentials. Check the access key ID, secret access key, and session token.", + "InvalidClientTokenId": "AWS rejected the access key. Check that the key is active.", + "InvalidSignatureException": "AWS rejected the signature. Check the secret access key and session token.", + "SignatureDoesNotMatch": "AWS rejected the signature. Check the secret access key.", + "ExpiredTokenException": "The AWS session token expired. Reconnect with new credentials.", + "ExpiredToken": "The AWS session token expired. Reconnect with new credentials.", + "MissingAuthenticationToken": "AWS requires credentials. Enter an access key ID and secret access key.", + "OptInRequiredException": "Enable AWS Compute Optimizer for this account before you sync recommendations.", + "SubscriptionRequiredException": "Enable AWS Compute Optimizer for this account before you sync recommendations.", +} diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/source.py b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/source.py index f4ae6e25fe35..b7dd6a20c73d 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/source.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/source.py @@ -1,8 +1,37 @@ from typing import cast -from products.warehouse_sources.backend.facade.source_config import DataWarehouseSourceCategory, SourceConfig -from products.warehouse_sources.backend.temporal.data_imports.sources.common.base import FieldType, SimpleSource +from products.warehouse_sources.backend.facade.source_config import ( + DataWarehouseSourceCategory, + ReleaseStatus, + SourceConfig, + SourceFieldInputConfig, + SourceFieldInputConfigType, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.aws_compute_optimizer import ( + AwsComputeOptimizerResumeConfig, + aws_compute_optimizer_source, + validate_credentials, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.canonical_descriptions import ( + CANONICAL_DESCRIPTIONS, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.settings import ( + API_DOCS_URL, + API_VERSION, + ENDPOINTS, + ERROR_MESSAGES, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.common.base import FieldType, ResumableSource +from products.warehouse_sources.backend.temporal.data_imports.sources.common.canonical_descriptions import ( + CanonicalDescriptions, +) from products.warehouse_sources.backend.temporal.data_imports.sources.common.registry import SourceRegistry +from products.warehouse_sources.backend.temporal.data_imports.sources.common.resumable import ResumableSourceManager +from products.warehouse_sources.backend.temporal.data_imports.sources.common.schema import ( + SourceSchema, + build_endpoint_schemas, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.common.typings import SourceInputs, SourceResponse from products.warehouse_sources.backend.temporal.data_imports.sources.generated_configs.awscomputeoptimizer import ( AwsComputeOptimizerSourceConfig, ) @@ -10,18 +39,117 @@ @SourceRegistry.register -class AwsComputeOptimizerSource(SimpleSource[AwsComputeOptimizerSourceConfig]): +class AwsComputeOptimizerSource(ResumableSource[AwsComputeOptimizerSourceConfig, AwsComputeOptimizerResumeConfig]): + lists_tables_without_credentials = True + supported_versions = (API_VERSION,) + default_version = API_VERSION + api_docs_url = API_DOCS_URL + @property def source_type(self) -> ExternalDataSourceType: return ExternalDataSourceType.AWSCOMPUTEOPTIMIZER + def get_non_retryable_errors(self) -> dict[str, str | None]: + return {f"AWS Compute Optimizer request failed: {code}": message for code, message in ERROR_MESSAGES.items()} + + def get_canonical_descriptions(self) -> CanonicalDescriptions: + return CANONICAL_DESCRIPTIONS + + def get_schemas( + self, + config: AwsComputeOptimizerSourceConfig, + team_id: int, + with_counts: bool = False, + names: list[str] | None = None, + force_refresh: bool = False, + api_version: str | None = None, + ) -> list[SourceSchema]: + return build_endpoint_schemas( + ENDPOINTS, {}, names, descriptions={name: endpoint.description for name, endpoint in ENDPOINTS.items()} + ) + + def validate_credentials( + self, + config: AwsComputeOptimizerSourceConfig, + team_id: int, + schema_name: str | None = None, + api_version: str | None = None, + ) -> tuple[bool, str | None]: + return validate_credentials(config, schema_name, self.resolve_api_version(api_version)) + + def get_resumable_source_manager( + self, inputs: SourceInputs + ) -> ResumableSourceManager[AwsComputeOptimizerResumeConfig]: + return ResumableSourceManager(inputs, AwsComputeOptimizerResumeConfig) + + def source_for_pipeline( + self, + config: AwsComputeOptimizerSourceConfig, + resumable_source_manager: ResumableSourceManager[AwsComputeOptimizerResumeConfig], + inputs: SourceInputs, + ) -> SourceResponse: + return aws_compute_optimizer_source( + config, inputs.schema_name, resumable_source_manager, self.resolve_api_version(inputs.api_version) + ) + @property def get_source_config(self) -> SourceConfig: return SourceConfig( name=ExternalDataSourceType.AWSCOMPUTEOPTIMIZER, category=DataWarehouseSourceCategory.ENGINEERING___MONITORING, - label="Amazon Web Services (AWS Compute Optimizer)", + label="AWS Compute Optimizer", + caption="""Sync AWS Compute Optimizer recommendations for the connected account in one AWS region. + +Enable Compute Optimizer for this account before you sync. +Grant `compute-optimizer:GetEnrollmentStatus` to check enrollment. +Grant these permissions for the tables you select: +- `compute-optimizer:GetEC2InstanceRecommendations` +- `compute-optimizer:GetAutoScalingGroupRecommendations` +- `compute-optimizer:GetLambdaFunctionRecommendations` +- `compute-optimizer:GetECSServiceRecommendations` +- `compute-optimizer:GetEBSVolumeRecommendations` +- `compute-optimizer:GetRecommendationSummaries` + +Each sync replaces the current recommendations. This source does not keep recommendation history.""", iconPath="/static/services/aws_compute_optimizer.png", - fields=cast(list[FieldType], []), - unreleasedSource=True, + releaseStatus=ReleaseStatus.ALPHA, + fields=cast( + list[FieldType], + [ + SourceFieldInputConfig( + name="aws_access_key_id", + label="AWS access key ID", + type=SourceFieldInputConfigType.TEXT, + required=True, + placeholder="AKIA...", + secret=False, + ), + SourceFieldInputConfig( + name="aws_secret_access_key", + label="AWS secret access key", + type=SourceFieldInputConfigType.PASSWORD, + required=True, + secret=True, + placeholder="", + ), + SourceFieldInputConfig( + name="aws_session_token", + label="AWS session token", + type=SourceFieldInputConfigType.PASSWORD, + required=False, + secret=True, + placeholder="Only needed for temporary credentials", + caption="Required for temporary credentials. Update the credentials before the token expires.", + ), + SourceFieldInputConfig( + name="region", + label="AWS region", + type=SourceFieldInputConfigType.TEXT, + required=False, + secret=False, + placeholder="us-east-1", + caption="Sync recommendations from this region. The default is us-east-1.", + ), + ], + ), ) diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/tests/test_aws_compute_optimizer.py b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/tests/test_aws_compute_optimizer.py new file mode 100644 index 000000000000..92c817a6bcbf --- /dev/null +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/aws_compute_optimizer/tests/test_aws_compute_optimizer.py @@ -0,0 +1,343 @@ +import json +from collections.abc import Generator, Iterable, Iterator +from datetime import UTC, datetime +from typing import Any, cast + +import pytest +from unittest.mock import MagicMock, patch + +import requests +from botocore.loaders import Loader + +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer import ( + aws_compute_optimizer as transport, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.aws_compute_optimizer import ( + AwsComputeOptimizerClient, + AwsComputeOptimizerError, + AwsComputeOptimizerResumeConfig, + aws_compute_optimizer_source, + validate_credentials, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.aws_compute_optimizer.source import ( + AwsComputeOptimizerSource, +) +from products.warehouse_sources.backend.temporal.data_imports.sources.generated_configs.awscomputeoptimizer import ( + AwsComputeOptimizerSourceConfig, +) + + +def config(region: str | None = "eu-west-1", token: str | None = None) -> AwsComputeOptimizerSourceConfig: + return AwsComputeOptimizerSourceConfig( + aws_access_key_id="AKIAEXAMPLE", + aws_secret_access_key="fake-secret", + aws_session_token=token, + region=region, + ) + + +def response(payload: object, status: int = 200, headers: dict[str, str] | None = None) -> requests.Response: + result = requests.Response() + result.status_code = status + result._content = json.dumps(payload).encode() + result.headers.update(headers or {}) + return result + + +@pytest.fixture +def session() -> Iterator[MagicMock]: + with patch.object(transport, "make_tracked_session") as factory: + yield factory.return_value + + +def manager(state: AwsComputeOptimizerResumeConfig | None = None) -> MagicMock: + result = MagicMock() + result.load_state.return_value = state + return result + + +@pytest.mark.parametrize( + "region,token,host", + [ + (None, None, "compute-optimizer.us-east-1.amazonaws.com"), + ("eu-west-1", "fake-session-token", "compute-optimizer.eu-west-1.amazonaws.com"), + ("cn-north-1", "fake-session-token", "compute-optimizer.cn-north-1.amazonaws.com.cn"), + ("us-gov-west-1", None, "compute-optimizer.us-gov-west-1.amazonaws.com"), + ], +) +def test_signed_json_request(session: MagicMock, region: str | None, token: str | None, host: str) -> None: + session.post.return_value = response({"status": "Active"}) + model = Loader().load_service_model("compute-optimizer", "service-2", api_version="2019-11-01") + client = AwsComputeOptimizerClient(config(region, token), "2019-11-01") + assert client.request("GetEnrollmentStatus", {}) == {"status": "Active"} + call = session.post.call_args + headers = call.kwargs["headers"] + assert call.args == (f"https://{host}/",) + assert call.kwargs["data"] == b"{}" + assert headers["Content-Type"] == f"application/x-amz-json-{model['metadata']['jsonVersion']}" + assert headers["X-Amz-Target"] == f"{model['metadata']['targetPrefix']}.GetEnrollmentStatus" + assert f"/{region or 'us-east-1'}/compute-optimizer/aws4_request" in headers["Authorization"] + assert "x-amz-target" in headers["Authorization"] + assert headers.get("X-Amz-Security-Token") == token + assert "X-Amz-Date" in headers + assert call.kwargs["allow_redirects"] is False + assert call.kwargs["timeout"] == 60 + + +@pytest.mark.parametrize( + "endpoint,operation,key,item,primary_keys", + [ + ( + "ec2_instance_recommendations", + "GetEC2InstanceRecommendations", + "instanceRecommendations", + {"instanceArn": "arn:aws:ec2:eu-west-1:123456789012:instance/i-example"}, + ["instance_arn"], + ), + ( + "auto_scaling_group_recommendations", + "GetAutoScalingGroupRecommendations", + "autoScalingGroupRecommendations", + {"autoScalingGroupArn": "arn:aws:autoscaling:eu-west-1:123456789012:autoScalingGroup:example"}, + ["auto_scaling_group_arn"], + ), + ( + "lambda_function_recommendations", + "GetLambdaFunctionRecommendations", + "lambdaFunctionRecommendations", + {"functionArn": "arn:aws:lambda:eu-west-1:123456789012:function:example", "functionVersion": "$LATEST"}, + ["function_arn", "function_version"], + ), + ( + "ecs_service_recommendations", + "GetECSServiceRecommendations", + "ecsServiceRecommendations", + {"serviceArn": "arn:aws:ecs:eu-west-1:123456789012:service/example/service"}, + ["service_arn"], + ), + ( + "ebs_volume_recommendations", + "GetEBSVolumeRecommendations", + "volumeRecommendations", + {"volumeArn": "arn:aws:ec2:eu-west-1:123456789012:volume/vol-example"}, + ["volume_arn"], + ), + ( + "recommendation_summaries", + "GetRecommendationSummaries", + "recommendationSummaries", + {"accountId": "123456789012", "recommendationResourceType": "Ec2Instance"}, + ["account_id", "recommendation_resource_type", "region"], + ), + ], +) +def test_full_refresh_pagination( + session: MagicMock, endpoint: str, operation: str, key: str, item: dict[str, Any], primary_keys: list[str] +) -> None: + session.post.side_effect = [ + response({key: [item], "nextToken": "page-2"}), + response({key: [], "nextToken": "page-3"}), + response({key: [item], "nextToken": None}), + ] + resume = manager() + source = aws_compute_optimizer_source(config(), endpoint, resume) + rows = list(cast(Iterable[Any], source.items())) + assert len(rows) == 2 + assert all(rows[0][0][column] for column in primary_keys) + assert source.primary_keys == primary_keys + assert rows[0][0]["region"] == "eu-west-1" + assert source.sort_mode is None + assert source.partition_keys is None + calls = session.post.call_args_list + assert [json.loads(call.kwargs["data"]) for call in calls] == [ + {"maxResults": 100}, + {"maxResults": 100, "nextToken": "page-2"}, + {"maxResults": 100, "nextToken": "page-3"}, + ] + assert all(call.kwargs["headers"]["X-Amz-Target"] == f"ComputeOptimizerService.{operation}" for call in calls) + assert [call.args[0] for call in resume.save_state.call_args_list] == [ + AwsComputeOptimizerResumeConfig(next_token="page-2"), + AwsComputeOptimizerResumeConfig(next_token="page-3"), + AwsComputeOptimizerResumeConfig(complete=True), + ] + assert resume.safe_point.call_count == 3 + resume.clear_state.assert_not_called() + assert source.on_complete is not None + source.on_complete() + resume.clear_state.assert_called_once() + session.close.assert_called_once() + + +@pytest.mark.parametrize("terminal", [{}, {"nextToken": None}, {"nextToken": ""}]) +def test_resume_and_terminal_empty_page(session: MagicMock, terminal: dict[str, object]) -> None: + resume = manager(AwsComputeOptimizerResumeConfig(next_token="saved-token")) + session.post.return_value = response(terminal) + source = aws_compute_optimizer_source(config(), "ec2_instance_recommendations", resume) + assert list(cast(Iterable[Any], source.items())) == [] + assert json.loads(session.post.call_args.kwargs["data"]) == {"maxResults": 100, "nextToken": "saved-token"} + resume.save_state.assert_called_once_with(AwsComputeOptimizerResumeConfig(complete=True)) + resume.safe_point.assert_called_once() + + +def test_completed_resume_does_not_refetch(session: MagicMock) -> None: + source = aws_compute_optimizer_source( + config(), "ec2_instance_recommendations", manager(AwsComputeOptimizerResumeConfig(complete=True)) + ) + assert list(cast(Iterable[Any], source.items())) == [] + session.post.assert_not_called() + + +def test_resume_is_staged_before_batch_and_not_cleared_early(session: MagicMock) -> None: + session.post.return_value = response({"instanceRecommendations": [{"instanceArn": "example"}], "nextToken": "next"}) + resume = manager() + source = aws_compute_optimizer_source(config(), "ec2_instance_recommendations", resume) + iterator = cast(Generator[list[dict[str, Any]]], source.items()) + assert next(iterator) == [{"instance_arn": "example", "region": "eu-west-1"}] + resume.save_state.assert_called_once_with(AwsComputeOptimizerResumeConfig(next_token="next")) + resume.clear_state.assert_not_called() + iterator.close() + session.close.assert_called_once() + + +def test_repeated_token_fails_without_marking_complete(session: MagicMock) -> None: + session.post.return_value = response({"nextToken": "same"}) + resume = manager(AwsComputeOptimizerResumeConfig(next_token="same")) + source = aws_compute_optimizer_source(config(), "ec2_instance_recommendations", resume) + with pytest.raises(ValueError, match="repeated a pagination token"): + list(cast(Iterable[Any], source.items())) + resume.save_state.assert_not_called() + session.close.assert_called_once() + + +def test_normalization_preserves_options_and_converts_timestamp() -> None: + options = [{"instanceType": "m7i.large", "rank": 1}] + row = transport.normalize_row( + { + "lastRefreshTimestamp": 1735689600, + "recommendationOptions": options, + "savingsOpportunity": {"estimatedMonthlySavings": {"value": 12.5, "currency": "USD"}}, + }, + "eu-west-1", + ) + assert row == { + "last_refresh_timestamp": datetime(2025, 1, 1, tzinfo=UTC), + "recommendation_options": options, + "savings_opportunity_estimated_monthly_savings_value": 12.5, + "savings_opportunity_estimated_monthly_savings_currency": "USD", + "region": "eu-west-1", + } + + +@pytest.mark.parametrize( + "code,status,attempts", + [ + ("ThrottlingException", 400, 4), + ("AccessDeniedException", 400, 1), + ("InternalServerException", 500, 1), + ("ThrottlingException", 429, 1), + ], +) +def test_only_body_throttling_has_application_retries( + session: MagicMock, code: str, status: int, attempts: int +) -> None: + session.post.return_value = response({"__type": f"com.amazonaws.computeoptimizer#{code}"}, status) + with pytest.raises(AwsComputeOptimizerError) as raised: + AwsComputeOptimizerClient(config()).request("GetEnrollmentStatus", {}) + assert raised.value.code == code + assert session.post.call_count == attempts + + +@pytest.mark.parametrize( + "payload,status,headers,code", + [ + ({"__type": "namespace#OptInRequiredException"}, 400, {}, "OptInRequiredException"), + ({"code": "InvalidClientTokenId"}, 403, {}, "InvalidClientTokenId"), + ({}, 403, {"x-amzn-ErrorType": "AccessDeniedException:extra"}, "AccessDeniedException"), + ("unavailable", 503, {}, "HTTP 503"), + ([], 502, {}, "HTTP 502"), + ], +) +def test_error_parsing(payload: object, status: int, headers: dict[str, str], code: str) -> None: + assert transport.error_for_response(response(payload, status, headers)).code == code + + +def test_partial_error_does_not_commit_incomplete_data(session: MagicMock) -> None: + session.post.return_value = response( + { + "instanceRecommendations": [{"instanceArn": "example"}], + "errors": [{"code": "AccessDeniedException", "message": "denied"}], + } + ) + resume = manager() + source = aws_compute_optimizer_source(config(), "ec2_instance_recommendations", resume) + with pytest.raises(AwsComputeOptimizerError, match="AccessDeniedException"): + list(cast(Iterable[Any], source.items())) + resume.save_state.assert_not_called() + session.close.assert_called_once() + + +@pytest.mark.parametrize("schema,valid", [(None, True), ("ec2_instance_recommendations", False)]) +def test_missing_permission_at_creation_and_per_table(session: MagicMock, schema: str | None, valid: bool) -> None: + session.post.return_value = response({"__type": "AccessDeniedException"}, 400) + result, message = validate_credentials(config(), schema) + assert result is valid + if schema is None: + assert message is None + assert session.post.call_args.kwargs["headers"]["X-Amz-Target"].endswith(".GetEnrollmentStatus") + assert json.loads(session.post.call_args.kwargs["data"]) == {} + else: + assert message and "compute-optimizer:GetEC2InstanceRecommendations" in message + assert json.loads(session.post.call_args.kwargs["data"]) == {"maxResults": 1} + session.close.assert_called_once() + + +@pytest.mark.parametrize("status,valid", [("Active", True), ("Inactive", False), ("Pending", False), ("Failed", False)]) +def test_enrollment_validation(session: MagicMock, status: str, valid: bool) -> None: + session.post.return_value = response({"status": status}) + result, message = validate_credentials(config()) + assert result is valid + assert (message is None) is valid + session.post.assert_called_once() + + +@pytest.mark.parametrize( + "code,expected", + [ + ("UnrecognizedClientException", "credentials"), + ("InvalidSignatureException", "signature"), + ("SignatureDoesNotMatch", "signature"), + ("ExpiredTokenException", "expired"), + ("OptInRequiredException", "Enable AWS Compute Optimizer"), + ("SubscriptionRequiredException", "Enable AWS Compute Optimizer"), + ], +) +def test_credentials_and_sync_errors_are_actionable(session: MagicMock, code: str, expected: str) -> None: + session.post.return_value = response({"__type": code}, 400) + valid, message = validate_credentials(config()) + assert not valid + assert message and expected in message + mappings = AwsComputeOptimizerSource().get_non_retryable_errors() + assert mappings[str(AwsComputeOptimizerError(code))] == message + + +@pytest.mark.parametrize("error", [requests.Timeout(), requests.ConnectionError(), ValueError("invalid JSON")]) +def test_validation_network_or_response_failure(session: MagicMock, error: Exception) -> None: + session.post.side_effect = error + assert validate_credentials(config()) == (False, "Could not reach the AWS Compute Optimizer API. Try again.") + session.close.assert_called_once() + + +@pytest.mark.parametrize( + "region", ["https://example.com", "us-east-1.example.com", "us-east-1/", "us-east-1@evil", " us-east-1"] +) +def test_rejects_invalid_region_before_request(session: MagicMock, region: str) -> None: + assert validate_credentials(config(region)) == (False, "Enter an AWS region, such as us-east-1.") + session.post.assert_not_called() + + +def test_unknown_schema_and_version_rejected(session: MagicMock) -> None: + assert validate_credentials(config(), "unknown")[0] is False + assert validate_credentials(config(), api_version="2000-01-01")[0] is False + with pytest.raises(ValueError, match="Unknown AWS Compute Optimizer table"): + aws_compute_optimizer_source(config(), "unknown", manager()) + session.post.assert_not_called() diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/awscomputeoptimizer.py b/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/awscomputeoptimizer.py index 7f774963446f..c5643e241a64 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/awscomputeoptimizer.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/awscomputeoptimizer.py @@ -6,4 +6,7 @@ @config.config class AwsComputeOptimizerSourceConfig(config.Config): - pass + aws_access_key_id: str + aws_secret_access_key: str + aws_session_token: str | None = None + region: str | None = None From f2a1c3494eb307aa162d34ff9d6e6bc2fbe833ce Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 17:22:59 +0200 Subject: [PATCH 2/4] feat(data-warehouse): implement the aws compute optimizer import source Re-run CI after the duplication lint job was cancelled without logs. Includes the conflict resolution that preserves both source catalog updates. From 939b93d1050cde054a7f272e9eb581198c5643d5 Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 17:41:32 +0200 Subject: [PATCH 3/4] feat(data-warehouse): implement the aws compute optimizer import source Retry CI after an unrelated ingestion warnings test failed with an inconsistent result count. The AWS Compute Optimizer source tree is unchanged. From 1ffc43af68c382cfd0e0cce3412b9f3830bfb1b4 Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 17:57:54 +0200 Subject: [PATCH 4/4] fix(tests): keep ingestion warning fixtures within TTL Move the frozen ingestion-warning fixture dates far enough into the future that ClickHouse's wall-clock TTL cannot remove the wider-window row during the test run. --- .../api/test/test_ingestion_warnings_v2.py | 22 +++++++++---------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/posthog/api/test/test_ingestion_warnings_v2.py b/posthog/api/test/test_ingestion_warnings_v2.py index 93783cb2a6c7..7f0ee50995c7 100644 --- a/posthog/api/test/test_ingestion_warnings_v2.py +++ b/posthog/api/test/test_ingestion_warnings_v2.py @@ -29,11 +29,11 @@ def create_warning(team_id: int, type: str, timestamp: str, details: dict, sourc ) -@time_machine.travel("2026-07-07T12:00:00.000Z", tick=False) +@time_machine.travel("2099-07-07T12:00:00.000Z", tick=False) class TestIngestionWarningsV2API(ClickhouseTestMixin, APIBaseTest): def setUp(self): super().setUp() - for hour_timestamp in ["2026-07-07 09:00:00", "2026-07-07 09:30:00", "2026-07-07 10:00:00"]: + for hour_timestamp in ["2099-07-07 09:00:00", "2099-07-07 09:30:00", "2099-07-07 10:00:00"]: create_warning( team_id=self.team.id, type="message_size_too_large", @@ -49,7 +49,7 @@ def setUp(self): create_warning( team_id=self.team.id, type="cannot_merge_already_identified", - timestamp="2026-07-07 11:00:00", + timestamp="2099-07-07 11:00:00", details={ "category": "merge", "severity": "warning", @@ -60,7 +60,7 @@ def setUp(self): create_warning( team_id=self.team.id, type="message_size_too_large", - timestamp="2026-07-04 10:00:00", + timestamp="2099-07-04 10:00:00", details={"category": "size", "severity": "error", "distinctId": "old-user"}, ) # Another team's warning must never leak @@ -68,7 +68,7 @@ def setUp(self): create_warning( team_id=other_team.id, type="message_size_too_large", - timestamp="2026-07-07 10:00:00", + timestamp="2099-07-07 10:00:00", details={"category": "size", "severity": "error", "distinctId": "other-team-user"}, ) @@ -88,7 +88,7 @@ def test_groups_warnings_by_type_with_counts_samples_and_sparkline(self): size_warnings = results[0] assert size_warnings["category"] == "size" assert size_warnings["severity"] == "error" - assert size_warnings["last_seen"].startswith("2026-07-07T10:00:00") + assert size_warnings["last_seen"].startswith("2099-07-07T10:00:00") sample_timestamps = [sample["timestamp"] for sample in size_warnings["samples"]] assert sample_timestamps == sorted(sample_timestamps, reverse=True) @@ -117,7 +117,7 @@ def test_samples_report_producer_and_pipeline_step( create_warning( team_id=self.team.id, type=warning_type, - timestamp="2026-07-07 11:30:00", + timestamp="2099-07-07 11:30:00", details={"category": "event", "severity": "error", **extra_details}, source="capture", ) @@ -149,7 +149,7 @@ def test_time_range_bounds_results(self): assert results[0]["count"] == 4 # Explicit ISO bounds narrow down to a single warning - _, results = self._list(since="2026-07-07T10:30:00Z", until="2026-07-07T11:30:00Z") + _, results = self._list(since="2099-07-07T10:30:00Z", until="2099-07-07T11:30:00Z") assert [(r["type"], r["count"]) for r in results] == [("cannot_merge_already_identified", 1)] @parameterized.expand( @@ -169,14 +169,14 @@ def test_samples_param_caps_returned_samples(self): create_warning( team_id=self.team.id, type="merge_race_condition", - timestamp=f"2026-07-07 11:{minute:02d}:00", + timestamp=f"2099-07-07 11:{minute:02d}:00", details={"category": "merge", "severity": "error"}, ) _, results = self._list(type="merge_race_condition") assert results[0]["count"] == 7 assert len(results[0]["samples"]) == 5 - assert results[0]["samples"][0]["timestamp"].startswith("2026-07-07T11:06:00") + assert results[0]["samples"][0]["timestamp"].startswith("2099-07-07T11:06:00") _, results = self._list(type="merge_race_condition", samples=2) assert len(results[0]["samples"]) == 2 @@ -184,7 +184,7 @@ def test_samples_param_caps_returned_samples(self): @parameterized.expand( [ ({"limit": 0},), - ({"since": "2026-07-07T11:00:00Z", "until": "2026-07-07T10:00:00Z"},), + ({"since": "2099-07-07T11:00:00Z", "until": "2099-07-07T10:00:00Z"},), ({"severity": "critical"},), ] )