Skip to content

Commit b9f732e

Browse files
nan-licursoragent
andcommitted
fix: [SDK-4874] single-flight UpdateSubscription to stop in-flight PATCH races
Gate concurrent updates per subscription model, drain pending after completion, and clear sentToClient on retryable failure so the gate does not stick. Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 82d0a16 commit b9f732e

3 files changed

Lines changed: 168 additions & 0 deletions

File tree

iOS_SDK/OneSignalSDK/OneSignalCoreMocks/MockOneSignalClient.swift

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,9 +36,15 @@ public class MockOneSignalClient: NSObject, IOneSignalClient {
3636
public var lastHTTPRequest: OneSignalRequest?
3737
public var networkRequestCount = 0
3838
public var executedRequests: [OneSignalRequest] = []
39+
/// Requests that have entered `execute` (including those still held / delayed).
40+
public private(set) var startedRequests: [OneSignalRequest] = []
3941
public var executeInstantaneously = false
4042
/// Set to true to make it unnecessary to setup mock responses for every request possible
4143
public var fireSuccessForAllRequests = false
44+
/// When true, `execute` records the request but does not complete until `releaseHeldResponses()`.
45+
public var holdResponses = false
46+
47+
private var heldExecutions: [(request: OneSignalRequest, onSuccess: OSResultSuccessBlock, onFailure: OSClientFailureBlock)] = []
4248

4349
var remoteParamsResponse: [String: Any]?
4450
var shouldUseProvisionalAuthorization = false // new in iOS 12 (aka Direct to History)
@@ -84,6 +90,9 @@ public class MockOneSignalClient: NSObject, IOneSignalClient {
8490
lastHTTPRequest = nil
8591
networkRequestCount = 0
8692
executedRequests.removeAll()
93+
startedRequests.removeAll()
94+
heldExecutions.removeAll()
95+
holdResponses = false
8796
executeInstantaneously = true
8897
remoteParamsResponse = nil
8998
shouldUseProvisionalAuthorization = false
@@ -93,6 +102,17 @@ public class MockOneSignalClient: NSObject, IOneSignalClient {
93102
public func execute(_ request: OneSignalRequest, onSuccess successBlock: @escaping OSResultSuccessBlock, onFailure failureBlock: @escaping OSClientFailureBlock) {
94103
print("🧪 MockOneSignalClient execute called")
95104

105+
lock.withLock {
106+
startedRequests.append(request)
107+
}
108+
109+
if holdResponses {
110+
lock.withLock {
111+
heldExecutions.append((request, successBlock, failureBlock))
112+
}
113+
return
114+
}
115+
96116
if executeInstantaneously {
97117
finishExecutingRequest(request, onSuccess: successBlock, onFailure: failureBlock)
98118
} else {
@@ -102,6 +122,18 @@ public class MockOneSignalClient: NSObject, IOneSignalClient {
102122
}
103123
}
104124

125+
/// Completes every request currently held by `holdResponses`.
126+
public func releaseHeldResponses() {
127+
let held: [(request: OneSignalRequest, onSuccess: OSResultSuccessBlock, onFailure: OSClientFailureBlock)] = lock.withLock {
128+
let copy = heldExecutions
129+
heldExecutions.removeAll()
130+
return copy
131+
}
132+
for item in held {
133+
finishExecutingRequest(item.request, onSuccess: item.onSuccess, onFailure: item.onFailure)
134+
}
135+
}
136+
105137
/// Helper method to stringify the name of a request for identification and comparison
106138
private func stringify(_ request: OneSignalRequest) -> String {
107139
var stringified = request.description

iOS_SDK/OneSignalSDK/OneSignalUser/Source/Executors/OSSubscriptionOperationExecutor.swift

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -397,6 +397,12 @@ class OSSubscriptionOperationExecutor: OSOperationExecutor {
397397
guard !request.sentToClient else {
398398
return
399399
}
400+
// Single-flight: don't send a second UpdateSubscription for the same model while one is in flight.
401+
// A later coalesced request stays queued and is drained when the in-flight one finishes.
402+
let modelId = request.subscriptionModel.modelId
403+
guard !updateRequestQueue.contains(where: { $0 !== request && $0.sentToClient && $0.subscriptionModel.modelId == modelId }) else {
404+
return
405+
}
400406
guard request.prepareForExecution(newRecordsState: newRecordsState) else {
401407
return
402408
}
@@ -413,6 +419,7 @@ class OSSubscriptionOperationExecutor: OSOperationExecutor {
413419
self.dispatchQueue.async {
414420
self.updateRequestQueue.removeAll(where: { $0 == request})
415421
OneSignalUserDefaults.initShared().saveCodeableData(forKey: OS_SUBSCRIPTION_EXECUTOR_UPDATE_REQUEST_QUEUE_KEY, withValue: self.updateRequestQueue)
422+
self.executeNextPendingUpdateSubscription(for: modelId, inBackground: inBackground)
416423
if inBackground {
417424
OSBackgroundTaskManager.endBackgroundTask(backgroundTaskIdentifier)
418425
}
@@ -440,11 +447,27 @@ class OSSubscriptionOperationExecutor: OSOperationExecutor {
440447
// Fail, no retry, remove from cache and queue
441448
self.updateRequestQueue.removeAll(where: { $0 == request})
442449
OneSignalUserDefaults.initShared().saveCodeableData(forKey: OS_SUBSCRIPTION_EXECUTOR_UPDATE_REQUEST_QUEUE_KEY, withValue: self.updateRequestQueue)
450+
self.executeNextPendingUpdateSubscription(for: modelId, inBackground: inBackground)
451+
} else {
452+
// Make the request eligible for the next flush
453+
request.sentToClient = false
443454
}
444455
if inBackground {
445456
OSBackgroundTaskManager.endBackgroundTask(backgroundTaskIdentifier)
446457
}
447458
}
448459
}
449460
}
461+
462+
}
463+
464+
extension OSSubscriptionOperationExecutor {
465+
/// Sends the oldest unsent UpdateSubscription for `modelId`, if any. Caller must be on `dispatchQueue`.
466+
private func executeNextPendingUpdateSubscription(for modelId: String, inBackground: Bool) {
467+
let pending = updateRequestQueue.filter { !$0.sentToClient && $0.subscriptionModel.modelId == modelId }
468+
guard let next = pending.min(by: { $0.timestamp < $1.timestamp }) else {
469+
return
470+
}
471+
executeUpdateSubscriptionRequest(next, inBackground: inBackground)
472+
}
450473
}

iOS_SDK/OneSignalSDK/OneSignalUserTests/Executors/SubscriptionUpdateRaceTests.swift

Lines changed: 113 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,119 @@ final class SubscriptionUpdateRaceTests: XCTestCase {
133133
XCTAssertEqual(updateRequests.count, 1, "Unsent updates for the same subscription should be coalesced")
134134
}
135135

136+
/**
137+
An UpdateSubscription already on the wire must not race a follow-up PATCH.
138+
The follow-up stays queued until the in-flight request completes, then sends live state.
139+
*/
140+
func testInFlightUpdateBlocksFollowUpUntilCompleteThenSendsLiveState() throws {
141+
let client = MockOneSignalClient()
142+
client.holdResponses = true
143+
client.fireSuccessForAllRequests = true
144+
OneSignalCoreImpl.setSharedClient(client)
145+
146+
let executor = OSSubscriptionOperationExecutor(newRecordsState: OSNewRecordsState())
147+
let model = makePushSubscriptionModel(notificationTypes: promptedNeverAnswered, subscriptionId: subscriptionId)
148+
let identityModelId = UUID().uuidString
149+
150+
executor.enqueueDelta(OSDelta(
151+
name: OS_UPDATE_SUBSCRIPTION_DELTA,
152+
identityModelId: identityModelId,
153+
model: model,
154+
property: "notificationTypes",
155+
value: promptedNeverAnswered
156+
))
157+
executor.processDeltaQueue(inBackground: false)
158+
OneSignalCoreMocks.waitForBackgroundThreads(seconds: 0.2)
159+
160+
XCTAssertEqual(client.startedRequests.count, 1, "First UpdateSubscription should be in flight")
161+
let firstPayload = try XCTUnwrap(
162+
(client.startedRequests[0] as? OSRequestUpdateSubscription)?.parameters?["subscription"] as? [String: Any]
163+
)
164+
XCTAssertEqual(firstPayload["notification_types"] as? Int, promptedNeverAnswered)
165+
XCTAssertEqual(firstPayload["enabled"] as? Bool, false)
166+
167+
// Permission granted while first PATCH is still in flight.
168+
model.notificationTypes = subscribedNotificationTypes
169+
XCTAssertTrue(model.enabled)
170+
171+
executor.enqueueDelta(OSDelta(
172+
name: OS_UPDATE_SUBSCRIPTION_DELTA,
173+
identityModelId: identityModelId,
174+
model: model,
175+
property: "notificationTypes",
176+
value: subscribedNotificationTypes
177+
))
178+
executor.processDeltaQueue(inBackground: false)
179+
OneSignalCoreMocks.waitForBackgroundThreads(seconds: 0.2)
180+
181+
XCTAssertEqual(client.startedRequests.count, 1, "Follow-up must wait for in-flight UpdateSubscription")
182+
183+
client.holdResponses = false
184+
client.releaseHeldResponses()
185+
OneSignalCoreMocks.waitForBackgroundThreads(seconds: 0.5)
186+
187+
XCTAssertEqual(client.startedRequests.count, 2, "Pending follow-up should send after in-flight completes")
188+
let secondPayload = try XCTUnwrap(
189+
(client.startedRequests[1] as? OSRequestUpdateSubscription)?.parameters?["subscription"] as? [String: Any]
190+
)
191+
XCTAssertEqual(secondPayload["notification_types"] as? Int, subscribedNotificationTypes)
192+
XCTAssertEqual(secondPayload["enabled"] as? Bool, true)
193+
}
194+
195+
/**
196+
A retryable failure (e.g. 500/timeout) must not leave the single-flight gate locked.
197+
The failed request becomes resendable and later updates for the model still go out.
198+
*/
199+
func testRetryableFailureDoesNotBlockSubsequentUpdates() throws {
200+
let client = MockOneSignalClient()
201+
client.fireSuccessForAllRequests = true
202+
OneSignalCoreImpl.setSharedClient(client)
203+
204+
let executor = OSSubscriptionOperationExecutor(newRecordsState: OSNewRecordsState())
205+
let model = makePushSubscriptionModel(notificationTypes: promptedNeverAnswered, subscriptionId: subscriptionId)
206+
let identityModelId = UUID().uuidString
207+
208+
// Fail the first update with a retryable error (mock responses are keyed by request description).
209+
let requestKey = "OSRequestUpdateSubscription with model: \(model.modelId)"
210+
client.setMockFailureResponseForRequest(
211+
request: requestKey,
212+
error: OneSignalClientError(code: 500, message: "retryable", responseHeaders: nil, response: nil, underlyingError: nil)
213+
)
214+
215+
executor.enqueueDelta(OSDelta(
216+
name: OS_UPDATE_SUBSCRIPTION_DELTA,
217+
identityModelId: identityModelId,
218+
model: model,
219+
property: "notificationTypes",
220+
value: promptedNeverAnswered
221+
))
222+
executor.processDeltaQueue(inBackground: false)
223+
OneSignalCoreMocks.waitForBackgroundThreads(seconds: 0.5)
224+
225+
XCTAssertEqual(client.executedRequests.count, 1, "First update should have been attempted and failed retryably")
226+
227+
// Server recovers; user accepts permission.
228+
client.setMockResponseForRequest(request: requestKey, response: [:])
229+
model.notificationTypes = subscribedNotificationTypes
230+
XCTAssertTrue(model.enabled)
231+
232+
executor.enqueueDelta(OSDelta(
233+
name: OS_UPDATE_SUBSCRIPTION_DELTA,
234+
identityModelId: identityModelId,
235+
model: model,
236+
property: "notificationTypes",
237+
value: subscribedNotificationTypes
238+
))
239+
executor.processDeltaQueue(inBackground: false)
240+
OneSignalCoreMocks.waitForBackgroundThreads(seconds: 0.5)
241+
242+
let updateRequests = client.executedRequests.compactMap { $0 as? OSRequestUpdateSubscription }
243+
XCTAssertEqual(updateRequests.count, 2, "Follow-up update must still send after a retryable failure")
244+
let lastPayload = try XCTUnwrap(updateRequests.last?.parameters?["subscription"] as? [String: Any])
245+
XCTAssertEqual(lastPayload["notification_types"] as? Int, subscribedNotificationTypes)
246+
XCTAssertEqual(lastPayload["enabled"] as? Bool, true)
247+
}
248+
136249
// MARK: - Helpers
137250

138251
private func makePushSubscriptionModel(notificationTypes: Int, subscriptionId: String?) -> OSSubscriptionModel {

0 commit comments

Comments
 (0)