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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,7 @@ the row lists both.
| aws_budgets | HTTP | requests | ✅ |
| aws_cloudtrail | HTTP | requests | ✅ |
| aws_compute_optimizer | HTTP | requests | ✅ |
| aws_config | HTTP | requests | ✅ |
| aws_cost_anomaly_detection | HTTP | requests | ✅ |
| aws_cost_explorer | HTTP | requests | ✅ |
| aws_glue_data_catalog | HTTP | requests | ✅ |
Expand Down Expand Up @@ -922,7 +923,6 @@ doesn't conflict with concurrent PRs.
- automox
- aws_athena
- aws_cloudformation
- aws_config
- aws_connect
- aws_cost_and_usage_report
- aws_guardduty
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,239 @@
import re
import json
import datetime as dt
from collections.abc import Generator
from typing import TYPE_CHECKING, 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_config.settings import (
AWS_CONFIG_ENDPOINTS,
ERROR_MESSAGES,
TARGET_PREFIXES,
AwsConfigEndpoint,
)
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.typings import SourceResponse

if TYPE_CHECKING:
from products.warehouse_sources.backend.temporal.data_imports.sources.common.resumable import ResumableSourceManager
from products.warehouse_sources.backend.temporal.data_imports.sources.generated_configs.awsconfig import (
AwsConfigSourceConfig,
)

_REGION = re.compile(r"[a-z]{2}(?:-[a-z]+)+-\d+")
_CAMEL_BOUNDARY = re.compile(r"(?<=[a-z0-9])(?=[A-Z])|(?<=[A-Z])(?=[A-Z][a-z])")
_THROTTLE_CODES = {"Throttling", "ThrottlingException", "TooManyRequestsException", "RequestLimitExceeded"}
_TIMESTAMP_COLUMNS = {"configuration_item_capture_time", "resource_creation_time", "last_update_requested_time"}

# Config reads use POST, so the transport must retry this method for transient HTTP failures.
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 AwsConfigResumeConfig:
next_token: str | None = None
complete: bool = False


class AwsConfigError(Exception):
def __init__(self, code: str, message: str) -> None:
super().__init__(f"AWS Config request failed: {code} - {message}")
self.code = code


class AwsConfigThrottledError(AwsConfigError):
pass


def error_for_response(response: requests.Response) -> AwsConfigError:
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]
message = str(body.get("message") or body.get("Message") or response.text)[:500]
# HTTP retries belong to the tracked transport. Only body-level throttles need another retry policy.
error_class = AwsConfigThrottledError if response.status_code == 400 and code in _THROTTLE_CODES else AwsConfigError
return error_class(code, message)


class AwsConfigClient:
def __init__(self, config: "AwsConfigSourceConfig", api_version: str) -> None:
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 "us-east-1"
if not _REGION.fullmatch(self.region):
raise ValueError("Enter an AWS region such as us-east-1.")
if api_version not in TARGET_PREFIXES:
raise ValueError(f"Unsupported AWS Config API version: {api_version}")
self.target_prefix = TARGET_PREFIXES[api_version]
suffix = "amazonaws.com.cn" if self.region.startswith("cn-") else "amazonaws.com"
self.url = f"https://config.{self.region}.{suffix}/"
self._signer = SigV4Auth(
Credentials(config.aws_access_key_id, config.aws_secret_access_key, config.aws_session_token or None),
"config",
self.region,
)
self._session = make_tracked_session(
retry=TRANSPORT_RETRY,
redact_values=tuple(value for value in (config.aws_secret_access_key, config.aws_session_token) if value),
)

def close(self) -> None:
self._session.close()

