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
83 changes: 61 additions & 22 deletions Sources/CmuxEventLogWriter.swift
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,8 @@ import os

nonisolated private let cmuxEventLogLogger = Logger(subsystem: "com.cmuxterm.app", category: "event-log")

// Sendable safety: pending state is protected by `lock`; file IO runs on `queue`.
// Sendable safety: pending state is protected by `lock`; file IO and
// `openLogHandle` are confined to `queue`.
final class CmuxEventLogWriter: @unchecked Sendable {
static let defaultMaxPendingLines = 1_024

Expand All @@ -17,6 +18,8 @@ final class CmuxEventLogWriter: @unchecked Sendable {
private var pendingLines: [String] = []
private var flushScheduled = false
private var droppedLineCount = 0
/// The append-only log descriptor kept open across flushes.
private var openLogHandle: FileHandle?
#if DEBUG
private var flushSuspendedForTesting = false
#endif
Expand Down Expand Up @@ -136,15 +139,9 @@ final class CmuxEventLogWriter: @unchecked Sendable {

private func append(_ lines: [String]) {
guard !lines.isEmpty else { return }
let fileManager = FileManager.default
do {
let fileManager = FileManager.default
// Every agent hook event flushes here, so the steady state is one
// open + seek + write. Directory and file creation run only when the
// open fails, and the end-of-file offset doubles as the current size
// instead of a separate attributesOfItem stat per flush.
var handle = try openForAppending(fileManager: fileManager)
defer { try? handle.close() }
var currentSize = try handle.seekToEnd()
var (handle, currentSize) = try logHandleForAppending(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()
Expand All @@ -159,33 +156,75 @@ final class CmuxEventLogWriter: @unchecked Sendable {
let lineBytes = UInt64(line.utf8.count) + 1
if currentSize + lineBytes > maxEventLogBytes {
try writeBatch()
try handle.close()
closeOpenLog()
try rotate(fileManager: fileManager)
handle = try FileHandle(forWritingTo: eventLogURL)
currentSize = 0
(handle, currentSize) = try logHandleForAppending(fileManager: fileManager)
}
batchData.append(contentsOf: line.utf8)
batchData.append(0x0a)
currentSize += lineBytes
}
try writeBatch()
} catch {
closeOpenLog()
cmuxEventLogLogger.error("Failed to append cmux event log: \(String(describing: error), privacy: .private)")
}
}

private func openForAppending(fileManager: FileManager) throws -> FileHandle {
if let handle = try? FileHandle(forWritingTo: eventLogURL) {
return handle
/// Returns an append-only handle for the log and the log's current size.
///
/// Events are flushed one burst at a time, often a single line. Reopening,
/// creating the directory, and stat'ing the file on every flush dominated the
/// event-log queue, so the handle stays open. Other cmux processes share this
/// log and may rotate, delete, or append to it:
/// - The open handle is reused only while the path still names the same file
/// (device and inode match).
/// - The descriptor is opened with `O_APPEND`, so every write lands at the
/// current end of file even when another process appended since our last
/// write. A seek-then-write handle could overwrite those lines.
/// - The size used for rotation comes from `fstat` on the open descriptor.
private func logHandleForAppending(fileManager: FileManager) throws -> (FileHandle, UInt64) {
let path = eventLogURL.path
var pathStatus = stat()
let pathExists = stat(path, &pathStatus) == 0
if let handle = openLogHandle {
var handleStatus = stat()
if pathExists,
fstat(handle.fileDescriptor, &handleStatus) == 0,
handleStatus.st_dev == pathStatus.st_dev,
handleStatus.st_ino == pathStatus.st_ino {
return (handle, UInt64(max(0, handleStatus.st_size)))
}
closeOpenLog()
}
try fileManager.createDirectory(
at: eventLogURL.deletingLastPathComponent(),
withIntermediateDirectories: true
)
if !fileManager.fileExists(atPath: eventLogURL.path) {
_ = fileManager.createFile(atPath: eventLogURL.path, contents: nil)

// Directory creation runs only when the open reports a missing parent.
var descriptor = open(path, O_WRONLY | O_APPEND | O_CREAT | O_CLOEXEC, 0o644)
if descriptor < 0, errno == ENOENT {
try fileManager.createDirectory(
at: eventLogURL.deletingLastPathComponent(),
withIntermediateDirectories: true
)
descriptor = open(path, O_WRONLY | O_APPEND | O_CREAT | O_CLOEXEC, 0o644)
}
guard descriptor >= 0 else {
throw POSIXError(POSIXErrorCode(rawValue: errno) ?? .EIO)
}
return try FileHandle(forWritingTo: eventLogURL)
let handle = FileHandle(fileDescriptor: descriptor, closeOnDealloc: true)
var handleStatus = stat()
guard fstat(descriptor, &handleStatus) == 0 else {
let code = POSIXErrorCode(rawValue: errno) ?? .EIO
try? handle.close()
throw POSIXError(code)
}
openLogHandle = handle
return (handle, UInt64(max(0, handleStatus.st_size)))
}

private func closeOpenLog() {
guard let handle = openLogHandle else { return }
try? handle.close()
openLogHandle = nil
}

private func rotate(fileManager: FileManager) throws {
Expand Down
15 changes: 14 additions & 1 deletion cmuxTests/CmuxEventLogWriteSpy.swift
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,14 @@ final class CmuxEventLogWriteSpy: @unchecked Sendable {
private let lock = NSLock()
private var sizes: [Int] = []
private var onMainThread = false
// Retained so a closed handle's address cannot be reused by a later one.
private var handles: [FileHandle] = []
private let failedCall: Int?
private let beforeWrite: (@Sendable () throws -> Void)?

init(failedCall: Int? = nil) {
init(failedCall: Int? = nil, beforeWrite: (@Sendable () throws -> Void)? = nil) {
self.failedCall = failedCall
self.beforeWrite = beforeWrite
}

var writeSizes: [Int] {
Expand All @@ -17,6 +21,13 @@ final class CmuxEventLogWriteSpy: @unchecked Sendable {
return sizes
}

/// Identity of the handle passed to each write, in call order.
var handleIdentities: [ObjectIdentifier] {
lock.lock()
defer { lock.unlock() }
return handles.map(ObjectIdentifier.init)
}

var wroteOnMainThread: Bool {
lock.lock()
defer { lock.unlock() }
Expand All @@ -26,10 +37,12 @@ final class CmuxEventLogWriteSpy: @unchecked Sendable {
func write(_ handle: FileHandle, data: Data) throws {
lock.lock()
sizes.append(data.count)
handles.append(handle)
onMainThread = onMainThread || Thread.isMainThread
let shouldFail = sizes.count == failedCall
lock.unlock()
if shouldFail { throw CocoaError(.fileWriteUnknown) }
try beforeWrite?()
try handle.write(contentsOf: data)
}
}
99 changes: 99 additions & 0 deletions cmuxTests/CmuxEventLogWriterTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,86 @@ struct CmuxEventLogWriterTests {
#expect(try Data(contentsOf: url) == jsonl([#"{"seq":9}"#]))
}

/// Every hook, feed, and sidebar event is flushed on its own. Reopening,
/// seeking, and stat'ing the log per flush made the event-log queue one of
/// the busiest background queues in an idle app sample.
@Test
func consecutiveFlushesReuseOneOpenHandle() throws {
let (writer, url, spy) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush([#"{"seq":1}"#], with: writer)
flush([#"{"seq":2}"#], with: writer)
flush([#"{"seq":3}"#], with: writer)

#expect(spy.writeSizes.count == 3)
#expect(Set(spy.handleIdentities).count == 1)
#expect(try Data(contentsOf: url) == jsonl([#"{"seq":1}"#, #"{"seq":2}"#, #"{"seq":3}"#]))
}

/// Another cmux process (a tagged dev build shares `~/.cmuxterm/events.jsonl`)
/// can rotate or delete the log between flushes. The next flush must land in
/// the file now at the path, not the renamed or unlinked inode.
@Test(arguments: [false, true])
func externalRotationOrDeletionBetweenFlushesWritesToCurrentPath(deleteInsteadOfRotate: Bool) throws {
let (writer, url, _) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }
let rotatedURL = url.appendingPathExtension("external")

flush([#"{"seq":1}"#], with: writer)
if deleteInsteadOfRotate {
try FileManager.default.removeItem(at: url)
} else {
try FileManager.default.moveItem(at: url, to: rotatedURL)
}
flush([#"{"seq":2}"#], with: writer)
Comment thread
coderabbitai[bot] marked this conversation as resolved.

#expect(try Data(contentsOf: url) == jsonl([#"{"seq":2}"#]))
if !deleteInsteadOfRotate {
#expect(try Data(contentsOf: rotatedURL) == jsonl([#"{"seq":1}"#]))
}
}

/// When another process rotates and has already created the next log, the
/// path exists but names a different inode. The writer must switch to it.
@Test
func externalRotationWithReplacementFileWritesToReplacement() throws {
let (writer, url, _) = makeWriter()
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }
let rotatedURL = url.appendingPathExtension("external")

flush([#"{"seq":1}"#], with: writer)
try FileManager.default.moveItem(at: url, to: rotatedURL)
#expect(FileManager.default.createFile(atPath: url.path, contents: nil))
flush([#"{"seq":2}"#], with: writer)

#expect(try Data(contentsOf: url) == jsonl([#"{"seq":2}"#]))
#expect(try Data(contentsOf: rotatedURL) == jsonl([#"{"seq":1}"#]))
}

/// Another cmux process can append to the shared log after this writer has
/// positioned its handle but before its write lands. The append-only
/// descriptor must keep both lines instead of overwriting the other writer's.
@Test
func concurrentExternalAppendIsNotOverwritten() throws {
let urlBox = CmuxEventLogURLBox()
let spy = CmuxEventLogWriteSpy(beforeWrite: {
guard let url = urlBox.url else { return }
let external = try FileHandle(forWritingTo: url)
defer { try? external.close() }
try external.seekToEnd()
try external.write(contentsOf: Data("{\"external\":true}\n".utf8))
})
let (writer, url, _) = makeWriter(spy: spy)
defer { try? FileManager.default.removeItem(at: url.deletingLastPathComponent()) }

flush([#"{"seq":1}"#], with: writer)
urlBox.url = url
flush([#"{"seq":2}"#], with: writer)

#expect(try Data(contentsOf: url) == jsonl([#"{"seq":1}"#, #"{"external":true}"#, #"{"seq":2}"#]))
}

private func makeWriter(
maxBytes: UInt64 = 16 * 1024 * 1024,
spy: CmuxEventLogWriteSpy = CmuxEventLogWriteSpy()
Expand Down Expand Up @@ -229,3 +309,22 @@ struct CmuxEventLogWriterTests {
}
}
}

/// Lets a `@Sendable` write hook see the log URL once the writer exists.
private final class CmuxEventLogURLBox: @unchecked Sendable {
private let lock = NSLock()
private var storedURL: URL?

var url: URL? {
get {
lock.lock()
defer { lock.unlock() }
return storedURL
}
set {
lock.lock()
storedURL = newValue
lock.unlock()
}
}
}
Loading