# Story 1.7: Shadow Mode Implementation & Metrics

**Epic:** Airflow Infrastructure & Parallel Deployment
**Status:** pending
**Priority:** Medium
**Estimate:** 6 hours

---

## User Story

As a **DevOps engineer**, I want to **implement shadow mode for Airflow** so that **we can process documents in parallel with Celery without affecting production data and compare results**.

---

## Acceptance Criteria

- [ ] Shadow mode flag disables actual storage operations
- [ ] Shadow mode logs what would be stored without persisting
- [ ] Metrics collected for processing time, chunk counts, embedding counts
- [ ] Comparison dashboard/report between Celery and Airflow outputs
- [ ] Shadow mode can be toggled via environment variable
- [ ] No production data affected during shadow testing

---

## Technical Details

### Shadow Mode Behavior

| Component | Normal Mode | Shadow Mode |
|-----------|-------------|-------------|
| Validate | Execute | Execute |
| Fetch | Execute | Execute |
| Parse | Execute | Execute |
| Chunk | Execute | Execute |
| Embed | Execute | Execute (or mock) |
| Store Chunks | Write to DB | Log only |
| Store Vectors | Write to Milvus | Log only |

### Files to Create/Modify

| File | Action | Description |
|------|--------|-------------|
| `dags/config.py` | Modify | Add shadow mode config |
| `plugins/operators/store_chunks.py` | Modify | Add shadow mode check |
| `plugins/operators/store_vectors.py` | Modify | Add shadow mode check |
| `plugins/metrics/extraction_metrics.py` | Create | Metrics collection |
| `scripts/compare_outputs.py` | Create | Output comparison tool |

---

## Implementation

### Shadow Mode Configuration

```python
# dags/config.py (additions)
import os

# Shadow mode settings
SHADOW_MODE = os.getenv('AIRFLOW_SHADOW_MODE', 'false').lower() == 'true'
SHADOW_MOCK_EMBEDDINGS = os.getenv('AIRFLOW_SHADOW_MOCK_EMBEDDINGS', 'false').lower() == 'true'

# Metrics settings
METRICS_ENABLED = True
METRICS_COLLECTION = 'extraction_metrics'
```

### Modified Store Operators

```python
# plugins/operators/store_chunks.py (shadow mode additions)
from dags.config import SHADOW_MODE

class StoreChunksOperator(BaseOperator):

    def execute(self, context):
        # ... existing code to get chunks ...

        if SHADOW_MODE:
            self.log.info(f"[SHADOW MODE] Would store {len(chunks)} chunks")
            self.log.info(f"[SHADOW MODE] Document: {document_id}")
            self.log.info(f"[SHADOW MODE] Sample chunk: {chunks[0] if chunks else 'none'}")

            # Record metrics
            self._record_metrics(context, {
                'operation': 'store_chunks',
                'chunk_count': len(chunks),
                'document_id': str(document_id),
                'mode': 'shadow',
            })

            return {
                'stored_count': len(chunks),
                'document_id': str(document_id),
                'status': 'shadow',
                'mode': 'shadow',
            }

        # ... normal execution ...

    def _record_metrics(self, context, metrics):
        """Record metrics to XCom for later analysis."""
        from plugins.metrics.extraction_metrics import MetricsCollector
        collector = MetricsCollector()
        collector.record(context['dag_run'].run_id, metrics)
```

```python
# plugins/operators/store_vectors.py (shadow mode additions)
from dags.config import SHADOW_MODE

class StoreVectorsOperator(BaseOperator):

    def execute(self, context):
        # ... existing code to get embeddings ...

        if SHADOW_MODE:
            self.log.info(f"[SHADOW MODE] Would store {len(embeddings)} vectors")
            self.log.info(f"[SHADOW MODE] Collection: {self.collection}")
            self.log.info(f"[SHADOW MODE] Document: {document_id}")

            # Record metrics
            self._record_metrics(context, {
                'operation': 'store_vectors',
                'vector_count': len(embeddings),
                'document_id': str(document_id),
                'collection': self.collection,
                'mode': 'shadow',
            })

            return {
                'stored_count': len(embeddings),
                'document_id': str(document_id),
                'collection': self.collection,
                'status': 'shadow',
                'mode': 'shadow',
            }

        # ... normal execution ...
```

### Metrics Collector

```python
# plugins/metrics/extraction_metrics.py
import json
import time
from datetime import datetime
from typing import Any
import redis

class MetricsCollector:
    """Collect and store extraction metrics for comparison."""

    def __init__(self, redis_url: str = None):
        self.redis_url = redis_url or os.getenv('REDIS_URL', 'redis://localhost:6379')
        self._redis = None

    @property
    def redis(self):
        if self._redis is None:
            self._redis = redis.from_url(self.redis_url)
        return self._redis

    def record(self, run_id: str, metrics: dict):
        """Record metrics for a DAG run."""
        key = f"airflow_metrics:{run_id}:{metrics.get('operation', 'unknown')}"
        metrics['timestamp'] = datetime.utcnow().isoformat()
        metrics['run_id'] = run_id

        self.redis.setex(
            key,
            86400 * 7,  # 7 day TTL
            json.dumps(metrics)
        )

    def get_run_metrics(self, run_id: str) -> list[dict]:
        """Get all metrics for a DAG run."""
        pattern = f"airflow_metrics:{run_id}:*"
        keys = self.redis.keys(pattern)
        metrics = []
        for key in keys:
            data = self.redis.get(key)
            if data:
                metrics.append(json.loads(data))
        return metrics

    def compare_runs(self, airflow_run_id: str, celery_task_id: str) -> dict:
        """Compare Airflow run with Celery task results."""
        airflow_metrics = self.get_run_metrics(airflow_run_id)

        # Get Celery metrics (stored separately during Celery processing)
        celery_key = f"celery_metrics:{celery_task_id}"
        celery_data = self.redis.get(celery_key)
        celery_metrics = json.loads(celery_data) if celery_data else {}

        return {
            'airflow': {
                'chunk_count': sum(m.get('chunk_count', 0) for m in airflow_metrics if m.get('operation') == 'store_chunks'),
                'vector_count': sum(m.get('vector_count', 0) for m in airflow_metrics if m.get('operation') == 'store_vectors'),
                'processing_time': self._calculate_duration(airflow_metrics),
            },
            'celery': {
                'chunk_count': celery_metrics.get('chunk_count', 0),
                'vector_count': celery_metrics.get('vector_count', 0),
                'processing_time': celery_metrics.get('processing_time', 0),
            },
            'match': self._compare_counts(airflow_metrics, celery_metrics),
        }

    def _calculate_duration(self, metrics: list) -> float:
        """Calculate total processing duration from metrics."""
        if not metrics:
            return 0
        timestamps = [m.get('timestamp') for m in metrics if m.get('timestamp')]
        if len(timestamps) < 2:
            return 0
        start = min(timestamps)
        end = max(timestamps)
        return (datetime.fromisoformat(end) - datetime.fromisoformat(start)).total_seconds()

    def _compare_counts(self, airflow_metrics: list, celery_metrics: dict) -> bool:
        """Compare chunk and vector counts between systems."""
        airflow_chunks = sum(m.get('chunk_count', 0) for m in airflow_metrics if m.get('operation') == 'store_chunks')
        celery_chunks = celery_metrics.get('chunk_count', 0)
        return airflow_chunks == celery_chunks
```

### Output Comparison Script

```python
# scripts/compare_outputs.py
#!/usr/bin/env python3
"""Compare Airflow and Celery extraction outputs."""

import argparse
import json
from plugins.metrics.extraction_metrics import MetricsCollector

def main():
    parser = argparse.ArgumentParser(description='Compare Airflow and Celery outputs')
    parser.add_argument('--airflow-run', required=True, help='Airflow DAG run ID')
    parser.add_argument('--celery-task', required=True, help='Celery task ID')
    parser.add_argument('--output', default='comparison.json', help='Output file')
    args = parser.parse_args()

    collector = MetricsCollector()
    comparison = collector.compare_runs(args.airflow_run, args.celery_task)

    print("\n=== Extraction Comparison ===\n")
    print(f"Airflow Run: {args.airflow_run}")
    print(f"Celery Task: {args.celery_task}")
    print()

    print("Airflow Results:")
    print(f"  Chunks: {comparison['airflow']['chunk_count']}")
    print(f"  Vectors: {comparison['airflow']['vector_count']}")
    print(f"  Duration: {comparison['airflow']['processing_time']:.2f}s")
    print()

    print("Celery Results:")
    print(f"  Chunks: {comparison['celery']['chunk_count']}")
    print(f"  Vectors: {comparison['celery']['vector_count']}")
    print(f"  Duration: {comparison['celery']['processing_time']:.2f}s")
    print()

    status = "MATCH" if comparison['match'] else "MISMATCH"
    print(f"Status: {status}")

    with open(args.output, 'w') as f:
        json.dump(comparison, f, indent=2)
    print(f"\nFull comparison saved to: {args.output}")

if __name__ == '__main__':
    main()
```

---

## Tasks

- [ ] Add shadow mode configuration to `dags/config.py`
- [ ] Modify `StoreChunksOperator` for shadow mode
- [ ] Modify `StoreVectorsOperator` for shadow mode
- [ ] Create `plugins/metrics/extraction_metrics.py`
- [ ] Create `scripts/compare_outputs.py`
- [ ] Add Celery metrics collection to existing task
- [ ] Test shadow mode execution
- [ ] Create comparison documentation

---

## Testing

```bash
# Enable shadow mode
export AIRFLOW_SHADOW_MODE=true

# Trigger shadow run
airflow dags trigger document_extraction_dag \
  --conf '{"document_id": "test-uuid", "file_path": "test.xlsx"}'

# Compare with Celery run
python scripts/compare_outputs.py \
  --airflow-run manual__2025-12-08 \
  --celery-task abc123
```

---

## Dependencies

- Story 1.6 (Main DAG)
- Story 1.5 (Storage operators)
- Redis for metrics storage

---

## Notes

- Shadow mode is safe for production - no data written
- Metrics stored in Redis with 7-day TTL
- Comparison tool helps validate migration correctness
- Can optionally mock embeddings to save API costs in shadow mode
