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
158 changes: 158 additions & 0 deletions Sources/CoalesceLatestPublisher.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,158 @@
import Combine
import Foundation

// CoalesceLatestPublisher and CoalesceLatestInner are one operator: the
// Inner is the publisher's subscription and shares its file-private access
// boundary. Splitting the Inner into its own file would force widening
// `private` to `internal` for an implementation detail, so the pair
// intentionally lives in this single file.

// MARK: - Leading-edge coalescing

extension Publisher where Failure == Never {
/// Coalesces bursts while keeping the leading edge synchronous.
///
/// Combine's `throttle` schedules every emission, including the first,
/// onto the scheduler, so even an isolated value is deferred to the next
/// run-loop turn and a subscriber never observes a synchronous emission.
/// The sidebar's immediate observation contract requires the opposite:
/// the current-state replay a subscriber receives from `@Published`
/// upstreams, and the first change after an idle period, must both arrive
/// in the same run-loop turn; only the tail of a burst may be deferred.
///
/// Semantics per subscription:
/// - The first value (the `@Published` replay of current state) is
/// forwarded synchronously and does not open a coalesce window, so a
/// change made right after subscribing is still synchronous.
/// - A value arriving when no window is open is forwarded synchronously
/// and opens a window of `interval`.
/// - Values arriving inside an open window are coalesced: the latest one
/// is emitted when the window closes (on `scheduler`), which opens the
/// next window.
///
/// Not thread-safe: intended for main-thread streams with `RunLoop.main`.
/// Downstream demand is ignored (sink-style subscribers only).
func coalesceLatest<Context: Scheduler>(
for interval: Context.SchedulerTimeType.Stride,
scheduler: Context
) -> AnyPublisher<Output, Never> {
CoalesceLatestPublisher(upstream: self, interval: interval, scheduler: scheduler)
.eraseToAnyPublisher()
}
}

private struct CoalesceLatestPublisher<Upstream: Publisher, Context: Scheduler>: Publisher
where Upstream.Failure == Never {
typealias Output = Upstream.Output
typealias Failure = Never

let upstream: Upstream
let interval: Context.SchedulerTimeType.Stride
let scheduler: Context

func receive<S: Subscriber>(subscriber: S) where S.Input == Output, S.Failure == Never {
upstream.subscribe(CoalesceLatestInner(
downstream: subscriber,
interval: interval,
scheduler: scheduler
))
}
}

private final class CoalesceLatestInner<Downstream: Subscriber, Context: Scheduler>: Subscriber, Subscription
where Downstream.Failure == Never {
typealias Input = Downstream.Input
typealias Failure = Never

private let downstream: Downstream
private let interval: Context.SchedulerTimeType.Stride
private let scheduler: Context
private var upstreamSubscription: Subscription?
private var hasReceivedReplay = false
private var windowStart: Context.SchedulerTimeType?
private var pendingValue: Input?
private var trailingScheduled = false
private var isCancelled = false

init(downstream: Downstream, interval: Context.SchedulerTimeType.Stride, scheduler: Context) {
self.downstream = downstream
self.interval = interval
self.scheduler = scheduler
}

func receive(subscription: Subscription) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: receive(subscription:) unconditionally requests .unlimited from upstream after calling downstream.receive(subscription:). If the downstream synchronously cancels during that call, the upstream subscription is cancelled and isCancelled is set, yet the method still requests demand on the now-cancelled subscription. Adding a guard !isCancelled else { return } before subscription.request(.unlimited) respects the cancellation and avoids driving unnecessary upstream work.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At Sources/CoalesceLatestPublisher.swift, line 83:

<comment>`receive(subscription:)` unconditionally requests `.unlimited` from upstream after calling `downstream.receive(subscription:)`. If the downstream synchronously cancels during that call, the upstream subscription is cancelled and `isCancelled` is set, yet the method still requests demand on the now-cancelled subscription. Adding a `guard !isCancelled else { return }` before `subscription.request(.unlimited)` respects the cancellation and avoids driving unnecessary upstream work.</comment>

<file context>
@@ -0,0 +1,158 @@
+        self.scheduler = scheduler
+    }
+
+    func receive(subscription: Subscription) {
+        upstreamSubscription = subscription
+        downstream.receive(subscription: self)
</file context>

upstreamSubscription = subscription
downstream.receive(subscription: self)
subscription.request(.unlimited)
}

