529 lines
17 KiB
Python

from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from typing import Dict, Any, List, Optional, Union
from pydantic import BaseModel
from datetime import datetime
import time
import json
from adapters.sqlite.database import get_db
from adapters.sqlite.models import (
Deal, Contact, Company, Pipeline, PipelineStage, User, Event, CustomField,
deal_contacts, deal_companies, contact_companies
)
router = APIRouter()
class CustomFieldValue(BaseModel):
field_id: int
field_name: str
field_type: str
values: List[Dict[str, Any]]
class EmbeddedRelation(BaseModel):
id: int
is_main: Optional[bool] = None
class DataItem(BaseModel):
id: int
name: Optional[str] = None
# Add other common fields as needed
created_at: Optional[int] = None
updated_at: Optional[int] = None
custom_fields_values: Optional[List[CustomFieldValue]] = None
_embedded: Optional[Dict[str, List[EmbeddedRelation]]] = None
# Store raw data for complete preservation
raw_data: Optional[Dict[str, Any]] = None
class DataIngestionRequest(BaseModel):
data: List[Dict[str, Any]]
sync_mode: str = "upsert" # "insert", "upsert", "replace"
def process_custom_fields(
custom_fields: List[Dict[str, Any]],
entity_type: str,
entity_id: int,
db: Session
):
"""Process and store custom fields"""
# Clear existing custom fields for this entity if replacing
db.query(CustomField).filter(
CustomField.entity_type == entity_type,
CustomField.entity_id == entity_id
).delete()
for field_data in custom_fields:
field_id = field_data.get('field_id')
field_name = field_data.get('field_name', '')
field_type = field_data.get('field_type', 'text')
values = field_data.get('values', [])
for value_data in values:
value = value_data.get('value', '')
custom_field = CustomField(
entity_type=entity_type,
entity_id=entity_id,
field_id=field_id,
field_name=field_name,
field_type=field_type,
field_value=str(value) if value else None,
is_custom=True,
created_at=int(time.time()),
updated_at=int(time.time())
)
# Handle different field types
if field_type == 'numeric' and value:
try:
custom_field.field_value_numeric = float(value)
except (ValueError, TypeError):
pass
elif field_type == 'date' and value:
try:
# AMO CRM returns dates as unix timestamps
custom_field.field_value_date = int(value)
except (ValueError, TypeError):
pass
elif field_type in ['select', 'multiselect']:
# For select fields, store the display value
custom_field.field_value = str(value)
db.add(custom_field)
def process_relationships(
embedded_data: Dict[str, Any],
entity_id: int,
entity_type: str,
db: Session
):
"""Process embedded relationships"""
if entity_type == "deals":
# Process deal-contact relationships
if 'contacts' in embedded_data:
# Clear existing relationships
db.execute(
deal_contacts.delete().where(deal_contacts.c.deal_id == entity_id)
)
for contact_data in embedded_data['contacts']:
contact_id = contact_data['id']
is_main = contact_data.get('is_main', False)
# Insert relationship
db.execute(
deal_contacts.insert().values(
deal_id=entity_id,
contact_id=contact_id,
is_main=is_main
)
)
# Process deal-company relationships
if 'companies' in embedded_data:
# Clear existing relationships
db.execute(
deal_companies.delete().where(deal_companies.c.deal_id == entity_id)
)
for company_data in embedded_data['companies']:
company_id = company_data['id']
is_main = company_data.get('is_main', False)
# Insert relationship
db.execute(
deal_companies.insert().values(
deal_id=entity_id,
company_id=company_id,
is_main=is_main
)
)
elif entity_type == "contacts":
# Process contact-company relationships
if 'companies' in embedded_data:
# Clear existing relationships
db.execute(
contact_companies.delete().where(contact_companies.c.contact_id == entity_id)
)
for company_data in embedded_data['companies']:
company_id = company_data['id']
is_main = company_data.get('is_main', False)
# Insert relationship
db.execute(
contact_companies.insert().values(
contact_id=entity_id,
company_id=company_id,
is_main=is_main
)
)
@router.post("/users")
async def put_users_data(
request: DataIngestionRequest,
db: Session = Depends(get_db)
) -> Dict[str, Any]:
"""Put users data into database"""
processed_count = 0
for user_data in request.data:
user_id = user_data['id']
# Check if user exists
existing_user = db.query(User).filter(User.id == user_id).first()
if request.sync_mode == "insert" and existing_user:
continue # Skip existing records in insert mode
user_values = {
'id': user_id,
'name': user_data.get('name'),
'email': user_data.get('email'),
'is_active': user_data.get('is_active', True),
'created_at': user_data.get('created_at'),
'updated_at': user_data.get('updated_at', int(time.time())),
'raw_data': user_data
}
if existing_user:
# Update existing user
for key, value in user_values.items():
if key != 'id': # Don't update ID
setattr(existing_user, key, value)
else:
# Create new user
new_user = User(**user_values)
db.add(new_user)
processed_count += 1
db.commit()
return {
"message": f"Processed {processed_count} users",
"processed_count": processed_count
}
@router.post("/pipelines")
async def put_pipelines_data(
request: DataIngestionRequest,
db: Session = Depends(get_db)
) -> Dict[str, Any]:
"""Put pipelines data into database"""
processed_count = 0
for pipeline_data in request.data:
pipeline_id = pipeline_data['id']
# Check if pipeline exists
existing_pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
if request.sync_mode == "insert" and existing_pipeline:
continue
pipeline_values = {
'id': pipeline_id,
'name': pipeline_data.get('name'),
'sort': pipeline_data.get('sort'),
'is_main': pipeline_data.get('is_main', False),
'is_unsorted': pipeline_data.get('is_unsorted', False),
'is_archive': pipeline_data.get('is_archive', False),
'account_id': pipeline_data.get('account_id'),
'created_at': pipeline_data.get('created_at'),
'updated_at': pipeline_data.get('updated_at', int(time.time())),
'raw_data': pipeline_data
}
if existing_pipeline:
for key, value in pipeline_values.items():
if key != 'id':
setattr(existing_pipeline, key, value)
else:
new_pipeline = Pipeline(**pipeline_values)
db.add(new_pipeline)
# Process pipeline stages
embedded_data = pipeline_data.get('_embedded', {})
if 'statuses' in embedded_data:
for status_data in embedded_data['statuses']:
status_id = status_data['id']
existing_stage = db.query(PipelineStage).filter(PipelineStage.id == status_id).first()
stage_values = {
'id': status_id,
'pipeline_id': pipeline_id,
'name': status_data.get('name'),
'sort': status_data.get('sort'),
'is_editable': status_data.get('is_editable', True),
'color': status_data.get('color'),
'created_at': status_data.get('created_at'),
'updated_at': status_data.get('updated_at', int(time.time())),
'raw_data': status_data
}
if existing_stage:
for key, value in stage_values.items():
if key != 'id':
setattr(existing_stage, key, value)
else:
new_stage = PipelineStage(**stage_values)
db.add(new_stage)
processed_count += 1
db.commit()
return {
"message": f"Processed {processed_count} pipelines",
"processed_count": processed_count
}
@router.post("/companies")
async def put_companies_data(
request: DataIngestionRequest,
db: Session = Depends(get_db)
) -> Dict[str, Any]:
"""Put companies data into database"""
processed_count = 0
for company_data in request.data:
company_id = company_data['id']
existing_company = db.query(Company).filter(Company.id == company_id).first()
if request.sync_mode == "insert" and existing_company:
continue
company_values = {
'id': company_id,
'name': company_data.get('name'),
'responsible_user_id': company_data.get('responsible_user_id'),
'group_id': company_data.get('group_id'),
'created_by': company_data.get('created_by'),
'updated_by': company_data.get('updated_by'),
'created_at': company_data.get('created_at'),
'updated_at': company_data.get('updated_at', int(time.time())),
'closest_task_at': company_data.get('closest_task_at'),
'is_deleted': company_data.get('is_deleted', False),
'raw_data': company_data
}
if existing_company:
for key, value in company_values.items():
if key != 'id':
setattr(existing_company, key, value)
else:
new_company = Company(**company_values)
db.add(new_company)
# Process custom fields
custom_fields = company_data.get('custom_fields_values', [])
if custom_fields:
process_custom_fields(custom_fields, "companies", company_id, db)
processed_count += 1
db.commit()
return {
"message": f"Processed {processed_count} companies",
"processed_count": processed_count
}
@router.post("/contacts")
async def put_contacts_data(
request: DataIngestionRequest,
db: Session = Depends(get_db)
) -> Dict[str, Any]:
"""Put contacts data into database"""
processed_count = 0
for contact_data in request.data:
contact_id = contact_data['id']
existing_contact = db.query(Contact).filter(Contact.id == contact_id).first()
if request.sync_mode == "insert" and existing_contact:
continue
contact_values = {
'id': contact_id,
'name': contact_data.get('name'),
'first_name': contact_data.get('first_name'),
'last_name': contact_data.get('last_name'),
'responsible_user_id': contact_data.get('responsible_user_id'),
'group_id': contact_data.get('group_id'),
'created_by': contact_data.get('created_by'),
'updated_by': contact_data.get('updated_by'),
'created_at': contact_data.get('created_at'),
'updated_at': contact_data.get('updated_at', int(time.time())),
'closest_task_at': contact_data.get('closest_task_at'),
'is_deleted': contact_data.get('is_deleted', False),
'raw_data': contact_data
}
if existing_contact:
for key, value in contact_values.items():
if key != 'id':
setattr(existing_contact, key, value)
else:
new_contact = Contact(**contact_values)
db.add(new_contact)
# Process custom fields
custom_fields = contact_data.get('custom_fields_values', [])
if custom_fields:
process_custom_fields(custom_fields, "contacts", contact_id, db)
# Process relationships
embedded_data = contact_data.get('_embedded', {})
if embedded_data:
process_relationships(embedded_data, contact_id, "contacts", db)
processed_count += 1
db.commit()
return {
"message": f"Processed {processed_count} contacts",
"processed_count": processed_count
}
@router.post("/deals")
async def put_deals_data(
request: DataIngestionRequest,
db: Session = Depends(get_db)
) -> Dict[str, Any]:
"""Put deals (leads) data into database"""
processed_count = 0
for deal_data in request.data:
deal_id = deal_data['id']
existing_deal = db.query(Deal).filter(Deal.id == deal_id).first()
if request.sync_mode == "insert" and existing_deal:
continue
deal_values = {
'id': deal_id,
'name': deal_data.get('name'),
'price': deal_data.get('price', 0),
'responsible_user_id': deal_data.get('responsible_user_id'),
'group_id': deal_data.get('group_id'),
'status_id': deal_data.get('status_id'),
'pipeline_id': deal_data.get('pipeline_id'),
'loss_reason_id': deal_data.get('loss_reason_id'),
'created_by': deal_data.get('created_by'),
'updated_by': deal_data.get('updated_by'),
'closed_at': deal_data.get('closed_at'),
'created_at': deal_data.get('created_at'),
'updated_at': deal_data.get('updated_at', int(time.time())),
'closest_task_at': deal_data.get('closest_task_at'),
'is_deleted': deal_data.get('is_deleted', False),
'raw_data': deal_data
}
if existing_deal:
for key, value in deal_values.items():
if key != 'id':
setattr(existing_deal, key, value)
else:
new_deal = Deal(**deal_values)
db.add(new_deal)
# Process custom fields
custom_fields = deal_data.get('custom_fields_values', [])
if custom_fields:
process_custom_fields(custom_fields, "deals", deal_id, db)
# Process relationships
embedded_data = deal_data.get('_embedded', {})
if embedded_data:
process_relationships(embedded_data, deal_id, "deals", db)
processed_count += 1
db.commit()
return {
"message": f"Processed {processed_count} deals",
"processed_count": processed_count
}
@router.post("/events")
async def put_events_data(
request: DataIngestionRequest,
db: Session = Depends(get_db)
) -> Dict[str, Any]:
"""Put events data into database"""
processed_count = 0
for event_data in request.data:
event_id = event_data['id']
event_type = event_data.get('type')
# Filter only supported event types
if event_type not in ['incoming_call', 'outgoing_call', 'lead_status_changed']:
continue
existing_event = db.query(Event).filter(Event.id == event_id).first()
if request.sync_mode == "insert" and existing_event:
continue
event_values = {
'id': event_id,
'type': event_type,
'entity_id': event_data.get('entity_id'),
'entity_type': event_data.get('entity_type'),
'created_by': event_data.get('created_by'),
'created_at': event_data.get('created_at'),
'value_after': event_data.get('value_after'),
'value_before': event_data.get('value_before'),
'account_id': event_data.get('account_id'),
'raw_data': event_data
}
if existing_event:
for key, value in event_values.items():
if key != 'id':
setattr(existing_event, key, value)
else:
new_event = Event(**event_values)
db.add(new_event)
processed_count += 1
db.commit()
return {
"message": f"Processed {processed_count} events",
"processed_count": processed_count
}