// // NetworkADSBDataSource.swift // LearnMapKit // // Network-based ADSB data source implementation // import Foundation import Combine @MainActor class NetworkADSBDataSource: ADSBDataSource { // MARK: - ADSBDataSource Protocol var dataStream: AnyPublisher { dataSubject.eraseToAnyPublisher() } var connectionState: AnyPublisher { stateSubject.eraseToAnyPublisher() } private(set) var configuration: DataSourceConfiguration // MARK: - Private Properties private let dataSubject = PassthroughSubject() private let stateSubject = CurrentValueSubject(.idle) private var netRunner: ADSBNetRunner? private var processingTask: Task? private var timer: Timer? private var reconnectAttempts = 0 // MARK: - Dependencies private let decoder: ADSBMessageDecoder private let tracker: AircraftTracker init(configuration: NetworkDataSourceConfiguration, decoder: ADSBMessageDecoder = DefaultADSBMessageDecoder(), tracker: AircraftTracker = DefaultAircraftTracker()) { self.configuration = configuration self.decoder = decoder self.tracker = tracker } // MARK: - ADSBDataSource Implementation func configure(with config: DataSourceConfiguration) async { guard let networkConfig = config as? NetworkDataSourceConfiguration else { stateSubject.send(.error("Invalid configuration type for NetworkADSBDataSource")) return } stateSubject.send(.configuring) self.configuration = networkConfig reconnectAttempts = 0 stateSubject.send(.idle) } func start() async throws { guard let networkConfig = configuration as? NetworkDataSourceConfiguration else { throw ADSBDataSourceError.invalidConfiguration } stateSubject.send(.connecting) do { // Initialize network runner netRunner = ADSBNetRunner(address: networkConfig.hostname, port: networkConfig.port) // Start network processing try await startNetworkProcessing() stateSubject.send(.connected) reconnectAttempts = 0 } catch { stateSubject.send(.error(error.localizedDescription)) await handleConnectionError(error: error) throw error } } func stop() async { timer?.invalidate() timer = nil processingTask?.cancel() processingTask = nil await netRunner?.stop() netRunner = nil stateSubject.send(.disconnected) reconnectAttempts = 0 } func reconnect() async throws { guard let networkConfig = configuration as? NetworkDataSourceConfiguration else { throw ADSBDataSourceError.invalidConfiguration } // Check if we should attempt reconnection guard reconnectAttempts < networkConfig.reconnectAttempts else { stateSubject.send(.error("Maximum reconnection attempts reached")) return } reconnectAttempts += 1 await stop() // Wait before reconnecting try await Task.sleep(nanoseconds: UInt64(2_000_000_000 * reconnectAttempts)) // Exponential backoff try await start() } // MARK: - Private Methods private func startNetworkProcessing() async throws { guard let runner = netRunner else { throw ADSBDataSourceError.initializationFailed } // Start network connection processingTask = Task { do { try await runner.start() await startDataConsumption(runner: runner) } catch { await MainActor.run { dataSubject.send(.error(error)) stateSubject.send(.error(error.localizedDescription)) } } } } private func startDataConsumption(runner: ADSBNetRunner) async { await MainActor.run { timer = Timer.scheduledTimer(withTimeInterval: 1.0, repeats: true) { [weak self] _ in self?.processDataBatch(from: runner) } } } private func processDataBatch(from runner: ADSBNetRunner) { guard runner.getCount() > 0 else { return } Task { for _ in 0..