Implement Lambda functions for PDF processing, suggestions, and library ingest

- pdf-extract: Parse vendor PDFs with pdfplumber, fallback to Claude multimodal
- pdf-generate: Generate branded proposal PDFs with reportlab Platypus
- library-ingest: Format approved proposals as markdown and sync to Bedrock KB
- suggestions: Query KB for similar proposals, generate line items via Claude
- All Lambdas use internal API key auth and cold-start secret caching
- Fix pdf_path unbound variable in pdf-extract error handling
This commit is contained in:
Adam Moussa 2026-05-16 22:10:21 -04:00
parent f9951d6dfa
commit 9d6612e337
9 changed files with 1621 additions and 8 deletions

View file

@ -1,14 +1,194 @@
"""Proposal System - Library Ingest Lambda.
Processes approved proposals into the Bedrock Knowledge Base.
Implementation in Phase 4.
Processes approved/sent proposals into the Bedrock Knowledge Base library.
Formats proposal data as structured markdown and uploads to the library bucket,
then triggers a KB sync.
"""
import json
import os
from datetime import datetime
import boto3
import httpx
LIBRARY_BUCKET = os.environ.get("LIBRARY_BUCKET", "")
KNOWLEDGE_BASE_ID = os.environ.get("KNOWLEDGE_BASE_ID", "")
DATA_SOURCE_ID = os.environ.get("DATA_SOURCE_ID", "")
API_BASE_URL = os.environ.get("API_BASE_URL", "")
INTERNAL_API_KEY_SECRET_ARN = os.environ.get("INTERNAL_API_KEY_SECRET_ARN", "")
s3 = boto3.client("s3")
bedrock_agent = boto3.client("bedrock-agent")
secrets_client = boto3.client("secretsmanager")
_cached_api_key: str | None = None
def _get_api_key() -> str:
global _cached_api_key
if _cached_api_key is None:
if INTERNAL_API_KEY_SECRET_ARN:
resp = secrets_client.get_secret_value(SecretId=INTERNAL_API_KEY_SECRET_ARN)
_cached_api_key = resp["SecretString"]
else:
_cached_api_key = ""
return _cached_api_key
def handler(event, context):
for record in event.get("Records", []):
body = json.loads(record["body"])
print(f"Processing library-ingest job: {body}")
payload = body.get("payload", body)
proposal_id = payload["proposalId"]
process_ingestion(proposal_id)
return {"statusCode": 200}
def process_ingestion(proposal_id: str):
proposal = fetch_proposal(proposal_id)
if not proposal:
print(f"Proposal {proposal_id} not found")
return
line_items = fetch_line_items(proposal_id)
document = format_proposal_document(proposal, line_items)
s3_key = upload_to_library(proposal, document)
if s3_key:
trigger_kb_sync()
def fetch_proposal(proposal_id: str) -> dict | None:
try:
resp = httpx.get(
f"{API_BASE_URL}/api/proposals/{proposal_id}",
headers=_api_headers(),
timeout=10,
)
if resp.status_code == 200:
return resp.json()
except Exception as e:
print(f"Error fetching proposal: {e}")
return None
def fetch_line_items(proposal_id: str) -> list[dict]:
try:
resp = httpx.get(
f"{API_BASE_URL}/api/proposals/{proposal_id}/line-items",
headers=_api_headers(),
timeout=10,
)
if resp.status_code == 200:
return resp.json()
except Exception as e:
print(f"Error fetching line items: {e}")
return []
def format_proposal_document(proposal: dict, line_items: list[dict]) -> str:
total = sum(li.get("totalPrice", 0) for li in line_items)
submitted_at = proposal.get("submittedAt", "")
if submitted_at:
try:
dt = datetime.fromisoformat(submitted_at.replace("Z", "+00:00"))
submitted_at = dt.strftime("%Y-%m-%d")
except (ValueError, TypeError):
pass
lines = [
f"# Proposal: {proposal['proposalNumber']}",
"",
f"**Customer:** {proposal['customerName']}",
f"**Address:** {proposal.get('customerAddress', '')}",
f"**Service Category:** {proposal['serviceCategory']}",
f"**Priority:** {proposal['priority']}",
f"**Date:** {submitted_at}",
f"**Total Bid Amount:** ${total:,.2f}",
f"**Work Order:** {proposal.get('workOrderNumber', '')}",
"",
"## Scope of Work",
"",
proposal.get("refinedScope") or proposal.get("scopeOfWork", ""),
"",
"## Line Items",
"",
"| # | Description | Qty | Unit | Unit Price | Total |",
"|---|---|---|---|---|---|",
]
for i, li in enumerate(line_items, 1):
desc = li.get("description", "")
qty = li.get("quantity", "")
unit = li.get("unit", "")
unit_price = li.get("unitPrice")
total_price = li.get("totalPrice", 0)
up_str = f"${unit_price:,.2f}" if unit_price else "-"
tp_str = f"${total_price:,.2f}"
lines.append(f"| {i} | {desc} | {qty} | {unit} | {up_str} | {tp_str} |")
lines.extend([
"",
f"**Total: ${total:,.2f}**",
])
return "\n".join(lines)
def upload_to_library(proposal: dict, document: str) -> str | None:
if not LIBRARY_BUCKET:
print("No library bucket configured")
return None
proposal_number = proposal["proposalNumber"]
category = proposal.get("serviceCategory", "General")
s3_key = f"proposals/{category.lower()}/{proposal_number}.md"
try:
s3.put_object(
Bucket=LIBRARY_BUCKET,
Key=s3_key,
Body=document.encode("utf-8"),
ContentType="text/markdown",
Metadata={
"service-category": category,
"proposal-number": proposal_number,
"customer-name": proposal.get("customerName", ""),
"total-amount": str(proposal.get("totalBidAmount", 0)),
"date-submitted": proposal.get("submittedAt", ""),
},
)
print(f"Uploaded {s3_key} to library bucket")
return s3_key
except Exception as e:
print(f"Error uploading to library: {e}")
return None
def trigger_kb_sync():
if not KNOWLEDGE_BASE_ID or not DATA_SOURCE_ID:
print("KB or data source ID not configured, skipping sync")
return
try:
response = bedrock_agent.start_ingestion_job(
knowledgeBaseId=KNOWLEDGE_BASE_ID,
dataSourceId=DATA_SOURCE_ID,
)
job_id = response.get("ingestionJob", {}).get("ingestionJobId", "")
print(f"Started KB ingestion job: {job_id}")
except Exception as e:
print(f"Error triggering KB sync: {e}")
def _api_headers() -> dict:
headers = {"Content-Type": "application/json"}
api_key = _get_api_key()
if api_key:
headers["X-Internal-Api-Key"] = api_key
return headers

View file

@ -0,0 +1,359 @@
"""Batch ingestion CLI for bootstrapping the Bedrock Knowledge Base.
Processes historical proposal PDFs from a local directory, extracts structured
data using pdfplumber + Claude, formats as markdown documents, and uploads to
the proposal-system-library S3 bucket. Triggers a KB sync after all uploads.
Usage:
python batch_ingest.py --input-dir ./historical-pdfs --bucket proposal-system-library-328440206208 --kb-id <KB_ID> --ds-id <DS_ID>
Prerequisites:
pip install pdfplumber boto3 rich
AWS credentials configured with access to S3, Bedrock, and the KB.
"""
import argparse
import json
import os
import sys
from pathlib import Path
import boto3
import pdfplumber
try:
from rich.console import Console
from rich.progress import Progress, SpinnerColumn, TextColumn, BarColumn, TaskProgressColumn
console = Console()
HAS_RICH = True
except ImportError:
HAS_RICH = False
s3 = boto3.client("s3")
bedrock_runtime = boto3.client("bedrock-runtime")
bedrock_agent = boto3.client("bedrock-agent")
MODEL_ID = "us.anthropic.claude-sonnet-4-5-20250929-v1:0"
SERVICE_CATEGORIES = ["HVAC", "Plumbing", "Electrical", "General", "Renovation"]
def main():
parser = argparse.ArgumentParser(description="Batch ingest historical proposals into Bedrock KB")
parser.add_argument("--input-dir", required=True, help="Directory containing historical proposal PDFs")
parser.add_argument("--bucket", required=True, help="S3 library bucket name")
parser.add_argument("--kb-id", required=True, help="Bedrock Knowledge Base ID")
parser.add_argument("--ds-id", required=True, help="Bedrock KB Data Source ID")
parser.add_argument("--dry-run", action="store_true", help="Parse and extract only, don't upload")
args = parser.parse_args()
input_dir = Path(args.input_dir)
if not input_dir.exists():
print(f"Error: {input_dir} does not exist")
sys.exit(1)
pdf_files = sorted(input_dir.glob("*.pdf"))
if not pdf_files:
print(f"No PDF files found in {input_dir}")
sys.exit(1)
print(f"Found {len(pdf_files)} PDF files to process")
results = {"success": 0, "failed": 0, "skipped": 0}
if HAS_RICH:
with Progress(
SpinnerColumn(),
TextColumn("[progress.description]{task.description}"),
BarColumn(),
TaskProgressColumn(),
console=console,
) as progress:
task = progress.add_task("Processing PDFs...", total=len(pdf_files))
for pdf_path in pdf_files:
status = process_pdf(pdf_path, args.bucket, args.dry_run)
results[status] += 1
progress.update(task, advance=1, description=f"[{'green' if status == 'success' else 'red'}]{pdf_path.name}")
else:
for i, pdf_path in enumerate(pdf_files, 1):
print(f"[{i}/{len(pdf_files)}] Processing {pdf_path.name}...", end=" ")
status = process_pdf(pdf_path, args.bucket, args.dry_run)
results[status] += 1
print(f"[{status.upper()}]")
print(f"\nResults: {results['success']} succeeded, {results['failed']} failed, {results['skipped']} skipped")
if not args.dry_run and results["success"] > 0:
print("\nTriggering Knowledge Base sync...")
trigger_kb_sync(args.kb_id, args.ds_id)
print("Done. KB ingestion job started — check AWS console for completion status.")
def process_pdf(pdf_path: Path, bucket: str, dry_run: bool) -> str:
try:
extracted = extract_proposal_data(pdf_path)
if not extracted.get("lineItems"):
extracted = extract_with_claude(pdf_path)
if not extracted.get("lineItems"):
print(f" Warning: No line items extracted from {pdf_path.name}")
return "skipped"
document = format_as_markdown(extracted, pdf_path.name)
if dry_run:
print(f"\n--- {pdf_path.name} ---")
print(f" Customer: {extracted.get('customerName', 'Unknown')}")
print(f" Category: {extracted.get('serviceCategory', 'General')}")
print(f" Line items: {len(extracted.get('lineItems', []))}")
print(f" Total: ${extracted.get('totalAmount', 0):,.2f}")
return "success"
s3_key = upload_to_library(bucket, extracted, document)
if s3_key:
return "success"
return "failed"
except Exception as e:
print(f" Error processing {pdf_path.name}: {e}")
return "failed"
def extract_proposal_data(pdf_path: Path) -> dict:
result = {
"customerName": "",
"serviceCategory": "General",
"lineItems": [],
"totalAmount": 0.0,
"scopeOfWork": "",
"date": "",
"proposalNumber": "",
}
with pdfplumber.open(str(pdf_path)) as pdf:
all_text = ""
all_tables = []
for page in pdf.pages:
text = page.extract_text() or ""
all_text += text + "\n"
tables = page.extract_tables()
all_tables.extend(tables)
result["scopeOfWork"] = all_text[:2000]
if all_tables:
result["lineItems"] = parse_tables(all_tables)
result["totalAmount"] = sum(
li.get("totalPrice", 0) for li in result["lineItems"]
)
return result
def extract_with_claude(pdf_path: Path) -> dict:
import base64
with open(pdf_path, "rb") as f:
pdf_b64 = base64.standard_b64encode(f.read()).decode("utf-8")
try:
response = bedrock_runtime.invoke_model(
modelId=MODEL_ID,
contentType="application/json",
accept="application/json",
body=json.dumps({
"anthropic_version": "bedrock-2023-05-31",
"max_tokens": 4096,
"messages": [{
"role": "user",
"content": [
{
"type": "document",
"source": {
"type": "base64",
"media_type": "application/pdf",
"data": pdf_b64,
},
},
{
"type": "text",
"text": f"""Extract structured data from this proposal PDF.
Return a JSON object with:
- customerName: the customer/client name
- serviceCategory: one of {SERVICE_CATEGORIES}
- proposalNumber: the proposal/quote number if visible
- date: the date on the proposal (YYYY-MM-DD format)
- scopeOfWork: brief description of the work (1-2 sentences)
- lineItems: array of objects with: description (string), quantity (number), unit (string), unitPrice (number or null), totalPrice (number)
- totalAmount: the grand total
Respond ONLY with the JSON object.""",
},
],
}],
"temperature": 0.1,
}),
)
response_body = json.loads(response["body"].read())
content = response_body["content"][0]["text"].strip()
if content.startswith("```"):
content = content.split("\n", 1)[1]
content = content.rsplit("```", 1)[0]
return json.loads(content)
except Exception as e:
print(f" Claude extraction failed: {e}")
return {}
def parse_tables(tables: list) -> list[dict]:
line_items = []
for table in tables:
if not table or len(table) < 2:
continue
header = [str(cell).lower().strip() if cell else "" for cell in table[0]]
desc_col = _find_column(header, ["description", "item", "service", "work", "scope"])
qty_col = _find_column(header, ["qty", "quantity", "count"])
price_col = _find_column(header, ["unit price", "rate", "price/unit", "unit cost"])
total_col = _find_column(header, ["total", "amount", "ext", "extended", "line total"])
if desc_col is None:
continue
for row in table[1:]:
if not row or len(row) <= desc_col:
continue
description = str(row[desc_col]).strip() if row[desc_col] else ""
if not description or description.lower() in ("", "total", "subtotal", "grand total"):
continue
quantity = _parse_number(row[qty_col]) if qty_col is not None and qty_col < len(row) else None
unit_price = _parse_number(row[price_col]) if price_col is not None and price_col < len(row) else None
total = _parse_number(row[total_col]) if total_col is not None and total_col < len(row) else None
if total is None and quantity and unit_price:
total = quantity * unit_price
if description and (total or unit_price):
line_items.append({
"description": description,
"quantity": quantity or 1,
"unit": "each",
"unitPrice": unit_price,
"totalPrice": total or 0,
})
return line_items
def _find_column(header: list[str], keywords: list[str]) -> int | None:
for i, col in enumerate(header):
for kw in keywords:
if kw in col:
return i
return None
def _parse_number(value) -> float | None:
if value is None:
return None
try:
cleaned = str(value).replace("$", "").replace(",", "").strip()
if not cleaned or cleaned == "-":
return None
return float(cleaned)
except (ValueError, TypeError):
return None
def format_as_markdown(extracted: dict, source_filename: str) -> str:
customer = extracted.get("customerName", "Unknown Customer")
category = extracted.get("serviceCategory", "General")
proposal_num = extracted.get("proposalNumber", source_filename.replace(".pdf", ""))
date = extracted.get("date", "")
scope = extracted.get("scopeOfWork", "")
total = extracted.get("totalAmount", 0)
line_items = extracted.get("lineItems", [])
lines = [
f"# Proposal: {proposal_num}",
"",
f"**Customer:** {customer}",
f"**Service Category:** {category}",
f"**Date:** {date}",
f"**Total Bid Amount:** ${total:,.2f}",
f"**Source File:** {source_filename}",
"",
"## Scope of Work",
"",
scope[:1000] if scope else "Not specified",
"",
"## Line Items",
"",
"| # | Description | Qty | Unit | Unit Price | Total |",
"|---|---|---|---|---|---|",
]
for i, li in enumerate(line_items, 1):
desc = li.get("description", "")
qty = li.get("quantity", "")
unit = li.get("unit", "each")
up = li.get("unitPrice")
tp = li.get("totalPrice", 0)
up_str = f"${up:,.2f}" if up else "-"
tp_str = f"${tp:,.2f}"
lines.append(f"| {i} | {desc} | {qty} | {unit} | {up_str} | {tp_str} |")
lines.extend(["", f"**Total: ${total:,.2f}**"])
return "\n".join(lines)
def upload_to_library(bucket: str, extracted: dict, document: str) -> str | None:
category = extracted.get("serviceCategory", "General")
proposal_num = extracted.get("proposalNumber", f"historical-{os.urandom(4).hex()}")
s3_key = f"proposals/{category.lower()}/{proposal_num}.md"
try:
s3.put_object(
Bucket=bucket,
Key=s3_key,
Body=document.encode("utf-8"),
ContentType="text/markdown",
Metadata={
"service-category": category,
"proposal-number": proposal_num,
"customer-name": extracted.get("customerName", ""),
"total-amount": str(extracted.get("totalAmount", 0)),
"date-submitted": extracted.get("date", ""),
"source": "batch-ingest",
},
)
return s3_key
except Exception as e:
print(f" Upload failed: {e}")
return None
def trigger_kb_sync(kb_id: str, ds_id: str):
try:
response = bedrock_agent.start_ingestion_job(
knowledgeBaseId=kb_id,
dataSourceId=ds_id,
)
job_id = response.get("ingestionJob", {}).get("ingestionJobId", "")
print(f" Ingestion job started: {job_id}")
except Exception as e:
print(f" Error triggering KB sync: {e}")
if __name__ == "__main__":
main()

View file

@ -1 +1,2 @@
boto3>=1.35.0,<2.0
httpx>=0.27.0,<1.0

View file

@ -1,14 +1,300 @@
"""Proposal System - PDF Extract Lambda.
Parses vendor proposal PDFs and extracts structured line item data.
Implementation in Phase 4.
Falls back to Claude multimodal for scanned/image-based PDFs.
"""
import json
import os
import tempfile
import base64
import boto3
import httpx
import pdfplumber
UPLOADS_BUCKET = os.environ.get("UPLOADS_BUCKET", "")
API_BASE_URL = os.environ.get("API_BASE_URL", "")
MODEL_ID = os.environ.get("MODEL_ID", "us.anthropic.claude-sonnet-4-5-20250929-v1:0")
INTERNAL_API_KEY_SECRET_ARN = os.environ.get("INTERNAL_API_KEY_SECRET_ARN", "")
s3 = boto3.client("s3")
bedrock_runtime = boto3.client("bedrock-runtime")
secrets_client = boto3.client("secretsmanager")
_cached_api_key: str | None = None
def _get_api_key() -> str:
global _cached_api_key
if _cached_api_key is None:
if INTERNAL_API_KEY_SECRET_ARN:
resp = secrets_client.get_secret_value(SecretId=INTERNAL_API_KEY_SECRET_ARN)
_cached_api_key = resp["SecretString"]
else:
_cached_api_key = ""
return _cached_api_key
def handler(event, context):
for record in event.get("Records", []):
body = json.loads(record["body"])
print(f"Processing pdf-extract job: {body}")
payload = body.get("payload", body)
proposal_id = payload["proposalId"]
s3_key = payload.get("s3Key", "")
vendor_proposal_id = payload.get("vendorProposalId", "")
if not s3_key:
print(f"No s3Key in payload for proposal {proposal_id}")
continue
process_pdf(proposal_id, s3_key, vendor_proposal_id)
return {"statusCode": 200}
def process_pdf(proposal_id: str, s3_key: str, vendor_proposal_id: str):
update_processing_status(vendor_proposal_id, "Processing")
pdf_path = None
try:
pdf_path = download_pdf(s3_key)
extracted = extract_with_pdfplumber(pdf_path)
if not extracted["lineItems"] and extracted["rawText"].strip():
extracted = extract_with_claude_multimodal(pdf_path)
if not extracted["lineItems"] and not extracted["rawText"].strip():
extracted = extract_with_claude_multimodal(pdf_path)
save_extraction(vendor_proposal_id, extracted)
except Exception as e:
print(f"Error processing PDF: {e}")
update_processing_status(vendor_proposal_id, "Failed")
finally:
if pdf_path:
try:
os.unlink(pdf_path)
except Exception:
pass
def download_pdf(s3_key: str) -> str:
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=".pdf")
s3.download_file(UPLOADS_BUCKET, s3_key, tmp.name)
tmp.close()
return tmp.name
def extract_with_pdfplumber(pdf_path: str) -> dict:
result = {
"vendorName": "",
"lineItems": [],
"rawText": "",
"totalVendorCost": 0.0,
}
try:
with pdfplumber.open(pdf_path) as pdf:
all_text = ""
all_tables = []
for page in pdf.pages:
text = page.extract_text() or ""
all_text += text + "\n"
tables = page.extract_tables()
for table in tables:
all_tables.append(table)
result["rawText"] = all_text.strip()
if all_tables:
result["lineItems"] = parse_tables(all_tables)
result["totalVendorCost"] = sum(
li.get("total", 0) for li in result["lineItems"]
)
if not result["vendorName"] and all_text:
lines = all_text.split("\n")
for line in lines[:5]:
stripped = line.strip()
if stripped and len(stripped) > 3 and not stripped[0].isdigit():
result["vendorName"] = stripped
break
except Exception as e:
print(f"pdfplumber extraction failed: {e}")
return result
def parse_tables(tables: list) -> list[dict]:
line_items = []
for table in tables:
if not table or len(table) < 2:
continue
header = [str(cell).lower().strip() if cell else "" for cell in table[0]]
desc_col = find_column(header, ["description", "item", "service", "work", "scope"])
qty_col = find_column(header, ["qty", "quantity", "count"])
price_col = find_column(header, ["unit price", "rate", "price/unit", "unit cost"])
total_col = find_column(header, ["total", "amount", "ext", "extended", "line total"])
if desc_col is None:
continue
for row in table[1:]:
if not row or len(row) <= desc_col:
continue
description = str(row[desc_col]).strip() if row[desc_col] else ""
if not description or description.lower() in ("", "total", "subtotal", "grand total"):
continue
quantity = parse_number(row[qty_col]) if qty_col is not None and qty_col < len(row) else None
unit_price = parse_number(row[price_col]) if price_col is not None and price_col < len(row) else None
total = parse_number(row[total_col]) if total_col is not None and total_col < len(row) else None
if total is None and quantity and unit_price:
total = quantity * unit_price
if description and (total or unit_price):
line_items.append({
"description": description,
"quantity": quantity,
"unitPrice": unit_price,
"total": total,
})
return line_items
def find_column(header: list[str], keywords: list[str]) -> int | None:
for i, col in enumerate(header):
for kw in keywords:
if kw in col:
return i
return None
def parse_number(value) -> float | None:
if value is None:
return None
try:
cleaned = str(value).replace("$", "").replace(",", "").strip()
if not cleaned or cleaned == "-":
return None
return float(cleaned)
except (ValueError, TypeError):
return None
def extract_with_claude_multimodal(pdf_path: str) -> dict:
try:
with open(pdf_path, "rb") as f:
pdf_bytes = f.read()
pdf_b64 = base64.standard_b64encode(pdf_bytes).decode("utf-8")
response = bedrock_runtime.invoke_model(
modelId=MODEL_ID,
contentType="application/json",
accept="application/json",
body=json.dumps({
"anthropic_version": "bedrock-2023-05-31",
"max_tokens": 4096,
"messages": [{
"role": "user",
"content": [
{
"type": "document",
"source": {
"type": "base64",
"media_type": "application/pdf",
"data": pdf_b64,
},
},
{
"type": "text",
"text": """Extract all line items from this vendor proposal PDF.
Return a JSON object with these fields:
- vendorName: the vendor/company name
- lineItems: array of objects with: description, quantity (number or null), unitPrice (number or null), total (number or null)
- totalVendorCost: the grand total amount
Respond ONLY with the JSON object, no additional text.""",
},
],
}],
"temperature": 0.1,
}),
)
response_body = json.loads(response["body"].read())
content = response_body["content"][0]["text"]
content = content.strip()
if content.startswith("```"):
content = content.split("\n", 1)[1]
content = content.rsplit("```", 1)[0]
parsed = json.loads(content)
return {
"vendorName": parsed.get("vendorName", ""),
"lineItems": parsed.get("lineItems", []),
"rawText": "",
"totalVendorCost": float(parsed.get("totalVendorCost", 0)),
}
except Exception as e:
print(f"Claude multimodal extraction failed: {e}")
return {"vendorName": "", "lineItems": [], "rawText": "", "totalVendorCost": 0.0}
def save_extraction(vendor_proposal_id: str, extracted: dict):
extracted_data = {
"lineItems": extracted["lineItems"],
"vendorName": extracted["vendorName"],
}
try:
resp = httpx.put(
f"{API_BASE_URL}/api/vendor-proposals/{vendor_proposal_id}",
json={
"vendorName": extracted["vendorName"],
"extractedData": json.dumps(extracted_data),
"totalVendorCost": extracted["totalVendorCost"],
"processingStatus": "Complete",
},
headers=_api_headers(),
timeout=10,
)
if resp.status_code not in (200, 204):
print(f"Failed to save extraction: {resp.status_code} {resp.text}")
except Exception as e:
print(f"Error saving extraction: {e}")
def update_processing_status(vendor_proposal_id: str, status: str):
if not vendor_proposal_id:
return
try:
httpx.put(
f"{API_BASE_URL}/api/vendor-proposals/{vendor_proposal_id}/status",
json={"processingStatus": status},
headers=_api_headers(),
timeout=10,
)
except Exception as e:
print(f"Error updating status: {e}")
def _api_headers() -> dict:
headers = {"Content-Type": "application/json"}
api_key = _get_api_key()
if api_key:
headers["X-Internal-Api-Key"] = api_key
return headers

