diff options
| author | Arturs Artamonovs <arturs.artamonovs@protonmail.com> | 2025-11-17 00:44:09 +0000 |
|---|---|---|
| committer | Arturs Artamonovs <arturs.artamonovs@protonmail.com> | 2025-11-17 00:44:09 +0000 |
| commit | e26830439dfdb97c1d685bf5cda05a5d2df9f21f (patch) | |
| tree | 85a16662d47787d0975c05a22727ad4301d0fb5b /LearnMapKit/DataSources/NetworkADSBDataSource.swift | |
| parent | 2b599499dc26f0f93537d1bea4733eae68333a1c (diff) | |
| download | ADSBDecoder-e26830439dfdb97c1d685bf5cda05a5d2df9f21f.tar.gz ADSBDecoder-e26830439dfdb97c1d685bf5cda05a5d2df9f21f.zip | |
refactor experiment
Diffstat (limited to 'LearnMapKit/DataSources/NetworkADSBDataSource.swift')
| -rw-r--r-- | LearnMapKit/DataSources/NetworkADSBDataSource.swift | 201 |
1 files changed, 201 insertions, 0 deletions
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 |
