Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion ContractTests/Source/Controllers/SdkController.swift
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,9 @@ final class SdkController: RouteCollection {
"client-prereq-events",
"client-prereq-cycle-detection",
"polling-gzip",
"client-per-context-summaries"
"client-per-context-summaries",
"retry-conformance-fdv1-streaming",
"retry-conformance-fdv1-polling"
]

return StatusResponse(
Expand Down
2 changes: 1 addition & 1 deletion LaunchDarkly.podspec
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,6 @@ Pod::Spec.new do |ld|
ld.swift_version = '5.0'

ld.subspec 'Core' do |es|
es.dependency 'LDSwiftEventSource', '3.3.1'
es.dependency 'LDSwiftEventSource', '3.4.0'
end
end
64 changes: 39 additions & 25 deletions LaunchDarkly.xcodeproj/project.pbxproj

Large diffs are not rendered by default.

19 changes: 5 additions & 14 deletions LaunchDarkly/LaunchDarkly/LDClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -676,37 +676,28 @@ public class LDClient {
case let .flagCollection((flagCollection, etag)):
os_log("%s: got flag collection with %d flags.", log: config.logger, type: .debug, typeName(and: #function), flagCollection.flags.count)
let oldStoredItems = flagStore.storedItems
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation, streamingMode: flagSynchronizer.streamingMode)
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation)
flagStore.replaceStore(newStoredItems: StoredItems(items: flagCollection.flags))
self.updateCacheAndReportChanges(context: self.context, oldStoredItems: oldStoredItems, etag: etag)
case let .patch(featureFlag):
let oldStoredItems = flagStore.storedItems
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation, streamingMode: flagSynchronizer.streamingMode)
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation)
flagStore.updateStore(updatedFlag: featureFlag)
self.updateCacheAndReportChanges(context: self.context, oldStoredItems: oldStoredItems, etag: nil)
case let .delete(deleteResponse):
let oldStoredItems = flagStore.storedItems
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation, streamingMode: flagSynchronizer.streamingMode)
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation)
flagStore.deleteFlag(deleteResponse: deleteResponse)
self.updateCacheAndReportChanges(context: self.context, oldStoredItems: oldStoredItems, etag: nil)
case .upToDate:
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation, streamingMode: flagSynchronizer.streamingMode)
connectionInformation = ConnectionInformation.checkEstablishingStreaming(connectionInformation: connectionInformation)
flagChangeNotifier.notifyUnchanged()
// If a polling request receives a 304 not modified, we still need
// to update the "last updated" field of the cache so subsequent
// restarts will honor the appropriate polling delay.
self.updateCacheFreshness(context: self.context)
case .error(let synchronizingError):
process(synchronizingError: synchronizingError, logPrefix: typeName(and: #function))
}
}

private func process(synchronizingError: SynchronizingError, logPrefix: String) {
connectionInformation = ConnectionInformation.synchronizingErrorCheck(synchronizingError: synchronizingError, connectionInformation: connectionInformation)
if synchronizingError.isTerminal {
os_log("%s data source terminal error; stopping flag delivery", log: config.logger, type: .debug, logPrefix)
flagSynchronizer.isOnline = false
initialized = true
connectionInformation = ConnectionInformation.recordSynchronizingError(synchronizingError, streamingMode: flagSynchronizer.streamingMode, connectionInformation: connectionInformation)
}
}

Expand Down
10 changes: 3 additions & 7 deletions LaunchDarkly/LaunchDarkly/Models/ConnectionInformation.swift
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ public struct ConnectionInformation: Codable, CustomStringConvertible {
}

