// // FileADSBDataSource.swift // LearnMapKit // // File-based ADSB data source implementation // import Foundation import Combine @MainActor class FileADSBDataSource: 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 fileRunner: ADSBFileRunner? private var processingTask: Task? private var timer: Timer? // MARK: - Dependencies private let decoder: ADSBMessageDecoder private let tracker: AircraftTracker init(configuration: FileDataSourceConfiguration, 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 fileConfig = config as? FileDataSourceConfiguration else { stateSubject.send(.error("Invalid configuration type for FileADSBDataSource")) return } stateSubject.send(.configuring) self.configuration = fileConfig stateSubject.send(.idle) } func start() async throws { guard let fileConfig = configuration as? FileDataSourceConfiguration else { throw ADSBDataSourceError.invalidConfiguration } stateSubject.send(.connecting) do { // Initialize file runner fileRunner = ADSBFileRunner(filename: fileConfig.filePath) // Start file processing try await startFileProcessing(with: fileConfig) stateSubject.send(.connected) } catch { stateSubject.send(.error(error.localizedDescription)) throw error } } func stop() async { timer?.invalidate() timer = nil processingTask?.cancel() processingTask = nil fileRunner = nil stateSubject.send(.disconnected) } func reconnect() async throws { await stop() try await start() } // MARK: - Private Methods private func startFileProcessing(with config: FileDataSourceConfiguration) async throws { guard let runner = fileRunner else { throw ADSBDataSourceError.initializationFailed } // Start background file processing processingTask = Task { do { runner.openFile() runner.readFile() // Start decoding in background await runner.decodeFromFile() // Start timer for data consumption await startDataConsumption(runner: runner, rate: config.processRate) } catch { await MainActor.run { dataSubject.send(.error(error)) stateSubject.send(.error(error.localizedDescription)) } } } } private func startDataConsumption(runner: ADSBFileRunner, rate: Int) async { await MainActor.run { timer = Timer.scheduledTimer(withTimeInterval: 1.0, repeats: true) { [weak self] _ in self?.processDataBatch(from: runner, batchSize: rate) } } } private func processDataBatch(from runner: ADSBFileRunner, batchSize: Int) { guard runner.jobDone() && runner.getCount() > 0 else { return } let batchSize = min(batchSize, runner.getCount()) let data = runner.getPlainData(batchSize) Task { for _ in 0..