Production-grade AI lead-intelligence pipeline for competitor dissatisfaction monitoring, scoring, and outbound assistance.
- Overview
- Architecture
- Features
- Repository Structure
- Getting Started
- Database Migrations
- Running
- Docker
- Testing
- CI/CD
- Deployment (Vercel + Cron)
- Troubleshooting
- Scaling Recommendations
A daily pipeline that collects competitor dissatisfaction signals from web sources, enriches them with OpenAI, deduplicates, scores, and ranks leads, then publishes a Slack digest and Excel reports for human review.
Daily flow:
- Scheduler (
Cron, GitHub Actions, or APScheduler) - Source collectors: OpenAI Web Research when
OPENAI_API_KEYis set (on by default; disable withENABLE_OPENAI_WEB_RESEARCH=false), optional OpenAI URL scrape, then Reddit and Google Alerts RSS only when those credentials or real feed URLs are configured; optional demo placeholders ifENABLE_PLACEHOLDER_SOURCES=true - Normalization + sanitization
- Deduplication (URL + content hash + fuzzy match window)
- OpenAI enrichment (second pass: structured classification, scores, suggested reply —
OpenAIEnricher) - Lead ranking and persistence (
PostgreSQL + SQLAlchemy) - Slack digest (
Block Kit, grouped by intent and competitor) - Human review and manual posting
- Multi-source collection — OpenAI web discovery, URL scraping, Reddit, Google Alerts RSS
- AI enrichment — Structured classification, intent scoring, and suggested reply generation via OpenAI
- Smart deduplication — URL + content hash + fuzzy match window
- Lead ranking — Scoring and persistence to PostgreSQL
- Slack digest — Block Kit formatted, grouped by intent and competitor
- Excel reports — Per-run and master reports with workflow tracking columns
- API — FastAPI endpoints for lead browsing, stats, reclassification, and reply generation
- Health checks —
/healthliveness and/readyPostgreSQL readiness probes - Non-root Docker — Container runs as
appuserwith built-in HEALTHCHECK
| Directory | Description |
|---|---|
app/api |
FastAPI routes (/leads, /leads/{id}, /stats, /reclassify, /generate-reply) |
app/collectors |
Source adapters |
app/classifiers |
OpenAI enrichment logic |
app/dedupe |
Duplicate detection engine |
app/ranking |
Lead scoring |
app/storage |
Repository layer |
app/slack |
Block Kit digest publishing |
app/scheduler |
Ingestion runners |
app/db + app/models |
SQLAlchemy models/session |
alembic |
Migration config and versions |
app/tests |
Unit/integration-ready tests |
python -m venv .venv
.venv\Scripts\activate # Windows
# source .venv/bin/activate # Linux/Mac
pip install -r requirements.txt
copy .env.example .env # Windows
# cp .env.example .env # Linux/MacMinimum to run ingestion: DATABASE_URL, OPENAI_API_KEY, and KEYWORDS. With only those, the pipeline uses OpenAI web discovery (no Reddit or Google setup required).
Also set SLACK_WEBHOOK_URL if you want the daily digest.
| Variable | Description |
|---|---|
ENABLE_OPENAI_WEB_RESEARCH |
Set to false to turn off web search discovery (saves API cost) |
OPENAI_SCRAPER_URLS |
Comma-separated https pages; OpenAI extracts lead-like snippets |
OPENAI_RESPONSES_MODEL |
Tune the enrichment model |
OPENAI_COLLECTION_MODEL |
Tune the collection model |
OPENAI_COLLECTION_MIN_RELEVANCE |
Relevance cutoff for collection |
REDDIT_CLIENT_ID / REDDIT_CLIENT_SECRET |
Adds Reddit collection when both are set |
GOOGLE_ALERT_RSS_URLS |
Comma-separated real Google Alert RSS feed URLs |
ENABLE_PLACEHOLDER_SOURCES |
Set to true for demo-only placeholder sources |
SENTRY_DSN |
Sentry error tracking |
SYNC_DATABASE_URL |
Alembic sync URL (often same as DATABASE_URL with sync driver) |
alembic upgrade headuvicorn app.main:app --reload --host 0.0.0.0 --port 8000make run
# or: python -m app.scheduler.runnerThis run also:
- Creates a per-run Excel report in
reports/new_leads_YYYY-MM-DD.xlsx - Refreshes a master Excel report in
reports/all_leads_master.xlsx(historical leads) - Writes report rows to Postgres table
lead_reportswith columns:Date | Source | Link | Competitor | Pain Point | Intent | Suggested Reply | Status - Both Excel reports include visibility columns:
Run Date | Captured At | ... | Status | Reviewed By | Reviewed At | Posted At - Posts daily digest to Slack webhook
- Optionally uploads both Excel files to Slack when
SLACK_BOT_TOKENandSLACK_CHANNEL_IDare set
python -m app.scheduler.apscheduler_jobdocker compose up --buildpytest -q --cov=app --cov-report=term-missingGitHub Actions workflows:
.github/workflows/ci.yml— Runs tests on push and pull requests tomain/master.github/workflows/daily-ingestion.yml— Runsalembic upgrade headthenpython -m app.scheduler.runneron a schedule (0 14 * * *UTC = 7:30 PM IST) and on manualworkflow_dispatch
| Secret | Required |
|---|---|
DATABASE_URL |
Yes (publicly reachable Postgres) |
OPENAI_API_KEY |
Yes |
KEYWORDS |
Yes |
SLACK_WEBHOOK_URL |
For Slack digest |
OPENAI_MODEL |
If using non-default model |
Optional: GOOGLE_ALERT_RSS_URLS, REDDIT_CLIENT_ID, REDDIT_CLIENT_SECRET, SYNC_DATABASE_URL, and others from .env.example.
- Enable Actions: Repo → Settings → Actions → General → allow actions
- Add secrets: Settings → Secrets and variables → Actions → New repository secret
- Confirm workflow: Actions → Daily lead ingestion (push any commit to
mainif tab is empty) - Test: Daily lead ingestion → Run workflow → branch
main→ Run workflow
The Run lead monitor step sets APP_ENV=production so ingestion validates OPENAI_API_KEY and KEYWORDS before work starts. Optional STRICT_STARTUP_VALIDATION=true tightens API boot and ingestion.
Health: GET /health liveness; GET /ready checks PostgreSQL (SELECT 1).
- Create a Vercel project from this repo. The FastAPI entrypoint is set in
pyproject.tomlasapp.main:app. - In Vercel Environment variables, add the same values as local (
DATABASE_URL,OPENAI_API_KEY,KEYWORDS, Slack, Reddit, etc.). SetCRON_SECRETto a long random string. vercel.jsonschedulesGET /api/cron/daily-ingestiondaily at 14:00 UTC (7:30 PM IST). That route runs the full ingestion pipeline.- Limits: Vercel Functions have a maximum duration (short on Hobby). For reliable daily jobs, prefer GitHub Actions or a VM/Docker scheduler. You can host the read API on Vercel and run the heavy job on GitHub.
ModuleNotFoundError: ensure virtual environment is active- DB connection errors: verify
DATABASE_URLand Postgres availability Name or service not knownwithneon. tech: connection string has a space in the hostname — fix the GitHub secret and local.env- Empty digests: verify
KEYWORDSandOPENAI_API_KEY; add Reddit or Google Alert feeds only if needed - OpenAI validation failures: inspect logs for malformed model output and retry behavior
- Add Celery/RQ worker pool for enrichment fan-out at high volume
- Batch API calls and include back-pressure controls per collector
- Add Redis cache for short-term dedupe candidate sets
- Add partitioning/index tuning as daily volume grows
- Extend to multi-tenant schema with org-level source and model configs