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 7a8225bd96..164c1fb50f 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 @@ -51,8 +51,7 @@ *
  • 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 - * dispatch of the current event if it has data.
  • + *
  • At end of stream, discards any incomplete line or event.
  • * * *

    Thread-safety

    @@ -164,11 +163,9 @@ private void processByte(final byte b) { } private void flushEndOfStream() { - if (lineLen > 0) { - handleLine(lineBuf, lineLen); - lineLen = 0; - } - handleLine(lineBuf, 0); + lineLen = 0; + data.setLength(0); + type = null; } private void appendByte(final byte b) { 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 3e8327b0bd..fc1e930553 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 @@ -113,12 +113,7 @@ public void data(final CharBuffer src, final boolean endOfStream) { } } if (endOfStream) { - if (partial.length() > 0) { - reader.line(partial.toString()); - partial.setLength(0); - } - // Flush any accumulated fields into a final event. - reader.line(""); + partial.setLength(0); } } 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 632409a6a5..10f8474c74 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 @@ -27,6 +27,7 @@ package org.apache.hc.client5.http.sse; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; import java.nio.ByteBuffer; @@ -129,4 +130,52 @@ void emitsRetry() throws Exception { assertEquals(Long.valueOf(2500L), cb.retry); } + + @Test + void doesNotDispatchIncompleteEventAtEndOfStream() throws Exception { + final Cb cb = new Cb(); + final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.consume(ByteBuffer.wrap("data: hello\n".getBytes(StandardCharsets.UTF_8))); + c.streamEnd(null); + + assertNull(cb.data); + } + + @Test + void dispatchesCompleteEventBeforeEndOfStream() throws Exception { + final Cb cb = new Cb(); + final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.consume(ByteBuffer.wrap("data: hello\n\n".getBytes(StandardCharsets.UTF_8))); + c.streamEnd(null); + + assertEquals("hello", cb.data); + } + + @Test + void doesNotProcessIncompleteLineAtEndOfStream() throws Exception { + final Cb cb = new Cb(); + final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.consume(ByteBuffer.wrap("retry: 2500".getBytes(StandardCharsets.UTF_8))); + c.streamEnd(null); + + assertNull(cb.retry); + } + + @Test + void processesCompleteLineBeforeEndOfStream() throws Exception { + final Cb cb = new Cb(); + final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.consume(ByteBuffer.wrap("retry: 2500\n".getBytes(StandardCharsets.UTF_8))); + c.streamEnd(null); + + assertEquals(Long.valueOf(2500L), cb.retry); + } } 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 5dd5a6a2f7..a1008b49f7 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 @@ -27,6 +27,7 @@ package org.apache.hc.client5.http.sse; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -101,4 +102,48 @@ void rejectsWrongContentType() { fail("Should have thrown"); } catch (final Exception expected) { /* ok */ } } + + @Test + void doesNotDispatchIncompleteEventAtEndOfStream() throws Exception { + final Cb cb = new Cb(); + final SseEntityConsumer c = new SseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.data(CharBuffer.wrap("data: hello\n"), true); + + assertNull(cb.data); + } + + @Test + void dispatchesCompleteEventBeforeEndOfStream() throws Exception { + final Cb cb = new Cb(); + final SseEntityConsumer c = new SseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.data(CharBuffer.wrap("data: hello\n\n"), true); + + assertEquals("hello", cb.data); + } + + @Test + void doesNotProcessIncompleteLineAtEndOfStream() throws Exception { + final Cb cb = new Cb(); + final SseEntityConsumer c = new SseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.data(CharBuffer.wrap("retry: 2500"), true); + + assertNull(cb.retry); + } + + @Test + void processesCompleteLineBeforeEndOfStream() throws Exception { + final Cb cb = new Cb(); + final SseEntityConsumer c = new SseEntityConsumer(cb); + c.streamStart(ContentType.parse("text/event-stream")); + + c.data(CharBuffer.wrap("retry: 2500\n"), true); + + assertEquals(Long.valueOf(2500L), cb.retry); + } }