mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-24 22:45:19 +00:00
Move GeoRelay fetch work off the main actor (#1060)
* Run GeoRelay fetch pipeline off main actor * Capture GeoRelay session before detached fetch --------- Co-authored-by: jack <jackjackbits@users.noreply.github.com>
This commit is contained in:
@@ -17,8 +17,8 @@ struct GeoRelayDirectoryDependencies {
|
||||
var refreshCheckInterval: TimeInterval
|
||||
var retryInitialSeconds: TimeInterval
|
||||
var retryMaxSeconds: TimeInterval
|
||||
var awaitTorReady: () async -> Bool
|
||||
var fetchData: (URLRequest) async throws -> Data
|
||||
var awaitTorReady: @Sendable () async -> Bool
|
||||
var makeFetchData: @MainActor @Sendable () -> (@Sendable (URLRequest) async throws -> Data)
|
||||
var readData: (URL) -> Data?
|
||||
var writeData: (Data, URL) throws -> Void
|
||||
var cacheURL: () -> URL?
|
||||
@@ -50,9 +50,12 @@ private extension GeoRelayDirectoryDependencies {
|
||||
retryInitialSeconds: TransportConfig.geoRelayRetryInitialSeconds,
|
||||
retryMaxSeconds: TransportConfig.geoRelayRetryMaxSeconds,
|
||||
awaitTorReady: { await TorManager.shared.awaitReady() },
|
||||
fetchData: { request in
|
||||
let (data, _) = try await TorURLSession.shared.session.data(for: request)
|
||||
makeFetchData: {
|
||||
let session = TorURLSession.shared.session
|
||||
return { request in
|
||||
let (data, _) = try await session.data(for: request)
|
||||
return data
|
||||
}
|
||||
},
|
||||
readData: { try? Data(contentsOf: $0) },
|
||||
writeData: { data, url in
|
||||
@@ -110,12 +113,19 @@ final class GeoRelayDirectory {
|
||||
}
|
||||
}
|
||||
|
||||
struct Entry: Hashable {
|
||||
struct Entry: Hashable, Sendable {
|
||||
let host: String
|
||||
let lat: Double
|
||||
let lon: Double
|
||||
}
|
||||
|
||||
private enum DetachedFetchOutcome: Sendable {
|
||||
case success(entries: [Entry], csv: String)
|
||||
case torNotReady
|
||||
case invalidData
|
||||
case network(String)
|
||||
}
|
||||
|
||||
static let shared = GeoRelayDirectory()
|
||||
|
||||
private(set) var entries: [Entry] = []
|
||||
@@ -212,40 +222,62 @@ final class GeoRelayDirectory {
|
||||
cachePolicy: .reloadIgnoringLocalCacheData,
|
||||
timeoutInterval: 15
|
||||
)
|
||||
let awaitTorReady = dependencies.awaitTorReady
|
||||
let fetchData = dependencies.makeFetchData()
|
||||
|
||||
Task { [weak self] in
|
||||
guard let self else { return }
|
||||
|
||||
let ready = await self.dependencies.awaitTorReady()
|
||||
if !ready {
|
||||
let outcome = await Self.fetchRemoteOutcome(
|
||||
request: request,
|
||||
awaitTorReady: awaitTorReady,
|
||||
fetchData: fetchData
|
||||
)
|
||||
|
||||
switch outcome {
|
||||
case .success(let parsed, let csv):
|
||||
self.handleFetchSuccess(entries: parsed, csv: csv)
|
||||
case .torNotReady:
|
||||
self.handleFetchFailure(.torNotReady)
|
||||
return
|
||||
case .invalidData:
|
||||
self.handleFetchFailure(.invalidData)
|
||||
case .network(let description):
|
||||
self.handleFetchFailure(.network(description))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
nonisolated private static func fetchRemoteOutcome(
|
||||
request: URLRequest,
|
||||
awaitTorReady: @escaping @Sendable () async -> Bool,
|
||||
fetchData: @escaping @Sendable (URLRequest) async throws -> Data
|
||||
) async -> DetachedFetchOutcome {
|
||||
await Task.detached(priority: .utility) {
|
||||
let ready = await awaitTorReady()
|
||||
guard ready else { return .torNotReady }
|
||||
|
||||
do {
|
||||
let data = try await self.dependencies.fetchData(request)
|
||||
let data = try await fetchData(request)
|
||||
guard let text = String(data: data, encoding: .utf8) else {
|
||||
self.handleFetchFailure(.invalidData)
|
||||
return
|
||||
return .invalidData
|
||||
}
|
||||
|
||||
let parsed = GeoRelayDirectory.parseCSV(text)
|
||||
let parsed = Self.parseCSV(text)
|
||||
guard !parsed.isEmpty else {
|
||||
self.handleFetchFailure(.invalidData)
|
||||
return
|
||||
return .invalidData
|
||||
}
|
||||
|
||||
self.handleFetchSuccess(entries: parsed, csv: text)
|
||||
return .success(entries: parsed, csv: text)
|
||||
} catch {
|
||||
self.handleFetchFailure(.network(error))
|
||||
}
|
||||
return .network(error.localizedDescription)
|
||||
}
|
||||
}.value
|
||||
}
|
||||
|
||||
private enum FetchFailure {
|
||||
case torNotReady
|
||||
case invalidData
|
||||
case network(Error)
|
||||
case network(String)
|
||||
}
|
||||
|
||||
@MainActor
|
||||
@@ -266,8 +298,8 @@ final class GeoRelayDirectory {
|
||||
SecureLogger.warning("GeoRelayDirectory: Tor not ready; scheduling retry", category: .session)
|
||||
case .invalidData:
|
||||
SecureLogger.warning("GeoRelayDirectory: remote fetch returned invalid data; scheduling retry", category: .session)
|
||||
case .network(let error):
|
||||
SecureLogger.warning("GeoRelayDirectory: remote fetch failed with error: \(error.localizedDescription)", category: .session)
|
||||
case .network(let errorDescription):
|
||||
SecureLogger.warning("GeoRelayDirectory: remote fetch failed with error: \(errorDescription)", category: .session)
|
||||
}
|
||||
isFetching = false
|
||||
scheduleRetry()
|
||||
|
||||
@@ -145,6 +145,34 @@ final class GeoRelayDirectoryTests: XCTestCase {
|
||||
XCTAssertEqual(forcedRequestCount, 1)
|
||||
}
|
||||
|
||||
func test_prefetchIfNeeded_runsRemoteFetchOffMainThread() async {
|
||||
var factoryThreadFlags: [Bool] = []
|
||||
let threadRecorder = MainThreadRecorder()
|
||||
let harness = makeHarness(
|
||||
fetchCSV: """
|
||||
relay url,lat,lon
|
||||
background.example,8,9
|
||||
""",
|
||||
fetchFactoryObserver: {
|
||||
factoryThreadFlags.append(isExecutingOnMainThread())
|
||||
},
|
||||
fetchObserver: {
|
||||
await threadRecorder.record(isExecutingOnMainThread())
|
||||
}
|
||||
)
|
||||
let directory = GeoRelayDirectory(dependencies: harness.dependencies)
|
||||
|
||||
directory.prefetchIfNeeded()
|
||||
|
||||
let refreshed = await waitUntil {
|
||||
directory.entries == [GeoRelayDirectory.Entry(host: "background.example", lat: 8, lon: 9)]
|
||||
}
|
||||
XCTAssertTrue(refreshed)
|
||||
XCTAssertEqual(factoryThreadFlags, [true])
|
||||
let recordedValues = await threadRecorder.recordedValues()
|
||||
XCTAssertEqual(recordedValues, [false])
|
||||
}
|
||||
|
||||
func test_prefetchIfNeeded_failureSchedulesRetryAndRecoversOnNextFetch() async {
|
||||
let csv = """
|
||||
relay url,lat,lon
|
||||
@@ -215,6 +243,8 @@ final class GeoRelayDirectoryTests: XCTestCase {
|
||||
workingDirectoryCSV: String? = nil,
|
||||
fetchCSV: String? = nil,
|
||||
fetchResults: [Result<Data, Error>] = [],
|
||||
fetchFactoryObserver: (@MainActor @Sendable () -> Void)? = nil,
|
||||
fetchObserver: (@Sendable () async -> Void)? = nil,
|
||||
autoStart: Bool = false,
|
||||
activeNotificationName: Notification.Name? = nil
|
||||
) -> GeoRelayHarness {
|
||||
@@ -254,8 +284,12 @@ final class GeoRelayDirectoryTests: XCTestCase {
|
||||
retryInitialSeconds: 5,
|
||||
retryMaxSeconds: 40,
|
||||
awaitTorReady: { true },
|
||||
fetchData: { request in
|
||||
try await fetcher.fetch(request)
|
||||
makeFetchData: {
|
||||
fetchFactoryObserver?()
|
||||
return { request in
|
||||
await fetchObserver?()
|
||||
return try await fetcher.fetch(request)
|
||||
}
|
||||
},
|
||||
readData: { url in
|
||||
fileStore.dataByURL[url]
|
||||
@@ -359,6 +393,22 @@ private actor RetryDelayRecorder {
|
||||
}
|
||||
}
|
||||
|
||||
private actor MainThreadRecorder {
|
||||
private var values: [Bool] = []
|
||||
|
||||
func record(_ value: Bool) {
|
||||
values.append(value)
|
||||
}
|
||||
|
||||
func recordedValues() -> [Bool] {
|
||||
values
|
||||
}
|
||||
}
|
||||
|
||||
private enum GeoRelayTestError: Error {
|
||||
case network
|
||||
}
|
||||
|
||||
private func isExecutingOnMainThread() -> Bool {
|
||||
Thread.isMainThread
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user