Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,10 @@
When updating the version here, ensure you match the correct aws-crt version below.
Get the correct version from: https://github.com/aws/aws-iot-device-sdk-java-v2/blob/main/sdk/pom.xml#L45
!-->
<version>1.33.0</version>
<!-- Snapshot uber jar from the GG Maven repo (greengrass-common CloudFront repository above)
built from tushar-aws/aws-iot-device-sdk-java-v2 @ fe364de (SubscriptionMode support).
Revert to the public release version once the SDK change ships to Maven Central. -->
<version>1.33.0-DIRECTMSG-SNAPSHOT</version>
<exclusions>
<exclusion>
<groupId>software.amazon.awssdk.crt</groupId>
Expand All @@ -249,7 +252,7 @@
<dependency>
<groupId>software.amazon.awssdk.crt</groupId>
<artifactId>aws-crt</artifactId>
<version>0.45.0</version>
<version>0.47.0</version>
<classifier>fips-where-available</classifier>
</dependency>
<dependency>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,14 @@
import org.mockito.ArgumentCaptor;
import software.amazon.awssdk.aws.greengrass.GreengrassCoreIPCClient;
import software.amazon.awssdk.aws.greengrass.SubscribeToIoTCoreResponseHandler;
import software.amazon.awssdk.aws.greengrass.model.InvalidArgumentsError;
import software.amazon.awssdk.aws.greengrass.model.IoTCoreMessage;
import software.amazon.awssdk.aws.greengrass.model.PublishToIoTCoreRequest;
import software.amazon.awssdk.aws.greengrass.model.QOS;
import software.amazon.awssdk.aws.greengrass.model.ServiceError;
import software.amazon.awssdk.aws.greengrass.model.SubscribeToIoTCoreRequest;
import software.amazon.awssdk.aws.greengrass.model.SubscriptionMode;
import software.amazon.awssdk.aws.greengrass.model.UnauthorizedError;
import software.amazon.awssdk.crt.io.SocketOptions;
import software.amazon.awssdk.eventstreamrpc.EventStreamRPCConnection;
import software.amazon.awssdk.eventstreamrpc.StreamResponseHandler;
Expand All @@ -57,11 +60,13 @@
import static com.aws.greengrass.testcommons.testutilities.ExceptionLogProtector.ignoreExceptionOfType;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.containsString;
import static org.hamcrest.Matchers.instanceOf;
import static org.hamcrest.Matchers.is;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
Expand Down Expand Up @@ -217,6 +222,199 @@
}
}

@Test
void GIVEN_receive_only_WHEN_subscribe_THEN_receive_only_registered_and_message_received() throws Exception {
lenient().when(mqttClient.unsubscribe(any(Unsubscribe.class)))
.thenReturn(CompletableFuture.completedFuture(null));
String subscribeTopic = "A/A/C/D"; // matches the mqttproxy.yaml grant "A/+/C/D*"
CountDownLatch messageLatch = new CountDownLatch(1);
GreengrassCoreIPCClient greengrassCoreIPCClient = new GreengrassCoreIPCClient(clientConnection);
SubscribeToIoTCoreRequest subscribeToIoTCoreRequest = new SubscribeToIoTCoreRequest();
subscribeToIoTCoreRequest.setTopicName(subscribeTopic);
// RECEIVE_ONLY mode with no qos.
subscribeToIoTCoreRequest.setSubscriptionMode(SubscriptionMode.RECEIVE_ONLY);

StreamResponseHandler<IoTCoreMessage> streamResponseHandler = new StreamResponseHandler<IoTCoreMessage>() {
@Override
public void onStreamEvent(IoTCoreMessage streamEvent) {
if (Arrays.equals(streamEvent.getMessage().getPayload(), TEST_PAYLOAD)
&& streamEvent.getMessage().getTopicName().equals(subscribeTopic)) {
messageLatch.countDown();
}
}

@Override
public boolean onStreamError(Throwable error) {
logger.atError().cause(error).log("Subscribe stream errored");
return false;
}

@Override
public void onStreamClosed() {
}
};

SubscribeToIoTCoreResponseHandler responseHandler = greengrassCoreIPCClient.subscribeToIoTCore(
subscribeToIoTCoreRequest, Optional.of(streamResponseHandler));
responseHandler.getResponse().get(TIMEOUT_FOR_MQTTPROXY_SECONDS, TimeUnit.SECONDS);

ArgumentCaptor<Subscribe> subscribeRequestArgumentCaptor = ArgumentCaptor.forClass(Subscribe.class);
verify(mqttClient).subscribe(subscribeRequestArgumentCaptor.capture());
Subscribe capturedSubscribeRequest = subscribeRequestArgumentCaptor.getValue();
assertThat(capturedSubscribeRequest.getTopic(), is(subscribeTopic));
// The captured Subscribe carries the receive-only flag.
assertThat(capturedSubscribeRequest.isSkipCloudSubscribe(), is(true));

// A message delivered to the captured callback streams back to the component over IPC.
Consumer<Publish> callback = capturedSubscribeRequest.getCallback();
callback.accept(Publish.builder().topic(subscribeTopic).payload(TEST_PAYLOAD).build());
assertTrue(messageLatch.await(TIMEOUT_FOR_MQTTPROXY_SECONDS, TimeUnit.SECONDS));

// Stream close unsubscribes the registration.
responseHandler.closeStream();
Thread.sleep(500);
ArgumentCaptor<Unsubscribe> unsubscribeRequestArgumentCaptor = ArgumentCaptor.forClass(Unsubscribe.class);
verify(mqttClient).unsubscribe(unsubscribeRequestArgumentCaptor.capture());
Unsubscribe capturedUnsubscribeRequest = unsubscribeRequestArgumentCaptor.getValue();
assertThat(capturedUnsubscribeRequest.getTopic(), is(subscribeTopic));
assertThat(capturedUnsubscribeRequest.getSubscriptionCallback(), is(callback));
}

@Test
void GIVEN_receive_only_with_qos_WHEN_subscribe_THEN_qos_ignored_and_receive_only_registered() throws Exception {
lenient().when(mqttClient.unsubscribe(any(Unsubscribe.class)))
.thenReturn(CompletableFuture.completedFuture(null));
String subscribeTopic = "A/A/C/D";
CountDownLatch messageLatch = new CountDownLatch(1);
GreengrassCoreIPCClient greengrassCoreIPCClient = new GreengrassCoreIPCClient(clientConnection);
SubscribeToIoTCoreRequest subscribeToIoTCoreRequest = new SubscribeToIoTCoreRequest();
subscribeToIoTCoreRequest.setTopicName(subscribeTopic);
// qos is set alongside RECEIVE_ONLY mode; it is ignored.
subscribeToIoTCoreRequest.setQos(QOS.AT_LEAST_ONCE);
subscribeToIoTCoreRequest.setSubscriptionMode(SubscriptionMode.RECEIVE_ONLY);

StreamResponseHandler<IoTCoreMessage> streamResponseHandler = new StreamResponseHandler<IoTCoreMessage>() {
@Override
public void onStreamEvent(IoTCoreMessage streamEvent) {
if (Arrays.equals(streamEvent.getMessage().getPayload(), TEST_PAYLOAD)
&& streamEvent.getMessage().getTopicName().equals(subscribeTopic)) {
messageLatch.countDown();
}
}

@Override
public boolean onStreamError(Throwable error) {
logger.atError().cause(error).log("Subscribe stream errored");
return false;
}

@Override
public void onStreamClosed() {
}
};

SubscribeToIoTCoreResponseHandler responseHandler = greengrassCoreIPCClient.subscribeToIoTCore(
subscribeToIoTCoreRequest, Optional.of(streamResponseHandler));
responseHandler.getResponse().get(TIMEOUT_FOR_MQTTPROXY_SECONDS, TimeUnit.SECONDS);

ArgumentCaptor<Subscribe> subscribeRequestArgumentCaptor = ArgumentCaptor.forClass(Subscribe.class);
verify(mqttClient).subscribe(subscribeRequestArgumentCaptor.capture());
Subscribe capturedSubscribeRequest = subscribeRequestArgumentCaptor.getValue();
assertThat(capturedSubscribeRequest.isSkipCloudSubscribe(), is(true));

Consumer<Publish> callback = capturedSubscribeRequest.getCallback();
callback.accept(Publish.builder().topic(subscribeTopic).payload(TEST_PAYLOAD).build());
assertTrue(messageLatch.await(TIMEOUT_FOR_MQTTPROXY_SECONDS, TimeUnit.SECONDS));

responseHandler.closeStream();
}

@Test
void GIVEN_receive_only_on_unauthorized_topic_WHEN_subscribe_THEN_client_gets_unauthorized(
ExtensionContext context) throws Exception {
ignoreExceptionOfType(context, ExecutionException.class);
ignoreExceptionOfType(context, UnauthorizedError.class);

// Topic outside the mqttproxy.yaml grant (only "A/+/C/D*" and "X/Y*Z/#" are allowed).
String unauthorizedTopic = "Z/unauthorized";
GreengrassCoreIPCClient greengrassCoreIPCClient = new GreengrassCoreIPCClient(clientConnection);
SubscribeToIoTCoreRequest subscribeToIoTCoreRequest = new SubscribeToIoTCoreRequest();
subscribeToIoTCoreRequest.setTopicName(unauthorizedTopic);
subscribeToIoTCoreRequest.setSubscriptionMode(SubscriptionMode.RECEIVE_ONLY);

// Subscribe is a streaming operation, so the client requires a stream handler.
StreamResponseHandler<IoTCoreMessage> streamResponseHandler = new StreamResponseHandler<IoTCoreMessage>() {
@Override
public void onStreamEvent(IoTCoreMessage streamEvent) {
}

@Override
public boolean onStreamError(Throwable error) {
return false;
}

@Override
public void onStreamClosed() {
}
};
SubscribeToIoTCoreResponseHandler responseHandler = greengrassCoreIPCClient.subscribeToIoTCore(
subscribeToIoTCoreRequest, Optional.of(streamResponseHandler));

boolean gotUnauthorized = false;
try {
responseHandler.getResponse().get(TIMEOUT_FOR_MQTTPROXY_SECONDS, TimeUnit.SECONDS);
} catch (ExecutionException e) {
gotUnauthorized = e.getCause() instanceof UnauthorizedError;
}
assertTrue(gotUnauthorized, "RECEIVE_ONLY on an unauthorized topic must surface UnauthorizedError");
// No subscription was attempted.
verify(mqttClient, never()).subscribe(any(Subscribe.class));
}

@Test
void GIVEN_unrecognized_subscription_mode_WHEN_subscribe_THEN_client_gets_invalid_arguments(
ExtensionContext context) throws Exception {
ignoreExceptionOfType(context, ExecutionException.class);
ignoreExceptionOfType(context, InvalidArgumentsError.class);

String subscribeTopic = "A/A/C/D"; // matches the mqttproxy.yaml grant "A/+/C/D*"
GreengrassCoreIPCClient greengrassCoreIPCClient = new GreengrassCoreIPCClient(clientConnection);
SubscribeToIoTCoreRequest subscribeToIoTCoreRequest = new SubscribeToIoTCoreRequest();
subscribeToIoTCoreRequest.setTopicName(subscribeTopic);
subscribeToIoTCoreRequest.setQos(QOS.AT_LEAST_ONCE);
subscribeToIoTCoreRequest.setSubscriptionMode("FUTURE_MODE");

// Subscribe is a streaming operation, so the client requires a stream handler.
StreamResponseHandler<IoTCoreMessage> streamResponseHandler = new StreamResponseHandler<IoTCoreMessage>() {
@Override
public void onStreamEvent(IoTCoreMessage streamEvent) {
}

@Override
public boolean onStreamError(Throwable error) {
return false;
}

@Override
public void onStreamClosed() {
}
};
SubscribeToIoTCoreResponseHandler responseHandler = greengrassCoreIPCClient.subscribeToIoTCore(
subscribeToIoTCoreRequest, Optional.of(streamResponseHandler));

Throwable cause = null;
try {
responseHandler.getResponse().get(TIMEOUT_FOR_MQTTPROXY_SECONDS, TimeUnit.SECONDS);
} catch (ExecutionException e) {
cause = e.getCause();
}
assertThat("an unrecognized subscriptionMode must surface InvalidArgumentsError, not UnmappedDataException",
cause, instanceOf(InvalidArgumentsError.class));
assertThat(cause.getMessage(), containsString("FUTURE_MODE"));
// No cloud subscription was created for a mode we could not interpret.
verify(mqttClient, never()).subscribe(any(Subscribe.class));
}

@Test
void GIVEN_MqttProxyEventStreamClient_WHEN_publish_throws_error_THEN_client_gets_error(ExtensionContext context)
throws InterruptedException, MqttRequestException, SpoolerStoreException {
Expand Down
Loading
Loading