func receive(_ input: Input) -> Subscribers.Demand {
guard !isCancelled else { return .none }
if !hasReceivedReplay {
hasReceivedReplay = true
_ = downstream.receive(input)
return .none
}
let now = scheduler.now
if let start = windowStart, now < start.advanced(by: interval) {
pendingValue = input
scheduleTrailingEmission(at: start.advanced(by: interval))
} else {
// If a trailing emission was scheduled but its callback is
// overdue (main run loop stalled past the deadline), this newer
// value supersedes the stale pending one; drop it so the late
// callback cannot emit it out of order after this value.
pendingValue = nil
windowStart = now
_ = downstream.receive(input)
}
return .none
}

func receive(completion: Subscribers.Completion<Never>) {
guard !isCancelled else { return }
if let value = pendingValue {
pendingValue = nil
_ = downstream.receive(value)
}
downstream.receive(completion: completion)
}

private func scheduleTrailingEmission(at deadline: Context.SchedulerTimeType) {
guard !trailingScheduled else { return }
trailingScheduled = true
scheduler.schedule(after: deadline) { [weak self] in
self?.emitTrailing()
}
}

private func emitTrailing() {
trailingScheduled = false
guard !isCancelled, let value = pendingValue else { return }
// An overdue callback may fire inside a window that a newer leading
// value opened; hold the pending value until that window's own
// deadline instead of emitting early.
if let start = windowStart {
let deadline = start.advanced(by: interval)
if scheduler.now < deadline {
scheduleTrailingEmission(at: deadline)
return
}
}
pendingValue = nil
windowStart = scheduler.now
_ = downstream.receive(value)
}

func request(_ demand: Subscribers.Demand) {
// Downstream demand is intentionally ignored; this operator backs
// sink-style Void observation streams with unlimited demand.
}

func cancel() {
isCancelled = true
pendingValue = nil
upstreamSubscription?.cancel()
upstreamSubscription = nil
}
}
7 changes: 2 additions & 5 deletions Sources/ContentView.swift
Original file line number Diff line number Diff line change
Expand Up @@ -10969,11 +10969,8 @@ struct VerticalTabsSidebar: View {
return
}

extensionSidebarImmediateObservationPublisher = Publishers.MergeMany(
tabs.map { $0.sidebarImmediateObservationPublisher }
)
.receive(on: RunLoop.main)
.eraseToAnyPublisher()
extensionSidebarImmediateObservationPublisher =
Workspace.mergedImmediateObservationPublisher(for: tabs)
extensionSidebarDebouncedObservationPublisher = Publishers.MergeMany(
tabs.map { $0.sidebarObservationPublisher }
)
Expand Down
33 changes: 33 additions & 0 deletions Sources/WorkspaceSidebarObservation.swift
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,19 @@ private struct SidebarObservationState: Equatable {
}

