Skip to content
Open
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 @@ -240,10 +240,15 @@ private void handleLine(final byte[] buf, final int len) {
}

private void dispatch() {
if (id != null) {
cb.onLastEventId(id);
}

if (data.length() == 0) {
type = null;
return;
}

final int n = data.length();
if (n > 0 && data.charAt(n - 1) == '\n') {
data.setLength(n - 1);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -323,7 +323,7 @@ private void doConnect() {
final SimpleRequestBuilder rb = SimpleRequestBuilder.get(uri);
rb.setHeader(HttpHeaders.ACCEPT, TEXT_EVENT_STREAM.getMimeType());
rb.setHeader(HttpHeaders.CACHE_CONTROL, "no-cache");
if (lastEventId != null) {
if (lastEventId != null && !lastEventId.isEmpty()) {
rb.setHeader("Last-Event-ID", lastEventId);
}
for (final Map.Entry<String, String> e : headers.entrySet()) {
Expand Down Expand Up @@ -393,11 +393,13 @@ public void onOpen() {
dispatch(listener::onOpen);
}

@Override
public void onLastEventId(final String id) {
lastEventId = id;
}

@Override
public void onEvent(final String id, final String type, final String data) {
if (id != null) {
lastEventId = id;
}
dispatch(() -> listener.onEvent(id, type, data));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,9 @@ public interface Callback {
void onComment(String comment);

void onRetryChange(long retryMs);

default void onLastEventId(final String id) {
}
}

private final Callback cb;
Expand Down Expand Up @@ -159,11 +162,16 @@ public void line(final String line) {
}

private void dispatch() {
if (id != null) {
cb.onLastEventId(id);
}

if (data.length() == 0) {
// spec: a blank line with no "data:" accumulates nothing -> just clear type
type = null;
return;
}

final int n = data.length();
if (n > 0 && data.charAt(n - 1) == '\n') {
data.setLength(n - 1);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,4 +64,12 @@ public interface SseCallbacks {
* @param retryMs new retry delay in milliseconds (non-negative)
*/
void onRetry(long retryMs);

/**
* Notifies that the last event ID has been committed by the SSE dispatch step.
*
* @param id the last event ID, possibly empty to reset it
*/
default void onLastEventId(final String id) {
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,11 @@ public void releaseResources() {
reader = null;
}

@Override
public void onLastEventId(final String id) {
cb.onLastEventId(id);
}

// ServerSentEventReader.Callback

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -43,6 +44,7 @@ static final class Cb implements SseCallbacks {
boolean opened;
String id, type, data;
Long retry;
String lastEventId;

@Override
public void onOpen() {
Expand All @@ -60,6 +62,11 @@ public void onEvent(final String id, final String type, final String data) {
public void onRetry(final long retryMs) {
retry = retryMs;
}

@Override
public void onLastEventId(final String id) {
lastEventId = id;
}
}


Expand Down Expand Up @@ -96,4 +103,32 @@ void emitsRetry() throws Exception {

assertEquals(Long.valueOf(2500L), cb.retry);
}

@Test
void updatesLastEventIdWithoutDispatchingEvent() throws Exception {
final Cb cb = new Cb();
final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb);

c.streamStart(ContentType.parse("text/event-stream"));

final byte[] p = "id: 42\n\n".getBytes(StandardCharsets.UTF_8);
c.consume(ByteBuffer.wrap(p));

assertEquals("42", cb.lastEventId);
assertNull(cb.data);
}

@Test
void resetsLastEventIdWithoutDispatchingEvent() throws Exception {
final Cb cb = new Cb();
final ByteSseEntityConsumer c = new ByteSseEntityConsumer(cb);

c.streamStart(ContentType.parse("text/event-stream"));

final byte[] p = "id: 42\n\nid:\n\n".getBytes(StandardCharsets.UTF_8);
c.consume(ByteBuffer.wrap(p));

assertEquals("", cb.lastEventId);
assertNull(cb.data);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ class ServerSentEventReaderTest {

static final class Capt implements ServerSentEventReader.Callback {
String id, type, data, comment;
String lastEventId;
Long retry;

@Override
Expand All @@ -45,6 +46,11 @@ public void onEvent(final String id, final String type, final String data) {
this.data = data;
}

@Override
public void onLastEventId(final String id) {
this.lastEventId = id;
}

@Override
public void onComment(final String comment) {
this.comment = comment;
Expand Down Expand Up @@ -101,4 +107,36 @@ void ignoresIdWithNul() {
assertNull(c.id);
assertEquals("d", c.data);
}

@Test
void updatesLastEventIdOnBlankLineWithoutDispatchingEvent() {
final Capt c = new Capt();
final ServerSentEventReader r = new ServerSentEventReader(c);

r.line("id: 42");

assertNull(c.lastEventId);

r.line("");

assertEquals("42", c.lastEventId);
assertNull(c.data);
}

@Test
void resetsLastEventIdOnBlankLineWithoutDispatchingEvent() {
final Capt c = new Capt();
final ServerSentEventReader r = new ServerSentEventReader(c);

r.line("id: 42");
r.line("");

assertEquals("42", c.lastEventId);

r.line("id:");
r.line("");

assertEquals("", c.lastEventId);
assertNull(c.data);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -43,6 +44,7 @@ static final class Cb implements SseCallbacks {
boolean opened;
String id, type, data;
Long retry;
String lastEventId;

@Override
public void onOpen() {
Expand All @@ -60,6 +62,11 @@ public void onEvent(final String id, final String type, final String data) {
public void onRetry(final long retryMs) {
retry = retryMs;
}

@Override
public void onLastEventId(final String id) {
lastEventId = id;
}
}

@Test
Expand All @@ -86,4 +93,28 @@ void rejectsWrongContentType() {
fail("Should have thrown");
} catch (final Exception expected) { /* ok */ }
}

@Test
void updatesLastEventIdWithoutDispatchingEvent() throws Exception {
final Cb cb = new Cb();
final SseEntityConsumer c = new SseEntityConsumer(cb);

c.streamStart(ContentType.parse("text/event-stream"));
c.data(CharBuffer.wrap("id: 42\n\n"), false);

assertEquals("42", cb.lastEventId);
assertNull(cb.data);
}

@Test
void resetsLastEventIdWithoutDispatchingEvent() throws Exception {
final Cb cb = new Cb();
final SseEntityConsumer c = new SseEntityConsumer(cb);

c.streamStart(ContentType.parse("text/event-stream"));
c.data(CharBuffer.wrap("id: 42\n\nid:\n\n"), false);

assertEquals("", cb.lastEventId);
assertNull(cb.data);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
*/
package org.apache.hc.client5.http.sse.impl;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;

Expand All @@ -49,6 +50,7 @@
import org.apache.hc.core5.concurrent.FutureCallback;
import org.apache.hc.core5.function.Supplier;
import org.apache.hc.core5.http.HttpHost;
import org.apache.hc.core5.http.HttpRequest;
import org.apache.hc.core5.http.nio.AsyncPushConsumer;
import org.apache.hc.core5.http.nio.AsyncRequestProducer;
import org.apache.hc.core5.http.nio.AsyncResponseConsumer;
Expand Down Expand Up @@ -127,6 +129,7 @@ void callerSchedulerIsNotShutdownOnCancel() {

static final class CapturingClient extends CloseableHttpAsyncClient {
volatile FutureCallback<Void> lastCallback;
volatile HttpRequest lastRequest;

@Override
public void start() { }
Expand Down Expand Up @@ -156,8 +159,19 @@ protected <T> Future<T> doExecute(
final HandlerFactory<AsyncPushConsumer> pushHandlerFactory,
final HttpContext context,
final FutureCallback<T> callback) {
@SuppressWarnings("unchecked") final FutureCallback<Void> cb = (FutureCallback<Void>) callback;

@SuppressWarnings("unchecked")
final FutureCallback<Void> cb = (FutureCallback<Void>) callback;
this.lastCallback = cb;

try {
requestProducer.sendRequest(
(request, entityDetails, requestContext) -> this.lastRequest = request,
context);
} catch (final Exception ex) {
throw new IllegalStateException(ex);
}

return new CompletableFuture<>();
}

Expand Down Expand Up @@ -265,4 +279,48 @@ public V get(final long timeout, final TimeUnit unit) {
return null;
}
}

@Test
void doesNotSendEmptyLastEventId() {
final RecordingScheduler scheduler = new RecordingScheduler();
final CapturingClient client = new CapturingClient();

final DefaultEventSource es = new DefaultEventSource(
client,
URI.create("http://example.org/sse"),
Collections.emptyMap(),
(id, type, data) -> { },
scheduler,
null,
null,
SseParser.CHAR);

es.setLastEventId("");
es.start();

assertFalse(client.lastRequest.containsHeader("Last-Event-ID"));
}

@Test
void sendsNonEmptyLastEventId() {
final RecordingScheduler scheduler = new RecordingScheduler();
final CapturingClient client = new CapturingClient();

final DefaultEventSource es = new DefaultEventSource(
client,
URI.create("http://example.org/sse"),
Collections.emptyMap(),
(id, type, data) -> { },
scheduler,
null,
null,
SseParser.CHAR);

es.setLastEventId("42");
es.start();

assertEquals(
"42",
client.lastRequest.getFirstHeader("Last-Event-ID").getValue());
}
}
Loading