From e26830439dfdb97c1d685bf5cda05a5d2df9f21f Mon Sep 17 00:00:00 2001 From: Arturs Artamonovs Date: Mon, 17 Nov 2025 00:44:09 +0000 Subject: refactor experiment --- LearnMapKit/DataSources/ADSBDataSource.swift | 92 ++++++++++ LearnMapKit/DataSources/FileADSBDataSource.swift | 185 +++++++++++++++++++ .../DataSources/NetworkADSBDataSource.swift | 201 +++++++++++++++++++++ 3 files changed, 478 insertions(+) create mode 100644 LearnMapKit/DataSources/ADSBDataSource.swift create mode 100644 LearnMapKit/DataSources/FileADSBDataSource.swift create mode 100644 LearnMapKit/DataSources/NetworkADSBDataSource.swift (limited to 'LearnMapKit/DataSources') 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 { get } + var connectionState: AnyPublisher { 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 { + 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.. { + 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..