IJPL-251781 read channel lines in chunks and strip the line terminator

`linesImpl` allocated a 1-byte buffer and issued one `receive` per byte. Each call
costs a coroutine dispatch plus a read syscall, so read cost scaled with payload size
rather than with the number of reads: measured at a flat ~22 KiB/s, where a 1 MB line
took 45.8 s and 1,049,004 `receive` calls. Reading agent stdout for AI Chat is the case
that made it visible. It now reads into a buffer and splits it on the line feeds it
holds.

Reading ahead means a read may take bytes past the line it completes. They are lost if
collection stops early, and the channel must not be read by other means afterwards; the
KDoc says so, and points at `PeekableEelReceiveChannel.readLine` for a caller that stops
after a handshake and needs the remainder back. That KDoc also loses a claim it should
never have made -- that CR and CRLF are separators "much like `BufferedReader`".

Which is now true, because the framing changed as well. `lines` kept the `\n` in every
emitted line and, because the end of the stream ended a line of its own, a stream ending
on `\n` finished with an empty string. No line reader on the JVM behaves that way:
`BufferedReader.readLine`, `Files.lines`, `Reader.readLines`, `Scanner.nextLine` and
Java's `String.lines` all strip the terminator and none of them reports a trailing empty
line -- nor does this module's own `readLine`, which the two APIs now agree with. Only
Kotlin's `CharSequence.lines` keeps it, and that one is a plain split rather than a
reader. A `\r` directly in front of the `\n` goes with it, so CRLF input reads the same
as LF input. A lone `\r` still does not end a line: command line tools use it to redraw
a line in place, and splitting there would tear that output apart.

Callers were audited one by one. Two relied on the terminator and now add it back where
they need it: `runAmperInBuildView` writes into a build console that appends the text as
is, and the `stdin echo` test feeds its output to a reader that needs a terminator to see
a line. `SolanaProcessUtils` gets a fix for free -- it collects lines with `appendLine`,
which until now doubled every newline and added a blank line at the end. `EmbarkCliExecutor`
loses a `trimEnd('\r', '\n')` that has nothing left to trim. The ACP transports feed the
lines to the library's `StdioTransport`, whose own line flow strips the terminator the
same way and whose output side appends the `\n` by hand, so the input side is now
symmetric with it and the empty line no longer reaches the JSON parser.

A caller can also choose the read buffer size now, since that is both how much a single
`receive` takes and how far the flow reads past the line it emits. One that has to keep
the read-ahead small, or that knows its lines are far longer than the 8 KiB default, can
say so; the existing callers keep the default. It is a separate overload rather than a
parameter with a default value: a default would change the signature of the existing
function and break every already-compiled caller, and `@JvmOverloads` does not help since
the synthesized overload is the one that changes.
https://kotlinlang.org/docs/api-guidelines-backward-compatibility.html#avoid-adding-arguments-to-existing-api-functions

The new framing test pins the shape explicitly, which the existing `testLines` could not
see because it trims every line before comparing.

