diff --git a/sumologic-app-utils/src/awsresource.py b/sumologic-app-utils/src/awsresource.py index 8d7cb5ce..cf956a70 100644 --- a/sumologic-app-utils/src/awsresource.py +++ b/sumologic-app-utils/src/awsresource.py @@ -513,6 +513,34 @@ def _batch_size_chunk(self, iterable, size=1): data = iterable[idx:min(idx + size, length)] yield data + @staticmethod + def _apply_bucket_policy(bucket_name, statements): + s3 = boto3.client('s3') + for attempt in range(5): + try: + response = s3.get_bucket_policy(Bucket=bucket_name) + existing_policy = json.loads(response["Policy"]) + except ClientError as e: + if e.response['Error']['Code'] == "NoSuchBucketPolicy": + existing_policy = {"Version": "2012-10-17", "Statement": []} + else: + raise + existing_sids = {s.get("Sid") for s in existing_policy["Statement"] if s.get("Sid")} + for stmt in statements: + if stmt.get("Sid") not in existing_sids: + existing_policy["Statement"].append(stmt) + try: + s3.put_bucket_policy(Bucket=bucket_name, Policy=json.dumps(existing_policy)) + print(f"put_bucket_policy succeeded for {bucket_name} on attempt {attempt + 1}") + return + except ClientError as e: + if e.response['Error']['Code'] == "OperationAborted": + print(f"OperationAborted on put_bucket_policy attempt {attempt + 1} for {bucket_name}, retrying in {2 ** attempt}s...") + time.sleep(2 ** attempt) + continue + raise + raise Exception(f"Failed to put bucket policy for {bucket_name} after 5 attempts") + class EC2Resources(AWSResourcesAbstract): @@ -847,82 +875,48 @@ def tag_resources_cloud_trail_event(self, arns, tags): tags.extend(tags_arn) self.client.add_tags_to_resource(ResourceName=arn, Tags=tags) -class LbResources(AWSResourcesAbstract): +class AlbResources(AWSResourcesAbstract): def add_bucket_policy(self, bucket_name): print("Adding policy to the bucket " + bucket_name) - s3 = boto3.client('s3') - try: - response = s3.get_bucket_policy(Bucket=bucket_name) - existing_policy = json.loads(response["Policy"]) - except ClientError as e: - if "Error" in e.response and "Code" in e.response["Error"] \ - and e.response['Error']['Code'] == "NoSuchBucketPolicy": - existing_policy = { - "Version": "2012-10-17", - "Statement": [ - ] - } - else: - raise e - - bucket_policy = [ + self._apply_bucket_policy(bucket_name, [ { "Sid": "AWSCloudTrailAclCheck", "Effect": "Allow", - "Principal": { - "Service": "cloudtrail.amazonaws.com" - }, + "Principal": {"Service": "cloudtrail.amazonaws.com"}, "Action": "s3:GetBucketAcl", "Resource": f"arn:{self.partition}:s3:::{bucket_name}" }, { "Sid": "AWSCloudTrailWrite", "Effect": "Allow", - "Principal": { - "Service": "cloudtrail.amazonaws.com" - }, + "Principal": {"Service": "cloudtrail.amazonaws.com"}, "Action": "s3:PutObject", "Resource": f"arn:{self.partition}:s3:::{bucket_name}/*", - "Condition": { - "StringEquals": { - "s3:x-amz-acl": "bucket-owner-full-control" - } - } + "Condition": {"StringEquals": {"s3:x-amz-acl": "bucket-owner-full-control"}} }, { "Sid": "AWSBucketExistenceCheck", "Effect": "Allow", - "Principal": { - "Service": "cloudtrail.amazonaws.com" - }, + "Principal": {"Service": "cloudtrail.amazonaws.com"}, "Action": "s3:ListBucket", "Resource": f"arn:{self.partition}:s3:::{bucket_name}" }, { - "Sid": "AWSAlbLogDeliveryAclCheck", + "Sid": "AWSALBLogDeliveryAclCheck", "Effect": "Allow", - "Principal": { - "Service": "delivery.logs.amazonaws.com" - }, + "Principal": {"Service": "delivery.logs.amazonaws.com"}, "Action": "s3:GetBucketAcl", "Resource": f"arn:{self.partition}:s3:::{bucket_name}" }, { - "Sid": "AddLBLogsStatement", + "Sid": "AddALBLogsStatement", "Effect": "Allow", - "Principal": { - "Service": "logdelivery.elasticloadbalancing.amazonaws.com" - }, + "Principal": {"Service": "logdelivery.elasticloadbalancing.amazonaws.com"}, "Action": "s3:PutObject", "Resource": f"arn:{self.partition}:s3:::{bucket_name}/*" } - ] - existing_policy["Statement"].extend(bucket_policy) - - s3.put_bucket_policy(Bucket=bucket_name, Policy=json.dumps(existing_policy)) - -class AlbResources(LbResources): + ]) def fetch_resources(self): resources = [] @@ -1148,47 +1142,23 @@ def enable_s3_logs(self, arns, s3_bucket, s3_prefix): def add_bucket_policy(self, bucket_name, prefix): print("Adding policy to the bucket " + bucket_name) - s3 = boto3.client('s3') - try: - response = s3.get_bucket_policy(Bucket=bucket_name) - existing_policy = json.loads(response["Policy"]) - except ClientError as e: - if "Error" in e.response and "Code" in e.response["Error"] \ - and e.response['Error']['Code'] == "NoSuchBucketPolicy": - existing_policy = { - "Version": "2012-10-17", - "Statement": [ - ] - } - else: - raise e - - bucket_policy = [{ - "Sid": "AWSLogDeliveryAclCheck", - "Effect": "Allow", - "Principal": { - "Service": "delivery.logs.amazonaws.com" + self._apply_bucket_policy(bucket_name, [ + { + "Sid": "AWSLogDeliveryAclCheck", + "Effect": "Allow", + "Principal": {"Service": "delivery.logs.amazonaws.com"}, + "Action": "s3:GetBucketAcl", + "Resource": f"arn:{self.partition}:s3:::{bucket_name}" }, - "Action": "s3:GetBucketAcl", - "Resource": f"arn:{self.partition}:s3:::{bucket_name}" - }, { "Sid": "AWSLogDeliveryWrite", "Effect": "Allow", - "Principal": { - "Service": "delivery.logs.amazonaws.com" - }, + "Principal": {"Service": "delivery.logs.amazonaws.com"}, "Action": "s3:PutObject", "Resource": f"arn:{self.partition}:s3:::{bucket_name}/{prefix}/AWSLogs/{self.account_id}/*", - "Condition": { - "StringEquals": { - "s3:x-amz-acl": "bucket-owner-full-control" - } - } - }] - existing_policy["Statement"].extend(bucket_policy) - - s3.put_bucket_policy(Bucket=bucket_name, Policy=json.dumps(existing_policy)) + "Condition": {"StringEquals": {"s3:x-amz-acl": "bucket-owner-full-control"}} + } + ]) def disable_s3_logs(self, arns, s3_bucket): if arns: @@ -1204,7 +1174,7 @@ def disable_s3_logs(self, arns, s3_bucket): self.client.delete_flow_logs(FlowLogIds=flow_ids) -class ElbResource(LbResources): +class ElbResource(AWSResourcesAbstract): def fetch_resources(self): resources = [] next_token = None @@ -1274,6 +1244,47 @@ def enable_s3_logs(self, names, s3_bucket, s3_prefix): else: raise e + def add_bucket_policy(self, bucket_name): + print("Adding policy to the bucket " + bucket_name) + self._apply_bucket_policy(bucket_name, [ + { + "Sid": "AWSCloudTrailAclCheck", + "Effect": "Allow", + "Principal": {"Service": "cloudtrail.amazonaws.com"}, + "Action": "s3:GetBucketAcl", + "Resource": f"arn:{self.partition}:s3:::{bucket_name}" + }, + { + "Sid": "AWSCloudTrailWrite", + "Effect": "Allow", + "Principal": {"Service": "cloudtrail.amazonaws.com"}, + "Action": "s3:PutObject", + "Resource": f"arn:{self.partition}:s3:::{bucket_name}/*", + "Condition": {"StringEquals": {"s3:x-amz-acl": "bucket-owner-full-control"}} + }, + { + "Sid": "AWSBucketExistenceCheck", + "Effect": "Allow", + "Principal": {"Service": "cloudtrail.amazonaws.com"}, + "Action": "s3:ListBucket", + "Resource": f"arn:{self.partition}:s3:::{bucket_name}" + }, + { + "Sid": "AWSELBLogDeliveryAclCheck", + "Effect": "Allow", + "Principal": {"Service": "delivery.logs.amazonaws.com"}, + "Action": "s3:GetBucketAcl", + "Resource": f"arn:{self.partition}:s3:::{bucket_name}" + }, + { + "Sid": "AddELBLogsStatement", + "Effect": "Allow", + "Principal": {"Service": "logdelivery.elasticloadbalancing.amazonaws.com"}, + "Action": "s3:PutObject", + "Resource": f"arn:{self.partition}:s3:::{bucket_name}/*" + } + ]) + def disable_s3_logs(self, names, s3_bucket): attributes = [{'Key': 'access_logs.s3.enabled', 'Value': 'false'}] @@ -1319,6 +1330,147 @@ def get_provider(cls, provider_name, region_value, account_id, *args, **kwargs): raise Exception(f"{provider_name} provider not found") +class AddBucketPolicy(AWSResource): + """Custom::AddBucketPolicy — appends S3 log-delivery policy statements to an existing bucket + without replacing any existing statements. Idempotent by Sid. + ServiceType controls which statements are added: 'alb', 'elb', or 'cloudtrail'.""" + + ALL_STATEMENTS = [ + {"Sid": "AWSCloudTrailAclCheck", "Effect": "Allow", + "Principal": {"Service": "cloudtrail.amazonaws.com"}, + "Action": ["s3:GetBucketAcl"], + "Resource": "arn:{partition}:s3:::{bucket}"}, + {"Sid": "AWSCloudTrailWrite", "Effect": "Allow", + "Principal": {"Service": "cloudtrail.amazonaws.com"}, + "Action": ["s3:PutObject"], + "Resource": "arn:{partition}:s3:::{bucket}/*", + "Condition": {"StringEquals": {"s3:x-amz-acl": "bucket-owner-full-control"}}}, + {"Sid": "AWSBucketExistenceCheck", "Effect": "Allow", + "Principal": {"Service": "cloudtrail.amazonaws.com"}, + "Action": ["s3:ListBucket"], + "Resource": "arn:{partition}:s3:::{bucket}"}, + {"Sid": "AWSALBLogDeliveryAclCheck", "Effect": "Allow", + "Principal": {"Service": "delivery.logs.amazonaws.com"}, + "Action": ["s3:GetBucketAcl"], + "Resource": "arn:{partition}:s3:::{bucket}"}, + {"Sid": "AddALBLogsStatement", "Effect": "Allow", + "Principal": {"Service": "logdelivery.elasticloadbalancing.amazonaws.com"}, + "Action": ["s3:PutObject"], + "Resource": "arn:{partition}:s3:::{bucket}/*"}, + {"Sid": "AWSELBLogDeliveryAclCheck", "Effect": "Allow", + "Principal": {"Service": "delivery.logs.amazonaws.com"}, + "Action": ["s3:GetBucketAcl"], + "Resource": "arn:{partition}:s3:::{bucket}"}, + {"Sid": "AddELBLogsStatement", "Effect": "Allow", + "Principal": {"Service": "logdelivery.elasticloadbalancing.amazonaws.com"}, + "Action": ["s3:PutObject"], + "Resource": "arn:{partition}:s3:::{bucket}/*"}, + ] + + SERVICE_SIDS = { + "cloudtrail": {"AWSCloudTrailAclCheck", "AWSCloudTrailWrite", "AWSBucketExistenceCheck"}, + "alb": {"AWSALBLogDeliveryAclCheck", "AddALBLogsStatement"}, + "elb": {"AWSELBLogDeliveryAclCheck", "AddELBLogsStatement"}, + } + + def __init__(self, props, *args, **kwargs): + self.props = props + + def _statements_for_service(self, service_type): + allowed = self.SERVICE_SIDS.get(service_type) + if not allowed: + return self.ALL_STATEMENTS + return [s for s in self.ALL_STATEMENTS if s["Sid"] in allowed] + + def _build_statements(self, bucket_name, partition, service_type): + statements = [] + for tmpl in self._statements_for_service(service_type): + stmt = json.loads(json.dumps(tmpl)) + stmt["Resource"] = stmt["Resource"].format(bucket=bucket_name, partition=partition) + statements.append(stmt) + return statements + + def _add_policy(self, bucket_name, partition, service_type): + s3 = boto3.client('s3') + expected_stmts = self._build_statements(bucket_name, partition, service_type) + expected_sids = {s["Sid"] for s in expected_stmts} + for attempt in range(5): + try: + response = s3.get_bucket_policy(Bucket=bucket_name) + existing_policy = json.loads(response["Policy"]) + except ClientError as e: + if e.response['Error']['Code'] == "NoSuchBucketPolicy": + existing_policy = {"Version": "2012-10-17", "Statement": []} + else: + raise + existing_sids = {s.get("Sid") for s in existing_policy["Statement"] if s.get("Sid")} + added = [] + for stmt in expected_stmts: + if stmt["Sid"] not in existing_sids: + existing_policy["Statement"].append(stmt) + added.append(stmt["Sid"]) + if added: + print(f"put_bucket_policy attempt {attempt + 1} for {bucket_name}: adding SIDs {added}") + try: + s3.put_bucket_policy(Bucket=bucket_name, Policy=json.dumps(existing_policy)) + except ClientError as e: + if e.response['Error']['Code'] == "OperationAborted": + print(f"OperationAborted on put_bucket_policy attempt {attempt + 1} for {bucket_name}, retrying in {2 ** attempt}s...") + time.sleep(2 ** attempt) + continue + raise + print(f"put_bucket_policy succeeded for {bucket_name} on attempt {attempt + 1}") + # Verify our SIDs survived — a concurrent write could have overwritten them + time.sleep(0.5 * (attempt + 1)) + verify = json.loads(s3.get_bucket_policy(Bucket=bucket_name)["Policy"]) + current_sids = {s.get("Sid") for s in verify["Statement"]} + if not expected_sids.issubset(current_sids): + print(f"Concurrent policy overwrite detected on attempt {attempt + 1} for {bucket_name}, retrying...") + continue + else: + print(f"put_bucket_policy skipped for {bucket_name} on attempt {attempt + 1}: all SIDs already present") + return added + raise Exception(f"Failed to persist bucket policy for {bucket_name} after 5 attempts — concurrent overwrite") + + def _remove_policy(self, bucket_name, service_type): + s3 = boto3.client('s3') + our_sids = {s["Sid"] for s in self._statements_for_service(service_type)} + try: + response = s3.get_bucket_policy(Bucket=bucket_name) + existing_policy = json.loads(response["Policy"]) + except ClientError as e: + if e.response['Error']['Code'] in ("NoSuchBucketPolicy", "NoSuchBucket"): + return + raise + existing_policy["Statement"] = [ + s for s in existing_policy["Statement"] if s.get("Sid") not in our_sids + ] + if existing_policy["Statement"]: + s3.put_bucket_policy(Bucket=bucket_name, Policy=json.dumps(existing_policy)) + else: + s3.delete_bucket_policy(Bucket=bucket_name) + + def create(self, bucket_name, partition, service_type, *args, **kwargs): + added = self._add_policy(bucket_name, partition, service_type) + return {"AddedSids": added, "BucketName": bucket_name}, bucket_name + + def update(self, bucket_name, partition, service_type, *args, **kwargs): + added = self._add_policy(bucket_name, partition, service_type) + return {"AddedSids": added, "BucketName": bucket_name}, bucket_name + + def delete(self, bucket_name, service_type, *args, **kwargs): + self._remove_policy(bucket_name, service_type) + return {"BucketName": bucket_name}, bucket_name + + def extract_params(self, event): + props = event.get("ResourceProperties", {}) + return { + "bucket_name": props.get("BucketName"), + "partition": props.get("Partition", "aws"), + "service_type": props.get("ServiceType", "").lower(), + } + + if __name__ == '__main__': params = {"AWSResource": "s3"} # value = ConfigDeliveryChannel() diff --git a/sumologic-app-utils/sumo-app-utils.zip b/sumologic-app-utils/sumo-app-utils.zip index 06887226..10aeda98 100644 Binary files a/sumologic-app-utils/sumo-app-utils.zip and b/sumologic-app-utils/sumo-app-utils.zip differ