View file

@ -1,2 +1,3 @@
pdfplumber>=0.11.0,<1.0
boto3>=1.35.0,<2.0
httpx>=0.27.0,<1.0

View file

@ -1,14 +1,496 @@
"""Proposal System - PDF Generate Lambda.
Generates professional branded proposal PDFs from approved proposals.
Implementation in Phase 5.
Generates professional branded proposal PDFs using reportlab Platypus.
Triggered via SQS when an admin requests PDF generation.
"""
import json
import os
import tempfile
from datetime import datetime
from io import BytesIO
import boto3
import httpx
from reportlab.lib import colors
from reportlab.lib.enums import TA_CENTER, TA_LEFT, TA_RIGHT
from reportlab.lib.pagesizes import letter
from reportlab.lib.styles import ParagraphStyle, getSampleStyleSheet
from reportlab.lib.units import inch, mm
from reportlab.platypus import (
Paragraph,
SimpleDocTemplate,
Spacer,
Table,
TableStyle,
)
GENERATED_BUCKET = os.environ.get("GENERATED_BUCKET", "")
API_BASE_URL = os.environ.get("API_BASE_URL", "")
INTERNAL_API_KEY_SECRET_ARN = os.environ.get("INTERNAL_API_KEY_SECRET_ARN", "")
s3 = boto3.client("s3")
secrets_client = boto3.client("secretsmanager")
_cached_api_key: str | None = None
COMPANY_NAME = "Sea Haven Industries"
COMPANY_ADDRESS = "Sea Haven Industries LLC"
COMPANY_PHONE = ""
COMPANY_EMAIL = "info@seahavenind.com"
TERMS_AND_CONDITIONS = """
1. This proposal is valid for 30 days from the date of issue.
2. Payment terms: Net 30 days from invoice date.
3. Any changes to the scope of work may result in additional charges.
4. Work will be scheduled upon acceptance of this proposal.
5. All materials and workmanship are guaranteed for one (1) year from completion.
6. Client is responsible for providing access to the work area.
7. This proposal does not include permits unless specifically noted in the line items.
""".strip()
def _get_api_key() -> str:
global _cached_api_key
if _cached_api_key is None:
if INTERNAL_API_KEY_SECRET_ARN:
resp = secrets_client.get_secret_value(SecretId=INTERNAL_API_KEY_SECRET_ARN)
_cached_api_key = resp["SecretString"]
else:
_cached_api_key = ""
return _cached_api_key
def handler(event, context):
for record in event.get("Records", []):
body = json.loads(record["body"])
print(f"Processing pdf-generate job: {body}")
payload = body.get("payload", body)
proposal_id = payload["proposalId"]
generate_pdf(proposal_id)
return {"statusCode": 200}
def generate_pdf(proposal_id: str):
proposal = fetch_proposal(proposal_id)
if not proposal:
print(f"Proposal {proposal_id} not found")
return
line_items = fetch_line_items(proposal_id)
pdf_bytes = build_pdf(proposal, line_items)
proposal_number = proposal["proposalNumber"]
revision = proposal.get("currentRevision", 1)
s3_key = f"{proposal_number}/rev-{revision}.pdf"
upload_pdf(s3_key, pdf_bytes)
register_pdf(proposal_id, s3_key)
print(f"Generated PDF: {s3_key} ({len(pdf_bytes)} bytes)")
def fetch_proposal(proposal_id: str) -> dict | None:
try:
resp = httpx.get(
f"{API_BASE_URL}/api/proposals/{proposal_id}",
headers=_api_headers(),
timeout=10,
)
if resp.status_code == 200:
return resp.json()
except Exception as e:
print(f"Error fetching proposal: {e}")
return None
def fetch_line_items(proposal_id: str) -> list[dict]:
try:
resp = httpx.get(
f"{API_BASE_URL}/api/proposals/{proposal_id}/line-items",
headers=_api_headers(),
timeout=10,
)
if resp.status_code == 200:
return resp.json()
except Exception as e:
print(f"Error fetching line items: {e}")
return []
def build_pdf(proposal: dict, line_items: list[dict]) -> bytes:
buffer = BytesIO()
doc = SimpleDocTemplate(
buffer,
pagesize=letter,
leftMargin=0.75 * inch,
rightMargin=0.75 * inch,
topMargin=0.75 * inch,
bottomMargin=0.75 * inch,
)
styles = _get_styles()
elements = []
# Header
elements.extend(_build_header(proposal, styles))
elements.append(Spacer(1, 0.3 * inch))
# Proposal metadata
elements.extend(_build_metadata(proposal, styles))
elements.append(Spacer(1, 0.3 * inch))
# Scope of work
elements.extend(_build_scope(proposal, styles))
elements.append(Spacer(1, 0.3 * inch))
# Line items table
elements.extend(_build_line_items_table(line_items, styles))
elements.append(Spacer(1, 0.4 * inch))
# Terms and conditions
elements.extend(_build_terms(styles))
doc.build(elements, onFirstPage=_page_footer, onLaterPages=_page_footer)
return buffer.getvalue()
def _get_styles():
styles = getSampleStyleSheet()
styles.add(ParagraphStyle(
"CompanyName",
parent=styles["Heading1"],
fontSize=18,
leading=22,
textColor=colors.HexColor("#1a237e"),
spaceAfter=2,
))
styles.add(ParagraphStyle(
"CompanyInfo",
parent=styles["Normal"],
fontSize=9,
leading=12,
textColor=colors.HexColor("#555555"),
))
styles.add(ParagraphStyle(
"ProposalTitle",
parent=styles["Heading2"],
fontSize=14,
leading=18,
textColor=colors.HexColor("#1a237e"),
spaceBefore=6,
spaceAfter=12,
))
styles.add(ParagraphStyle(
"SectionHeader",
parent=styles["Heading3"],
fontSize=11,
leading=14,
textColor=colors.HexColor("#1a237e"),
spaceBefore=8,
spaceAfter=6,
borderWidth=0,
))
styles.add(ParagraphStyle(
"MetaLabel",
parent=styles["Normal"],
fontSize=9,
leading=12,
textColor=colors.HexColor("#666666"),
))
styles.add(ParagraphStyle(
"MetaValue",
parent=styles["Normal"],
fontSize=10,
leading=13,
fontName="Helvetica-Bold",
))
styles.add(ParagraphStyle(
"ScopeText",
parent=styles["Normal"],
fontSize=10,
leading=14,
spaceBefore=4,
))
styles.add(ParagraphStyle(
"TermsText",
parent=styles["Normal"],
fontSize=8,
leading=11,
textColor=colors.HexColor("#555555"),
))
styles.add(ParagraphStyle(
"TotalLabel",
parent=styles["Normal"],
fontSize=11,
leading=14,
fontName="Helvetica-Bold",
alignment=TA_RIGHT,
))
styles.add(ParagraphStyle(
"FooterText",
parent=styles["Normal"],
fontSize=8,
leading=10,
textColor=colors.HexColor("#888888"),
alignment=TA_CENTER,
))
return styles
def _build_header(proposal: dict, styles) -> list:
revision = proposal.get("currentRevision", 1)
revision_text = f" | Rev {revision}" if revision > 1 else ""
header_data = [
[
Paragraph(COMPANY_NAME, styles["CompanyName"]),
Paragraph(f"PROPOSAL{revision_text}", styles["ProposalTitle"]),
],
[
Paragraph(f"{COMPANY_EMAIL}", styles["CompanyInfo"]),
Paragraph(f"#{proposal['proposalNumber']}", styles["MetaValue"]),
],
]
header_table = Table(header_data, colWidths=[3.5 * inch, 3.5 * inch])
header_table.setStyle(TableStyle([
("VALIGN", (0, 0), (-1, -1), "TOP"),
("ALIGN", (1, 0), (1, -1), "RIGHT"),
("LINEBELOW", (0, -1), (-1, -1), 1.5, colors.HexColor("#1a237e")),
("BOTTOMPADDING", (0, -1), (-1, -1), 8),
]))
return [header_table]
def _build_metadata(proposal: dict, styles) -> list:
submitted_at = proposal.get("submittedAt", "")
if submitted_at:
try:
dt = datetime.fromisoformat(submitted_at.replace("Z", "+00:00"))
submitted_at = dt.strftime("%B %d, %Y")
except (ValueError, TypeError):
pass
approved_at = proposal.get("approvedAt", "")
if approved_at:
try:
dt = datetime.fromisoformat(approved_at.replace("Z", "+00:00"))
approved_at = dt.strftime("%B %d, %Y")
except (ValueError, TypeError):
pass
meta_data = [
[
Paragraph("Customer", styles["MetaLabel"]),
Paragraph("Site Address", styles["MetaLabel"]),
],
[
Paragraph(proposal.get("customerName", ""), styles["MetaValue"]),
Paragraph(proposal.get("customerAddress", ""), styles["MetaValue"]),
],
[
Paragraph("Work Order #", styles["MetaLabel"]),
Paragraph("Date", styles["MetaLabel"]),
],
[
Paragraph(proposal.get("workOrderNumber", ""), styles["MetaValue"]),
Paragraph(approved_at or submitted_at, styles["MetaValue"]),
],
[
Paragraph("Category", styles["MetaLabel"]),
Paragraph("Priority", styles["MetaLabel"]),
],
[
Paragraph(proposal.get("serviceCategory", ""), styles["MetaValue"]),
Paragraph(proposal.get("priority", ""), styles["MetaValue"]),
],
]
meta_table = Table(meta_data, colWidths=[3.5 * inch, 3.5 * inch])
meta_table.setStyle(TableStyle([
("VALIGN", (0, 0), (-1, -1), "TOP"),
("TOPPADDING", (0, 0), (-1, -1), 2),
("BOTTOMPADDING", (0, 0), (-1, -1), 2),
]))
return [meta_table]
def _build_scope(proposal: dict, styles) -> list:
scope = proposal.get("refinedScope") or proposal.get("scopeOfWork", "")
if not scope:
return []
return [
Paragraph("Scope of Work", styles["SectionHeader"]),
Paragraph(scope, styles["ScopeText"]),
]
def _build_line_items_table(line_items: list[dict], styles) -> list:
if not line_items:
return [Paragraph("No line items", styles["Normal"])]
elements = [Paragraph("Itemized Pricing", styles["SectionHeader"])]
header = ["#", "Description", "Qty", "Unit", "Unit Price", "Total"]
table_data = [header]
subtotal = 0.0
for i, li in enumerate(line_items, 1):
qty = li.get("quantity", "")
unit = li.get("unit", "")
unit_price = li.get("unitPrice")
total_price = li.get("totalPrice", 0)
pricing_mode = li.get("pricingMode", "TotalPrice")
subtotal += float(total_price or 0)
if pricing_mode == "TotalPrice":
up_str = "-"
elif unit_price is not None:
up_str = f"${float(unit_price):,.2f}"
else:
up_str = "-"
tp_str = f"${float(total_price):,.2f}" if total_price else "-"
row = [
str(i),
Paragraph(li.get("description", ""), styles["Normal"]),
str(qty) if qty else "",
unit,
up_str,
tp_str,
]
table_data.append(row)
col_widths = [0.35 * inch, 3.15 * inch, 0.55 * inch, 0.7 * inch, 1.0 * inch, 1.0 * inch]
table = Table(table_data, colWidths=col_widths, repeatRows=1)
table.setStyle(TableStyle([
# Header row
("BACKGROUND", (0, 0), (-1, 0), colors.HexColor("#1a237e")),
("TEXTCOLOR", (0, 0), (-1, 0), colors.white),
("FONTNAME", (0, 0), (-1, 0), "Helvetica-Bold"),
("FONTSIZE", (0, 0), (-1, 0), 9),
("BOTTOMPADDING", (0, 0), (-1, 0), 6),
("TOPPADDING", (0, 0), (-1, 0), 6),
# Data rows
("FONTSIZE", (0, 1), (-1, -1), 9),
("TOPPADDING", (0, 1), (-1, -1), 4),
("BOTTOMPADDING", (0, 1), (-1, -1), 4),
("VALIGN", (0, 0), (-1, -1), "MIDDLE"),
# Alignment
("ALIGN", (0, 0), (0, -1), "CENTER"),
("ALIGN", (2, 0), (2, -1), "CENTER"),
("ALIGN", (3, 0), (3, -1), "CENTER"),
("ALIGN", (4, 0), (4, -1), "RIGHT"),
("ALIGN", (5, 0), (5, -1), "RIGHT"),
# Grid
("LINEBELOW", (0, 0), (-1, 0), 1, colors.HexColor("#1a237e")),
("LINEBELOW", (0, 1), (-1, -2), 0.5, colors.HexColor("#e0e0e0")),
("LINEBELOW", (0, -1), (-1, -1), 1, colors.HexColor("#1a237e")),
# Alternating row colors
*[("BACKGROUND", (0, i), (-1, i), colors.HexColor("#f5f5f5"))
for i in range(2, len(table_data), 2)],
]))
elements.append(table)
elements.append(Spacer(1, 0.15 * inch))
# Total row
total_data = [
["", "", "", "", "TOTAL:", f"${subtotal:,.2f}"],
]
total_table = Table(total_data, colWidths=col_widths)
total_table.setStyle(TableStyle([
("FONTNAME", (0, 0), (-1, -1), "Helvetica-Bold"),
("FONTSIZE", (0, 0), (-1, -1), 11),
("ALIGN", (4, 0), (4, 0), "RIGHT"),
("ALIGN", (5, 0), (5, 0), "RIGHT"),
("TOPPADDING", (0, 0), (-1, -1), 4),
("LINEABOVE", (4, 0), (5, 0), 1.5, colors.HexColor("#1a237e")),
]))
elements.append(total_table)
return elements
def _build_terms(styles) -> list:
elements = [
Spacer(1, 0.2 * inch),
Paragraph("Terms & Conditions", styles["SectionHeader"]),
]
for line in TERMS_AND_CONDITIONS.split("\n"):
elements.append(Paragraph(line, styles["TermsText"]))
return elements
def _page_footer(canvas, doc):
canvas.saveState()
page_num = canvas.getPageNumber()
footer_text = f"Page {page_num}"
canvas.setFont("Helvetica", 8)
canvas.setFillColor(colors.HexColor("#888888"))
canvas.drawCentredString(letter[0] / 2, 0.4 * inch, footer_text)
canvas.drawString(
0.75 * inch,
0.4 * inch,
f"{COMPANY_NAME} — Confidential",
)
canvas.restoreState()
def upload_pdf(s3_key: str, pdf_bytes: bytes):
try:
s3.put_object(
Bucket=GENERATED_BUCKET,
Key=s3_key,
Body=pdf_bytes,
ContentType="application/pdf",
)
except Exception as e:
print(f"Error uploading PDF: {e}")
raise
def register_pdf(proposal_id: str, s3_key: str):
try:
resp = httpx.post(
f"{API_BASE_URL}/api/generated-pdfs",
json={"proposalId": proposal_id, "s3Key": s3_key},
headers=_api_headers(),
timeout=10,
)
if resp.status_code not in (200, 201):
print(f"Failed to register PDF: {resp.status_code} {resp.text}")
except Exception as e:
print(f"Error registering PDF: {e}")
def _api_headers() -> dict:
headers = {"Content-Type": "application/json"}
api_key = _get_api_key()
if api_key:
headers["X-Internal-Api-Key"] = api_key
return headers

