-
Notifications
You must be signed in to change notification settings - Fork 94
feat: Retry data sources on terminal errors instead of stopping #533
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
09a7f82
a879ebf
8a46f12
544b132
dd1b646
d30eb21
f116264
4d62600
3808b60
a4d838e
80a95f4
147d849
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Large diffs are not rendered by default.
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -67,11 +67,18 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
|
|
||
| let service: DarklyServiceProvider | ||
| private var eventSource: DarklyStreamingProvider? | ||
| // Only accessed on isOnlineQueue. | ||
| private var flagRequestTimer: TimeResponding? | ||
| var onSyncComplete: FlagSyncCompleteClosure? | ||
|
|
||
| let streamingMode: LDStreamingMode | ||
|
|
||
| // Only accessed on the event source callback queue. | ||
| private let streamingRetry: StreamingRetryState | ||
| // Only accessed on isOnlineQueue. | ||
| private let pollingRetry: PollingRetryState | ||
| private static let healthyResetThreshold: TimeInterval = 60 | ||
|
|
||
| var isOnline: Bool { | ||
| get { isOnlineQueue.sync { _isOnline } } | ||
| set { | ||
|
|
@@ -91,6 +98,10 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
|
|
||
| private var syncQueue = DispatchQueue(label: Constants.queueName, qos: .utility) | ||
| private var eventSourceStarted: Date? | ||
| // Only accessed on the event source callback queue. | ||
| private var connectedAt: Date? | ||
| // Only accessed on isOnlineQueue. | ||
| private var reconnectTimer: TimeResponding? | ||
|
|
||
| init(streamingMode: LDStreamingMode, | ||
| pollingInterval: TimeInterval, | ||
|
|
@@ -100,6 +111,8 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| onSyncComplete: FlagSyncCompleteClosure?) { | ||
| self.streamingMode = streamingMode | ||
| self.pollingInterval = pollingInterval | ||
| self.streamingRetry = StreamingRetryState() | ||
| self.pollingRetry = PollingRetryState(pollInterval: pollingInterval) | ||
| self.useReport = useReport | ||
| self.lastCachedRequestedTime = lastUpdated | ||
| self.service = service | ||
|
|
@@ -142,6 +155,8 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| } | ||
|
|
||
| private func stopEventSource() { | ||
| reconnectTimer?.cancel() | ||
| reconnectTimer = nil | ||
| guard eventSource != nil | ||
| else { | ||
| os_log("%s aborted. Clientstream is not connected.", log: service.config.logger, type: .debug, typeName(and: #function)) | ||
|
|
@@ -153,6 +168,19 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| eventSource = nil | ||
| } | ||
|
|
||
| private func scheduleReconnect(after delay: TimeInterval) { | ||
| reconnectTimer?.cancel() | ||
| reconnectTimer = LDTimer(withTimeInterval: delay, fireQueue: isOnlineQueue, repeats: false) { [weak self] in | ||
| self?.reconnect() | ||
| } | ||
| } | ||
|
|
||
| private func reconnect() { | ||
| guard _isOnline | ||
| else { return } | ||
| startEventSource() | ||
| } | ||
|
|
||
| // MARK: Polling | ||
|
|
||
| private func startPolling() { | ||
|
|
@@ -162,20 +190,18 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| return | ||
| } | ||
|
|
||
| // We should fire right away, unless we know how fresh the cache is and can | ||
| // adjust accordingly. | ||
| var fireAt = Date.distantPast | ||
| var initialDelay: TimeInterval = 0 | ||
| if let lastTime = self.lastCachedRequestedTime { | ||
| fireAt = lastTime.addingTimeInterval(pollingInterval) | ||
| // If we do consider the cached values already fresh enough, we should | ||
| // signal completion immediately | ||
| // The cache is still fresh, so delay the first poll until the cached flags would go stale. | ||
| initialDelay = max(0, lastTime.addingTimeInterval(pollingInterval).timeIntervalSinceNow) | ||
| // We loaded cached flags, so report completion immediately rather than blocking on the first poll. | ||
| syncQueue.async { [self] in | ||
| guard isOnline | ||
| else { return } | ||
| reportSyncComplete(.upToDate) | ||
| } | ||
| } | ||
| flagRequestTimer = LDTimer(withTimeInterval: pollingInterval, fireQueue: syncQueue, fireAt: fireAt, execute: processTimer) | ||
| scheduleNextPoll(after: initialDelay) | ||
| os_log("%s", log: service.config.logger, type: .debug, typeName(and: #function)) | ||
| } | ||
|
|
||
|
|
@@ -191,6 +217,24 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| flagRequestTimer = nil | ||
| } | ||
|
|
||
| private func scheduleNextPoll(after delay: TimeInterval) { | ||
| flagRequestTimer?.cancel() | ||
| flagRequestTimer = LDTimer(withTimeInterval: delay, fireQueue: syncQueue, repeats: false, execute: processTimer) | ||
| } | ||
|
|
||
| private func pollDidComplete(failed: Bool, unexpected: Bool) { | ||
| isOnlineQueue.async { [weak self] in | ||
| guard let self = self, self._isOnline, self.streamingMode == .polling | ||
| else { return } | ||
| if failed { | ||
| self.pollingRetry.recordFailure(unexpected: unexpected) | ||
| } else { | ||
| self.pollingRetry.recordSuccess() | ||
| } | ||
| self.scheduleNextPoll(after: self.pollingRetry.nextDelay()) | ||
| } | ||
| } | ||
|
tanderson-ld marked this conversation as resolved.
|
||
|
|
||
| @objc private func processTimer() { | ||
| makeFlagRequest(isOnline: isOnline) | ||
| } | ||
|
|
@@ -238,9 +282,14 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| service.resetFlagResponseCache(etag: nil) | ||
| return | ||
| } | ||
| var failed = false | ||
| var unexpected = false | ||
| defer { pollDidComplete(failed: failed, unexpected: unexpected) } | ||
|
|
||
| if let serviceResponseError = serviceResponse.error { | ||
| os_log("%s error: %s", log: service.config.logger, type: .debug, typeName(and: #function), String(describing: serviceResponseError)) | ||
| reportSyncComplete(.error(.request(serviceResponseError))) | ||
| failed = true | ||
| return | ||
| } | ||
| if serviceResponse.urlResponse?.httpStatusCode == HTTPURLResponse.StatusCodes.notModified { | ||
|
|
@@ -250,13 +299,17 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| guard serviceResponse.urlResponse?.httpStatusCode == HTTPURLResponse.StatusCodes.ok | ||
| else { | ||
| os_log("%s response: %s", log: service.config.logger, type: .debug, typeName(and: #function), String(describing: serviceResponse.urlResponse)) | ||
| reportSyncComplete(.error(.response(serviceResponse.urlResponse))) | ||
| let syncError = SynchronizingError.response(serviceResponse.urlResponse) | ||
| reportSyncComplete(.error(syncError)) | ||
| failed = true | ||
| unexpected = syncError.isTerminal | ||
| return | ||
| } | ||
| guard let data = serviceResponse.data, | ||
| let flagCollection = try? JSONDecoder().decode(FeatureFlagCollection.self, from: data) | ||
| else { | ||
| reportDataError(serviceResponse.data) | ||
| failed = true | ||
| return | ||
| } | ||
| reportSyncComplete(.flagCollection((flagCollection, serviceResponse.etag))) | ||
|
|
@@ -290,16 +343,23 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
| } | ||
| eventSourceStarted = now | ||
|
|
||
| guard let unsuccessfulResponseError = error as? UnsuccessfulResponseError | ||
| else { return .proceed } | ||
| // Now we know that we received an error HTTP response code | ||
| let responseCode: Int = unsuccessfulResponseError.responseCode | ||
| if HTTPURLResponse.StatusCodes.isTerminalStatusCode(responseCode) { | ||
| reportSyncComplete(.error(.streamError(error))) | ||
| return .shutdown | ||
| // A stream that stayed up long enough before failing clears the backoff. | ||
| if let connectedAt = connectedAt, now.timeIntervalSince(connectedAt) >= FlagSynchronizer.healthyResetThreshold { | ||
| streamingRetry.reset() | ||
| } | ||
| connectedAt = nil | ||
|
cursor[bot] marked this conversation as resolved.
|
||
|
|
||
| streamingRetry.recordFailure(unexpected: SynchronizingError.streamError(error).isTerminal) | ||
| let delay = streamingRetry.nextDelay() | ||
| os_log("%s stream error; reconnecting in %.3fs. error: %s", log: service.config.logger, type: .debug, typeName(and: #function), delay, String(describing: error)) | ||
| reportSyncComplete(.error(.streamError(error))) | ||
|
|
||
| // Tear down the failed connection now, then reconnect after the backoff. | ||
| isOnlineQueue.async { [weak self] in | ||
| self?.stopEventSource() | ||
| self?.scheduleReconnect(after: delay) | ||
| } | ||
| // Otherwise we will retry | ||
| return .proceed | ||
| return .shutdown | ||
| } | ||
|
|
||
| func shouldAbortStreamUpdate() -> Bool { | ||
|
|
@@ -335,13 +395,19 @@ class FlagSynchronizer: LDFlagSynchronizing, EventHandler { | |
|
|
||
| public func onClosed() { | ||
| os_log("%s EventSource closed", log: service.config.logger, type: .debug, typeName(and: #function)) | ||
| connectedAt = nil | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Close wipes healthy-stream backoff markerMedium Severity Clearing Reviewed by Cursor Bugbot for commit 4d62600. Configure here.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I don't think this is actually a bug. |
||
| NotificationCenter.default.post(name: Notification.Name(FlagSynchronizer.Constants.didCloseEventSourceName), object: nil) | ||
| } | ||
|
|
||
| public func onMessage(eventType: String, messageEvent: MessageEvent) { | ||
| guard !shouldAbortStreamUpdate() | ||
| else { return } | ||
|
|
||
| // The first payload on a fresh stream marks healthy operation. | ||
| if connectedAt == nil { | ||
| connectedAt = Date() | ||
| } | ||
|
|
||
| switch eventType { | ||
| case "ping": makeFlagRequest(isOnline: isOnline) | ||
| case "put": | ||
|
|
@@ -411,6 +477,35 @@ extension FlagSynchronizer { | |
| func testProcessFlagResponse(serviceResponse: ServiceResponse) { | ||
| processFlagResponse(serviceResponse: serviceResponse) | ||
| } | ||
|
|
||
| // connectedAt is set on the event source callback queue. | ||
| // A test uses this on the thread where it drives eventSourceErrorHandler. | ||
| var testConnectedAt: Date? { | ||
| get { connectedAt } | ||
| set { connectedAt = newValue } | ||
| } | ||
|
|
||
| // Marks the synchronizer online without starting a data source, so a test can drive pollDidComplete without real polls racing the assertions. | ||
| func testForceOnline() { | ||
| isOnlineQueue.sync { _isOnline = true } | ||
| } | ||
|
|
||
| func testPollDidComplete(failed: Bool, unexpected: Bool) { | ||
| pollDidComplete(failed: failed, unexpected: unexpected) | ||
| } | ||
|
|
||
| // Reads the poll delay on isOnlineQueue so it serializes after a pending pollDidComplete. | ||
| var testNextPollDelay: TimeInterval { | ||
| isOnlineQueue.sync { pollingRetry.nextDelay() } | ||
| } | ||
|
|
||
| func testReconnect() { | ||
| reconnect() | ||
| } | ||
|
|
||
| var testReconnectFireDate: Date? { | ||
| isOnlineQueue.sync { reconnectTimer?.fireDate } | ||
| } | ||
| } | ||
|
|
||
| #endif | ||


Uh oh!
There was an error while loading. Please reload this page.