@@ -120,22 +120,37 @@ final class AccountProviderFactory: AccountProviderFactoryProtocol {
120120/// RobinHood's generic observable drops mapper failures for inserts and updates, but emits deletes
121121/// directly from the raw identifier. A corrupt row that duplicates a healthy wallet identifier can
122122/// therefore remove the healthy item from a stream when that corrupt row is deleted. Tracking
123- /// successfully mapped rows by object ID makes deletes unambiguous and also turns a valid-to-corrupt
124- /// update into removal of that one previously valid object.
125- private final class TolerantMetaAccountContextObservable < Model: Identifiable > :
123+ /// successfully mapped rows by object ID makes deletes unambiguous. Aggregating every notification
124+ /// by persistent identifier also makes same-identifier object handoffs independent of Core Data's
125+ /// unordered notification sets.
126+ final class TolerantMetaAccountContextObservable < Model: Identifiable > :
126127 DataProviderRepositoryObservable {
127128 private struct MappedObject {
128129 let objectId : NSManagedObjectID
129130 let model : Model ?
130131 }
131132
133+ private enum ObjectMutation {
134+ case replace( Model ? )
135+ case delete
136+ }
137+
138+ private struct InFlightDelivery {
139+ let observerIdentifier : ObjectIdentifier
140+ let identifiers : Set < String >
141+ }
142+
132143 private let service : CoreDataServiceProtocol
133144 private let mapper : AnyCoreDataMapper < Model , CDMetaAccount >
134145 private let predicate : ( CDMetaAccount ) -> Bool
135146 private let processingQueue : DispatchQueue
136147
137148 private var observers : [ RepositoryObserver < Model > ] = [ ]
138149 private var trackedModels : [ NSManagedObjectID : Model ] = [ : ]
150+ private var pendingReconciliationIdentifiers = Set < String > ( )
151+ // StreamableProvider can queue source callbacks behind its user-observer teardown.
152+ // Reconcile identifiers whose serial source deliveries were not acknowledged first.
153+ private var inFlightDeliveries : [ UUID : InFlightDelivery ] = [ : ]
139154 private var notificationToken : NSObjectProtocol ?
140155
141156 init (
@@ -192,6 +207,8 @@ private final class TolerantMetaAccountContextObservable<Model: Identifiable>:
192207 return ( $0. objectId, model)
193208 }
194209 )
210+ self . pendingReconciliationIdentifiers. removeAll ( )
211+ self . inFlightDeliveries. removeAll ( )
195212 completionBlock ( nil )
196213 }
197214 } catch {
@@ -216,6 +233,8 @@ private final class TolerantMetaAccountContextObservable<Model: Identifiable>:
216233
217234 self . processingQueue. async {
218235 self . trackedModels. removeAll ( )
236+ self . pendingReconciliationIdentifiers. removeAll ( )
237+ self . inFlightDeliveries. removeAll ( )
219238 completionBlock ( optionalError)
220239 }
221240 }
@@ -240,14 +259,44 @@ private final class TolerantMetaAccountContextObservable<Model: Identifiable>:
240259 updateBlock: updateBlock
241260 )
242261 )
262+
263+ let reconciliationChanges = self . reconciliationChanges (
264+ for: self . pendingReconciliationIdentifiers
265+ )
266+ self . pendingReconciliationIdentifiers. removeAll ( )
267+ self . deliver ( reconciliationChanges)
243268 }
244269 }
245270
246271 func removeObserver( _ observer: AnyObject ) {
247272 processingQueue. async {
273+ let observerIdentifier = ObjectIdentifier ( observer)
274+ let wasRegistered = self . observers. contains {
275+ $0. observer === observer
276+ }
248277 self . observers = self . observers. filter {
249278 $0. observer != nil && $0. observer !== observer
250279 }
280+
281+ if wasRegistered, self . observers. isEmpty {
282+ let deliveryTokens = self . inFlightDeliveries. compactMap {
283+ token, delivery in
284+ delivery. observerIdentifier == observerIdentifier
285+ ? token
286+ : nil
287+ }
288+ for token in deliveryTokens {
289+ guard let delivery = self . inFlightDeliveries. removeValue (
290+ forKey: token
291+ ) else {
292+ continue
293+ }
294+
295+ self . pendingReconciliationIdentifiers. formUnion (
296+ delivery. identifiers
297+ )
298+ }
299+ }
251300 }
252301 }
253302
@@ -274,72 +323,12 @@ private final class TolerantMetaAccountContextObservable<Model: Identifiable>:
274323 }
275324
276325 processingQueue. async {
277- var changes : [ DataProviderChange < Model > ] = [ ]
278-
279- for object in updatedObjects {
280- self . applyMappedObject (
281- object,
282- preferredChange: . update,
283- to: & changes
284- )
285- }
286- for objectId in deletedObjectIds {
287- if let previousModel = self . trackedModels. removeValue ( forKey: objectId) {
288- changes. append (
289- . delete( deletedIdentifier: previousModel. identifier)
290- )
291- }
292- }
293- for object in insertedObjects {
294- self . applyMappedObject (
295- object,
296- preferredChange: . insert,
297- to: & changes
298- )
299- }
300-
301- self . deliver ( changes)
302- }
303- }
304-
305- private enum PreferredChange {
306- case insert
307- case update
308- }
309-
310- private func applyMappedObject(
311- _ object: MappedObject ,
312- preferredChange: PreferredChange ,
313- to changes: inout [ DataProviderChange < Model > ]
314- ) {
315- guard let model = object. model else {
316- if let previousModel = trackedModels. removeValue ( forKey: object. objectId) {
317- changes. append (
318- . delete( deletedIdentifier: previousModel. identifier)
319- )
320- }
321- return
322- }
323-
324- if let previousModel = trackedModels [ object. objectId] {
325- if previousModel. identifier == model. identifier {
326- changes. append ( . update( newItem: model) )
327- } else {
328- changes. append (
329- . delete( deletedIdentifier: previousModel. identifier)
330- )
331- changes. append ( . insert( newItem: model) )
332- }
333- } else {
334- switch preferredChange {
335- case . insert, . update:
336- // A previously corrupt/nonmatching object becoming valid is new to the stream,
337- // even though Core Data reports it as an update.
338- changes. append ( . insert( newItem: model) )
339- }
326+ self . apply (
327+ updatedObjects: updatedObjects,
328+ deletedObjectIds: deletedObjectIds,
329+ insertedObjects: insertedObjects
330+ )
340331 }
341-
342- trackedModels [ object. objectId] = model
343332 }
344333
345334 private func mappedObject( for entity: CDMetaAccount ) -> MappedObject {
@@ -367,20 +356,180 @@ private final class TolerantMetaAccountContextObservable<Model: Identifiable>:
367356 return objects. allObjects. compactMap { $0 as? CDMetaAccount }
368357 }
369358
359+ private func apply(
360+ updatedObjects: [ MappedObject ] ,
361+ deletedObjectIds: [ NSManagedObjectID ] ,
362+ insertedObjects: [ MappedObject ]
363+ ) {
364+ var mutations : [ NSManagedObjectID : ObjectMutation ] = [ : ]
365+ for object in updatedObjects {
366+ mutations [ object. objectId] = . replace( object. model)
367+ }
368+ for object in insertedObjects {
369+ mutations [ object. objectId] = . replace( object. model)
370+ }
371+ for objectId in deletedObjectIds {
372+ mutations [ objectId] = . delete
373+ }
374+
375+ var affectedIdentifiers = Set < String > ( )
376+ for (objectId, mutation) in mutations {
377+ if let previousModel = trackedModels [ objectId] {
378+ affectedIdentifiers. insert ( previousModel. identifier)
379+ }
380+
381+ if case let . replace( model) = mutation, let model {
382+ affectedIdentifiers. insert ( model. identifier)
383+ }
384+ }
385+
386+ let modelsBefore = aggregateModels ( for: affectedIdentifiers)
387+
388+ for (objectId, mutation) in mutations {
389+ switch mutation {
390+ case let . replace( model) :
391+ trackedModels [ objectId] = model
392+ case . delete:
393+ trackedModels. removeValue ( forKey: objectId)
394+ }
395+ }
396+
397+ let modelsAfter = aggregateModels ( for: affectedIdentifiers)
398+ let changes = affectedIdentifiers. sorted ( ) . compactMap { identifier in
399+ aggregateChange (
400+ identifier: identifier,
401+ before: modelsBefore [ identifier] ,
402+ after: modelsAfter [ identifier]
403+ )
404+ }
405+ deliver ( changes)
406+ }
407+
408+ private func aggregateModels(
409+ for identifiers: Set < String >
410+ ) -> [ String : Model ] {
411+ guard !identifiers. isEmpty else {
412+ return [ : ]
413+ }
414+
415+ var models : [ String : Model ] = [ : ]
416+ for (_, model) in trackedModels. sorted ( by: objectIdOrder) {
417+ guard
418+ identifiers. contains ( model. identifier) ,
419+ models [ model. identifier] == nil
420+ else {
421+ continue
422+ }
423+
424+ models [ model. identifier] = model
425+ }
426+
427+ return models
428+ }
429+
430+ private func aggregateChange(
431+ identifier: String ,
432+ before: Model ? ,
433+ after: Model ?
434+ ) -> DataProviderChange < Model > ? {
435+ switch ( before, after) {
436+ case ( nil , nil ) :
437+ return nil
438+ case ( nil , let model? ) :
439+ return . insert( newItem: model)
440+ case ( . some, nil ) :
441+ return . delete( deletedIdentifier: identifier)
442+ case let ( . some, model? ) :
443+ return . update( newItem: model)
444+ }
445+ }
446+
447+ private func reconciliationChanges(
448+ for identifiers: Set < String >
449+ ) -> [ DataProviderChange < Model > ] {
450+ let currentModels = aggregateModels ( for: identifiers)
451+
452+ return identifiers. sorted ( ) . flatMap { identifier in
453+ var changes : [ DataProviderChange < Model > ] = [
454+ . delete( deletedIdentifier: identifier)
455+ ]
456+ if let model = currentModels [ identifier] {
457+ changes. append ( . insert( newItem: model) )
458+ }
459+
460+ return changes
461+ }
462+ }
463+
370464 private func deliver( _ changes: [ DataProviderChange < Model > ] ) {
371465 guard !changes. isEmpty else {
372466 return
373467 }
374468
375469 observers = observers. filter { $0. observer != nil }
470+ guard !observers. isEmpty else {
471+ pendingReconciliationIdentifiers. formUnion (
472+ changes. map ( changeIdentifier)
473+ )
474+ return
475+ }
476+
376477 for observer in observers {
377478 if processingQueue == observer. queue {
378479 observer. updateBlock ( changes)
379480 } else {
380- observer. queue. async {
481+ guard let observerObject = observer. observer else {
482+ continue
483+ }
484+
485+ let observerIdentifier = ObjectIdentifier ( observerObject)
486+ let identifiers = Set ( changes. map ( changeIdentifier) )
487+ let deliveryToken = beginDelivery (
488+ identifiers: identifiers,
489+ to: observerIdentifier
490+ )
491+ observer. queue. async { [ weak self] in
381492 observer. updateBlock ( changes)
493+ self ? . finishDelivery ( deliveryToken)
382494 }
383495 }
384496 }
385497 }
498+
499+ private func beginDelivery(
500+ identifiers: Set < String > ,
501+ to observerIdentifier: ObjectIdentifier
502+ ) -> UUID {
503+ let token = UUID ( )
504+ inFlightDeliveries [ token] = InFlightDelivery (
505+ observerIdentifier: observerIdentifier,
506+ identifiers: identifiers
507+ )
508+ return token
509+ }
510+
511+ private func finishDelivery( _ token: UUID ) {
512+ processingQueue. async {
513+ self . inFlightDeliveries. removeValue ( forKey: token)
514+ }
515+ }
516+
517+ private func changeIdentifier(
518+ _ change: DataProviderChange < Model >
519+ ) -> String {
520+ switch change {
521+ case let . insert( model) , let . update( model) :
522+ return model. identifier
523+ case let . delete( identifier) :
524+ return identifier
525+ }
526+ }
527+
528+ private func objectIdOrder(
529+ _ lhs: ( key: NSManagedObjectID , value: Model ) ,
530+ _ rhs: ( key: NSManagedObjectID , value: Model )
531+ ) -> Bool {
532+ lhs. key. uriRepresentation ( ) . absoluteString <
533+ rhs. key. uriRepresentation ( ) . absoluteString
534+ }
386535}
0 commit comments