Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@
* <li>Validates {@code Content-Type} equals {@code text/event-stream}
* in {@link #streamStart(ContentType)}; otherwise throws {@link HttpException}.</li>
* <li>Strips a UTF-8 BOM if present in the first chunk.</li>
* <li>Accepts LF and CRLF line endings; tolerates CRLF split across buffers.</li>
* <li>Accepts CR, LF and CRLF line endings per the SSE grammar; tolerates CRLF split across buffers.</li>
* <li>Implements WHATWG SSE fields: {@code data}, {@code id}, {@code event}, {@code retry}.
* Unknown fields and malformed {@code retry} values are ignored.</li>
* <li>At end of stream, flushes any partially accumulated line and forces a final
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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'
Expand All @@ -134,31 +135,37 @@ 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) {
flushEndOfStream();
}
}

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);
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ public final class SseEntityConsumer extends AbstractCharAsyncEntityConsumer<Voi
private final StringBuilder partial = new StringBuilder(256);
private ServerSentEventReader reader;
private boolean firstChunk = true;
private boolean pendingCr = false;

public SseEntityConsumer(final SseCallbacks callbacks) {
this.cb = callbacks;
Expand Down Expand Up @@ -94,13 +95,20 @@ public void data(final CharBuffer src, final boolean endOfStream) {
while (src.hasRemaining()) {
final char c = src.get();
if (c == '\n') {
final int len = partial.length();
if (len > 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);
}
}
Expand All @@ -127,6 +135,7 @@ protected Void generateContent() {
@Override
public void releaseResources() {
partial.setLength(0);
pendingCr = false;
reader = null;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Loading