Skip to content

[BUG] Since #59, concurrent sends reach consumers in different orders and current-value/replay consumers can end on a stale value #61

Description

@PaulTaykalo

Hi @twittemb, a follow-up to #52 and #59.

Problem

Since #59, send delivers outside the lock, so two send calls made at the same time can reach consumers in different orders:

  • Passthrough: consumer A gets 1, 2 while consumer B gets 2, 1.
  • CurrentValue and Replay: a consumer can end on 1 while value is 2, and it stays stale until the next send.

Repro: two consumers, and two threads sending at the same time with DispatchQueue.concurrentPerform. On main, the consumers' arrays differ in almost every run. Before #59, they never did.

Suggested fix

Queue each delivery while holding the state lock, then deliver the queue in order without holding any lock:

final class OrderedDelivery: @unchecked Sendable {
  private let state = ManagedCriticalState((pending: Deque<() -> Void>(), isDelivering: false))

  /// Call while holding the subject's state lock. Returns true if the caller must `drain()` afterwards.
  func enqueue(_ delivery: @escaping () -> Void) -> Bool {
    self.state.withCriticalRegion { state in
      state.pending.append(delivery)
      defer { state.isDelivering = true }
      return !state.isDelivering
    }
  }

  /// Runs queued deliveries in order. Call without holding any lock.
  func drain() {
    while let delivery = self.state.withCriticalRegion({ state -> (() -> Void)? in
      if state.pending.isEmpty { state.isDelivering = false }
      return state.pending.popFirst()
    }) {
      delivery()
    }
  }
}

// In each subject's send(_:) and send(_ termination:):
let shouldDrain = self.state.withCriticalRegion { state -> Bool in
  // existing state updates stay here
  let channels = Array(state.channels.values)
  return self.delivery.enqueue { for channel in channels { channel.send(element) } }
}
if shouldDrain { self.delivery.drain() }

Holding a lock during delivery instead doesn't work. If a consumer's cancellation handler calls send, it deadlocks the same way as #52.

One trade-off: when several threads send at once, send can return before its value is delivered.

With this change, #59's AsyncSubjectCancellationTests and the full suite pass (152 tests). It also passes new tests for ordering and for a cancellation handler that calls send. I can open a PR with the fix and the tests.

No activity

Activity on this issue will appear here.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions