import time import boto3 from opensearchpy import OpenSearch, RequestsHttpConnection from requests_aws4auth import AWS4Auth def handler(event, context): if event["RequestType"] == "Delete": 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"] session = boto3.Session() credentials = session.get_credentials().get_frozen_credentials() region = session.region_name 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: client.indices.create(index=index_name, body=index_body) return {"PhysicalResourceId": index_name} except Exception as e: error_str = str(e) if "resource_already_exists_exception" in error_str: return {"PhysicalResourceId": index_name} if "403" in error_str and attempt < 29: time.sleep(10) continue raise raise Exception("Timeout waiting for AOSS access policy propagation")