Expand sales leads and customer contracts

This commit is contained in:
Magzhan Zhumabayev
2026-05-12 19:09:31 +05:00
parent ce3ba270aa
commit fbc77ba8a7
3 changed files with 459 additions and 11 deletions
+369 -11
View File
@@ -29,6 +29,8 @@ from services.shared.sales_models import (
SalesConditionUpsertIn,
SalesCounterpartyOut,
SalesCounterpartyUpsertIn,
SalesCustomerOut,
SalesCustomerUpdate,
SalesDashboardStageOut,
SalesDashboardSummaryOut,
SalesDealCloseIn,
@@ -663,10 +665,11 @@ def _publish_lead_entered_event(
)
def _lead_to_out(row: SalesLeadRow) -> SalesLeadOut:
def _lead_to_out(row: SalesLeadRow, deal: SalesDealRow | None = None) -> SalesLeadOut:
return SalesLeadOut(
lead_id=row.lead_id,
tenant_id=row.tenant_id,
deal_id=deal.deal_id if deal is not None else None,
source_type=row.source_type,
source_channel=row.source_channel,
source_campaign_id=row.source_campaign_id,
@@ -689,6 +692,18 @@ def _lead_to_out(row: SalesLeadRow) -> SalesLeadOut:
)
def _customer_to_out(row: Customer, deal_count: int = 0) -> SalesCustomerOut:
return SalesCustomerOut(
customer_id=row.customer_id,
display_name=row.display_name,
phones=_json_list(row.phones_json),
preferred_phone=row.preferred_phone,
tags=_json_list(row.tags_json),
deal_count=deal_count,
created_at=row.created_at,
)
def _deal_to_out(
row: SalesDealRow,
*,
@@ -963,12 +978,17 @@ def _escalation_to_out(row: SalesEscalationRow) -> SalesEscalationOut:
)
def _task_to_out(row: SalesAutomationTaskRow) -> SalesAutomationTaskOut:
def _task_to_out(row: SalesAutomationTaskRow, deal: SalesDealRow | None = None) -> SalesAutomationTaskOut:
payload = _json_dict(row.payload_json)
return SalesAutomationTaskOut(
task_id=row.task_id,
deal_id=row.deal_id,
deal_title=deal.title if deal is not None else None,
deal_status=deal.status if deal is not None else None, # type: ignore[arg-type]
deal_stage_id=deal.stage_id if deal is not None else None,
deal_pipeline_id=deal.pipeline_id if deal is not None else None,
task_type=row.task_type,
payload=_json_dict(row.payload_json),
payload=payload,
run_at=row.run_at,
status=row.status, # type: ignore[arg-type]
retry_count=row.retry_count,
@@ -978,6 +998,8 @@ def _task_to_out(row: SalesAutomationTaskRow) -> SalesAutomationTaskOut:
completed_at=row.completed_at,
failed_at=row.failed_at,
last_error=row.last_error,
recommended_to_channel=payload.get("recommended_to_channel") or payload.get("to_channel"),
reason_code=payload.get("reason_code") or payload.get("reason"),
created_at=row.created_at,
updated_at=row.updated_at,
)
@@ -1057,6 +1079,29 @@ def _get_lead(session, lead_id: str, tenant_id: str) -> SalesLeadRow:
return row
def _latest_deal_for_lead(session, lead_id: str, tenant_id: str) -> SalesDealRow | None:
return session.execute(
select(SalesDealRow)
.where(SalesDealRow.lead_id == lead_id, SalesDealRow.tenant_id == tenant_id)
.order_by(SalesDealRow.id.desc())
).scalars().first()
def _latest_deals_by_lead(session, lead_ids: list[str], tenant_id: str) -> dict[str, SalesDealRow]:
if not lead_ids:
return {}
rows = session.execute(
select(SalesDealRow)
.where(SalesDealRow.lead_id.in_(lead_ids), SalesDealRow.tenant_id == tenant_id)
.order_by(SalesDealRow.id.desc())
).scalars().all()
latest: dict[str, SalesDealRow] = {}
for row in rows:
if row.lead_id and row.lead_id not in latest:
latest[row.lead_id] = row
return latest
def _get_deal(session, deal_id: str, tenant_id: str) -> SalesDealRow:
row = session.execute(
select(SalesDealRow).where(SalesDealRow.deal_id == deal_id, SalesDealRow.tenant_id == tenant_id)
@@ -1176,6 +1221,44 @@ def _find_customer_by_customer_id(session, customer_id: str | None) -> Customer
return session.execute(select(Customer).where(Customer.customer_id == normalized)).scalar_one_or_none()
def _sales_customer_ids_for_tenant(session, tenant_id: str) -> set[str]:
deal_ids = session.execute(
select(SalesDealRow.customer_id).where(
SalesDealRow.tenant_id == tenant_id,
SalesDealRow.customer_id.is_not(None),
)
).scalars().all()
lead_ids = session.execute(
select(SalesLeadRow.crm_customer_id).where(
SalesLeadRow.tenant_id == tenant_id,
SalesLeadRow.crm_customer_id.is_not(None),
)
).scalars().all()
return {str(customer_id) for customer_id in [*deal_ids, *lead_ids] if customer_id}
def _get_sales_customer(session, customer_id: str, tenant_id: str) -> Customer:
normalized = str(customer_id or "").strip()
if not normalized or normalized not in _sales_customer_ids_for_tenant(session, tenant_id):
raise HTTPException(status_code=404, detail="Customer not found")
row = _find_customer_by_customer_id(session, normalized)
if row is None:
raise HTTPException(status_code=404, detail="Customer not found")
return row
def _customer_deal_count(session, customer_id: str, tenant_id: str) -> int:
return int(
session.execute(
select(func.count()).select_from(SalesDealRow).where(
SalesDealRow.tenant_id == tenant_id,
SalesDealRow.customer_id == customer_id,
)
).scalar_one()
or 0
)
def _create_lead_and_deal_for_contact(
session,
*,
@@ -3047,7 +3130,7 @@ def _build_workspace(session, deal: SalesDealRow) -> SalesWorkspaceOut:
active_escalation = next((row for row in escalations if row.status in {"open", "assigned", "in_progress"}), None)
return SalesWorkspaceOut(
lead=_lead_to_out(lead) if lead else None,
lead=_lead_to_out(lead, deal) if lead else None,
deal=_deal_to_out(deal, pipeline=pipeline, stage=stage),
pipeline=_pipeline_to_out(pipeline) if pipeline else None,
stage=_stage_to_ref(stage),
@@ -3472,7 +3555,7 @@ def create_lead(
)
_publish_lead_entered_event(session, lead=lead, deal=deal, stage=stage, actor_type="human", actor_id=actor["user"])
session.commit()
return _lead_to_out(lead)
return _lead_to_out(lead, deal)
finally:
session.close()
@@ -3480,8 +3563,15 @@ def create_lead(
@app.get("/api/v1/leads", response_model=list[SalesLeadOut])
def list_leads(
query: str | None = None,
search: str | None = None,
status: str | None = None,
lead_temperature: str | None = None,
source_channel: str | None = None,
source_type: str | None = None,
customer_type: str | None = None,
preferred_channel: str | None = None,
limit: int = Query(default=100, ge=1, le=200),
offset: int = Query(default=0, ge=0),
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesLeadOut]:
session = get_session()
@@ -3494,18 +3584,31 @@ def list_leads(
)
if status:
stmt = stmt.where(SalesLeadRow.status == status)
if query:
pattern = f"%{query.strip().lower()}%"
if lead_temperature:
stmt = stmt.where(SalesLeadRow.lead_temperature == lead_temperature)
if source_channel:
stmt = stmt.where(SalesLeadRow.source_channel == source_channel)
if source_type:
stmt = stmt.where(SalesLeadRow.source_type == source_type)
if customer_type:
stmt = stmt.where(SalesLeadRow.customer_type == customer_type)
if preferred_channel:
stmt = stmt.where(SalesLeadRow.preferred_channel == preferred_channel)
search_term = (search or query or "").strip()
if search_term:
pattern = f"%{search_term.lower()}%"
stmt = stmt.where(
or_(
func.lower(func.coalesce(SalesLeadRow.full_name, "")).like(pattern),
func.lower(func.coalesce(SalesLeadRow.company_name, "")).like(pattern),
func.lower(func.coalesce(SalesLeadRow.phone, "")).like(pattern),
func.lower(func.coalesce(SalesLeadRow.email, "")).like(pattern),
func.lower(func.coalesce(SalesLeadRow.lead_id, "")).like(pattern),
)
)
rows = session.execute(stmt.limit(limit)).scalars().all()
return [_lead_to_out(row) for row in rows]
rows = session.execute(stmt.limit(limit).offset(offset)).scalars().all()
deals_by_lead = _latest_deals_by_lead(session, [row.lead_id for row in rows], tenant_id)
return [_lead_to_out(row, deals_by_lead.get(row.lead_id)) for row in rows]
finally:
session.close()
@@ -3517,7 +3620,9 @@ def get_lead(
) -> SalesLeadOut:
session = get_session()
try:
return _lead_to_out(_get_lead(session, lead_id, _tenant_id(actor)))
tenant_id = _tenant_id(actor)
lead = _get_lead(session, lead_id, tenant_id)
return _lead_to_out(lead, _latest_deal_for_lead(session, lead.lead_id, tenant_id))
finally:
session.close()
@@ -3539,7 +3644,7 @@ def update_lead(
setattr(lead, key, value)
lead.updated_at = utc_now_iso()
session.commit()
return _lead_to_out(lead)
return _lead_to_out(lead, _latest_deal_for_lead(session, lead.lead_id, tenant_id))
finally:
session.close()
@@ -3631,6 +3736,186 @@ def convert_lead_to_deal(
session.close()
@app.get("/api/v1/customers", response_model=list[SalesCustomerOut])
def list_customers(
query: str | None = None,
search: str | None = None,
limit: int = Query(default=100, ge=1, le=200),
offset: int = Query(default=0, ge=0),
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesCustomerOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
customer_ids = _sales_customer_ids_for_tenant(session, tenant_id)
if not customer_ids:
return []
stmt = select(Customer).where(Customer.customer_id.in_(customer_ids))
search_term = (search or query or "").strip()
if search_term:
pattern = f"%{search_term.lower()}%"
stmt = stmt.where(
or_(
func.lower(func.coalesce(Customer.customer_id, "")).like(pattern),
func.lower(func.coalesce(Customer.display_name, "")).like(pattern),
func.lower(func.coalesce(Customer.preferred_phone, "")).like(pattern),
)
)
rows = session.execute(stmt.order_by(Customer.display_name.asc(), Customer.id.desc()).limit(limit).offset(offset)).scalars().all()
counts = dict(
session.execute(
select(SalesDealRow.customer_id, func.count())
.where(
SalesDealRow.tenant_id == tenant_id,
SalesDealRow.customer_id.in_([row.customer_id for row in rows] or [""]),
)
.group_by(SalesDealRow.customer_id)
).all()
)
return [_customer_to_out(row, int(counts.get(row.customer_id) or 0)) for row in rows]
finally:
session.close()
@app.get("/api/v1/customers/{customer_id}", response_model=SalesCustomerOut)
def get_customer(
customer_id: str,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> SalesCustomerOut:
session = get_session()
try:
tenant_id = _tenant_id(actor)
customer = _get_sales_customer(session, customer_id, tenant_id)
return _customer_to_out(customer, _customer_deal_count(session, customer.customer_id, tenant_id))
finally:
session.close()
@app.patch("/api/v1/customers/{customer_id}", response_model=SalesCustomerOut)
def update_customer(
customer_id: str,
payload: SalesCustomerUpdate,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR)),
) -> SalesCustomerOut:
session = get_session()
try:
tenant_id = _tenant_id(actor)
customer = _get_sales_customer(session, customer_id, tenant_id)
updates = payload.model_dump(exclude_unset=True)
if "display_name" in updates and updates["display_name"] is not None:
customer.display_name = str(updates["display_name"])
if "phones" in updates and updates["phones"] is not None:
customer.phones_json = json.dumps(updates["phones"], ensure_ascii=False)
if "preferred_phone" in updates:
customer.preferred_phone = updates["preferred_phone"]
if "tags" in updates and updates["tags"] is not None:
customer.tags_json = json.dumps(updates["tags"], ensure_ascii=False)
session.commit()
return _customer_to_out(customer, _customer_deal_count(session, customer.customer_id, tenant_id))
finally:
session.close()
@app.get("/api/v1/customers/{customer_id}/deals", response_model=list[SalesDealOut])
def list_customer_deals(
customer_id: str,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesDealOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
_get_sales_customer(session, customer_id, tenant_id)
rows = session.execute(
select(SalesDealRow)
.where(SalesDealRow.tenant_id == tenant_id, SalesDealRow.customer_id == customer_id)
.order_by(SalesDealRow.updated_at.desc(), SalesDealRow.id.desc())
).scalars().all()
return [_deal_to_out_with_refs(session, row) for row in rows]
finally:
session.close()
@app.get("/api/v1/customers/{customer_id}/communications", response_model=list[SalesCommunicationOut])
def list_customer_communications(
customer_id: str,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesCommunicationOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
_get_sales_customer(session, customer_id, tenant_id)
rows = session.execute(
select(SalesCommunicationSessionRow)
.where(SalesCommunicationSessionRow.tenant_id == tenant_id, SalesCommunicationSessionRow.customer_id == customer_id)
.order_by(SalesCommunicationSessionRow.started_at.desc(), SalesCommunicationSessionRow.id.desc())
).scalars().all()
return [_communication_to_out(row) for row in rows]
finally:
session.close()
@app.get("/api/v1/customers/{customer_id}/documents", response_model=list[SalesDocumentOut])
def list_customer_documents(
customer_id: str,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesDocumentOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
_get_sales_customer(session, customer_id, tenant_id)
rows = session.execute(
select(SalesDocumentRow)
.where(SalesDocumentRow.tenant_id == tenant_id, SalesDocumentRow.customer_id == customer_id)
.order_by(SalesDocumentRow.updated_at.desc(), SalesDocumentRow.id.desc())
).scalars().all()
return [_document_to_out(row) for row in rows]
finally:
session.close()
@app.get("/api/v1/customers/{customer_id}/invoices", response_model=list[SalesInvoiceOut])
def list_customer_invoices(
customer_id: str,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesInvoiceOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
_get_sales_customer(session, customer_id, tenant_id)
rows = session.execute(
select(SalesInvoiceRow)
.where(SalesInvoiceRow.tenant_id == tenant_id, SalesInvoiceRow.customer_id == customer_id)
.order_by(SalesInvoiceRow.updated_at.desc(), SalesInvoiceRow.id.desc())
).scalars().all()
return [_invoice_to_out(row) for row in rows]
finally:
session.close()
@app.get("/api/v1/customers/{customer_id}/payments", response_model=list[SalesPaymentOut])
def list_customer_payments(
customer_id: str,
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesPaymentOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
_get_sales_customer(session, customer_id, tenant_id)
deal_ids = session.execute(
select(SalesDealRow.deal_id).where(SalesDealRow.tenant_id == tenant_id, SalesDealRow.customer_id == customer_id)
).scalars().all()
rows = session.execute(
select(SalesPaymentRow)
.where(SalesPaymentRow.tenant_id == tenant_id, SalesPaymentRow.deal_id.in_(deal_ids or [""]))
.order_by(SalesPaymentRow.updated_at.desc(), SalesPaymentRow.id.desc())
).scalars().all()
return [_payment_to_out(row) for row in rows]
finally:
session.close()
@app.post("/api/v1/deals", response_model=SalesDealOut)
def create_deal(
payload: SalesDealCreate,
@@ -6494,6 +6779,79 @@ def list_automation_tasks(
session.close()
@app.get("/api/v1/automation-tasks", response_model=list[SalesAutomationTaskOut])
def list_tenant_automation_tasks(
status: str | None = Query(default=None),
task_type: str | None = Query(default=None),
run_at: str | None = Query(default=None),
deal_search: str | None = Query(default=None),
failed_only: bool = Query(default=False),
pending_only: bool = Query(default=False),
limit: int = Query(default=100, ge=1, le=500),
offset: int = Query(default=0, ge=0),
actor: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.OPERATOR, Role.ANALYST)),
) -> list[SalesAutomationTaskOut]:
session = get_session()
try:
tenant_id = _tenant_id(actor)
conditions = [SalesAutomationTaskRow.tenant_id == tenant_id]
if status:
conditions.append(SalesAutomationTaskRow.status == status)
if task_type:
conditions.append(SalesAutomationTaskRow.task_type == task_type)
if failed_only:
conditions.append(
or_(
SalesAutomationTaskRow.status == "failed",
SalesAutomationTaskRow.failed_at.is_not(None),
SalesAutomationTaskRow.last_error.is_not(None),
)
)
elif pending_only:
conditions.append(SalesAutomationTaskRow.status == "pending")
now = datetime.now(tz=timezone.utc).replace(microsecond=0)
if run_at == "today":
day_start = now.replace(hour=0, minute=0, second=0)
day_end = day_start + timedelta(days=1)
conditions.append(SalesAutomationTaskRow.run_at >= day_start.isoformat())
conditions.append(SalesAutomationTaskRow.run_at < day_end.isoformat())
elif run_at == "overdue":
conditions.append(SalesAutomationTaskRow.status == "pending")
conditions.append(SalesAutomationTaskRow.run_at < now.isoformat())
elif run_at == "future":
conditions.append(SalesAutomationTaskRow.run_at >= now.isoformat())
if deal_search:
term = f"%{deal_search.strip().lower()}%"
if term != "%%":
conditions.append(
or_(
func.lower(SalesAutomationTaskRow.deal_id).like(term),
func.lower(SalesDealRow.title).like(term),
func.lower(SalesDealRow.customer_id).like(term),
)
)
rows = session.execute(
select(SalesAutomationTaskRow, SalesDealRow)
.join(
SalesDealRow,
(SalesDealRow.deal_id == SalesAutomationTaskRow.deal_id)
& (SalesDealRow.tenant_id == SalesAutomationTaskRow.tenant_id),
isouter=True,
)
.where(*conditions)
.order_by(SalesAutomationTaskRow.run_at.asc(), SalesAutomationTaskRow.id.desc())
.limit(limit)
.offset(offset)
).all()
return [_task_to_out(task, deal) for task, deal in rows]
finally:
session.close()
@app.post("/api/v1/automation-tasks/{task_id}/cancel", response_model=SalesAutomationTaskOut)
def cancel_automation_task(
task_id: str,
+25
View File
@@ -162,6 +162,7 @@ class SalesLeadEnrichIn(BaseModel):
class SalesLeadOut(BaseModel):
lead_id: str
tenant_id: str
deal_id: str | None = None
source_type: str
source_channel: str
source_campaign_id: str | None = None
@@ -183,6 +184,23 @@ class SalesLeadOut(BaseModel):
updated_at: str
class SalesCustomerOut(BaseModel):
customer_id: str
display_name: str
phones: list = Field(default_factory=list)
preferred_phone: str | None = None
tags: list = Field(default_factory=list)
deal_count: int = 0
created_at: str
class SalesCustomerUpdate(BaseModel):
display_name: str | None = Field(default=None, min_length=2)
phones: list | None = None
preferred_phone: str | None = None
tags: list | None = None
class SalesDealCreate(BaseModel):
lead_id: str | None = None
customer_id: str | None = None
@@ -684,6 +702,11 @@ class SalesEscalationCancelIn(BaseModel):
class SalesAutomationTaskOut(BaseModel):
task_id: str
deal_id: str
deal_title: str | None = None
deal_status: SalesDealStatus | None = None
deal_stage_id: str | None = None
deal_stage_code: str | None = None
deal_pipeline_id: str | None = None
task_type: str
payload: dict = Field(default_factory=dict)
run_at: str
@@ -695,6 +718,8 @@ class SalesAutomationTaskOut(BaseModel):
completed_at: str | None = None
failed_at: str | None = None
last_error: str | None = None
recommended_to_channel: SalesSwitchChannel | str | None = None
reason_code: str | None = None
created_at: str
updated_at: str
+65
View File
@@ -1,3 +1,5 @@
from fastapi.testclient import TestClient
import services.sales_service.app as sales_module
@@ -22,10 +24,73 @@ def test_sales_frontend_contract_routes_are_registered():
"/api/v1/deals/{deal_id}/conditions": {"GET", "POST", "PATCH", "PUT"},
"/api/v1/deals/{deal_id}/conditions/confirm": {"POST"},
"/api/v1/deals/{deal_id}/escalations": {"GET", "POST"},
"/api/v1/automation-tasks": {"GET"},
"/api/v1/deals/{deal_id}/automation-tasks": {"GET"},
"/api/v1/automation-tasks/{task_id}/cancel": {"POST"},
"/api/v1/automation-tasks/{task_id}/run-now": {"POST"},
"/api/v1/leads": {"GET", "POST"},
"/api/v1/leads/{lead_id}": {"GET", "PATCH"},
"/api/v1/leads/{lead_id}/enrich": {"POST"},
"/api/v1/leads/{lead_id}/convert-to-deal": {"POST"},
"/api/v1/customers": {"GET"},
"/api/v1/customers/{customer_id}": {"GET", "PATCH"},
"/api/v1/customers/{customer_id}/deals": {"GET"},
"/api/v1/customers/{customer_id}/communications": {"GET"},
"/api/v1/customers/{customer_id}/documents": {"GET"},
"/api/v1/customers/{customer_id}/invoices": {"GET"},
"/api/v1/customers/{customer_id}/payments": {"GET"},
}
for path, methods in expected.items():
assert methods.issubset(_methods(path))
def test_tenant_automation_tasks_route_accepts_registry_filters():
client = TestClient(sales_module.app)
response = client.get(
"/api/v1/automation-tasks",
params={
"status": "pending",
"run_at": "future",
"pending_only": "true",
"deal_search": "demo",
},
headers={"X-User": "admin", "X-Role": "admin", "X-Tenant-ID": "tenant_contract"},
)
assert response.status_code == 200
assert isinstance(response.json(), list)
def test_leads_route_accepts_registry_filters():
client = TestClient(sales_module.app)
response = client.get(
"/api/v1/leads",
params={
"search": "demo",
"status": "new_qualified_lead",
"lead_temperature": "warm",
"source_channel": "telegram",
"source_type": "crm",
"customer_type": "b2b",
"preferred_channel": "telegram",
"limit": "25",
"offset": "0",
},
headers={"X-User": "admin", "X-Role": "admin", "X-Tenant-ID": "tenant_contract"},
)
assert response.status_code == 200
assert isinstance(response.json(), list)
def test_customers_route_accepts_registry_filters():
client = TestClient(sales_module.app)
response = client.get(
"/api/v1/customers",
params={"search": "demo", "limit": "25", "offset": "0"},
headers={"X-User": "admin", "X-Role": "admin", "X-Tenant-ID": "tenant_contract"},
)
assert response.status_code == 200
assert isinstance(response.json(), list)