From 30f56a7bcac1ce97eee06ddc7963c899eb8000d5 Mon Sep 17 00:00:00 2001 From: Oliver Drobnik Date: Wed, 23 Sep 2026 12:17:07 +0200 Subject: [PATCH 1/3] Framing: an optional byte limit, and a way to report exceeding it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `LineFraming.push` appended to an unbounded buffer and only emitted on a newline, so a peer that sent a huge line — or never sent one at all — grew it without limit. There was also nowhere to report that: `push` returned `[Data]` with no failure path, so a transport could not distinguish "no messages yet" from "this stream is unusable". - `MessageFraming.push` is now `throws`, and `FramingError.messageTooLarge` carries the limit and how much had accumulated. - `LineFraming(maxBytes:)` and `ContentLengthFraming(maxBytes:)` default to `0`, unlimited — existing behaviour is unchanged until a caller opts in. Policy (what the limit should be, which environment variable names it) belongs to the application, not here. - `ContentLengthFraming` refuses on the *declared* length, before a byte of the body is buffered; it also caps headers that never reach a separator. `LineFraming` cannot know a size in advance, so it checks a completed line and an unterminated remainder that has already passed the limit. - A failure drops the buffer: there is no boundary left to resynchronise on, so the framing does not keep answering with the same error. Transports answer by finishing their inbound stream with the error, which is already an `AsyncThrowingStream` — the stdio pump propagates through its existing `catch`, and `startFramedReaderThread` gained an `onFailure` callback that its three callers wire to `finish(throwing:)`. Source-breaking for out-of-tree conformers and callers: a conformer adds `throws` without behaviour change, a caller adds `try`. Co-Authored-By: Claude Opus 5 --- Sources/JSONRPCStdio/ProcessTransport.swift | 3 +- .../StdioMessageTransport.swift | 5 +- Sources/JSONRPCTCP/TCPClientTransport.swift | 3 +- Sources/JSONRPCWire/FramedReaderThread.swift | 14 ++- Sources/JSONRPCWire/MessageFraming.swift | 75 ++++++++++-- Tests/JSONRPCWireTests/JSONRPCWireTests.swift | 110 +++++++++++++----- 6 files changed, 166 insertions(+), 44 deletions(-) diff --git a/Sources/JSONRPCStdio/ProcessTransport.swift b/Sources/JSONRPCStdio/ProcessTransport.swift index 4ae9dcc..49026ee 100644 --- a/Sources/JSONRPCStdio/ProcessTransport.swift +++ b/Sources/JSONRPCStdio/ProcessTransport.swift @@ -125,7 +125,8 @@ public final class ProcessTransport: JSONRPCMessageTran continuation.yield(message) } }, - onEOF: { continuation.finish() }) + onEOF: { continuation.finish() }, + onFailure: { continuation.finish(throwing: $0) }) continuation.onTermination = { [weak self] _ in self?.close() diff --git a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift index 6ab134e..a93343c 100644 --- a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift +++ b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift @@ -142,7 +142,7 @@ public final class StdioTransport: 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)) { + for body in try decoder.push(Data(bytes)) { for message in (try? JSONRPCMessage.decodeMessages(from: body)) ?? [] { inbound.yield(message) } @@ -188,7 +188,8 @@ public final class StdioTransport: JSONRPCMessageTransp inbound.yield(message) } }, - onEOF: { inbound.finish() }) + onEOF: { inbound.finish() }, + onFailure: { inbound.finish(throwing: $0) }) // Writer: a single task drains outbound to our stdout (no lock needed). return Task { diff --git a/Sources/JSONRPCTCP/TCPClientTransport.swift b/Sources/JSONRPCTCP/TCPClientTransport.swift index 26a0590..7972a1a 100644 --- a/Sources/JSONRPCTCP/TCPClientTransport.swift +++ b/Sources/JSONRPCTCP/TCPClientTransport.swift @@ -162,7 +162,8 @@ public final class TCPClientTransport: JSONRPCMessageTr continuation.yield(message) } }, - onEOF: { continuation.finish() }) + onEOF: { continuation.finish() }, + onFailure: { continuation.finish(throwing: $0) }) continuation.onTermination = { [weak self] _ in self?.close() diff --git a/Sources/JSONRPCWire/FramedReaderThread.swift b/Sources/JSONRPCWire/FramedReaderThread.swift index a817518..662e204 100644 --- a/Sources/JSONRPCWire/FramedReaderThread.swift +++ b/Sources/JSONRPCWire/FramedReaderThread.swift @@ -23,15 +23,23 @@ package func startFramedReaderThread( framing: some MessageFraming, readChunk: @escaping @Sendable () -> Data, onBody: @escaping @Sendable (Data) -> Void, - onEOF: @escaping @Sendable () -> Void + onEOF: @escaping @Sendable () -> Void, + onFailure: @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 { + for body in try decoder.push(chunk) { + onBody(body) + } + } catch { + // A framing failure leaves no boundary to resynchronise on, so the + // read ends here and the caller fails its stream rather than looping. + onFailure(error) + return } } onEOF() diff --git a/Sources/JSONRPCWire/MessageFraming.swift b/Sources/JSONRPCWire/MessageFraming.swift index b416ea7..708361c 100644 --- a/Sources/JSONRPCWire/MessageFraming.swift +++ b/Sources/JSONRPCWire/MessageFraming.swift @@ -16,7 +16,28 @@ public protocol MessageFraming: Sendable { 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] + /// + /// Throws ``FramingError`` when the bytes cannot yield a message — 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. + mutating func push(_ bytes: Data) throws -> [Data] +} + +/// 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: \r\n\r\n`. @@ -26,8 +47,15 @@ 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) @@ -35,14 +63,24 @@ public struct ContentLengthFraming: MessageFraming { return out } - public mutating func push(_ bytes: Data) -> [Data] { + public mutating func push(_ bytes: Data) throws -> [Data] { buffer.append(bytes) var messages: [Data] = [] - while let message = next() { messages.append(message) } + while let message = try next() { messages.append(message) } + // Headers with no separator in sight would otherwise buffer without bound. + if maxBytes > 0, expectedLength == nil, buffer.count > maxBytes { + throw drop(pending: buffer.count) + } return messages } - private mutating func next() -> Data? { + private mutating func drop(pending: Int) -> FramingError { + buffer.removeAll(keepingCapacity: false) + expectedLength = nil + return FramingError.messageTooLarge(limit: maxBytes, pending: pending) + } + + private mutating func next() throws -> Data? { if let length = expectedLength { guard buffer.count >= length else { return nil } let body = Data(buffer.prefix(length)) @@ -54,9 +92,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) @@ -82,8 +122,17 @@ public struct ContentLengthFraming: MessageFraming { /// Newline-delimited JSON framing (ACP and MCP-over-stdio): `\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 @@ -91,14 +140,22 @@ public struct LineFraming: MessageFraming { return out } - public mutating func push(_ bytes: Data) -> [Data] { + public mutating func push(_ bytes: Data) throws -> [Data] { 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 maxBytes > 0, size > maxBytes { throw drop(pending: size) } if !line.isEmpty { messages.append(Data(line)) } } + if maxBytes > 0, buffer.count > maxBytes { throw drop(pending: buffer.count) } return messages } + + private mutating func drop(pending: Int) -> FramingError { + buffer.removeAll(keepingCapacity: false) + return FramingError.messageTooLarge(limit: maxBytes, pending: pending) + } } diff --git a/Tests/JSONRPCWireTests/JSONRPCWireTests.swift b/Tests/JSONRPCWireTests/JSONRPCWireTests.swift index 02d94de..5a15298 100644 --- a/Tests/JSONRPCWireTests/JSONRPCWireTests.swift +++ b/Tests/JSONRPCWireTests/JSONRPCWireTests.swift @@ -7,40 +7,40 @@ private func text(_ data: Data) -> String? { String(data: data, encoding: .utf8) // MARK: - ContentLengthFraming -@Test func contentLengthRoundTripsOneMessage() { +@Test func contentLengthRoundTripsOneMessage() throws { var framing = ContentLengthFraming() - let out = framing.push(framing.frame(body(#"{"jsonrpc":"2.0","id":1}"#))) + let out = try framing.push(framing.frame(body(#"{"jsonrpc":"2.0","id":1}"#))) #expect(out.count == 1) #expect(text(out[0]) == #"{"jsonrpc":"2.0","id":1}"#) } -@Test func contentLengthSplitsTwoMessagesInOneChunk() { +@Test func contentLengthSplitsTwoMessagesInOneChunk() throws { var framing = ContentLengthFraming() let chunk = framing.frame(body(#"{"a":1}"#)) + framing.frame(body(#"{"b":2}"#)) - let out = framing.push(chunk) + let out = try framing.push(chunk) #expect(out.count == 2) #expect(text(out[1]) == #"{"b":2}"#) } -@Test func contentLengthReassemblesAcrossChunks() { +@Test func contentLengthReassemblesAcrossChunks() throws { var framing = ContentLengthFraming() var emitted: [Data] = [] - for byte in framing.frame(body(#"{"hello":"world"}"#)) { emitted += framing.push(Data([byte])) } + for byte in framing.frame(body(#"{"hello":"world"}"#)) { emitted += try framing.push(Data([byte])) } #expect(emitted.count == 1) #expect(text(emitted[0]) == #"{"hello":"world"}"#) } -@Test func contentLengthCountsBytesNotCharacters() { +@Test func contentLengthCountsBytesNotCharacters() throws { var framing = ContentLengthFraming() - let out = framing.push(framing.frame(body(#"{"v":"café"}"#))) // 5 UTF-8 bytes, 4 chars + let out = try framing.push(framing.frame(body(#"{"v":"café"}"#))) // 5 UTF-8 bytes, 4 chars #expect(out.count == 1) #expect(text(out[0]) == #"{"v":"café"}"#) } -@Test func contentLengthRejectsNegativeLength() { +@Test func contentLengthRejectsNegativeLength() throws { // A malformed `Content-Length: -1` must be dropped, not used as a frame size. var framing = ContentLengthFraming() - let out = framing.push(Data("Content-Length: -1\r\n\r\n{}".utf8)) + let out = try framing.push(Data("Content-Length: -1\r\n\r\n{}".utf8)) #expect(out.isEmpty) } @@ -50,74 +50,128 @@ private func text(_ data: Data) -> String? { String(data: data, encoding: .utf8) #expect(LineFraming().frame(body(#"{"id":1}"#)).last == 0x0A) } -@Test func lineFramingSplitsMultipleLines() { +@Test func lineFramingSplitsMultipleLines() throws { var framing = LineFraming() let chunk = framing.frame(body(#"{"a":1}"#)) + framing.frame(body(#"{"b":2}"#)) - let out = framing.push(chunk) + let out = try framing.push(chunk) #expect(out.count == 2) #expect(text(out[0]) == #"{"a":1}"#) } -@Test func lineFramingReassemblesAcrossChunks() { +@Test func lineFramingReassemblesAcrossChunks() throws { var framing = LineFraming() var emitted: [Data] = [] - for byte in framing.frame(body(#"{"x":42}"#)) { emitted += framing.push(Data([byte])) } + for byte in framing.frame(body(#"{"x":42}"#)) { emitted += try framing.push(Data([byte])) } #expect(emitted.count == 1) #expect(text(emitted[0]) == #"{"x":42}"#) } // MARK: - SSEEventDecoder -@Test func sseDecodesOneEvent() { +@Test func sseDecodesOneEvent() throws { var decoder = SSEEventDecoder() - let out = decoder.push(body("data: {\"id\":1}\n\n")) + let out = try decoder.push(body("data: {\"id\":1}\n\n")) #expect(out.count == 1) #expect(text(out[0]) == "{\"id\":1}") } -@Test func sseIgnoresCommentsAndNonDataFields() { +@Test func sseIgnoresCommentsAndNonDataFields() throws { var decoder = SSEEventDecoder() - let out = decoder.push(body(": keep-alive\nevent: message\nid: 7\ndata: {\"x\":1}\n\n")) + let out = try decoder.push(body(": keep-alive\nevent: message\nid: 7\ndata: {\"x\":1}\n\n")) #expect(out.count == 1) #expect(text(out[0]) == "{\"x\":1}") } -@Test func sseDecodesTwoEventsInOneChunk() { +@Test func sseDecodesTwoEventsInOneChunk() throws { var decoder = SSEEventDecoder() - let out = decoder.push(body("data: {\"a\":1}\n\ndata: {\"b\":2}\n\n")) + let out = try decoder.push(body("data: {\"a\":1}\n\ndata: {\"b\":2}\n\n")) #expect(out.count == 2) #expect(text(out[1]) == "{\"b\":2}") } -@Test func sseToleratesCRLF() { +@Test func sseToleratesCRLF() throws { var decoder = SSEEventDecoder() - let out = decoder.push(body("data: {\"id\":5}\r\n\r\n")) + let out = try decoder.push(body("data: {\"id\":5}\r\n\r\n")) #expect(out.count == 1) #expect(text(out[0]) == "{\"id\":5}") } -@Test func sseJoinsMultipleDataLines() { +@Test func sseJoinsMultipleDataLines() throws { // Doc-promised: multiple `data:` lines of one event join with "\n". var decoder = SSEEventDecoder() - let out = decoder.push(body("data: {\"a\":\ndata: 1}\n\n")) + let out = try decoder.push(body("data: {\"a\":\ndata: 1}\n\n")) #expect(out.count == 1) #expect(text(out[0]) == "{\"a\":\n1}") } -@Test func sseReassemblesAcrossChunks() { +@Test func sseReassemblesAcrossChunks() throws { // The SSE path is the one that actually sees arbitrary network chunking. var decoder = SSEEventDecoder() var emitted: [Data] = [] - for byte in body("data: {\"id\":9}\n\n") { emitted += decoder.push(Data([byte])) } + for byte in body("data: {\"id\":9}\n\n") { emitted += try decoder.push(Data([byte])) } #expect(emitted.count == 1) #expect(text(emitted[0]) == "{\"id\":9}") } -@Test func sseTreatsBareFieldNameAsEmptyValue() { +@Test func sseTreatsBareFieldNameAsEmptyValue() throws { // Doc-promised: a line with no colon is a field name with an empty value, so // a bare `data` line contributes an empty payload — still a dispatched event. var decoder = SSEEventDecoder() - let out = decoder.push(body("data\n\n")) + let out = try decoder.push(body("data\n\n")) #expect(out.count == 1) #expect(text(out[0]) == "") } + +// MARK: - Byte limits + +@Test func lineFramingIsUnlimitedByDefault() throws { + var framing = LineFraming() + let huge = body(String(repeating: "x", count: 1 << 20)) + #expect(try framing.push(framing.frame(huge)).count == 1) +} + +@Test func lineFramingRejectsAnOversizedMessage() throws { + var framing = LineFraming(maxBytes: 16) + #expect(throws: FramingError.messageTooLarge(limit: 16, pending: 32)) { + try framing.push(framing.frame(body(String(repeating: "x", count: 32)))) + } + // The buffer is dropped, so the framing does not keep answering with the failure. + #expect(try framing.push(framing.frame(body("{}"))).count == 1) +} + +/// The case a limit exists for: a peer that never terminates its line would otherwise +/// grow the buffer without bound. +@Test func lineFramingRejectsAnUnterminatedFlood() throws { + var framing = LineFraming(maxBytes: 8) + #expect(try framing.push(body("12345")).isEmpty) + #expect(throws: FramingError.messageTooLarge(limit: 8, pending: 10)) { + try framing.push(body("67890")) + } +} + +@Test func lineFramingAcceptsAMessageExactlyAtTheLimit() throws { + var framing = LineFraming(maxBytes: 4) + let out = try framing.push(framing.frame(body("abcd"))) + #expect(text(out.first ?? Data()) == "abcd") +} + +@Test func contentLengthRejectsByTheDeclaredLengthBeforeTheBody() throws { + var framing = ContentLengthFraming(maxBytes: 16) + // Only the header is fed: the length alone is enough to refuse it. + #expect(throws: FramingError.messageTooLarge(limit: 16, pending: 4096)) { + try framing.push(body("Content-Length: 4096\r\n\r\n")) + } +} + +@Test func contentLengthRejectsHeadersThatNeverEnd() throws { + var framing = ContentLengthFraming(maxBytes: 8) + #expect(throws: (any Error).self) { + try framing.push(body(String(repeating: "X-Pad: 1\r\n", count: 4))) + } +} + +@Test func contentLengthIsUnlimitedByDefault() throws { + var framing = ContentLengthFraming() + let big = body(String(repeating: "y", count: 1 << 16)) + #expect(try framing.push(framing.frame(big)).count == 1) +} From 9334ed2138cdb5289a8cb9fb504e887511912cba Mon Sep 17 00:00:00 2001 From: Oliver Drobnik Date: Wed, 23 Sep 2026 12:18:20 +0200 Subject: [PATCH 2/3] Reader thread: one finish callback instead of two SwiftLint's parameter-count limit pushed back on `onEOF` + `onFailure`, and it was right to: `AsyncThrowingStream.Continuation.finish(throwing:)` already takes an optional error, so a single `onFinish: (any Error?) -> Void` covers both endings and every caller collapses to onFinish: { continuation.finish(throwing: $0) } which is shorter than what was there before this branch. Co-Authored-By: Claude Opus 5 --- Sources/JSONRPCStdio/ProcessTransport.swift | 3 +-- Sources/JSONRPCSubprocess/StdioMessageTransport.swift | 3 +-- Sources/JSONRPCTCP/TCPClientTransport.swift | 3 +-- Sources/JSONRPCWire/FramedReaderThread.swift | 7 +++---- 4 files changed, 6 insertions(+), 10 deletions(-) diff --git a/Sources/JSONRPCStdio/ProcessTransport.swift b/Sources/JSONRPCStdio/ProcessTransport.swift index 49026ee..c3be9d1 100644 --- a/Sources/JSONRPCStdio/ProcessTransport.swift +++ b/Sources/JSONRPCStdio/ProcessTransport.swift @@ -125,8 +125,7 @@ public final class ProcessTransport: JSONRPCMessageTran continuation.yield(message) } }, - onEOF: { continuation.finish() }, - onFailure: { continuation.finish(throwing: $0) }) + onFinish: { continuation.finish(throwing: $0) }) continuation.onTermination = { [weak self] _ in self?.close() diff --git a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift index a93343c..0a9c461 100644 --- a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift +++ b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift @@ -188,8 +188,7 @@ public final class StdioTransport: JSONRPCMessageTransp inbound.yield(message) } }, - onEOF: { inbound.finish() }, - onFailure: { inbound.finish(throwing: $0) }) + onFinish: { inbound.finish(throwing: $0) }) // Writer: a single task drains outbound to our stdout (no lock needed). return Task { diff --git a/Sources/JSONRPCTCP/TCPClientTransport.swift b/Sources/JSONRPCTCP/TCPClientTransport.swift index 7972a1a..8543d97 100644 --- a/Sources/JSONRPCTCP/TCPClientTransport.swift +++ b/Sources/JSONRPCTCP/TCPClientTransport.swift @@ -162,8 +162,7 @@ public final class TCPClientTransport: JSONRPCMessageTr continuation.yield(message) } }, - onEOF: { continuation.finish() }, - onFailure: { continuation.finish(throwing: $0) }) + onFinish: { continuation.finish(throwing: $0) }) continuation.onTermination = { [weak self] _ in self?.close() diff --git a/Sources/JSONRPCWire/FramedReaderThread.swift b/Sources/JSONRPCWire/FramedReaderThread.swift index 662e204..8619601 100644 --- a/Sources/JSONRPCWire/FramedReaderThread.swift +++ b/Sources/JSONRPCWire/FramedReaderThread.swift @@ -23,8 +23,7 @@ package func startFramedReaderThread( framing: some MessageFraming, readChunk: @escaping @Sendable () -> Data, onBody: @escaping @Sendable (Data) -> Void, - onEOF: @escaping @Sendable () -> Void, - onFailure: @escaping @Sendable (any Error) -> Void + onFinish: @escaping @Sendable (any Error?) -> Void ) { let thread = Thread { var decoder = framing @@ -38,11 +37,11 @@ package func startFramedReaderThread( } catch { // A framing failure leaves no boundary to resynchronise on, so the // read ends here and the caller fails its stream rather than looping. - onFailure(error) + onFinish(error) return } } - onEOF() + onFinish(nil) } thread.name = name thread.stackSize = 4 << 20 From d162dc76c3ae603ea07c2341744bd3d968b318b6 Mon Sep 17 00:00:00 2001 From: Oliver Drobnik Date: Wed, 23 Sep 2026 13:50:07 +0200 Subject: [PATCH 3/3] Framing: deliver as you decode, and cap headers on their own limit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI plus two review findings. **`(any Error)?`.** `any Error?` parses as `any (Error?)` and is rejected by the toolchains CI builds with, though the one here accepted it. Spelled properly now. **Completed messages are no longer lost with the failure.** One read can carry a valid message *and* an oversized one; returning `[Data]` meant the caller discarded the valid message along with the throw, so a good response could be lost because of what followed it in the same buffer. `push` now takes an `emit` closure and hands over each body as it is decoded — delivery and failure both happen, in that order, by construction rather than by care. The array-returning form stays as an extension for callers that treat any framing failure as fatal, with its loss documented. **A body limit is not a header limit.** `ContentLengthFraming` applied `maxBytes` to an incomplete header, so `maxBytes: 16` accepted a 2-byte frame in one push and rejected the same frame split after `Content-Length: 2` — and a transport read may split anywhere. Headers now have their own generous cap (8 KiB), independent of the body limit, which still stops a peer that never sends a separator. Tests: a good message emitted before an oversized one throws (both framings), a header split under a small body limit accepted, and the header cap tripped only by an actually unbounded header. Co-Authored-By: Claude Opus 5 --- .../StdioMessageTransport.swift | 2 +- Sources/JSONRPCWire/FramedReaderThread.swift | 6 +-- Sources/JSONRPCWire/MessageFraming.swift | 54 +++++++++++++------ Tests/JSONRPCWireTests/JSONRPCWireTests.swift | 45 +++++++++++++++- 4 files changed, 84 insertions(+), 23 deletions(-) diff --git a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift index 0a9c461..5d086b9 100644 --- a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift +++ b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift @@ -142,7 +142,7 @@ public final class StdioTransport: JSONRPCMessageTransp var decoder = framing // value copy → fresh buffer for try await buffer in execution.standardOutput { let bytes = buffer.withUnsafeBytes { Array($0) } - for body in try decoder.push(Data(bytes)) { + try decoder.push(Data(bytes)) { body in for message in (try? JSONRPCMessage.decodeMessages(from: body)) ?? [] { inbound.yield(message) } diff --git a/Sources/JSONRPCWire/FramedReaderThread.swift b/Sources/JSONRPCWire/FramedReaderThread.swift index 8619601..17d7629 100644 --- a/Sources/JSONRPCWire/FramedReaderThread.swift +++ b/Sources/JSONRPCWire/FramedReaderThread.swift @@ -23,7 +23,7 @@ package func startFramedReaderThread( framing: some MessageFraming, readChunk: @escaping @Sendable () -> Data, onBody: @escaping @Sendable (Data) -> Void, - onFinish: @escaping @Sendable (any Error?) -> Void + onFinish: @escaping @Sendable ((any Error)?) -> Void ) { let thread = Thread { var decoder = framing @@ -31,9 +31,7 @@ package func startFramedReaderThread( let chunk = readChunk() if chunk.isEmpty { break } do { - for body in try decoder.push(chunk) { - onBody(body) - } + 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. diff --git a/Sources/JSONRPCWire/MessageFraming.swift b/Sources/JSONRPCWire/MessageFraming.swift index 708361c..afd7cdd 100644 --- a/Sources/JSONRPCWire/MessageFraming.swift +++ b/Sources/JSONRPCWire/MessageFraming.swift @@ -14,14 +14,29 @@ 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. + /// 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 a message — 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. - mutating func push(_ bytes: Data) throws -> [Data] + /// 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. @@ -63,17 +78,22 @@ public struct ContentLengthFraming: MessageFraming { return out } - public mutating func push(_ bytes: Data) throws -> [Data] { + public mutating func push(_ bytes: Data, emit: (Data) -> Void) throws { buffer.append(bytes) - var messages: [Data] = [] - while let message = try next() { messages.append(message) } - // Headers with no separator in sight would otherwise buffer without bound. - if maxBytes > 0, expectedLength == nil, buffer.count > maxBytes { + 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) } - return messages } + /// 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 @@ -140,18 +160,18 @@ public struct LineFraming: MessageFraming { return out } - public mutating func push(_ bytes: Data) throws -> [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) + // 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) } - if !line.isEmpty { messages.append(Data(line)) } + if !line.isEmpty { emit(Data(line)) } } if maxBytes > 0, buffer.count > maxBytes { throw drop(pending: buffer.count) } - return messages } private mutating func drop(pending: Int) -> FramingError { diff --git a/Tests/JSONRPCWireTests/JSONRPCWireTests.swift b/Tests/JSONRPCWireTests/JSONRPCWireTests.swift index 5a15298..e512e93 100644 --- a/Tests/JSONRPCWireTests/JSONRPCWireTests.swift +++ b/Tests/JSONRPCWireTests/JSONRPCWireTests.swift @@ -163,10 +163,15 @@ private func text(_ data: Data) -> String? { String(data: data, encoding: .utf8) } } +/// Headers are capped on their own, generous limit — not on `maxBytes`, which governs +/// a body — so this takes far more than a small body limit to trip. @Test func contentLengthRejectsHeadersThatNeverEnd() throws { var framing = ContentLengthFraming(maxBytes: 8) + // Well under the header cap: still fine, however small the body limit is. + #expect(try framing.push(body(String(repeating: "X-Pad: 1\r\n", count: 4))).isEmpty) + // Past it: a peer that never sends a separator cannot buffer without bound. #expect(throws: (any Error).self) { - try framing.push(body(String(repeating: "X-Pad: 1\r\n", count: 4))) + try framing.push(body(String(repeating: "X-Pad: 1\r\n", count: 1024))) } } @@ -175,3 +180,41 @@ private func text(_ data: Data) -> String? { String(data: data, encoding: .utf8) let big = body(String(repeating: "y", count: 1 << 16)) #expect(try framing.push(framing.frame(big)).count == 1) } + +/// A read can carry a good message and an oversized one together. The good one is +/// already delivered when the failure is reported — losing it because of what followed +/// it in the same buffer would be a bug in the framing, not in the peer. +@Test func lineFramingEmitsCompletedMessagesBeforeFailing() { + var framing = LineFraming(maxBytes: 8) + var emitted: [String] = [] + let chunk = framing.frame(body("{\"a\":1}")) + framing.frame(body(String(repeating: "x", count: 32))) + + #expect(throws: FramingError.messageTooLarge(limit: 8, pending: 32)) { + try framing.push(chunk) { emitted.append(text($0) ?? "") } + } + #expect(emitted == ["{\"a\":1}"]) +} + +@Test func contentLengthEmitsCompletedMessagesBeforeFailing() { + var framing = ContentLengthFraming(maxBytes: 8) + var emitted: [String] = [] + let chunk = framing.frame(body("{\"a\":1}")) + framing.frame(body(String(repeating: "y", count: 64))) + + #expect(throws: FramingError.messageTooLarge(limit: 8, pending: 64)) { + try framing.push(chunk) { emitted.append(text($0) ?? "") } + } + #expect(emitted == ["{\"a\":1}"]) +} + +/// A transport read can split anywhere, including inside a header that is longer than a +/// small body limit. The body limit must not be applied to the header. +@Test func contentLengthAcceptsAHeaderSplitUnderASmallLimit() throws { + var framing = ContentLengthFraming(maxBytes: 16) + let frame = framing.frame(body("{}")) + let split = frame.count - 1 + + #expect(try framing.push(frame.prefix(split)).isEmpty) + let out = try framing.push(frame.suffix(from: split)) + #expect(out.count == 1) + #expect(text(out[0]) == "{}") +}