diff --git a/lambdas/library-ingest/app.py b/lambdas/library-ingest/app.py index 6c4eaee..9c903e1 100644 --- a/lambdas/library-ingest/app.py +++ b/lambdas/library-ingest/app.py @@ -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 diff --git a/lambdas/library-ingest/batch_ingest.py b/lambdas/library-ingest/batch_ingest.py new file mode 100644 index 0000000..b50a462 --- /dev/null +++ b/lambdas/library-ingest/batch_ingest.py @@ -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 --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() diff --git a/lambdas/library-ingest/requirements.txt b/lambdas/library-ingest/requirements.txt index 97c4921..996956a 100644 --- a/lambdas/library-ingest/requirements.txt +++ b/lambdas/library-ingest/requirements.txt @@ -1 +1,2 @@ boto3>=1.35.0,<2.0 +httpx>=0.27.0,<1.0 diff --git a/lambdas/pdf-extract/app.py b/lambdas/pdf-extract/app.py index 42430aa..77728cf 100644 --- a/lambdas/pdf-extract/app.py +++ b/lambdas/pdf-extract/app.py @@ -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 diff --git a/lambdas/pdf-extract/requirements.txt b/lambdas/pdf-extract/requirements.txt index 1d64c16..8545875 100644 --- a/lambdas/pdf-extract/requirements.txt +++ b/lambdas/pdf-extract/requirements.txt @@ -1,2 +1,3 @@ pdfplumber>=0.11.0,<1.0 boto3>=1.35.0,<2.0 +httpx>=0.27.0,<1.0 diff --git a/lambdas/pdf-generate/app.py b/lambdas/pdf-generate/app.py index b59ee6f..6d421f5 100644 --- a/lambdas/pdf-generate/app.py +++ b/lambdas/pdf-generate/app.py @@ -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 diff --git a/lambdas/pdf-generate/requirements.txt b/lambdas/pdf-generate/requirements.txt index 64cd117..82cc9af 100644 --- a/lambdas/pdf-generate/requirements.txt +++ b/lambdas/pdf-generate/requirements.txt @@ -1,2 +1,3 @@ reportlab>=4.2.0,<5.0 boto3>=1.35.0,<2.0 +httpx>=0.27.0,<1.0 diff --git a/lambdas/suggestions/app.py b/lambdas/suggestions/app.py new file mode 100644 index 0000000..aff8d49 --- /dev/null +++ b/lambdas/suggestions/app.py @@ -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 diff --git a/lambdas/suggestions/requirements.txt b/lambdas/suggestions/requirements.txt new file mode 100644 index 0000000..996956a --- /dev/null +++ b/lambdas/suggestions/requirements.txt @@ -0,0 +1,2 @@ +boto3>=1.35.0,<2.0 +httpx>=0.27.0,<1.0