diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionController.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionController.java index 354d1f06c..d70838af6 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionController.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionController.java @@ -1,22 +1,30 @@ package ru.spcex.clearing.backendapi.controller.queue.payment; import io.swagger.annotations.ApiOperation; +import io.swagger.annotations.ApiParam; import io.swagger.annotations.ApiResponse; import io.swagger.annotations.ApiResponses; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.http.MediaType; import org.springframework.stereotype.Controller; +import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; import org.springframework.web.bind.annotation.ResponseBody; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.spcex.clearing.backendapi.controller.queue.AbstractQueueController; +import ru.spcex.clearing.backendapi.controller.request.cud.payment.PIClearingOutbondActionNew; +import ru.spcex.clearing.backendapi.controller.response.BasicSpcexResponse; +import ru.spcex.clearing.backendapi.controller.response.cud.CudResponse; import ru.spcex.clearing.backendapi.controller.response.entity.CommonGetAllResponse; import ru.spcex.clearing.backendapi.service.IOperator; import ru.spcex.clearing.backendapi.service.IStateLoader; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; import java.util.Collection; import java.util.Map; +import java.util.concurrent.ExecutionException; @Controller @RequestMapping("/payment-instructions") @@ -39,4 +47,14 @@ public class PaymentInstructionController extends AbstractQueueController { response.fromEntity(all); return response; } + + @ApiOperation(value = "Досрочное изъятие части денежных средств по договору") + @ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = CudResponse.class), @ApiResponse(code = 400, message = "Ошибка валидации", response = BasicSpcexResponse.class)}) + @RequestMapping(value = "/clearing-outbound", method = RequestMethod.POST, consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_JSON_VALUE) + @ResponseBody + public CudResponse clearingOutboundAction( + @ApiParam(value = "Параметры команды в JSON формате.", required = true) + @RequestBody PIClearingOutbondActionNew piClearingOutbondAction) throws ExecutionException, InterruptedException { + return processRequest(Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION, piClearingOutbondAction); + } } diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/register/RegistryController.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/register/RegistryController.java new file mode 100644 index 000000000..b77031539 --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/register/RegistryController.java @@ -0,0 +1,60 @@ +package ru.spcex.clearing.backendapi.controller.queue.register; + +import io.swagger.annotations.ApiOperation; +import io.swagger.annotations.ApiParam; +import io.swagger.annotations.ApiResponse; +import io.swagger.annotations.ApiResponses; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.http.MediaType; +import org.springframework.stereotype.Controller; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestMethod; +import org.springframework.web.bind.annotation.ResponseBody; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.backendapi.controller.queue.AbstractQueueController; +import ru.spcex.clearing.backendapi.controller.request.cud.registry.RSplitDepositActionNew; +import ru.spcex.clearing.backendapi.controller.response.BasicSpcexResponse; +import ru.spcex.clearing.backendapi.controller.response.cud.CudResponse; +import ru.spcex.clearing.backendapi.controller.response.entity.CommonGetAllResponse; +import ru.spcex.clearing.backendapi.service.IOperator; +import ru.spcex.clearing.backendapi.service.IStateLoader; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; + +import java.util.Collection; +import java.util.Map; +import java.util.concurrent.ExecutionException; + +@Controller +@RequestMapping("/registries") +public class RegistryController extends AbstractQueueController { + private final IStateLoader stateLoader; + + @Autowired + public RegistryController(IOperator operator, IStateLoader stateLoader) { + super(operator); + this.stateLoader = stateLoader; + } + + @ApiOperation(value = "get all registries.") + @ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = CommonGetAllResponse.class)}) + @RequestMapping(method = RequestMethod.GET) + @ResponseBody + public CommonGetAllResponse getAll() { + Collection> all = stateLoader.getAllMetaTransform(IMDGDistributedNames.Map_Registry, Registry.class); + CommonGetAllResponse response = new CommonGetAllResponse(); + response.fromEntity(all); + return response; + } + + @ApiOperation(value = "Досрочное изъятие части денежных средств по договору") + @ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = CudResponse.class), @ApiResponse(code = 400, message = "Ошибка валидации", response = BasicSpcexResponse.class)}) + @RequestMapping(value = "/splitDeposit", method = RequestMethod.POST, consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_JSON_VALUE) + @ResponseBody + public CudResponse splitDepositAction( + @ApiParam(value = "Параметры команды в JSON формате.", required = true) + @RequestBody RSplitDepositActionNew rSplitDepositActionNew) throws ExecutionException, InterruptedException { + return processRequest(Consts.REGISTRY_SPLIT_DEPOSIT_ACTION, rSplitDepositActionNew); + } +} diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/payment/PIClearingOutbondActionNew.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/payment/PIClearingOutbondActionNew.java new file mode 100644 index 000000000..2c8559991 --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/payment/PIClearingOutbondActionNew.java @@ -0,0 +1,114 @@ +package ru.spcex.clearing.backendapi.controller.request.cud.payment; + +import com.fasterxml.jackson.annotation.JsonProperty; +import io.swagger.annotations.ApiModelProperty; +import ru.spcex.clearing.backendapi.domain.actions.IAction; +import ru.spcex.clearing.backendapi.errors.BackEndError; +import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest; +import ru.spcex.platform.utils.enumeration.EnumMessage; + +import java.math.BigDecimal; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.List; + +public class PIClearingOutbondActionNew implements IAction { + + @ApiModelProperty(value = "Участник отправитель", example = "123") + @JsonProperty + public Long senderId; + @ApiModelProperty(value = "Участник получатель", example = "123") + @JsonProperty + public Long addresseeId; + @ApiModelProperty(value = "Назначение платежа", example = "Text ...") + @JsonProperty + public String paymentPurpose; + @ApiModelProperty(value = "Сумма отправителя", example = "123.45") + @JsonProperty + public BigDecimal creditLeg_amount; + @ApiModelProperty(value = "Наименование счета отправителя", example = "2001") + @JsonProperty + public Long creditLeg_accountId; + @ApiModelProperty(value = "Наименование счета получателя", example = "2002") + @JsonProperty + public Long debitLeg_accountId; + + + @Override + public Collection validate() { + List errors = new ArrayList<>(); + if (this.senderId == null) + errors.add(new EnumMessage(BackEndError.ValidationError, "senderId")); + if (this.addresseeId == null) + errors.add(new EnumMessage(BackEndError.ValidationError, "addresseeId")); + return errors.isEmpty() ? Collections.emptyList() : errors; + } + + @Override + public PIClearingOutbondActionNewRequest toRequest() { + var req = new PIClearingOutbondActionNewRequest(); + req.setSenderId(this.senderId); + req.setAddresseeId(this.addresseeId); + req.setPaymentPurpose(this.paymentPurpose); + req.setCreditLeg_amount(this.creditLeg_amount); + req.setCreditLeg_accountId(this.creditLeg_accountId); + req.setDebitLeg_accountId(this.debitLeg_accountId); + return req; + } + + @ApiModelProperty(hidden = true) + @Override + public ActionType getActionType() { + return ActionType.NEW; + } + + public Long getSenderId() { + return senderId; + } + + public void setSenderId(Long senderId) { + this.senderId = senderId; + } + + public Long getAddresseeId() { + return addresseeId; + } + + public void setAddresseeId(Long addresseeId) { + this.addresseeId = addresseeId; + } + + public String getPaymentPurpose() { + return paymentPurpose; + } + + public void setPaymentPurpose(String paymentPurpose) { + this.paymentPurpose = paymentPurpose; + } + + public BigDecimal getCreditLeg_amount() { + return creditLeg_amount; + } + + public void setCreditLeg_amount(BigDecimal creditLeg_amount) { + this.creditLeg_amount = creditLeg_amount; + } + + public Long getCreditLeg_accountId() { + return creditLeg_accountId; + } + + public void setCreditLeg_accountId(Long creditLeg_accountId) { + this.creditLeg_accountId = creditLeg_accountId; + } + + public Long getDebitLeg_accountId() { + return debitLeg_accountId; + } + + public void setDebitLeg_accountId(Long debitLeg_accountId) { + this.debitLeg_accountId = debitLeg_accountId; + } +} diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/registry/RSplitDepositActionNew.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/registry/RSplitDepositActionNew.java new file mode 100644 index 000000000..1abe96c9f --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/request/cud/registry/RSplitDepositActionNew.java @@ -0,0 +1,99 @@ +package ru.spcex.clearing.backendapi.controller.request.cud.registry; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import io.swagger.annotations.ApiModelProperty; +import org.apache.commons.lang3.StringUtils; +import ru.spcex.clearing.backendapi.domain.actions.IAction; +import ru.spcex.clearing.backendapi.errors.BackEndError; +import ru.spcex.clearing.platform.messaging.domain.ActionType; +import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistrySplitDepositActionRequest; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; +import ru.spcex.platform.utils.enumeration.EnumMessage; + +import java.math.BigDecimal; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.List; + +public class RSplitDepositActionNew implements IAction { + + @ApiModelProperty(value = "Основание", example = "Основание ...") + @JsonProperty + public String comment; + @ApiModelProperty(value = "Сумма для вывода", example = "1600.20") + @JsonProperty + public BigDecimal outboundAmount; + @ApiModelProperty(value = "Дата возврата", example = "2022-12-25") + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty + public LocalDate refundDate; + @ApiModelProperty(value = "Договор для разделения", example = "ABCD") + @JsonProperty + public String contract; + + @Override + public Collection validate() { + List errors = new ArrayList<>(); + if (StringUtils.isEmpty(this.contract)) + errors.add(new EnumMessage(BackEndError.ValidationError, "contract")); + if (this.refundDate == null) + errors.add(new EnumMessage(BackEndError.ValidationError, "refundDate")); + if (this.outboundAmount == null) + errors.add(new EnumMessage(BackEndError.ValidationError, "outboundAmount")); + return errors.isEmpty() ? Collections.emptyList() : errors; + } + + @Override + public RegistrySplitDepositActionRequest toRequest() { + var req = new RegistrySplitDepositActionRequest(); + req.setComment(this.comment); + req.setOutboundAmount(this.outboundAmount); + req.setRefundDate(this.refundDate); + req.setContract(this.contract); + return req; + } + + @ApiModelProperty(hidden = true) + @Override + public ActionType getActionType() { + return ActionType.NEW; + } + + public String getComment() { + return comment; + } + + public void setComment(String comment) { + this.comment = comment; + } + + public BigDecimal getOutboundAmount() { + return outboundAmount; + } + + public void setOutboundAmount(BigDecimal outboundAmount) { + this.outboundAmount = outboundAmount; + } + + public LocalDate getRefundDate() { + return refundDate; + } + + public void setRefundDate(LocalDate refundDate) { + this.refundDate = refundDate; + } + + public String getContract() { + return contract; + } + + public void setContract(String contract) { + this.contract = contract; + } +} diff --git a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java index b4292d047..7a4783b30 100644 --- a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java +++ b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/AbstractControllerTest.java @@ -120,6 +120,7 @@ import static ru.spcex.clearing.test.json.MatcherFactoryWithJson.usingIgnoringFi ContractRegisterController.class, ReportRegisterController.class, //registry + RegistryController.class, TradingClearingRegistryController.class, //scheduler ClearingCalendarController.class, diff --git a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java index e86246a80..80cdd1ad8 100644 --- a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java +++ b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/payment/PaymentInstructionControllerTest.java @@ -6,8 +6,10 @@ import org.springframework.http.MediaType; import org.springframework.test.web.servlet.request.MockMvcRequestBuilders; import ru.clearing.classes.statics.data.payment.PaymentInstruction; import ru.spcex.clearing.backendapi.controller.queue.AbstractControllerTest; +import ru.spcex.clearing.backendapi.controller.request.cud.payment.PIClearingOutbondActionNew; import ru.spcex.clearing.backendapi.controller.response.entity.CommonGetAllResponse; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; import java.math.BigDecimal; import java.time.Instant; @@ -79,4 +81,20 @@ class PaymentInstructionControllerTest extends AbstractControllerTest { .andExpect(content().contentTypeCompatibleWith(MediaType.APPLICATION_JSON)) .andExpect(content().json(writeValue(expected))); } + + @Test + void clearingOutboundAction() throws Exception { + //ARRANGE + PIClearingOutbondActionNew piClearingOutbondActionNew = new PIClearingOutbondActionNew(); + piClearingOutbondActionNew.setSenderId(111L); + piClearingOutbondActionNew.setAddresseeId(222L); + piClearingOutbondActionNew.setPaymentPurpose("Check now."); + piClearingOutbondActionNew.setCreditLeg_amount(BigDecimal.valueOf(120.99)); + piClearingOutbondActionNew.setCreditLeg_accountId(333L); + piClearingOutbondActionNew.setDebitLeg_accountId(335L); + //ACT and ASSERT + checkAddingByRestApi("/payment-instructions/clearing-outbound/", piClearingOutbondActionNew); + checkSendedMessegeFromKafka(Consts.PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION, piClearingOutbondActionNew); + } + } \ No newline at end of file diff --git a/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/register/RegistryControllerTest.java b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/register/RegistryControllerTest.java new file mode 100644 index 000000000..a6b7843df --- /dev/null +++ b/clearing-parent/backend-api/src/test/java/ru/spcex/clearing/backendapi/controller/queue/register/RegistryControllerTest.java @@ -0,0 +1,92 @@ +package ru.spcex.clearing.backendapi.controller.queue.register; + +import org.junit.jupiter.api.Test; +import ru.clearing.classes.statics.data.register.ContractRegister; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.backendapi.controller.queue.AbstractControllerTest; +import ru.spcex.clearing.backendapi.controller.request.cud.registry.RSplitDepositActionNew; +import ru.spcex.clearing.backendapi.controller.request.cud.registry.TradingClearingRegistryNewAction; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.Consts; + +import java.math.BigDecimal; +import java.time.LocalDate; + +import static org.junit.jupiter.api.Assertions.*; + +class RegistryControllerTest extends AbstractControllerTest { + private static final String REST_URL = "/registries/"; + + /** + * {@link RegistryController#getAll()}
+ * Тест проверяет получение запроса по REST API.
+ * Входной запрос /registries/
+ * Ответ CommonGetAllResponse
+ */ + @Test + void getAll() throws Exception { + //ARRANGE + Registry registry = new Registry(); + registry.setId(currentId.get()); + registry.setCompanyId(currentId.get()); + registry.setTradingCode("TCDE"); + registry.setClearingCode("CCDE"); + registry.setShortName("S name"); + registry.setFullName("Fullly names"); + registry.setAccountId(currentId.get()); + registry.setAccountType("TCR"); + registry.setAccount("ACCT-111222"); + registry.setRegistryDesignation("DST"); + registry.setRegistryInstrumentType("BOND"); + registry.setRegistryCapacity("capasity"); + registry.setRegistryUnit("unit"); + registry.setRegistryCode("code"); + registry.setTradingClearingRegistryId(currentId.get()); + registry.setTradingClearingRegistry("String"); + registry.setRegistryStatus("String"); + registry.setSecurityId(currentId.get()); + registry.setSecuritySymbol("CYMBL"); + registry.setBalance(BigDecimal.valueOf(111)); + registry.setOpenBalance(BigDecimal.valueOf(222)); + registry.setCloseBalance(BigDecimal.valueOf(333)); + registry.setCredit(BigDecimal.valueOf(444)); + registry.setDebit(BigDecimal.valueOf(123)); + registry.setSettledCredit(BigDecimal.valueOf(234)); + registry.setSettledDebit(BigDecimal.valueOf(3456)); + registry.setCheckBalance(BigDecimal.valueOf(789)); + registry.setDiffBalance(BigDecimal.valueOf(10.11)); + registry.setPlanBalance(BigDecimal.valueOf(12.1516)); + registry.setBalanceDimension("dimm"); + registry.setSettlementDate(LocalDate.now()); + registry.setSettlementCode("SETTL"); + registry.setTradingDate(LocalDate.now()); + registry.setClearingDate(LocalDate.now()); + registry.setRefundDate(LocalDate.now()); + registry.setValueDate(LocalDate.now()); + registry.setPrice(BigDecimal.valueOf(111222)); + registry.setContract("CONTRACT-1"); + registry.setCounterPartyId(currentId.get()); + registry.setComment("String test"); + registry.setParentId(currentId.get()); + registry.setGroupId(currentId.get()); + registry.setSessionId(currentId.get()); + registry.setSessionType("CLR1"); + + //ACT and ASSERT + checkGettingAllFromRestApi(IMDGDistributedNames.Map_Registry, registry, REST_URL); + } + + @Test + void splitDepositAction() throws Exception { + //ARRANGE + RSplitDepositActionNew rSplitDepositActionNew=new RSplitDepositActionNew(); + rSplitDepositActionNew.setComment("comment"); + rSplitDepositActionNew.setContract("CONTRACT-1"); + rSplitDepositActionNew.setRefundDate(LocalDate.now()); + rSplitDepositActionNew.setOutboundAmount(BigDecimal.ONE); + //ACT and ASSERT + checkAddingByRestApi("/registries/splitDeposit/", rSplitDepositActionNew); + checkSendedMessegeFromKafka(Consts.REGISTRY_SPLIT_DEPOSIT_ACTION, rSplitDepositActionNew); + + } +} \ No newline at end of file diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java index 0708fa158..79d7c78a5 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/config/KafkaConfig.java @@ -1,15 +1,25 @@ package ru.spcex.clearing.trade.importer.config; import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.producer.Producer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.trade.importer.config.settings.ImportTradeServiceSettings; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgId; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.util.function.Supplier; @Configuration public class KafkaConfig { @@ -21,9 +31,30 @@ public class KafkaConfig { return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } + @Bean + public ProducerFactory pf(ImportTradeServiceSettings settings) { + KafkaProducerSettings kafkaSettings = settings.getKafkaProducer(); + return KafkaProducerFactory.producerFactory(kafkaSettings); + } + + @Bean("kafkaTemplate") + public KafkaTemplate kafkaTemplate(ProducerFactory pf) { + return new KafkaTemplate<>(pf); + } + @Autowired @Bean - public Producer createProducer(ImportTradeServiceSettings settings) { - return KafkaProducerFactory.producer(settings.getKafkaProducer()); + public Supplier kafkaSender(KafkaTemplate kafkaTemplate, + ImdgProvider imdgProvider) { + ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator(); + return () -> KafkaSender + .setup() + .setKafkaTemplate(kafkaTemplate) + .idGenerator(imdgIdGenerator::nextId) + .imdgProvider(s -> { + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + return imdg::insert; + }) + .build(); } } diff --git a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java index 496b4b391..7b89a2c00 100644 --- a/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java +++ b/clearing-parent/trade-importer/src/main/java/ru/spcex/clearing/trade/importer/services/TradeImporterService.java @@ -1,9 +1,8 @@ package ru.spcex.clearing.trade.importer.services; -import org.apache.kafka.clients.producer.Producer; -import org.apache.kafka.clients.producer.ProducerRecord; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Value; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; @@ -12,6 +11,7 @@ import org.springframework.util.StringUtils; import ru.clearing.classes.statics.data.misc.STrades; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.enumeration.EnumMessage; @@ -26,6 +26,7 @@ import java.time.Instant; import java.time.LocalDate; import java.util.Collection; import java.util.Map; +import java.util.function.Supplier; import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; import static ru.spcex.clearing.trade.importer.error.TradeImporterError.sTradesNotValid; @@ -36,14 +37,18 @@ public class TradeImporterService { private final Logger log = LoggerFactory.getLogger(getClass()); private final Imdg sTradesImdg; private final JdbcTemplate jdbcTemplate; - private final Producer producer; + private final Supplier kafka; private final IMessageResolver messageResolver; + @Value("${trade-importer.database.schema}") + private String schema; - public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Producer producer, IMessageResolver messageResolver) { + + public TradeImporterService(ImdgProvider imdgProvider, JdbcTemplate jdbcTemplate, Supplier kafka, IMessageResolver messageResolver) { + imdgProvider.waitAvailable(); this.sTradesImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class); this.jdbcTemplate = jdbcTemplate; - this.producer = producer; + this.kafka = kafka; this.messageResolver = messageResolver; } @@ -54,7 +59,7 @@ public class TradeImporterService { public void process(boolean byCommand) { log.debug("Start process import STrades from DB"); - Collection tradesFromDB = jdbcTemplate.query("SELECT * FROM Trades", (resultSet, i) -> readSTrades(resultSet)); + Collection tradesFromDB = jdbcTemplate.query(String.format("SELECT * FROM %s.Trades", schema), (resultSet, i) -> readSTrades(resultSet)); Instant currentInstant = Instant.now(); log.debug("load from DB STrades: {}", tradesFromDB.size()); @@ -79,7 +84,7 @@ public class TradeImporterService { log.debug("Successfully import STrades from DB by scheduled. Send to kafka command, topic={}", S_TRADES_IMPORTED); } STradesImportedRequest sTradesImportedRequest = new STradesImportedRequest(); - producer.send(new ProducerRecord<>(S_TRADES_IMPORTED, sTradesImportedRequest)); + kafka.get().sendRequestToQueue(S_TRADES_IMPORTED, sTradesImportedRequest); } public STrades getSTradesFromImdg(STrades tradesDb, Imdg sTradesImdg) { diff --git a/clearing-parent/trade-importer/src/main/resources/application.properties b/clearing-parent/trade-importer/src/main/resources/application.properties index e2cfdbd94..7b7f8f668 100644 --- a/clearing-parent/trade-importer/src/main/resources/application.properties +++ b/clearing-parent/trade-importer/src/main/resources/application.properties @@ -5,7 +5,8 @@ trade-importer.cron.load-from-db-cron=0 0/5 * * * ? trade-importer.database.login=sa trade-importer.database.password=Aa123456 -trade-importer.database.url=jdbc:sqlserver://localhost:1433;database=SPVB_TS;schema=dbo +trade-importer.database.schema=SPVB_TS +trade-importer.database.url=jdbc:sqlserver://10.200.200.144:1433;database=ni; trade-importer.database.driver=com.microsoft.sqlserver.jdbc.SQLServerDriver trade-importer.hazelcast.cluster-members=127.0.0.1:5701 diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 1991be59e..cace0e3ba 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -118,6 +118,9 @@ public interface Consts { String REQUEST_INFO_UPDATE = "request-info-update"; + String PAYMENT_INSTRUCTION_CLEARING_OUTBOUND_ACTION = "clearing-outbound-action"; + + String REGISTRY_SPLIT_DEPOSIT_ACTION = "registry-split-deposit-action"; String REGISTRY_COVERED_DEAL_REGISTER_NEW = "registry-covered-deal-register-new"; String REGISTRY_ADMITTED_DEAL_REGISTER_NEW = "registry-admitted-deal-register-new"; String REGISTRY_DEAL_REGISTER_NEW = "registry-deal-register-new"; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/payment/PIClearingOutbondActionNewRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/payment/PIClearingOutbondActionNewRequest.java new file mode 100644 index 000000000..a850c598e --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/payment/PIClearingOutbondActionNewRequest.java @@ -0,0 +1,69 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.payment; + +import com.fasterxml.jackson.annotation.JsonProperty; + +import java.math.BigDecimal; + +public class PIClearingOutbondActionNewRequest { + + @JsonProperty + private Long senderId; + @JsonProperty + private Long addresseeId; + @JsonProperty + private String paymentPurpose; + @JsonProperty + private BigDecimal creditLeg_amount; + @JsonProperty + private Long creditLeg_accountId; + @JsonProperty + private Long debitLeg_accountId; + + public Long getSenderId() { + return senderId; + } + + public void setSenderId(Long senderId) { + this.senderId = senderId; + } + + public Long getAddresseeId() { + return addresseeId; + } + + public void setAddresseeId(Long addresseeId) { + this.addresseeId = addresseeId; + } + + public String getPaymentPurpose() { + return paymentPurpose; + } + + public void setPaymentPurpose(String paymentPurpose) { + this.paymentPurpose = paymentPurpose; + } + + public BigDecimal getCreditLeg_amount() { + return creditLeg_amount; + } + + public void setCreditLeg_amount(BigDecimal creditLeg_amount) { + this.creditLeg_amount = creditLeg_amount; + } + + public Long getCreditLeg_accountId() { + return creditLeg_accountId; + } + + public void setCreditLeg_accountId(Long creditLeg_accountId) { + this.creditLeg_accountId = creditLeg_accountId; + } + + public Long getDebitLeg_accountId() { + return debitLeg_accountId; + } + + public void setDebitLeg_accountId(Long debitLeg_accountId) { + this.debitLeg_accountId = debitLeg_accountId; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/RegistrySplitDepositActionRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/RegistrySplitDepositActionRequest.java new file mode 100644 index 000000000..231d82d69 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/registry/RegistrySplitDepositActionRequest.java @@ -0,0 +1,56 @@ +package ru.spcex.clearing.platform.messaging.domain.cud.registry; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import ru.spcex.clearing.platform.messaging.domain.json.deserialize.LocalDateDeserializer; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalDateSerializer; + +import java.math.BigDecimal; +import java.time.LocalDate; + +public class RegistrySplitDepositActionRequest { + + @JsonProperty + private String comment; + @JsonProperty + private BigDecimal outboundAmount; + @JsonSerialize(using = LocalDateSerializer.class) + @JsonDeserialize(using = LocalDateDeserializer.class) + @JsonProperty + private LocalDate refundDate; + @JsonProperty + private String contract; + + public String getComment() { + return comment; + } + + public void setComment(String comment) { + this.comment = comment; + } + + public BigDecimal getOutboundAmount() { + return outboundAmount; + } + + public void setOutboundAmount(BigDecimal outboundAmount) { + this.outboundAmount = outboundAmount; + } + + public LocalDate getRefundDate() { + return refundDate; + } + + public void setRefundDate(LocalDate refundDate) { + this.refundDate = refundDate; + } + + public String getContract() { + return contract; + } + + public void setContract(String contract) { + this.contract = contract; + } +}