proposal-system/lambdas/library-ingest/batch_ingest.py
Adam Moussa 9d6612e337 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
2026-05-16 22:10:21 -04:00

359 lines
12 KiB
Python

"""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()