import ComposableArchitecture
import Foundation

/// Manages the real-time connection to Orpheus.
///
/// This component handles subscribing to real-time events and emitting them as a stream.
/// It does not manage session state - that's handled by OrpheusSessionStore.
public actor OrpheusConnectionManager {
    @Dependency(OrpheusRealtimeClient.self) private var realtimeClient
    private var subscriptionTask: Task<Void, Never>?
    private var eventContinuations: [UUID: AsyncStream<OrpheusRealtimeEvent>.Continuation] = [:]
    
    /// Real-time events stream. Emits events from the connected session.
    public func realtimeEvents() -> AsyncStream<OrpheusRealtimeEvent> {
        AsyncStream { continuation in
            let id = UUID()
            eventContinuations[id] = continuation
            
            continuation.onTermination = { [weak self] _ in
                Task { [weak self] in
                    await self?.removeContinuation(id: id)
                }
            }
        }
    }
    
    public init() {}
    
    /// Connects to real-time events for the given session.
    public func connect(sessionId: String) {
        subscriptionTask?.cancel()
        
        subscriptionTask = Task { [weak self] in
            guard let self = self else { return }
            
            do {
                let eventStream = await realtimeClient.subscribe(sessionId: sessionId)
                
                for try await event in eventStream {
                    // Emit to all continuations
                    let continuations = await self.eventContinuations
                    for continuation in continuations.values {
                        continuation.yield(event)
                    }
                    
                    if Task.isCancelled {
                        break
                    }
                }
                
                print("📡 Realtime connection ended for session: \(sessionId)")
            } catch is CancellationError {
                print("📡 Realtime connection cancelled for session: \(sessionId)")
            } catch {
                print("❌ Error in realtime connection for session: \(sessionId), error: \(error)")
                let continuations = await self.eventContinuations
                for continuation in continuations.values {
                    continuation.yield(.error(error))
                }
            }
        }
        
        print("📡 Started realtime connection for session: \(sessionId)")
    }
    
    /// Disconnects from real-time events.
    public func disconnect() {
        subscriptionTask?.cancel()
        subscriptionTask = nil
        
        // Finish all continuations
        for continuation in eventContinuations.values {
            continuation.finish()
        }
        eventContinuations.removeAll()
        
        print("📡 Stopped realtime connection")
    }
    
    private func removeContinuation(id: UUID) {
        eventContinuations.removeValue(forKey: id)
    }
}
