This commit is contained in:
Christopher Schnick committed 2022-11-27 14:59:36 +01:00
1 parent 696dc036ac
commit de70b0d5b0
26 files changed
+344 -327

No files matched your search

@@ -1,6 +1,6 @@
package io.xpipe.api;
import io.xpipe.api.connector.XPipeConnection;
import io.xpipe.api.connector.XPipeApiConnection;
import io.xpipe.api.util.QuietDialogHandler;
import io.xpipe.beacon.exchange.cli.StoreAddExchange;
import io.xpipe.core.store.DataStore;
@@ -10,7 +10,7 @@ import java.util.Map;
public class DataStores {
public static void addNamedStore(DataStore store, String name) {
XPipeConnection.execute(con -> {
XPipeApiConnection.execute(con -> {
var req = StoreAddExchange.Request.builder()
.storeInput(store)
.name(name)
@@ -1,23 +1,27 @@
package io.xpipe.api.connector;
import io.xpipe.beacon.*;
import io.xpipe.beacon.BeaconClient;
import io.xpipe.beacon.BeaconConnection;
import io.xpipe.beacon.BeaconException;
import io.xpipe.beacon.BeaconServer;
import io.xpipe.beacon.exchange.cli.DialogExchange;
import io.xpipe.core.dialog.DialogReference;
import io.xpipe.core.util.XPipeInstallation;
import java.util.Optional;
public final class XPipeConnection extends BeaconConnection {
public final class XPipeApiConnection extends BeaconConnection {
private XPipeConnection() {}
private XPipeApiConnection() {}
public static XPipeConnection open() {
var con = new XPipeConnection();
public static XPipeApiConnection open() {
var con = new XPipeApiConnection();
con.constructSocket();
return con;
}
public static void finishDialog(DialogReference reference) {
try (var con = new XPipeConnection()) {
try (var con = new XPipeApiConnection()) {
con.constructSocket();
var element = reference.getStart();
while (true) {
@@ -41,7 +45,7 @@ public final class XPipeConnection extends BeaconConnection {
}
public static void execute(Handler handler) {
try (var con = new XPipeConnection()) {
try (var con = new XPipeApiConnection()) {
con.constructSocket();
handler.handle(con);
} catch (BeaconException e) {
@@ -52,7 +56,7 @@ public final class XPipeConnection extends BeaconConnection {
}
public static <T> T execute(Mapper<T> mapper) {
try (var con = new XPipeConnection()) {
try (var con = new XPipeApiConnection()) {
con.constructSocket();
return mapper.handle(con);
} catch (BeaconException e) {
@@ -73,7 +77,10 @@ public final class XPipeConnection extends BeaconConnection {
} catch (InterruptedException ignored) {
}
var s = BeaconClient.tryConnect(BeaconClient.ApiClientInformation.builder().version("?").language("Java").build());
var s = BeaconClient.tryConnect(BeaconClient.ApiClientInformation.builder()
.version("?")
.language("Java")
.build());
if (s.isPresent()) {
return s;
}
@@ -114,17 +121,18 @@ public final class XPipeConnection extends BeaconConnection {
}
try {
beaconClient = BeaconClient.connect(BeaconClient.ApiClientInformation.builder().version("?").language("Java").build());
beaconClient = BeaconClient.connect(BeaconClient.ApiClientInformation.builder()
.version("?")
.language("Java")
.build());
} catch (Exception ex) {
throw new BeaconException("Unable to connect to running xpipe daemon", ex);
}
}
private void start() throws Exception {
if (BeaconServer.tryStart() == null) {
throw new UnsupportedOperationException("Unable to determine xpipe daemon launch command");
}
;
var installation = XPipeInstallation.getDefaultInstallationBasePath();
BeaconServer.start(installation);
}
@FunctionalInterface
@@ -2,7 +2,7 @@ package io.xpipe.api.impl;
import io.xpipe.api.DataSource;
import io.xpipe.api.DataSourceConfig;
import io.xpipe.api.connector.XPipeConnection;
import io.xpipe.api.connector.XPipeApiConnection;
import io.xpipe.beacon.exchange.*;
import io.xpipe.core.source.DataSourceId;
import io.xpipe.core.source.DataSourceReference;
@@ -25,7 +25,7 @@ public abstract class DataSourceImpl implements DataSource {
}
public static DataSource get(DataSourceReference ds) {
return XPipeConnection.execute(con -> {
return XPipeApiConnection.execute(con -> {
var req = QueryDataSourceExchange.Request.builder().ref(ds).build();
QueryDataSourceExchange.Response res = con.performSimpleExchange(req);
var config = new DataSourceConfig(res.getProvider(), res.getConfig());
@@ -53,7 +53,7 @@ public abstract class DataSourceImpl implements DataSource {
public static DataSource create(DataSourceId id, io.xpipe.core.source.DataSource<?> source) {
var startReq =
AddSourceExchange.Request.builder().source(source).target(id).build();
var returnedId = XPipeConnection.execute(con -> {
var returnedId = XPipeApiConnection.execute(con -> {
AddSourceExchange.Response r = con.performSimpleExchange(startReq);
return r.getId();
});
@@ -64,7 +64,7 @@ public abstract class DataSourceImpl implements DataSource {
public static DataSource create(DataSourceId id, String type, DataStore store) {
if (store instanceof StreamDataStore s && s.isContentExclusivelyAccessible()) {
var res = XPipeConnection.execute(con -> {
var res = XPipeApiConnection.execute(con -> {
var req = StoreStreamExchange.Request.builder().build();
StoreStreamExchange.Response r = con.performOutputExchange(req, out -> {
try (InputStream inputStream = s.openInput()) {
@@ -83,20 +83,20 @@ public abstract class DataSourceImpl implements DataSource {
.target(id)
.configureAll(false)
.build();
var startRes = XPipeConnection.execute(con -> {
var startRes = XPipeApiConnection.execute(con -> {
ReadExchange.Response r = con.performSimpleExchange(startReq);
return r;
});
var configInstance = startRes.getConfig();
XPipeConnection.finishDialog(configInstance);
XPipeApiConnection.finishDialog(configInstance);
var ref = id != null ? DataSourceReference.id(id) : DataSourceReference.latest();
return get(ref);
}
public static DataSource create(DataSourceId id, String type, InputStream in) {
var res = XPipeConnection.execute(con -> {
var res = XPipeApiConnection.execute(con -> {
var req = StoreStreamExchange.Request.builder().build();
StoreStreamExchange.Response r = con.performOutputExchange(req, out -> in.transferTo(out));
return r;
@@ -110,13 +110,13 @@ public abstract class DataSourceImpl implements DataSource {
.target(id)
.configureAll(false)
.build();
var startRes = XPipeConnection.execute(con -> {
var startRes = XPipeApiConnection.execute(con -> {
ReadExchange.Response r = con.performSimpleExchange(startReq);
return r;
});
var configInstance = startRes.getConfig();
XPipeConnection.finishDialog(configInstance);
XPipeApiConnection.finishDialog(configInstance);
var ref = id != null ? DataSourceReference.id(id) : DataSourceReference.latest();
return get(ref);
@@ -124,7 +124,7 @@ public abstract class DataSourceImpl implements DataSource {
@Override
public void forwardTo(DataSource target) {
XPipeConnection.execute(con -> {
XPipeApiConnection.execute(con -> {
var req = ForwardExchange.Request.builder()
.source(DataSourceReference.id(sourceId))
.target(DataSourceReference.id(target.getId()))
@@ -135,7 +135,7 @@ public abstract class DataSourceImpl implements DataSource {
@Override
public void appendTo(DataSource target) {
XPipeConnection.execute(con -> {
XPipeApiConnection.execute(con -> {
var req = ForwardExchange.Request.builder()
.source(DataSourceReference.id(sourceId))
.target(DataSourceReference.id(target.getId()))
@@ -3,7 +3,7 @@ package io.xpipe.api.impl;
import io.xpipe.api.DataSource;
import io.xpipe.api.DataTable;
import io.xpipe.api.DataTableAccumulator;
import io.xpipe.api.connector.XPipeConnection;
import io.xpipe.api.connector.XPipeApiConnection;
import io.xpipe.api.util.TypeDescriptor;
import io.xpipe.beacon.BeaconException;
import io.xpipe.beacon.exchange.ReadExchange;
@@ -22,7 +22,7 @@ import java.nio.charset.StandardCharsets;
public class DataTableAccumulatorImpl implements DataTableAccumulator {
private final XPipeConnection connection;
private final XPipeApiConnection connection;
private final TupleType type;
private int rows;
private TupleType writtenDescriptor;
@@ -30,7 +30,7 @@ public class DataTableAccumulatorImpl implements DataTableAccumulator {
public DataTableAccumulatorImpl(TupleType type) {
this.type = type;
connection = XPipeConnection.open();
connection = XPipeApiConnection.open();
connection.sendRequest(StoreStreamExchange.Request.builder().build());
bodyOutput = connection.sendBody();
}
@@ -52,12 +52,12 @@ public class DataTableAccumulatorImpl implements DataTableAccumulator {
.provider("xpbt")
.configureAll(false)
.build();
ReadExchange.Response response = XPipeConnection.execute(con -> {
ReadExchange.Response response = XPipeApiConnection.execute(con -> {
return con.performSimpleExchange(req);
});
var configInstance = response.getConfig();
XPipeConnection.finishDialog(configInstance);
XPipeApiConnection.finishDialog(configInstance);
return DataSource.get(DataSourceReference.id(id)).asTable();
}
@@ -2,7 +2,7 @@ package io.xpipe.api.impl;
import io.xpipe.api.DataSourceConfig;
import io.xpipe.api.DataTable;
import io.xpipe.api.connector.XPipeConnection;
import io.xpipe.api.connector.XPipeApiConnection;
import io.xpipe.beacon.BeaconConnection;
import io.xpipe.beacon.BeaconException;
import io.xpipe.beacon.exchange.api.QueryTableDataExchange;
@@ -55,7 +55,7 @@ public class DataTableImpl extends DataSourceImpl implements DataTable {
@Override
public ArrayNode read(int maxRows) {
List<DataStructureNode> nodes = new ArrayList<>();
XPipeConnection.execute(con -> {
XPipeApiConnection.execute(con -> {
var req = QueryTableDataExchange.Request.builder()
.ref(DataSourceReference.id(getId()))
.maxRows(maxRows)
@@ -83,7 +83,7 @@ public class DataTableImpl extends DataSourceImpl implements DataTable {
private TupleNode node;
{
connection = XPipeConnection.open();
connection = XPipeApiConnection.open();
var req = QueryTableDataExchange.Request.builder()
.ref(DataSourceReference.id(getId()))
.maxRows(Integer.MAX_VALUE)
@@ -2,7 +2,7 @@ package io.xpipe.api.impl;
import io.xpipe.api.DataSourceConfig;
import io.xpipe.api.DataText;
import io.xpipe.api.connector.XPipeConnection;
import io.xpipe.api.connector.XPipeApiConnection;
import io.xpipe.beacon.BeaconConnection;
import io.xpipe.beacon.BeaconException;
import io.xpipe.beacon.exchange.api.QueryTextDataExchange;
@@ -62,7 +62,7 @@ public class DataTextImpl extends DataSourceImpl implements DataText {
private String nextValue;
{
connection = XPipeConnection.open();
connection = XPipeApiConnection.open();
var req = QueryTextDataExchange.Request.builder()
.ref(DataSourceReference.id(getId()))
.maxLines(-1)
@@ -1,46 +0,0 @@
package io.xpipe.api.util;
import io.xpipe.api.connector.XPipeConnection;
import io.xpipe.beacon.BeaconClient;
import io.xpipe.beacon.BeaconServer;
public class XPipeDaemonController {
private static boolean alreadyStarted;
public static void start() throws Exception {
if (BeaconServer.isRunning()) {
alreadyStarted = true;
return;
}
Process process = null;
if ((process = BeaconServer.tryStartCustom()) != null) {
} else {
if ((process = BeaconServer.tryStart()) == null) {
throw new AssertionError();
}
}
XPipeConnection.waitForStartup(process).orElseThrow();
if (!BeaconServer.isRunning()) {
throw new AssertionError();
}
}
public static void stop() throws Exception {
if (alreadyStarted) {
return;
}
if (!BeaconServer.isRunning()) {
return;
}
var client = BeaconClient.connect(BeaconClient.ApiClientInformation.builder().version("?").language("Java API Test").build());
if (!BeaconServer.tryStop(client)) {
throw new AssertionError();
}
XPipeConnection.waitForShutdown();
}
}