View file

@ -1,2 +1,3 @@
reportlab>=4.2.0,<5.0
boto3>=1.35.0,<2.0
httpx>=0.27.0,<1.0

301
lambdas/suggestions/app.py Normal file
View file

@ -0,0 +1,301 @@
"""Proposal System - Suggestion Engine Lambda.
Queries Bedrock Knowledge Base for similar proposals and invokes Claude
to generate line item suggestions for new proposals.
"""
import json
import os
import boto3
import httpx
KNOWLEDGE_BASE_ID = os.environ.get("KNOWLEDGE_BASE_ID", "")
MODEL_ID = os.environ.get("MODEL_ID", "us.anthropic.claude-sonnet-4-5-20250929-v1:0")
API_BASE_URL = os.environ.get("API_BASE_URL", "")
INTERNAL_API_KEY_SECRET_ARN = os.environ.get("INTERNAL_API_KEY_SECRET_ARN", "")
bedrock_agent = boto3.client("bedrock-agent-runtime")
bedrock_runtime = boto3.client("bedrock-runtime")
secrets_client = boto3.client("secretsmanager")
_cached_api_key: str | None = None
def _get_api_key() -> str:
global _cached_api_key
if _cached_api_key is None:
if INTERNAL_API_KEY_SECRET_ARN:
resp = secrets_client.get_secret_value(SecretId=INTERNAL_API_KEY_SECRET_ARN)
_cached_api_key = resp["SecretString"]
else:
_cached_api_key = ""
return _cached_api_key
def handler(event, context):
for record in event.get("Records", []):
body = json.loads(record["body"])
payload = body.get("payload", body)
proposal_id = payload["proposalId"]
trigger = payload.get("trigger", "generate")
process_suggestion(proposal_id, trigger)
return {"statusCode": 200}
def process_suggestion(proposal_id: str, trigger: str):
proposal = fetch_proposal(proposal_id)
if not proposal:
print(f"Proposal {proposal_id} not found")
return
scope = proposal.get("refinedScope") or proposal.get("scopeOfWork", "")
category = proposal.get("serviceCategory", "")
priority = proposal.get("priority", "")
existing_items = fetch_line_items(proposal_id)
similar_proposals = retrieve_similar(scope, category)
suggested_items = generate_line_items(scope, category, priority, similar_proposals)
post_line_items(proposal_id, suggested_items, existing_items)
store_similar_references(proposal_id, similar_proposals)
update_status_to_in_review(proposal_id)
def fetch_line_items(proposal_id: str) -> list[dict]:
try:
resp = httpx.get(
f"{API_BASE_URL}/api/proposals/{proposal_id}/line-items",
headers=_api_headers(),
timeout=10,
)
if resp.status_code == 200:
return resp.json()
except Exception as e:
print(f"Error fetching line items: {e}")
return []
def fetch_proposal(proposal_id: str) -> dict | None:
try:
resp = httpx.get(
f"{API_BASE_URL}/api/proposals/{proposal_id}",
headers=_api_headers(),
timeout=10,
)
if resp.status_code == 200:
return resp.json()
except Exception as e:
print(f"Error fetching proposal: {e}")
return None
def retrieve_similar(scope: str, category: str) -> list[dict]:
if not KNOWLEDGE_BASE_ID:
print("No Knowledge Base configured, skipping retrieval")
return []
try:
filter_config = {
"equals": {"key": "service_category", "value": category}
} if category else None
params = {
"knowledgeBaseId": KNOWLEDGE_BASE_ID,
"retrievalQuery": {"text": scope},
"retrievalConfiguration": {
"vectorSearchConfiguration": {
"numberOfResults": 10,
}
},
}
if filter_config:
params["retrievalConfiguration"]["vectorSearchConfiguration"]["filter"] = filter_config
response = bedrock_agent.retrieve(**params)
results = []
for result in response.get("retrievalResults", []):
content = result.get("content", {}).get("text", "")
score = result.get("score", 0.0)
metadata = result.get("metadata", {})
source_uri = result.get("location", {}).get("s3Location", {}).get("uri", "")
results.append({
"content": content,
"score": score,
"metadata": metadata,
"sourceUri": source_uri,
})
return results
except Exception as e:
print(f"Error retrieving from KB: {e}")
return []
def generate_line_items(
scope: str,
category: str,
priority: str,
similar_proposals: list[dict],
) -> list[dict]:
context_block = ""
if similar_proposals:
context_block = "Here are similar historical proposals and their line items for reference:\n\n"
for i, sp in enumerate(similar_proposals[:5], 1):
context_block += f"--- Similar Proposal {i} (relevance: {sp['score']:.2f}) ---\n"
context_block += sp["content"] + "\n\n"
prompt = f"""You are a construction/facilities proposal estimator for Sea Haven Industries.
Based on the scope of work and similar historical proposals, generate a detailed list of line items
with quantities, units, and estimated pricing.
Service Category: {category}
Priority: {priority}
Scope of Work:
{scope}
{context_block}
Generate line items as a JSON array. Each item should have:
- description: clear description of the work/material
- quantity: numeric quantity
- unit: unit of measurement (e.g., "sq ft", "hours", "each", "linear ft")
- unitPrice: price per unit in dollars (or null if lump sum)
- totalPrice: total price for this line item in dollars
- pricingMode: "UnitPrice" if unit price provided, "TotalPrice" if lump sum
Respond ONLY with the JSON array, no additional text."""
try:
response = bedrock_runtime.invoke_model(
modelId=MODEL_ID,
contentType="application/json",
accept="application/json",
body=json.dumps({
"anthropic_version": "bedrock-2023-05-31",
"max_tokens": 4096,
"messages": [{"role": "user", "content": prompt}],
"temperature": 0.3,
}),
)
response_body = json.loads(response["body"].read())
content = response_body["content"][0]["text"]
content = content.strip()
if content.startswith("```"):
content = content.split("\n", 1)[1]
content = content.rsplit("```", 1)[0]
line_items = json.loads(content)
return line_items if isinstance(line_items, list) else []
except Exception as e:
print(f"Error generating line items: {e}")
return []
def post_line_items(proposal_id: str, items: list[dict], existing_items: list[dict]):
if not items and not existing_items:
return
line_items_payload = []
# Preserve non-AI items (Manual, Vendor, Historical)
preserved = [li for li in existing_items if li.get("source") != "AI"]
for i, li in enumerate(preserved):
line_items_payload.append({
"id": li.get("id"),
"description": li["description"],
"quantity": float(li.get("quantity", 1)),
"unit": li.get("unit", "each"),
"unitPrice": li.get("unitPrice"),
"totalPrice": float(li.get("totalPrice", 0)),
"pricingMode": li.get("pricingMode", "TotalPrice"),
"sortOrder": i + 1,
"source": li.get("source", "Manual"),
})
# Add new AI-generated items after preserved ones
offset = len(line_items_payload)
for i, item in enumerate(items):
pricing_mode = item.get("pricingMode", "TotalPrice")
if pricing_mode not in ("UnitPrice", "TotalPrice", "Both"):
pricing_mode = "UnitPrice" if item.get("unitPrice") else "TotalPrice"
line_items_payload.append({
"id": None,
"description": item["description"],
"quantity": float(item.get("quantity", 1)),
"unit": item.get("unit", "each"),
"unitPrice": item.get("unitPrice"),
"totalPrice": float(item.get("totalPrice", 0)),
"pricingMode": pricing_mode,
"sortOrder": offset + i + 1,
"source": "AI",
})
try:
resp = httpx.put(
f"{API_BASE_URL}/api/proposals/{proposal_id}/line-items",
json={"lineItems": line_items_payload},
headers=_api_headers(),
timeout=15,
)
if resp.status_code not in (200, 201):
print(f"Failed to post line items: {resp.status_code} {resp.text}")
except Exception as e:
print(f"Error posting line items: {e}")
def store_similar_references(proposal_id: str, similar_proposals: list[dict]):
if not similar_proposals:
return
for sp in similar_proposals[:5]:
source_uri = sp.get("sourceUri", "")
library_item_id = source_uri.split("/")[-1] if source_uri else ""
if not library_item_id:
continue
try:
httpx.post(
f"{API_BASE_URL}/api/proposals/{proposal_id}/similar-references",
json={
"referencedLibraryItemId": library_item_id,
"similarityScore": sp["score"],
},
headers=_api_headers(),
timeout=10,
)
except Exception as e:
print(f"Error storing similar reference: {e}")
def update_status_to_in_review(proposal_id: str):
try:
httpx.put(
f"{API_BASE_URL}/api/proposals/{proposal_id}",
json={"status": "InReview"},
headers=_api_headers(),
timeout=10,
)
except Exception as e:
print(f"Error updating status: {e}")
def _api_headers() -> dict:
headers = {"Content-Type": "application/json"}
api_key = _get_api_key()
if api_key:
headers["X-Internal-Api-Key"] = api_key
return headers

View file

@ -0,0 +1,2 @@
boto3>=1.35.0,<2.0
httpx>=0.27.0,<1.0