From 6531b3c065277c5536f5fa429316f3cab353c7f5 Mon Sep 17 00:00:00 2001 From: Arturo Bernal Date: Tue, 6 Oct 2026 13:54:17 +0200 Subject: [PATCH] SSE: Update Last-Event-ID independently of event dispatch --- .../http/sse/impl/ByteSseEntityConsumer.java | 5 ++ .../http/sse/impl/DefaultEventSource.java | 10 ++-- .../http/sse/impl/ServerSentEventReader.java | 8 +++ .../client5/http/sse/impl/SseCallbacks.java | 8 +++ .../http/sse/impl/SseEntityConsumer.java | 5 ++ .../http/sse/ByteSseEntityConsumerTest.java | 35 +++++++++++ .../http/sse/ServerSentEventReaderTest.java | 38 ++++++++++++ .../http/sse/SseEntityConsumerTest.java | 31 ++++++++++ .../http/sse/impl/DefaultEventSourceTest.java | 60 ++++++++++++++++++- 9 files changed, 195 insertions(+), 5 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..dcbe49c251 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 @@ -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); diff --git a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/DefaultEventSource.java b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/DefaultEventSource.java index 41169a7c3a..6d59dff0ec 100644 --- a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/DefaultEventSource.java +++ b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/DefaultEventSource.java @@ -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 e : headers.entrySet()) { @@ -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)); } diff --git a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ServerSentEventReader.java b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ServerSentEventReader.java index d674cd1c47..3dd7b1ebb4 100644 --- a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ServerSentEventReader.java +++ b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/ServerSentEventReader.java @@ -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; @@ -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); diff --git a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseCallbacks.java b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseCallbacks.java index 4d480f2393..75c1653365 100644 --- a/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseCallbacks.java +++ b/httpclient5-sse/src/main/java/org/apache/hc/client5/http/sse/impl/SseCallbacks.java @@ -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) { + } } 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..c3d8caa1a4 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 @@ -130,6 +130,11 @@ public void releaseResources() { reader = null; } + @Override + public void onLastEventId(final String id) { + cb.onLastEventId(id); + } + // ServerSentEventReader.Callback @Override 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..ad71724336 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; @@ -43,6 +44,7 @@ static final class Cb implements SseCallbacks { boolean opened; String id, type, data; Long retry; + String lastEventId; @Override public void onOpen() { @@ -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; + } } @@ -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); + } } diff --git a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ServerSentEventReaderTest.java b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ServerSentEventReaderTest.java index 3e1951508d..32360698ba 100644 --- a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ServerSentEventReaderTest.java +++ b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/ServerSentEventReaderTest.java @@ -36,6 +36,7 @@ class ServerSentEventReaderTest { static final class Capt implements ServerSentEventReader.Callback { String id, type, data, comment; + String lastEventId; Long retry; @Override @@ -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; @@ -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); + } } 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..8a08feb215 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; @@ -43,6 +44,7 @@ static final class Cb implements SseCallbacks { boolean opened; String id, type, data; Long retry; + String lastEventId; @Override public void onOpen() { @@ -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 @@ -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); + } } diff --git a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/impl/DefaultEventSourceTest.java b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/impl/DefaultEventSourceTest.java index 894d6d35e9..43d2d1fd50 100644 --- a/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/impl/DefaultEventSourceTest.java +++ b/httpclient5-sse/src/test/java/org/apache/hc/client5/http/sse/impl/DefaultEventSourceTest.java @@ -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; @@ -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; @@ -127,6 +129,7 @@ void callerSchedulerIsNotShutdownOnCancel() { static final class CapturingClient extends CloseableHttpAsyncClient { volatile FutureCallback lastCallback; + volatile HttpRequest lastRequest; @Override public void start() { } @@ -156,8 +159,19 @@ protected Future doExecute( final HandlerFactory pushHandlerFactory, final HttpContext context, final FutureCallback callback) { - @SuppressWarnings("unchecked") final FutureCallback cb = (FutureCallback) callback; + + @SuppressWarnings("unchecked") + final FutureCallback cb = (FutureCallback) 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<>(); } @@ -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()); + } }