Skip to content

Throttle: hold input received without downstream demand - #25

Closed
piercifani wants to merge 1 commit into
OpenSwiftUIProject:mainfrom
theleftbit:fix/throttle-hold-input-without-demand
Closed

piercifani wants to merge 1 commit into
OpenSwiftUIProject:mainfrom
theleftbit:fix/throttle-hold-input-without-demand

Conversation

@piercifani

Copy link
Copy Markdown

Problem

Publishers.Throttle requests .unlimited from 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 when latest: true, the first when latest: 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))):

Scenario Apple Combine OpenCombine 0.15.1 This PR
send(1,2,3), latest: true [3] [] [3]
send(1,2,3), latest: false [1] [] [1]
demand used up, then send(2,3), request [1, 3] [1] [1, 3]

Fix

A value received with zero demand is stored in heldInput, following the latest rule. request(_:) schedules it through the normal emission path, which still honors the throttle interval, when no emission is already pending.

Tests

  • New testInputWithoutDemandIsHeldUntilDemand covers both latest modes, at subscribe time and after demand is used up. It fails on main and passes with this change.
  • It also passes against Apple's Combine (swift test -Xswiftc -DOPENCOMBINE_COMPATIBILITY_TEST), so the expected values are Apple's behavior.
  • Full suite passes: 1567 tests, 0 failures.

🤖 Generated with Claude Code

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>
@Kyle-Ye
Kyle-Ye self-requested a review September 25, 2026 13:46
@Kyle-Ye

Kyle-Ye commented Sep 25, 2026

Copy link
Copy Markdown
Member

I'll review this after I fixed the recent CI issue on this repo.

@codecov

codecov Bot commented Sep 27, 2026

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 96.07%. Comparing base (d81fffe) to head (9c22a83).
⚠️ Report is 5 commits behind head on main.

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.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@Kyle-Ye

Kyle-Ye commented Sep 27, 2026 •

Copy link
Copy Markdown
Member

Your PR has the following issue and does not align with Combine:

The separate heldInput buffer can keep an old value after newer input has been delivered. A later demand request can then emit that old value, causing output such as [1, 3, 2].

The examples below use a one-second interval and the repository's VirtualTimeScheduler. The subscriber returns .none from receive(_:). Running scheduled actions drains the scheduler, including actions scheduled for a later time.

1. A held value can arrive after a newer value

Use latest: false and initial demand .max(1):

Steps Apple Combine This PR
Send 1, 2; request .max(1) before running scheduled actions; then run them [1] [1]
Send 3; run scheduled actions [1, 3] [1, 3]
Request .max(1); run scheduled actions [1, 3] [1, 3, 2]

Save the following as Tests/OpenCombineTests/PublisherTests/ThrottleDemandReproTests.swift in this repository. It uses the existing VirtualTimeScheduler helper.

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 ThrottleDemandReproTests

The last assertion fails, and the output is Received: [1, 3, 2].

To compare the same test with Apple Combine:

OPENCOMBINE_COMPATIBILITY_TEST=1 swift test \
  --scratch-path .build-apple-throttle-repro \
  --filter ThrottleDemandReproTests

This passes and prints Received: [1, 3].

The first input consumes demand and becomes pendingInput. Input 2 becomes heldInput. The next request adds demand but leaves heldInput in place because output is already pending. Input 3 can then use that demand, while 2 remains available for a later request.

2. One pending output becomes two outputs

With initial demand .max(1), send 1, 2, 3 before running any scheduled actions:

Mode After running scheduled actions: Apple / this PR After another .max(1) request and running scheduled actions: Apple / this PR
latest: true [3] / [1] [3] / [1, 3]
latest: false [1] / [1] [1] / [1, 2]

The wrong first value with latest: true already occurs on main. The extra output after the second request is introduced by retaining a separate heldInput.

3. The held input does not follow the throttle window

With latest: false and no initial demand:

  • Send 1 and request .max(1) at t = 0. Apple Combine emits 1 at t = 1s; this PR emits it at t = 0.
  • In a separate run, send 1 at t = 0, 2 at t = 0.2s, and 3 at t = 1.2s. Request .max(1) at t = 1.2s. Apple Combine emits 3 at t = 2.2s; this PR emits 1 at t = 1.2s.

@Kyle-Ye

Kyle-Ye commented Sep 27, 2026

Copy link
Copy Markdown
Member

Close in favor of #28

@Kyle-Ye Kyle-Ye closed this Sep 27, 2026
@Kyle-Ye

Kyle-Ye commented Sep 27, 2026

Copy link
Copy Markdown
Member

If you spot any diff result compared with Combine, an new issue is always welcomed. eg. #29

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants