1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
|
//
// MutuallyExclusive.swift
// MullvadVPN
//
// Created by pronebird on 24/10/2019.
// Copyright © 2019 Mullvad VPN AB. All rights reserved.
//
import Combine
import Foundation
extension Publishers {
/// A publisher that blocks the given DispatchQueue until the produced publisher reported the
/// completion.
final class MutuallyExclusive<PublisherType, Context>: Publisher
where
PublisherType: Publisher,
Context: Scheduler
{
typealias MakePublisherBlock = () -> PublisherType
typealias Output = PublisherType.Output
typealias Failure = PublisherType.Failure
private let exclusivityQueue: Context
private let executionQueue: Context
private let createPublisher: MakePublisherBlock
init(exclusivityQueue: Context, executionQueue: Context, createPublisher: @escaping MakePublisherBlock) {
self.exclusivityQueue = exclusivityQueue
self.executionQueue = executionQueue
self.createPublisher = createPublisher
}
func receive<S>(subscriber: S) where S : Subscriber, S.Failure == Failure, S.Input == Output {
let subscription = MutuallyExclusive.Subscription(
subscriber: subscriber,
createPublisher: createPublisher,
exclusivityQueue: exclusivityQueue,
executionQueue: executionQueue)
subscriber.receive(subscription: subscription)
}
}
}
private extension Publishers.MutuallyExclusive {
/// A subscription used by `MutuallyExclusive` publisher
final class Subscription<SubscriberType, PublisherType, Context>: Combine.Subscription
where
SubscriberType: Subscriber, PublisherType: Publisher,
PublisherType.Output == SubscriberType.Input,
PublisherType.Failure == SubscriberType.Failure,
Context: Scheduler
{
typealias MakePublisherBlock = () -> PublisherType
private let subscriber: SubscriberType
private var innerSubscriber: AnyCancellable?
private let createPublisher: MakePublisherBlock
private let exclusivityQueue: Context
private let executionQueue: Context
private let sema = DispatchSemaphore(value: 0)
private let cancelLock = NSLock()
private var isCancelled = false
init(subscriber: SubscriberType,
createPublisher: @escaping MakePublisherBlock,
exclusivityQueue: Context,
executionQueue: Context)
{
self.subscriber = subscriber
self.createPublisher = createPublisher
self.exclusivityQueue = exclusivityQueue
self.executionQueue = executionQueue
}
func request(_ demand: Subscribers.Demand) {
self.exclusivityQueue.schedule {
self.executionQueue.schedule {
self.cancelLock.withCriticalBlock {
guard !self.isCancelled else { return }
self.innerSubscriber = self.createPublisher()
.sink(receiveCompletion: { [weak self] (completion) in
guard let self = self else { return }
self.subscriber.receive(completion: completion)
self.signalSemaphore()
}, receiveValue: { [weak self] (output) in
_ = self?.subscriber.receive(output)
})
}
}
self.sema.wait()
}
}
func cancel() {
cancelLock.withCriticalBlock {
guard !isCancelled else { return }
isCancelled = true
innerSubscriber?.cancel()
innerSubscriber = nil
signalSemaphore()
}
}
private func signalSemaphore() {
_ = sema.signal()
}
}
}
typealias MutuallyExclusive = Publishers.MutuallyExclusive
|