Repository navigation
Drain buffered coalesced values before completion - #8247
Conversation
Bugbot is paused — on-demand spend limit reachedBugbot uses usage-based billing for this team and has hit its on-demand spend limit. A team admin can raise the spend limit in the Cursor dashboard, or wait for the next billing cycle to continue. |
📝 WalkthroughWalkthrough
ChangesDemand-aware coalescing
Estimated code review effort: 3 (Moderate) | ~25 minutes Sequence Diagram(s)sequenceDiagram
participant Upstream
participant CoalesceLatestInner
participant DemandControlledSubscriber
Upstream->>CoalesceLatestInner: send latest value
CoalesceLatestInner->>CoalesceLatestInner: buffer readyValue
DemandControlledSubscriber->>CoalesceLatestInner: request demand
CoalesceLatestInner->>DemandControlledSubscriber: deliver buffered value
Upstream->>CoalesceLatestInner: complete
CoalesceLatestInner->>DemandControlledSubscriber: forward completion after draining
Possibly related PRs
Suggested reviewers: Important Pre-merge checks failedPlease resolve all errors before merging. Addressing warnings is optional. ❌ Failed checks (1 error, 1 warning)
✅ Passed checks (23 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: b271e51381
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
|
@codex review |
98e76f4 to
acc82f4
Compare
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@cmuxTests/WorkspaceSidebarObservationTests.swift`:
- Around line 257-281: Add a regression test alongside
coalesceLatestDrainsReentrantValueBeforeCompletionWithUnlimitedDemand that uses
finite demand, sends a reentrant value followed by completion when demand is
exhausted, verifies the value and completion are buffered, then requests
additional demand and asserts the value drains before exactly one completion.
In `@Sources/CoalesceLatestPublisher.swift`:
- Around line 187-202: Update finishPendingCompletionIfPossible() so it returns
while readyValue is non-nil, retaining both the buffered value and
pendingCompletion until demand allows the value to be delivered. Remove the
logic that clears readyValue and completes immediately with zero demand; forward
completion only after the buffered value has been drained, while preserving the
existing cancellation and downstream delivery behavior.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 8cb9dfff-5311-49ea-a294-5a3da9e323d1
📒 Files selected for processing (2)
Sources/CoalesceLatestPublisher.swiftcmuxTests/WorkspaceSidebarObservationTests.swift
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: acc82f4050
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| if let value = pendingValue { | ||
| pendingValue = nil | ||
| _ = downstream.receive(value) | ||
| enqueueForDelivery(value) |
There was a problem hiding this comment.
Preserve the reentrant leading value at completion
With unlimited demand, if a downstream receiving the replay reentrantly sends a leading value, then another value inside that new window, then completes (send(2); send(3); send(completion:)), this line enqueues the pending trailing value through readyValue = value while isDelivering is still true, overwriting the held leading value before the outer drain can emit it. Fresh evidence beyond the prior fixed thread is the extra in-window value before completion; the old sink behavior would emit [1, 2, 3], but this change emits [1, 3], so leading-edge semantics are still broken for completing reentrant streams.
Useful? React with 👍 / 👎.
Greptile SummaryFixes a deterministic crash introduced by a prior commit where
Confidence Score: 5/5Safe to merge; the fix is narrowly scoped to the coalescing operator's demand accounting and completion ordering, and both new regression tests pass against the corrected implementation. The drain loop, reentrancy guard, and deferred-completion logic all trace correctly for unlimited-demand and demand-limited subscribers, including reentrant upstream emissions mid-delivery. No blocking primitives, actor isolation issues, test seams in production source, or other rule violations were found. No files require special attention. Important Files Changed
Reviews (2): Last reviewed commit: "fix(sidebar): drain reentrant values bef..." | Re-trigger Greptile |
| private func finishPendingCompletionIfPossible() { | ||
| guard !isCancelled, | ||
| !isDelivering, | ||
| let completion = pendingCompletion | ||
| else { return } | ||
|
|
||
| pendingCompletion = nil | ||
| // Combine completion is not demand-gated. Values that reentered an | ||
| // active delivery drain first while demand remains; if demand is zero, | ||
| // the buffered latest value cannot legally precede completion and drops. | ||
| readyValue = nil | ||
| downstreamDemand = .none | ||
| isCancelled = true | ||
| downstream.receive(completion: completion) | ||
| upstreamSubscription = nil | ||
| } |
There was a problem hiding this comment.
Silent value drop when demand is zero at completion
finishPendingCompletionIfPossible unconditionally clears readyValue before forwarding completion. In the unlimited-demand path this is always nil by the time the drain exits, so there is no loss. For a demand-limited subscriber, however, if the upstream completes while downstreamDemand == .none and a value is sitting in readyValue (e.g. the leading replay was delivered but the subscriber has not yet called request again), that buffered value is silently dropped — the subscriber receives completion without ever seeing the latest coalesced value.
For the Publisher.values / AsyncSequence bridge motivating this fix, next() always requests .max(1) before suspending, so demand is always in-flight and this case cannot arise in practice. Worth noting here for future demand-limited consumers that are slow to re-request: the contract is "latest value only when demand exists", and a completion that races a pending request wins.
acc82f4 to
0bf8747
Compare
Bugbot is paused — on-demand spend limit reachedBugbot uses usage-based billing for this team and has hit its on-demand spend limit. A team admin can raise the spend limit in the Cursor dashboard, or wait for the next billing cycle to continue. |
|
@codex review |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 0bf8747f56
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| } | ||
|
|
||
| private func enqueueForDelivery(_ value: Input) { | ||
| readyValue = value |
There was a problem hiding this comment.
Don't overwrite demanded reentrant emissions
When downstream still has outstanding demand but enqueueForDelivery is re-entered before the current downstream.receive returns, this single slot replaces an earlier ready emission. For example, with .coalesceLatest(for: .zero, ...) (or a callback that stalls past the interval), an unlimited-demand sink that sends 2 and then 3 from its receive(1) only receives [1, 3]; neither value is inside a coalescing window and the old unlimited-demand behavior delivered both. Keep a FIFO of ready emission points, or otherwise avoid conflating values while demand remains.
Useful? React with 👍 / 👎.
Ported from upstream manaflow-ai#8247 (7d62e38): "Drain buffered coalesced values before completion (manaflow-ai#8247)".
Follow-up to #8236.
That PR restored downstream demand accounting, but completion can still overtake a value buffered before completion. This happens when downstream delivery causes a reentrant value and completion, or when the final value arrives while downstream demand is zero.
The operator now retains pending completion until its buffered value drains. It prevents reentrant delivery with one drain loop, forwards the value once demand resumes, then forwards completion exactly once.
Regression structure:
aeb9aba971adds unlimited-demand reentrancy and finite-demand final-value tests.0bf8747f56adds the completion-aware drain loop and makes them pass.Prior tagged build of the full demand and completion implementation:
code2, https://github.com/manaflow-ai/cmux/actions/runs/29474218386No user-facing strings changed, so localization files are unaffected.
Summary by CodeRabbit
coalesceLatestto strictly honor downstream demand before emitting buffered values.