// Used for parsing SynchronizingError in LDClient.process
static func synchronizingErrorCheck(synchronizingError: SynchronizingError, connectionInformation: ConnectionInformation) -> ConnectionInformation {
static func recordSynchronizingError(_ synchronizingError: SynchronizingError, streamingMode: LDStreamingMode, connectionInformation: ConnectionInformation) -> ConnectionInformation {
var connectionInformationVar = connectionInformation
if synchronizingError.isClientUnauthorized {
connectionInformationVar.lastConnectionFailureReason = .unauthorized
Expand All @@ -112,18 +112,14 @@ public struct ConnectionInformation: Codable, CustomStringConvertible {
}
connectionInformationVar.lastFailedConnection = Date()
if synchronizingError.isTerminal {
connectionInformationVar.currentConnectionMode = .offline
connectionInformationVar.currentConnectionMode = (streamingMode == .streaming) ? .establishingStreamingConnection : .polling
}
return connectionInformationVar
}

// Reconciles the connection mode and flag validity after a successful sync.
static func checkEstablishingStreaming(connectionInformation: ConnectionInformation, streamingMode: LDStreamingMode) -> ConnectionInformation {
static func checkEstablishingStreaming(connectionInformation: ConnectionInformation) -> ConnectionInformation {
var connectionInformationVar = connectionInformation
// Recover a terminal-error offline into the connecting mode. The checks below finish the transition.
if connectionInformationVar.currentConnectionMode == .offline {
connectionInformationVar.currentConnectionMode = (streamingMode == .streaming) ? .establishingStreamingConnection : .polling
}
if connectionInformationVar.currentConnectionMode == .establishingStreamingConnection {
connectionInformationVar.currentConnectionMode = .streaming
connectionInformationVar.lastKnownFlagValidity = nil
Expand Down
129 changes: 112 additions & 17 deletions LaunchDarkly/LaunchDarkly/ServiceObjects/FlagSynchronizer.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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))
Expand All @@ -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() {
Expand All @@ -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)
Comment thread
tanderson-ld marked this conversation as resolved.
// 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))
}

Expand All @@ -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())
}
}
Comment thread
tanderson-ld marked this conversation as resolved.

@objc private func processTimer() {
makeFlagRequest(isOnline: isOnline)
}
Expand Down Expand Up @@ -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 {
Expand All @@ -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)))
Expand Down Expand Up @@ -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
Comment thread
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 {
Expand Down Expand Up @@ -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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Close wipes healthy-stream backoff marker

Medium Severity

Clearing connectedAt in onClosed can run before eventSourceErrorHandler when an open stream fails. The error handler then sees a nil marker and skips the healthy-stream reset, so a connection that stayed up long enough still keeps the previous backoff.

Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 4d62600. Configure here.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't think this is actually a bug. swift-eventsource calls our error handler first and onClosed() only after it returns (LDSwiftEventSource.swift:282-285, both on one serial queue), so the healthy check always reads connectedAt before it's cleared. And the clear is needed anyway, to stop a stale marker leaking across an offline/online cycle.

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":
Expand Down Expand Up @@ -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
8 changes: 4 additions & 4 deletions LaunchDarkly/LaunchDarkly/ServiceObjects/LDTimer.swift
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ import Foundation
protocol TimeResponding {
var fireDate: Date? { get }

init(withTimeInterval: TimeInterval, fireQueue: DispatchQueue, fireAt: Date?, execute: @escaping () -> Void)
init(withTimeInterval: TimeInterval, fireQueue: DispatchQueue, fireAt: Date?, repeats: Bool, execute: @escaping () -> Void)
func cancel()
}

Expand All @@ -15,16 +15,16 @@ final class LDTimer: TimeResponding {
private(set) var isCancelled: Bool = false
var fireDate: Date? { timer?.fireDate }

init(withTimeInterval timeInterval: TimeInterval, fireQueue: DispatchQueue = DispatchQueue.main, fireAt: Date? = nil, execute: @escaping () -> Void) {
init(withTimeInterval timeInterval: TimeInterval, fireQueue: DispatchQueue = DispatchQueue.main, fireAt: Date? = nil, repeats: Bool = true, execute: @escaping () -> Void) {
self.fireQueue = fireQueue
self.execute = execute

// the run loop retains the timer, so the property is weak to avoid a retain cycle. Setting the timer to a strong reference is important so that the timer doesn't get nil'd before it's added to the run loop.
let timer: Timer
if let at = fireAt {
timer = Timer(fireAt: at, interval: timeInterval, target: self, selector: #selector(timerFired), userInfo: nil, repeats: true)
timer = Timer(fireAt: at, interval: timeInterval, target: self, selector: #selector(timerFired), userInfo: nil, repeats: repeats)
} else {
timer = Timer(timeInterval: timeInterval, target: self, selector: #selector(timerFired), userInfo: nil, repeats: true)
timer = Timer(timeInterval: timeInterval, target: self, selector: #selector(timerFired), userInfo: nil, repeats: repeats)
}
self.timer = timer
RunLoop.main.add(timer, forMode: RunLoop.Mode.default)
Expand Down
Loading
Loading