extension Workspace {
// Leading-edge coalescing for the immediate sidebar observation stream.
// Every subscription (a sidebar row, the MergeMany extension-sidebar
// aggregate) fires a full makeWorkspaceSnapshot() rebuild per emission.
// Agents (e.g. Codex) rewrite a workspace title every turn, and
// removeDuplicates() cannot collapse distinct titles, so without coalescing
// each rewrite drives a snapshot rebuild per consumer per workspace.
// coalesceLatest (below) keeps the first change in a burst synchronous
// (a user pin/color/title edit stays immediate, which Combine's throttle
// cannot guarantee because it schedules every emission onto the scheduler)
// and collapses the tail of the burst into one trailing emission per window.
// See https://github.com/manaflow-ai/cmux/issues/4127.
static let sidebarImmediateObservationCoalesceInterval: RunLoop.SchedulerTimeType.Stride = .milliseconds(50)

func makeSidebarImmediateObservationPublisher() -> AnyPublisher<Void, Never> {
let workspaceFields = Publishers.CombineLatest4(
$title,
Expand Down Expand Up @@ -87,10 +100,30 @@ extension Workspace {
)
}
.removeDuplicates()
.coalesceLatest(
for: Self.sidebarImmediateObservationCoalesceInterval,
scheduler: RunLoop.main
)
.map { _ in () }
.eraseToAnyPublisher()
}

/// Merged immediate observation across workspaces for the extension
/// sidebar. Coalesced again across the merge: per-workspace coalescing
/// caps each stream, but N workspaces bursting concurrently would still
/// re-render the whole extension sidebar once per workspace per window.
/// The leading edge stays synchronous, so a lone change is as immediate
/// as before.
static func mergedImmediateObservationPublisher(for workspaces: [Workspace]) -> AnyPublisher<Void, Never> {
Publishers.MergeMany(workspaces.map { $0.sidebarImmediateObservationPublisher })
.receive(on: RunLoop.main)
.coalesceLatest(
for: sidebarImmediateObservationCoalesceInterval,
scheduler: RunLoop.main
)
.eraseToAnyPublisher()
}

func makeSidebarObservationPublisher() -> AnyPublisher<Void, Never> {
let workspaceFields = Publishers.CombineLatest4(
$currentDirectory,
Expand Down
4 changes: 4 additions & 0 deletions cmux.xcodeproj/project.pbxproj
Original file line number Diff line number Diff line change
Expand Up @@ -407,6 +407,7 @@
5E55200000000000000000B3 /* CmuxWindowing in Frameworks */ = {isa = PBXBuildFile; productRef = 5E55200000000000000000B2 /* CmuxWindowing */; };
E30750000000000000000004 /* CmuxWorkspaceDefinition.swift in Sources */ = {isa = PBXBuildFile; fileRef = E30750000000000000000003 /* CmuxWorkspaceDefinition.swift */; };
E3B7A3000000000000000103 /* CmuxWorkspaces in Frameworks */ = {isa = PBXBuildFile; productRef = E3B7A3000000000000000102 /* CmuxWorkspaces */; };
C0A1E5CE01C0A1E5CE01A001 /* CoalesceLatestPublisher.swift in Sources */ = {isa = PBXBuildFile; fileRef = C0A1E5CE01C0A1E5CE01A002 /* CoalesceLatestPublisher.swift */; };
A9F200000000000000000016 /* CodexAppServerQueuedInput.swift in Sources */ = {isa = PBXBuildFile; fileRef = A9F100000000000000000016 /* CodexAppServerQueuedInput.swift */; };
A9E02000000000000000000E /* CodexAppServerSession.swift in Sources */ = {isa = PBXBuildFile; fileRef = A9E01000000000000000000E /* CodexAppServerSession.swift */; };
A9E040000000000000000002 /* CodexAppServerSessionTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = A9E040000000000000000001 /* CodexAppServerSessionTests.swift */; };
Expand Down Expand Up @@ -1651,6 +1652,7 @@
C0DE58990000000000000002 /* CmuxWebViewKeyDownReentryTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxWebViewKeyDownReentryTests.swift; sourceTree = "<group>"; };
43F90FAF3FD44F11BF547BE9 /* CmuxWebViewMouseNavigationButtonTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxWebViewMouseNavigationButtonTests.swift; sourceTree = "<group>"; };
E30750000000000000000003 /* CmuxWorkspaceDefinition.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxWorkspaceDefinition.swift; sourceTree = "<group>"; };
C0A1E5CE01C0A1E5CE01A002 /* CoalesceLatestPublisher.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CoalesceLatestPublisher.swift; sourceTree = "<group>"; };
A9F100000000000000000016 /* CodexAppServerQueuedInput.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Panels/CodexAppServerQueuedInput.swift; sourceTree = "<group>"; };
A9E01000000000000000000E /* CodexAppServerSession.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Panels/CodexAppServerSession.swift; sourceTree = "<group>"; };
A9E040000000000000000001 /* CodexAppServerSessionTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CodexAppServerSessionTests.swift; sourceTree = "<group>"; };
Expand Down Expand Up @@ -2919,6 +2921,7 @@
A5001520 /* PostHogAnalytics.swift */,
A5001416 /* Workspace.swift */,
6799A0026799A0026799A002 /* WorkspaceSidebarAgentRuntimeObservationModel.swift */,
C0A1E5CE01C0A1E5CE01A002 /* CoalesceLatestPublisher.swift */,
5659A0025659A0025659A002 /* WorkspaceSidebarObservation.swift */,
5659B0025659B0025659B002 /* WorkspaceSidebarLogEntryLimitProvider.swift */,
D7AB3605C10DEF0000000004 /* WorkspaceCloseTabsBatching.swift */,
Expand Down Expand Up @@ -4438,6 +4441,7 @@
C67540030000000000000001 /* CmuxWebView+SubframeDownloadIntentScript.swift in Sources */,
A5001500 /* CmuxWebView.swift in Sources */,
E30750000000000000000004 /* CmuxWorkspaceDefinition.swift in Sources */,
C0A1E5CE01C0A1E5CE01A001 /* CoalesceLatestPublisher.swift in Sources */,
A9F200000000000000000016 /* CodexAppServerQueuedInput.swift in Sources */,
A9E02000000000000000000E /* CodexAppServerSession.swift in Sources */,
C4041001000000000000001B /* CommandClickFileOpenRouter.swift in Sources */,
Expand Down
Loading
Loading