@ -18,19 +18,16 @@ package org.thingsboard.server.queue.provider;
import lombok.extern.slf4j.Slf4j ;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression ;
import org.springframework.stereotype.Component ;
import org.thingsboard.server.common.msg.queue.ServiceType ;
import org.thingsboard.server.gen.js.JsInvokeProtos ;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.queue.TbQueueConsumer ;
import org.thingsboard.server.queue.TbQueueProducer ;
import org.thingsboard.server.queue.TbQueueRequestTemplate ;
import org.thingsboard.server.queue.common.TbProtoJsQueueMsg ;
import org.thingsboard.server.queue.common.TbProtoQueueMsg ;
import org.thingsboard.server.queue.discovery.PartitionService ;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider ;
import org.thingsboard.server.queue.memory.InMemoryTbQueueConsumer ;
import org.thingsboard.server.queue.memory.InMemoryTbQueueProducer ;
import org.thingsboard.server.queue.settings.TbQueueCoreSettings ;
@ -44,74 +41,79 @@ import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration;
@ConditionalOnExpression ( "'${queue.type:null}'=='in-memory' && '${service.type:null}'=='monolith'" )
public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory , TbRuleEngineQueueFactory {
private final PartitionService partitionService ;
private final TbQueueCoreSettings coreSettings ;
private final TbServiceInfoProvider serviceInfoProvider ;
private final TbQueueRuleEngineSettings ruleEngineSettings ;
private final TbQueueTransportApiSettings transportApiSettings ;
private final TbQueueTransportNotificationSettings notificationSettings ;
private final TbQueueTransportNotificationSettings tra nsportN otificationSettings;
public InMemoryMonolithQueueFactory ( TbQueueCoreSettings coreSettings ,
public InMemoryMonolithQueueFactory ( PartitionService partitionService , TbQueueCoreSettings coreSettings ,
TbQueueRuleEngineSettings ruleEngineSettings ,
TbServiceInfoProvider serviceInfoProvider ,
TbQueueTransportApiSettings transportApiSettings ,
TbQueueTransportNotificationSettings notificationSettings ) {
TbQueueTransportNotificationSettings transportNotificationSettings ) {
this . partitionService = partitionService ;
this . coreSettings = coreSettings ;
this . serviceInfoProvider = serviceInfoProvider ;
this . ruleEngineSettings = ruleEngineSettings ;
this . transportApiSettings = transportApiSettings ;
this . notificationSettings = notificationSettings ;
this . tra nsportN otificationSettings = tra nsportN otificationSettings;
}
@Override
public TbQueueProducer < TbProtoQueueMsg < ToTransportMsg > > createTransportNotificationsMsgProducer ( ) {
return new InMemoryTbQueueProducer < > ( notificationSettings . getNotificationsTopic ( ) ) ;
public TbQueueProducer < TbProtoQueueMsg < TransportProtos . T oTransportMsg > > createTransportNotificationsMsgProducer ( ) {
return new InMemoryTbQueueProducer < > ( tra nsportN otificationSettings. getNotificationsTopic ( ) ) ;
}
@Override
public TbQueueProducer < TbProtoQueueMsg < ToRuleEngineMsg > > createRuleEngineMsgProducer ( ) {
public TbQueueProducer < TbProtoQueueMsg < TransportProtos . T oRuleEngineMsg > > createRuleEngineMsgProducer ( ) {
return new InMemoryTbQueueProducer < > ( ruleEngineSettings . getTopic ( ) ) ;
}
@Override
public TbQueueProducer < TbProtoQueueMsg < ToCoreMsg > > createTbCore MsgProducer ( ) {
return new InMemoryTbQueueProducer < > ( co reSettings. getTopic ( ) ) ;
public TbQueueProducer < TbProtoQueueMsg < TransportProtos . ToRuleEngineNotificationMsg > > createRuleEngineNotifications MsgProducer ( ) {
return new InMemoryTbQueueProducer < > ( ruleEngin eSettings . getTopic ( ) ) ;
}
@Override
public TbQueueConsum er < TbProtoQueueMsg < ToRuleEngineMsg > > createToRuleEngineMsgConsumer ( TbRuleEngineQueueConfiguration configuration ) {
return new InMemoryTbQueueConsum er < > ( ruleEngin eSettings . getTopic ( ) ) ;
public TbQueueProduc er < TbProtoQueueMsg < TransportProtos . ToCoreMsg > > createTbCoreMsgProducer ( ) {
return new InMemoryTbQueueProduc er < > ( co reSettings. getTopic ( ) ) ;
}
@Override
public TbQueueConsum er < TbProtoQueueMsg < ToCoreMsg > > createToCoreMsgConsum er ( ) {
return new InMemoryTbQueueConsum er < > ( coreSettings . getTopic ( ) ) ;
public TbQueueProduc er < TbProtoQueueMsg < TransportProtos . T oCoreNotification Msg > > createTbCoreNotificationsMsgProduc er ( ) {
return new InMemoryTbQueueProduc er < > ( coreSettings . getTopic ( ) ) ;
}
@Override
public TbQueueConsumer < TbProtoQueueMsg < TransportApiRequestMsg > > createTransportApiRequestConsumer ( ) {
return new InMemoryTbQueueConsumer < > ( transportApi Settings. getRequests Topic ( ) ) ;
public TbQueueConsumer < TbProtoQueueMsg < TransportProtos . ToRuleEngineMsg > > createToRuleEngineMsgConsumer ( TbRuleEngineQueueConfiguration configuration ) {
return new InMemoryTbQueueConsumer < > ( ruleEngine Settings. getTopic ( ) ) ;
}
@Override
public TbQueueProduc er < TbProtoQueueMsg < TransportApiResponseMsg > > createTransportApiResponseProduc er ( ) {
return new InMemoryTbQueueProduc er < > ( transportApiSettings . getResponsesTopic ( ) ) ;
public TbQueueConsum er < TbProtoQueueMsg < TransportProtos . ToRuleEngineNotificationMsg > > createToRuleEngineNotificationsMsgConsum er ( ) {
return new InMemoryTbQueueConsum er < > ( partitionService . getNotificationsTopic ( ServiceType . TB_RULE_ENGINE , serviceInfoProvider . getServiceId ( ) ) . getFullTopicName ( ) ) ;
}
@Override
public TbQueueProduc er < TbProtoQueueMsg < ToRuleEngineNotificationMsg > > createRuleEngineNotificationsMsgProduc er ( ) {
return new InMemoryTbQueueProduc er < > ( ruleEngin eSettings . getTopic ( ) + ".notifications" ) ;
public TbQueueConsum er < TbProtoQueueMsg < TransportProtos . ToCoreMsg > > createToCoreMsgConsum er ( ) {
return new InMemoryTbQueueConsum er < > ( co reSettings. getTopic ( ) ) ;
}
@Override
public TbQueueProduc er < TbProtoQueueMsg < ToCoreNotificationMsg > > createTbCoreNotificationsMsgProduc er ( ) {
return new InMemoryTbQueueProduc er < > ( coreSettings . getTopic ( ) + ".notifications" ) ;
public TbQueueConsum er < TbProtoQueueMsg < TransportProtos . T oCoreNotificationMsg > > createToCoreNotificationsMsgConsum er ( ) {
return new InMemoryTbQueueConsum er < > ( partitionService . getNotificationsTopic ( ServiceType . TB_CORE , serviceInfoProvider . getServiceId ( ) ) . getFullTopicName ( ) ) ;
}
@Override
public TbQueueConsumer < TbProtoQueueMsg < ToCoreNotificationMsg > > createToCoreNotificationsMsg Consumer ( ) {
return new InMemoryTbQueueConsumer < > ( core Settings. getTopic ( ) + ".notifications" ) ;
public TbQueueConsumer < TbProtoQueueMsg < TransportProtos . TransportApiRequestMsg > > createTransportApiRequest Consumer ( ) {
return new InMemoryTbQueueConsumer < > ( transportApi Settings. getRequests Topic ( ) ) ;
}
@Override
public TbQueueConsum er < TbProtoQueueMsg < ToRuleEngineNotificationMsg > > createToRuleEngineNotificationsMsgConsum er ( ) {
return new InMemoryTbQueueConsum er < > ( ruleEngineSettings . getTopic ( ) + ".notifications" ) ;
public TbQueueProduc er < TbProtoQueueMsg < TransportProtos . TransportApiResponseMsg > > createTransportApiResponseProduc er ( ) {
return new InMemoryTbQueueProduc er < > ( transportApiSettings . getResponsesTopic ( ) ) ;
}
@Override