Refactor various components

This commit is contained in:
Christopher Schnick
2022-01-01 20:11:13 +01:00
parent b54d8ad362
commit d63882c5ff
45 changed files with 771 additions and 190 deletions
@@ -12,9 +12,6 @@ import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.beacon.message.ServerErrorMessage;
import io.xpipe.core.util.JacksonHelper;
import org.apache.commons.lang3.function.FailableBiConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.InputStream;
@@ -23,13 +20,34 @@ import java.net.InetAddress;
import java.net.Socket;
import java.util.Arrays;
import java.util.Optional;
import java.util.function.Consumer;
import static io.xpipe.beacon.BeaconConfig.BODY_SEPARATOR;
public class BeaconClient {
private static final Logger logger = LoggerFactory.getLogger(BeaconClient.class);
@FunctionalInterface
public interface FailableBiConsumer<T, U, E extends Throwable> {
void accept(T var1, U var2) throws E;
}
@FunctionalInterface
public interface FailableConsumer<T, E extends Throwable> {
void accept(T var1) throws E;
}
public static Optional<BeaconClient> tryConnect() {
if (BeaconConfig.debugEnabled()) {
System.out.println("Attempting connection to server at port " + BeaconConfig.getUsedPort());
}
try {
return Optional.of(new BeaconClient());
} catch (IOException ex) {
return Optional.empty();
}
}
private final Socket socket;
private final InputStream in;
@@ -51,29 +69,27 @@ public class BeaconClient {
public <REQ extends RequestMessage, RES extends ResponseMessage> void exchange(
REQ req,
Consumer<OutputStream> output,
FailableBiConsumer<RES, InputStream, IOException> responseConsumer,
boolean keepOpen) throws ConnectorException, ClientException, ServerException {
FailableConsumer<OutputStream, IOException> reqWriter,
FailableBiConsumer<RES, InputStream, IOException> resReader)
throws ConnectorException, ClientException, ServerException {
try {
sendRequest(req);
if (output != null) {
if (reqWriter != null) {
out.write(BODY_SEPARATOR);
output.accept(out);
reqWriter.accept(out);
}
var res = this.<RES>receiveResponse();
var sep = in.readNBytes(BODY_SEPARATOR.length);
if (!Arrays.equals(BODY_SEPARATOR, sep)) {
if (sep.length != 0 && !Arrays.equals(BODY_SEPARATOR, sep)) {
throw new ConnectorException("Invalid body separator");
}
responseConsumer.accept(res, in);
resReader.accept(res, in);
} catch (IOException ex) {
throw new ConnectorException("Couldn't communicate with socket", ex);
} finally {
if (!keepOpen) {
close();
}
close();
}
}
@@ -2,52 +2,42 @@ package io.xpipe.beacon;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import org.apache.commons.lang3.function.FailableBiConsumer;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.function.Consumer;
import java.util.concurrent.atomic.AtomicReference;
public abstract class BeaconConnector {
protected abstract void waitForStartup();
protected abstract BeaconClient constructSocket() throws ConnectorException;
protected BeaconClient constructSocket() throws ConnectorException {
if (!BeaconServer.isRunning()) {
try {
BeaconServer.start();
waitForStartup();
if (!BeaconServer.isRunning()) {
throw new ConnectorException("Unable to start xpipe daemon");
}
} catch (Exception ex) {
throw new ConnectorException("Unable to start xpipe daemon: " + ex.getMessage());
}
}
try {
return new BeaconClient();
} catch (Exception ex) {
throw new ConnectorException("Unable to connect to running xpipe daemon: " + ex.getMessage());
}
}
protected <REQ extends RequestMessage, RES extends ResponseMessage> void performExchange(
protected <REQ extends RequestMessage, RES extends ResponseMessage> void performInputExchange(
BeaconClient socket,
REQ req,
FailableBiConsumer<RES, InputStream, IOException> responseConsumer,
boolean keepOpen) throws ServerException, ConnectorException, ClientException {
performExchange(socket, req, null, responseConsumer, keepOpen);
BeaconClient.FailableBiConsumer<RES, InputStream, IOException> responseConsumer) throws ServerException, ConnectorException, ClientException {
performInputOutputExchange(socket, req, null, responseConsumer);
}
protected <REQ extends RequestMessage, RES extends ResponseMessage> void performExchange(
protected <REQ extends RequestMessage, RES extends ResponseMessage> void performInputOutputExchange(
BeaconClient socket,
REQ req,
Consumer<OutputStream> output,
FailableBiConsumer<RES, InputStream, IOException> responseConsumer,
boolean keepOpen) throws ServerException, ConnectorException, ClientException {
socket.exchange(req, output, responseConsumer, keepOpen);
BeaconClient.FailableConsumer<OutputStream, IOException> reqWriter,
BeaconClient.FailableBiConsumer<RES, InputStream, IOException> responseConsumer)
throws ServerException, ConnectorException, ClientException {
socket.exchange(req, reqWriter, responseConsumer);
}
protected <REQ extends RequestMessage, RES extends ResponseMessage> RES performOutputExchange(
BeaconClient socket,
REQ req,
BeaconClient.FailableConsumer<OutputStream, IOException> reqWriter)
throws ServerException, ConnectorException, ClientException {
AtomicReference<RES> response = new AtomicReference<>();
socket.exchange(req, reqWriter, (RES res, InputStream in) -> {
response.set(res);
});
return response.get();
}
protected <REQ extends RequestMessage, RES extends ResponseMessage> RES performSimpleExchange(
@@ -10,6 +10,8 @@ public interface BeaconHandler {
void prepareBody() throws IOException;
void startBodyRead() throws IOException;
public <T extends ResponseMessage> void sendResponse(T obj) throws Exception;
public void sendClientErrorResponse(String message) throws Exception;
@@ -20,19 +20,19 @@ public class BeaconServer {
return !isPortAvailable(port);
}
public static void start() throws Exception {
public static boolean tryStart() throws Exception {
if (BeaconConfig.shouldStartInProcess()) {
startInProcess();
return;
return true;
}
var custom = BeaconConfig.getCustomExecCommand();
if (custom != null) {
Runtime.getRuntime().exec(System.getenv(custom));
return;
var proc = new ProcessBuilder("cmd", "/c", "CALL", "C:\\Projects\\xpipe\\xpipe\\gradlew.bat", ":app:run").inheritIO().start();
return true;
}
throw new IllegalArgumentException("Unable to start xpipe daemon");
return false;
}
private static void startInProcess() throws Exception {
@@ -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.data.type.DataType;
import io.xpipe.core.source.DataSourceId;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class CliOptionPageExchange implements MessageExchange<CliOptionPageExchange.Request, CliOptionPageExchange.Response> {
@Override
public String getId() {
return "cliOptionPage";
}
@Override
public Class<CliOptionPageExchange.Request> getRequestClass() {
return CliOptionPageExchange.Request.class;
}
@Override
public Class<CliOptionPageExchange.Response> getResponseClass() {
return CliOptionPageExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
DataSourceId newSourceId;
String type;
boolean hasInputStream;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
DataSourceId sourceId;
DataType dataType;
int rowCount;
}
}
@@ -2,11 +2,14 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
import java.time.Instant;
import java.util.List;
public abstract class ListCollectionsExchange implements MessageExchange<ListCollectionsExchange.Request, ListCollectionsExchange.Response> {
public class ListCollectionsExchange implements MessageExchange<ListCollectionsExchange.Request, ListCollectionsExchange.Response> {
@Override
public String getId() {
@@ -23,7 +26,10 @@ public abstract class ListCollectionsExchange implements MessageExchange<ListCol
return Response.class;
}
public static record Request() implements RequestMessage {
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
}
@@ -31,7 +37,10 @@ public abstract class ListCollectionsExchange implements MessageExchange<ListCol
}
public static record Response(List<Entry> entries) implements ResponseMessage {
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
List<Entry> entries;
}
}
@@ -2,11 +2,14 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
import java.time.Instant;
import java.util.List;
public abstract class ListEntriesExchange implements MessageExchange<ListEntriesExchange.Request, ListEntriesExchange.Response> {
public class ListEntriesExchange implements MessageExchange<ListEntriesExchange.Request, ListEntriesExchange.Response> {
@Override
public String getId() {
@@ -23,15 +26,21 @@ public abstract class ListEntriesExchange implements MessageExchange<ListEntries
return Response.class;
}
public static record Request(String collection) implements RequestMessage {
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
String collection;
}
public static record Entry(String name, String type, String description, Instant lastUsed) {
}
public static record Response(List<Entry> entries) implements ResponseMessage {
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
List<Entry> entries;
}
}
@@ -2,9 +2,6 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.beacon.BeaconHandler;
import java.io.InputStream;
public interface MessageExchange<RQ extends RequestMessage, RP extends ResponseMessage> {
@@ -13,6 +10,4 @@ public interface MessageExchange<RQ extends RequestMessage, RP extends ResponseM
Class<RQ> getRequestClass();
Class<RP> getResponseClass();
void handleRequest(BeaconHandler handler, RQ msg, InputStream body) throws Exception;
}
@@ -2,8 +2,11 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public abstract class ModeExchange implements MessageExchange<ModeExchange.Request, ModeExchange.Response> {
public class ModeExchange implements MessageExchange<ModeExchange.Request, ModeExchange.Response> {
@Override
public String getId() {
@@ -20,11 +23,17 @@ public abstract class ModeExchange implements MessageExchange<ModeExchange.Reque
return ModeExchange.Response.class;
}
public static record Request(String modeId) implements RequestMessage {
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
String modeId;
}
public static record Response() implements ResponseMessage {
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
}
}
@@ -0,0 +1,56 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.data.type.DataType;
import io.xpipe.core.source.DataSourceId;
import io.xpipe.core.source.DataSourceType;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class ReadInfoExchange implements MessageExchange<ReadInfoExchange.Request, ReadInfoExchange.Response> {
@Override
public String getId() {
return "readTableInfo";
}
@Override
public Class<ReadInfoExchange.Request> getRequestClass() {
return ReadInfoExchange.Request.class;
}
@Override
public Class<ReadInfoExchange.Response> getResponseClass() {
return ReadInfoExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
DataSourceId sourceId;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
DataSourceId sourceId;
DataSourceType type;
Object data;
public TableData getTableData() {
return (TableData) data;
}
}
@Jacksonized
@Builder
@Value
public static class TableData {
DataType dataType;
int rowCount;
}
}
@@ -4,7 +4,7 @@ import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.source.DataSourceId;
public abstract class ReadStructureExchange implements MessageExchange<ReadStructureExchange.Request, ReadStructureExchange.Response> {
public class ReadStructureExchange implements MessageExchange<ReadStructureExchange.Request, ReadStructureExchange.Response> {
@Override
public String getId() {
@@ -3,8 +3,11 @@ 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 abstract class ReadTableDataExchange implements MessageExchange<ReadTableDataExchange.Request, ReadTableDataExchange.Response> {
public class ReadTableDataExchange implements MessageExchange<ReadTableDataExchange.Request, ReadTableDataExchange.Response> {
@Override
public String getId() {
@@ -21,11 +24,18 @@ public abstract class ReadTableDataExchange implements MessageExchange<ReadTable
return ReadTableDataExchange.Response.class;
}
public static record Request(DataSourceId sourceId, int maxRows) implements RequestMessage {
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
DataSourceId sourceId;
int maxRows;
}
public static record Response() implements ResponseMessage {
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
}
}
@@ -1,32 +0,0 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.core.data.type.DataType;
import io.xpipe.core.source.DataSourceId;
public abstract class ReadTableInfoExchange implements MessageExchange<ReadTableInfoExchange.Request, ReadTableInfoExchange.Response> {
@Override
public String getId() {
return "readTableInfo";
}
@Override
public Class<ReadTableInfoExchange.Request> getRequestClass() {
return ReadTableInfoExchange.Request.class;
}
@Override
public Class<ReadTableInfoExchange.Response> getResponseClass() {
return ReadTableInfoExchange.Response.class;
}
public static record Request(DataSourceId sourceId) implements RequestMessage {
}
public static record Response(DataSourceId sourceId, DataType dataType, int rowCount) implements ResponseMessage {
}
}
@@ -2,8 +2,11 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public abstract class StatusExchange implements MessageExchange<StatusExchange.Request, StatusExchange.Response> {
public class StatusExchange implements MessageExchange<StatusExchange.Request, StatusExchange.Response> {
@Override
public String getId() {
@@ -20,11 +23,16 @@ public abstract class StatusExchange implements MessageExchange<StatusExchange.R
return Response.class;
}
public static record Request() implements RequestMessage {
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
}
public static record Response(String mode) implements ResponseMessage {
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
String mode;
}
}
@@ -2,8 +2,11 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public abstract class StopExchange implements MessageExchange<StopExchange.Request, StopExchange.Response> {
public class StopExchange implements MessageExchange<StopExchange.Request, StopExchange.Response> {
@Override
public String getId() {
@@ -20,11 +23,17 @@ public abstract class StopExchange implements MessageExchange<StopExchange.Reque
return StopExchange.Response.class;
}
public static record Request(boolean force) implements RequestMessage {
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
boolean force;
}
public static record Response(boolean success) implements ResponseMessage {
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
boolean success;
}
}
@@ -0,0 +1,44 @@
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;
import java.util.Map;
import java.util.UUID;
public class StoreEndExchange implements MessageExchange<StoreEndExchange.Request, StoreEndExchange.Response> {
@Override
public String getId() {
return "storeEnd";
}
@Override
public Class<StoreEndExchange.Request> getRequestClass() {
return StoreEndExchange.Request.class;
}
@Override
public Class<StoreEndExchange.Response> getResponseClass() {
return StoreEndExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
UUID entryId;
Map<String, String> values;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
DataSourceId sourceId;
}
}
@@ -0,0 +1,41 @@
package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import io.xpipe.extension.cli.CliOptionPage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public class StoreStartExchange implements MessageExchange<StoreStartExchange.Request, StoreStartExchange.Response> {
@Override
public String getId() {
return "storeStart";
}
@Override
public Class<StoreStartExchange.Request> getRequestClass() {
return StoreStartExchange.Request.class;
}
@Override
public Class<StoreStartExchange.Response> getResponseClass() {
return StoreStartExchange.Response.class;
}
@Jacksonized
@Builder
@Value
public static class Request implements RequestMessage {
String type;
boolean hasInputStream;
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
CliOptionPage page;
}
}
@@ -2,8 +2,11 @@ package io.xpipe.beacon.exchange;
import io.xpipe.beacon.message.RequestMessage;
import io.xpipe.beacon.message.ResponseMessage;
import lombok.Builder;
import lombok.Value;
import lombok.extern.jackson.Jacksonized;
public abstract class VersionExchange implements MessageExchange<VersionExchange.Request, VersionExchange.Response> {
public class VersionExchange implements MessageExchange<VersionExchange.Request, VersionExchange.Response> {
@Override
public String getId() {
@@ -20,18 +23,20 @@ public abstract class VersionExchange implements MessageExchange<VersionExchange
return VersionExchange.Response.class;
}
public static record Request() implements RequestMessage {
@lombok.extern.jackson.Jacksonized
@lombok.Builder
@lombok.Value
public static class Request implements RequestMessage {
}
@Jacksonized
@Builder
@Value
public static class Response implements ResponseMessage {
public final String version;
public final String jvmVersion;
public Response(String version, String jvmVersion) {
this.version = version;
this.jvmVersion = jvmVersion;
}
String version;
String buildVersion;
String jvmVersion;
}
}
+11 -3
View File
@@ -1,22 +1,30 @@
import io.xpipe.beacon.exchange.MessageExchange;
import io.xpipe.beacon.exchange.*;
module io.xpipe.beacon {
exports io.xpipe.beacon;
exports io.xpipe.beacon.exchange;
exports io.xpipe.beacon.message;
requires org.slf4j;
requires org.slf4j.simple;
requires com.fasterxml.jackson.core;
requires com.fasterxml.jackson.databind;
requires com.fasterxml.jackson.module.paramnames;
requires io.xpipe.core;
requires io.xpipe.extension;
opens io.xpipe.beacon;
opens io.xpipe.beacon.exchange;
opens io.xpipe.beacon.message;
requires org.apache.commons.lang;
requires static lombok;
uses MessageExchange;
provides io.xpipe.beacon.exchange.MessageExchange with
ListCollectionsExchange,
ListEntriesExchange,
ReadTableDataExchange,
ReadInfoExchange,
StatusExchange,
StopExchange,
VersionExchange;
}