FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

Backport fix: race condition in ResponseSubscribers into v1.1.x by nickcaballero · Pull Request #1154 · modelcontextprotocol/java-sdk · GitHub

Repository navigation

Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension .java  (2) All 1 file type selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -320,7 +320,7 @@ static class AggregateSubscriber extends BaseSubscriber<String> {
* The response information from the HTTP response. Send with each event to
* provide context.
*/
private ResponseInfo responseInfo;
private final ResponseInfo responseInfo;

volatile boolean hasRequestedDemand = false;

Expand Down Expand Up @@ -348,16 +348,15 @@ public AggregateSubscriber(ResponseInfo responseInfo, FluxSink<ResponseEvent> si

@Override
protected void hookOnSubscribe(Subscription subscription) {
// Register disposal callback to cancel subscription when Flux is disposed
sink.onDispose(subscription::cancel);

sink.onRequest(n -> {
if (!hasRequestedDemand) {
hasRequestedDemand = true;
subscription.request(Long.MAX_VALUE);
}
hasRequestedDemand = true;
});

// Register disposal callback to cancel subscription when Flux is disposed
sink.onDispose(subscription::cancel);
}

@Override
Expand Down Expand Up @@ -410,17 +409,14 @@ public BodilessResponseLineSubscriber(ResponseInfo responseInfo, FluxSink<Respon

@Override
protected void hookOnSubscribe(Subscription subscription) {
// Register disposal callback to cancel subscription when Flux is disposed
sink.onDispose(subscription::cancel);

sink.onRequest(n -> {
if (!hasRequestedDemand) {
hasRequestedDemand = true;
subscription.request(Long.MAX_VALUE);
}
hasRequestedDemand = true;
});

// Register disposal callback to cancel subscription when Flux is disposed
sink.onDispose(() -> {
subscription.cancel();
});
}

Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
/*
* Copyright 2024 - 2024 the original author or authors.
*/

package io.modelcontextprotocol.client.transport;

import java.net.http.HttpResponse.ResponseInfo;

import org.junit.jupiter.api.Test;

import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;

import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;

class ResponseSubscribersTests {

@Test
void aggregateSubscriberEmitsResponseWhenRequestCompletesSynchronously() {
ResponseInfo responseInfo = mock(ResponseInfo.class);

Flux<ResponseSubscribers.ResponseEvent> response = Flux.create(sink -> {
var subscriber = new ResponseSubscribers.AggregateSubscriber(responseInfo, sink, Integer.MAX_VALUE);
Flux.just("payload").subscribe(subscriber);
});

StepVerifier.create(response).assertNext(event -> {
var aggregate = (ResponseSubscribers.AggregateResponseEvent) event;
assertThat(aggregate.responseInfo()).isSameAs(responseInfo);
assertThat(aggregate.data()).isEqualTo("payload\n");
}).verifyComplete();
}

@Test
void bodilessSubscriberEmitsResponseWhenRequestCompletesSynchronously() {
ResponseInfo responseInfo = mock(ResponseInfo.class);

Flux<ResponseSubscribers.ResponseEvent> response = Flux.create(sink -> {
var subscriber = new ResponseSubscribers.BodilessResponseLineSubscriber(responseInfo, sink);
Flux.<String>empty().subscribe(subscriber);
});

StepVerifier.create(response).assertNext(event -> {
var dummy = (ResponseSubscribers.DummyEvent) event;
assertThat(dummy.responseInfo()).isSameAs(responseInfo);
}).verifyComplete();
}

}

Back | FazBrowse Home | New Git URL