diff --git a/src/main/java/org/sensorhub/oshconnect/OSHSystem.java b/src/main/java/org/sensorhub/oshconnect/OSHSystem.java index ca64a81..809bea1 100644 --- a/src/main/java/org/sensorhub/oshconnect/OSHSystem.java +++ b/src/main/java/org/sensorhub/oshconnect/OSHSystem.java @@ -15,6 +15,8 @@ import org.sensorhub.oshconnect.util.DataStreamsQueryBuilder; import org.sensorhub.oshconnect.util.Utilities; +import java.time.Duration; +import java.time.Instant; import java.util.*; import java.util.concurrent.ExecutionException; @@ -99,6 +101,14 @@ public List discoverDataStreams(String query) throws ExecutionExc return result; } + public boolean isLive(Duration window) throws ExecutionException, InterruptedException { + var cutoff = Instant.now().minus(window); + return getConnectedSystemsApiClientExtras() + .getLastObservationTimes(getId(), "").get() + .values().stream() + .anyMatch(t -> t != null && t.isAfter(cutoff)); + } + /** * Add or update a data stream in the list of data streams and notify listeners. * diff --git a/src/main/java/org/sensorhub/oshconnect/net/ConSysApiClientExtras.java b/src/main/java/org/sensorhub/oshconnect/net/ConSysApiClientExtras.java index 854116c..438a3ee 100644 --- a/src/main/java/org/sensorhub/oshconnect/net/ConSysApiClientExtras.java +++ b/src/main/java/org/sensorhub/oshconnect/net/ConSysApiClientExtras.java @@ -24,12 +24,17 @@ import java.io.*; import java.net.*; +import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import java.util.function.Function; +import java.time.Instant; +import java.time.format.DateTimeParseException; +import java.util.LinkedHashMap; +import java.util.Map; public class ConSysApiClientExtras { static final String JSON_ARRAY_ITEMS = "items"; @@ -136,6 +141,45 @@ public CompletableFuture> getDataStreamIds(String systemID, String }); } + public CompletableFuture> getLastObservationTimes(String systemID, String queryString) { + if (queryString == null) + queryString = ""; + + return sendGetRequest(endpoint.resolve(SYSTEMS_COLLECTION + "/" + systemID + "/" + DATASTREAMS_COLLECTION + queryString), ResourceFormat.JSON, body -> { + try { + var ctx = new RequestContext(body); + JsonObject bodyJson = JsonParser.parseReader(new InputStreamReader(ctx.getInputStream())).getAsJsonObject(); + JsonArray features = bodyJson.getAsJsonArray(JSON_ARRAY_ITEMS); + + Map times = new LinkedHashMap<>(); + for (var feature : features) { + var obj = feature.getAsJsonObject(); + if (!obj.has("id")) + throw new IOException("No id found in feature"); + times.put(obj.get("id").getAsString(), parseExtentEnd(obj.getAsJsonArray("phenomenonTime"))); + } + return times; + } catch (IOException e) { + throw new CompletionException(e); + } + }); + } + + private static Instant parseExtentEnd(JsonArray extent) { + if (extent == null || extent.size() < 2) + return null; + var end = extent.get(1).getAsString(); + if ("0".equals(end)) // no observations yet + return null; + if ("now".equals(end)) // defensive: your node uses this on validTime, not phenomenonTime + return Instant.now(); + try { + return Instant.parse(end); + } catch (DateTimeParseException e) { + return null; + } + } + /** * Update a data stream. * @@ -290,7 +334,7 @@ public CompletableFuture> getObservations(String dataStrea } return observations; - } catch (IOException e) { + } catch (Exception e) { throw new CompletionException(e); } }); diff --git a/src/main/java/org/sensorhub/oshconnect/net/websocket/WebSocketConnection.java b/src/main/java/org/sensorhub/oshconnect/net/websocket/WebSocketConnection.java index df77218..52ad1fd 100644 --- a/src/main/java/org/sensorhub/oshconnect/net/websocket/WebSocketConnection.java +++ b/src/main/java/org/sensorhub/oshconnect/net/websocket/WebSocketConnection.java @@ -78,6 +78,8 @@ public void checkServerTrusted(X509Certificate[] chain, String authType) {} sslContextFactory.getSslContext().getClientSessionContext().setSessionCacheSize(0); client = new WebSocketClient(sslContextFactory); + client.getPolicy().setMaxBinaryMessageSize(4 * 1024 * 1024); // Jetty 9.x +// client.setMaxBinaryMessageSize(4 * 1024 * 1024); // Jetty 10/11+ client.start(); client.connect(this, new URI(urlString), clientUpgradeRequest);