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
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
import javasabr.mqtt.service.message.handler.MqttInMessageHandler;
import javasabr.mqtt.service.message.handler.impl.ConnectInMqttInMessageHandler;
import javasabr.mqtt.service.message.handler.impl.DisconnectMqttInMessageHandler;
import javasabr.mqtt.service.message.handler.impl.PingRequestMqttInMessageHandler;
import javasabr.mqtt.service.message.handler.impl.PublishAckMqttInMessageHandler;
import javasabr.mqtt.service.message.handler.impl.PublishCompleteMqttInMessageHandler;
import javasabr.mqtt.service.message.handler.impl.PublishMqttInMessageHandler;
Expand Down Expand Up @@ -246,6 +247,11 @@ MqttInMessageHandler unsubscribeMqttInMessageHandler(
topicService);
}

@Bean
PingRequestMqttInMessageHandler pingRequestMqttInMessageHandler(MessageOutFactoryService messageOutFactoryService) {
return new PingRequestMqttInMessageHandler(messageOutFactoryService);
}

@Bean
ConnectionService externalMqttConnectionService(Collection<? extends MqttInMessageHandler> inMessageHandlers) {
return new DefaultConnectionService(ExternalNetworkMqttUser.class, inMessageHandlers);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package javasabr.mqtt.service.message.handler.impl;

import javasabr.mqtt.model.message.MqttMessageType;
import javasabr.mqtt.network.MqttConnection;
import javasabr.mqtt.network.impl.ExternalNetworkMqttUser;
import javasabr.mqtt.network.message.in.PingRequestMqttInMessage;
import javasabr.mqtt.network.message.out.MqttOutMessage;
import javasabr.mqtt.network.session.NetworkMqttSession;
import javasabr.mqtt.service.MessageOutFactoryService;

public class PingRequestMqttInMessageHandler extends
AbstractMqttInMessageHandler<ExternalNetworkMqttUser, PingRequestMqttInMessage> {

public PingRequestMqttInMessageHandler(MessageOutFactoryService messageOutFactoryService) {
super(ExternalNetworkMqttUser.class, PingRequestMqttInMessage.class, messageOutFactoryService);
}

@Override
public MqttMessageType messageType() {
return MqttMessageType.PING_REQUEST;
}

@Override
protected void processValidMessage(
MqttConnection connection,
ExternalNetworkMqttUser user,
NetworkMqttSession session,
PingRequestMqttInMessage message) {
MqttOutMessage pingResponse = messageOutFactoryService
.resolveFactory(user)
.newPingResponse();
user.sendInBackground(pingResponse);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
package javasabr.mqtt.service.message.handler.impl

import javasabr.mqtt.model.exception.MalformedProtocolMqttException
import javasabr.mqtt.model.message.MqttMessageType
import javasabr.mqtt.model.reason.code.PublishReceivedReasonCode
import javasabr.mqtt.network.MqttConnection
import javasabr.mqtt.network.NetworkUnitSpecification
import javasabr.mqtt.network.impl.ExternalNetworkMqttUser
import javasabr.mqtt.network.message.in.PingRequestMqttInMessage
import javasabr.mqtt.network.message.in.PingResponseMqttInMessage
import javasabr.mqtt.network.message.out.PingResponseMqtt311OutMessage
import javasabr.mqtt.service.MessageOutFactoryService
import javasabr.mqtt.service.message.out.factory.MqttMessageOutFactory
import javasabr.mqtt.service.session.impl.InMemoryNetworkMqttSession
import javasabr.rlib.common.util.BufferUtils

class PingRequestMqttInMessageHandlerTest extends NetworkUnitSpecification {

def "should send PINGRESP in response to PINGREQ"() {
given:
def messageOutFactoryService = Mock(MessageOutFactoryService)
def pintRequestHandler = new PingRequestMqttInMessageHandler(messageOutFactoryService)
and:
def session = Mock(InMemoryNetworkMqttSession)
def user = Mock(ExternalNetworkMqttUser)
def connection = Mock(MqttConnection)
def pingResponse = Mock(PingResponseMqtt311OutMessage)
def messageOutFactory = Mock(MqttMessageOutFactory)
def pingRequest = new PingRequestMqttInMessage(PingResponseMqttInMessage.MESSAGE_FLAGS)

when:
pintRequestHandler.processValidMessage(connection, pingRequest)

then:
1 * messageOutFactory.newPingResponse() >> pingResponse
1 * messageOutFactoryService.resolveFactory(user) >> messageOutFactory
1 * connection.user() >> user
1 * user.session() >> session
1 * user.sendInBackground(_) >> { args ->
assert args[0] instanceof PingResponseMqtt311OutMessage
}
}

def "should not allow invalid message flags"() {
given:
def dataBuffer = BufferUtils.prepareBuffer(512) {
it.putShort(testMessageId)
it.put(PublishReceivedReasonCode.SUCCESS)
it.putMbi(0)
}
when:
def inMessage = new PingRequestMqttInMessage(0b0101_0101 as byte)
def result = inMessage.read(defaultMqtt5Connection, dataBuffer, dataBuffer.limit())
then:
!result
with(inMessage) {
exception() instanceof MalformedProtocolMqttException
exception().message == "Unexpected message flags:[0b0101_0101] in message:[$MqttMessageType.PING_REQUEST]"
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
*/
public class PingRequestMqttInMessage extends MqttInMessage {

public static final byte MESSAGE_FLAGS = 0b0000_0000;
public static final byte MESSAGE_TYPE = (byte) MqttMessageType.PING_REQUEST.ordinal();

public PingRequestMqttInMessage(byte messageFlags) {
Expand All @@ -22,4 +23,9 @@ public byte messageTypeId() {
public MqttMessageType messageType() {
return MqttMessageType.PING_REQUEST;
}

@Override
protected boolean validMessageFlags(byte messageFlags) {
return messageFlags == MESSAGE_FLAGS;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
*/
public class PingResponseMqttInMessage extends MqttInMessage {

public static final byte MESSAGE_FLAGS = 0b0000_0000;
public static final byte MESSAGE_TYPE = (byte) MqttMessageType.PING_RESPONSE.ordinal();

public PingResponseMqttInMessage(byte messageFlags) {
Expand All @@ -22,4 +23,9 @@ public byte messageTypeId() {
public MqttMessageType messageType() {
return MqttMessageType.PING_RESPONSE;
}

@Override
protected boolean validMessageFlags(byte messageFlags) {
return messageFlags == MESSAGE_FLAGS;
}
}
Loading