-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmain.py
More file actions
160 lines (131 loc) · 5.4 KB
/
Copy pathmain.py
File metadata and controls
160 lines (131 loc) · 5.4 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
"""
Lead Processing API — main entry point.
Endpoints:
POST /webhook/lead — accept a lead form submission
GET /health — liveness check
GET / — basic info page
"""
import asyncio
import logging
import os
from contextlib import asynccontextmanager
from dotenv import load_dotenv
from fastapi import FastAPI, Request, status
from fastapi.responses import JSONResponse, HTMLResponse
# Кожен модуль читає свої ключі через os.getenv на момент виклику, тож .env
# треба підвантажити тут, до першого запиту. На платформах змінні задані в
# оточенні — там load_dotenv просто не знаходить файлу і нічого не робить.
load_dotenv()
from ai_service import analyze_lead
from models import LeadRequest, ProcessedLead
from airtable import append_lead_to_airtable
from telegram import send_telegram_notification
# ---------------------------------------------------------------------------
# Logging
# ---------------------------------------------------------------------------
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)-8s %(name)s — %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# App
# ---------------------------------------------------------------------------
@asynccontextmanager
async def lifespan(app: FastAPI):
logger.info("🚀 Lead Processor starting up")
yield
logger.info("👋 Lead Processor shutting down")
app = FastAPI(
title="Lead Processor API",
description="Pipeline: receive → normalize → AI summary → classify → Airtable + Telegram",
version="1.0.0",
lifespan=lifespan,
)
# ---------------------------------------------------------------------------
# Routes
# ---------------------------------------------------------------------------
@app.get("/", response_class=HTMLResponse, include_in_schema=False)
async def root():
return """
<html><body style="font-family:sans-serif;padding:2rem">
<h2>🎯 Lead Processor API</h2>
<p>Send a <code>POST /webhook/lead</code> with a JSON body to process a lead.</p>
<p><a href="/docs">📖 Interactive docs (Swagger UI)</a></p>
</body></html>
"""
@app.get("/health")
async def health():
return {"status": "ok"}
@app.post(
"/webhook/lead",
status_code=status.HTTP_200_OK,
summary="Process an incoming lead form submission",
response_description="Processing result with lead ID and classification",
)
async def process_lead(payload: LeadRequest):
"""
Full pipeline:
1. Validate & normalize the incoming JSON (handled by Pydantic).
2. Send to Groq (llama-3.3-70b) for AI summary + classification.
3. Persist to Airtable.
4. Send Telegram notification.
5. Return a summary response.
All downstream failures (Airtable / Telegram) are non-fatal — the endpoint
always returns 200 if the lead was received and processed by AI.
"""
logger.info("New lead received: %s <%s>", payload.name, payload.email)
# Step 1: AI analysis
summary, classification, reason = await analyze_lead(payload)
logger.info("Lead classified as %s", classification)
# Step 2: Build enriched lead object
lead = ProcessedLead(
**payload.model_dump(),
ai_summary=summary,
classification=classification,
classification_reason=reason,
)
# Step 3 & 4: Airtable + Telegram run concurrently to save time
sheet_task = asyncio.create_task(append_lead_to_airtable(lead))
telegram_task = asyncio.create_task(send_telegram_notification(lead))
sheet_result, telegram_result = await asyncio.gather(
sheet_task, telegram_task, return_exceptions=True
)
# Log but don't fail on downstream errors
if isinstance(sheet_result, Exception):
logger.error("Airtable task raised: %s", sheet_result)
if isinstance(telegram_result, Exception):
logger.error("Telegram task raised: %s", telegram_result)
return {
"success": True,
"lead_id": lead.lead_id,
"received_at": lead.received_at,
"classification": lead.classification,
"ai_summary": lead.ai_summary,
"destinations": {
"airtable": bool(sheet_result and not isinstance(sheet_result, Exception)),
"telegram": bool(telegram_result and not isinstance(telegram_result, Exception)),
},
}
# ---------------------------------------------------------------------------
# Global error handler — always return JSON, never expose stack traces
# ---------------------------------------------------------------------------
@app.exception_handler(Exception)
async def global_exception_handler(request: Request, exc: Exception):
logger.exception("Unhandled exception on %s %s", request.method, request.url)
return JSONResponse(
status_code=500,
content={"success": False, "error": "Internal server error"},
)
# ---------------------------------------------------------------------------
# Dev runner
# ---------------------------------------------------------------------------
if __name__ == "__main__":
import uvicorn
uvicorn.run(
"main:app",
host="0.0.0.0",
port=int(os.getenv("PORT", 8000)),
reload=True,
)