1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
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
}
}
|