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..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 @@ -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); @@ -115,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' @@ -134,17 +135,7 @@ 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--; - } - handleLine(lineBuf, len); - lineLen = 0; - } else { - appendByte(b); - } + processByte(src.get()); } if (endOfStream) { @@ -152,13 +143,29 @@ 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) { - int len = lineLen; - if (lineBuf[len - 1] == CR) { - len--; - } - handleLine(lineBuf, len); + handleLine(lineBuf, lineLen); lineLen = 0; } handleLine(lineBuf, 0); @@ -182,6 +189,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..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 @@ -84,6 +84,39 @@ 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 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(); 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();