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
9 changes: 9 additions & 0 deletions src/main/java/ch/rasc/sse/eventbus/ClientEvent.java
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,15 @@ public SseEvent getSseEvent() {
return this.event;
}

/**
* The already converted data of the event, or {@code null} when the data is
* converted at send time.
* @return the converted data
*/
public @Nullable String getConvertedValue() {
return this.convertedValue;
}

public SseEventBuilder createSseEventBuilder() {

SseEventBuilder sseBuilder = SseEmitter.event();
Expand Down
51 changes: 49 additions & 2 deletions src/main/java/ch/rasc/sse/eventbus/ClientSendBuffer.java
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
package ch.rasc.sse.eventbus;

import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
Expand Down Expand Up @@ -49,6 +51,11 @@
* Delivery accounting happens on the dispatcher thread: {@link DeliveryListener} is
* invoked after the event was actually written to the {@code SseEmitter} (or after a send
* failure), so an enqueued event is never reported as delivered before the sink has run.
* <p>
* When an {@link EventCoalescer} is configured, consecutive buffered events that can be
* merged are written as a single SSE frame: the dispatcher drains a bounded batch from
* the queue and merges adjacent events. This keeps a high-frequency stream of small
* events (for example LLM token streaming) from filling up the queue of a slow client.
*/
final class ClientSendBuffer implements AutoCloseable {

Expand Down Expand Up @@ -130,6 +137,8 @@

private final @Nullable SlowClientListener listener;

private final @Nullable EventCoalescer coalescer;

private final @Nullable SseBackpressureMetrics metrics;

private final AtomicBoolean closed = new AtomicBoolean();
Expand All @@ -142,7 +151,8 @@

ClientSendBuffer(String clientId, int capacity, OverflowPolicy overflowPolicy, EventSink sink,
DeliveryListener deliveryListener, @Nullable Consumer<ClientSendBuffer> onDisconnect,
@Nullable SlowClientListener listener, @Nullable SseBackpressureMetrics metrics) {
@Nullable SlowClientListener listener, @Nullable EventCoalescer coalescer,
@Nullable SseBackpressureMetrics metrics) {
this.clientId = Objects.requireNonNull(clientId, "clientId");
if (capacity <= 0) {
throw new IllegalArgumentException("clientSendBufferCapacity must be > 0");
Expand All @@ -154,6 +164,7 @@
this.deliveryListener = Objects.requireNonNull(deliveryListener, "deliveryListener");
this.onDisconnect = onDisconnect;
this.listener = listener;
this.coalescer = coalescer;
this.metrics = metrics;
}

Expand Down Expand Up @@ -243,6 +254,11 @@
}
}

/**
* Maximum number of events drained from the queue for a single coalescing pass.
*/
private static final int MAX_COALESCE_BATCH = 32;

private void dispatchLoop() {
try {
while (!this.closed.get() && !Thread.currentThread().isInterrupted()) {
Expand All @@ -258,7 +274,14 @@
continue;
}
if (!this.closed.get()) {
deliver(event);
EventCoalescer coalescer = this.coalescer;
if (coalescer == null) {
deliver(event);
continue;
}
List<BufferedEvent> rest = new ArrayList<>(MAX_COALESCE_BATCH - 1);
this.queue.drainTo(rest, MAX_COALESCE_BATCH - 1);
deliverCoalesced(coalescer, event, rest);
}
}
}
Expand All @@ -267,6 +290,30 @@
}
}

private void deliverCoalesced(EventCoalescer coalescer, BufferedEvent first, List<BufferedEvent> rest) {
BufferedEvent current = first;
for (BufferedEvent next : rest) {
if (current.heartbeat() || next.heartbeat()) {
// heartbeats are sent as-is, they are not merged with regular events
deliver(current);
current = next;
continue;
}
ClientEvent merged = coalescer.coalesce(current.event(), next.event());
if (merged != null) {
if (this.metrics != null) {
this.metrics.recordCoalesced(this.clientId);
}
current = new BufferedEvent(merged, false);
}
else {
deliver(current);
current = next;
}
}
deliver(current);
}

