"""CloudFormation custom-resource Lambda: initialize the Aurora pgvector store for the Bedrock Knowledge Base (PR3 — replaces the OpenSearch index creator). Runs the one-time bootstrap SQL via the RDS Data API as the master user: enables pgvector, creates the `bedrock_integration` schema + `bedrock_kb` table (with the column/index layout Bedrock requires), and creates/owns the dedicated `bedrock_user` role whose credentials Bedrock uses to query the table. Idempotent: safe to run on stack create/update (IF NOT EXISTS throughout). Only boto3 (rds-data, secretsmanager) is used — both ship in the Lambda runtime. """ import json import logging import os import time import boto3 from botocore.exceptions import ClientError logger = logging.getLogger(__name__) logger.setLevel(os.environ.get("LOG_LEVEL", "INFO")) rds_data = boto3.client("rds-data") secrets = boto3.client("secretsmanager") CLUSTER_ARN = os.environ["CLUSTER_ARN"] MASTER_SECRET_ARN = os.environ["MASTER_SECRET_ARN"] BEDROCK_SECRET_ARN = os.environ["BEDROCK_SECRET_ARN"] DATABASE = os.environ.get("DATABASE", "proposals") EMBED_DIM = int(os.environ.get("EMBED_DIM", "1024")) # Titan Embed v2 PHYSICAL_ID = "aurora-pgvector-init" TABLE = "bedrock_integration.bedrock_kb" def _exec(sql: str, attempts: int = 6) -> None: """Run one statement via the RDS Data API as the master user. Retries transient errors while a Serverless v2 cluster is still becoming reachable right after deploy (Data API can briefly report the cluster as unavailable / resuming). """ for attempt in range(attempts): try: rds_data.execute_statement( resourceArn=CLUSTER_ARN, secretArn=MASTER_SECRET_ARN, database=DATABASE, sql=sql, ) return except ClientError as exc: code = exc.response.get("Error", {}).get("Code", "") message = str(exc) transient = code == "DatabaseResumingException" or any( s in message for s in ("not currently available", "Communication link failure") ) if transient and attempt < attempts - 1: logger.warning( "Transient Data API error (attempt %d/%d): %s", attempt + 1, attempts, code or message, ) time.sleep(10) continue raise def handler(event, context): request_type = event.get("RequestType") logger.info("RequestType=%s", request_type) # Leave the data in place on stack delete — nothing to undo. if request_type == "Delete": return {"PhysicalResourceId": PHYSICAL_ID} # The password Bedrock will use to log in as bedrock_user. Generated by # Secrets Manager with SQL-unsafe characters excluded (see CDK). secret = json.loads( secrets.get_secret_value(SecretId=BEDROCK_SECRET_ARN)["SecretString"] ) password = secret["password"] # Defense-in-depth: the password is inlined into CREATE/ALTER ROLE SQL (DDL can't # bind parameters). The secret is generated with excludePunctuation=true, so it must # be strictly alphanumeric — refuse anything else rather than risk SQL breakage. if not password.isalnum(): raise ValueError( "bedrock_user password is not alphanumeric; refusing to inline" ) statements = [ "CREATE EXTENSION IF NOT EXISTS vector;", "CREATE SCHEMA IF NOT EXISTS bedrock_integration;", # Create the role if missing, then (re)set its password to match the secret. "DO $$ BEGIN " "IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname = 'bedrock_user') " f"THEN CREATE ROLE bedrock_user LOGIN PASSWORD '{password}'; END IF; END $$;", f"ALTER ROLE bedrock_user WITH LOGIN PASSWORD '{password}';", "GRANT ALL ON SCHEMA bedrock_integration TO bedrock_user;", f"CREATE TABLE IF NOT EXISTS {TABLE} (" "id uuid PRIMARY KEY, " f"embedding vector({EMBED_DIM}), " "chunks text, " "metadata json, " "custom_metadata jsonb);", f"ALTER TABLE {TABLE} OWNER TO bedrock_user;", f"CREATE INDEX IF NOT EXISTS bedrock_kb_embedding_idx ON {TABLE} " "USING hnsw (embedding vector_cosine_ops) WITH (ef_construction = 256);", f"CREATE INDEX IF NOT EXISTS bedrock_kb_chunks_idx ON {TABLE} " "USING gin (to_tsvector('simple', chunks));", f"CREATE INDEX IF NOT EXISTS bedrock_kb_custom_metadata_idx ON {TABLE} " "USING gin (custom_metadata);", "GRANT ALL ON ALL TABLES IN SCHEMA bedrock_integration TO bedrock_user;", ] # Log only the position — the statements contain the bedrock_user password # (CREATE/ALTER ROLE), so the SQL text itself must never reach the logs. for i, sql in enumerate(statements, 1): logger.info("executing bootstrap statement %d/%d", i, len(statements)) _exec(sql) logger.info("pgvector store initialized: %s (dim=%d)", TABLE, EMBED_DIM) return {"PhysicalResourceId": PHYSICAL_ID, "Data": {"TableName": TABLE}}