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
32 changes: 27 additions & 5 deletions Sources/CmuxEventLogWriter.swift
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ final class CmuxEventLogWriter: @unchecked Sendable {
private let eventLogURL: URL
private let maxEventLogBytes: UInt64
private let maxPendingLines: Int
private let writeData: @Sendable (FileHandle, Data) throws -> Void
private let lock = NSLock()
private var pendingLines: [String] = []
private var flushScheduled = false
Expand All @@ -20,10 +21,18 @@ final class CmuxEventLogWriter: @unchecked Sendable {
private var flushSuspendedForTesting = false
#endif

init(eventLogURL: URL, maxEventLogBytes: UInt64, maxPendingLines: Int) {
init(
eventLogURL: URL,
maxEventLogBytes: UInt64,
maxPendingLines: Int,
writeData: @escaping @Sendable (FileHandle, Data) throws -> Void = { handle, data in
try handle.write(contentsOf: data)
}
) {
self.eventLogURL = eventLogURL
self.maxEventLogBytes = max(1, maxEventLogBytes)
self.maxPendingLines = max(1, maxPendingLines)
self.writeData = writeData
}

func enqueue(_ line: String) {
Expand Down Expand Up @@ -140,17 +149,30 @@ final class CmuxEventLogWriter: @unchecked Sendable {
defer { try? handle.close() }
try handle.seekToEnd()
var currentSize = Self.fileSize(at: eventLogURL, fileManager: fileManager)
// One file segment at a time bounds the extra buffer to the rotation limit,
// except for an indivisible oversized record (the existing write-whole policy).
var batchData = Data()

func writeBatch() throws {
guard !batchData.isEmpty else { return }
try writeData(handle, batchData)
batchData.removeAll(keepingCapacity: true)
}

for line in lines {
let data = Data((line + "\n").utf8)
if currentSize + UInt64(data.count) > maxEventLogBytes {
let lineBytes = UInt64(line.utf8.count) + 1
if currentSize + lineBytes > maxEventLogBytes {
try writeBatch()
try handle.close()
try rotate(fileManager: fileManager)
handle = try FileHandle(forWritingTo: eventLogURL)
currentSize = 0
}
try handle.write(contentsOf: data)
currentSize += UInt64(data.count)
batchData.append(contentsOf: line.utf8)
batchData.append(0x0a)
currentSize += lineBytes
}
try writeBatch()
} catch {
cmuxEventLogLogger.error("Failed to append cmux event log: \(String(describing: error), privacy: .private)")
}
Expand Down
8 changes: 8 additions & 0 deletions cmux.xcodeproj/project.pbxproj
Original file line number Diff line number Diff line change
Expand Up @@ -1254,6 +1254,8 @@
E7E000000000000000000005 /* CmuxEventBus.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E000000000000000000006 /* CmuxEventBus.swift */; };
E7E000000000000000000003 /* CmuxEventBusTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E000000000000000000004 /* CmuxEventBusTests.swift */; };
E7E00000000000000000000D /* CmuxEventLogWriter.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E00000000000000000000E /* CmuxEventLogWriter.swift */; };
E7E00000000000000000000F /* CmuxEventLogWriterTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E000000000000000000010 /* CmuxEventLogWriterTests.swift */; };
E7E000000000000000000011 /* CmuxEventLogWriteSpy.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E000000000000000000012 /* CmuxEventLogWriteSpy.swift */; };
E7E000000000000000000007 /* CmuxEventPublishing.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E000000000000000000008 /* CmuxEventPublishing.swift */; };
E7E000000000000000000001 /* CmuxEventStream.swift in Sources */ = {isa = PBXBuildFile; fileRef = E7E000000000000000000002 /* CmuxEventStream.swift */; };
C0DE45000000000000000001 /* CmuxExtensionKit in Frameworks */ = {isa = PBXBuildFile; productRef = C0DE45000000000000000002 /* CmuxExtensionKit */; };
Expand Down Expand Up @@ -5178,6 +5180,8 @@
E7E000000000000000000006 /* CmuxEventBus.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventBus.swift; sourceTree = "<group>"; };
E7E000000000000000000004 /* CmuxEventBusTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventBusTests.swift; sourceTree = "<group>"; };
E7E00000000000000000000E /* CmuxEventLogWriter.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventLogWriter.swift; sourceTree = "<group>"; };
E7E000000000000000000010 /* CmuxEventLogWriterTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventLogWriterTests.swift; sourceTree = "<group>"; };
E7E000000000000000000012 /* CmuxEventLogWriteSpy.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventLogWriteSpy.swift; sourceTree = "<group>"; };
E7E000000000000000000008 /* CmuxEventPublishing.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventPublishing.swift; sourceTree = "<group>"; };
E7E000000000000000000002 /* CmuxEventStream.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = CmuxEventStream.swift; sourceTree = "<group>"; };
C57B00050000000000000002 /* CmuxExtensionSidebarSelection+CustomSidebarName.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = "CmuxExtensionSidebarSelection+CustomSidebarName.swift"; sourceTree = "<group>"; };
Expand Down Expand Up @@ -11467,6 +11471,8 @@
A7984AD10000000000000002 /* SocketConfigurationLifecycleTests.swift */,
9C1BEA3D2E6F49709A71C021 /* TerminalControllerTerminalTextTests.swift */,
E7E000000000000000000004 /* CmuxEventBusTests.swift */,
E7E000000000000000000010 /* CmuxEventLogWriterTests.swift */,
E7E000000000000000000012 /* CmuxEventLogWriteSpy.swift */,
A10649100000000000000001 /* AutomationRuleTests.swift */,
A10649340000000000000001 /* AutomationProcessSessionTests.swift */,
7837E0027837E0027837E002 /* CmuxSocketEventMapperTests.swift */,
Expand Down Expand Up @@ -15464,6 +15470,8 @@
A5FB1205 /* CmuxConfigWorkspaceActionTests.swift in Sources */,
C54860040000000000000001 /* CmuxDurableDeepLinkRestoreTests.swift in Sources */,
E7E000000000000000000003 /* CmuxEventBusTests.swift in Sources */,
E7E00000000000000000000F /* CmuxEventLogWriterTests.swift in Sources */,
E7E000000000000000000011 /* CmuxEventLogWriteSpy.swift in Sources */,
A11C00040000000000000001 /* CmuxHostedSystemSymbolImageTests.swift in Sources */,
D36090010000000000000005 /* CmuxMainWindowConstrainFrameTests.swift in Sources */,
D36090020000000000000005 /* CmuxMainWindowFullScreenCapabilityTests.swift in Sources */,
Expand Down
35 changes: 35 additions & 0 deletions cmuxTests/CmuxEventLogWriteSpy.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
import Foundation

/// A synchronous FileHandle dependency; the lock protects observations shared with the utility queue.
final class CmuxEventLogWriteSpy: @unchecked Sendable {
private let lock = NSLock()
private var sizes: [Int] = []
private var onMainThread = false
private let failedCall: Int?

init(failedCall: Int? = nil) {
self.failedCall = failedCall
}

var writeSizes: [Int] {
lock.lock()
defer { lock.unlock() }
return sizes
}

var wroteOnMainThread: Bool {
lock.lock()
defer { lock.unlock() }
return onMainThread
}

func write(_ handle: FileHandle, data: Data) throws {
lock.lock()
sizes.append(data.count)
onMainThread = onMainThread || Thread.isMainThread
let shouldFail = sizes.count == failedCall
lock.unlock()
if shouldFail { throw CocoaError(.fileWriteUnknown) }
try handle.write(contentsOf: data)
}
}
231 changes: 231 additions & 0 deletions cmuxTests/CmuxEventLogWriterTests.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,231 @@
import Foundation
import Testing

#if canImport(cmux_DEV)
@testable import cmux_DEV
#elseif canImport(cmux)
@testable import cmux
#endif

@Suite("Durable event-log batch writes", .serialized)
struct CmuxEventLogWriterTests {
private let logLimit = 16 * 1024 * 1024

@Test(arguments: [0, 1, 32, 256, 1_024])
func burstUsesOneWriteAndPreservesJSONL(count: Int) throws {
let lines = (0..<count).map { jsonLine(index: $0) }
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush(lines, with: writer)

#expect(spy.writeSizes.count == (count == 0 ? 0 : 1))
#expect(!spy.wroteOnMainThread)
if count > 0 {
let stored = try Data(contentsOf: url)
#expect(stored == jsonl(lines))
try expectJSONRecords(stored, count: count)
} else {
#expect(!FileManager.default.fileExists(atPath: url.path))
}
#expect(writer.backlogSnapshotForTesting().pending == 0)
#expect(writer.backlogSnapshotForTesting().dropped == 0)
}

@Test(arguments: [0, 1])
func batchAtOrBelowSixteenMiBLimitDoesNotRotate(spareBytes: Int) throws {
let lines = (0..<32).map { jsonLine(index: $0) }
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }
let seed = try seedLog(url, bytes: logLimit - jsonl(lines).count - spareBytes)

flush(lines, with: writer)

#expect(spy.writeSizes == [jsonl(lines).count])
#expect(try Data(contentsOf: url) == seed + jsonl(lines))
#expect(!FileManager.default.fileExists(atPath: url.appendingPathExtension("1").path))
}

@Test
func batchCrossingSixteenMiBLimitWritesOncePerFile() throws {
let lines = (0..<32).map { jsonLine(index: $0) }
let prefix = jsonl(Array(lines.prefix(13)))
let suffix = jsonl(Array(lines.dropFirst(13)))
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }
let seed = try seedLog(url, bytes: logLimit - prefix.count)

flush(lines, with: writer)

#expect(spy.writeSizes == [prefix.count, suffix.count])
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == seed + prefix)
#expect(try Data(contentsOf: url) == suffix)
try expectJSONRecords(suffix, count: 19)
}

@Test
func maximumPendingBatchBoundsWritesAtSixteenMiB() throws {
// 1,024 maximum-sized producer records plus JSONL delimiters straddle the cap.
let line = "{\"text\":\"" + String(repeating: "x", count: 16_384 - 11) + "\"}"
#expect(line.utf8.count == 16_384)
let lines = Array(repeating: line, count: 1_024)
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush(lines, with: writer)

#expect(spy.writeSizes == [16_385 * 1_023, 16_385])
#expect(spy.writeSizes.allSatisfy { $0 <= logLimit })
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == jsonl(Array(lines.prefix(1_023))))
#expect(try Data(contentsOf: url) == jsonl([line]))
}

@Test
func fullExistingLogRotatesBeforeWritingTheBatch() throws {
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }
let seed = try seedLog(url, bytes: logLimit)
let lines = (0..<32).map { jsonLine(index: $0) }

flush(lines, with: writer)

#expect(spy.writeSizes == [jsonl(lines).count])
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == seed)
#expect(try Data(contentsOf: url) == jsonl(lines))
}

@Test
func multipleRotationsRetainTheSameLastTwoFiles() throws {
let lines = (0..<8).map { "{\"seq\":\($0)}" }
let (writer, url, spy) = makeWriter(maxBytes: 30)
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush(lines, with: writer)

#expect(spy.writeSizes == [30, 30, 20])
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == jsonl(Array(lines[3..<6])))
#expect(try Data(contentsOf: url) == jsonl(Array(lines[6..<8])))
}

@Test
func utf8BytesAndNewlinesDetermineTheBoundary() throws {
let lines = [#"{"text":"🌍"}"#, #"{"text":"é"}"#, #"{"seq":2}"#]
let prefix = jsonl(Array(lines.prefix(2)))
let (writer, url, spy) = makeWriter(maxBytes: UInt64(prefix.count))
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush(lines, with: writer)

#expect(spy.writeSizes == [prefix.count, jsonl([lines[2]]).count])
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == prefix)
#expect(try Data(contentsOf: url) == jsonl([lines[2]]))
}

@Test
func oversizedRecordKeepsExistingWholeRecordAndCleanupPolicy() throws {
let oversized = jsonLine(index: 1)
let seed = #"{"seq":0}"#
let tail = #"{"seq":2}"#
let (writer, url, spy) = makeWriter(maxBytes: 32)
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush([seed, oversized], with: writer)
#expect(try Data(contentsOf: url) == jsonl([oversized]))
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == jsonl([seed]))

