Skip to content

Commit ee5aeaf

Browse files
committed
Merge remote-tracking branch 'review-pr/29' into review/integration
2 parents d98045e + 62a9212 commit ee5aeaf

8 files changed

Lines changed: 200 additions & 19 deletions

File tree

‎apps/api/routers/leads.py‎

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818
from pydantic import BaseModel
1919

2020
from apps.api.services.leadgen.db import LeadDB
21-
from apps.api.services.leadgen.models import Lead, LEAD_STATUSES
21+
from apps.api.services.leadgen.models import Lead, LEAD_STATUSES, LeadAlreadyExistsError
2222
from apps.api.core.tenancy import (
2323
WorkspaceCtx,
2424
current_workspace,
@@ -350,12 +350,23 @@ class AddLeadRequest(BaseModel):
350350

351351
@router.post("/lead")
352352
def add_lead(body: AddLeadRequest, ctx: WorkspaceCtx = Depends(require_editor)):
353+
"""Create one company/city lead; collisions return 409 without modifying it.
354+
355+
Use PUT /api/lead/{lead_id} for explicit updates, including field clearing.
356+
"""
353357
lead = Lead.from_dict(body.model_dump())
354358
lead.source = body.source
355359
lead.workspace_id = ctx.workspace_id
356360
db = ctx.lead_db()
357-
lead_id = db.upsert_lead(lead)
358-
db.close()
361+
try:
362+
lead_id = db.upsert_lead(lead, create_only=True)
363+
except LeadAlreadyExistsError as exc:
364+
raise HTTPException(status_code=409, detail={
365+
"message": "A lead already exists for this company and city. Update the existing lead instead.",
366+
"lead_id": exc.lead_id,
367+
}) from exc
368+
finally:
369+
db.close()
359370
return {"ok": True, "id": lead_id}
360371

361372

‎apps/api/services/leadgen/db.py‎

Lines changed: 21 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ def _utcnow() -> datetime:
1616
"""Return current UTC time (non-deprecated alternative to datetime.utcnow())."""
1717
return datetime.now(timezone.utc)
1818

19-
from apps.api.services.leadgen.models import Lead
19+
from apps.api.services.leadgen.models import Lead, LeadAlreadyExistsError
2020

2121
# Thread-local singleton so each thread reuses one connection
2222
_thread_local = threading.local()
@@ -319,8 +319,8 @@ def _migrate(self):
319319

320320
# ── CRUD ───────────────────────────────────────────────────────────
321321

322-
def upsert_lead(self, lead: Lead) -> int:
323-
"""Insert or update a lead. Deduplicates by (workspace_id, company, city)."""
322+
def upsert_lead(self, lead: Lead, *, create_only: bool = False) -> int:
323+
"""Insert or update by (workspace_id, company, city); create_only refuses collisions."""
324324
lead.updated_at = _utcnow().isoformat()
325325

326326
existing = self.conn.execute(
@@ -330,6 +330,8 @@ def upsert_lead(self, lead: Lead) -> int:
330330
).fetchone()
331331

332332
if existing:
333+
if create_only:
334+
raise LeadAlreadyExistsError(existing["id"])
333335
lead.id = existing["id"]
334336
fields = {k: v for k, v in lead.to_dict().items()
335337
if k != "id" and k != "created_at"}
@@ -343,10 +345,22 @@ def upsert_lead(self, lead: Lead) -> int:
343345
d.pop("id", None)
344346
cols = ", ".join(d.keys())
345347
placeholders = ", ".join("?" for _ in d)
346-
cursor = self.conn.execute(
347-
f"INSERT INTO leads ({cols}) VALUES ({placeholders})",
348-
list(d.values())
349-
)
348+
try:
349+
cursor = self.conn.execute(
350+
f"INSERT INTO leads ({cols}) VALUES ({placeholders})",
351+
list(d.values())
352+
)
353+
except sqlite3.IntegrityError:
354+
if not create_only:
355+
raise
356+
self.conn.rollback()
357+
winner = self.conn.execute(
358+
"SELECT id FROM leads WHERE workspace_id = ? AND company = ? AND city = ?",
359+
(lead.workspace_id, lead.company, lead.city),
360+
).fetchone()
361+
if winner is None:
362+
raise
363+
raise LeadAlreadyExistsError(winner["id"]) from None
350364
lead.id = cursor.lastrowid
351365

352366
self.conn.commit()

‎apps/api/services/leadgen/models.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,14 @@
77
from typing import Optional
88

99

10+
class LeadAlreadyExistsError(ValueError):
11+
"""A create-only request collided with the existing company/city lead."""
12+
13+
def __init__(self, lead_id: int):
14+
self.lead_id = lead_id
15+
super().__init__("A lead already exists for this company and city")
16+
17+
1018
@dataclass
1119
class Lead:
1220
"""A single lead (company) in the pipeline."""

‎apps/api/services/leadgen/store.py‎

Lines changed: 21 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323

2424
from apps.api.core.config import settings
2525
from apps.api.database import IS_SQLITE, SessionLocal
26-
from apps.api.services.leadgen.models import Lead
26+
from apps.api.services.leadgen.models import Lead, LeadAlreadyExistsError
2727
from apps.api.services.leadgen.orm_models import LeadRow, SignalRow, LLMUsageRow
2828

2929
logger = logging.getLogger("leadgen.store")
@@ -180,7 +180,7 @@ def _lead_payload(self, lead: Lead) -> Dict[str, Any]:
180180
return {k: v for k, v in d.items() if k in _LEAD_COLUMNS}
181181

182182
# ── CRUD ──
183-
def upsert_lead(self, lead: Lead) -> int:
183+
def upsert_lead(self, lead: Lead, *, create_only: bool = False) -> int:
184184
from apps.api.services.leadgen.db import _utcnow
185185

186186
lead.updated_at = _utcnow().isoformat()
@@ -196,15 +196,32 @@ def upsert_lead(self, lead: Lead) -> int:
196196
)
197197
payload = self._lead_payload(lead)
198198
if existing:
199+
if create_only:
200+
raise LeadAlreadyExistsError(existing.id)
199201
for k, v in payload.items():
200202
if k == "created_at":
201203
continue
202204
setattr(existing, k, v)
203205
lead.id = existing.id
204206
else:
205207
row = LeadRow(**payload)
206-
s.add(row)
207-
s.flush()
208+
if create_only:
209+
from sqlalchemy.exc import IntegrityError
210+
211+
try:
212+
with s.begin_nested():
213+
s.add(row)
214+
s.flush()
215+
except IntegrityError:
216+
winner = s.query(LeadRow).filter_by(
217+
workspace_id=self.workspace_id, company=lead.company, city=lead.city,
218+
).first()
219+
if winner is None:
220+
raise
221+
raise LeadAlreadyExistsError(winner.id) from None
222+
else:
223+
s.add(row)
224+
s.flush()
208225
lead.id = row.id
209226
return lead.id
210227

‎apps/web/src/components/add-lead-dialog.tsx‎

Lines changed: 15 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -31,16 +31,24 @@ function FormField({ label, name, required, type = "text", span }: {
3131

3232
export function AddLeadDialog({ open, onOpenChange, onAdded }: Props) {
3333
const [loading, setLoading] = useState(false)
34+
const [error, setError] = useState("")
3435

3536
const handleSubmit = async (e: React.FormEvent<HTMLFormElement>) => {
3637
e.preventDefault()
38+
const form = e.currentTarget
3739
setLoading(true)
38-
const fd = new FormData(e.currentTarget)
40+
setError("")
41+
const fd = new FormData(form)
3942
const data = Object.fromEntries(fd)
40-
await addLead(data as Record<string, string>)
41-
setLoading(false)
42-
e.currentTarget.reset()
43-
onAdded()
43+
try {
44+
await addLead(data as Record<string, string>)
45+
form.reset()
46+
onAdded()
47+
} catch (err) {
48+
setError(err instanceof Error ? err.message : "Could not add lead")
49+
} finally {
50+
setLoading(false)
51+
}
4452
}
4553

4654
return (
@@ -88,6 +96,8 @@ export function AddLeadDialog({ open, onOpenChange, onAdded }: Props) {
8896

8997
<Separator />
9098

99+
{error && <p role="alert" className="text-xs text-destructive">{error}</p>}
100+
91101
{/* Actions */}
92102
<div className="flex justify-end gap-2 pt-1">
93103
<Button

‎apps/web/src/lib/api.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -477,6 +477,11 @@ export async function addLead(data: Partial<Lead>): Promise<{ ok: boolean; id: n
477477
headers: { "Content-Type": "application/json" },
478478
body: JSON.stringify(data),
479479
})
480+
if (!res.ok) {
481+
const body: { detail?: string | { message?: string } } = await res.json().catch(() => ({}))
482+
const detail = body.detail
483+
throw new Error(typeof detail === "string" ? detail : detail?.message || `Could not add lead (${res.status})`)
484+
}
480485
return res.json()
481486
}
482487

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
import { afterEach, expect, test } from "bun:test"
2+
3+
import { addLead } from "../src/lib/api"
4+
5+
const originalFetch = globalThis.fetch
6+
7+
afterEach(() => {
8+
globalThis.fetch = originalFetch
9+
})
10+
11+
test("lead creation refuses the API collision response with its actionable message", async () => {
12+
globalThis.fetch = (async () => new Response(JSON.stringify({
13+
detail: { message: "Update the existing lead instead.", lead_id: 7 },
14+
}), { status: 409, headers: { "Content-Type": "application/json" } })) as typeof fetch
15+
16+
await expect(addLead({ company: "Fixture" })).rejects.toThrow("Update the existing lead instead.")
17+
})
18+
19+
test("a successful lead creation still returns the created identifier", async () => {
20+
globalThis.fetch = (async () => new Response(JSON.stringify({ ok: true, id: 7 }), {
21+
status: 200, headers: { "Content-Type": "application/json" },
22+
})) as typeof fetch
23+
24+
await expect(addLead({ company: "Fixture" })).resolves.toEqual({ ok: true, id: 7 })
25+
})
Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,91 @@
1+
"""Creating a company lead must never silently replace its existing contact."""
2+
3+
from concurrent.futures import ThreadPoolExecutor
4+
import threading
5+
from types import SimpleNamespace
6+
7+
import pytest
8+
from fastapi import FastAPI
9+
from fastapi.testclient import TestClient
10+
from sqlalchemy import create_engine
11+
from sqlalchemy.orm import sessionmaker
12+
13+
from apps.api.routers import leads
14+
from apps.api.services.leadgen.db import LeadDB
15+
from apps.api.services.leadgen.orm_models import LeadRow
16+
from apps.api.services.leadgen import store
17+
18+
19+
@pytest.fixture(params=["sqlite", "orm"])
20+
def client(request, tmp_path, monkeypatch):
21+
if request.param == "sqlite":
22+
path = tmp_path / "leads.db"
23+
factory = lambda: LeadDB(str(path))
24+
else:
25+
engine = create_engine(f"sqlite:///{tmp_path / 'orm.db'}", connect_args={"check_same_thread": False, "timeout": 30})
26+
LeadRow.__table__.create(engine)
27+
monkeypatch.setattr(store, "SessionLocal", sessionmaker(bind=engine))
28+
factory = lambda: store.PgLeadStore("ws-fixture")
29+
ctx = SimpleNamespace(workspace_id="ws-fixture", lead_db=factory)
30+
app = FastAPI()
31+
app.include_router(leads.router)
32+
app.dependency_overrides[leads.require_editor] = lambda: ctx
33+
app.dependency_overrides[leads.current_workspace] = lambda: ctx
34+
with TestClient(app) as http:
35+
yield http
36+
if request.param == "orm":
37+
engine.dispose()
38+
39+
40+
def test_create_collision_preserves_contact_and_omitted_enrichment(client):
41+
first = client.post("/api/lead", json={
42+
"company": "Fixture", "city": "Repro City", "contact_person": "Person One",
43+
"email": "one@example.invalid", "website": "https://example.invalid",
44+
})
45+
assert first.status_code == 200, first.text
46+
lead_id = first.json()["id"]
47+
assert client.put(f"/api/lead/{lead_id}", json={"phone": "+15555550100", "specialization": "B2B SaaS"}).status_code == 200
48+
before = client.get(f"/api/lead/{lead_id}").json()
49+
collision = client.post("/api/lead", json={
50+
"company": "Fixture", "city": "Repro City", "contact_person": "Person Two",
51+
"email": "two@example.invalid",
52+
})
53+
assert collision.status_code == 409, collision.text
54+
assert collision.json()["detail"]["lead_id"] == lead_id
55+
assert client.get(f"/api/lead/{lead_id}").json() == before
56+
# An explicit edit still supports replacement and deliberate clearing.
57+
changed = client.put(f"/api/lead/{lead_id}", json={"contact_person": "Person Two", "email": "two@example.invalid", "phone": ""})
58+
assert changed.status_code == 200, changed.text
59+
updated = client.get(f"/api/lead/{lead_id}").json()
60+
assert updated["contact_person"] == "Person Two"
61+
assert updated["email"] == "two@example.invalid" and updated["phone"] == ""
62+
assert updated["website"] == "https://example.invalid"
63+
assert updated["specialization"] == "B2B SaaS"
64+
65+
66+
def test_concurrent_create_has_one_winner_without_contact_replacement(client):
67+
# Initialize the schema before the simultaneous HTTP requests.
68+
assert client.post("/api/lead", json={"company": "Other Fixture"}).status_code == 200
69+
barrier = threading.Barrier(6)
70+
71+
def create(index):
72+
barrier.wait(timeout=10)
73+
response = client.post("/api/lead", json={
74+
"company": "Concurrent Fixture", "city": "Repro City",
75+
"contact_person": f"Person {index}", "email": f"person{index}@example.invalid",
76+
})
77+
return index, response
78+
79+
with ThreadPoolExecutor(max_workers=6) as pool:
80+
responses = list(pool.map(create, range(6)))
81+
winners = [(index, response) for index, response in responses if response.status_code == 200]
82+
assert len(winners) == 1, [(index, response.status_code, response.text) for index, response in responses]
83+
winner, response = winners[0]
84+
lead_id = response.json()["id"]
85+
for index, result in responses:
86+
if index != winner:
87+
assert result.status_code == 409, result.text
88+
assert result.json()["detail"]["lead_id"] == lead_id
89+
saved = client.get(f"/api/lead/{lead_id}").json()
90+
assert saved["contact_person"] == f"Person {winner}"
91+
assert saved["email"] == f"person{winner}@example.invalid"

0 commit comments

Comments
 (0)