Skip to content

PR 14: Pipeline Integration (Final Wiring) #15

Description

@gilmanb1

Overview

Wire the new UnifiedCategorizationService into the main SortAIPipeline and AppState. This is the final PR that enables the new provider cascade in production.

Dependencies

Note: This should be the last PR merged to enable the feature.

Files to Modify

File Action Description
Sources/SortAI/Core/Pipeline/SortAIPipeline.swift Modify Use UnifiedCategorizationService
Sources/SortAI/Core/Types/FileSignature.swift Modify Add provider field to BrainResult
Sources/SortAI/App/AppState.swift Modify Pass AI config to pipeline

Implementation Details

1. Update BrainResult

// In FileSignature.swift

/// Result from LLM categorization
struct BrainResult: Sendable, Codable {
    let suggestedCategory: String
    let confidence: Double
    let rationale: String
    let extractedKeywords: [String]
    let provider: String?  // NEW: Track which provider generated this
    let timestamp: Date
    
    init(
        suggestedCategory: String,
        confidence: Double,
        rationale: String,
        extractedKeywords: [String],
        provider: String? = nil,
        timestamp: Date = Date()
    ) {
        self.suggestedCategory = suggestedCategory
        self.confidence = confidence
        self.rationale = rationale
        self.extractedKeywords = extractedKeywords
        self.provider = provider
        self.timestamp = timestamp
    }
    
    /// Convert from CategorizationResult
    init(from result: CategorizationResult) {
        self.suggestedCategory = result.categoryPath.path
        self.confidence = result.confidence
        self.rationale = result.rationale
        self.extractedKeywords = result.extractedKeywords
        self.provider = result.provider
        self.timestamp = Date()
    }
}

2. Update SortAIPipeline

// In SortAIPipeline.swift

actor SortAIPipeline {
    // MARK: - Dependencies
    
    private var categorizationService: UnifiedCategorizationService?
    private let aiConfig: AIProviderConfiguration
    
    // Existing dependencies...
    
    // MARK: - Initialization
    
    init(configuration: AppConfiguration) async {
        self.aiConfig = configuration.aiProvider
        
        // Initialize the unified categorization service
        self.categorizationService = await UnifiedCategorizationService()
        await categorizationService?.setPreference(aiConfig.preference)
        await categorizationService?.setEscalationThreshold(aiConfig.escalationThreshold)
        await categorizationService?.setAutoInstallOllama(aiConfig.autoInstallOllama)
        
        // ... existing initialization
    }
    
    // MARK: - Categorization
    
    /// Perform brain categorization using the unified service
    func performBrainCategorization(_ signature: FileSignature) async throws -> BrainResult {
        guard let service = categorizationService else {
            throw PipelineError.serviceNotInitialized
        }
        
        do {
            let result = try await service.categorize(signature: signature)
            
            // Log which provider was used
            NSLog("🧠 [Pipeline] Categorized '%@' via %@ (confidence: %.2f)",
                  signature.url.lastPathComponent,
                  result.provider,
                  result.confidence)
            
            return BrainResult(from: result)
            
        } catch {
            NSLog("❌ [Pipeline] Categorization failed: %@", error.localizedDescription)
            throw error
        }
    }
    
    /// Update AI configuration at runtime
    func updateAIConfiguration(_ config: AIProviderConfiguration) async {
        await categorizationService?.setPreference(config.preference)
        await categorizationService?.setEscalationThreshold(config.escalationThreshold)
        await categorizationService?.setAutoInstallOllama(config.autoInstallOllama)
    }
    
    // MARK: - Provider Status
    
    /// Get current active provider
    func activeProvider() async -> String {
        await categorizationService?.activeProvider ?? "unknown"
    }
    
    /// Check if escalation is in progress
    func isEscalating() async -> Bool {
        await categorizationService?.isEscalating ?? false
    }
    
    /// Get available providers
    func availableProviders() async -> [String] {
        await categorizationService?.getAvailableProviders() ?? []
    }
}

enum PipelineError: LocalizedError {
    case serviceNotInitialized
    case categorizationFailed(Error)
    
    var errorDescription: String? {
        switch self {
        case .serviceNotInitialized:
            return "Categorization service has not been initialized"
        case .categorizationFailed(let error):
            return "Categorization failed: \(error.localizedDescription)"
        }
    }
}

3. Update AppState

// In AppState.swift

@MainActor
class AppState: ObservableObject {
    // MARK: - Published State
    
    @Published var activeProvider: String = "initializing"
    @Published var isEscalating: Bool = false
    
    // MARK: - Pipeline
    
    private var pipeline: SortAIPipeline?
    private var providerObservationTask: Task<Void, Never>?
    
    // MARK: - Initialization
    
    func initialize() async {
        let config = AppConfiguration.shared
        
        // Initialize pipeline with AI config
        pipeline = await SortAIPipeline(configuration: config)
        
        // Observe provider changes
        observeProviderStatus()
        
        // Update initial state
        await refreshProviderStatus()
    }
    
    // MARK: - Configuration
    
    func updateAIConfiguration(_ config: AIProviderConfiguration) {
        Task {
            await pipeline?.updateAIConfiguration(config)
            await refreshProviderStatus()
        }
    }
    
    // MARK: - Provider Status
    
    private func observeProviderStatus() {
        providerObservationTask?.cancel()
        
        providerObservationTask = Task {
            while !Task.isCancelled {
                await refreshProviderStatus()
                try? await Task.sleep(nanoseconds: 5_000_000_000) // 5 seconds
            }
        }
    }
    
    private func refreshProviderStatus() async {
        guard let pipeline = pipeline else { return }
        
        let provider = await pipeline.activeProvider()
        let escalating = await pipeline.isEscalating()
        
        await MainActor.run {
            self.activeProvider = provider
            self.isEscalating = escalating
        }
    }
    
    // MARK: - Processing
    
    func processFile(_ url: URL) async throws {
        guard let pipeline = pipeline else {
            throw AppError.notInitialized
        }
        
        // ... existing file processing logic ...
        
        // Get categorization result (now includes provider info)
        let signature = try await pipeline.analyzeFile(url)
        let brainResult = try await pipeline.performBrainCategorization(signature)
        
        // Provider info is now in brainResult.provider
        NSLog("📊 Result from %@: %@ (%.0f%%)",
              brainResult.provider ?? "unknown",
              brainResult.suggestedCategory,
              brainResult.confidence * 100)
        
        // ... rest of processing ...
    }
}

4. Feature Flag (Optional Safety)

// In AppConfiguration.swift

extension AppConfiguration {
    /// Feature flag to enable new provider cascade
    var useNewProviderCascade: Bool {
        get {
            UserDefaults.standard.bool(forKey: "useNewProviderCascade")
        }
        set {
            UserDefaults.standard.set(newValue, forKey: "useNewProviderCascade")
        }
    }
}

// In SortAIPipeline.swift
func performBrainCategorization(_ signature: FileSignature) async throws -> BrainResult {
    if AppConfiguration.shared.useNewProviderCascade {
        // New cascade
        return try await performBrainCategorizationV2(signature)
    } else {
        // Legacy Ollama-only path
        return try await performLegacyBrainCategorization(signature)
    }
}

Migration Notes

Breaking Changes

None. The BrainResult.provider field is optional and defaults to nil for backward compatibility.

Gradual Rollout

  1. Merge PR with feature flag disabled by default
  2. Enable for internal testing
  3. Roll out to users via settings toggle
  4. Remove feature flag once stable

Integration Flow

User drops file
       │
       ▼
  AppState.processFile()
       │
       ▼
  SortAIPipeline.analyzeFile()
       │
       ▼
  SortAIPipeline.performBrainCategorization()
       │
       ▼
  UnifiedCategorizationService.categorize()
       │
       ├─► Apple Intelligence (try first)
       │         │
       │         ├─ Success (confidence > threshold) → Return result
       │         │
       │         └─ Low confidence → Escalate
       │                    │
       ├──────────────────◄─┘
       │
       ├─► Ollama (if available)
       │         │
       │         ├─ Success → Return result
       │         │
       │         └─ Failure → Continue
       │
       ├─► Cloud (if configured)
       │         │
       │         └─ ...
       │
       └─► Local ML (final fallback)
                 │
                 └─ Always returns result
       │
       ▼
  BrainResult (includes provider field)
       │
       ▼
  UI displays result with ProviderBadge

Acceptance Criteria

  • Pipeline uses UnifiedCategorizationService
  • BrainResult includes provider field
  • Provider info flows to UI
  • Configuration updates apply at runtime
  • Feature flag allows gradual rollout
  • Existing tests continue to pass
  • New integration tests verify cascade
  • Logging shows which provider was used

Testing

func testPipelineUsesUnifiedService() async throws {
    let config = AppConfiguration()
    let pipeline = await SortAIPipeline(configuration: config)
    
    let signature = FileSignature.mock(filename: "test.pdf")
    let result = try await pipeline.performBrainCategorization(signature)
    
    // Result should have provider info
    XCTAssertNotNil(result.provider)
}

func testProviderInfoInBrainResult() async throws {
    let config = AppConfiguration()
    config.aiProvider.preference = .appleIntelligenceOnly
    
    let pipeline = await SortAIPipeline(configuration: config)
    let signature = FileSignature.mock(filename: "test.txt")
    
    let result = try await pipeline.performBrainCategorization(signature)
    
    // Should use Apple Intelligence or Local ML
    XCTAssertTrue(
        result.provider == "apple-intelligence" || result.provider == "local-ml"
    )
}

func testConfigurationUpdateAtRuntime() async throws {
    let config = AppConfiguration()
    let pipeline = await SortAIPipeline(configuration: config)
    
    // Initial preference
    XCTAssertEqual(config.aiProvider.preference, .automatic)
    
    // Update preference
    var newConfig = config.aiProvider
    newConfig.preference = .preferOllama
    await pipeline.updateAIConfiguration(newConfig)
    
    // Pipeline should now prefer Ollama
    // (Would need to verify through categorization)
}

func testFeatureFlagDisablesNewCascade() async throws {
    AppConfiguration.shared.useNewProviderCascade = false
    
    let pipeline = await SortAIPipeline(configuration: AppConfiguration.shared)
    let signature = FileSignature.mock(filename: "test.txt")
    
    let result = try await pipeline.performBrainCategorization(signature)
    
    // Should use legacy path (Ollama only)
    // Provider might be nil or "ollama"
}

Estimated Size

~100-150 lines of code

Risk Assessment

Medium - Integration PR affects core pipeline. Mitigation:

  • Feature flag for gradual rollout
  • Comprehensive testing
  • Legacy path preserved as fallback
  • Monitoring/logging for provider usage

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or requestinfrastructureCore infrastructure and protocolsphase-3Phase 3 - Integration

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions