# Django Cronjob Utils - Architecture Document

## Overview

This document outlines the architecture and design for a Django package that provides robust cronjob handling between OS-level cronjobs and Django applications. The package addresses core requirements for execution tracking, duplicate prevention, stakeholder notifications, configurable retry, and comprehensive database logging.

## Core Requirements

1. **Execution Tracking**: Record all cronjob executions with start/end times, success/failure status
2. **Duplicate Prevention**: Prevent concurrent execution of the same cronjob
3. **Stakeholder Notification**: Alert stakeholders when cronjobs fail or encounter errors
4. **Configurable Retry**: Ability to automatically rerun failed cronjobs (configurable per task)
5. **Database Logging**: All executions logged to database for monitoring and audit trail

## Design Principles

### 1. Database-Level Concurrency Control

**Approach**: Use database transactions with `SELECT FOR UPDATE` to prevent race conditions

```python
# Atomic duplicate check and record creation
with transaction.atomic():
    # Lock any existing records for this task/date
    existing = CronExecution.objects.select_for_update().filter(
        code=task_code,
        odate=execution_date,
        completed=False
    ).first()
    
    if existing:
        raise ConcurrentExecutionError("Task already running")
    
    # Create new record
    execution = CronExecution.objects.create(...)
```

### 2. Task Registry System

**Approach**: Decorative-based task registration, separating package from application code

```python
@register_task('calc-commission', 'A001', 
               retry_on_failure=True, 
               max_retries=3,
               timeout=3600)
class CalcCommissionTask(CronTask):
    def execute(self, date: date) -> dict:
        # Business logic
        return {'error': False, 'message': 'Success'}
```

### 3. Pluggable Notification System

**Approach**: Signal-based notification with multiple backends

```python
# Settings
CRONJOB_NOTIFICATIONS = {
    'on_failure': ['email', 'slack', 'telegram'],
    'on_success': [],  # Optional
    'email': {
        'recipients': ['admin@example.com'],
    },
    'slack': {
        'webhook_url': 'https://hooks.slack.com/...',
    },
    'telegram': {
        'bot_token': 'your-bot-token',
        'chat_id': 'your-chat-id',
    }
}
```

### 4. Configurable Retry Mechanism

**Approach**: Task-level configuration with automatic retry scheduling

```python
@register_task('calc-commission', 'A001', 
               retry_on_failure=True,
               max_retries=3,
               retry_delay=300)  # 5 minutes
class CalcCommissionTask(CronTask):
    pass
```

### 5. Timeout Handling

**Approach**: Per-task timeout configuration with automatic termination

```python
@register_task('long-running-task', 'A002',
               timeout=7200)  # 2 hours
class LongRunningTask(CronTask):
    pass
```

### 6. Execution Patterns

**Approach**: Built-in execution patterns with clear semantics

- **STANDARD**: Check if already executed, skip if yes
- **ALWAYS**: Always execute (no duplicate check)
- **RERUN_ON_FAILURE**: Only skip if previous execution succeeded
- **RATE_LIMITED**: Use locking to prevent concurrent execution

## Architecture

### Package Structure

```
django_cronjob_utils/
├── __init__.py
├── models.py              # CronExecution model
├── base.py                # CronTask base class
├── registry.py            # Task registry
├── decorators.py          # @register_task decorator
├── exceptions.py          # Custom exceptions
├── notifications.py       # Notification system
├── locks.py               # Concurrency control
├── management/
│   └── commands/
│       └── run_cron_task.py
├── admin.py               # Django admin integration
└── migrations/
    └── 0001_initial.py
```

### Component Overview

#### 1. Models (`models.py`)

**CronExecution**: Main execution tracking model

```python
class CronExecution(models.Model):
    task_code = models.CharField(max_length=50, db_index=True)
    task_name = models.CharField(max_length=100, db_index=True)
    started = models.DateTimeField(auto_now_add=True)
    ended = models.DateTimeField(null=True, blank=True)
    completed = models.BooleanField(default=False, db_index=True)
    success = models.BooleanField(default=False)
    message = models.TextField(blank=True)
    error_code = models.CharField(max_length=50, blank=True)
    execution_date = models.DateField(db_index=True)  # odate equivalent
    retry_count = models.IntegerField(default=0)
    pid = models.IntegerField(null=True, blank=True)  # Process ID for timeout handling
    
    class Meta:
        indexes = [
            models.Index(fields=['task_code', 'execution_date']),
            models.Index(fields=['task_code', 'success', 'execution_date']),
            models.Index(fields=['completed']),
        ]
```

**Key Features**:
- Boolean fields for clear status tracking
- Added `task_name` for human-readable identification
- Added `retry_count` for tracking retries
- Added `pid` for process management (timeout handling)
- Optimized indexing strategy for query performance

#### 2. Base Class (`base.py`)

**CronTask**: Abstract base class for all cron tasks

```python
class CronTask(ABC):
    task_name: str
    task_code: str
    execution_pattern: str = ExecutionPattern.STANDARD
    retry_on_failure: bool = False
    max_retries: int = 0
    retry_delay: int = 300  # seconds
    timeout: Optional[int] = None  # seconds
    
    def __init__(self, execution_date: date, **options):
        self.execution_date = execution_date
        self.options = options
        self.execution = None
    
    def run(self) -> ExecutionResult:
        """Main execution method with full lifecycle management"""
        # 1. Validate input
        self.validate()
        
        # 2. Check if should run (duplicate check)
        if not self.should_run():
            return ExecutionResult(skipped=True, reason="Already executed")
        
        # 3. Acquire lock (prevent concurrent execution)
        with self.acquire_lock():
            # 4. Create execution record
            self.execution = self.create_execution_record()
            
            try:
                # 5. Execute with timeout
                result = self.execute_with_timeout()
                
                # 6. Update record
                self.execution.mark_completed(
                    success=result.success,
                    message=result.message,
                    error_code=result.error_code
                )
                
                # 7. Handle notifications
                if not result.success:
                    self.notify_failure(result)
                
                return result
                
            except TimeoutError:
                self.execution.mark_failed(
                    message="Execution timeout",
                    error_code="TIMEOUT"
                )
                self.notify_failure(...)
                raise
            except Exception as e:
                self.execution.mark_failed(
                    message=str(e),
                    error_code=self.get_error_code(e)
                )
                self.notify_failure(...)
                
                # Handle retry
                if self.should_retry():
                    self.schedule_retry()
                
                raise
    
    @abstractmethod
    def execute(self, date: date) -> dict:
        """Business logic - must be implemented by subclasses"""
        pass
    
    def should_run(self) -> bool:
        """Check if task should run (duplicate prevention)"""
        if self.execution_pattern == ExecutionPattern.ALWAYS:
            return True
        
        if self.execution_pattern == ExecutionPattern.RERUN_ON_FAILURE:
            return not CronExecution.objects.filter(
                task_code=self.task_code,
                execution_date=self.execution_date,
                success=True
            ).exists()
        
        # STANDARD pattern
        return not CronExecution.objects.filter(
            task_code=self.task_code,
            execution_date=self.execution_date
        ).exists()
    
    def acquire_lock(self):
        """Acquire database lock to prevent concurrent execution"""
        return DatabaseLock(self.task_code, self.execution_date)
    
    def execute_with_timeout(self) -> ExecutionResult:
        """Execute with optional timeout"""
        if self.timeout:
            return timeout(self.timeout)(self._execute)()
        return self._execute()
    
    def _execute(self) -> ExecutionResult:
        """Internal execution method"""
        try:
            result = self.execute(self.execution_date)
            return ExecutionResult(
                success=not result.get('error', False),
                message=result.get('message', ''),
                error_code=result.get('error_code', '')
            )
        except Exception as e:
            return ExecutionResult(
                success=False,
                message=str(e),
                error_code=self.get_error_code(e)
            )
    
    def should_retry(self) -> bool:
        """Check if task should be retried"""
        if not self.retry_on_failure:
            return False
        
        if self.execution.retry_count >= self.max_retries:
            return False
        
        return True
    
    def schedule_retry(self):
        """Schedule retry execution"""
        # Could use Celery, Django-Q, or simple delay
        # For now, log and notify
        logger.info(f"Scheduling retry for {self.task_name}")
```

**Key Features**:
- Clear separation of concerns
- Proper error handling
- Timeout support
- Retry mechanism
- Lock-based concurrency control

#### 3. Task Registry (`registry.py`)

**TaskRegistry**: Centralized task registration and lookup

```python
class TaskRegistry:
    _tasks: Dict[str, Type[CronTask]] = {}
    _codes: Dict[str, str] = {}  # code -> name mapping
    
    @classmethod
    def register(cls, name: str, code: str, task_class: Type[CronTask], **config):
        """Register a task"""
        cls._tasks[name] = task_class
        cls._codes[code] = name
        
        # Apply configuration
        task_class.task_name = name
        task_class.task_code = code
        task_class.execution_pattern = config.get('execution_pattern', ExecutionPattern.STANDARD)
        task_class.retry_on_failure = config.get('retry_on_failure', False)
        task_class.max_retries = config.get('max_retries', 0)
        task_class.timeout = config.get('timeout')
    
    @classmethod
    def get_task(cls, name: str) -> Type[CronTask]:
        """Get task class by name"""
        if name not in cls._tasks:
            raise TaskNotFoundError(f"Task '{name}' not found")
        return cls._tasks[name]
    
    @classmethod
    def list_tasks(cls) -> List[str]:
        """List all registered task names"""
        return list(cls._tasks.keys())
```

#### 4. Decorator (`decorators.py`)

**@register_task**: Decorator for registering tasks

```python
def register_task(name: str, code: str, **config):
    """Register a cron task"""
    def decorator(cls):
        TaskRegistry.register(name, code, cls, **config)
        return cls
    return decorator
```

**Usage**:
```python
@register_task('calc-commission', 'A001',
               execution_pattern=ExecutionPattern.RERUN_ON_FAILURE,
               retry_on_failure=True,
               max_retries=3,
               timeout=3600)
class CalcCommissionTask(CronTask):
    def execute(self, date: date) -> dict:
        # Business logic
        return {'error': False, 'message': 'Success'}
```

#### 5. Concurrency Control (`locks.py`)

**DatabaseLock**: Database-level locking to prevent concurrent execution

```python
class DatabaseLock:
    """Context manager for database-level locking"""
    
    def __init__(self, task_code: str, execution_date: date):
        self.task_code = task_code
        self.execution_date = execution_date
    
    def __enter__(self):
        with transaction.atomic():
            # Lock any running executions
            running = CronExecution.objects.select_for_update().filter(
                task_code=self.task_code,
                execution_date=self.execution_date,
                completed=False
            ).first()
            
            if running:
                raise ConcurrentExecutionError(
                    f"Task {self.task_code} already running for {self.execution_date}"
                )
        
        return self
    
    def __exit__(self, exc_type, exc_val, exc_tb):
        pass
```

#### 6. Notification System (`notifications.py`)

**NotificationManager**: Pluggable notification system

```python
class NotificationManager:
    """Manages notifications for cronjob events"""
    
    def __init__(self):
        self.backends = self._load_backends()
    
    def notify_failure(self, execution: CronExecution, error: str):
        """Notify stakeholders of failure"""
        for backend in self.backends:
            try:
                backend.notify_failure(execution, error)
            except Exception as e:
                logger.error(f"Notification backend {backend} failed: {e}")
    
    def _load_backends(self) -> List[NotificationBackend]:
        """Load configured notification backends"""
        backends = []
        config = settings.CRONJOB_NOTIFICATIONS
        
        if 'email' in config.get('on_failure', []):
            backends.append(EmailBackend(config.get('email', {})))
        
        if 'slack' in config.get('on_failure', []):
            backends.append(SlackBackend(config.get('slack', {})))
        
        if 'telegram' in config.get('on_failure', []):
            backends.append(TelegramBackend(config.get('telegram', {})))
        
        return backends
```

**TelegramBackend**: Telegram notification backend

```python
class TelegramBackend(NotificationBackend):
    """Telegram notification backend"""
    
    def __init__(self, config: dict):
        self.bot_token = config.get('bot_token')
        self.chat_id = config.get('chat_id')
        self.api_url = f"https://api.telegram.org/bot{self.bot_token}/sendMessage"
    
    def notify_failure(self, execution: CronExecution, error_message: str):
        """Send failure notification via Telegram"""
        message = (
            f"❌ Cronjob Failed\n\n"
            f"Task: {execution.task_name} ({execution.task_code})\n"
            f"Date: {execution.execution_date}\n"
            f"Error: {error_message}\n"
            f"Started: {execution.started}\n"
            f"Error Code: {execution.error_code or 'N/A'}"
        )
        
        payload = {
            'chat_id': self.chat_id,
            'text': message,
            'parse_mode': 'HTML'
        }
        
        try:
            response = requests.post(self.api_url, json=payload, timeout=10)
            response.raise_for_status()
        except Exception as e:
            logger.error(f"Failed to send Telegram notification: {e}")
```

#### 7. Management Command (`management/commands/run_cron_task.py`)

**RunCronTaskCommand**: Entry point for OS cronjobs

```python
class Command(BaseCommand):
    help = 'Run a registered cron task'
    
    def add_arguments(self, parser):
        parser.add_argument('task_name', type=str, help='Task name')
        parser.add_argument('date', type=str, help='Execution date (YYYY-MM-DD)')
        parser.add_argument('--force', action='store_true', help='Force execution even if already ran')
        parser.add_argument('--rerun', action='store_true', help='Rerun even if successful')
    
    def handle(self, *args, **options):
        task_name = options['task_name']
        execution_date = date.fromisoformat(options['date'])
        
        try:
            task_class = TaskRegistry.get_task(task_name)
            task = task_class(execution_date, force=options['force'], rerun=options['rerun'])
            result = task.run()
            
            if result.success:
                self.stdout.write(self.style.SUCCESS(f"Task '{task_name}' completed successfully"))
                return 0
            else:
                self.stdout.write(self.style.ERROR(f"Task '{task_name}' failed: {result.message}"))
                return 1
                
        except TaskNotFoundError:
            self.stdout.write(self.style.ERROR(f"Task '{task_name}' not found"))
            return 1
        except ConcurrentExecutionError as e:
            self.stdout.write(self.style.WARNING(str(e)))
            return 2
```

## Execution Flow

```
1. OS Cronjob Trigger
   ↓
2. Management Command (run_cron_task)
   ↓
3. Task Registry Lookup
   ↓
4. Task Instantiation
   ↓
5. Validation (date format, etc.)
   ↓
6. Duplicate Check (based on execution pattern)
   ↓
7. Acquire Database Lock (prevent concurrent execution)
   ↓
8. Create Execution Record (status: running)
   ↓
9. Execute Business Logic (with optional timeout)
   ↓
10. Update Execution Record (success/failure)
    ↓
11. Handle Notifications (if failure)
    ↓
12. Handle Retry (if configured and failed)
    ↓
13. Release Lock
```

## Configuration

### Django Settings

```python
# settings.py

INSTALLED_APPS = [
    # ...
    'django_cronjob_utils',
]

# Cronjob Utils Configuration
CRONJOB_UTILS = {
    # Notification settings
    'NOTIFICATIONS': {
        'on_failure': ['email', 'slack', 'telegram'],  # Backends to use on failure
        'on_success': [],  # Backends to use on success (optional)
        'email': {
            'recipients': ['admin@example.com'],
            'from_email': 'noreply@example.com',
        },
        'slack': {
            'webhook_url': 'https://hooks.slack.com/services/...',
            'channel': '#cronjobs',
        },
        'telegram': {
            'bot_token': 'your-bot-token',
            'chat_id': 'your-chat-id',
        },
    },
    
    # Default task settings
    'DEFAULT_TIMEOUT': 3600,  # 1 hour
    'DEFAULT_RETRY_ON_FAILURE': False,
    'DEFAULT_MAX_RETRIES': 0,
    
    # Admin settings
    'ENABLE_ADMIN': True,
    
    # Logging
    'LOG_LEVEL': 'INFO',
}
```

## Usage Examples

### Example 1: Standard Task (Run Once Per Day)

```python
from django_cronjob_utils import CronTask, register_task
from datetime import date

@register_task('calc-commission', 'A001')
class CalcCommissionTask(CronTask):
    def execute(self, date: date) -> dict:
        # Business logic here
        commission_service = CommissionService()
        result = commission_service.calculate(date)
        
        if result.success:
            return {'error': False, 'message': f'Processed {result.count} records'}
        else:
            return {'error': True, 'message': result.error, 'error_code': 'CALC_ERROR'}
```

**Crontab entry**:
```bash
0 1 * * * cd /path/to/project && python manage.py run_cron_task calc-commission $(date +\%Y-\%m-\%d)
```

### Example 2: Task with Retry on Failure

```python
@register_task('sync-external-api', 'A002',
               execution_pattern=ExecutionPattern.RERUN_ON_FAILURE,
               retry_on_failure=True,
               max_retries=3,
               retry_delay=300)  # 5 minutes between retries
class SyncExternalAPITask(CronTask):
    def execute(self, date: date) -> dict:
        api_client = ExternalAPIClient()
        result = api_client.sync(date)
        
        if result.success:
            return {'error': False, 'message': 'Sync completed'}
        else:
            return {'error': True, 'message': result.error}
```

### Example 3: Always Execute Task (No Duplicate Check)

```python
@register_task('update-cache', 'A003',
               execution_pattern=ExecutionPattern.ALWAYS)
class UpdateCacheTask(CronTask):
    def execute(self, date: date) -> dict:
        cache_service = CacheService()
        cache_service.refresh()
        return {'error': False, 'message': 'Cache updated'}
```

### Example 4: Rate-Limited Task (Prevent Concurrent Execution)

```python
@register_task('process-batch', 'A004',
               execution_pattern=ExecutionPattern.RATE_LIMITED,
               timeout=7200)  # 2 hour timeout
class ProcessBatchTask(CronTask):
    def execute(self, date: date) -> dict:
        batch_processor = BatchProcessor()
        count = 0
        
        while True:
            batch = batch_processor.get_next_batch()
            if not batch:
                break
            
            batch_processor.process(batch)
            count += 1
        
        return {'error': False, 'message': f'Processed {count} batches'}
```

## Monitoring and Admin

### Django Admin Integration

The package provides Django admin integration for monitoring:

```python
# admin.py
from django_cronjob_utils.admin import CronExecutionAdmin
from django_cronjob_utils.models import CronExecution

admin.site.register(CronExecution, CronExecutionAdmin)
```

**Features**:
- View execution history
- Filter by task, date, success status
- View error messages and codes
- Manual retry capability
- Export to CSV

### Querying Executions

```python
from django_cronjob_utils.models import CronExecution
from datetime import date, timedelta

# Recent failures
failures = CronExecution.objects.filter(
    success=False,
    execution_date__gte=date.today() - timedelta(days=7)
)

# Stuck jobs (running > 1 hour)
from django.utils import timezone
stuck = CronExecution.objects.filter(
    completed=False,
    started__lt=timezone.now() - timedelta(hours=1)
)

# Task success rate
from django.db.models import Count, Q
stats = CronExecution.objects.filter(
    execution_date__gte=date.today() - timedelta(days=30)
).values('task_name').annotate(
    total=Count('id'),
    successful=Count('id', filter=Q(success=True)),
    failed=Count('id', filter=Q(success=False))
)
```

## Testing

### Unit Tests

```python
from django.test import TestCase
from django_cronjob_utils import CronTask, register_task
from datetime import date

@register_task('test-task', 'T001')
class TestTask(CronTask):
    def execute(self, date: date) -> dict:
        return {'error': False, 'message': 'Test success'}

class CronTaskTest(TestCase):
    def test_task_execution(self):
        task = TestTask(date.today())
        result = task.run()
        self.assertTrue(result.success)
    
    def test_duplicate_prevention(self):
        task = TestTask(date.today())
        task.run()
        
        # Second run should be skipped
        task2 = TestTask(date.today())
        result = task2.run()
        self.assertTrue(result.skipped)
```

## Future Enhancements

1. **Celery Integration**: Use Celery for retry scheduling
2. **Distributed Locks**: Redis-based locking for multi-server deployments
3. **Metrics**: Prometheus metrics export
4. **Dashboard**: Web-based monitoring dashboard
5. **Scheduling**: Built-in scheduling (alternative to crontab)
6. **Dependencies**: Task dependencies and execution order
7. **Parallel Execution**: Run multiple tasks in parallel

## Summary

This architecture addresses all core requirements:

✅ **Execution Tracking**: Comprehensive database logging  
✅ **Duplicate Prevention**: Database-level locking  
✅ **Stakeholder Notification**: Pluggable notification system (Email, Slack, Telegram)  
✅ **Configurable Retry**: Per-task retry configuration  
✅ **Database Logging**: All executions logged with full details  

**Key Features**:
- Eliminates race conditions with database locks
- Proper error handling and type safety
- Extensible task registration system
- Configurable retry and timeout mechanisms
- Pluggable notification system with multiple backends
- Better monitoring and admin interface
