1313import java .util .concurrent .CompletionStage ;
1414import java .util .concurrent .Flow ;
1515import java .util .concurrent .atomic .AtomicReference ;
16- import java .util .regex .Pattern ;
1716
1817import org .reactivestreams .FlowAdapters ;
1918import org .reactivestreams .Subscription ;
@@ -147,21 +146,6 @@ static BodyHandler<String> boundedStringBodyHandler(int maxSize) {
147146
148147 static class SseLineSubscriber extends BaseSubscriber <String > {
149148
150- /**
151- * Pattern to extract data content from SSE "data:" lines.
152- */
153- private static final Pattern EVENT_DATA_PATTERN = Pattern .compile ("^data:(.+)$" , Pattern .MULTILINE );
154-
155- /**
156- * Pattern to extract event ID from SSE "id:" lines.
157- */
158- private static final Pattern EVENT_ID_PATTERN = Pattern .compile ("^id:(.+)$" , Pattern .MULTILINE );
159-
160- /**
161- * Pattern to extract event type from SSE "event:" lines.
162- */
163- private static final Pattern EVENT_TYPE_PATTERN = Pattern .compile ("^event:(.+)$" , Pattern .MULTILINE );
164-
165149 /**
166150 * The sink for emitting parsed response events.
167151 */
@@ -227,49 +211,75 @@ protected void hookOnSubscribe(Subscription subscription) {
227211 });
228212 }
229213
214+ /**
215+ * Extracts the value of an SSE field from a line, per the <a href=
216+ * "https://html.spec.whatwg.org/multipage/server-sent-events.html#event-stream-interpretation">
217+ * SSE specification</a>: the characters after the colon with a single leading
218+ * space removed.
219+ *
220+ * <p>
221+ * A value may legally contain U+2028 (LINE SEPARATOR), U+2029 (PARAGRAPH
222+ * SEPARATOR) and U+0085 (NEXT LINE). Those are not SSE line terminators, so
223+ * extracting the value with a {@code MULTILINE} regex instead of this method
224+ * silently truncates it there.
225+ * @param line the SSE line, already stripped of its terminator by the line
226+ * subscriber
227+ * @param field the field prefix, e.g. {@code "data:"}
228+ * @return the field value with a single leading space removed, never truncated
229+ * @see #hookOnNext(String)
230+ */
231+ private static String fieldValue (String line , String field ) {
232+ String value = line .substring (field .length ());
233+ if (value .startsWith (" " )) {
234+ value = value .substring (1 );
235+ }
236+ return value ;
237+ }
238+
239+ /**
240+ * Returns the buffered data lines joined by the separators the data: handler
241+ * appended, with only the final separator removed. Trimming the whole buffer
242+ * instead would also strip significant leading/trailing whitespace from the first
243+ * and last data lines, which the SSE field rules explicitly preserve.
244+ */
245+ private String concatenatedDataLines () {
246+ String buffered = this .eventBuilder .toString ();
247+ return buffered .substring (0 , buffered .length () - 1 );
248+ }
249+
230250 @ Override
231251 protected void hookOnNext (String line ) {
232252 if (line .isEmpty ()) {
233253 // Empty line means end of event
234254 if (this .eventBuilder .length () > 0 ) {
235- String eventData = this . eventBuilder . toString ();
236- SseEvent sseEvent = new SseEvent (currentEventId .get (), currentEventType .get (), eventData . trim () );
255+ String eventData = concatenatedDataLines ();
256+ SseEvent sseEvent = new SseEvent (currentEventId .get (), currentEventType .get (), eventData );
237257
238258 this .sink .next (new SseResponseEvent (responseInfo , sseEvent ));
239259 this .eventBuilder .setLength (0 );
240260 }
241261 }
242262 else {
243263 if (line .startsWith ("data:" )) {
244- var matcher = EVENT_DATA_PATTERN .matcher (line );
245- if (matcher .find ()) {
246- String data = matcher .group (1 ).trim ();
247- // Measured before appending, so that an event carrying exactly
248- // maxSize of data is accepted: the trailing separator below is
249- // stripped again before the event is emitted.
250- if (this .eventBuilder .length () + data .length () > this .maxSize ) {
251- upstream ().cancel ();
252- this .sink .error (
253- new McpTransportException ("Inbound SSE event exceeds the maximum allowed size of "
254- + this .maxSize + " bytes" ));
255- return ;
256- }
257- this .eventBuilder .append (data ).append ("\n " );
264+ String data = fieldValue (line , "data:" );
265+ // Measured before appending, so that an event carrying exactly
266+ // maxSize of data is accepted: the trailing separator below is
267+ // stripped again before the event is emitted.
268+ if (this .eventBuilder .length () + data .length () > this .maxSize ) {
269+ upstream ().cancel ();
270+ this .sink .error (new McpTransportException (
271+ "Inbound SSE event exceeds the maximum allowed size of " + this .maxSize + " bytes" ));
272+ return ;
258273 }
274+ this .eventBuilder .append (data ).append ("\n " );
259275 upstream ().request (1 );
260276 }
261277 else if (line .startsWith ("id:" )) {
262- var matcher = EVENT_ID_PATTERN .matcher (line );
263- if (matcher .find ()) {
264- this .currentEventId .set (matcher .group (1 ).trim ());
265- }
278+ this .currentEventId .set (fieldValue (line , "id:" ));
266279 upstream ().request (1 );
267280 }
268281 else if (line .startsWith ("event:" )) {
269- var matcher = EVENT_TYPE_PATTERN .matcher (line );
270- if (matcher .find ()) {
271- this .currentEventType .set (matcher .group (1 ).trim ());
272- }
282+ this .currentEventType .set (fieldValue (line , "event:" ));
273283 upstream ().request (1 );
274284 }
275285 else if (line .startsWith (":" )) {
@@ -290,8 +300,8 @@ else if (line.startsWith(":")) {
290300 @ Override
291301 protected void hookOnComplete () {
292302 if (this .eventBuilder .length () > 0 ) {
293- String eventData = this . eventBuilder . toString ();
294- SseEvent sseEvent = new SseEvent (currentEventId .get (), currentEventType .get (), eventData . trim () );
303+ String eventData = concatenatedDataLines ();
304+ SseEvent sseEvent = new SseEvent (currentEventId .get (), currentEventType .get (), eventData );
295305 this .sink .next (new SseResponseEvent (responseInfo , sseEvent ));
296306 }
297307 this .sink .complete ();
0 commit comments