summaryrefslogtreecommitdiff
path: root/LearnMapKit/DataSources/NetworkADSBDataSource.swift
diff options
context:
space:
mode:
authorArturs Artamonovs <arturs.artamonovs@protonmail.com>2025-11-17 00:44:09 +0000
committerArturs Artamonovs <arturs.artamonovs@protonmail.com>2025-11-17 00:44:09 +0000
commite26830439dfdb97c1d685bf5cda05a5d2df9f21f (patch)
tree85a16662d47787d0975c05a22727ad4301d0fb5b /LearnMapKit/DataSources/NetworkADSBDataSource.swift
parent2b599499dc26f0f93537d1bea4733eae68333a1c (diff)
downloadADSBDecoder-e26830439dfdb97c1d685bf5cda05a5d2df9f21f.tar.gz
ADSBDecoder-e26830439dfdb97c1d685bf5cda05a5d2df9f21f.zip
refactor experiment
Diffstat (limited to 'LearnMapKit/DataSources/NetworkADSBDataSource.swift')
-rw-r--r--LearnMapKit/DataSources/NetworkADSBDataSource.swift201
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