diff --git a/Sources/CmuxEventBus.swift b/Sources/CmuxEventBus.swift index a25b9c70fcc1..364ce022b6d0 100644 --- a/Sources/CmuxEventBus.swift +++ b/Sources/CmuxEventBus.swift @@ -1,3 +1,4 @@ +import Darwin import Foundation struct CmuxEventSubscriptionSnapshot { @@ -6,6 +7,90 @@ struct CmuxEventSubscriptionSnapshot { let ack: [String: Any] } +/// Allocates durable sequences across cmux processes that share the event log. +/// The sidecar is intentionally separate from the JSONL stream so recovery can still +/// preserve a high-water mark when one of the bounded log segments is unreadable. +private final class CmuxEventSequenceStore: @unchecked Sendable { + private let stateURL: URL + private let lockURL: URL + + init(eventLogURL: URL) { + self.stateURL = eventLogURL.appendingPathExtension("seq") + self.lockURL = eventLogURL.appendingPathExtension("seq.lock") + } + + func current() -> Int64 { + withLock { readState() } ?? 0 + } + + func raiseHighWater(to minimum: Int64) { + guard minimum > 0 else { return } + _ = withLock { + let current = readState() + if minimum > current { + try? writeState(minimum) + } + } + } + + #if DEBUG + func resetForTesting() { + try? FileManager.default.removeItem(at: stateURL) + } + #endif + + func allocate(minimum: Int64) -> Int64? { + guard let result = withLock({ () -> Int64? in + let current = readState() + let next = max(current, minimum) + 1 + guard next > 0 else { return nil } + do { + try writeState(next) + return next + } catch { + return nil + } + }) else { return nil } + return result + } + + private func withLock(_ body: () -> T) -> T? { + let fileManager = FileManager.default + do { + try fileManager.createDirectory( + at: lockURL.deletingLastPathComponent(), + withIntermediateDirectories: true + ) + let handle = try FileHandle(forWritingTo: lockURL) + defer { try? handle.close() } + guard flock(handle.fileDescriptor, LOCK_EX) == 0 else { return nil } + defer { flock(handle.fileDescriptor, LOCK_UN) } + return body() + } catch { + _ = fileManager.createFile(atPath: lockURL.path, contents: nil) + guard let handle = try? FileHandle(forWritingTo: lockURL) else { return nil } + defer { try? handle.close() } + guard flock(handle.fileDescriptor, LOCK_EX) == 0 else { return nil } + defer { flock(handle.fileDescriptor, LOCK_UN) } + return body() + } + } + + private func readState() -> Int64 { + guard let data = try? Data(contentsOf: stateURL), + let string = String(data: data, encoding: .utf8), + let value = Int64(string.trimmingCharacters(in: .whitespacesAndNewlines)), + value > 0 else { + return 0 + } + return value + } + + private func writeState(_ value: Int64) throws { + try Data("\(value)\n".utf8).write(to: stateURL, options: .atomic) + } +} + // Sendable safety: every mutable field is protected by `lock`; `semaphore` only wakes `next(timeout:)`. final class CmuxEventSubscription: @unchecked Sendable { let id: UUID @@ -19,6 +104,7 @@ final class CmuxEventSubscription: @unchecked Sendable { private var asyncWaiters: [CheckedContinuation<[String: Any]?, Never>] = [] private var closed = false private var closedReason: String? + private var replayQueueCount = 0 init(id: UUID = UUID(), names: Set, categories: Set, maxPendingEvents: Int) { self.id = id @@ -63,10 +149,11 @@ final class CmuxEventSubscription: @unchecked Sendable { shouldSignal = false accepted = true waiter = nextWaiter - } else if queue.count >= maxPendingEvents { + } else if queue.count - replayQueueCount >= maxPendingEvents { closed = true closedReason = "pending event buffer exceeded \(maxPendingEvents) events" queue.removeAll() + replayQueueCount = 0 shouldSignal = true accepted = false waiter = nil @@ -85,6 +172,28 @@ final class CmuxEventSubscription: @unchecked Sendable { return accepted } + /// Enqueues the bounded restored window without applying the live-event + /// backpressure limit. Restored history is finite and must be delivered in + /// full before a subscription is considered a slow consumer. + func enqueueReplay(_ event: [String: Any]) -> Bool { + lock.lock() + guard !closed else { + lock.unlock() + return false + } + if let waiter = asyncWaiters.first { + asyncWaiters.removeFirst() + lock.unlock() + waiter.resume(returning: event) + return true + } + queue.append(event) + replayQueueCount += 1 + lock.unlock() + semaphore.signal() + return true + } + /// Awaits the next event without tying up a thread in a semaphore wait. /// Cancellation closes the subscription so the waiter is always resumed. func nextAsync() async -> [String: Any]? { @@ -93,6 +202,7 @@ final class CmuxEventSubscription: @unchecked Sendable { lock.lock() if !queue.isEmpty { let event = queue.removeFirst() + if replayQueueCount > 0 { replayQueueCount -= 1 } lock.unlock() continuation.resume(returning: event) } else if closed { @@ -112,6 +222,7 @@ final class CmuxEventSubscription: @unchecked Sendable { lock.lock() if !queue.isEmpty { let event = queue.removeFirst() + if replayQueueCount > 0 { replayQueueCount -= 1 } lock.unlock() return event } @@ -127,7 +238,9 @@ final class CmuxEventSubscription: @unchecked Sendable { lock.lock() defer { lock.unlock() } guard !queue.isEmpty else { return nil } - return queue.removeFirst() + let event = queue.removeFirst() + if replayQueueCount > 0 { replayQueueCount -= 1 } + return event } func close(reason: String? = nil) { @@ -137,6 +250,7 @@ final class CmuxEventSubscription: @unchecked Sendable { closedReason = reason } queue.removeAll() + replayQueueCount = 0 let waiters = asyncWaiters asyncWaiters.removeAll(keepingCapacity: true) lock.unlock() @@ -147,6 +261,14 @@ final class CmuxEventSubscription: @unchecked Sendable { // Sendable safety: event state is protected by `lock`; disk appends are delegated to `CmuxEventLogWriter`. final class CmuxEventBus: @unchecked Sendable { + private struct PersistedEventRestore { + let events: [[String: Any]] + let allEvents: [[String: Any]] + let nextSequence: Int64 + let gap: Bool + let needsRewrite: Bool + } + private struct PendingPublish { let name: String let category: String @@ -190,9 +312,12 @@ final class CmuxEventBus: @unchecked Sendable { private let lock = NSLock() private let retainedEventLimit: Int + private let eventLogURL: URL? + private let maxEventLogBytes: UInt64 private let maxEventLineBytes: Int private let maxPendingEventsPerSubscription: Int private let eventLogWriter: CmuxEventLogWriter? + private let sequenceStore: CmuxEventSequenceStore? private let bootId = UUID().uuidString private var restorePending: Bool private var restoreGap = false @@ -212,10 +337,13 @@ final class CmuxEventBus: @unchecked Sendable { maxPendingEventsPerSubscription: Int = CmuxEventBus.defaultMaxPendingEventsPerSubscription ) { self.retainedEventLimit = max(1, retainedEventLimit) + self.eventLogURL = eventLogURL + self.maxEventLogBytes = max(1, maxEventLogBytes) self.maxEventLineBytes = max(1, maxEventLineBytes) self.maxPendingEventsPerSubscription = max(1, maxPendingEventsPerSubscription) self.restorePending = eventLogURL != nil self.restoreTask = nil + self.sequenceStore = eventLogURL.map(CmuxEventSequenceStore.init(eventLogURL:)) self.eventLogWriter = eventLogURL.map { CmuxEventLogWriter( eventLogURL: $0, @@ -225,6 +353,7 @@ final class CmuxEventBus: @unchecked Sendable { } if let eventLogURL { + let maxEventLogBytes = self.maxEventLogBytes self.restoreTask = Task.detached(priority: .utility) { [weak self] in let restored = Self.loadPersistedEvents( eventLogURL: eventLogURL, @@ -296,8 +425,8 @@ final class CmuxEventBus: @unchecked Sendable { /// The caller must hold ``lock`` while appending the event. private func appendEventLocked(_ pending: PendingPublish) -> EventPublication { - let sequence = nextSequence - nextSequence += 1 + let sequence = sequenceStore?.allocate(minimum: nextSequence - 1) ?? nextSequence + nextSequence = sequence + 1 var event: [String: Any] = [ "type": "event", @@ -445,7 +574,7 @@ final class CmuxEventBus: @unchecked Sendable { } private func completeRestore( - _ restored: (events: [[String: Any]], nextSequence: Int64, gap: Bool) + _ restored: PersistedEventRestore ) { lock.lock() guard restorePending else { @@ -457,12 +586,22 @@ final class CmuxEventBus: @unchecked Sendable { restoreGap = restored.gap lock.unlock() + sequenceStore?.raiseHighWater(to: restored.nextSequence - 1) + if restored.needsRewrite, let eventLogURL { + Self.rewritePersistedEvents( + restored.allEvents, + eventLogURL: eventLogURL, + maxEventLogBytes: maxEventLogBytes + ) + } + while true { lock.lock() let subscriptionsToReplay = Array(pendingSubscriptions.values) pendingSubscriptions.removeAll() let publishesToFlush = pendingPublishes pendingPublishes.removeAll() + let replayWindow = retained let publicationsToFlush = publishesToFlush.map { appendEventLocked($0) } if subscriptionsToReplay.isEmpty, publicationsToFlush.isEmpty { restorePending = false @@ -475,11 +614,11 @@ final class CmuxEventBus: @unchecked Sendable { let afterSequence = pending.liveOnly ? restored.nextSequence - 1 : (pending.afterSequence ?? restored.nextSequence - 1) - let replay = restored.events.filter { event in + let replay = replayWindow.filter { event in let sequence = Self.int64(event["seq"]) ?? 0 return sequence > afterSequence && pending.subscription.accepts(event) } - for event in replay where !pending.subscription.enqueue(event) { + for event in replay where !pending.subscription.enqueueReplay(event) { removeSubscriptionIfStillActive(pending.subscription) break } @@ -538,6 +677,7 @@ final class CmuxEventBus: @unchecked Sendable { pendingSubscriptions.removeAll() lock.unlock() active.forEach { $0.close() } + sequenceStore?.resetForTesting() eventLogWriter?.resetForTesting() } @@ -570,10 +710,11 @@ final class CmuxEventBus: @unchecked Sendable { eventLogURL: URL, maxEventLogBytes: UInt64, retainedEventLimit: Int - ) -> (events: [[String: Any]], nextSequence: Int64, gap: Bool) { + ) -> PersistedEventRestore { let rotatedURL = eventLogURL.appendingPathExtension("1") - var loaded: [[String: Any]] = [] + var segments: [[[String: Any]]] = [] var gap = false + var needsRewrite = false for url in [rotatedURL, eventLogURL] { guard FileManager.default.fileExists(atPath: url.path) else { @@ -588,6 +729,7 @@ final class CmuxEventBus: @unchecked Sendable { continue } + var segment: [[String: Any]] = [] for line in text.split(whereSeparator: \.isNewline) { guard let lineData = String(line).data(using: .utf8), let object = try? JSONSerialization.jsonObject(with: lineData) as? [String: Any], @@ -595,12 +737,15 @@ final class CmuxEventBus: @unchecked Sendable { let sequence = int64(object["seq"]), sequence > 0 else { gap = true + needsRewrite = true continue } - loaded.append(object) + segment.append(object) } + segments.append(segment) } + var loaded = segments.flatMap { $0 } var highestSequence: Int64 = 0 for index in loaded.indices { guard let originalSequence = int64(loaded[index]["seq"]) else { continue } @@ -610,17 +755,70 @@ final class CmuxEventBus: @unchecked Sendable { if normalizedSequence != originalSequence { loaded[index]["legacy_seq"] = NSNumber(value: originalSequence) loaded[index]["seq"] = NSNumber(value: normalizedSequence) + needsRewrite = true } highestSequence = normalizedSequence } - return ( - Array(loaded.suffix(max(1, retainedEventLimit))), - highestSequence + 1, - gap + let persistedHighWater = CmuxEventSequenceStore(eventLogURL: eventLogURL).current() + let nextSequence = max(highestSequence, persistedHighWater) + 1 + + return PersistedEventRestore( + events: Array(loaded.suffix(max(1, retainedEventLimit))), + allEvents: loaded, + nextSequence: nextSequence, + gap: gap || persistedHighWater > highestSequence, + needsRewrite: needsRewrite ) } + private static func rewritePersistedEvents( + _ events: [[String: Any]], + eventLogURL: URL, + maxEventLogBytes: UInt64 + ) { + let fileManager = FileManager.default + var segments: [[String]] = [[]] + var currentSize: UInt64 = 0 + + for event in events { + var persistedEvent = event + persistedEvent.removeValue(forKey: "legacy_seq") + guard let line = encodeLine(persistedEvent) else { return } + let lineBytes = UInt64(line.utf8.count) + 1 + guard lineBytes <= maxEventLogBytes else { return } + if currentSize > 0, currentSize + lineBytes > maxEventLogBytes { + segments.append([]) + currentSize = 0 + } + segments[segments.count - 1].append(line) + currentSize += lineBytes + } + + if segments.count > 2 { + segments = Array(segments.suffix(2)) + } + + do { + try fileManager.createDirectory( + at: eventLogURL.deletingLastPathComponent(), + withIntermediateDirectories: true + ) + let rotatedURL = eventLogURL.appendingPathExtension("1") + if segments.count == 2 { + try Data((segments[0].joined(separator: "\n") + "\n").utf8) + .write(to: rotatedURL, options: .atomic) + } else if fileManager.fileExists(atPath: rotatedURL.path) { + try fileManager.removeItem(at: rotatedURL) + } + let current = segments.last ?? [] + try Data((current.isEmpty ? "" : current.joined(separator: "\n") + "\n").utf8) + .write(to: eventLogURL, options: .atomic) + } catch { + // Recovery remains usable in memory; the next restore will retry the rewrite. + } + } + static func encodeLine(_ object: [String: Any]) -> String? { let clean = sanitizedJSONValue(object) guard JSONSerialization.isValidJSONObject(clean), diff --git a/cmuxTests/CmuxEventBusTests.swift b/cmuxTests/CmuxEventBusTests.swift index e9bb52a5389a..fb8c703b3bef 100644 --- a/cmuxTests/CmuxEventBusTests.swift +++ b/cmuxTests/CmuxEventBusTests.swift @@ -127,6 +127,23 @@ final class CmuxEventBusTests: XCTestCase { XCTAssertEqual(replay.replay.compactMap { CmuxEventBus.int64($0["seq"]) }, [1, 2, 3]) XCTAssertEqual(CmuxEventBus.int64(replay.replay[2]["legacy_seq"]), 1) XCTAssertEqual(bus.latestSequence, 3) + + let persistedSequences = try String(contentsOf: logURL, encoding: .utf8) + .split(whereSeparator: \.isNewline) + .compactMap { line -> Int64? in + guard let data = String(line).data(using: .utf8), + let object = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else { + return nil + } + return CmuxEventBus.int64(object["seq"]) + } + XCTAssertEqual(persistedSequences, [1, 2, 3]) + + let secondBus = CmuxEventBus(retainedEventLimit: 4, eventLogURL: logURL) + await secondBus.waitUntilRestored() + let secondReplay = secondBus.subscribe(afterSequence: 0, names: [], categories: []) + defer { secondBus.unsubscribe(secondReplay.subscription) } + XCTAssertEqual(secondReplay.replay.compactMap { CmuxEventBus.int64($0["seq"]) }, [1, 2, 3]) } func testDurableReplayReportsUnreadableRecordsAsAGap() async throws { @@ -186,6 +203,57 @@ final class CmuxEventBusTests: XCTestCase { XCTAssertNil(snapshot.subscription.next(timeout: 0.05)) } + func testRestoredReplayDoesNotUseLivePendingLimit() { + let subscription = CmuxEventSubscription(names: [], categories: [], maxPendingEvents: 1) + defer { subscription.close() } + let events = (1...3).map { sequence in + ["type": "event", "seq": sequence, "name": "event", "category": "test"] as [String: Any] + } + + for event in events { + XCTAssertTrue(subscription.enqueueReplay(event)) + } + XCTAssertFalse(subscription.isClosed) + let received = (1...3).compactMap { _ in + subscription.next(timeout: 0.05).flatMap { CmuxEventBus.int64($0["seq"]) } + } + XCTAssertEqual(received, [1, 2, 3]) + } + + func testUnreadableSegmentUsesPersistedSequenceHighWater() async throws { + let directory = FileManager.default.temporaryDirectory + .appendingPathComponent("cmux-event-replay-high-water-\(UUID().uuidString)", isDirectory: true) + let logURL = directory.appendingPathComponent("events.jsonl") + defer { try? FileManager.default.removeItem(at: directory) } + + let older = [ + "type": "event", "seq": 1, "id": "old-1", + "name": "old", "category": "test", "source": "test", "payload": [:] + ] as [String: Any] + let olderLine = try XCTUnwrap(CmuxEventBus.encodeLine(older)) + try FileManager.default.createDirectory(at: directory, withIntermediateDirectories: true) + try olderLine.appending("\n").write( + to: logURL.appendingPathExtension("1"), + atomically: true, + encoding: .utf8 + ) + try Data(repeating: 0x78, count: 512).write(to: logURL) + try "100\n".write( + to: logURL.appendingPathExtension("seq"), + atomically: true, + encoding: .utf8 + ) + + let bus = CmuxEventBus(retainedEventLimit: 4, eventLogURL: logURL, maxEventLogBytes: 256) + await bus.waitUntilRestored() + bus.publish(name: "new", category: "test", source: "test") + + XCTAssertEqual(bus.latestSequence, 101) + let snapshot = bus.subscribe(afterSequence: 100, names: [], categories: []) + defer { bus.unsubscribe(snapshot.subscription) } + XCTAssertEqual(snapshot.replay.compactMap { CmuxEventBus.int64($0["seq"]) }, [101]) + } + func testEventEncodingIsSingleLineJSON() throws { let bus = CmuxEventBus(retainedEventLimit: 4) bus.publish(