|
|
@@ -17,7 +17,6 @@ import org.elasticsearch.action.support.master.AcknowledgedResponse;
|
|
|
import org.elasticsearch.common.Strings;
|
|
|
import org.elasticsearch.common.compress.CompressorFactory;
|
|
|
import org.elasticsearch.common.io.stream.ByteBufferStreamInput;
|
|
|
-import org.elasticsearch.common.io.stream.InputStreamStreamInput;
|
|
|
import org.elasticsearch.common.io.stream.NamedWriteableAwareStreamInput;
|
|
|
import org.elasticsearch.common.io.stream.NamedWriteableRegistry;
|
|
|
import org.elasticsearch.common.io.stream.StreamInput;
|
|
|
@@ -41,7 +40,6 @@ import org.hamcrest.BaseMatcher;
|
|
|
import org.hamcrest.Description;
|
|
|
import org.junit.After;
|
|
|
|
|
|
-import java.io.InputStream;
|
|
|
import java.nio.ByteBuffer;
|
|
|
import java.util.ArrayList;
|
|
|
import java.util.Base64;
|
|
|
@@ -312,8 +310,12 @@ public class AsyncEqlSearchActionIT extends AbstractEqlBlockingIntegTestCase {
|
|
|
String value = doc.getSource().get("result").toString();
|
|
|
try (ByteBufferStreamInput buf = new ByteBufferStreamInput(ByteBuffer.wrap(Base64.getDecoder().decode(value)))) {
|
|
|
TransportVersion version = TransportVersion.readVersion(buf);
|
|
|
- final InputStream compressedIn = CompressorFactory.COMPRESSOR.threadLocalInputStream(buf);
|
|
|
- try (StreamInput in = new NamedWriteableAwareStreamInput(new InputStreamStreamInput(compressedIn), registry)) {
|
|
|
+ try (
|
|
|
+ StreamInput in = new NamedWriteableAwareStreamInput(
|
|
|
+ CompressorFactory.COMPRESSOR.threadLocalStreamInput(buf),
|
|
|
+ registry
|
|
|
+ )
|
|
|
+ ) {
|
|
|
in.setTransportVersion(version);
|
|
|
return new StoredAsyncResponse<>(EqlSearchResponse::new, in);
|
|
|
}
|