diff --git a/allure-okhttp3/build.gradle.kts b/allure-okhttp3/build.gradle.kts index 477d717f3..a88fb6775 100644 --- a/allure-okhttp3/build.gradle.kts +++ b/allure-okhttp3/build.gradle.kts @@ -7,6 +7,7 @@ dependencies { compileOnly("com.squareup.okhttp3:okhttp:$okhttpVersion") testImplementation("org.wiremock:wiremock") testImplementation("com.squareup.okhttp3:okhttp:$okhttpVersion") + testImplementation("com.squareup.okhttp3:okhttp-sse:$okhttpVersion") testImplementation("org.assertj:assertj-core") testImplementation(project(":allure-assertj")) testImplementation("org.jboss.resteasy:resteasy-client") diff --git a/allure-okhttp3/src/main/java/io/qameta/allure/okhttp3/AllureOkHttp3.java b/allure-okhttp3/src/main/java/io/qameta/allure/okhttp3/AllureOkHttp3.java index 2a1338fd1..7dbecba55 100644 --- a/allure-okhttp3/src/main/java/io/qameta/allure/okhttp3/AllureOkHttp3.java +++ b/allure-okhttp3/src/main/java/io/qameta/allure/okhttp3/AllureOkHttp3.java @@ -89,7 +89,7 @@ public Response intercept(final Chain chain) throws IOException { final Response.Builder okHttpResponseBuilder = response.newBuilder(); final ResponseBody responseBody = response.body(); - if (Objects.nonNull(responseBody)) { + if (Objects.nonNull(responseBody) && !isEventStream(responseBody.contentType())) { final byte[] bytes = responseBody.bytes(); responseBuilder.setBody(body(responseBody.contentType(), new String(bytes, StandardCharsets.UTF_8))); okHttpResponseBuilder.body(ResponseBody.create(responseBody.contentType(), bytes)); @@ -129,6 +129,16 @@ private static List toNameValues(final Map#1036 + */ +@IsolatedLifecycle +class AllureOkHttp3SseTest { + + private static final List EVENTS = List.of("first", "second"); + + private HttpServer server; + private final CountDownLatch connectionHold = new CountDownLatch(1); + + @BeforeEach + void setUp() throws IOException { + server = HttpServer.create(new InetSocketAddress("localhost", 0), 0); + server.createContext("/sse", exchange -> { + exchange.getResponseHeaders().add("Content-Type", "text/event-stream"); + exchange.sendResponseHeaders(200, 0); + try (OutputStream os = exchange.getResponseBody()) { + for (final String event : EVENTS) { + os.write(("data: " + event + "\n\n").getBytes(StandardCharsets.UTF_8)); + os.flush(); + } + // hold the connection open: if the body completes, buffering passes too + // and the test stops catching the regression + try { + connectionHold.await(15, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + }); + server.start(); + } + + @AfterEach + void tearDown() { + connectionHold.countDown(); + if (Objects.nonNull(server)) { + server.stop(0); + } + } + + @Test + void shouldDeliverSseEventsWhileStreamIsOpen() throws InterruptedException { + final List received = new CopyOnWriteArrayList<>(); + final CountDownLatch eventsArrived = new CountDownLatch(EVENTS.size()); + + final OkHttpClient client = new OkHttpClient.Builder() + .addInterceptor(new AllureOkHttp3()) + .build(); + final Request request = new Request.Builder() + .url("http://localhost:" + server.getAddress().getPort() + "/sse") + .build(); + + final EventSourceListener listener = new EventSourceListener() { + @Override + public void onEvent(final EventSource eventSource, final String id, + final String type, final String data) { + received.add(data); + eventsArrived.countDown(); + } + + @Override + public void onFailure(final EventSource eventSource, final Throwable t, + final okhttp3.Response response) { + // cancel() lands here on every run, so don't fail the test - just stop waiting + while (eventsArrived.getCount() > 0) { + eventsArrived.countDown(); + } + } + }; + + runWithinTestContext(() -> { + final EventSource eventSource = EventSources.createFactory(client) + .newEventSource(request, listener); + try { + eventsArrived.await(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + eventSource.cancel(); + } + }); + + assertThat(received).containsExactlyElementsOf(EVENTS); + } +}