From cc4a1c2842fd46787643b604c4039768a80596fd Mon Sep 17 00:00:00 2001 From: Javid Khan Date: Tue, 6 Oct 2026 01:13:33 +0530 Subject: [PATCH 1/2] treat a lone CR as an SSE line terminator --- .../http/sse/impl/ByteSseEntityConsumer.java | 26 ++++++++++++------- .../http/sse/impl/SseEntityConsumer.java | 15 ++++++++--- .../http/sse/ByteSseEntityConsumerTest.java | 17 ++++++++++++ .../http/sse/SseEntityConsumerTest.java | 15 +++++++++++ 4 files changed, 60 insertions(+), 13 deletions(-) diff --git a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java index d3caa0867f..2d7ff361bb 100644 --- a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java +++ b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java @@ -48,7 +48,7 @@ *
  • Validates {@code Content-Type} equals {@code text/event-stream} * in {@link #streamStart(ContentType)}; otherwise throws {@link HttpException}.
  • *
  • Strips a UTF-8 BOM if present in the first chunk.
  • - *
  • Accepts LF and CRLF line endings; tolerates CRLF split across buffers.
  • + *
  • Accepts CR, LF and CRLF line endings per the SSE grammar; tolerates CRLF split across buffers.
  • *
  • Implements WHATWG SSE fields: {@code data}, {@code id}, {@code event}, {@code retry}. * Unknown fields and malformed {@code retry} values are ignored.
  • *
  • At end of stream, flushes any partially accumulated line and forces a final @@ -76,6 +76,7 @@ public final class ByteSseEntityConsumer extends AbstractBinAsyncEntityConsumer< // line accumulator private byte[] lineBuf = new byte[256]; private int lineLen = 0; + private boolean pendingCr = false; // event accumulator private final StringBuilder data = new StringBuilder(256); @@ -136,13 +137,20 @@ protected void data(final ByteBuffer src, final boolean endOfStream) { while (src.hasRemaining()) { final byte b = src.get(); if (b == LF) { - int len = lineLen; - if (len > 0 && lineBuf[len - 1] == CR) { - len--; + if (pendingCr) { + // LF completing a CRLF pair; the line was already emitted on the CR. + pendingCr = false; + } else { + handleLine(lineBuf, lineLen); + lineLen = 0; } - handleLine(lineBuf, len); + } else if (b == CR) { + // A lone CR is a line terminator per the SSE grammar (CR / LF / CRLF). + pendingCr = true; + handleLine(lineBuf, lineLen); lineLen = 0; } else { + pendingCr = false; appendByte(b); } } @@ -154,11 +162,7 @@ protected void data(final ByteBuffer src, final boolean endOfStream) { private void flushEndOfStream() { if (lineLen > 0) { - int len = lineLen; - if (lineBuf[len - 1] == CR) { - len--; - } - handleLine(lineBuf, len); + handleLine(lineBuf, lineLen); lineLen = 0; } handleLine(lineBuf, 0); @@ -182,6 +186,8 @@ protected Void generateContent() { @Override public void releaseResources() { lineBuf = new byte[0]; + lineLen = 0; + pendingCr = false; data.setLength(0); id = null; type = null; diff --git a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseEntityConsumer.java b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseEntityConsumer.java index 5b90d4e72a..3e8327b0bd 100644 --- a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseEntityConsumer.java +++ b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseEntityConsumer.java @@ -66,6 +66,7 @@ public final class SseEntityConsumer extends AbstractCharAsyncEntityConsumer 0 && partial.charAt(len - 1) == '\r') { - partial.setLength(len - 1); + if (pendingCr) { + // LF completing a CRLF pair; the line was already emitted on the CR. + pendingCr = false; + } else { + reader.line(partial.toString()); + partial.setLength(0); } + } else if (c == '\r') { + // A lone CR is a line terminator per the SSE grammar (CR / LF / CRLF). + pendingCr = true; reader.line(partial.toString()); partial.setLength(0); } else { + pendingCr = false; partial.append(c); } } @@ -127,6 +135,7 @@ protected Void generateContent() { @Override public void releaseResources() { partial.setLength(0); + pendingCr = false; reader = null; } diff --git a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java index 4b0ceefc73..da9b6ab1fa 100644 --- a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java +++ b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java @@ -84,6 +84,23 @@ void handlesBomCrLfAndDispatch() throws Exception { assertEquals("hi", cb.data); } + @Test + void treatsLoneCrAsLineTerminator() throws Exception { + final Cb cb = new Cb(); + final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + // A stream delimited with lone CR (a valid SSE separator). A CR in the id value must + // terminate the line rather than be retained, otherwise it would later be copied into + // the Last-Event-ID request header on reconnect. + final byte[] p = "id: 1\rdata: v\r\r".getBytes(StandardCharsets.UTF_8); + c.consume(ByteBuffer.wrap(p)); + c.streamEnd(null); + + assertEquals("1", cb.id); + assertEquals("v", cb.data); + } + @Test void emitsRetry() throws Exception { final Cb cb = new Cb(); diff --git a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/SseEntityConsumerTest.java b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/SseEntityConsumerTest.java index bafdd4ed12..5dd5a6a2f7 100644 --- a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/SseEntityConsumerTest.java +++ b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/SseEntityConsumerTest.java @@ -77,6 +77,21 @@ void parsesLinesAndFlushesOnEndOfStream() throws Exception { assertEquals("v", cb.data); } + @Test + void treatsLoneCrAsLineTerminator() throws Exception { + final Cb cb = new Cb(); + final SseEntityConsumer c = new SseEntityConsumer(cb); + + c.streamStart(ContentType.parse("text/event-stream")); + // A stream delimited with lone CR (a valid SSE separator). A CR in the id value must + // terminate the line rather than be retained, otherwise it would later be copied into + // the Last-Event-ID request header on reconnect. + c.data(CharBuffer.wrap("id: 1\rdata: v\r\r"), true); + + assertEquals("1", cb.id); + assertEquals("v", cb.data); + } + @Test void rejectsWrongContentType() { final Cb cb = new Cb(); From cec5b93aec3934189190830b90b508c171863795 Mon Sep 17 00:00:00 2001 From: Javid Khan Date: Wed, 7 Oct 2026 09:47:52 +0530 Subject: [PATCH 2/2] route the first non-BOM byte through the CR/LF handling in ByteSseEntityConsumer Signed-off-by: Javid Khan --- .../http/sse/impl/ByteSseEntityConsumer.java | 45 ++++++++++--------- .../http/sse/ByteSseEntityConsumerTest.java | 16 +++++++ 2 files changed, 40 insertions(+), 21 deletions(-) diff --git a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java index 2d7ff361bb..7a8225bd96 100644 --- a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java +++ b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ByteSseEntityConsumer.java @@ -116,12 +116,12 @@ protected void data(final ByteBuffer src, final boolean endOfStream) { continue; } if (bomMatched > 0) { - appendByte((byte) 0xEF); + processByte((byte) 0xEF); if (bomMatched >= 2) { - appendByte((byte) 0xBB); + processByte((byte) 0xBB); } } - appendByte((byte) b); + processByte((byte) b); bomMatched = 0; bomDone = true; break; // drop into normal loop below for the rest of 'src' @@ -135,24 +135,7 @@ protected void data(final ByteBuffer src, final boolean endOfStream) { } while (src.hasRemaining()) { - final byte b = src.get(); - if (b == LF) { - if (pendingCr) { - // LF completing a CRLF pair; the line was already emitted on the CR. - pendingCr = false; - } else { - handleLine(lineBuf, lineLen); - lineLen = 0; - } - } else if (b == CR) { - // A lone CR is a line terminator per the SSE grammar (CR / LF / CRLF). - pendingCr = true; - handleLine(lineBuf, lineLen); - lineLen = 0; - } else { - pendingCr = false; - appendByte(b); - } + processByte(src.get()); } if (endOfStream) { @@ -160,6 +143,26 @@ protected void data(final ByteBuffer src, final boolean endOfStream) { } } + private void processByte(final byte b) { + if (b == LF) { + if (pendingCr) { + // LF completing a CRLF pair; the line was already emitted on the CR. + pendingCr = false; + } else { + handleLine(lineBuf, lineLen); + lineLen = 0; + } + } else if (b == CR) { + // A lone CR is a line terminator per the SSE grammar (CR / LF / CRLF). + pendingCr = true; + handleLine(lineBuf, lineLen); + lineLen = 0; + } else { + pendingCr = false; + appendByte(b); + } + } + private void flushEndOfStream() { if (lineLen > 0) { handleLine(lineBuf, lineLen); diff --git a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java index da9b6ab1fa..632409a6a5 100644 --- a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java +++ b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ByteSseEntityConsumerTest.java @@ -101,6 +101,22 @@ void treatsLoneCrAsLineTerminator() throws Exception { assertEquals("v", cb.data); } + @Test + void treatsLeadingLoneCrAsLineTerminator() throws Exception { + final Cb cb = new Cb(); + final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + // The first byte is a lone CR, which is resolved during BOM detection. It must still be + // routed through the normal CR/LF handling so the empty leading line is terminated rather + // than the CR being retained inside the buffer. + final byte[] p = "\rdata: v\r\r".getBytes(StandardCharsets.UTF_8); + c.consume(ByteBuffer.wrap(p)); + c.streamEnd(null); + + assertEquals("v", cb.data); + } + @Test void emitsRetry() throws Exception { final Cb cb = new Cb();