// The pre-existing policy discards an oversized active log on the next rotation.
flush([tail], with: writer)
#expect(spy.writeSizes == [jsonl([seed]).count, jsonl([oversized]).count, jsonl([tail]).count])
#expect(try Data(contentsOf: url) == jsonl([tail]))
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == jsonl([seed]))
}

@Test
func suspendedBurstKeepsNewestLinesAndDropAccounting() throws {
let lines = (0..<1_152).map { jsonLine(index: $0) }
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }
lines.forEach(writer.enqueue)

#expect(writer.backlogSnapshotForTesting().pending == 1_024)
#expect(writer.backlogSnapshotForTesting().dropped == 128)
writer.setFlushSuspendedForTesting(false)
writer.flushForTesting()

#expect(spy.writeSizes.count == 1)
#expect(try Data(contentsOf: url) == jsonl(Array(lines.suffix(1_024))))
#expect(writer.backlogSnapshotForTesting().pending == 0)
#expect(writer.backlogSnapshotForTesting().dropped == 0)
}

@Test(arguments: [1, 2])
func failedWriteStopsTheBatchAndNextFlushRecovers(failedCall: Int) throws {
let lines = (0..<8).map { "{\"seq\":\($0)}" }
let spy = CmuxEventLogWriteSpy(failedCall: failedCall)
let (writer, url, _) = makeWriter(maxBytes: 30, spy: spy)
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush(lines, with: writer)

#expect(spy.writeSizes.count == failedCall)
#expect(try Data(contentsOf: url).isEmpty)
if failedCall == 2 {
#expect(try Data(contentsOf: url.appendingPathExtension("1")) == jsonl(Array(lines.prefix(3))))
} else {
#expect(!FileManager.default.fileExists(atPath: url.appendingPathExtension("1").path))
}
// As before, failed I/O is logged and abandoned, not retried or counted as queue drops.
#expect(writer.backlogSnapshotForTesting().dropped == 0)
flush([#"{"seq":9}"#], with: writer)
#expect(try Data(contentsOf: url) == jsonl([#"{"seq":9}"#]))
}

private func makeWriter(
maxBytes: UInt64 = 16 * 1024 * 1024,
spy: CmuxEventLogWriteSpy = CmuxEventLogWriteSpy()
) -> (CmuxEventLogWriter, URL, CmuxEventLogWriteSpy) {
let url = FileManager.default.temporaryDirectory
.appendingPathComponent("cmux-event-log-writer-\(UUID().uuidString)", isDirectory: true)
.appendingPathComponent("events.jsonl")
let writer = CmuxEventLogWriter(
eventLogURL: url,
maxEventLogBytes: maxBytes,
maxPendingLines: 1_024,
writeData: { try spy.write($0, data: $1) }
)
writer.setFlushSuspendedForTesting(true)
return (writer, url, spy)
}

private func flush(_ lines: [String], with writer: CmuxEventLogWriter) {
writer.setFlushSuspendedForTesting(true)
lines.forEach(writer.enqueue)
writer.setFlushSuspendedForTesting(false)
writer.flushForTesting()
}

private func jsonLine(index: Int) -> String {
"{\"seq\":\(index),\"name\":\"agent.hook.PreToolUse\",\"payload\":\"\(String(repeating: "x", count: 900))\"}"
}

private func jsonl(_ lines: [String]) -> Data {
Data(lines.map { $0 + "\n" }.joined().utf8)
}

private func seedLog(_ url: URL, bytes: Int) throws -> Data {
try FileManager.default.createDirectory(at: url.deletingLastPathComponent(), withIntermediateDirectories: true)
let data = Data(("{\"seed\":\"" + String(repeating: "x", count: bytes - 12) + "\"}\n").utf8)
#expect(data.count == bytes)
try data.write(to: url)
return data
}

private func expectJSONRecords(_ data: Data, count: Int) throws {
#expect(data.last == 0x0a)
let records = data.split(separator: 0x0a)
#expect(records.count == count)
for record in records {
#expect(try JSONSerialization.jsonObject(with: Data(record)) is [String: Any])
}
}
}
Loading
Loading