import json import time import boto3 from opensearchpy import OpenSearch, RequestsHttpConnection from requests_aws4auth import AWS4Auth def handler(event, context): print(f"Event: {json.dumps(event, default=str)}") if event["RequestType"] == "Delete": print("Delete request, returning") return {"PhysicalResourceId": event.get("PhysicalResourceId", "none")} props = event["ResourceProperties"] endpoint = props["Endpoint"].replace("https://", "") index_name = props["IndexName"] vector_field = props["VectorField"] text_field = props["TextField"] metadata_field = props["MetadataField"] print(f"Endpoint: {endpoint}") print(f"Index: {index_name}") session = boto3.Session() credentials = session.get_credentials().get_frozen_credentials() region = session.region_name print(f"Region: {region}") awsauth = AWS4Auth( credentials.access_key, credentials.secret_key, region, "aoss", session_token=credentials.token, ) client = OpenSearch( hosts=[{"host": endpoint, "port": 443}], http_auth=awsauth, use_ssl=True, verify_certs=True, connection_class=RequestsHttpConnection, timeout=30, ) index_body = { "settings": { "index": {"knn": True, "knn.algo_param.ef_search": 512} }, "mappings": { "properties": { vector_field: { "type": "knn_vector", "dimension": 1024, "method": { "engine": "faiss", "name": "hnsw", "space_type": "l2", }, }, text_field: {"type": "text"}, metadata_field: {"type": "text"}, } }, } for attempt in range(30): try: print(f"Attempt {attempt}: creating index...") response = client.indices.create(index=index_name, body=index_body) print(f"Index created successfully: {response}") # Verify the index exists exists = client.indices.exists(index=index_name) print(f"Index exists check: {exists}") return {"PhysicalResourceId": index_name} except Exception as e: error_str = str(e) print(f"Attempt {attempt} error: {error_str}") if "resource_already_exists_exception" in error_str: print("Index already exists, returning success") return {"PhysicalResourceId": index_name} if "403" in error_str and attempt < 29: print(f"403 error, retrying in 10s (attempt {attempt}/29)") time.sleep(10) continue print(f"Fatal error on attempt {attempt}: {error_str}") raise raise Exception("Timeout waiting for AOSS access policy propagation")