Add more exchanges + move some files

This commit is contained in:
Christopher Schnick
2022-01-18 09:55:43 +01:00
parent 774689e42a
commit caafb9d850
28 changed files with 775 additions and 76 deletions
@@ -23,7 +23,7 @@ import java.util.Optional;
import static io.xpipe.beacon.BeaconConfig.BODY_SEPARATOR;
public class BeaconClient {
public class BeaconClient implements AutoCloseable {
@FunctionalInterface
public interface FailableBiConsumer<T, U, E extends Throwable> {
@@ -76,7 +76,7 @@ public class BeaconClient {
public <REQ extends RequestMessage, RES extends ResponseMessage> void exchange(
REQ req,
FailableConsumer<OutputStream, IOException> reqWriter,
FailableBiPredicate<RES, InputStream, IOException> resReader)
FailableBiConsumer<RES, InputStream, IOException> resReader)
throws ConnectorException, ClientException, ServerException {
try {
sendRequest(req);
@@ -91,23 +91,16 @@ public class BeaconClient {
throw new ConnectorException("Invalid body separator");
}
if (resReader.test(res, in)) {
close();
}
resReader.accept(res, in);
} catch (IOException ex) {
close();
throw new ConnectorException("Couldn't communicate with socket", ex);
}
}
public <REQ extends RequestMessage, RES extends ResponseMessage> RES simpleExchange(REQ req)
throws ServerException, ConnectorException, ClientException {
try {
sendRequest(req);
return this.receiveResponse();
} finally {
close();
}
throws ServerException, ConnectorException, ClientException {
sendRequest(req);
return this.receiveResponse();
}
private <T extends RequestMessage> void sendRequest(T req) throws ClientException, ConnectorException {
@@ -16,7 +16,7 @@ public abstract class BeaconConnector {
protected <REQ extends RequestMessage, RES extends ResponseMessage> void performInputExchange(
BeaconClient socket,
REQ req,
BeaconClient.FailableBiPredicate<RES, InputStream, IOException> responseConsumer) throws ServerException, ConnectorException, ClientException {
BeaconClient.FailableBiConsumer<RES, InputStream, IOException> responseConsumer) throws ServerException, ConnectorException, ClientException {
performInputOutputExchange(socket, req, null, responseConsumer);
}
@@ -24,7 +24,7 @@ public abstract class BeaconConnector {
BeaconClient socket,
REQ req,
BeaconClient.FailableConsumer<OutputStream, IOException> reqWriter,
BeaconClient.FailableBiPredicate<RES, InputStream, IOException> responseConsumer)
BeaconClient.FailableBiConsumer<RES, InputStream, IOException> responseConsumer)
throws ServerException, ConnectorException, ClientException {
socket.exchange(req, reqWriter, responseConsumer);
}
@@ -37,7 +37,6 @@ public abstract class BeaconConnector {
AtomicReference<RES> response = new AtomicReference<>();
socket.exchange(req, reqWriter, (RES res, InputStream in) -> {
response.set(res);
return true;
});
return response.get();
}
@@ -0,0 +1,43 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceConfigInstance;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class DialogExchange implements MessageExchange<DialogExchange.Request, DialogExchange.Response> {
@Override
public String getId() {
return "dialog";
}
@Override
public Class<DialogExchange.Request> getRequestClass() {
return DialogExchange.Request.class;
}
@Override
public Class<DialogExchange.Response> getResponseClass() {
return DialogExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
DataSourceConfigInstance instance;
String key;
String value;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
DataSourceConfigInstance instance;
String errorMsg;
}
}
@@ -0,0 +1,45 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceConfigInstance;
import io.xpipe.core.source.DataSourceId;
import io.xpipe.core.source.DataSourceInfo;
import io.xpipe.core.store.DataStore;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class InfoExchange implements MessageExchange<InfoExchange.Request, InfoExchange.Response> {
@Override
public String getId() {
return "info";
}
@Override
public Class<InfoExchange.Request> getRequestClass() {
return InfoExchange.Request.class;
}
@Override
public Class<InfoExchange.Response> getResponseClass() {
return InfoExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
DataSourceId id;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
DataSourceInfo info;
DataStore store;
DataSourceConfigInstance config;
}
}
@@ -0,0 +1,47 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceConfigInstance;
import io.xpipe.core.source.DataSourceId;
import io.xpipe.core.store.DataStore;
import lombok.Builder;
import lombok.NonNull;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class ReadExecuteExchange implements MessageExchange<ReadExecuteExchange.Request, ReadExecuteExchange.Response> {
@Override
public String getId() {
return "readExecute";
}
@Override
public Class<ReadExecuteExchange.Request> getRequestClass() {
return ReadExecuteExchange.Request.class;
}
@Override
public Class<ReadExecuteExchange.Response> getResponseClass() {
return ReadExecuteExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
@NonNull
DataStore dataStore;
@NonNull
DataSourceConfigInstance config;
@NonNull
DataSourceId targetId;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
}
}
@@ -0,0 +1,43 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceConfigInstance;
import io.xpipe.core.store.DataStore;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class ReadPreparationExchange implements MessageExchange<ReadPreparationExchange.Request, ReadPreparationExchange.Response> {
@Override
public String getId() {
return "readPreparation";
}
@Override
public Class<ReadPreparationExchange.Request> getRequestClass() {
return ReadPreparationExchange.Request.class;
}
@Override
public Class<ReadPreparationExchange.Response> getResponseClass() {
return ReadPreparationExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
String providerType;
String dataStore;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
DataSourceConfigInstance config;
DataStore dataStore;
}
}
@@ -0,0 +1,39 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceId;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class SelectExchange implements MessageExchange<SelectExchange.Request, SelectExchange.Response> {
@Override
public String getId() {
return "select";
}
@Override
public Class<SelectExchange.Request> getRequestClass() {
return SelectExchange.Request.class;
}
@Override
public Class<SelectExchange.Response> getResponseClass() {
return SelectExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
DataSourceId id;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
}
}
@@ -0,0 +1,47 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceConfigInstance;
import io.xpipe.core.source.DataSourceId;
import io.xpipe.core.store.DataStore;
import lombok.Builder;
import lombok.NonNull;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class WriteExecuteExchange implements MessageExchange<WriteExecuteExchange.Request, WriteExecuteExchange.Response> {
@Override
public String getId() {
return "writeExecute";
}
@Override
public Class<WriteExecuteExchange.Request> getRequestClass() {
return WriteExecuteExchange.Request.class;
}
@Override
public Class<WriteExecuteExchange.Response> getResponseClass() {
return WriteExecuteExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
@NonNull
DataSourceId sourceId;
@NonNull
DataStore dataStore;
@NonNull
DataSourceConfigInstance config;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
}
}
@@ -0,0 +1,50 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceConfigInstance;
import io.xpipe.core.source.DataSourceId;
import io.xpipe.core.store.DataStore;
import lombok.Builder;
import lombok.NonNull;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class WritePreparationExchange implements MessageExchange<WritePreparationExchange.Request, WritePreparationExchange.Response> {
@Override
public String getId() {
return "writePreparation";
}
@Override
public Class<WritePreparationExchange.Request> getRequestClass() {
return WritePreparationExchange.Request.class;
}
@Override
public Class<WritePreparationExchange.Response> getResponseClass() {
return WritePreparationExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
String providerType;
String output;
@NonNull
DataSourceId sourceId;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
@NonNull
DataStore dataStore;
@NonNull
DataSourceConfigInstance config;
}
}
+7
View File
@@ -25,5 +25,12 @@ module io.xpipe.beacon {
StatusExchange,
StopExchange,
StoreResourceExchange,
WritePreparationExchange,
WriteExecuteExchange,
SelectExchange,
ReadPreparationExchange,
ReadExecuteExchange,
DialogExchange,
InfoExchange,
VersionExchange;
}