GitOrigin-RevId: 5b21d2579e768b28b183ad9fa3893b54d345d063
This commit is contained in:
Alex Plate
2026-08-04 17:56:41 +00:00
committed by intellij-monorepo-bot
parent 4b4f2a0d85
commit 83ec867866
3 changed files with 119 additions and 19 deletions
@@ -79,9 +79,16 @@ fun CoroutineScope.consumeReceiveChannelAsKotlin(receiveChannel: EelReceiveChann
/**
* Collect data from channel line-by-line using [charset] to convert bytes to chars.
* Much like [java.io.BufferedReader], we consider CR or CRLF as a new line chars.
* This API might be slow (as it reads one byte per time) so you might prefer to [readAllBytes] first, then decode it and split by lines.
* However, for interactive source you can't read till the end, so you use this api.
* A line ends at `\n`, and neither it nor a `\r` in front of it is part of the emitted line, much like
* [java.io.BufferedReader.readLine]. A lone `\r` does not end a line: command line tools use it to
* redraw a line in place, and splitting there would tear that output apart.
* The end of the stream ends the last line, and emits nothing when there is no unterminated line left.
*
* For a non-interactive source you might prefer to [readAllBytes] first, then decode it and split by lines.
*
* This reads ahead, so it fits a channel that is read to its end. To read only the first lines -- a
* handshake, a header -- and leave the rest readable, use [com.intellij.platform.eel.channels.readLine]
* of [com.intellij.platform.eel.channels.PeekableEelReceiveChannel], which puts the remainder back.
*
* As soon as channel gets closed -- flow finishes.
* ```kotlin
@@ -99,7 +106,17 @@ fun CoroutineScope.consumeReceiveChannelAsKotlin(receiveChannel: EelReceiveChann
*/
@ApiStatus.Internal
fun EelReceiveChannel.lines(charset: Charset): Flow<String> =
linesImpl(charset)
linesImpl(charset, DEFAULT_BUFFER_SIZE)
/**
* Please read the documentation for the other overload of [lines].
*
* [bufferSize] is how many bytes a single `receive` may take, so it bounds the read-ahead described
* there. It does not limit the length of an emitted line.
*/
@ApiStatus.Internal
fun EelReceiveChannel.lines(charset: Charset, bufferSize: Int): Flow<String> =
linesImpl(charset, bufferSize)
@ApiStatus.Internal
fun EelReceiveChannel.lines(): Flow<String> = lines(Charset.defaultCharset())
@@ -310,26 +310,57 @@ internal fun Socket.consumeAsEelChannelImpl(): EelReceiveChannel =
internal fun Socket.asEelChannelImpl(): EelSendChannel =
channel?.asEelChannel() ?: outputStream.asEelChannel()
internal fun EelReceiveChannel.linesImpl(charset: Charset): Flow<String> = flow {
val tmpBuffer = ByteBuffer.allocate(1)
var result = ByteArrayOutputStream()
while (true) {
tmpBuffer.rewind()
suspend fun emitBuffer() {
emit(charset.decode(ByteBuffer.wrap(result.toByteArray())).toString())
}
private const val LINE_FEED = '\n'.code.toByte()
private const val CARRIAGE_RETURN = '\r'.code.toByte()
val r = receive(tmpBuffer)
internal fun EelReceiveChannel.linesImpl(charset: Charset, bufferSize: Int): Flow<String> = flow {
val chunk = ByteBuffer.allocate(bufferSize)
var result = ByteArrayOutputStream()
// The terminator is dropped by decoding fewer bytes rather than by trimming the decoded line, which
// would allocate a second string for every line that has one.
suspend fun emitBuffer() {
val bytes = result.toByteArray()
var size = bytes.size
if (size > 0 && bytes[size - 1] == LINE_FEED) {
size--
if (size > 0 && bytes[size - 1] == CARRIAGE_RETURN) {
size--
}
}
emit(charset.decode(ByteBuffer.wrap(bytes, 0, size)).toString())
}
while (true) {
chunk.clear()
val r = receive(chunk)
if (r == ReadResult.EOF) {
emitBuffer()
// A stream ending right after a line feed has no unterminated last line to report.
if (result.size() > 0) {
emitBuffer()
}
return@flow
}
val b = tmpBuffer.flip().get().toInt()
result.write(b)
if (b == 10) {
emitBuffer()
result = ByteArrayOutputStream()
chunk.flip()
while (chunk.hasRemaining()) {
val start = chunk.position()
val limit = chunk.limit()
var end = limit
var lineComplete = false
for (i in start until limit) {
if (chunk.get(i) == LINE_FEED) {
end = i + 1
lineComplete = true
break
}
}
result.write(chunk.array(), chunk.arrayOffset() + start, end - start)
chunk.position(end)
if (lineComplete) {
emitBuffer()
result = ByteArrayOutputStream()
}
}
}
}
@@ -700,6 +700,53 @@ class EelChannelToolsTest {
}
Assertions.assertArrayEquals(lines.toTypedArray(), result.toTypedArray(), "Wrong lines collected")
}
/**
* The terminator is not part of the line, and a stream ending on one has no last line to report.
*/
@CartesianTest
fun testLinesFraming(@IntRangeSource(from = 1, to = 4) blockSize: Int): Unit = timeoutRunBlocking {
suspend fun collect(text: String): List<String> {
val channel = ByteArrayInputStreamLimited(text.encodeToByteArray(), blockSize).consumeAsEelChannel()
val result = mutableListOf<String>()
channel.lines(Charsets.UTF_8).collect(result::add)
return result
}
assertEquals(listOf("one", "two"), collect("one\ntwo\n"), "Terminated stream")
assertEquals(listOf("one", "two"), collect("one\ntwo"), "Unterminated last line")
assertEquals(listOf("a", "", "b"), collect("a\n\nb\n"), "Empty line in the middle")
assertEquals(listOf<String>(), collect(""), "Empty stream")
assertEquals(listOf(""), collect("\n"), "Nothing but a terminator")
assertEquals(listOf("a", "b"), collect("a\r\nb\r\n"), "CRLF, possibly split between reads")
assertEquals(listOf("50%\r100%"), collect("50%\r100%\n"), "A lone CR does not end a line")
assertEquals(listOf("a\r"), collect("a\r"), "A CR is only dropped as the first half of a CRLF")
}
/**
* Guards against reading byte by byte, which makes the read cost scale with the payload size.
*/
@Test
fun testLinesReadInChunks(): Unit = timeoutRunBlocking {
val bytes = ("a".repeat(512 * 1024) + "\n").encodeToByteArray()
val stream = ByteArrayInputStreamLimited(bytes, bytes.size)
stream.consumeAsEelChannel().lines(Charsets.UTF_8).collect { }
assertTrue(stream.reads.get() < bytes.size / 1024, "One line of ${bytes.size} bytes took ${stream.reads} reads")
}
/**
* The buffer size bounds what a single read takes from the channel, not the length of a line.
*/
@Test
fun testLinesBufferSize(): Unit = timeoutRunBlocking {
val line = "a".repeat(4096)
val bytes = "$line\n".encodeToByteArray()
val stream = ByteArrayInputStreamLimited(bytes, bytes.size)
val result = mutableListOf<String>()
stream.consumeAsEelChannel().lines(Charsets.UTF_8, 64).collect(result::add)
assertEquals(listOf(line), result, "Wrong lines collected")
assertTrue(stream.reads.get() >= bytes.size / 64, "${bytes.size} bytes in 64-byte reads took ${stream.reads} reads")
}
}
/**
@@ -733,8 +780,13 @@ private class ByteArrayChannelLimited(private val blockSize: Int) : WritableByte
*/
private class ByteArrayInputStreamLimited(data: ByteArray, private val blockSize: Int) : InputStream() {
private val iterator = data.iterator()
/** One read per `receive`, so this also counts the reads of the channel above the stream. */
val reads: AtomicInteger = AtomicInteger()
override fun read(): Int = if (iterator.hasNext()) iterator.next().toInt() else -1
override fun read(b: ByteArray, off: Int, len: Int): Int {
reads.incrementAndGet()
if (!iterator.hasNext()) return -1
val bytesToWrite = min(len, blockSize)
var i = 0