summaryrefslogtreecommitdiff
path: root/LearnMapKit/DataSources
diff options
context:
space:
mode:
Diffstat (limited to 'LearnMapKit/DataSources')
-rw-r--r--LearnMapKit/DataSources/ADSBDataSource.swift92
-rw-r--r--LearnMapKit/DataSources/FileADSBDataSource.swift185
-rw-r--r--LearnMapKit/DataSources/NetworkADSBDataSource.swift201
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