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
2 changes: 1 addition & 1 deletion Sources/JSONRPCStdio/ProcessTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ public final class ProcessTransport<Framing: MessageFraming>: JSONRPCMessageTran
continuation.yield(message)
}
},
onEOF: { continuation.finish() })
onFinish: { continuation.finish(throwing: $0) })

continuation.onTermination = { [weak self] _ in
self?.close()
Expand Down
4 changes: 2 additions & 2 deletions Sources/JSONRPCSubprocess/StdioMessageTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ public final class StdioTransport<Framing: MessageFraming>: JSONRPCMessageTransp
var decoder = framing // value copy → fresh buffer
for try await buffer in execution.standardOutput {
let bytes = buffer.withUnsafeBytes { Array($0) }
for body in decoder.push(Data(bytes)) {
try decoder.push(Data(bytes)) { body in
for message in (try? JSONRPCMessage.decodeMessages(from: body)) ?? [] {
inbound.yield(message)
}
Expand Down Expand Up @@ -188,7 +188,7 @@ public final class StdioTransport<Framing: MessageFraming>: JSONRPCMessageTransp
inbound.yield(message)
}
},
onEOF: { inbound.finish() })
onFinish: { inbound.finish(throwing: $0) })

// Writer: a single task drains outbound to our stdout (no lock needed).
return Task {
Expand Down
2 changes: 1 addition & 1 deletion Sources/JSONRPCTCP/TCPClientTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ public final class TCPClientTransport<Framing: MessageFraming>: JSONRPCMessageTr
continuation.yield(message)
}
},
onEOF: { continuation.finish() })
onFinish: { continuation.finish(throwing: $0) })

continuation.onTermination = { [weak self] _ in
self?.close()
Expand Down
13 changes: 9 additions & 4 deletions Sources/JSONRPCWire/FramedReaderThread.swift
Original file line number Diff line number Diff line change
Expand Up @@ -23,18 +23,23 @@ package func startFramedReaderThread(
framing: some MessageFraming,
readChunk: @escaping @Sendable () -> Data,
onBody: @escaping @Sendable (Data) -> Void,
onEOF: @escaping @Sendable () -> Void
onFinish: @escaping @Sendable ((any Error)?) -> Void
) {
let thread = Thread {
var decoder = framing
while true {
let chunk = readChunk()
if chunk.isEmpty { break }
for body in decoder.push(chunk) {
onBody(body)
do {
try decoder.push(chunk) { onBody($0) }
} catch {
// A framing failure leaves no boundary to resynchronise on, so the
// read ends here and the caller fails its stream rather than looping.
onFinish(error)
return
}
}
onEOF()
onFinish(nil)
}
thread.name = name
thread.stackSize = 4 << 20
Expand Down
109 changes: 93 additions & 16 deletions Sources/JSONRPCWire/MessageFraming.swift
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,45 @@ import Foundation
public protocol MessageFraming: Sendable {
/// Wrap one message body for the wire (prepend a header / append a terminator).
func frame(_ body: Data) -> Data
/// Feed newly-read bytes; return every complete message body they now yield
/// (header/terminator stripped), buffering any partial remainder.
mutating func push(_ bytes: Data) -> [Data]
/// Feed newly-read bytes, handing each complete message body (header/terminator
/// stripped) to `emit` as it is decoded, and buffering any partial remainder.
///
/// Throws ``FramingError`` when the bytes cannot yield further messages — a peer
/// sending more than the configured limit, say. Such a stream cannot be
/// resynchronised, so a transport answers by finishing its inbound stream with the
/// error rather than reading on.
///
/// Delivery is a callback rather than a return value precisely because of that
/// throw: one read can carry a complete message *and* an oversized one, and the
/// complete message has already been emitted by the time the failure is reported.
mutating func push(_ bytes: Data, emit: (Data) -> Void) throws
}

extension MessageFraming {
/// Collects into an array instead of emitting. A failure discards whatever the same
/// call had already decoded, so transports should prefer the emitting form; this is
/// for callers that treat any framing failure as fatal.
public mutating func push(_ bytes: Data) throws -> [Data] {
var messages: [Data] = []
try push(bytes) { messages.append($0) }
return messages
}
}

/// Why a framing could not turn the bytes it was given into messages.
public enum FramingError: Error, Equatable, Sendable, CustomStringConvertible {
/// One message exceeded the framing's `maxBytes`. `pending` is how many bytes had
/// accumulated when the limit was passed — for an unterminated flood that is what
/// arrived before the buffer was dropped, not the message's real size, which is
/// unknowable.
case messageTooLarge(limit: Int, pending: Int)

public var description: String {
switch self {
case .messageTooLarge(let limit, let pending):
return "Message exceeded the \(limit)-byte framing limit (\(pending) bytes buffered)."
}
}
}

/// LSP base-protocol framing: `Content-Length: <n>\r\n\r\n<n bytes of JSON>`.
Expand All @@ -26,23 +62,45 @@ public protocol MessageFraming: Sendable {
public struct ContentLengthFraming: MessageFraming {
private var buffer = Data()
private var expectedLength: Int?
/// Largest single message to accept, in bytes; `0` (the default) is unlimited.
///
/// A declared `Content-Length` is checked *before* its body is buffered, so an
/// oversized message costs nothing but its header.
public let maxBytes: Int

public init() {}
public init(maxBytes: Int = 0) {
self.maxBytes = maxBytes
}

public func frame(_ body: Data) -> Data {
var out = Data("Content-Length: \(body.count)\r\n\r\n".utf8)
out.append(body)
return out
}

public mutating func push(_ bytes: Data) -> [Data] {
public mutating func push(_ bytes: Data, emit: (Data) -> Void) throws {
buffer.append(bytes)
var messages: [Data] = []
while let message = next() { messages.append(message) }
return messages
while let message = try next() { emit(message) }
// Headers with no separator in sight would otherwise buffer without bound. This
// is *not* `maxBytes`: that limits a message body, while a read can split
// anywhere — including part-way through a header longer than a small body limit.
if expectedLength == nil, buffer.count > Self.maxHeaderBytes {
throw drop(pending: buffer.count)
}
}

/// How much unterminated header to tolerate. Generous next to any real header, and
/// independent of `maxBytes` so a small body limit never rejects a legal header that
/// a read happened to split.
private static let maxHeaderBytes = 8 * 1024

private mutating func drop(pending: Int) -> FramingError {
buffer.removeAll(keepingCapacity: false)
expectedLength = nil
return FramingError.messageTooLarge(limit: maxBytes, pending: pending)
}

private mutating func next() -> Data? {
private mutating func next() throws -> Data? {
if let length = expectedLength {
guard buffer.count >= length else { return nil }
let body = Data(buffer.prefix(length))
Expand All @@ -54,9 +112,11 @@ public struct ContentLengthFraming: MessageFraming {
let headerBytes = buffer[buffer.startIndex ..< separator.lowerBound]
let length = Self.contentLength(in: headerBytes)
buffer.removeSubrange(buffer.startIndex ..< separator.upperBound)
guard let length else { return next() }
guard let length else { return try next() }
// Checked before the body arrives: the header already says how big it is.
if maxBytes > 0, length > maxBytes { throw drop(pending: length) }
expectedLength = length
return next()
return try next()
}

private static let headerSeparator = Data("\r\n\r\n".utf8)
Expand All @@ -82,23 +142,40 @@ public struct ContentLengthFraming: MessageFraming {
/// Newline-delimited JSON framing (ACP and MCP-over-stdio): `<json>\n`.
public struct LineFraming: MessageFraming {
private var buffer = Data()
/// Largest single message to accept, in bytes; `0` (the default) is unlimited.
///
/// Newline framing cannot know a message's size in advance, so the limit is
/// applied twice: to a completed line, and to an unterminated remainder that has
/// already passed it — which is what stops a peer that never sends a newline from
/// growing the buffer without bound.
public let maxBytes: Int

public init() {}
public init(maxBytes: Int = 0) {
self.maxBytes = maxBytes
}

public func frame(_ body: Data) -> Data {
var out = body
out.append(0x0A) // newline terminator
return out
}

public mutating func push(_ bytes: Data) -> [Data] {
public mutating func push(_ bytes: Data, emit: (Data) -> Void) throws {
buffer.append(bytes)
var messages: [Data] = []
while let newline = buffer.firstIndex(of: 0x0A) {
let line = buffer[buffer.startIndex ..< newline]
let size = line.count
buffer.removeSubrange(buffer.startIndex ... newline)
if !line.isEmpty { messages.append(Data(line)) }
// Emitted before the check on the *next* line, so a message that arrived in
// the same read as an oversized one is still delivered.
if maxBytes > 0, size > maxBytes { throw drop(pending: size) }
Comment thread
odrobnik marked this conversation as resolved.
if !line.isEmpty { emit(Data(line)) }
}
return messages
if maxBytes > 0, buffer.count > maxBytes { throw drop(pending: buffer.count) }
}

private mutating func drop(pending: Int) -> FramingError {
buffer.removeAll(keepingCapacity: false)
return FramingError.messageTooLarge(limit: maxBytes, pending: pending)
}
}
Loading
Loading