proposal-system/lambdas/aurora-pgvector-init/app.py

127 lines
5.1 KiB
Python
Raw Permalink Normal View History

"""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}}