diff options
Diffstat (limited to 'LearnMapKit/DataSources')
| -rw-r--r-- | LearnMapKit/DataSources/ADSBDataSource.swift | 92 | ||||
| -rw-r--r-- | LearnMapKit/DataSources/FileADSBDataSource.swift | 185 | ||||
| -rw-r--r-- | LearnMapKit/DataSources/NetworkADSBDataSource.swift | 201 |
3 files changed, 478 insertions, 0 deletions
diff --git a/LearnMapKit/DataSources/ADSBDataSource.swift b/LearnMapKit/DataSources/ADSBDataSource.swift new file mode 100644 index 0000000..b30e344 --- /dev/null +++ b/LearnMapKit/DataSources/ADSBDataSource.swift @@ -0,0 +1,92 @@ +// +// ADSBDataSource.swift +// LearnMapKit +// +// Created for MVVM refactoring +// + +import Foundation +import Combine + +// MARK: - Data Source Protocol +protocol ADSBDataSource: AnyObject { + var dataStream: AnyPublisher<ADSBDataEvent, Never> { get } + var connectionState: AnyPublisher<DataSourceState, Never> { get } + var configuration: DataSourceConfiguration { get } + + func configure(with config: DataSourceConfiguration) async + func start() async throws + func stop() async + func reconnect() async throws +} + +// MARK: - Data Events +enum ADSBDataEvent { + case aircraftIdentification(address: Int, icaoName: String) + case aircraftPosition(address: Int, latitude: Double, longitude: Double) + case aircraftAltitude(address: Int, altitude: Int) + case rawMessage(String) + case error(Error) +} + +// MARK: - Data Source State +enum DataSourceState: Equatable { + case idle + case configuring + case connecting + case connected + case disconnected + case error(String) + + static func == (lhs: DataSourceState, rhs: DataSourceState) -> Bool { + switch (lhs, rhs) { + case (.idle, .idle), (.configuring, .configuring), + (.connecting, .connecting), (.connected, .connected), + (.disconnected, .disconnected): + return true + case let (.error(lhsError), .error(rhsError)): + return lhsError == rhsError + default: + return false + } + } +} + +// MARK: - Configuration +protocol DataSourceConfiguration { + var sourceType: DataSourceType { get } +} + +enum DataSourceType { + case file(path: String) + case network(host: String, port: Int) + case test(mockData: [String]) +} + +struct FileDataSourceConfiguration: DataSourceConfiguration { + let sourceType: DataSourceType + let filePath: String + let processRate: Int // messages per second + + init(filePath: String, processRate: Int = 120) { + self.filePath = filePath + self.processRate = processRate + self.sourceType = .file(path: filePath) + } +} + +struct NetworkDataSourceConfiguration: DataSourceConfiguration { + let sourceType: DataSourceType + let hostname: String + let port: Int + let reconnectAttempts: Int + let timeoutInterval: TimeInterval + + init(hostname: String, port: Int, reconnectAttempts: Int = 3, timeoutInterval: TimeInterval = 10.0) { + self.hostname = hostname + self.port = port + self.reconnectAttempts = reconnectAttempts + self.timeoutInterval = timeoutInterval + self.sourceType = .network(host: hostname, port: port) + } +}
\ No newline at end of file diff --git a/LearnMapKit/DataSources/FileADSBDataSource.swift b/LearnMapKit/DataSources/FileADSBDataSource.swift new file mode 100644 index 0000000..448721b --- /dev/null +++ b/LearnMapKit/DataSources/FileADSBDataSource.swift @@ -0,0 +1,185 @@ +// +// FileADSBDataSource.swift +// LearnMapKit +// +// File-based ADSB data source implementation +// + +import Foundation +import Combine + +@MainActor +class FileADSBDataSource: ADSBDataSource { + + // MARK: - ADSBDataSource Protocol + var dataStream: AnyPublisher<ADSBDataEvent, Never> { + dataSubject.eraseToAnyPublisher() + } + + var connectionState: AnyPublisher<DataSourceState, Never> { + stateSubject.eraseToAnyPublisher() + } + + private(set) var configuration: DataSourceConfiguration + + // MARK: - Private Properties + private let dataSubject = PassthroughSubject<ADSBDataEvent, Never>() + private let stateSubject = CurrentValueSubject<DataSourceState, Never>(.idle) + + private var fileRunner: ADSBFileRunner? + private var processingTask: Task<Void, Error>? + 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..<data.getCount() { + let nextTag = data.getNextTag() + + switch nextTag { + case .ADSB_ICAO: + let icao = data.getIcaoName() + await tracker.updateIdentification(address: icao.address, icaoName: icao.ICAOname) + dataSubject.send(.aircraftIdentification(address: icao.address, icaoName: icao.ICAOname)) + + case .ADSB_LOCATION: + let location = data.getLocation() + await tracker.updatePosition(address: location.address, + latitude: location.lat, + longitude: location.long) + dataSubject.send(.aircraftPosition(address: location.address, + latitude: location.lat, + longitude: location.long)) + + case .ADSB_ALTITUDE: + let altitude = data.getAltitude() + await tracker.updateAltitude(address: altitude.address, altitude: altitude.altitude) + dataSubject.send(.aircraftAltitude(address: altitude.address, altitude: altitude.altitude)) + + case .EMPTY: + break + } + } + } + } +} + +// MARK: - Errors +enum ADSBDataSourceError: Error, LocalizedError { + case invalidConfiguration + case fileNotFound(String) + case initializationFailed + case processingError(Error) + + var errorDescription: String? { + switch self { + case .invalidConfiguration: + return "Invalid data source configuration" + case .fileNotFound(let path): + return "File not found at path: \(path)" + case .initializationFailed: + return "Failed to initialize data source" + case .processingError(let error): + return "Processing error: \(error.localizedDescription)" + } + } +}
\ No newline at end of file diff --git a/LearnMapKit/DataSources/NetworkADSBDataSource.swift b/LearnMapKit/DataSources/NetworkADSBDataSource.swift new file mode 100644 index 0000000..a2c3890 --- /dev/null +++ b/LearnMapKit/DataSources/NetworkADSBDataSource.swift @@ -0,0 +1,201 @@ +// +// NetworkADSBDataSource.swift +// LearnMapKit +// +// Network-based ADSB data source implementation +// + +import Foundation +import Combine + +@MainActor +class NetworkADSBDataSource: ADSBDataSource { + + // MARK: - ADSBDataSource Protocol + var dataStream: AnyPublisher<ADSBDataEvent, Never> { + dataSubject.eraseToAnyPublisher() + } + + var connectionState: AnyPublisher<DataSourceState, Never> { + stateSubject.eraseToAnyPublisher() + } + + private(set) var configuration: DataSourceConfiguration + + // MARK: - Private Properties + private let dataSubject = PassthroughSubject<ADSBDataEvent, Never>() + private let stateSubject = CurrentValueSubject<DataSourceState, Never>(.idle) + + private var netRunner: ADSBNetRunner? + private var processingTask: Task<Void, Error>? + 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..<runner.adsb_tag_stream.getCount() { + let nextTag = runner.adsb_tag_stream.getNextTag() + + switch nextTag { + case .ADSB_ICAO: + let icao = runner.adsb_tag_stream.getIcaoName() + await tracker.updateIdentification(address: icao.address, icaoName: icao.ICAOname) + dataSubject.send(.aircraftIdentification(address: icao.address, icaoName: icao.ICAOname)) + + case .ADSB_LOCATION: + let location = runner.adsb_tag_stream.getLocation() + await tracker.updatePosition(address: location.address, + latitude: location.lat, + longitude: location.long) + dataSubject.send(.aircraftPosition(address: location.address, + latitude: location.lat, + longitude: location.long)) + + case .ADSB_ALTITUDE: + let altitude = runner.adsb_tag_stream.getAltitude() + await tracker.updateAltitude(address: altitude.address, altitude: altitude.altitude) + dataSubject.send(.aircraftAltitude(address: altitude.address, altitude: altitude.altitude)) + + case .EMPTY: + break + } + } + } + } + + private func handleConnectionError(error: Error) async { + guard let networkConfig = configuration as? NetworkDataSourceConfiguration else { return } + + // Attempt automatic reconnection for network errors + if reconnectAttempts < networkConfig.reconnectAttempts { + Task { + do { + try await reconnect() + } catch { + // Final failure after all attempts + stateSubject.send(.error("Connection failed after \(networkConfig.reconnectAttempts) attempts")) + } + } + } + } +} + +// MARK: - Extensions +extension ADSBNetRunner { + func stop() async { + // Add proper cleanup method to ADSBNetRunner + timer?.invalidate() + timer = nil + } +}
\ No newline at end of file |
