Integrations - Services Documentation

Location: api/nextango/apps/integrations/sync/ Last Updated: 2025-12-21

Overview

The Integrations domain provides a comprehensive sync framework for bidirectional data synchronization between Django models and Sanity CMS. The service layer consists of:

  1. Universal Sync Framework - Generic sync infrastructure for all model types
  2. Domain Sync Handlers - Domain-specific sync implementations (users, products, stores, etc.)
  3. Service Layer - Shared services for field transformation, reference resolution, and orchestration

This architecture enables consistent, maintainable sync operations across all domains while allowing domain-specific business logic where needed.


Architecture Patterns

Layered Architecture

┌─────────────────────────────────────────────────┐
│         Domain Signal Handlers                  │
│    (users/signals.py, products/signals.py)      │
└──────────────────┬──────────────────────────────┘

┌──────────────────▼──────────────────────────────┐
│      Domain Sync Handlers                       │
│  (user_sync.py, product_sync.py, store_sync.py) │
│         ↓ inherits from ↓                       │
│      BaseSyncHandler (base.py)                  │
└──────────────────┬──────────────────────────────┘

┌──────────────────▼──────────────────────────────┐
│      Universal Sync Framework                   │
│       (universal_sync.py)                       │
└──────────────────┬──────────────────────────────┘

┌──────────────────▼──────────────────────────────┐
│         Service Layer                           │
│  • FieldTransformer                             │
│  • ReferenceResolver                            │
│  • SyncOrchestrator                             │
└──────────────────┬──────────────────────────────┘

┌──────────────────▼──────────────────────────────┐
│      Sanity Client & Infrastructure             │
│  (sanity_client.py, models.py)                  │
└─────────────────────────────────────────────────┘

Key Design Principles

  1. Inheritance-Based Domain Handlers - All domain sync handlers inherit from BaseSyncHandler
  2. Dependency Inversion - Services register business logic functions rather than containing domain knowledge
  3. Thread-Local Context - Loop prevention using threading.local() for webhook detection
  4. Race Condition Handling - update_or_create() with idempotent operations
  5. Audit Trail - Every sync operation logged to SyncAuditLog

Universal Sync Framework

File: api/nextango/apps/integrations/sync/universal_sync.py

UniversalSyncFramework Class

Generic sync framework handling all model types automatically.

Initialization

class UniversalSyncFramework:
    def __init__(self):
        self.sanity_client = SanityAPIClient()
        self._model_registry = {}
        self._field_transformers = {}
        self._sync_config = self._load_sync_config()
        self._initialize_transformers()
        self._use_business_references = self._sync_config.get('reference_strategy') == 'business_meaningful'

Source: api/nextango/apps/integrations/sync/universal_sync.py:28-39

Configuration Loading

Loads sync configuration with fallback chain:

def _load_sync_config(self) -> Dict[str, Any]:
    """Load sync configuration from file or settings"""
    try:
        # Try to load schema-driven config first
        with open('/app/schema_driven_sync_config.json', 'r') as f:
            config = json.load(f)
            logger.debug("Loaded schema-driven sync configuration with business references")
            return config
    except FileNotFoundError:
        try:
            # Fallback to comprehensive sync config
            with open('comprehensive_sync_config.json', 'r') as f:
                return json.load(f)
        except FileNotFoundError:
            # Fallback to hardcoded config
            return self._get_default_config()

Source: api/nextango/apps/integrations/sync/universal_sync.py:169-184

Configuration structure includes:

  • Model Mappings - Django model to Sanity type mappings
  • Field Mappings - snake_case to camelCase field conversions
  • Reference Fields - ForeignKey and ManyToMany reference handling
  • Global Settings - Batch sizes, conflict resolution, excluded fields

Model Registration

def register_model(self, model_class: Type[models.Model],
                  sanity_type: Optional[str] = None,
                  field_mappings: Optional[Dict[str, str]] = None):
    """Register a Django model for sync"""
    model_key = f"{model_class._meta.app_label}.{model_class.__name__}"

    # Check if model is in config
    if model_key in self._sync_config['model_mappings']:
        config = self._sync_config['model_mappings'][model_key].copy()
        if not config.get('sync_enabled', True):
            logger.info(f"Sync disabled for {model_key}")
            return

Source: api/nextango/apps/integrations/sync/universal_sync.py:501-512

Registered models stored in _model_registry with:

  • model_class - Django model class reference
  • config - Sync configuration (sanity_type, field_mappings, sync_direction)

Sync to Sanity (Outbound)

Primary outbound sync method with comprehensive loop prevention:

def sync_to_sanity(self, instance: models.Model, created: bool = False):
    """Sync a Django instance to Sanity"""
    # Check for LOAD_TESTING_MODE first
    if os.getenv('LOAD_TESTING_MODE') == 'true':
        logger.info(f"LOAD TESTING MODE: Skipping Sanity sync for {instance._meta.app_label}.{instance.__class__.__name__} {instance.pk}")
        return

    model_key = f"{instance._meta.app_label}.{instance.__class__.__name__}"

    if model_key not in self._model_registry:
        logger.warning(f"Model {model_key} not registered for sync")
        return

    # ENHANCED LOOP PREVENTION: Check thread-local context first, then instance attribute
    from .sync_context import should_skip_sync_to_sanity, SyncContext

    if should_skip_sync_to_sanity():
        sync_source = SyncContext.get_sync_source()
        logger.info(f"Skipping Django→Sanity sync for {model_key} {instance.pk} (thread-local source: {sync_source})")
        return

Source: api/nextango/apps/integrations/sync/universal_sync.py:560-578

Sync Flow:

  1. Check LOAD_TESTING_MODE environment variable
  2. Verify model is registered
  3. Primary loop prevention - Check thread-local context via should_skip_sync_to_sanity()
  4. Fallback check - Instance _sync_source attribute
  5. Deduplication - Cache-based locking (60s TTL)
  6. Model-specific locking - UserRole: 5s, User models: 10s, Others: 30s
  7. Transform - Django instance → Sanity document
  8. Sync - Create or update in Sanity
  9. Audit - Log to SyncAuditLog

Business Reference Support

Enhanced sync with meaningful business references (vs database IDs):

def transform_to_sanity_with_business_refs(self, instance: models.Model) -> Dict[str, Any]:
    """
    Transform Django instance to Sanity document using business-meaningful references.
    Consolidated from enhanced_universal_sync functionality.
    """
    model_key = f"{instance._meta.app_label}.{instance.__class__.__name__}"
    config = self._model_registry.get(model_key, {}).get('config')

    if not config:
        logger.warning(f"No configuration found for {model_key}")
        return {}

    sanity_doc = {'_type': config['sanity_type']}

    # Map basic fields
    field_mappings = config.get('field_mappings', {})
    for django_field, sanity_field in field_mappings.items():
        if hasattr(instance, django_field):
            value = getattr(instance, django_field)
            if django_field not in ['role', 'store_assignments'] and value is not None:
                sanity_doc[sanity_field] = self._transform_field_value_enhanced(value)

    # Handle business-meaningful references
    reference_fields = config.get('reference_fields', {})
    for django_field, ref_config in reference_fields.items():
        if hasattr(instance, django_field):
            ref_value = getattr(instance, django_field)
            if ref_value:
                sanity_field = ref_config['sanity_field']
                if ref_config.get('is_array'):
                    # Handle many-to-many relationships
                    sanity_doc[sanity_field] = [{'_key': str(item.id), '_ref': str(item.sanity_id)} for item in ref_value.all() if hasattr(item, 'sanity_id') and item.sanity_id]
                else:
                    # Handle foreign key relationships
                    if hasattr(ref_value, 'sanity_id') and ref_value.sanity_id:
                        sanity_doc[sanity_field] = {'_ref': str(ref_value.sanity_id)}

    return sanity_doc

Source: api/nextango/apps/integrations/sync/universal_sync.py:41-78

Field Transformation

Comprehensive field transformer registry:

def _initialize_transformers(self):
    """Initialize field transformation functions"""
    from decimal import Decimal
    import json

    self._field_transformers = {
        # Date transformations
        'date_to_string': lambda d: d.isoformat() if d else None,
        'string_to_date': lambda s: datetime.fromisoformat(s) if s else None,

        # Decimal transformations
        'decimal_to_float': lambda d: float(d) if isinstance(d, Decimal) else d,
        'decimal_to_string': lambda d: str(d) if isinstance(d, Decimal) else d,

        # Reference transformations
        'django_to_reference': self._django_to_reference_safe,

        'reference_to_django': lambda ref, model_class: (
            model_class.objects.filter(sanity_id=ref.get('_ref')).first()
            if ref and '_ref' in ref else None
        ),

        # Array transformations - handle RelatedManager and QuerySet
        'queryset_to_array': lambda qs: [
            ref for ref in [
                self._field_transformers['django_to_reference'](obj)
                for obj in (qs.all() if hasattr(qs, 'all') else qs)
            ] if ref is not None
        ] if qs else [],

        # Generic JSON serializable converter
        'make_serializable': self._make_json_serializable,
    }

Source: api/nextango/apps/integrations/sync/universal_sync.py:371-412

User Model Special Handling

Users require custom userRoles array construction:

def _add_user_roles_to_sanity_doc(self, user_instance, sanity_doc):
    """
    Add the userRoles array to the Sanity document for User models.

    UPDATED: userRoles array is built from user_roles_data JSONField and existing role models.
    This provides a complete view of the user's roles for Sanity Studio display.
    """
    from decimal import Decimal

    # Initialize userRoles array
    user_roles = []

    # Get all roles from the user's computed methods (which use user_roles_data)
    all_roles = user_instance.get_all_roles()

    if all_roles:
        logger.info(f"Building userRoles array from computed roles: {all_roles}")

        # Map role types to Sanity _type
        role_type_mapping = {
            'employee': 'roleEmployee',
            'customer': 'roleCustomer',
            'manager': 'roleManager',
            'influencer': 'roleInfluencer',
            'vip': 'roleVip',
        }

        for role_type in all_roles:
            sanity_role_type = role_type_mapping.get(role_type)
            if not sanity_role_type:
                continue

            # Try to get the actual role instance for this type
            role_instance = None
            try:
                if role_type == 'employee' and hasattr(user_instance, 'employee'):
                    role_instance = user_instance.employee
                elif role_type == 'customer' and hasattr(user_instance, 'customer'):
                    role_instance = user_instance.customer
                # ... etc for other role types

Source: api/nextango/apps/integrations/sync/universal_sync.py:755-806

Bulk Sync Operations

def bulk_sync_to_sanity(self, model_key: str, queryset: Optional[models.QuerySet] = None):
    """Bulk sync Django instances to Sanity"""
    if model_key not in self._model_registry:
        logger.error(f"Model {model_key} not registered")
        return

    model_class = self._model_registry[model_key]['model_class']
    config = self._model_registry[model_key]['config']

    if queryset is None:
        queryset = model_class.objects.all()

    batch_size = self._sync_config['global_settings']['batch_size']
    total = queryset.count()
    synced = 0
    failed = 0

    logger.info(f"Starting bulk sync of {total} {model_key} instances to Sanity")

    for batch_start in range(0, total, batch_size):
        batch = queryset[batch_start:batch_start + batch_size]

        for instance in batch:
            try:
                self.sync_to_sanity(instance, created=False)
                synced += 1
            except Exception as e:
                logger.error(f"Failed to sync {instance.pk}: {e}")
                failed += 1

    logger.info(f"Bulk sync completed: {synced} synced, {failed} failed")
    return {'synced': synced, 'failed': failed, 'total': total}

Source: api/nextango/apps/integrations/sync/universal_sync.py:1026-1057

Singleton Pattern

# Singleton instance
_universal_sync_instance = None

def get_universal_sync():
    global _universal_sync_instance
    if _universal_sync_instance is None:
        _universal_sync_instance = UniversalSyncFramework()
    return _universal_sync_instance

# For backward compatibility
universal_sync = get_universal_sync()

Source: api/nextango/apps/integrations/sync/universal_sync.py:1102-1113


Domain Sync Handlers

Domain-specific sync handlers inherit from BaseSyncHandler and implement custom business logic.

UserSyncHandler

File: api/nextango/apps/integrations/sync/user_sync.py

Example domain handler showing the standard pattern.

Class Structure

class UserSyncHandler(BaseSyncHandler):
    """
    Consolidated user domain sync handler with BaseSyncHandler inheritance.
    Uses targeted service integration only where complexity justifies abstraction.
    """

    def __init__(self, sanity_client: Optional[SanityAPIClient] = None):
        super().__init__(sanity_client)

        # Use global shared reference resolver for cross-domain reference handling
        self.reference_resolver = get_global_reference_resolver()
        logger.info(f"User handler using global reference resolver: {id(self.reference_resolver)}")

        # Register domain-specific reference logic for architectural purity
        register_user_reference_logic(self.reference_resolver)

        self._load_user_sync_config()
        self._register_user_models()

Source: api/nextango/apps/integrations/sync/user_sync.py:24-42

Configuration

Domain-specific field mappings with snake_case standardization:

def _load_user_sync_config(self):
    """Load user-specific sync configuration with systematic snake_case field mappings."""
    self.user_sync_config = {
        'model_mappings': {
            'users.User': {
                'sanity_type': 'user',
                'sync_enabled': True,
                'sync_direction': 'bidirectional',
                'field_mappings': {
                    'name': 'name',
                    'email': 'email',
                    'phone': 'phone',
                    'status': 'status',
                    'created_date': 'createdDate',
                    'user_roles_data': 'userRoles'  # JSONField -> special handling
                },
                'business_logic': ['user_roles_array_construction']
            },
            'users.Employee': {
                'sanity_type': 'userEmployee',
                'sync_enabled': True,
                'sync_direction': 'bidirectional',
                'field_mappings': {
                    'employee_id': 'employeeId',
                    'pin_code': 'pinCode',
                    'hourly_rate': 'hourlyRate',
                    'commission_rate': 'commissionRate',
                    'hire_date': 'hireDate',
                    'store_assignments': 'storeAssignments',
                    'user': 'user'  # ForeignKey reference - uses ReferenceResolver
                }
            },
            # ... other user role models

Source: api/nextango/apps/integrations/sync/user_sync.py:43-74

Outbound Sync (sync_to_sanity)

Enhanced with webhook detection and context preservation:

@monitor_sync_operation
def sync_to_sanity(self, django_instance, created: bool = False) -> Optional[str]:
    """
    Single code path for syncing user domain instances to Sanity.
    Direct field conversion with targeted service usage for complex business logic.
    WEBHOOK PROCESSING: This should NEVER be called during webhook processing (inbound only).
    """
    # CRITICAL: Prevent outbound sync during webhook processing - FIRST CHECK
    from .sync_context import SyncContext

    sync_source = SyncContext.get_sync_source()
    full_context = SyncContext.get_context_data()
    is_webhook = SyncContext.is_webhook_processing()

    # ENHANCED DEBUGGING: Log every sync_to_sanity call attempt
    import traceback
    caller_info = ''.join(traceback.format_stack()[-3:-1])  # Get calling context
    logger.info(f"sync_to_sanity CALLED: sync_source={sync_source}, full_context={full_context}, is_webhook_processing={is_webhook}")
    logger.debug(f"sync_to_sanity called from:\n{caller_info}")

    # ABSOLUTE PREVENTION: If any inbound sync context is detected, immediately return
    if sync_source:
        # Check for any inbound sync indicators
        inbound_indicators = [
            'from_sanity',
            'domain_sync_from_sanity',
            'user_domain_sync_from_sanity',
            'product_domain_sync_from_sanity',
            'store_domain_sync_from_sanity',
            'promotional_campaign_domain_sync_from_sanity',
            'webhook',
            'sanity_webhook',
            'webhook_processing'
        ]

        for indicator in inbound_indicators:
            if indicator in sync_source:
                logger.info(f"ABSOLUTE PREVENTION: Skipping outbound sync to Sanity - inbound context detected ({sync_source}) - INBOUND ONLY")
                return None

Source: api/nextango/apps/integrations/sync/user_sync.py:158-195

Key Features:

  • Multi-layer loop prevention - Thread-local context + instance attributes
  • Context preservation - Maintains webhook processing state across call stack
  • Debugging support - Stack traces for troubleshooting circular sync
  • Early return - Fails fast when inbound context detected

User Role Array Construction

Business logic for nested userRoles array:

def _construct_user_roles_array(self, user_instance) -> List[Dict[str, Any]]:
    """
    User domain business logic: construct userRoles array from role models.
    Direct implementation for production reliability.
    """
    user_roles = []

    try:
        # Get all role types from user's computed methods
        all_role_types = user_instance.get_all_roles() if hasattr(user_instance, 'get_all_roles') else []

        role_type_mapping = {
            'employee': ('roleEmployee', 'employee_roles'),
            'customer': ('roleCustomer', 'customer_roles'),
            'manager': ('roleManager', 'manager_roles'),
            'influencer': ('roleInfluencer', 'influencer_roles'),
            'vip': ('roleVip', 'vip_roles')
        }

        for role_type in all_role_types:
            if role_type not in role_type_mapping:
                continue

            sanity_role_type, relation_name = role_type_mapping[role_type]

            # Direct role instance access - no service overhead
            if hasattr(user_instance, relation_name):
                role_instance = getattr(user_instance, relation_name).first()
                if role_instance:
                    role_object = self._build_user_role_object_direct(
                        role_type, role_instance, sanity_role_type, user_instance
                    )
                    if role_object:
                        user_roles.append(role_object)

Source: api/nextango/apps/integrations/sync/user_sync.py:390-422

Inbound Sync (sync_from_sanity)

Unified webhook processing through serializers:

@monitor_sync_operation
def sync_from_sanity(self, sanity_document: Dict[str, Any], model_class_key: str) -> Optional[Any]:
    """
    Unified sync method: Routes webhook data through UserSerializer for consistent processing.
    Eliminates architectural mismatch between webhook and API processing paths.
    Follows proven store_sync serializer delegation pattern.
    WEBHOOK PROCESSING: This method is INBOUND ONLY during webhook processing.
    """
    # CRITICAL: Detect webhook processing context at method entry
    from .sync_context import SyncContext

    is_webhook_processing = SyncContext.is_webhook_processing()
    if is_webhook_processing:
        logger.info(f"sync_from_sanity called during webhook processing - INBOUND ONLY mode for {model_class_key}")

    # ... validation checks ...

    try:
        with sync_context('user_domain_sync_from_sanity'):
            # Get model class
            app_label, model_name = model_class_key.split('.')
            model_class = apps.get_model(app_label, model_name)

            # Extract Sanity ID
            sanity_id = sanity_document.get('_id')

            # Resolve references before serialization
            normalized_data = self._resolve_webhook_references(sanity_document)

            # Map camelCase webhook fields to snake_case Django fields
            normalized_data = self._map_webhook_fields_to_django(normalized_data)

            # Use update_or_create() for idempotent, race-condition-safe record handling
            django_instance, is_created = model_class.objects.update_or_create(
                sanity_id=sanity_id,
                defaults={}  # Minimal defaults - serializer will populate fields next
            )

Source: api/nextango/apps/integrations/sync/user_sync.py:607-660

Inbound Flow:

  1. Context detection - Check if called during webhook processing
  2. Set sync context - 'user_domain_sync_from_sanity'
  3. Resolve references - Convert Sanity references to Django lookups
  4. Field mapping - camelCase → snake_case
  5. Idempotent create - update_or_create() prevents race conditions
  6. Serializer routing - Delegate to UserSerializer for validation
  7. Role sync - Call sync_roles_to_model_records() for nested data
  8. Admin access - Process adminAccess settings from manager roles

Reference Resolution

def _resolve_webhook_references(self, sanity_document: Dict[str, Any]) -> Dict[str, Any]:
    """
    Resolve Sanity references in webhook data before serialization.
    Handles userRoles array and reference objects following store_sync pattern.
    """
    normalized_data = sanity_document.copy()

    # Handle userRoles references if present
    if 'userRoles' in normalized_data and isinstance(normalized_data['userRoles'], list):
        resolved_user_roles = []
        for role in normalized_data['userRoles']:
            if isinstance(role, dict):
                # Handle storeAssignments references within roles
                # Preserve full Sanity reference structure for proper round-trip sync
                if 'storeAssignments' in role and isinstance(role['storeAssignments'], list):
                    resolved_assignments = []
                    for assignment in role['storeAssignments']:
                        if isinstance(assignment, dict) and '_ref' in assignment:
                            # Preserve full reference object with _type, _ref, and _key
                            resolved_assignments.append(assignment)
                        elif isinstance(assignment, str):
                            # Convert string ID to proper Sanity reference format
                            resolved_assignments.append({
                                '_type': 'reference',
                                '_ref': assignment
                            })
                    role['storeAssignments'] = resolved_assignments
                resolved_user_roles.append(role)
        normalized_data['userRoles'] = resolved_user_roles

    return normalized_data

Source: api/nextango/apps/integrations/sync/user_sync.py:807-838

Singleton Instance

# Singleton instance for direct operational usage
user_sync_handler = UserSyncHandler()

Source: api/nextango/apps/integrations/sync/user_sync.py:920

Other Domain Handlers

Similar patterns exist for all domains:

  • product_sync.py - Products and variants
  • store_sync.py - Store locations
  • transaction_sync.py - Sales transactions
  • payment_sync.py - Payment processing
  • promotional_campaign_sync.py - Marketing campaigns
  • zones_sync.py - Inventory zones
  • task_sync.py - Background tasks
  • hardware_sync.py - POS hardware integration

All follow the same architecture:

  1. Inherit from BaseSyncHandler
  2. Load domain-specific configuration
  3. Register reference logic with global ReferenceResolver
  4. Implement sync_to_sanity() and sync_from_sanity()
  5. Provide singleton instance

Service Layer

FieldTransformer Service

File: api/nextango/apps/integrations/sync/services/field_transformer.py

Stateless service for field transformations between Django and Sanity formats.

Type Transformations

class FieldTransformer:
    """
    Stateless service for field transformations between Django and Sanity formats.
    All methods are static to avoid dependency issues and improve testability.
    """

    @staticmethod
    def transform_decimal_to_float(decimal_value: Optional[Decimal]) -> Optional[float]:
        """Transform Decimal to float for Sanity serialization."""
        if decimal_value is None:
            return None
        return float(decimal_value)

    @staticmethod
    def transform_datetime_to_iso_string(datetime_value: Optional[datetime]) -> Optional[str]:
        """Transform datetime to ISO string for Sanity serialization."""
        if datetime_value is None:
            return None
        return datetime_value.isoformat()

    @staticmethod
    def transform_iso_string_to_datetime(iso_string: Optional[str]) -> Optional[datetime]:
        """Transform ISO string to datetime for Django model storage."""
        if not iso_string:
            return None
        try:
            return datetime.fromisoformat(iso_string.replace('Z', '+00:00'))
        except ValueError as e:
            logger.warning(f"Failed to parse datetime string '{iso_string}': {e}")
            return None

Source: api/nextango/apps/integrations/sync/services/field_transformer.py:18-64

Supported Transformations:

  • Decimal ↔ Float - Financial precision
  • Datetime/Date ↔ ISO String - Temporal data
  • Boolean ↔ String - Status fields
  • JSONField - Pass-through for dict/list

Field Mapping Transformations

@staticmethod
def transform_django_fields_to_sanity(django_data: Dict[str, Any],
                                     field_mappings: Dict[str, str]) -> Dict[str, Any]:
    """
    Transform Django model data to Sanity document format using field mappings.

    Args:
        django_data: Dictionary of Django field names and values
        field_mappings: Mapping from Django field names (snake_case) to Sanity field names (camelCase)

    Returns:
        Dictionary with Sanity field names and transformed values
    """
    sanity_data = {}

    for django_field_name, value in django_data.items():
        # Get Sanity field name from mapping, or use original name if no mapping
        sanity_field_name = field_mappings.get(django_field_name, django_field_name)

        # Apply appropriate transformation based on value type
        transformed_value = FieldTransformer._transform_value_for_sanity(value)

        if transformed_value is not None:
            sanity_data[sanity_field_name] = transformed_value

    return sanity_data

Source: api/nextango/apps/integrations/sync/services/field_transformer.py:110-134

Automatic Field Mapping Generation

@staticmethod
def create_field_mappings_from_model(model_class: type,
                                    excluded_fields: Optional[List[str]] = None) -> Dict[str, str]:
    """
    Create snake_case to camelCase field mappings automatically from Django model.

    Args:
        model_class: Django model class
        excluded_fields: List of field names to exclude from mappings

    Returns:
        Dictionary mapping Django field names to Sanity field names
    """
    excluded_fields = excluded_fields or []
    field_mappings = {}

    # Default excluded field patterns
    default_excluded = [
        'id', 'created_at', 'updated_at', 'deleted_at',
        'sanity_id', 'sync_status', 'last_synced_at'
    ]
    all_excluded = set(excluded_fields + default_excluded)

    for field in model_class._meta.get_fields():
        if (field.name not in all_excluded and
            not field.name.startswith('_') and
            not field.name.endswith('_set') and
            not field.name.endswith('Rel')):

            # Convert snake_case Django field name to camelCase Sanity field name
            camel_case_name = FieldTransformer._convert_snake_to_camel_case(field.name)
            field_mappings[field.name] = camel_case_name

    return field_mappings

Source: api/nextango/apps/integrations/sync/services/field_transformer.py:167-199

Case Conversion Utilities

@staticmethod
def _convert_snake_to_camel_case(snake_case_str: str) -> str:
    """Convert snake_case string to camelCase."""
    if not snake_case_str:
        return snake_case_str

    components = snake_case_str.split('_')
    return components[0] + ''.join(word.capitalize() for word in components[1:])

@staticmethod
def _convert_camel_to_snake_case(camel_case_str: str) -> str:
    """Convert camelCase string to snake_case."""
    if not camel_case_str:
        return camel_case_str

    import re
    snake_case = re.sub('(.)([A-Z][a-z]+)', r'\1_\2', camel_case_str)
    snake_case = re.sub('([a-z0-9])([A-Z])', r'\1_\2', snake_case)
    return snake_case.lower()

Source: api/nextango/apps/integrations/sync/services/field_transformer.py:231-249


ReferenceResolver Service

File: api/nextango/apps/integrations/sync/services/reference_resolver.py

Pure orchestrator service for cross-domain reference handling using dependency inversion.

Architecture

class ReferenceResolver:
    """
    Domain-agnostic reference resolver service using dependency inversion.
    Domain handlers register business logic functions for reference operations.
    """

    def __init__(self, skip_reference_patterns: Optional[List[str]] = None):
        self.skip_reference_patterns = skip_reference_patterns or []

        # Default patterns to avoid common circular reference issues
        self.default_skip_patterns = [
            'audit_',
            'log_',
            'sync_',
            'lock_',
            'drafts.'  # Skip Sanity draft documents
        ]

        # Registry for domain-specific reference logic (dependency injection)
        self.reference_generators = {}  # schema_type -> generator_function
        self.reference_resolvers = {}   # schema_type -> resolver_function
        self.registered_domains = set()

Source: api/nextango/apps/integrations/sync/services/reference_resolver.py:26-53

Design Philosophy:

  • No domain knowledge - Pure orchestrator
  • Explicit registration required - Fails loudly if logic not registered
  • Dependency inversion - Business logic injected from domain handlers

Reference Logic Registration

def register_reference_logic(self, schema_type: str, generator_func: Callable, resolver_func: Callable, domain: str = None):
    """
    Register domain-specific reference generation and resolution logic.

    Args:
        schema_type: Sanity schema type (e.g., 'userRole', 'store', 'product')
        generator_func: Function(django_instance) -> reference_string
        resolver_func: Function(reference_string, model_class) -> django_instance
        domain: Optional domain identifier for tracking
    """
    self.reference_generators[schema_type] = generator_func
    self.reference_resolvers[schema_type] = resolver_func

    if domain:
        self.registered_domains.add(domain)

    logger.debug(f"Registered reference logic for schema '{schema_type}' from domain '{domain}'")

Source: api/nextango/apps/integrations/sync/services/reference_resolver.py:55-71

Example Registration (from user_reference_logic.py):

def register_user_reference_logic(resolver: ReferenceResolver):
    """Register user domain reference generation and resolution logic."""

    def generate_user_reference(user_instance) -> str:
        """Generate Sanity reference ID for User model."""
        return f"user_{user_instance.pk}"

    def resolve_user_reference(reference_id: str, model_class) -> Optional[Any]:
        """Resolve Sanity reference ID to User model instance."""
        if not reference_id.startswith('user_'):
            return None
        pk = reference_id.split('_', 1)[1]
        return model_class.objects.filter(pk=pk).first()

    resolver.register_reference_logic(
        schema_type='user',
        generator_func=generate_user_reference,
        resolver_func=resolve_user_reference,
        domain='users'
    )

Creating Sanity References

def create_sanity_reference(self, django_instance: models.Model,
                           sanity_type_mapping: Optional[Dict[str, str]] = None) -> Dict[str, str]:
    """
    Create Sanity reference object using registered domain-specific logic.
    Fails explicitly if no business logic registered for schema type.

    Raises:
        ReferenceLogicNotRegisteredError: When no generator registered for schema type
    """
    if not django_instance or not hasattr(django_instance, '_meta'):
        raise ValueError("Invalid Django instance provided for reference creation")

    # Determine schema type for this instance
    model_key = f"{django_instance._meta.app_label}.{django_instance.__class__.__name__}"
    schema_type = None

    if sanity_type_mapping:
        schema_type = sanity_type_mapping.get(model_key)

    if not schema_type:
        raise ReferenceLogicNotRegisteredError(
            f"No schema type mapping found for model {model_key}. "
            f"Register domain logic with register_reference_logic() method."
        )

    # Require registered generators - no fallback
    if schema_type not in self.reference_generators:
        raise ReferenceLogicNotRegisteredError(
            f"No reference generator registered for schema type '{schema_type}'. "
            f"Available schemas: {list(self.reference_generators.keys())}"
        )

    reference_id = self.reference_generators[schema_type](django_instance)

    if self._should_skip_reference(reference_id):
        logger.debug(f"Skipping reference creation for {reference_id} due to skip patterns")
        return None

    return {
        "_type": "reference",
        "_ref": reference_id
    }

Source: api/nextango/apps/integrations/sync/services/reference_resolver.py:81-169

Resolving Sanity References

def resolve_sanity_reference(self, reference_data: Union[Dict[str, Any], str],
                           target_model_class: Type[models.Model]) -> models.Model:
    """
    Resolve Sanity reference to Django model instance using registered logic.

    Raises:
        ReferenceResolverNotRegisteredError: When no resolver registered for schema type
    """
    # Extract reference ID from reference object or use string directly
    if isinstance(reference_data, dict):
        reference_id = reference_data.get('_ref')
    elif isinstance(reference_data, str):
        reference_id = reference_data
    else:
        raise ValueError(f"Invalid reference data format: {type(reference_data)}")

    # Determine schema type from reference ID or model
    schema_type = self._determine_schema_type_from_reference(reference_id, target_model_class)

    if not schema_type:
        raise ReferenceResolverNotRegisteredError(
            f"Cannot determine schema type for reference '{reference_id}'. "
            f"Register domain logic with register_reference_logic() method."
        )

    # Require registered resolvers - no fallback
    if schema_type not in self.reference_resolvers:
        raise ReferenceResolverNotRegisteredError(
            f"No reference resolver registered for schema type '{schema_type}'. "
            f"Available schemas: {list(self.reference_resolvers.keys())}"
        )

    resolved_instance = self.reference_resolvers[schema_type](reference_id, target_model_class)

    if not resolved_instance:
        raise ReferenceResolverNotRegisteredError(
            f"Registered resolver for '{schema_type}' could not resolve reference '{reference_id}'"
        )

    return resolved_instance

Source: api/nextango/apps/integrations/sync/services/reference_resolver.py:177-236

Reference Arrays

def create_reference_array(self, django_instances: List[models.Model],
                          sanity_type_mapping: Optional[Dict[str, str]] = None) -> List[Dict[str, str]]:
    """
    Create array of Sanity references from list of Django instances.

    Returns:
        List of Sanity reference objects (excludes None values)
    """
    if not django_instances:
        return []

    references = []
    for instance in django_instances:
        reference = self.create_sanity_reference(instance, sanity_type_mapping)
        if reference:
            # Add _key for Sanity array items - use the _ref value as a unique key
            reference['_key'] = reference['_ref']
            references.append(reference)

    return references

Source: api/nextango/apps/integrations/sync/services/reference_resolver.py:263-284

Global Singleton

# Global singleton instance for cross-domain reference resolution
_global_reference_resolver = None

def get_global_reference_resolver() -> ReferenceResolver:
    """Get the global singleton reference resolver for cross-domain reference handling."""
    global _global_reference_resolver
    if _global_reference_resolver is None:
        _global_reference_resolver = ReferenceResolver()
        logger.info("Created global reference resolver singleton")
    return _global_reference_resolver

Source: api/nextango/apps/integrations/sync/services/reference_resolver.py:410-420


SyncOrchestrator Service

File: api/nextango/apps/integrations/sync/services/sync_orchestrator.py

Workflow coordinator for sync operations with transaction management and error recovery.

Orchestration Result

class SyncOperationResult:
    """
    Container for sync operation results with consistent structure.
    """

    def __init__(self, success: bool, message: str,
                 instance_id: Optional[Any] = None,
                 sanity_id: Optional[str] = None,
                 error_details: Optional[Dict[str, Any]] = None):
        self.success = success
        self.message = message
        self.instance_id = instance_id
        self.sanity_id = sanity_id
        self.error_details = error_details or {}
        self.timestamp = timezone.now()

    def to_dict(self) -> Dict[str, Any]:
        """Convert result to dictionary for logging/serialization."""
        return {
            'success': self.success,
            'message': self.message,
            'instance_id': self.instance_id,
            'sanity_id': self.sanity_id,
            'error_details': self.error_details,
            'timestamp': self.timestamp.isoformat()
        }

Source: api/nextango/apps/integrations/sync/services/sync_orchestrator.py:22-47

Orchestrator Configuration

class SyncOrchestrator:
    """
    Orchestrates sync operations across domain handlers.
    Handles workflow coordination, transaction management, and error recovery.
    """

    def __init__(self, enable_transaction_rollback: bool = True,
                 enable_audit_logging: bool = True):
        """
        Initialize sync orchestrator with configuration options.

        Args:
            enable_transaction_rollback: Whether to use database transactions for rollback
            enable_audit_logging: Whether to create audit log entries
        """
        self.enable_transaction_rollback = enable_transaction_rollback
        self.enable_audit_logging = enable_audit_logging
        self.registered_handlers = {}

Source: api/nextango/apps/integrations/sync/services/sync_orchestrator.py:50-67

Outbound Orchestration

@monitor_sync_operation
def orchestrate_sync_to_sanity(self, django_instance: models.Model,
                               domain_handler: Any,
                               created: bool = False,
                               sync_context_source: str = "orchestrated_sync") -> SyncOperationResult:
    """
    Orchestrate sync operation from Django to Sanity with error recovery.

    Returns:
        SyncOperationResult with success/failure details
    """
    model_key = f"{django_instance._meta.app_label}.{django_instance.__class__.__name__}"
    instance_id = django_instance.pk

    try:
        with self._sync_transaction_context():
            with sync_context(sync_context_source):
                # Delegate to domain handler
                sanity_id = domain_handler.sync_to_sanity(django_instance, created=created)

                # Create success result
                result = SyncOperationResult(
                    success=True,
                    message=f"Successfully synced {model_key} instance {instance_id}",
                    instance_id=instance_id,
                    sanity_id=sanity_id
                )

                # Log success if enabled
                if self.enable_audit_logging:
                    SyncAuditLogger.log_sync_success(
                        model_key=model_key,
                        instance_id=instance_id,
                        direction='outbound',
                        sanity_id=sanity_id,
                        additional_message=sync_context_source
                    )

                return result

    except Exception as e:
        # Handle sync failure
        error_message = str(e)
        logger.error(f"Sync orchestration failed for {model_key} {instance_id}: {error_message}")

        result = SyncOperationResult(
            success=False,
            message=f"Failed to sync {model_key} instance {instance_id}: {error_message}",
            instance_id=instance_id,
            error_details={'exception_type': type(e).__name__, 'exception_message': error_message}
        )

        # Log failure if enabled
        if self.enable_audit_logging:
            SyncAuditLogger.log_sync_failure(
                model_key=model_key,
                instance_id=instance_id,
                direction='outbound',
                error_message=error_message
            )

        return result

Source: api/nextango/apps/integrations/sync/services/sync_orchestrator.py:84-152

Bulk Orchestration

def orchestrate_bulk_sync(self, instances_and_handlers: List[Dict[str, Any]],
                         sync_context_source: str = "bulk_orchestrated_sync") -> Dict[str, List[SyncOperationResult]]:
    """
    Orchestrate bulk sync operations with error isolation.

    Args:
        instances_and_handlers: List of dicts with 'instance', 'handler', and optional 'created' keys

    Returns:
        Dictionary with 'successful' and 'failed' result lists
    """
    successful_results = []
    failed_results = []

    for item in instances_and_handlers:
        django_instance = item.get('instance')
        domain_handler = item.get('handler')
        created = item.get('created', False)

        # Orchestrate individual sync (with its own transaction/error handling)
        result = self.orchestrate_sync_to_sanity(
            django_instance=django_instance,
            domain_handler=domain_handler,
            created=created,
            sync_context_source=f"{sync_context_source}_item"
        )

        if result.success:
            successful_results.append(result)
        else:
            failed_results.append(result)

    logger.info(f"Bulk sync completed: {len(successful_results)} successful, {len(failed_results)} failed")

    return {
        'successful': successful_results,
        'failed': failed_results
    }

Source: api/nextango/apps/integrations/sync/services/sync_orchestrator.py:225-272

Conflict Resolution

Hardware operations take precedence over all other operations:

def orchestrate_sync_with_conflict_resolution(self, django_instance: models.Model,
                                             domain_handler: Any,
                                             operation: str = 'update',
                                             source: str = 'django') -> SyncOperationResult:
    """
    Orchestrate sync with conflict resolution capabilities.

    Resolution Priority:
    1. Hardware operations always win
    2. Most recent timestamp wins
    3. Source system wins (if timestamps very close)
    """
    # Check for hardware precedence
    is_hardware_op = self._is_hardware_operation(django_instance, source)

    try:
        with self._sync_transaction_context():
            # Get current states for conflict detection
            django_state = self._get_django_state(django_instance)
            sanity_state = self._get_sanity_state(django_instance)

            # Detect conflicts
            conflicts = self._detect_conflicts(django_state, sanity_state, is_hardware_op)

            if conflicts:
                # Resolve conflicts using business rules
                resolution = self._resolve_conflicts(conflicts, django_state, sanity_state, is_hardware_op, source)
                logger.info(f"Resolved {len(conflicts)} conflicts: {resolution['reason']}")

Source: api/nextango/apps/integrations/sync/services/sync_orchestrator.py:274-339

Conflict Resolution Rules:

def _resolve_conflicts(self, conflicts: List[Dict[str, Any]],
                      django_state: Dict[str, Any],
                      sanity_state: Dict[str, Any],
                      is_hardware_operation: bool,
                      source: str) -> Dict[str, Any]:
    """
    Resolution Priority:
    1. Hardware operations always win
    2. Most recent timestamp wins
    3. Source system wins (if timestamps very close)
    """
    resolution = {
        'winner': None,
        'reason': '',
        'merged_state': None
    }

    # Rule 1: Hardware operations always take precedence
    if is_hardware_operation:
        resolution['winner'] = 'django'
        resolution['reason'] = 'Hardware operation takes precedence'
        return resolution

    # Rule 2: Most recent modification wins
    django_timestamp = django_state.get('last_modified')
    sanity_timestamp = sanity_state.get('last_modified')

    if django_timestamp and sanity_timestamp:
        time_diff = (django_timestamp - sanity_timestamp).total_seconds()

        if abs(time_diff) > 5:  # Clear winner by timestamp
            if time_diff > 0:
                resolution['winner'] = 'django'
                resolution['reason'] = f'Django more recent by {time_diff:.1f}s'
            else:
                resolution['winner'] = 'sanity'
                resolution['reason'] = f'Sanity more recent by {abs(time_diff):.1f}s'
        else:
            # Rule 3: Source system wins if timestamps are very close
            resolution['winner'] = source if source in ['django', 'sanity'] else 'django'
            resolution['reason'] = f'Source system ({source}) wins on close timestamps'

    return resolution

Source: api/nextango/apps/integrations/sync/services/sync_orchestrator.py:407-451


Domain-Specific Reference Logic

Location: api/nextango/apps/integrations/sync/services/

Each domain provides reference logic registration for the global ReferenceResolver:

Available Reference Logic Modules

  • user_reference_logic.py - User, Employee, Customer, Manager references
  • product_reference_logic.py - Product and ProductVariant references
  • store_reference_logic.py - Store location references
  • transaction_reference_logic.py - Transaction references
  • zones_reference_logic.py - Inventory zone references
  • task_reference_logic.py - Background task references

Pattern

Each module exports a register_*_reference_logic() function:

def register_user_reference_logic(resolver: ReferenceResolver):
    """Register user domain reference logic with global resolver."""

    # Define generator
    def generate_user_reference(user_instance) -> str:
        return f"user_{user_instance.pk}"

    # Define resolver
    def resolve_user_reference(reference_id: str, model_class):
        pk = reference_id.split('_', 1)[1]
        return model_class.objects.filter(pk=pk).first()

    # Register with resolver
    resolver.register_reference_logic(
        schema_type='user',
        generator_func=generate_user_reference,
        resolver_func=resolve_user_reference,
        domain='users'
    )

This pattern ensures:

  • Domain isolation - Business logic stays in domain handlers
  • Global coordination - Single resolver handles cross-domain references
  • Explicit registration - No magic, fails loudly if not registered
  • Testability - Pure functions easy to unit test

Integration Patterns

Signal → Sync Handler Flow

# Domain signal (users/signals.py)
@receiver(post_save, sender=User)
def sync_user_to_sanity(sender, instance, created, **kwargs):
    """Sync User to Sanity after save."""
    if should_skip_sync_to_sanity():
        return

    try:
        user_sync_handler.sync_to_sanity(instance, created=created)
    except Exception as e:
        logger.error(f"Failed to sync user {instance.pk}: {e}")

Webhook → Sync Handler Flow

# Webhook view (integrations/views.py)
def process_webhook(request):
    """Process incoming Sanity webhook."""
    payload = json.loads(request.body)

    with webhook_processing_context():
        # Determine document type and route to handler
        doc_type = payload.get('_type')

        if doc_type == 'user':
            user_sync_handler.sync_from_sanity(payload, 'users.User')
        elif doc_type == 'product':
            product_sync_handler.sync_from_sanity(payload, 'products.Product')
        # ... etc

Service Layer Usage

# Domain handler using services
class ProductSyncHandler(BaseSyncHandler):
    def __init__(self):
        super().__init__()
        self.field_transformer = FieldTransformer()
        self.reference_resolver = get_global_reference_resolver()

        # Register domain reference logic
        register_product_reference_logic(self.reference_resolver)

    def sync_to_sanity(self, instance, created=False):
        # Transform fields
        sanity_data = self.field_transformer.transform_django_fields_to_sanity(
            django_data={'name': instance.name, 'price': instance.price},
            field_mappings={'name': 'name', 'price': 'price'}
        )

        # Create references
        if instance.category:
            sanity_data['category'] = self.reference_resolver.create_sanity_reference(
                instance.category,
                sanity_type_mapping={'products.Category': 'category'}
            )

Usage Examples

Registering a New Model

from nextango.apps.integrations.sync.universal_sync import get_universal_sync

universal_sync = get_universal_sync()

universal_sync.register_model(
    model_class=MyModel,
    sanity_type='myModelType',
    field_mappings={
        'my_field': 'myField',
        'another_field': 'anotherField'
    }
)

Bulk Sync Operation

from nextango.apps.integrations.sync.universal_sync import get_universal_sync

universal_sync = get_universal_sync()

result = universal_sync.bulk_sync_to_sanity(
    model_key='users.User',
    queryset=User.objects.filter(is_active=True)
)

print(f"Synced: {result['synced']}, Failed: {result['failed']}")

Using Field Transformer

from nextango.apps.integrations.sync.services.field_transformer import FieldTransformer

# Transform single value
iso_string = FieldTransformer.transform_datetime_to_iso_string(user.created_at)

# Transform dictionary
sanity_data = FieldTransformer.transform_django_fields_to_sanity(
    django_data={'created_at': user.created_at, 'price': Decimal('19.99')},
    field_mappings={'created_at': 'createdAt', 'price': 'price'}
)
# Result: {'createdAt': '2025-11-17T10:30:00', 'price': 19.99}

Using Reference Resolver

from nextango.apps.integrations.sync.services.reference_resolver import get_global_reference_resolver

resolver = get_global_reference_resolver()

# Create reference
user_ref = resolver.create_sanity_reference(
    user_instance,
    sanity_type_mapping={'users.User': 'user'}
)
# Result: {'_type': 'reference', '_ref': 'user_123'}

# Resolve reference
user = resolver.resolve_sanity_reference(
    reference_data={'_ref': 'user_123'},
    target_model_class=User
)

Using Sync Orchestrator

from nextango.apps.integrations.sync.services.sync_orchestrator import SyncOrchestrator
from nextango.apps.integrations.sync.user_sync import user_sync_handler

orchestrator = SyncOrchestrator(
    enable_transaction_rollback=True,
    enable_audit_logging=True
)

result = orchestrator.orchestrate_sync_to_sanity(
    django_instance=user,
    domain_handler=user_sync_handler,
    created=True,
    sync_context_source='api_create'
)

if result.success:
    print(f"Synced successfully: {result.sanity_id}")
else:
    print(f"Sync failed: {result.message}")

Cross-Domain Integration

Global Reference Resolver

All domain handlers share a single ReferenceResolver instance:

# Each domain handler initialization
from nextango.apps.integrations.sync.services.reference_resolver import get_global_reference_resolver
from nextango.apps.integrations.sync.services.user_reference_logic import register_user_reference_logic

class UserSyncHandler(BaseSyncHandler):
    def __init__(self):
        super().__init__()
        self.reference_resolver = get_global_reference_resolver()
        register_user_reference_logic(self.reference_resolver)

This enables:

  • Cross-domain references - Users can reference Products, Stores reference Users, etc.
  • Single source of truth - One resolver handles all reference operations
  • Domain independence - Each domain registers its own logic

Shared Audit Trail

All sync operations use SyncAuditLog from models.py:

from nextango.apps.integrations.sync.base import SyncAuditLogger

SyncAuditLogger.log_sync_success(
    model_key='users.User',
    instance_id=user.pk,
    direction='outbound',
    sanity_id='user_123'
)

Provides centralized monitoring across all domains.


  • Models - Infrastructure models (WebhookEvent, SyncAuditLog, SyncLock, OutboundMutationLog)
  • SyncContext - Thread-local loop prevention mechanism
  • Views - Webhook endpoint handlers
  • Users Domain Services - User-specific business logic
  • README - Domain overview

Note: This documentation reflects the consolidated architecture with BaseSyncHandler inheritance and targeted service usage. All domain handlers follow this pattern for consistency and maintainability.

Was this page helpful?