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:
- Universal Sync Framework - Generic sync infrastructure for all model types
- Domain Sync Handlers - Domain-specific sync implementations (users, products, stores, etc.)
- 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
- Inheritance-Based Domain Handlers - All domain sync handlers inherit from
BaseSyncHandler - Dependency Inversion - Services register business logic functions rather than containing domain knowledge
- Thread-Local Context - Loop prevention using
threading.local()for webhook detection - Race Condition Handling -
update_or_create()with idempotent operations - 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 referenceconfig- 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:
- Check LOAD_TESTING_MODE environment variable
- Verify model is registered
- Primary loop prevention - Check thread-local context via
should_skip_sync_to_sanity() - Fallback check - Instance
_sync_sourceattribute - Deduplication - Cache-based locking (60s TTL)
- Model-specific locking - UserRole: 5s, User models: 10s, Others: 30s
- Transform - Django instance → Sanity document
- Sync - Create or update in Sanity
- 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:
- Context detection - Check if called during webhook processing
- Set sync context -
'user_domain_sync_from_sanity' - Resolve references - Convert Sanity references to Django lookups
- Field mapping - camelCase → snake_case
- Idempotent create -
update_or_create()prevents race conditions - Serializer routing - Delegate to UserSerializer for validation
- Role sync - Call
sync_roles_to_model_records()for nested data - Admin access - Process
adminAccesssettings 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:
- Inherit from
BaseSyncHandler - Load domain-specific configuration
- Register reference logic with global
ReferenceResolver - Implement
sync_to_sanity()andsync_from_sanity() - 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.
Related Documentation
- 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.