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