@retry(
retry=retry_if_exception_type(AwsConfigThrottledError),
stop=stop_after_attempt(5),
wait=wait_exponential_jitter(initial=1, max=30),
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.1",
"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 response.status_code != 200:
raise error_for_response(response)
parsed = response.json()
if not isinstance(parsed, dict):
raise ValueError("AWS Config returned an invalid response. Try the sync again.")
return parsed


def request_payload(endpoint: AwsConfigEndpoint, *, probe: bool = False) -> dict[str, Any]:
payload: dict[str, Any] = {}
if endpoint.page_size is not None:
payload["Limit"] = 1 if probe else endpoint.page_size
if endpoint.expression is not None:
payload["Expression"] = endpoint.expression
return payload


def normalize_row(item: dict[str, Any], region: str) -> dict[str, Any]:
row = {_CAMEL_BOUNDARY.sub("_", key).lower(): value for key, value in item.items()}
row["region"] = region
for key in _TIMESTAMP_COLUMNS & row.keys():
value = row[key]
if isinstance(value, int | float) and not isinstance(value, bool):
row[key] = dt.datetime.fromtimestamp(value, tz=dt.UTC)
elif isinstance(value, str) and value:
row[key] = dt.datetime.fromisoformat(value.replace("Z", "+00:00"))
return row


def get_rows(

Check warning on line 152 in products/warehouse_sources/backend/temporal/data_imports/sources/aws_config/aws_config.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`get_rows` has cyclomatic complexity 12 (warn >10)

Check warning on line 152 in products/warehouse_sources/backend/temporal/data_imports/sources/aws_config/aws_config.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`get_rows` has cyclomatic complexity 12 (warn >10)
config: "AwsConfigSourceConfig",
endpoint: AwsConfigEndpoint,
api_version: str,
manager: "ResumableSourceManager[AwsConfigResumeConfig]",
) -> Generator[list[dict[str, Any]]]:
resume = manager.load_state() if manager.can_resume() else None
if resume is not None and resume.complete:
manager.clear_state()
return
next_token = resume.next_token if resume else None
restarting = next_token is not None
client = AwsConfigClient(config, api_version)
try:
while True:
payload = request_payload(endpoint)
if next_token:
payload["NextToken"] = next_token
try:
body = client.request(endpoint.operation, payload)
except AwsConfigError as error:
if restarting and error.code == "InvalidNextTokenException":
# A fresh retry must restart the full-refresh destination as well as extraction.
manager.clear_state()
raise
restarting = False
rows = []
for item in body.get(endpoint.result_key) or []:
if endpoint.expression is not None:
item = json.loads(item)
if not isinstance(item, dict):
raise ValueError("AWS Config returned an invalid resource. Try the sync again.")
rows.append(normalize_row(item, client.region))
token = body.get("NextToken") or None
if token is not None and (not isinstance(token, str) or token == next_token):
raise ValueError("AWS Config returned an invalid page token. Try the sync again.")
next_token = token
manager.save_state(AwsConfigResumeConfig(next_token=next_token, complete=next_token is None))
if rows:
yield rows
manager.safe_point()
if next_token is None:
break
manager.clear_state()
finally:
client.close()


def validate_credentials(
config: "AwsConfigSourceConfig", api_version: str, schema_name: str | None = None
) -> tuple[bool, str | None]:
endpoint = AWS_CONFIG_ENDPOINTS.get(schema_name or "config_rules")
if endpoint is None:
return False, f"Unknown AWS Config table: {schema_name}"
try:
client = AwsConfigClient(config, api_version)
except ValueError as error:
return False, str(error)
try:
client.request(endpoint.operation, request_payload(endpoint, probe=True))
except AwsConfigError as error:
if error.code in {"AccessDenied", "AccessDeniedException"}:
if schema_name is None:
return True, None
return False, f"Grant config:{endpoint.operation} to this IAM user or role to sync this table."
return False, ERROR_MESSAGES.get(error.code, "Could not read AWS Config. Check the region and try again.")
except (requests.RequestException, ValueError):
return False, "Could not reach the AWS Config API. Check the region and try again."
finally:
client.close()
return True, None


def aws_config_source(
config: "AwsConfigSourceConfig",
endpoint: str,
api_version: str,
manager: "ResumableSourceManager[AwsConfigResumeConfig]",
) -> SourceResponse:
endpoint_config = AWS_CONFIG_ENDPOINTS.get(endpoint)
if endpoint_config is None:
raise ValueError(f"Unknown AWS Config table: {endpoint}")
return SourceResponse(
name=endpoint,
items=lambda: get_rows(config, endpoint_config, api_version, manager),
primary_keys=list(endpoint_config.primary_key),
sort_mode=None,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
from products.warehouse_sources.backend.temporal.data_imports.sources.common.canonical_descriptions import (
CanonicalDescriptions,
)

CANONICAL_DESCRIPTIONS: CanonicalDescriptions = {
"resources": {
"description": "Current configurations of recorded AWS resources in the selected region. Excludes deleted resources.",
"docs_url": "https://docs.aws.amazon.com/config/latest/APIReference/API_SelectResourceConfig.html",
"columns": {
"account_id": "AWS account that owns the resource.",
"aws_region": "AWS region reported for the resource.",
"region": "AWS region selected for this import.",
"arn": "Amazon Resource Name of the resource.",
"resource_id": "Resource identifier assigned by its AWS service.",
"resource_type": "AWS resource type, such as AWS::EC2::Instance.",
"resource_name": "Resource name, when the service provides one.",
"configuration": "Configuration properties recorded for the resource.",
"configuration_item_capture_time": "Time when AWS Config recorded this configuration.",
"resource_creation_time": "Time when the resource was created, when available.",
"tags": "Tags attached to the resource.",
},
},
"config_rules": {
"description": "AWS Config rules and their evaluation settings in the selected region.",
"docs_url": "https://docs.aws.amazon.com/config/latest/APIReference/API_DescribeConfigRules.html",
"columns": {
"config_rule_arn": "Amazon Resource Name of the rule.",
"config_rule_id": "Identifier assigned to the rule by AWS Config.",
"config_rule_name": "Name of the AWS Config rule.",
"config_rule_state": "Current state of the rule.",
"description": "Description of the rule.",
"scope": "Resource types, identifiers, or tags that restrict the rule's evaluations.",
"source": "Rule owner, identifier, and evaluation triggers.",
"input_parameters": "Parameters passed to the rule as a JSON string.",
"maximum_execution_frequency": "Maximum frequency for periodic evaluations.",
"region": "AWS region selected for this import.",
},
},
"rule_compliance": {
"description": "Current compliance status for AWS Config rules, including counts of noncompliant resources.",
"docs_url": "https://docs.aws.amazon.com/config/latest/APIReference/API_DescribeComplianceByConfigRule.html",
"columns": {
"config_rule_name": "Name of the evaluated AWS Config rule.",
"compliance": "Compliance status and a capped count of contributing resources, with an indicator when the cap is exceeded.",
"region": "AWS region selected for this import.",
},
},
"conformance_packs": {
"description": "Conformance packs and their deployment settings in the selected region.",
"docs_url": "https://docs.aws.amazon.com/config/latest/APIReference/API_DescribeConformancePacks.html",
"columns": {
"conformance_pack_arn": "Amazon Resource Name of the conformance pack.",
"conformance_pack_id": "Identifier assigned to the conformance pack.",
"conformance_pack_name": "Name of the conformance pack.",
"conformance_pack_input_parameters": "Parameters used to deploy the conformance pack.",
"delivery_s3_bucket": "S3 bucket used for the conformance pack template.",
"delivery_s3_key_prefix": "S3 prefix used for the conformance pack template.",
"last_update_requested_time": "Time when the most recent update was requested.",
"created_by": "AWS service that created the conformance pack.",
"region": "AWS region selected for this import.",
},
},
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
from posthog.dataclasses import frozen

CONFIG_API_VERSION = "2014-11-12"
TARGET_PREFIXES = {CONFIG_API_VERSION: "StarlingDoveService"}

RESOURCE_EXPRESSION = (
"SELECT accountId, awsRegion, arn, resourceId, resourceType, resourceName, "
"configuration, configurationItemCaptureTime, resourceCreationTime, tags"
)


@frozen
class AwsConfigEndpoint:
operation: str
result_key: str
primary_key: tuple[str, ...]
description: str
page_size: int | None = None
expression: str | None = None


AWS_CONFIG_ENDPOINTS: dict[str, AwsConfigEndpoint] = {
"resources": AwsConfigEndpoint(
operation="SelectResourceConfig",
result_key="Results",
primary_key=("account_id", "aws_region", "resource_type", "resource_id"),
description="Current configurations of recorded resources in the selected region. Excludes deleted resources.",
page_size=100,
expression=RESOURCE_EXPRESSION,
),
"config_rules": AwsConfigEndpoint(
operation="DescribeConfigRules",
result_key="ConfigRules",
primary_key=("config_rule_arn",),
description="AWS Config rules, their evaluation settings, and resource scopes.",
),
"rule_compliance": AwsConfigEndpoint(
operation="DescribeComplianceByConfigRule",
result_key="ComplianceByConfigRules",
primary_key=("region", "config_rule_name"),
description="Current compliance status and counts of noncompliant resources for each rule.",
),
"conformance_packs": AwsConfigEndpoint(
operation="DescribeConformancePacks",
result_key="ConformancePackDetails",
primary_key=("conformance_pack_arn",),
description="Conformance packs and their deployment settings in the selected region.",
page_size=20,
),
}

ENDPOINTS = tuple(AWS_CONFIG_ENDPOINTS)
ENDPOINT_DESCRIPTIONS = {name: endpoint.description for name, endpoint in AWS_CONFIG_ENDPOINTS.items()}

ERROR_MESSAGES = {
"AccessDenied": "AWS denied access. Grant the config read permission for the selected table to this IAM user or role.",
"AccessDeniedException": "AWS denied access. Grant the config read permission for the selected table to this IAM user or role.",
"UnrecognizedClientException": "AWS rejected the credentials. Check the access key ID, secret access key, and session token.",
"InvalidClientTokenId": "AWS rejected the access key ID. 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 and session token.",
"ExpiredTokenException": "The AWS session token has expired. Enter new temporary credentials.",
"ExpiredToken": "The AWS session token has expired. Enter new temporary credentials.",
"MissingAuthenticationToken": "AWS did not receive valid credentials. Enter the access key ID and secret access key.",
"OptInRequired": "Enable AWS Config in the selected account and region, then try again.",
"SubscriptionRequiredException": "Enable AWS Config in the selected account and region, then try again.",
"NoAvailableConfigurationRecorderException": "Set up an AWS Config recorder in the selected region, then try again.",
}
Loading
Loading