send error msg to gateway [4]
This commit is contained in:
parent
9147c6cf87
commit
49f54f4b39
1 changed files with 12 additions and 3 deletions
|
|
@ -27,6 +27,7 @@ import ru.clearing.classes.statics.data.company.relation.Relation;
|
|||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistryList;
|
||||
import ru.spcex.clearing.account.errors.AccountError;
|
||||
import ru.spcex.clearing.account.service.v2.GatewayRequestCreator;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
|
|
@ -91,6 +92,7 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
|||
private final IMessageResolver messageResolver;
|
||||
private final Producer<String, Object> kafkaProducer;
|
||||
private final KafkaSender kafkaSender;
|
||||
private final GatewayRequestCreator gatewayRequestCreator;
|
||||
|
||||
public TradingClearingRegistryService(Consumer<String, Object> kafkaQueue,
|
||||
Producer<String, Object> kafkaProducer,
|
||||
|
|
@ -105,7 +107,7 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
|||
@Qualifier("tradingClearingRegistryUpdateRequest")
|
||||
Function<TradingClearingRegistryUpdateRequest, IValidator> tradingClearingRegistryUpdateRequestValidator,
|
||||
@Qualifier("tradingClearingRegistryBlockRequest")
|
||||
Function<CommonIdRequest, IValidator> tradingClearingRegistryBlockRequestValidator) {
|
||||
Function<CommonIdRequest, IValidator> tradingClearingRegistryBlockRequestValidator, GatewayRequestCreator gatewayRequestCreator) {
|
||||
super(kafkaQueue, kafkaProducer);
|
||||
this.kafkaProducer = kafkaProducer;
|
||||
this.kafkaSender = kafkaSender;
|
||||
|
|
@ -128,6 +130,7 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
|||
this.tradingClearingRegistryUpdateRequestValidator = tradingClearingRegistryUpdateRequestValidator;
|
||||
this.tradingClearingRegistryBlockRequestValidator = tradingClearingRegistryBlockRequestValidator;
|
||||
this.messageResolver = messageResolver;
|
||||
this.gatewayRequestCreator = gatewayRequestCreator;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -519,12 +522,18 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini
|
|||
Map.of("code", accountForCheck.getTkrCode())
|
||||
);
|
||||
Optional<Tkr> tkr = createRequestToGateway(tradingClearingRegistry);
|
||||
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||
if (tkr.isPresent()) {
|
||||
SendTkrRequest sendTkrRequest = new SendTkrRequest();
|
||||
sendTkrRequest.getTkrs().add(tkr.get());
|
||||
sendTkrRequest.setRequestId(req.getRequestId());
|
||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||
} else {
|
||||
Optional<String> errorMsg = Optional.of(
|
||||
messageResolver.resolve(AccountError.TradingClearingRegistryNotFound, accountForCheck.getTkrCode()));
|
||||
Tkr errorTkr = gatewayRequestCreator.crateErrorTkrToGateway(accountForCheck, errorMsg);
|
||||
sendTkrRequest.getTkrs().add(errorTkr);
|
||||
sendTkrRequest.setRequestId(req.getRequestId());
|
||||
}
|
||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue