Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
133 changes: 126 additions & 7 deletions .github/copilot-instructions.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,15 +10,75 @@ The codebase follows a clear separation of concerns:

```
paper/
├── core.py # PaperMatrix - disk-backed matrix using memory-mapped files
├── plan.py # Lazy evaluation plan tree (EagerNode, AddNode, MultiplyNode, etc.)
├── optimizer.py # Plan inspection, I/O trace generation, and fusion rules
├── backend.py # High-performance execution kernels for matrix operations
├── buffer.py # BufferManager with LRU and Belady's optimal eviction
├── config.py # Centralized configuration (TILE_SIZE, cache sizes)
└── numpy_api.py # NumPy-compatible API layer (pnp.array, pnp.zeros, etc.)
├── core.py # PaperMatrix - disk-backed matrix using memory-mapped files
├── plan.py # Lazy evaluation plan tree (EagerNode, AddNode, MultiplyNode, etc.)
├── optimizer.py # Three-stage optimizer: analyze(), rewrite(), execute_plan()
├── backend.py # High-performance execution kernels for matrix operations
├── buffer.py # BufferManager with LRU and Belady's optimal eviction
├── config.py # Centralized configuration (TILE_SIZE, cache sizes)
├── numpy_api.py # NumPy-compatible API layer (pnp.array, pnp.zeros, etc.)
└── observability.py # Logging, profiling, tracing utilities
```

## Core Principles

These principles guide all architectural decisions in the Paper framework:

1. **Make optimizer purely analytical**: Never materialize data during pattern matching.
2. **Separate concerns**: Analysis (trace + match) → IR rewrite (fuse) → Execution (apply kernels).
3. **Define small, stable contracts**: Each layer (Plan/Node metadata, BufferManager, Backend) has clear interfaces.
4. **Make fusion selection deterministic**: Testable and fast (no I/O during optimization).
5. **Treat heavy tests and benchmarks as gated**: Cached jobs in CI for efficiency.

## Priority Roadmap

### **Priority 1 (Critical)** - Stop side-effects in optimizer ✅ DONE
- Remove any calls to `node.execute()` inside optimizer/analysis
- Replace with metadata-only inspection APIs
- **Why**: Prevents data materialization during plan analysis, critical for performance
- **Status**: ✅ Pattern detection in `_detect_fusion_pattern()` uses only `isinstance()` checks and metadata access. No `execute()` calls during analysis phase

### **Priority 2 (Critical)** - Create two-stage optimizer ✅ DONE
- `analyze(plan)` → Trace + MatchResults
- `rewrite(plan, match)` → FusedPlan (new node types)
- `execute(fused_plan, backend, buffer_mgr)`
- **Why**: Clean separation enables testability and independent optimization
- **Status**: ✅ Three-stage pipeline implemented in `optimizer.py`: `analyze()` generates I/O trace and detects patterns without execution, `rewrite()` prepares optimized plan, `execute_plan()` runs with fusion. Legacy `execute()` maintained for backward compatibility

### **Priority 3 (High)** - Introduce immutable, hashable Plan representation ❌ NOT DONE
- Deterministic hashing for plan diffs, caching, and baseline comparisons
- **Why**: Enables plan comparison, caching, and regression detection
- **Status**: No `__hash__()` or `__eq__()` methods found in Plan or Node classes

### **Priority 4 (High)** - Add a small, fast unit test surface ✅ DONE
- Optimizer tests: trace generation, pattern detection, rewrite correctness (use mocked backend)
- **Why**: Fast feedback loop for optimizer development
- **Status**: `test_plan_optimizer.py` has tests for plan construction and I/O trace generation

### **Priority 5 (High)** - Add integration smoke tests ✅ DONE
- Run fused kernels on tiny matrices (CI quick job)
- **Why**: Catch integration bugs without heavy computation
- **Status**: `test_fusion_operations.py` tests fused kernels with small (128×128) matrices. CI runs all tests in `ci.yml`

### **Priority 6 (Medium)** - Add benchmarks & regression checks ⚠️ PARTIAL
- Baselines stored as CI artifacts
- Regress only if delta > threshold
- **Why**: Prevent performance regressions systematically
- **Status**: `benchmarks/benchmark_dask.py` exists with benchmarking framework. CI uploads test artifacts but no baseline comparison or regression checks

### **Priority 7 (Medium)** - Plugin backend API ✅ DONE
- Simple interface: `(inputs, params, output_path, buffer_mgr)`
- Swapping implementations is trivial
- **Why**: Enables experimentation with different execution strategies
- **Status**: All backend kernels follow consistent signature `(A, B, ..., output_path, buffer_manager)`. Interface is clean and swappable

### **Priority 8 (Low)** - Observability ✅ DONE
- Trace-level logging
- Cost model hooks
- Per-plan flame profiles
- **Why**: Debugging and performance analysis
- **Status**: ✅ Comprehensive observability in `observability.py`: structured logging with `configure_logging()`, `ExecutionProfiler` for timing/flame graphs, `TraceLogger` for execution traces, `CostEstimate` dataclass for I/O cost modeling in optimizer

## Key Design Patterns

### 1. Lazy Evaluation
Expand All @@ -41,6 +101,65 @@ paper/

## Coding Conventions

### Three-Stage Optimizer Usage

The new three-stage optimizer provides clean separation between analysis, rewriting, and execution:

```python
from paper.optimizer import analyze, rewrite, execute_plan, estimate_cost
from paper.buffer import BufferManager

# Stage 1: Analyze (no execution - metadata only)
io_trace, match_result = analyze(plan)

# Check what pattern was detected
if match_result.is_fusable:
print(f"Detected pattern: {match_result.pattern.value}")
print(f"Parameters: {match_result.parameters}")

# Estimate execution cost
cost = estimate_cost(plan, match_result)
print(f"Predicted I/O ops: {cost.io_operations}")
print(f"Total cost: {cost.total_cost}")

# Stage 2: Rewrite (prepare optimized plan)
rewritten_plan = rewrite(plan, match_result)

# Stage 3: Execute
buffer_manager = BufferManager(max_cache_size_tiles=64, io_trace=io_trace)
result = execute_plan(rewritten_plan, match_result, output_path, buffer_manager)
```

### Observability Features

Enable structured logging, profiling, and tracing:

```python
from paper.observability import configure_logging, get_profiler, TraceLogger

# Configure logging
logger = configure_logging(level="DEBUG", log_file="paper.log")

# Use global profiler
profiler = get_profiler()

with profiler.profile("my_operation"):
# ... do work ...
pass

# Print profiling summary
profiler.print_summary()
profiler.save_json("profile.json")
profiler.save_flame_graph("flame.json")

# Execution tracing
trace = TraceLogger()
trace.begin("operation")
trace.log("step 1 complete")
trace.end()
trace.print()
```

### Python Style
- Use type hints for function parameters: `def add(A: PaperMatrix, B: PaperMatrix, ...)`
- Use `np.float32` as the default dtype
Expand Down
238 changes: 238 additions & 0 deletions examples/two_stage_optimizer_demo.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,238 @@
"""
Example demonstrating the new two-stage optimizer and observability features.

This example shows:
1. Three-stage optimizer pipeline (analyze → rewrite → execute)
2. Structured logging
3. Performance profiling
4. Cost estimation
5. Execution tracing
"""

import numpy as np
import os
import sys
import tempfile

# Add parent directory to path
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), '..')))

from paper.core import PaperMatrix
from paper.plan import Plan, EagerNode
from paper.optimizer import analyze, rewrite, execute_plan, estimate_cost
from paper.buffer import BufferManager
from paper.observability import (
configure_logging, get_profiler, TraceLogger
)


def main():
"""Run the example with observability features."""

# ========================================================================
# Setup: Configure logging
# ========================================================================
print("\n" + "="*80)
print("EXAMPLE: Two-Stage Optimizer with Observability")
print("="*80 + "\n")

# Configure structured logging
logger = configure_logging(level="INFO", log_file="paper_example.log")
logger.info("Starting two-stage optimizer example")

# Get global profiler
profiler = get_profiler()

# Create trace logger
trace = TraceLogger()

# ========================================================================
# Step 1: Create test data
# ========================================================================
print("Step 1: Creating test matrices...")
trace.begin("Data Creation")

with profiler.profile("create_test_data"):
test_dir = tempfile.mkdtemp()
shape = (2048, 2048)
dtype = np.float32

# Create matrices A and B
A_path = os.path.join(test_dir, "A.bin")
B_path = os.path.join(test_dir, "B.bin")

np.random.seed(42)
data_A = np.random.rand(*shape).astype(dtype)
data_B = np.random.rand(*shape).astype(dtype)

data_A.tofile(A_path)
data_B.tofile(B_path)

A = PaperMatrix(A_path, shape, dtype=dtype, mode='r')
B = PaperMatrix(B_path, shape, dtype=dtype, mode='r')

trace.log(f"Created matrices: {shape}")

trace.end()

# ========================================================================
# Step 2: Build computation plan
# ========================================================================
print("\nStep 2: Building computation plan...")
trace.begin("Plan Construction")

with profiler.profile("build_plan"):
plan_A = Plan(EagerNode(A))
plan_B = Plan(EagerNode(B))

# Build a complex plan: (A + B) * 2.5
plan = (plan_A + plan_B) * 2.5

trace.log(f"Plan: (A + B) * 2.5")

trace.end()

# ========================================================================
# Step 3: STAGE 1 - Analyze (no execution)
# ========================================================================
print("\nStep 3: STAGE 1 - Analyzing plan (no execution)...")
trace.begin("Stage 1: Analyze")

with profiler.profile("analyze_plan"):
io_trace, match_result = analyze(plan)

print(f"\n ✓ Pattern detected: {match_result.pattern.value}")
print(f" ✓ Fusion available: {match_result.is_fusable}")
print(f" ✓ I/O trace length: {len(io_trace)} tile accesses")
print(f" ✓ Input shapes: {match_result.input_shapes}")
print(f" ✓ Parameters: {match_result.parameters}")

trace.log(f"Pattern: {match_result.pattern.value}")
trace.log(f"I/O trace: {len(io_trace)} accesses")

trace.end()

# ========================================================================
# Step 4: Cost Estimation
# ========================================================================
print("\nStep 4: Estimating execution cost...")
trace.begin("Cost Estimation")

with profiler.profile("estimate_cost"):
cost = estimate_cost(plan, match_result)

print(f"\n ✓ I/O operations: {cost.io_operations}")
print(f" ✓ Compute operations: {cost.compute_operations}")
print(f" ✓ Estimated I/O bytes: {cost.estimated_io_bytes:,}")
print(f" ✓ Cache benefit: {cost.cache_benefit:.1%}")
print(f" ✓ Total cost: {cost.total_cost:.0f}")

trace.log(f"Total cost: {cost.total_cost:.0f}")

trace.end()

# ========================================================================
# Step 5: STAGE 2 - Rewrite plan
# ========================================================================
print("\nStep 5: STAGE 2 - Rewriting plan...")
trace.begin("Stage 2: Rewrite")

with profiler.profile("rewrite_plan"):
rewritten_plan = rewrite(plan, match_result)
trace.log("Plan rewritten for fusion")

trace.end()

# ========================================================================
# Step 6: STAGE 3 - Execute
# ========================================================================
print("\nStep 6: STAGE 3 - Executing plan...")
trace.begin("Stage 3: Execute")

output_path = os.path.join(test_dir, "result.bin")

with profiler.profile("execute_plan"):
buffer_manager = BufferManager(max_cache_size_tiles=32, io_trace=io_trace)
result = execute_plan(rewritten_plan, match_result, output_path, buffer_manager)

trace.log("Execution complete")

trace.end()

# ========================================================================
# Step 7: Verify results
# ========================================================================
print("\nStep 7: Verifying results...")
trace.begin("Verification")

with profiler.profile("verify_result"):
result_data = np.fromfile(output_path, dtype=dtype).reshape(shape)
expected = (data_A + data_B) * 2.5

max_diff = np.max(np.abs(result_data - expected))
print(f"\n ✓ Max difference from expected: {max_diff:.2e}")

trace.log(f"Verification: max_diff={max_diff:.2e}")

trace.end()

# ========================================================================
# Step 8: Print profiling results
# ========================================================================
print("\n" + "="*80)
print("PROFILING RESULTS")
print("="*80)

profiler.print_summary()

# Save profiling data
profiler.save_json("profile_results.json")
profiler.save_flame_graph("flame_graph.json")
print("✓ Profiling data saved to profile_results.json and flame_graph.json")

# ========================================================================
# Step 9: Print execution trace
# ========================================================================
trace.print()

# Save trace
trace.save("execution_trace.json")
print("✓ Execution trace saved to execution_trace.json")

# ========================================================================
# Step 10: Show cache statistics
# ========================================================================
print("\n" + "="*80)
print("CACHE STATISTICS")
print("="*80)

cache_log = buffer_manager.get_log()

hits = sum(1 for event in cache_log if event[1] == 'HIT')
misses = sum(1 for event in cache_log if event[1] == 'MISS')
evictions = sum(1 for event in cache_log if event[1] == 'EVICT')
total = hits + misses

print(f"Cache hits: {hits:6d} ({100*hits/total if total > 0 else 0:.1f}%)")
print(f"Cache misses: {misses:6d} ({100*misses/total if total > 0 else 0:.1f}%)")
print(f"Evictions: {evictions:6d}")
print(f"Total accesses: {total:6d}")
print("="*80 + "\n")

# ========================================================================
# Cleanup
# ========================================================================
A.close()
B.close()
result.close()

import shutil
shutil.rmtree(test_dir)

print("✓ Example complete!")
print(f"✓ Log file: paper_example.log")
logger.info("Example completed successfully")


if __name__ == "__main__":
main()
Loading
Loading