Merge pull request 'clients/java: SDK methods for AcknowledgeAlarm + QueryActiveAlarms (PR E.5)' (#109) from track-e5-java-alarm-sdk into main
This commit was merged in pull request #109.
This commit is contained in:
+2
-2
@@ -32,7 +32,7 @@ final class MxGatewayCliTests {
|
|||||||
assertEquals(0, run.exitCode());
|
assertEquals(0, run.exitCode());
|
||||||
assertEquals("", run.errors());
|
assertEquals("", run.errors());
|
||||||
assertTrue(run.output().contains("mxgateway-java 0.1.0"));
|
assertTrue(run.output().contains("mxgateway-java 0.1.0"));
|
||||||
assertTrue(run.output().contains("gatewayProtocolVersion=1"));
|
assertTrue(run.output().contains("gatewayProtocolVersion=3"));
|
||||||
assertTrue(run.output().contains("workerProtocolVersion=1"));
|
assertTrue(run.output().contains("workerProtocolVersion=1"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -42,7 +42,7 @@ final class MxGatewayCliTests {
|
|||||||
|
|
||||||
assertEquals(0, run.exitCode());
|
assertEquals(0, run.exitCode());
|
||||||
assertTrue(run.output().contains("\"clientVersion\":\"0.1.0\""));
|
assertTrue(run.output().contains("\"clientVersion\":\"0.1.0\""));
|
||||||
assertTrue(run.output().contains("\"gatewayProtocolVersion\":1"));
|
assertTrue(run.output().contains("\"gatewayProtocolVersion\":3"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
+67
@@ -0,0 +1,67 @@
|
|||||||
|
package com.dohertylan.mxgateway.client;
|
||||||
|
|
||||||
|
import io.grpc.stub.ClientCallStreamObserver;
|
||||||
|
import io.grpc.stub.ClientResponseObserver;
|
||||||
|
import io.grpc.stub.StreamObserver;
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
|
import mxaccess_gateway.v1.MxaccessGateway.ActiveAlarmSnapshot;
|
||||||
|
import mxaccess_gateway.v1.MxaccessGateway.QueryActiveAlarmsRequest;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Cancellable handle returned by {@code queryActiveAlarms}.
|
||||||
|
*
|
||||||
|
* <p>Wraps a caller-supplied {@link StreamObserver} and exposes a
|
||||||
|
* {@link #cancel()} entry point that aborts the underlying gRPC call. The
|
||||||
|
* subscription also implements {@link AutoCloseable} so it can participate in
|
||||||
|
* try-with-resources blocks.
|
||||||
|
*/
|
||||||
|
public final class MxGatewayActiveAlarmsSubscription implements AutoCloseable {
|
||||||
|
private final AtomicReference<ClientCallStreamObserver<QueryActiveAlarmsRequest>> requestStream = new AtomicReference<>();
|
||||||
|
private final AtomicBoolean cancelled = new AtomicBoolean();
|
||||||
|
|
||||||
|
ClientResponseObserver<QueryActiveAlarmsRequest, ActiveAlarmSnapshot> wrap(StreamObserver<ActiveAlarmSnapshot> observer) {
|
||||||
|
return new ClientResponseObserver<>() {
|
||||||
|
@Override
|
||||||
|
public void beforeStart(ClientCallStreamObserver<QueryActiveAlarmsRequest> stream) {
|
||||||
|
requestStream.set(stream);
|
||||||
|
if (cancelled.get()) {
|
||||||
|
stream.cancel("client cancelled active-alarms query", null);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onNext(ActiveAlarmSnapshot value) {
|
||||||
|
observer.onNext(value);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onError(Throwable error) {
|
||||||
|
observer.onError(error);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onCompleted() {
|
||||||
|
observer.onCompleted();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Cancels the underlying gRPC call. Safe to invoke before the call has
|
||||||
|
* started; cancellation is recorded and applied as soon as the stream
|
||||||
|
* attaches.
|
||||||
|
*/
|
||||||
|
public void cancel() {
|
||||||
|
cancelled.set(true);
|
||||||
|
ClientCallStreamObserver<QueryActiveAlarmsRequest> stream = requestStream.get();
|
||||||
|
if (stream != null) {
|
||||||
|
stream.cancel("client cancelled active-alarms query", null);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void close() {
|
||||||
|
cancel();
|
||||||
|
}
|
||||||
|
}
|
||||||
+61
@@ -15,6 +15,9 @@ import java.util.concurrent.CompletableFuture;
|
|||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
import javax.net.ssl.SSLException;
|
import javax.net.ssl.SSLException;
|
||||||
import mxaccess_gateway.v1.MxAccessGatewayGrpc;
|
import mxaccess_gateway.v1.MxAccessGatewayGrpc;
|
||||||
|
import mxaccess_gateway.v1.MxaccessGateway.AcknowledgeAlarmReply;
|
||||||
|
import mxaccess_gateway.v1.MxaccessGateway.AcknowledgeAlarmRequest;
|
||||||
|
import mxaccess_gateway.v1.MxaccessGateway.ActiveAlarmSnapshot;
|
||||||
import mxaccess_gateway.v1.MxaccessGateway.CloseSessionReply;
|
import mxaccess_gateway.v1.MxaccessGateway.CloseSessionReply;
|
||||||
import mxaccess_gateway.v1.MxaccessGateway.CloseSessionRequest;
|
import mxaccess_gateway.v1.MxaccessGateway.CloseSessionRequest;
|
||||||
import mxaccess_gateway.v1.MxaccessGateway.MxCommandReply;
|
import mxaccess_gateway.v1.MxaccessGateway.MxCommandReply;
|
||||||
@@ -23,6 +26,7 @@ import mxaccess_gateway.v1.MxaccessGateway.MxEvent;
|
|||||||
import mxaccess_gateway.v1.MxaccessGateway.OpenSessionReply;
|
import mxaccess_gateway.v1.MxaccessGateway.OpenSessionReply;
|
||||||
import mxaccess_gateway.v1.MxaccessGateway.OpenSessionRequest;
|
import mxaccess_gateway.v1.MxaccessGateway.OpenSessionRequest;
|
||||||
import mxaccess_gateway.v1.MxaccessGateway.ProtocolStatusCode;
|
import mxaccess_gateway.v1.MxaccessGateway.ProtocolStatusCode;
|
||||||
|
import mxaccess_gateway.v1.MxaccessGateway.QueryActiveAlarmsRequest;
|
||||||
import mxaccess_gateway.v1.MxaccessGateway.StreamEventsRequest;
|
import mxaccess_gateway.v1.MxaccessGateway.StreamEventsRequest;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -259,6 +263,63 @@ public final class MxGatewayClient implements AutoCloseable {
|
|||||||
return subscription;
|
return subscription;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Acknowledges an active MXAccess alarm condition through the gateway.
|
||||||
|
*
|
||||||
|
* <p>The gateway authenticates the request against the API key's
|
||||||
|
* {@code invoke:alarm-ack} scope and forwards the acknowledge to the
|
||||||
|
* worker's MXAccess session; the resulting native MxStatus is returned
|
||||||
|
* in the reply. Acks are idempotent at the MxAccess layer.
|
||||||
|
*
|
||||||
|
* @param request the {@code AcknowledgeAlarmRequest}
|
||||||
|
* @return the raw acknowledge reply
|
||||||
|
* @throws MxGatewayException on transport or protocol failure
|
||||||
|
*/
|
||||||
|
public AcknowledgeAlarmReply acknowledgeAlarm(AcknowledgeAlarmRequest request) {
|
||||||
|
try {
|
||||||
|
AcknowledgeAlarmReply reply = rawBlockingStub().acknowledgeAlarm(request);
|
||||||
|
MxGatewayErrors.ensureProtocolSuccess("acknowledge alarm", reply.getProtocolStatus(), null);
|
||||||
|
return reply;
|
||||||
|
} catch (RuntimeException error) {
|
||||||
|
if (error instanceof MxGatewayException) {
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
throw MxGatewayErrors.fromGrpc("acknowledge alarm", error);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Acknowledges an active MXAccess alarm condition asynchronously.
|
||||||
|
*
|
||||||
|
* @param request the {@code AcknowledgeAlarmRequest}
|
||||||
|
* @return a future completed with the raw reply, or completed exceptionally
|
||||||
|
* with {@link MxGatewayException} on failure
|
||||||
|
*/
|
||||||
|
public CompletableFuture<AcknowledgeAlarmReply> acknowledgeAlarmAsync(AcknowledgeAlarmRequest request) {
|
||||||
|
CompletableFuture<AcknowledgeAlarmReply> future = toCompletable(rawFutureStub().acknowledgeAlarm(request));
|
||||||
|
return future.thenApply(reply -> {
|
||||||
|
MxGatewayErrors.ensureProtocolSuccess("acknowledge alarm", reply.getProtocolStatus(), null);
|
||||||
|
return reply;
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Streams a snapshot of all alarms currently Active or ActiveAcked — the
|
||||||
|
* gateway's ConditionRefresh equivalent. Used after reconnect to seed
|
||||||
|
* local Part 9 state.
|
||||||
|
*
|
||||||
|
* @param request the {@code QueryActiveAlarmsRequest}, optionally scoped by
|
||||||
|
* alarm-reference prefix
|
||||||
|
* @param observer caller-supplied observer that receives snapshots and completion
|
||||||
|
* @return a cancellable subscription handle
|
||||||
|
*/
|
||||||
|
public MxGatewayActiveAlarmsSubscription queryActiveAlarms(
|
||||||
|
QueryActiveAlarmsRequest request, StreamObserver<ActiveAlarmSnapshot> observer) {
|
||||||
|
MxGatewayActiveAlarmsSubscription subscription = new MxGatewayActiveAlarmsSubscription();
|
||||||
|
withStreamDeadline(rawAsyncStub()).queryActiveAlarms(request, subscription.wrap(observer));
|
||||||
|
return subscription;
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void close() {
|
public void close() {
|
||||||
if (ownedChannel != null) {
|
if (ownedChannel != null) {
|
||||||
|
|||||||
+1
-1
@@ -7,7 +7,7 @@ package com.dohertylan.mxgateway.client;
|
|||||||
* worker speak the same protocol version as the client.
|
* worker speak the same protocol version as the client.
|
||||||
*/
|
*/
|
||||||
public final class MxGatewayClientVersion {
|
public final class MxGatewayClientVersion {
|
||||||
private static final int GATEWAY_PROTOCOL_VERSION = 2;
|
private static final int GATEWAY_PROTOCOL_VERSION = 3;
|
||||||
private static final int WORKER_PROTOCOL_VERSION = 1;
|
private static final int WORKER_PROTOCOL_VERSION = 1;
|
||||||
private static final String CLIENT_VERSION = "0.1.0";
|
private static final String CLIENT_VERSION = "0.1.0";
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user