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
126 changes: 74 additions & 52 deletions Sources/OpenCombine/Publishers/Publishers.Throttle.swift
Original file line number Diff line number Diff line change
Expand Up @@ -144,15 +144,14 @@ extension Publishers.Throttle {
private var state: State
private let downstreamLock = UnfairRecursiveLock.allocate()

private var lastEmissionTime: Context.SchedulerTimeType?
private var nextEmissionTime: Context.SchedulerTimeType
private var hasScheduledOutput = false

private var pendingInput: Input?
private var pendingCompletion: Subscribers.Completion<Failure>?

private var demand: Subscribers.Demand = .none

private var lastTime: Context.SchedulerTimeType

init(interval: Context.SchedulerTimeType.Stride,
scheduler: Context,
latest: Bool,
Expand All @@ -162,7 +161,7 @@ extension Publishers.Throttle {
self.scheduler = scheduler
self.latest = latest

self.lastTime = scheduler.now
self.nextEmissionTime = scheduler.now
}

deinit {
Expand All @@ -177,16 +176,16 @@ extension Publishers.Throttle {
subscription.cancel()
return
}
self.lastTime = scheduler.now
self.nextEmissionTime = scheduler.now

state = .subscribed(subscription, downstream)
lock.unlock()

subscription.request(.unlimited)

downstreamLock.lock()
downstream.receive(subscription: self)
downstreamLock.unlock()

subscription.request(.unlimited)
}

func receive(_ input: Input) -> Subscribers.Demand {
Expand All @@ -196,42 +195,36 @@ extension Publishers.Throttle {
return .none
}

let lastTime = scheduler.now
self.lastTime = lastTime
let now = scheduler.now
let emitImmediately = now >= nextEmissionTime

guard demand > .none else {
lock.unlock()
return .none
}

let hasScheduledOutput = (pendingInput != nil || pendingCompletion != nil)

if hasScheduledOutput && latest {
if latest {
pendingInput = input
lock.unlock()
} else if !hasScheduledOutput {
let minimumEmissionTime =
lastEmissionTime.map { $0.advanced(by: interval) }

let emissionTime =
minimumEmissionTime.map { Swift.max(lastTime, $0) } ?? lastTime

demand -= 1

} else if emitImmediately {
// Each new window selects its first input, even without demand.
nextEmissionTime = now.advanced(by: interval)
pendingInput = input
} else if pendingInput == nil {
pendingInput = input
}

guard !hasScheduledOutput, demand > .none else {
lock.unlock()
return .none
}

let action: () -> Void = { [weak self] in
self?.scheduledEmission()
}
hasScheduledOutput = true
let emissionTime = nextEmissionTime
lock.unlock()

if emissionTime == lastTime {
scheduler.schedule(action)
} else {
scheduler.schedule(after: emissionTime, action)
if emitImmediately {
scheduler.schedule {
self.scheduledEmission()
}
} else {
lock.unlock()
scheduler.schedule(after: emissionTime) {
self.scheduledEmission()
}
}

return .none
Expand All @@ -240,27 +233,25 @@ extension Publishers.Throttle {
func receive(completion: Subscribers.Completion<Failure>) {
lock.lock()
guard case let .subscribed(subscription, downstream) = state else {
if !hasScheduledOutput {
state = .terminal
}
lock.unlock()
return
}
let lastTime = scheduler.now
self.lastTime = lastTime
nextEmissionTime = scheduler.now
state = .pendingTerminal(subscription, downstream)
pendingCompletion = completion

let hasScheduledOutput = (pendingInput != nil || pendingCompletion != nil)

if hasScheduledOutput && pendingCompletion == nil {
pendingCompletion = completion
if hasScheduledOutput {
lock.unlock()
} else if !hasScheduledOutput {
pendingCompletion = completion
} else {
hasScheduledOutput = true
lock.unlock()

scheduler.schedule { [weak self] in
self?.scheduledEmission()
scheduler.schedule {
self.scheduledEmission()
}
} else {
lock.unlock()
}
}

Expand All @@ -278,15 +269,27 @@ extension Publishers.Throttle {
downstream = foundDownstream
}

if self.pendingInput != nil && self.pendingCompletion == nil {
lastEmissionTime = scheduler.now
guard hasScheduledOutput else {
lock.unlock()
return
}

let pendingInput: Input?
if self.pendingInput != nil && demand > .none {
// Scheduled input can change until demand is consumed here.
demand -= 1
pendingInput = self.pendingInput.take()
} else {
pendingInput = nil
}

let pendingInput = self.pendingInput.take()
hasScheduledOutput = false
let pendingCompletion = self.pendingCompletion.take()

if pendingCompletion != nil {
state = .terminal
} else {
nextEmissionTime = scheduler.now.advanced(by: interval)
}

lock.unlock()
Expand All @@ -305,21 +308,37 @@ extension Publishers.Throttle {
}
downstreamLock.unlock()

guard newDemand > 0 else { return }
guard newDemand > 0, pendingCompletion == nil else { return }
self.lock.lock()
demand += newDemand
self.lock.unlock()
}

func request(_ demand: Subscribers.Demand) {
guard demand > 0 else { return }
lock.lock()
guard case .subscribed = state else {
lock.unlock()
return
}
self.demand += demand
guard pendingInput != nil, !hasScheduledOutput else {
lock.unlock()
return
}
hasScheduledOutput = true
let now = scheduler.now
let emissionTime = nextEmissionTime
lock.unlock()

if now >= emissionTime {
scheduler.schedule {
self.scheduledEmission()
}
} else {
scheduler.schedule(after: emissionTime) {
self.scheduledEmission()
}
}
}

func cancel() {
Expand All @@ -330,6 +349,9 @@ extension Publishers.Throttle {
case let .subscribed(existingSubscription, _),
let .pendingTerminal(existingSubscription, _):
subscription = existingSubscription
pendingInput = nil
pendingCompletion = nil
demand = .none
case .awaitingSubscription, .terminal:
subscription = nil
}
Expand Down
Loading