private void deliver(BufferedEvent bufferedEvent) {
ClientEvent event = bufferedEvent.event();
try {
Expand Down Expand Up @@ -312,7 +359,7 @@
*/
void awaitDrained(long deadlineNanos) {
Thread thread = this.dispatcherThread;
if (thread == null || thread == Thread.currentThread()) {

Check warning on line 362 in src/main/java/ch/rasc/sse/eventbus/ClientSendBuffer.java

View workflow job for this annotation

GitHub Actions / build (21)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 362 in src/main/java/ch/rasc/sse/eventbus/ClientSendBuffer.java

View workflow job for this annotation

GitHub Actions / build (25)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 362 in src/main/java/ch/rasc/sse/eventbus/ClientSendBuffer.java

View workflow job for this annotation

GitHub Actions / build (26)

[ReferenceEquality] Comparison using reference equality instead of value equality
return;
}
try {
Expand Down
71 changes: 71 additions & 0 deletions src/main/java/ch/rasc/sse/eventbus/DefaultEventCoalescer.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
/*
* Copyright the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package ch.rasc.sse.eventbus;

import org.jspecify.annotations.Nullable;

/**
* Default {@link EventCoalescer} that merges plain data events.
* <p>
* Two events are merged when both have no event id, no retry, no comment, no JSON view
* and the same event name. Their payload is taken from the converted value when present,
* or directly from a String payload (the send path writes String payloads as-is). The
* data is joined with a line break, which the SSE protocol encodes as multiple
* {@code data:} lines of a single event, so a client that concatenates the data lines
* receives the original data in order. Unconverted non-String payloads are never merged.
*/
public class DefaultEventCoalescer implements EventCoalescer {

@Override
public @Nullable ClientEvent coalesce(ClientEvent first, ClientEvent second) {
SseEvent firstEvent = first.getSseEvent();
SseEvent secondEvent = second.getSseEvent();
if (firstEvent.event() != null && secondEvent.event() != null
&& !firstEvent.event().equals(secondEvent.event())) {
return null;
}
if (firstEvent.id().isPresent() || secondEvent.id().isPresent() || firstEvent.retry().isPresent()
|| secondEvent.retry().isPresent() || firstEvent.comment().isPresent()
|| secondEvent.comment().isPresent() || firstEvent.jsonView().isPresent()
|| secondEvent.jsonView().isPresent()) {
return null;
}
String firstData = eventData(first);
String secondData = eventData(second);
if (firstData == null || secondData == null) {
return null;
}
String mergedData = firstData + "\n" + secondData;
SseEvent merged = SseEvent.builder()
.event(firstEvent.event())
.data(mergedData)
.build();
return new ClientEvent(first.getClient(), merged, mergedData);
}

private static @Nullable String eventData(ClientEvent event) {
String converted = event.getConvertedValue();
if (converted != null) {
return converted;
}
// A String payload is written as-is by the send path (the converted value stays
// null), so it can be merged directly. Unconverted non-String payloads are never
// merged: stringifying them here could differ from the converter's output.
Object data = event.getSseEvent().data();
return data instanceof String string ? string : null;
}

}
47 changes: 47 additions & 0 deletions src/main/java/ch/rasc/sse/eventbus/EventCoalescer.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
/*
* Copyright the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package ch.rasc.sse.eventbus;

import org.jspecify.annotations.Nullable;

/**
* Merges two consecutive events of the same client into a single event before it is
* written to the connection.
* <p>
* This is useful for slow clients that receive a high-frequency stream of small events,
* for example LLM token streaming. Coalescing reduces the number of writes and network
* frames and keeps the per-client send buffer from filling up, so fewer events are
* dropped or trigger a disconnect.
* <p>
* Coalescing is opt-in: return {@code null} from {@link #coalesce(ClientEvent, ClientEvent)}
* for pairs that must be sent separately. The default implementation is
* {@link DefaultEventCoalescer}, which merges plain data events by joining their data
* with a line break.
*/
@FunctionalInterface
public interface EventCoalescer {

/**
* Merges two consecutive buffered events of the same client.
* @param first the first (older) event
* @param second the second (newer) event
* @return the merged event, or {@code null} when the two events cannot be merged and
* must be written separately
*/
@Nullable
ClientEvent coalesce(ClientEvent first, ClientEvent second);

}
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ final class SseBackpressureMetrics {

private final Counter slowClientNotifications;

private final Counter coalescedEvents;

private final Map<String, Gauge> queueGauges = new ConcurrentHashMap<>();

private final Map<String, Supplier<Number>> strongSuppliers = new ConcurrentHashMap<>();
Expand All @@ -75,6 +77,9 @@ final class SseBackpressureMetrics {
this.slowClientNotifications = Counter.builder(metric("slow.client.notifications"))
.description("Slow client callback invocations")
.register(registry);
this.coalescedEvents = Counter.builder(metric("coalesced.events"))
.description("Events merged into a single SSE frame by the event coalescer")
.register(registry);
}

private String metric(String name) {
Expand All @@ -97,6 +102,10 @@ void recordNotification(String clientId) {
this.slowClientNotifications.increment();
}

void recordCoalesced(String clientId) {
this.coalescedEvents.increment();
}

void registerQueueGauge(String clientId, Supplier<Number> supplier) {
this.strongSuppliers.put("queue:" + clientId, supplier);
this.queueGauges.computeIfAbsent(clientId,
Expand Down
5 changes: 4 additions & 1 deletion src/main/java/ch/rasc/sse/eventbus/SseEventBus.java
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,8 @@

private final SlowClientListener slowClientListener;

private final @Nullable EventCoalescer eventCoalescer;

private final @Nullable SseBackpressureMetrics backpressureMetrics;

private final @Nullable ReplayStore replayStore;
Expand Down Expand Up @@ -199,6 +201,7 @@
this.clientSendBufferCapacity = configurer.clientSendBufferCapacity();
this.overflowPolicy = configurer.overflowPolicy();
this.slowClientListener = configurer.slowClientListener();
this.eventCoalescer = configurer.eventCoalescer();
MeterRegistry meterRegistry = configurer.meterRegistry();
this.backpressureMetrics = meterRegistry != null ? new SseBackpressureMetrics(meterRegistry) : null;
if (this.backpressureMetrics != null) {
Expand Down Expand Up @@ -470,7 +473,7 @@
if (this.replayEnabled) {
this.replayLocks.computeIfAbsent(clientId, k -> new ReentrantLock());
}
if (oldEmitter.get() != null && oldEmitter.get() != emitter) {

Check warning on line 476 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (21)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 476 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (25)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 476 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (26)

[ReferenceEquality] Comparison using reference equality instead of value equality
try {
oldEmitter.get().complete();
}
Expand Down Expand Up @@ -804,7 +807,7 @@
break;
}
String clientId = clientEvent.getClient().getId();
if (this.clients.get(clientId) != clientEvent.getClient() || !this.subscriptionRegistry

Check warning on line 810 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (21)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 810 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (25)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 810 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (26)

[ReferenceEquality] Comparison using reference equality instead of value equality
.isClientSubscribedToEvent(clientId, clientEvent.getSseEvent().event())) {
continue;
}
Expand Down Expand Up @@ -846,7 +849,7 @@
while (!Thread.currentThread().isInterrupted()) {
try {
ClientEvent clientEvent = this.sendQueue.take();
if (this.clients.get(clientEvent.getClient().getId()) != clientEvent.getClient()) {

Check warning on line 852 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (21)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 852 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (25)

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 852 in src/main/java/ch/rasc/sse/eventbus/SseEventBus.java

View workflow job for this annotation

GitHub Actions / build (26)

[ReferenceEquality] Comparison using reference equality instead of value equality
continue;
}
if (clientEvent.getErrorCounter() < this.noOfSendResponseTries) {
Expand Down Expand Up @@ -1044,7 +1047,7 @@
notifyAfterEventSent(event, exception);
}
}, disconnectedBuffer -> unregisterClient(client.getId(), disconnectedBuffer),
this.slowClientListener, this.backpressureMetrics);
this.slowClientListener, this.eventCoalescer, this.backpressureMetrics);
// Register the queue gauge before the buffer becomes visible to the send
// workers to avoid a registration race with an early disconnect
if (this.backpressureMetrics != null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@

import ch.rasc.sse.eventbus.Client;
import ch.rasc.sse.eventbus.ClientEvent;
import ch.rasc.sse.eventbus.EventCoalescer;
import ch.rasc.sse.eventbus.OverflowPolicy;
import ch.rasc.sse.eventbus.ReplayStore;
import ch.rasc.sse.eventbus.SlowClientListener;
Expand Down Expand Up @@ -223,4 +224,19 @@ default SlowClientListener slowClientListener() {
return null;
}

/**
* Optional coalescer that merges consecutive buffered events of the same client into a
* single SSE frame before they are written to the connection. This is useful for slow
* clients that receive a high-frequency stream of small events (for example LLM token
* streaming): fewer, larger writes reduce system call and network overhead and keep
* the per-client send buffer from filling up.
* <p>
* Only used when {@link #clientSendBufferCapacity()} is greater than zero.
* <p>
* Default: {@code null} (coalescing disabled, each event is sent separately)
*/
default @Nullable EventCoalescer eventCoalescer() {
return null;
}

}
Loading
Loading