Repository navigation
Throttle: hold input received without downstream demand - #25
piercifani wants to merge 1 commit into
Conversation
Apple's Combine keeps a value that arrives while downstream demand is zero (the latest when latest == true, otherwise the first) and emits it once demand arrives. OpenCombine dropped it, so a subscriber whose demand lands after the upstream's first value (e.g. a replayed snapshot) never saw it. testInputWithoutDemandIsHeldUntilDemand fails before this change and passes both here and against Apple's Combine (OPENCOMBINE_COMPATIBILITY_TEST). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
I'll review this after I fixed the recent CI issue on this repo. |
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #25 +/- ##
==========================================
- Coverage 96.58% 96.07% -0.51%
==========================================
Files 107 107
Lines 8311 8310 -1
==========================================
- Hits 8027 7984 -43
- Misses 284 326 +42 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
Your PR has the following issue and does not align with Combine: The separate The examples below use a one-second interval and the repository's 1. A held value can arrive after a newer valueUse
Save the following as import XCTest
#if OPENCOMBINE_COMPATIBILITY_TEST
import Combine
#else
import OpenCombine
#endif
@available(macOS 10.15, iOS 13.0, tvOS 13.0, watchOS 6.0, *)
final class ThrottleDemandReproTests: XCTestCase {
func testHeldInputDoesNotArriveAfterNewerInput() throws {
let scheduler = VirtualTimeScheduler()
let subject = PassthroughSubject<Int, Never>()
var downstreamSubscription: Subscription?
var values: [Int] = []
let subscriber = AnySubscriber<Int, Never>(
receiveSubscription: {
downstreamSubscription = $0
$0.request(.max(1))
},
receiveValue: {
values.append($0)
return .none
},
receiveCompletion: { _ in }
)
subject.throttle(for: .seconds(1), scheduler: scheduler, latest: false)
.subscribe(subscriber)
let subscription = try XCTUnwrap(downstreamSubscription)
defer { subscription.cancel() }
subject.send(1)
subject.send(2)
// Add demand while the first output is still pending.
subscription.request(.max(1))
scheduler.executeScheduledActions()
XCTAssertEqual(values, [1])
subject.send(3)
scheduler.executeScheduledActions()
XCTAssertEqual(values, [1, 3])
// A new request must not emit an old input.
subscription.request(.max(1))
scheduler.executeScheduledActions()
print("Received:", values)
XCTAssertEqual(values, [1, 3])
}
}Run this from the repository root on this PR's branch: swift test --filter ThrottleDemandReproTestsThe last assertion fails, and the output is To compare the same test with Apple Combine: OPENCOMBINE_COMPATIBILITY_TEST=1 swift test \
--scratch-path .build-apple-throttle-repro \
--filter ThrottleDemandReproTestsThis passes and prints The first input consumes demand and becomes 2. One pending output becomes two outputsWith initial demand
The wrong first value with 3. The held input does not follow the throttle windowWith
|
|
Close in favor of #28 |
|
If you spot any diff result compared with Combine, an new issue is always welcomed. eg. #29 |
Problem
Publishers.Throttlerequests.unlimitedfrom upstream as soon as it's subscribed. If a value arrives while downstream demand is still zero, OpenCombine drops it. Apple's Combine holds it: the latest value whenlatest: true, the first whenlatest: false. Apple emits the held value once downstream requests demand.This is visible whenever a subscriber's demand arrives after the upstream's first value. That covers a snapshot replayed on subscribe, a
prepend, or a quiet stream whose only value lands while the subscriber is still wiring up. It also happens mid-stream after downstream demand is used up.Side-by-side on macOS (value(s) sent before demand, then
request(.max(1))):send(1,2,3),latest: true[3][][3]send(1,2,3),latest: false[1][][1]send(2,3), request[1, 3][1][1, 3]Fix
A value received with zero demand is stored in
heldInput, following thelatestrule.request(_:)schedules it through the normal emission path, which still honors the throttle interval, when no emission is already pending.Tests
testInputWithoutDemandIsHeldUntilDemandcovers bothlatestmodes, at subscribe time and after demand is used up. It fails onmainand passes with this change.swift test -Xswiftc -DOPENCOMBINE_COMPATIBILITY_TEST), so the expected values are Apple's behavior.🤖 Generated with Claude Code