From 3b2476851fcfd26ade78b3336664ccf6b106054e Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Tue, 30 Jun 2026 21:27:52 -0700 Subject: [PATCH 1/8] perf(sidebar): leading-edge throttle on immediate observation publisher Agents (e.g. Codex) rewrite a workspace title every turn. The immediate sidebar observation publisher had removeDuplicates() but no burst coalescing, and removeDuplicates() cannot collapse distinct titles. Every rewrite therefore drove a full makeWorkspaceSnapshot() rebuild (git/PR/directory lookups) in each downstream consumer, for every workspace, every turn. With many open workspaces this is a sustained main-thread CPU spike. The publisher now fans out to two consumers on main: the per-row TabItemView subscription and the MergeMany extension-sidebar aggregate. Placing the throttle in the publisher coalesces both. A 50ms leading-edge throttle (latest: true) keeps the first change in a burst instant (user pin/color/title edits still feel immediate) while collapsing the rest to at most one emission per window. Mirrors the existing 40ms debounce on the slower sidebarObservationPublisher. Verified on a tagged build: 60 unique title changes in 0.5s coalesced to 2 snapshot refreshes; 3 edits spaced 250ms apart stayed at 3 (instant single-edit feedback preserved). Refs https://github.com/manaflow-ai/cmux/issues/4127 Co-Authored-By: Claude Opus 4.8 (1M context) --- Sources/WorkspaceSidebarObservation.swift | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/Sources/WorkspaceSidebarObservation.swift b/Sources/WorkspaceSidebarObservation.swift index 8b2dd957402b..cbe378c9d4ec 100644 --- a/Sources/WorkspaceSidebarObservation.swift +++ b/Sources/WorkspaceSidebarObservation.swift @@ -59,6 +59,19 @@ private struct SidebarObservationState: Equatable { } extension Workspace { + // Leading-edge coalescing for the immediate sidebar observation stream. + // The publisher fans out to every sidebar consumer (per-row rows and the + // MergeMany extension-sidebar aggregate), and each downstream fires a full + // makeWorkspaceSnapshot() rebuild. 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. A leading-edge throttle keeps the first change in a burst + // instant (so a user pin/color/title edit still feels immediate) while + // collapsing the rest to at most one emission per window. Mirrors the 40ms + // debounce already applied to the slower observation publisher. + // See https://github.com/manaflow-ai/cmux/issues/4127. + static let sidebarImmediateObservationCoalesceInterval: RunLoop.SchedulerTimeType.Stride = .milliseconds(50) + func makeSidebarImmediateObservationPublisher() -> AnyPublisher { let workspaceFields = Publishers.CombineLatest4( $title, @@ -86,6 +99,11 @@ extension Workspace { ) } .removeDuplicates() + .throttle( + for: Self.sidebarImmediateObservationCoalesceInterval, + scheduler: RunLoop.main, + latest: true + ) .map { _ in () } .eraseToAnyPublisher() } From 36df8c1ce1a7387fb7d6ed5e759c4bfd092d7a58 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 1 Jul 2026 17:24:32 -0700 Subject: [PATCH 2/8] test(sidebar): assert the immediate publisher's synchronous leading-edge contract The replay a late subscriber receives, and the first change after idle, must arrive in the same run-loop turn; only a burst tail may coalesce. These fail on the current throttle-based head: Combine's throttle schedules every emission onto the scheduler, so nothing is synchronous. Co-Authored-By: Claude Fable 5 --- .../WorkspaceSidebarObservationTests.swift | 46 +++++++++++++++++++ 1 file changed, 46 insertions(+) diff --git a/cmuxTests/WorkspaceSidebarObservationTests.swift b/cmuxTests/WorkspaceSidebarObservationTests.swift index 791bda6b9417..e976ce161d95 100644 --- a/cmuxTests/WorkspaceSidebarObservationTests.swift +++ b/cmuxTests/WorkspaceSidebarObservationTests.swift @@ -118,6 +118,52 @@ struct WorkspaceSidebarObservationTests { ) } + @Test func sidebarImmediateObservationPublisherDeliversFirstChangeSynchronously() { + let workspace = Workspace() + + var publishCount = 0 + let cancellable = workspace.sidebarImmediateObservationPublisher.sink { + publishCount += 1 + } + defer { cancellable.cancel() } + publishCount = 0 + + workspace.title = "User Edit" + + #expect( + publishCount == 1, + "The first immediate-field change after subscribing must reach the sidebar in the same run-loop turn; coalescing may only defer the tail of a burst." + ) + } + + @Test func sidebarImmediateObservationPublisherCoalescesTitleBursts() { + let workspace = Workspace() + + var publishCount = 0 + let cancellable = workspace.sidebarImmediateObservationPublisher.sink { + publishCount += 1 + } + defer { cancellable.cancel() } + publishCount = 0 + + for turn in 0..<20 { + workspace.title = "Agent Turn \(turn)" + } + + #expect( + publishCount == 1, + "A synchronous burst of distinct titles must deliver only its leading edge immediately." + ) + + // Generous pump so the 50ms trailing emission fires deterministically. + RunLoop.main.run(until: Date().addingTimeInterval(0.3)) + + #expect( + publishCount == 2, + "A coalesced burst must settle with exactly one trailing emission carrying the latest state." + ) + } + @Test func sidebarObservationPublisherIgnoresRemoteHeartbeatOnlyChanges() { let workspace = Workspace() From e68322b62a1bd25921a68796628d6442ad45fc4d Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 1 Jul 2026 17:24:50 -0700 Subject: [PATCH 3/8] perf(sidebar): synchronous-leading coalesce for immediate observation, per workspace and across the extension aggregate Replace Combine's throttle with coalesceLatest, a custom operator whose leading edge is synchronous: throttle schedules every emission onto the scheduler, so the @Published replay and the first change after idle were deferred to the next run-loop turn, breaking the immediate-invalidation contract the tests in the previous commit assert. coalesceLatest forwards the replay and any post-idle change in the same run-loop turn and defers only the tail of a burst, emitting the latest value once per 50ms window. Also coalesce across the extension-sidebar MergeMany aggregate: per- workspace coalescing caps each stream, but N workspaces bursting concurrently still re-rendered the whole extension sidebar once per workspace per window. Refs https://github.com/manaflow-ai/cmux/issues/4127 Co-Authored-By: Claude Fable 5 --- Sources/ContentView.swift | 12 ++ Sources/WorkspaceSidebarObservation.swift | 158 ++++++++++++++++-- .../WorkspaceSidebarObservationTests.swift | 32 ++++ 3 files changed, 190 insertions(+), 12 deletions(-) diff --git a/Sources/ContentView.swift b/Sources/ContentView.swift index 4af7a08c9d33..25815e0db9e4 100644 --- a/Sources/ContentView.swift +++ b/Sources/ContentView.swift @@ -10969,10 +10969,22 @@ struct VerticalTabsSidebar: View { return } + // Per-workspace coalescing (coalesceLatest inside each workspace's + // immediate publisher) caps each stream, but this merge fans in every + // workspace: with N workspaces bursting concurrently (many agents + // finishing turns), the aggregate would still re-render the whole + // extension sidebar once per workspace per window. Coalesce again + // across the merge so a cross-workspace burst settles into one + // re-render per window; the leading edge stays synchronous so a lone + // change is as immediate as before. extensionSidebarImmediateObservationPublisher = Publishers.MergeMany( tabs.map { $0.sidebarImmediateObservationPublisher } ) .receive(on: RunLoop.main) + .coalesceLatest( + for: Workspace.sidebarImmediateObservationCoalesceInterval, + scheduler: RunLoop.main + ) .eraseToAnyPublisher() extensionSidebarDebouncedObservationPublisher = Publishers.MergeMany( tabs.map { $0.sidebarObservationPublisher } diff --git a/Sources/WorkspaceSidebarObservation.swift b/Sources/WorkspaceSidebarObservation.swift index cbe378c9d4ec..60c413a93127 100644 --- a/Sources/WorkspaceSidebarObservation.swift +++ b/Sources/WorkspaceSidebarObservation.swift @@ -60,15 +60,15 @@ private struct SidebarObservationState: Equatable { extension Workspace { // Leading-edge coalescing for the immediate sidebar observation stream. - // The publisher fans out to every sidebar consumer (per-row rows and the - // MergeMany extension-sidebar aggregate), and each downstream fires a full - // makeWorkspaceSnapshot() rebuild. 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. A leading-edge throttle keeps the first change in a burst - // instant (so a user pin/color/title edit still feels immediate) while - // collapsing the rest to at most one emission per window. Mirrors the 40ms - // debounce already applied to the slower observation publisher. + // 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) @@ -99,10 +99,9 @@ extension Workspace { ) } .removeDuplicates() - .throttle( + .coalesceLatest( for: Self.sidebarImmediateObservationCoalesceInterval, - scheduler: RunLoop.main, - latest: true + scheduler: RunLoop.main ) .map { _ in () } .eraseToAnyPublisher() @@ -172,3 +171,138 @@ extension Workspace { .eraseToAnyPublisher() } } + +// 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( + for interval: Context.SchedulerTimeType.Stride, + scheduler: Context + ) -> AnyPublisher { + CoalesceLatestPublisher(upstream: self, interval: interval, scheduler: scheduler) + .eraseToAnyPublisher() + } +} + +private struct CoalesceLatestPublisher: 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(subscriber: S) where S.Input == Output, S.Failure == Never { + upstream.subscribe(CoalesceLatestInner( + downstream: subscriber, + interval: interval, + scheduler: scheduler + )) + } +} + +private final class CoalesceLatestInner: 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) { + 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 { + windowStart = now + _ = downstream.receive(input) + } + return .none + } + + func receive(completion: Subscribers.Completion) { + 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 } + 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 + } +} diff --git a/cmuxTests/WorkspaceSidebarObservationTests.swift b/cmuxTests/WorkspaceSidebarObservationTests.swift index e976ce161d95..6929a2e2b43f 100644 --- a/cmuxTests/WorkspaceSidebarObservationTests.swift +++ b/cmuxTests/WorkspaceSidebarObservationTests.swift @@ -164,6 +164,38 @@ struct WorkspaceSidebarObservationTests { ) } + @Test func coalesceLatestKeepsLeadingEdgeSynchronousAndEmitsLatestTrailing() { + let subject = PassthroughSubject() + var received: [Int] = [] + let cancellable = subject + .coalesceLatest(for: .milliseconds(50), scheduler: RunLoop.main) + .sink { received.append($0) } + defer { cancellable.cancel() } + + // First value models the @Published current-state replay: forwarded + // synchronously without opening a coalesce window. + subject.send(1) + #expect(received == [1]) + + // First change is the synchronous leading edge and opens the window. + subject.send(2) + #expect(received == [1, 2]) + + // Burst inside the window coalesces to the latest value. + subject.send(3) + subject.send(4) + subject.send(5) + #expect(received == [1, 2]) + + RunLoop.main.run(until: Date().addingTimeInterval(0.3)) + #expect(received == [1, 2, 5]) + + // After the window closes and the trailing window expires, the next + // value is synchronous again. + subject.send(6) + #expect(received == [1, 2, 5, 6]) + } + @Test func sidebarObservationPublisherIgnoresRemoteHeartbeatOnlyChanges() { let workspace = Workspace() From 00ae6fdd37ce0ffed8b4044ef55a3d55d2a73a76 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 1 Jul 2026 19:37:38 -0700 Subject: [PATCH 4/8] sidebar: move merged immediate aggregate into WorkspaceSidebarObservation Keeps ContentView.swift inside the Swift file length budget and puts the aggregate coalescing next to the operator and interval it uses. Co-Authored-By: Claude Fable 5 --- Sources/ContentView.swift | 19 ++----------------- Sources/WorkspaceSidebarObservation.swift | 16 ++++++++++++++++ 2 files changed, 18 insertions(+), 17 deletions(-) diff --git a/Sources/ContentView.swift b/Sources/ContentView.swift index eb82778702ea..789b5226db5b 100644 --- a/Sources/ContentView.swift +++ b/Sources/ContentView.swift @@ -10969,23 +10969,8 @@ struct VerticalTabsSidebar: View { return } - // Per-workspace coalescing (coalesceLatest inside each workspace's - // immediate publisher) caps each stream, but this merge fans in every - // workspace: with N workspaces bursting concurrently (many agents - // finishing turns), the aggregate would still re-render the whole - // extension sidebar once per workspace per window. Coalesce again - // across the merge so a cross-workspace burst settles into one - // re-render per window; the leading edge stays synchronous so a lone - // change is as immediate as before. - extensionSidebarImmediateObservationPublisher = Publishers.MergeMany( - tabs.map { $0.sidebarImmediateObservationPublisher } - ) - .receive(on: RunLoop.main) - .coalesceLatest( - for: Workspace.sidebarImmediateObservationCoalesceInterval, - scheduler: RunLoop.main - ) - .eraseToAnyPublisher() + extensionSidebarImmediateObservationPublisher = + Workspace.mergedImmediateObservationPublisher(for: tabs) extensionSidebarDebouncedObservationPublisher = Publishers.MergeMany( tabs.map { $0.sidebarObservationPublisher } ) diff --git a/Sources/WorkspaceSidebarObservation.swift b/Sources/WorkspaceSidebarObservation.swift index 6e9d4b209062..51d3dd2f3df2 100644 --- a/Sources/WorkspaceSidebarObservation.swift +++ b/Sources/WorkspaceSidebarObservation.swift @@ -108,6 +108,22 @@ extension Workspace { .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 { + Publishers.MergeMany(workspaces.map { $0.sidebarImmediateObservationPublisher }) + .receive(on: RunLoop.main) + .coalesceLatest( + for: sidebarImmediateObservationCoalesceInterval, + scheduler: RunLoop.main + ) + .eraseToAnyPublisher() + } + func makeSidebarObservationPublisher() -> AnyPublisher { let workspaceFields = Publishers.CombineLatest4( $currentDirectory, From cee98e4def82ebf9db1751caff3eb5358e341482 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 1 Jul 2026 21:31:57 -0700 Subject: [PATCH 5/8] coalesceLatest: drop stale pending value when an overdue leading emission supersedes it Codex autoreview P2: if the trailing callback is delayed past its deadline (main run-loop stall) and a newer value arrives first, the new value took the leading branch but left pendingValue set, so the late callback emitted the stale value out of order. The leading branch now clears the superseded pending value, and an overdue callback firing inside a newer window reschedules to that window's deadline instead of emitting early. Co-Authored-By: Claude Fable 5 --- Sources/WorkspaceSidebarObservation.swift | 15 ++++++++++ .../WorkspaceSidebarObservationTests.swift | 30 +++++++++++++++++++ 2 files changed, 45 insertions(+) diff --git a/Sources/WorkspaceSidebarObservation.swift b/Sources/WorkspaceSidebarObservation.swift index 51d3dd2f3df2..2cfc30211c0a 100644 --- a/Sources/WorkspaceSidebarObservation.swift +++ b/Sources/WorkspaceSidebarObservation.swift @@ -281,6 +281,11 @@ private final class CoalesceLatestInner() + var received: [Int] = [] + let cancellable = subject + .coalesceLatest(for: .milliseconds(50), scheduler: RunLoop.main) + .sink { received.append($0) } + defer { cancellable.cancel() } + + subject.send(1) // replay: forwarded, no window + subject.send(2) // leading edge: opens window + subject.send(3) // pending trailing value for the open window + #expect(received == [1, 2]) + + // Stall the main run loop past the trailing deadline WITHOUT pumping + // it, so the scheduled callback is overdue when the next value lands. + Thread.sleep(forTimeInterval: 0.12) + subject.send(4) // deadline passed: new leading edge must supersede 3 + + #expect( + received == [1, 2, 4], + "A newer leading value after an overdue deadline must drop the stale pending value." + ) + + RunLoop.main.run(until: Date().addingTimeInterval(0.3)) + #expect( + received == [1, 2, 4], + "The overdue trailing callback must not emit the superseded stale value out of order." + ) + } + @Test func sidebarObservationPublisherIgnoresRemoteHeartbeatOnlyChanges() { let workspace = Workspace() From 0046754620ae61f5ed0613cda1522a682ada1587 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 1 Jul 2026 22:29:32 -0700 Subject: [PATCH 6/8] Move coalesceLatest operator into Sources/CoalesceLatestPublisher.swift cmux policy: major added types get their own TypeName.swift file. Wires the new file into the pbxproj alongside WorkspaceSidebarObservation.swift. Co-Authored-By: Claude Fable 5 --- Sources/CoalesceLatestPublisher.swift | 152 ++++++++++++++++++++++ Sources/WorkspaceSidebarObservation.swift | 150 --------------------- cmux.xcodeproj/project.pbxproj | 4 + 3 files changed, 156 insertions(+), 150 deletions(-) create mode 100644 Sources/CoalesceLatestPublisher.swift diff --git a/Sources/CoalesceLatestPublisher.swift b/Sources/CoalesceLatestPublisher.swift new file mode 100644 index 000000000000..c197cfa3e934 --- /dev/null +++ b/Sources/CoalesceLatestPublisher.swift @@ -0,0 +1,152 @@ +import Combine +import Foundation + +// 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( + for interval: Context.SchedulerTimeType.Stride, + scheduler: Context + ) -> AnyPublisher { + CoalesceLatestPublisher(upstream: self, interval: interval, scheduler: scheduler) + .eraseToAnyPublisher() + } +} + +private struct CoalesceLatestPublisher: 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(subscriber: S) where S.Input == Output, S.Failure == Never { + upstream.subscribe(CoalesceLatestInner( + downstream: subscriber, + interval: interval, + scheduler: scheduler + )) + } +} + +private final class CoalesceLatestInner: 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) { + 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) { + 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 + } +} diff --git a/Sources/WorkspaceSidebarObservation.swift b/Sources/WorkspaceSidebarObservation.swift index 2cfc30211c0a..58fe421d6ac6 100644 --- a/Sources/WorkspaceSidebarObservation.swift +++ b/Sources/WorkspaceSidebarObservation.swift @@ -189,153 +189,3 @@ extension Workspace { .eraseToAnyPublisher() } } - -// 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( - for interval: Context.SchedulerTimeType.Stride, - scheduler: Context - ) -> AnyPublisher { - CoalesceLatestPublisher(upstream: self, interval: interval, scheduler: scheduler) - .eraseToAnyPublisher() - } -} - -private struct CoalesceLatestPublisher: 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(subscriber: S) where S.Input == Output, S.Failure == Never { - upstream.subscribe(CoalesceLatestInner( - downstream: subscriber, - interval: interval, - scheduler: scheduler - )) - } -} - -private final class CoalesceLatestInner: 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) { - 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) { - 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 - } -} diff --git a/cmux.xcodeproj/project.pbxproj b/cmux.xcodeproj/project.pbxproj index dd9f6338c557..63a05d2d17fd 100644 --- a/cmux.xcodeproj/project.pbxproj +++ b/cmux.xcodeproj/project.pbxproj @@ -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 */; }; @@ -1651,6 +1652,7 @@ C0DE58990000000000000002 /* CmuxWebViewKeyDownReentryTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxWebViewKeyDownReentryTests.swift; sourceTree = ""; }; 43F90FAF3FD44F11BF547BE9 /* CmuxWebViewMouseNavigationButtonTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxWebViewMouseNavigationButtonTests.swift; sourceTree = ""; }; E30750000000000000000003 /* CmuxWorkspaceDefinition.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxWorkspaceDefinition.swift; sourceTree = ""; }; + C0A1E5CE01C0A1E5CE01A002 /* CoalesceLatestPublisher.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CoalesceLatestPublisher.swift; sourceTree = ""; }; A9F100000000000000000016 /* CodexAppServerQueuedInput.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Panels/CodexAppServerQueuedInput.swift; sourceTree = ""; }; A9E01000000000000000000E /* CodexAppServerSession.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Panels/CodexAppServerSession.swift; sourceTree = ""; }; A9E040000000000000000001 /* CodexAppServerSessionTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CodexAppServerSessionTests.swift; sourceTree = ""; }; @@ -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 */, @@ -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 */, From 8dac5f4abaca5afa0e30e6242279f8441eb00099 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 1 Jul 2026 22:43:34 -0700 Subject: [PATCH 7/8] Document why CoalesceLatestPublisher and its Inner share one file Co-Authored-By: Claude Fable 5 --- Sources/CoalesceLatestPublisher.swift | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/Sources/CoalesceLatestPublisher.swift b/Sources/CoalesceLatestPublisher.swift index c197cfa3e934..5299ef4fb719 100644 --- a/Sources/CoalesceLatestPublisher.swift +++ b/Sources/CoalesceLatestPublisher.swift @@ -1,6 +1,12 @@ 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 { From 9de6d4a334fb30cc4065dd195990eea5caf0f11f Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Thu, 2 Jul 2026 00:41:23 -0700 Subject: [PATCH 8/8] Make the overdue-trailing regression test deterministic with a virtual scheduler The test determinism gate rejects sleep-then-assert. coalesceLatest is generic over Scheduler, so drive the stall interleaving exactly: advance now past the deadline without running the scheduled callback, assert the supersede, then run the overdue callback and assert no stale emission. Co-Authored-By: Claude Fable 5 --- .../WorkspaceSidebarObservationTests.swift | 60 +++++++++++++++++-- 1 file changed, 55 insertions(+), 5 deletions(-) diff --git a/cmuxTests/WorkspaceSidebarObservationTests.swift b/cmuxTests/WorkspaceSidebarObservationTests.swift index aad130876bb5..fa140cae821d 100644 --- a/cmuxTests/WorkspaceSidebarObservationTests.swift +++ b/cmuxTests/WorkspaceSidebarObservationTests.swift @@ -197,10 +197,11 @@ struct WorkspaceSidebarObservationTests { } @Test func coalesceLatestDropsStalePendingValueWhenLeadingSupersedesOverdueTrailing() { + let scheduler = VirtualCoalesceScheduler() let subject = PassthroughSubject() var received: [Int] = [] let cancellable = subject - .coalesceLatest(for: .milliseconds(50), scheduler: RunLoop.main) + .coalesceLatest(for: .milliseconds(50), scheduler: scheduler) .sink { received.append($0) } defer { cancellable.cancel() } @@ -208,10 +209,11 @@ struct WorkspaceSidebarObservationTests { subject.send(2) // leading edge: opens window subject.send(3) // pending trailing value for the open window #expect(received == [1, 2]) + #expect(scheduler.scheduledActionCount == 1) - // Stall the main run loop past the trailing deadline WITHOUT pumping - // it, so the scheduled callback is overdue when the next value lands. - Thread.sleep(forTimeInterval: 0.12) + // The deadline passes WITHOUT the scheduled callback running, + // modeling a stalled main run loop with an overdue timer. + scheduler.advance(by: 0.12) subject.send(4) // deadline passed: new leading edge must supersede 3 #expect( @@ -219,7 +221,7 @@ struct WorkspaceSidebarObservationTests { "A newer leading value after an overdue deadline must drop the stale pending value." ) - RunLoop.main.run(until: Date().addingTimeInterval(0.3)) + scheduler.runScheduledActions() #expect( received == [1, 2, 4], "The overdue trailing callback must not emit the superseded stale value out of order." @@ -254,3 +256,51 @@ private final class ObservationChangeFlag: @unchecked Sendable { fired = true } } + +// Deterministic Combine scheduler for coalesceLatest tests: `now` only moves +// via advance(by:), and scheduled actions run only when runScheduledActions() +// is called, so overdue-timer interleavings are exact instead of wall-clock. +private final class VirtualCoalesceScheduler: Scheduler { + typealias SchedulerTimeType = RunLoop.SchedulerTimeType + typealias SchedulerOptions = Never + + private(set) var now = SchedulerTimeType(Date(timeIntervalSinceReferenceDate: 0)) + var minimumTolerance: SchedulerTimeType.Stride { .seconds(0) } + private var scheduledActions: [() -> Void] = [] + + var scheduledActionCount: Int { scheduledActions.count } + + func advance(by seconds: TimeInterval) { + now = SchedulerTimeType(now.date.addingTimeInterval(seconds)) + } + + func runScheduledActions() { + let actions = scheduledActions + scheduledActions = [] + actions.forEach { $0() } + } + + func schedule(options: Never?, _ action: @escaping () -> Void) { + action() + } + + func schedule( + after date: SchedulerTimeType, + tolerance: SchedulerTimeType.Stride, + options: Never?, + _ action: @escaping () -> Void + ) { + scheduledActions.append(action) + } + + func schedule( + after date: SchedulerTimeType, + interval: SchedulerTimeType.Stride, + tolerance: SchedulerTimeType.Stride, + options: Never?, + _ action: @escaping () -> Void + ) -> Cancellable { + scheduledActions.append(action) + return AnyCancellable {} + } +}