"""Physical mail queue + LetterStream send integration. Letters are the paper collection notices mailed to debtors. One row per mailed letter in `letters` (lifecycle DRAFT -> APPROVED -> PREAUTH -> SENT, plus REJECTED / CANCELLED / ERROR), with tracking scan events in `letter_events` pushed by LetterStream's callback. Send flow is human-in-the-loop on cost (production mode, real money): 1. staff creates a draft letter (content defaults to claims.letter_*, recipient address supplied by staff after client-data + Super Search verification) 2. staff approves -> APPROVED 3. staff clicks "price & queue" -> LetterStream preauth=1, store authcode + cost (status PREAUTH); nothing is released or mailed yet 4. staff confirms -> LetterStream doauth, status SENT, tracking begins Signatory: Debt Recovery Experts LLC (LETTERSTREAM_SIGNATORY). Return address: LETTERSTREAM_RETURN_ADDRESS (colon-delimited addr1:addr2:city:state:zip). """ from __future__ import annotations import json import logging import os import uuid from fastapi import APIRouter, Depends, Request, status from fastapi.responses import JSONResponse from pydantic import BaseModel, ConfigDict, Field, ValidationError, field_validator from . import auth as authmod from . import letterstream from .db import get_conn, new_uuid, utcnow_iso logger = logging.getLogger("dre.letters") router = APIRouter() SIGNATORY = os.environ.get("LETTERSTREAM_SIGNATORY", "Debt Recovery Experts LLC") RETURN_ADDR = os.environ.get("LETTERSTREAM_RETURN_ADDRESS", "") # addr1:addr2:city:state:zip CALLBACK_KEY = os.environ.get("LETTERSTREAM_CALLBACK_KEY", "") LETTERS_DIR = os.environ.get("DRE_LETTERS_DIR", "/opt/dre-portal/data/letters") MAILTYPES = ("firstclass", "firstclass_hse", "certified", "certnoerr", "postcard", "flat", "propostcard") LETTER_STATUSES = ("DRAFT", "APPROVED", "PREAUTH", "SENT", "REJECTED", "CANCELLED", "ERROR") class LetterStreamError(Exception): """Raised on any LetterStream or letter-lifecycle failure.""" # --------------------------------------------------------------- # Request models (letter endpoints only — kept local to this router) # --------------------------------------------------------------- class RecipientAddress(BaseModel): model_config = ConfigDict(extra="forbid") name: str = Field(..., min_length=1, max_length=200) name2: str | None = Field(None, max_length=200) addr1: str = Field(..., min_length=1, max_length=200) addr2: str | None = Field(None, max_length=200) city: str = Field(..., min_length=1, max_length=100) state: str = Field(..., min_length=2, max_length=2) zip: str = Field(..., min_length=5, max_length=10) @field_validator("name", "name2", "addr1", "addr2", "city", "state", "zip") @classmethod def _v(cls, v): return v class LetterCreate(BaseModel): model_config = ConfigDict(extra="forbid") recipient: RecipientAddress letter_type: str = Field("demand", max_length=40) subject: str | None = Field(None, max_length=200) body: str | None = Field(None, max_length=20000) mailtype: str = "firstclass" coversheet: bool = True @field_validator("mailtype") @classmethod def _mt(cls, v): if v not in MAILTYPES: raise ValueError(f"mailtype must be one of {MAILTYPES}") return v class LetterAction(BaseModel): model_config = ConfigDict(extra="forbid") reason: str | None = Field(None, max_length=1000) # --------------------------------------------------------------- # Helpers # --------------------------------------------------------------- def sender_parts() -> list[str]: """Return [name1, name2, addr1, addr2, city, state, zip] for the `from` field.""" if not RETURN_ADDR: raise LetterStreamError("LETTERSTREAM_RETURN_ADDRESS is not configured") parts = RETURN_ADDR.split(":") if len(parts) != 5: raise LetterStreamError( "LETTERSTREAM_RETURN_ADDRESS must be addr1:addr2:city:state:zip") addr1, addr2, city, state, zipc = parts return [SIGNATORY, "", addr1, addr2, city, state.upper(), zipc] def sender_string() -> str: return ":".join(sender_parts()) def recipient_string(doc_id: str, r: dict) -> str: return ":".join([ doc_id, r.get("name", ""), r.get("name2") or "", r.get("addr1", ""), r.get("addr2") or "", r.get("city", ""), r.get("state", "").upper(), r.get("zip", ""), ]) def render_letter_pdf(*, sender_block: list[str], recipient_block: list[str], date_str: str, subject: str, body: str) -> tuple[bytes, int]: """Render a single-page (or multi-page) letter PDF. Returns (bytes, page_count). Address blocks are placed in the standard #10 double-window positions (verified via LetterStream preflight before first production send). """ from fpdf import FPDF def _safe(s: str) -> str: return (s or "").encode("latin-1", "replace").decode("latin-1") pdf = FPDF(unit="mm", format="Letter") pdf.set_auto_page_break(auto=True, margin=15) pdf.add_page() # Return-address window (top-left) pdf.set_font("Helvetica", size=10) y = 14.0 for line in sender_block: pdf.set_xy(14.3, y) pdf.cell(0, 4.2, _safe(line)) y += 4.2 # Recipient window y = 50.8 for line in recipient_block: pdf.set_xy(14.3, y) pdf.cell(0, 4.2, _safe(line)) y += 4.2 # Date pdf.set_xy(14.3, 76.0) pdf.cell(0, 5, _safe(date_str)) # Subject (bold) pdf.set_font("Helvetica", "B", size=12) pdf.set_xy(14.3, 84.0) pdf.multi_cell(180, 6, _safe(subject)) # Body pdf.set_font("Helvetica", size=10) pdf.set_xy(14.3, 94.0) pdf.multi_cell(180, 4.6, _safe(body)) out = pdf.output(dest="S") return out, len(pdf.pages) def _sender_block() -> list[str]: parts = sender_parts() name1, _name2, addr1, addr2, city, state, zipc = parts lines = [name1] if addr1: lines.append(addr1) if addr2: lines.append(addr2) lines.append(f"{city}, {state} {zipc}") return lines def _recipient_block(r: dict) -> list[str]: lines = [r.get("name", "")] if r.get("name2"): lines.append(r["name2"]) lines.append(r.get("addr1", "")) if r.get("addr2"): lines.append(r["addr2"]) lines.append(f"{r.get('city', '')}, {r.get('state', '')} {r.get('zip', '')}") return lines def _letter_row(conn, letter_id: str): return conn.execute( "SELECT l.*, c.claim_number FROM letters l " "JOIN claims c ON c.id = l.claim_id WHERE l.id = ?", (letter_id,), ).fetchone() def _serialize(row) -> dict: return { "id": row["id"], "claim_id": row["claim_id"], "claim_number": row["claim_number"], "letter_type": row["letter_type"], "subject": row["subject"], "body": row["body"], "recipient": json.loads(row["recipient_json"]), "sender": json.loads(row["sender_json"]), "mailtype": row["mailtype"], "status": row["status"], "job_id": row["job_id"], "batch_id": row["batch_id"], "doc_id": row["doc_id"], "tracking_no": row["tracking_no"], "cost_cents": row["cost_cents"], "pages": row["pages"], "error": row["error"], "note": row["note"], "created_at": row["created_at"], "updated_at": row["updated_at"], "sent_at": row["sent_at"], } # --------------------------------------------------------------- # POST /api/staff/claims/{claim_number}/letters (create draft) # --------------------------------------------------------------- @router.post("/api/staff/claims/{claim_number}/letters") async def create_letter(claim_number: str, request: Request, staff_name: str = Depends(authmod.require_staff)): try: body = await request.json() except Exception: return _err("validation_error", "Invalid JSON body.", status.HTTP_422_UNPROCESSABLE_ENTITY) try: lc = LetterCreate.model_validate(body) except ValidationError as exc: parts = [f"{'.'.join(str(x) for x in e['loc'])}: {e['msg']}" for e in exc.errors()] return _err("validation_error", "; ".join(parts), status.HTTP_422_UNPROCESSABLE_ENTITY) with get_conn() as conn: row = conn.execute( "SELECT c.id, c.letter_subject, c.letter_body, d.name AS debtor_name " "FROM claims c JOIN debtors d ON d.id = c.debtor_id WHERE c.claim_number = ?", (claim_number,), ).fetchone() if row is None: return _err("not_found", "Claim not found.", status.HTTP_404_NOT_FOUND) subject = lc.subject or row["letter_subject"] or "Notice from Debt Recovery Experts, LLC" body_text = lc.body or row["letter_body"] or "" if not body_text.strip(): return _err("validation_error", "Claim has no letter body; generate a letter first.", status.HTTP_422_UNPROCESSABLE_ENTITY) recipient = lc.recipient.model_dump() # Sender address is global config (LETTERSTREAM_RETURN_ADDRESS), resolved at # send time. A draft records only the signatory; drafting must not block on # the return address, which may not be configured yet. sender = {"name": SIGNATORY, "name2": "", "addr1": "", "addr2": "", "city": "", "state": "", "zip": ""} now = utcnow_iso() letter_id = new_uuid() conn.execute( "INSERT INTO letters (id, claim_id, letter_type, subject, body, recipient_json, " "sender_json, mailtype, status, created_at, updated_at) " "VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'DRAFT', ?, ?)", (letter_id, row["id"], lc.letter_type, subject, body_text, json.dumps(recipient), json.dumps(sender), lc.mailtype, now, now), ) conn.execute( "INSERT INTO audit_log (id, entity_type, entity_id, action, field, old_value, new_value, actor, created_at) " "VALUES (?, 'letter', ?, 'create', 'status', NULL, 'DRAFT', ?, ?)", (new_uuid(), letter_id, staff_name, now), ) conn.commit() return await _get_letter(letter_id) # --------------------------------------------------------------- # GET /api/staff/letters # --------------------------------------------------------------- @router.get("/api/staff/letters") async def list_letters(request: Request, _staff=Depends(authmod.require_staff)): status_filter = request.query_params.get("status") limit = min(int(request.query_params.get("limit", "50")), 200) offset = max(int(request.query_params.get("offset", "0")), 0) where = "" params: list = [] if status_filter: where = "WHERE l.status = ?" params.append(status_filter) from_clause = "FROM letters l JOIN claims c ON c.id = l.claim_id" with get_conn() as conn: rows = conn.execute( f"SELECT l.*, c.claim_number {from_clause} {where} " f"ORDER BY l.created_at DESC LIMIT ? OFFSET ?", tuple(params + [limit, offset]), ).fetchall() total = conn.execute( f"SELECT COUNT(*) AS n {from_clause} {where}", tuple(params), ).fetchone()["n"] return {"letters": [_serialize(r) for r in rows], "total": total, "limit": limit, "offset": offset} # --------------------------------------------------------------- # GET /api/staff/letters/{letter_id} # --------------------------------------------------------------- @router.get("/api/staff/letters/{letter_id}") async def get_letter(letter_id: str, _staff=Depends(authmod.require_staff)): return await _get_letter(letter_id) async def _get_letter(letter_id: str): with get_conn() as conn: row = _letter_row(conn, letter_id) if row is None: return _err("not_found", "Letter not found.", status.HTTP_404_NOT_FOUND) events = conn.execute( "SELECT scan_code, scan_status, scan_date, scan_zip, scan_facility, tracking_id, " "created_at FROM letter_events WHERE letter_id = ? ORDER BY created_at ASC", (letter_id,), ).fetchall() data = _serialize(row) data["events"] = [dict(e) for e in events] return data # --------------------------------------------------------------- # POST /api/staff/letters/{letter_id}/approve (DRAFT -> APPROVED) # --------------------------------------------------------------- @router.post("/api/staff/letters/{letter_id}/approve") async def approve_letter(letter_id: str, staff_name: str = Depends(authmod.require_staff)): with get_conn() as conn: row = _letter_row(conn, letter_id) if row is None: return _err("not_found", "Letter not found.", status.HTTP_404_NOT_FOUND) if row["status"] != "DRAFT": return _err("conflict", f"Letter is {row['status']}, not DRAFT.", status.HTTP_409_CONFLICT) now = utcnow_iso() conn.execute("UPDATE letters SET status = 'APPROVED', updated_at = ? WHERE id = ?", (now, letter_id)) conn.execute( "INSERT INTO audit_log (id, entity_type, entity_id, action, field, old_value, new_value, actor, created_at) " "VALUES (?, 'letter', ?, 'approve', 'status', 'DRAFT', 'APPROVED', ?, ?)", (new_uuid(), letter_id, staff_name, now), ) conn.commit() return await _get_letter(letter_id) # --------------------------------------------------------------- # POST /api/staff/letters/{letter_id}/send (preauth — price only) # --------------------------------------------------------------- @router.post("/api/staff/letters/{letter_id}/send") async def send_letter(letter_id: str, staff_name: str = Depends(authmod.require_staff)): with get_conn() as conn: row = _letter_row(conn, letter_id) if row is None: return _err("not_found", "Letter not found.", status.HTTP_404_NOT_FOUND) if row["status"] not in ("APPROVED", "PREAUTH", "ERROR"): return _err("conflict", f"Letter is {row['status']}; approve it first.", status.HTTP_409_CONFLICT) recipient = json.loads(row["recipient_json"]) subject = row["subject"] or "" body = row["body"] or "" try: sender = {k: v for k, v in zip( ("name", "name2", "addr1", "addr2", "city", "state", "zip"), sender_parts())} except LetterStreamError as exc: return _err("config_error", str(exc), status.HTTP_409_CONFLICT) # Render + persist PDF doc_id = new_uuid().replace("-", "")[:16] job = "DRE" + doc_id try: pdf_bytes, pages = render_letter_pdf( sender_block=_sender_block(), recipient_block=_recipient_block(recipient), date_str=utcnow_iso()[:10], subject=subject, body=body, ) except Exception as exc: return _err("render_error", f"PDF render failed: {exc}", status.HTTP_500_INTERNAL_SERVER_ERROR) os.makedirs(LETTERS_DIR, exist_ok=True) pdf_path = os.path.join(LETTERS_DIR, f"{letter_id}.pdf") with open(pdf_path, "wb") as fh: fh.write(pdf_bytes) try: resp = letterstream.send_single( pdf_bytes, f"{letter_id}.pdf", job, sender_string(), [recipient_string(doc_id, recipient)], pages, mailtype=row["mailtype"], preauth=True, ) except letterstream.LetterStreamError as exc: now = utcnow_iso() conn.execute( "UPDATE letters SET status = 'ERROR', error = ?, updated_at = ? WHERE id = ?", (str(exc), now, letter_id), ) conn.commit() return _err("letterstream_error", str(exc), status.HTTP_502_BAD_GATEWAY) m = letterstream._first(resp) if resp.get("messages") else {} authcode = m.get("authcode", "") cost = m.get("cost") batch = m.get("batch") docs = m.get("docs", []) ls_doc = docs[0].get("id") if docs else None cost_cents = int(round(float(cost) * 100)) if cost else None now = utcnow_iso() conn.execute( "UPDATE letters SET status = 'PREAUTH', job_id = ?, batch_id = ?, doc_id = ?, " "sender_json = ?, authcode = ?, cost_cents = ?, pages = ?, pdf_path = ?, error = NULL, updated_at = ? " "WHERE id = ?", (job, batch, ls_doc, json.dumps(sender), authcode, cost_cents, pages, pdf_path, now, letter_id), ) conn.execute( "INSERT INTO audit_log (id, entity_type, entity_id, action, field, old_value, new_value, actor, created_at) " "VALUES (?, 'letter', ?, 'preauth', 'status', NULL, 'PREAUTH', ?, ?)", (new_uuid(), letter_id, staff_name, now), ) conn.commit() return await _get_letter(letter_id) # --------------------------------------------------------------- # POST /api/staff/letters/{letter_id}/confirm (doauth — release) # --------------------------------------------------------------- @router.post("/api/staff/letters/{letter_id}/confirm") async def confirm_letter(letter_id: str, staff_name: str = Depends(authmod.require_staff)): with get_conn() as conn: row = _letter_row(conn, letter_id) if row is None: return _err("not_found", "Letter not found.", status.HTTP_404_NOT_FOUND) if row["status"] != "PREAUTH" or not row["authcode"]: return _err("conflict", "Letter has no pending preauth to confirm.", status.HTTP_409_CONFLICT) try: letterstream.doauth(row["authcode"]) except letterstream.LetterStreamError as exc: now = utcnow_iso() conn.execute( "UPDATE letters SET status = 'ERROR', error = ?, updated_at = ? WHERE id = ?", (str(exc), now, letter_id), ) conn.commit() return _err("letterstream_error", str(exc), status.HTTP_502_BAD_GATEWAY) now = utcnow_iso() conn.execute( "UPDATE letters SET status = 'SENT', sent_at = ?, updated_at = ? WHERE id = ?", (now, now, letter_id), ) conn.execute( "INSERT INTO audit_log (id, entity_type, entity_id, action, field, old_value, new_value, actor, created_at) " "VALUES (?, 'letter', ?, 'confirm', 'status', 'PREAUTH', 'SENT', ?, ?)", (new_uuid(), letter_id, staff_name, now), ) conn.commit() return await _get_letter(letter_id) # --------------------------------------------------------------- # POST /api/staff/letters/{letter_id}/reject (terminal — REJECTED) # --------------------------------------------------------------- @router.post("/api/staff/letters/{letter_id}/reject") async def reject_letter(letter_id: str, request: Request, staff_name: str = Depends(authmod.require_staff)): reason = "" try: body = await request.json() if isinstance(body, dict): reason = str(body.get("reason") or "").strip() except Exception: pass if not reason: return _err("validation_error", "A reason is required to reject a letter.", status.HTTP_422_UNPROCESSABLE_ENTITY) reason = reason[:1000] with get_conn() as conn: row = _letter_row(conn, letter_id) if row is None: return _err("not_found", "Letter not found.", status.HTTP_404_NOT_FOUND) if row["status"] in ("SENT", "REJECTED", "CANCELLED"): return _err("conflict", f"Letter is {row['status']}; cannot reject.", status.HTTP_409_CONFLICT) old = row["status"] now = utcnow_iso() conn.execute( "UPDATE letters SET status = 'REJECTED', note = ?, error = NULL, updated_at = ? WHERE id = ?", (reason, now, letter_id), ) conn.execute( "INSERT INTO audit_log (id, entity_type, entity_id, action, field, old_value, new_value, actor, created_at) " "VALUES (?, 'letter', ?, 'reject', 'status', ?, 'REJECTED', ?, ?)", (new_uuid(), letter_id, old, staff_name, now), ) conn.commit() return await _get_letter(letter_id) # --------------------------------------------------------------- # POST /api/staff/letters/{letter_id}/cancel (terminal — CANCELLED) # --------------------------------------------------------------- @router.post("/api/staff/letters/{letter_id}/cancel") async def cancel_letter(letter_id: str, request: Request, staff_name: str = Depends(authmod.require_staff)): reason = "" try: body = await request.json() if isinstance(body, dict): reason = str(body.get("reason") or "").strip() except Exception: pass reason = reason[:1000] with get_conn() as conn: row = _letter_row(conn, letter_id) if row is None: return _err("not_found", "Letter not found.", status.HTTP_404_NOT_FOUND) if row["status"] in ("SENT", "REJECTED", "CANCELLED"): return _err("conflict", f"Letter is {row['status']}; cannot cancel.", status.HTTP_409_CONFLICT) old = row["status"] now = utcnow_iso() conn.execute( "UPDATE letters SET status = 'CANCELLED', note = ?, error = NULL, updated_at = ? WHERE id = ?", (reason, now, letter_id), ) conn.execute( "INSERT INTO audit_log (id, entity_type, entity_id, action, field, old_value, new_value, actor, created_at) " "VALUES (?, 'letter', ?, 'cancel', 'status', ?, 'CANCELLED', ?, ?)", (new_uuid(), letter_id, old, staff_name, now), ) conn.commit() return await _get_letter(letter_id) # --------------------------------------------------------------- # POST /api/letters/webhook (PUBLIC — LetterStream tracking callback) # --------------------------------------------------------------- @router.post("/api/letters/webhook") async def letter_webhook(request: Request): ctype = request.headers.get("content-type", "") data: dict = {} if "application/json" in ctype: try: data = await request.json() except Exception: data = {} else: try: form = await request.form() data = {k: v for k, v in form.items()} except Exception: data = {} key = data.get("key", "") if not CALLBACK_KEY or key != CALLBACK_KEY: return JSONResponse(status_code=401, content={"success": False, "reason": "Invalid key"}) raw_json = data.get("json", "") payload = None try: payload = json.loads(raw_json) if isinstance(raw_json, str) else raw_json except (ValueError, TypeError): payload = None events: list[dict] = [] if isinstance(payload, list): events = [e for e in payload if isinstance(e, dict)] elif isinstance(payload, dict): if isinstance(payload.get("data"), list): events = [e for e in payload["data"] if isinstance(e, dict)] elif payload: events = [payload] with get_conn() as conn: for ev in events: doc_id = str(ev.get("doc_id", "")) job_id = str(ev.get("job_id", "")) letter = None if doc_id: letter = conn.execute("SELECT id FROM letters WHERE doc_id = ?", (doc_id,)).fetchone() if letter is None and job_id: letter = conn.execute("SELECT id FROM letters WHERE job_id = ?", (job_id,)).fetchone() if letter is None: logger.warning("letter webhook: no letter for doc_id=%s job_id=%s", doc_id, job_id) continue conn.execute( "INSERT INTO letter_events (id, letter_id, scan_code, scan_status, scan_date, " "scan_zip, scan_facility, tracking_id, batch_id, job_id, doc_id, raw_json, created_at) " "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", (new_uuid(), letter["id"], ev.get("scan_code"), ev.get("scan_status"), ev.get("scan_date"), ev.get("scan_zip"), ev.get("scan_facility"), ev.get("tracking_id"), ev.get("batch_id"), job_id, doc_id, json.dumps(ev), utcnow_iso()), ) if ev.get("tracking_id"): conn.execute( "UPDATE letters SET tracking_no = COALESCE(tracking_no, ?), updated_at = ? WHERE id = ?", (str(ev["tracking_id"]), utcnow_iso(), letter["id"]), ) conn.commit() return {"success": True, "reason": "Received data"} def _err(code: str, message: str, status_code: int): return JSONResponse(status_code=status_code, content={"error": {"code": code, "message": message}})