Merge remote-tracking branch 'origin/dev' into dev
This commit is contained in:
commit
f70d36d7bd
16 changed files with 255 additions and 116 deletions
|
|
@ -4,26 +4,43 @@ import io.swagger.annotations.ApiOperation;
|
||||||
import io.swagger.annotations.ApiResponse;
|
import io.swagger.annotations.ApiResponse;
|
||||||
import io.swagger.annotations.ApiResponses;
|
import io.swagger.annotations.ApiResponses;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.http.MediaType;
|
||||||
|
import org.springframework.security.core.Authentication;
|
||||||
|
import org.springframework.security.core.context.SecurityContextHolder;
|
||||||
import org.springframework.stereotype.Controller;
|
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.RequestMapping;
|
||||||
import org.springframework.web.bind.annotation.RequestMethod;
|
import org.springframework.web.bind.annotation.RequestMethod;
|
||||||
import org.springframework.web.bind.annotation.ResponseBody;
|
import org.springframework.web.bind.annotation.ResponseBody;
|
||||||
import ru.clearing.classes.statics.data.user.User;
|
import ru.clearing.classes.statics.data.user.User;
|
||||||
|
import ru.spcex.clearing.backendapi.controller.queue.AbstractQueueController;
|
||||||
|
import ru.spcex.clearing.backendapi.controller.request.cud.utilities.UserAuthAction;
|
||||||
|
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.controller.response.entity.CommonGetAllResponse;
|
||||||
|
import ru.spcex.clearing.backendapi.security.KeycloakUtils;
|
||||||
|
import ru.spcex.clearing.backendapi.service.IOperator;
|
||||||
import ru.spcex.clearing.backendapi.service.IStateLoader;
|
import ru.spcex.clearing.backendapi.service.IStateLoader;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.ExecutionException;
|
||||||
|
|
||||||
@Controller
|
@Controller
|
||||||
@RequestMapping("/users")
|
@RequestMapping("/users")
|
||||||
public class UserController {
|
public class UserController extends AbstractQueueController {
|
||||||
private final IStateLoader stateLoader;
|
private final IStateLoader stateLoader;
|
||||||
|
private final Imdg<User> userImdg;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public UserController(IStateLoader stateLoader) {
|
public UserController(IOperator operator, IStateLoader stateLoader, ImdgProvider imdgProvider) {
|
||||||
|
super(operator);
|
||||||
this.stateLoader = stateLoader;
|
this.stateLoader = stateLoader;
|
||||||
|
this.userImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_User, User.class);
|
||||||
}
|
}
|
||||||
|
|
||||||
@ApiOperation(value = "get all users.")
|
@ApiOperation(value = "get all users.")
|
||||||
|
|
@ -36,4 +53,22 @@ public class UserController {
|
||||||
response.fromEntity(all);
|
response.fromEntity(all);
|
||||||
return response;
|
return response;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@ApiOperation(value = "update user")
|
||||||
|
@ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = CudResponse.class), @ApiResponse(code = 400, message = "Ошибка валидации", response = BasicSpcexResponse.class)})
|
||||||
|
@RequestMapping(method = RequestMethod.PUT, consumes = MediaType.APPLICATION_JSON_VALUE)
|
||||||
|
@ResponseBody
|
||||||
|
public CudResponse update(
|
||||||
|
@RequestBody UserAuthAction userAuthAction) throws ExecutionException, InterruptedException {
|
||||||
|
Authentication authentication = SecurityContextHolder.getContext().getAuthentication();
|
||||||
|
String username = KeycloakUtils.getUserNameFromAuthentication(authentication);
|
||||||
|
User user = userImdg.getSingleObjectByFieldValues(Map.of("identifier", username));
|
||||||
|
if (user == null) {
|
||||||
|
return processRequest(Consts.USER_AUTH_SUCCESS, userAuthAction);
|
||||||
|
}
|
||||||
|
throw new IllegalStateException("cannot create user cause it is exists: " + username);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -5,20 +5,16 @@ server.ssl.key-store=classpath:keystore/client.p12
|
||||||
server.ssl.key-store-password=Aa123456
|
server.ssl.key-store-password=Aa123456
|
||||||
server.ssl.enabled=true
|
server.ssl.enabled=true
|
||||||
spring.main.web-application-type=servlet
|
spring.main.web-application-type=servlet
|
||||||
|
|
||||||
backend-api.example-setting=test
|
backend-api.example-setting=test
|
||||||
|
|
||||||
backend-api.hazelcast.cluster-members=127.0.0.1:5701
|
backend-api.hazelcast.cluster-members=127.0.0.1:5701
|
||||||
backend-api.hazelcast.login=dev
|
backend-api.hazelcast.login=dev
|
||||||
backend-api.hazelcast.password=dev-pass
|
backend-api.hazelcast.password=dev-pass
|
||||||
|
|
||||||
backend-api.kafka-producer.bootstrap-servers=localhost:9092
|
backend-api.kafka-producer.bootstrap-servers=localhost:9092
|
||||||
backend-api.kafka-producer.acks=all
|
backend-api.kafka-producer.acks=all
|
||||||
backend-api.kafka-producer.retries=0
|
backend-api.kafka-producer.retries=0
|
||||||
backend-api.kafka-producer.batch-size=16384
|
backend-api.kafka-producer.batch-size=16384
|
||||||
backend-api.kafka-producer.linger-ms=1
|
backend-api.kafka-producer.linger-ms=1
|
||||||
backend-api.kafka-producer.buffer-memory=33554432
|
backend-api.kafka-producer.buffer-memory=33554432
|
||||||
|
|
||||||
backend-api.kafka-consumer.bootstrap-servers=localhost:9092
|
backend-api.kafka-consumer.bootstrap-servers=localhost:9092
|
||||||
backend-api.kafka-consumer.group-id=dev-group-backend-api
|
backend-api.kafka-consumer.group-id=dev-group-backend-api
|
||||||
backend-api.kafka-consumer.enable-auto-commit=true
|
backend-api.kafka-consumer.enable-auto-commit=true
|
||||||
|
|
@ -26,10 +22,7 @@ backend-api.kafka-consumer.session-timeout-ms=30000
|
||||||
backend-api.kafka-consumer.auto-offset-reset=latest
|
backend-api.kafka-consumer.auto-offset-reset=latest
|
||||||
backend-api.kafka-consumer.linger-ms=1
|
backend-api.kafka-consumer.linger-ms=1
|
||||||
backend-api.kafka-consumer.buffer-memory=33554432
|
backend-api.kafka-consumer.buffer-memory=33554432
|
||||||
|
|
||||||
backend-api.security.authorization-disabled=false
|
backend-api.security.authorization-disabled=false
|
||||||
|
|
||||||
|
|
||||||
##keycloak
|
##keycloak
|
||||||
##keycloak.auth-server-url=http://10.200.200.147:8080/
|
##keycloak.auth-server-url=http://10.200.200.147:8080/
|
||||||
##keycloak.realm=master
|
##keycloak.realm=master
|
||||||
|
|
|
||||||
|
|
@ -112,7 +112,7 @@ public class Sdf09Executor extends AbstractExecutor<SDf09> {
|
||||||
statement.setSenderId(Sender.Prc.getId());
|
statement.setSenderId(Sender.Prc.getId());
|
||||||
statement.setCreated(Instant.now());
|
statement.setCreated(Instant.now());
|
||||||
statement.setClearingDate(LocalDate.now());
|
statement.setClearingDate(LocalDate.now());
|
||||||
statement.setStatementType(StatementType.full.getKey());
|
statement.setStatementType(StatementType.incr.getKey());
|
||||||
statement.setAccountId(account.getId());
|
statement.setAccountId(account.getId());
|
||||||
statement.setAccount(sdf09.getAccount());
|
statement.setAccount(sdf09.getAccount());
|
||||||
statement.setInOutDirection(InOutDirection.in.getKey());
|
statement.setInOutDirection(InOutDirection.in.getKey());
|
||||||
|
|
|
||||||
|
|
@ -6,118 +6,33 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.scheduling.annotation.EnableScheduling;
|
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||||
import org.springframework.scheduling.annotation.Scheduled;
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf03;
|
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf11;
|
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
|
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
|
||||||
import ru.spcex.platform.enumeration.TransactionStatus;
|
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgId;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
|
||||||
|
|
||||||
import java.util.AbstractMap;
|
|
||||||
import java.util.Comparator;
|
|
||||||
import java.util.List;
|
|
||||||
import java.util.Map;
|
|
||||||
import java.util.concurrent.ExecutorService;
|
import java.util.concurrent.ExecutorService;
|
||||||
import java.util.concurrent.Executors;
|
import java.util.concurrent.Executors;
|
||||||
import java.util.stream.Collectors;
|
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
@EnableScheduling
|
@EnableScheduling
|
||||||
public class ClearingService {
|
public class ClearingService {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
private final Imdg<PaymentInstruction> paymentImdgs;
|
|
||||||
private final ExecutorService executor;
|
private final ExecutorService executor;
|
||||||
private final ImdgId idGenerator;
|
private final SdfCreatorBySTLDPayment sdfCreator;
|
||||||
private final PaymentInstructionSorter senderGroupSorter;
|
private final PaymentUpdateBySdf04 paymentUpdater;
|
||||||
private final Imdg<SDf03> sdf03Imdg;
|
|
||||||
private final Imdg<SDf11> sdf11Imdg;
|
|
||||||
private final KafkaSender kafkaSender;
|
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public ClearingService(ImdgProvider imdgProvider, PaymentInstructionSorter senderGroupSorter, KafkaSender kafkaSender) {
|
public ClearingService(SdfCreatorBySTLDPayment sdfCreator, PaymentUpdateBySdf04 paymentUpdater) {
|
||||||
this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
|
this.sdfCreator = sdfCreator;
|
||||||
this.sdf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class);
|
this.paymentUpdater = paymentUpdater;
|
||||||
this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class);
|
|
||||||
this.idGenerator = imdgProvider.getImdgIdGenerator();
|
|
||||||
this.senderGroupSorter = senderGroupSorter;
|
|
||||||
this.kafkaSender = kafkaSender;
|
|
||||||
this.executor = Executors.newSingleThreadExecutor();
|
this.executor = Executors.newSingleThreadExecutor();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Scheduled(cron = "${clearing-service.scheduler.check-payment-instruction}")
|
@Scheduled(cron = "${clearing-service.scheduler.check-payment-instruction}")
|
||||||
public void run() {
|
public void sdfCreate() {
|
||||||
executor.execute(this::createSdfFromPaymentInstructionSTLD);
|
log.info("creating sdf03/11 from STLD payments task added to queue");
|
||||||
|
executor.execute(sdfCreator::createSdfFromPaymentInstructionSTLD);
|
||||||
}
|
}
|
||||||
|
|
||||||
private void createSdfFromPaymentInstructionSTLD() {
|
public void paymentUpdateBySdf04(Long sdf04GroupId) {
|
||||||
final boolean[] anyError = {false};
|
log.info("adding sdf03/11 from STLD payments task to queue");
|
||||||
//generationId для созадаваемых Sdf03/Sdf11
|
executor.execute(() -> paymentUpdater.updatePayments(sdf04GroupId));
|
||||||
Long generationId = idGenerator.nextId();
|
|
||||||
//выгружаем PaymentInstructions с нужным статусом
|
|
||||||
Map<Long, PaymentBatchInfo> paymentBySender = paymentImdgs.getCollectionObjectsByFieldValues(
|
|
||||||
Map.of("transactionStatus", TransactionStatus.stld.getKey()))
|
|
||||||
.stream()
|
|
||||||
//группируем по компаниям (fixme sorted убрать?)
|
|
||||||
.sorted(Comparator.comparing(PaymentInstruction::getSenderId))
|
|
||||||
.collect(Collectors.groupingBy(PaymentInstruction::getSenderId))
|
|
||||||
.entrySet()
|
|
||||||
.stream()
|
|
||||||
//результатом работы senderGroupSorter будет Map<senderId -> PaymentBatchInfo>
|
|
||||||
//PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment
|
|
||||||
//тип ClearingMemberCategory
|
|
||||||
.map(entry -> {
|
|
||||||
Long senderId = entry.getKey();
|
|
||||||
List<PaymentInstruction> pmtInstrcs = entry.getValue();
|
|
||||||
PaymentBatchInfo senderInfo = senderGroupSorter.sortCompanyPayments(generationId, senderId, pmtInstrcs);
|
|
||||||
if (senderInfo.getError() != null) {
|
|
||||||
anyError[0] = true;
|
|
||||||
}
|
|
||||||
return new AbstractMap.SimpleEntry<>(senderId, senderInfo);
|
|
||||||
})
|
|
||||||
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
|
|
||||||
//save SDF03/SDF11
|
|
||||||
for (var entry : paymentBySender.entrySet()) {
|
|
||||||
PaymentBatchInfo senderPayments = entry.getValue();
|
|
||||||
saveSdfAnSendToKafka(senderPayments, generationId);
|
|
||||||
}
|
|
||||||
//update PaymentInstruction.transactionStatus
|
|
||||||
for (var entry : paymentBySender.entrySet()) {
|
|
||||||
PaymentBatchInfo batch = entry.getValue();
|
|
||||||
//все PaymentInstruction.transactionStatus в batch с error != null
|
|
||||||
//уже проапдейтились в методе sortCompanyPayments
|
|
||||||
if (batch.getError() == null) {
|
|
||||||
TransactionStatus stat = anyError[0] ? TransactionStatus.notSent : TransactionStatus.sent;
|
|
||||||
batch.getOrderedPaymentInstructions()
|
|
||||||
.forEach(paymentInstruction -> {
|
|
||||||
paymentInstruction.setTransactionStatus(stat.getKey());
|
|
||||||
paymentImdgs.update(paymentInstruction);
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
private void saveSdfAnSendToKafka(PaymentBatchInfo batch, Long generationId) {
|
|
||||||
SdfClearingRequest kafkaMessage = new SdfClearingRequest();
|
|
||||||
kafkaMessage.setGroupId(generationId);
|
|
||||||
switch (batch.getCategoryD()) {
|
|
||||||
case I -> {
|
|
||||||
batch.getOrderedPaymentInstructions()
|
|
||||||
.map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId))
|
|
||||||
.forEach(sdf03Imdg::insert);
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage);
|
|
||||||
}
|
|
||||||
case B -> {
|
|
||||||
batch.getOrderedPaymentInstructions()
|
|
||||||
.map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId))
|
|
||||||
.forEach(sdf11Imdg::insert);
|
|
||||||
kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,47 @@
|
||||||
|
package ru.spcex.clearing.service;
|
||||||
|
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||||
|
import ru.clearing.classes.statics.data.sdf.SDf04;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.platform.enumeration.TransactionStatus;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
import java.util.Collection;
|
||||||
|
import java.util.Map;
|
||||||
|
|
||||||
|
@Component
|
||||||
|
public class PaymentUpdateBySdf04 {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final Imdg<PaymentInstruction> paymentImdgs;
|
||||||
|
private final Imdg<SDf04> sdf04Imdg;
|
||||||
|
|
||||||
|
public PaymentUpdateBySdf04(ImdgProvider imdgProvider) {
|
||||||
|
this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
|
||||||
|
this.sdf04Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf04, SDf04.class);
|
||||||
|
}
|
||||||
|
|
||||||
|
public void updatePayments(Long sdf04GroupId) {
|
||||||
|
Collection<SDf04> sDf04s = sdf04Imdg.getCollectionObjectsByFieldValues(Map.of("generationId", sdf04GroupId));
|
||||||
|
for (SDf04 sDf04 : sDf04s) {
|
||||||
|
String docnm_ref = sDf04.getDocnm_ref();
|
||||||
|
long paymentId;
|
||||||
|
try {
|
||||||
|
paymentId = Long.parseLong(docnm_ref);
|
||||||
|
} catch (NumberFormatException e) {
|
||||||
|
log.warn("sdf04.id={} cannot parse docnm_ref={} as payment id", sDf04.getId(), sDf04.getDocnm_ref());
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
PaymentInstruction payment = paymentImdgs.getSingleObjectByFieldValues(Map.of("id", paymentId));
|
||||||
|
if ("OK!".equals(sDf04.getImp_result())) {
|
||||||
|
payment.setTransactionStatus(TransactionStatus.ok.getKey());
|
||||||
|
} else {
|
||||||
|
payment.setTransactionStatus(TransactionStatus.fail.getKey());
|
||||||
|
}
|
||||||
|
paymentImdgs.update(payment);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,28 @@
|
||||||
|
package ru.spcex.clearing.service;
|
||||||
|
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class Sdf04Receiver extends QueueConsumer implements InitializingBean {
|
||||||
|
private final ClearingService clearingService;
|
||||||
|
public Sdf04Receiver(Consumer<String, Object> kafkaQueue, ClearingService clearingService) {
|
||||||
|
super(kafkaQueue);
|
||||||
|
this.clearingService = clearingService;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() {
|
||||||
|
callback(Sdf04Request.class)
|
||||||
|
.setConsumer(event -> {
|
||||||
|
Sdf04Request requestPayload = event.getRequestPayload();
|
||||||
|
clearingService.paymentUpdateBySdf04(requestPayload.getGroupId());
|
||||||
|
})
|
||||||
|
.forDestination(Consts.SDF04_PROCESS, callbacks::put);
|
||||||
|
init();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,112 @@
|
||||||
|
package ru.spcex.clearing.service;
|
||||||
|
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||||
|
import ru.clearing.classes.statics.data.sdf.SDf03;
|
||||||
|
import ru.clearing.classes.statics.data.sdf.SDf11;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SdfClearingRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
import ru.spcex.clearing.service.order.PaymentBatchInfo;
|
||||||
|
import ru.spcex.clearing.service.order.PaymentInstructionSorter;
|
||||||
|
import ru.spcex.clearing.service.order.Sdf03Builder;
|
||||||
|
import ru.spcex.clearing.service.order.Sdf11Builder;
|
||||||
|
import ru.spcex.platform.enumeration.TransactionStatus;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgId;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
|
import java.util.AbstractMap;
|
||||||
|
import java.util.Comparator;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
|
@Component
|
||||||
|
public class SdfCreatorBySTLDPayment {
|
||||||
|
private final Imdg<PaymentInstruction> paymentImdgs;
|
||||||
|
private final ImdgId idGenerator;
|
||||||
|
private final PaymentInstructionSorter senderGroupSorter;
|
||||||
|
private final Imdg<SDf03> sdf03Imdg;
|
||||||
|
private final Imdg<SDf11> sdf11Imdg;
|
||||||
|
private final KafkaSender kafkaSender;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public SdfCreatorBySTLDPayment(ImdgProvider imdgProvider, PaymentInstructionSorter senderGroupSorter, KafkaSender kafkaSender) {
|
||||||
|
this.paymentImdgs = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
|
||||||
|
this.sdf03Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf03, SDf03.class);
|
||||||
|
this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class);
|
||||||
|
this.idGenerator = imdgProvider.getImdgIdGenerator();
|
||||||
|
this.senderGroupSorter = senderGroupSorter;
|
||||||
|
this.kafkaSender = kafkaSender;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void createSdfFromPaymentInstructionSTLD() {
|
||||||
|
final boolean[] anyError = {false};
|
||||||
|
//generationId для созадаваемых Sdf03/Sdf11
|
||||||
|
Long generationId = idGenerator.nextId();
|
||||||
|
//выгружаем PaymentInstructions с нужным статусом
|
||||||
|
Map<Long, PaymentBatchInfo> paymentBySender = paymentImdgs.getCollectionObjectsByFieldValues(
|
||||||
|
Map.of("transactionStatus", TransactionStatus.stld.getKey()))
|
||||||
|
.stream()
|
||||||
|
//группируем по компаниям (fixme sorted убрать?)
|
||||||
|
.sorted(Comparator.comparing(PaymentInstruction::getSenderId))
|
||||||
|
.collect(Collectors.groupingBy(PaymentInstruction::getSenderId))
|
||||||
|
.entrySet()
|
||||||
|
.stream()
|
||||||
|
//результатом работы senderGroupSorter будет Map<senderId -> PaymentBatchInfo>
|
||||||
|
//PaymentBatchInfo содержит возможную ошибку, при необходимости отсортированные Payment
|
||||||
|
//тип ClearingMemberCategory
|
||||||
|
.map(entry -> {
|
||||||
|
Long senderId = entry.getKey();
|
||||||
|
List<PaymentInstruction> pmtInstrcs = entry.getValue();
|
||||||
|
PaymentBatchInfo senderInfo = senderGroupSorter.sortCompanyPayments(generationId, senderId, pmtInstrcs);
|
||||||
|
if (senderInfo.getError() != null) {
|
||||||
|
anyError[0] = true;
|
||||||
|
}
|
||||||
|
return new AbstractMap.SimpleEntry<>(senderId, senderInfo);
|
||||||
|
})
|
||||||
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
|
||||||
|
//save SDF03/SDF11
|
||||||
|
for (var entry : paymentBySender.entrySet()) {
|
||||||
|
PaymentBatchInfo senderPayments = entry.getValue();
|
||||||
|
saveSdfAnSendToKafka(senderPayments, generationId);
|
||||||
|
}
|
||||||
|
//update PaymentInstruction.transactionStatus
|
||||||
|
for (var entry : paymentBySender.entrySet()) {
|
||||||
|
PaymentBatchInfo batch = entry.getValue();
|
||||||
|
//все PaymentInstruction.transactionStatus в batch с error != null
|
||||||
|
//уже проапдейтились в методе sortCompanyPayments
|
||||||
|
if (batch.getError() == null) {
|
||||||
|
TransactionStatus stat = anyError[0] ? TransactionStatus.notSent : TransactionStatus.sent;
|
||||||
|
batch.getOrderedPaymentInstructions()
|
||||||
|
.forEach(paymentInstruction -> {
|
||||||
|
paymentInstruction.setTransactionStatus(stat.getKey());
|
||||||
|
paymentImdgs.update(paymentInstruction);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
private void saveSdfAnSendToKafka(PaymentBatchInfo batch, Long generationId) {
|
||||||
|
SdfClearingRequest kafkaMessage = new SdfClearingRequest();
|
||||||
|
kafkaMessage.setGroupId(generationId);
|
||||||
|
switch (batch.getCategoryD()) {
|
||||||
|
case I -> {
|
||||||
|
batch.getOrderedPaymentInstructions()
|
||||||
|
.map(paymentInstruction -> Sdf03Builder.buildSdf03(paymentInstruction, generationId))
|
||||||
|
.forEach(sdf03Imdg::insert);
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.SDF03_PROCESS, kafkaMessage);
|
||||||
|
}
|
||||||
|
case B -> {
|
||||||
|
batch.getOrderedPaymentInstructions()
|
||||||
|
.map(paymentInstruction -> Sdf11Builder.buildSdf11(paymentInstruction, generationId))
|
||||||
|
.forEach(sdf11Imdg::insert);
|
||||||
|
kafkaSender.sendRequestToQueue(Consts.SDF11_PROCESS, kafkaMessage);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service.order;
|
||||||
|
|
||||||
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||||
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
|
import ru.spcex.platform.enumeration.ClearingMemberCategoryD;
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service.order;
|
||||||
|
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
@ -40,7 +40,7 @@ public class PaymentInstructionSorter {
|
||||||
this.messageResolver = messageResolver;
|
this.messageResolver = messageResolver;
|
||||||
}
|
}
|
||||||
|
|
||||||
PaymentBatchInfo sortCompanyPayments(Long generationId, Long senderId, List<PaymentInstruction> payments) {
|
public PaymentBatchInfo sortCompanyPayments(Long generationId, Long senderId, List<PaymentInstruction> payments) {
|
||||||
PaymentBatchInfo batchInfo = new PaymentBatchInfo();
|
PaymentBatchInfo batchInfo = new PaymentBatchInfo();
|
||||||
log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size());
|
log.info("processing PaymentInstruction's generationId={} senderId={} size={}", generationId, senderId, payments.size());
|
||||||
ClearingMemberCategory category = clrngMmbrImdg.getSingleObjectByFieldValues(Map.of("companyId", senderId));
|
ClearingMemberCategory category = clrngMmbrImdg.getSingleObjectByFieldValues(Map.of("companyId", senderId));
|
||||||
|
|
@ -1,7 +1,8 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service.order;
|
||||||
|
|
||||||
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf03;
|
import ru.clearing.classes.statics.data.sdf.SDf03;
|
||||||
|
import ru.spcex.clearing.service.SpecifUtil;
|
||||||
import ru.spcex.platform.enumeration.Sender;
|
import ru.spcex.platform.enumeration.Sender;
|
||||||
import ru.spcex.platform.utils.time.TimeUtil;
|
import ru.spcex.platform.utils.time.TimeUtil;
|
||||||
|
|
||||||
|
|
@ -1,7 +1,8 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service.order;
|
||||||
|
|
||||||
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
import ru.clearing.classes.statics.data.payment.PaymentInstruction;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf11;
|
import ru.clearing.classes.statics.data.sdf.SDf11;
|
||||||
|
import ru.spcex.clearing.service.SpecifUtil;
|
||||||
import ru.spcex.platform.enumeration.Sender;
|
import ru.spcex.platform.enumeration.Sender;
|
||||||
import ru.spcex.platform.utils.time.TimeUtil;
|
import ru.spcex.platform.utils.time.TimeUtil;
|
||||||
|
|
||||||
|
|
@ -36,10 +36,11 @@ public class UserRoleSessionMapStore extends TemplateMapStore<UserRoleSession> {
|
||||||
@Override
|
@Override
|
||||||
public UserRoleSession objectReader(ResultSet resultSet) throws SQLException {
|
public UserRoleSession objectReader(ResultSet resultSet) throws SQLException {
|
||||||
UserRoleSession object = new UserRoleSession();
|
UserRoleSession object = new UserRoleSession();
|
||||||
|
Long companyId = resultSet.getObject("COMPANY_ID", Long.class);
|
||||||
object.setId(resultSet.getObject("ID", Long.class));
|
object.setId(resultSet.getObject("ID", Long.class));
|
||||||
object.setUserId(resultSet.getObject("USER_ID", Long.class));
|
object.setUserId(resultSet.getObject("USER_ID", Long.class));
|
||||||
object.setUserRole(resultSet.getObject("USER_ROLE", String.class));
|
object.setUserRole(resultSet.getObject("USER_ROLE", String.class));
|
||||||
object.setCompanyId(resultSet.getObject("COMPANY_ID", Long.class));
|
object.setCompanyId(companyId.equals(0L) ? 1 : companyId);
|
||||||
object.setStatus(resultSet.getObject("STATUS", String.class));
|
object.setStatus(resultSet.getObject("STATUS", String.class));
|
||||||
return object;
|
return object;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,4 +4,4 @@ imdg.hazelcast.password=dev-pass
|
||||||
imdg.hazelcast.cluster-members[0]=127.0.0.1
|
imdg.hazelcast.cluster-members[0]=127.0.0.1
|
||||||
imdg.database.login=clearing
|
imdg.database.login=clearing
|
||||||
imdg.database.password=Aa111111
|
imdg.database.password=Aa111111
|
||||||
imdg.database.url=jdbc:postgresql://10.200.200.133:5432/postgres
|
imdg.database.url=jdbc:postgresql://10.200.200.133:5432/postgres?currentSchema=clearing_prod
|
||||||
|
|
@ -47,7 +47,7 @@ public class KeyRateService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
||||||
private void newKeyRate(BaseRequest<KeyRateNewRequest> userRequest) {
|
private void newKeyRate(BaseRequest<KeyRateNewRequest> userRequest) {
|
||||||
KeyRateNewRequest req = userRequest.getRequestPayload();
|
KeyRateNewRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("MoneyMarketSecurityNewRequest received");
|
log.debug("KeyRateNewRequest received");
|
||||||
KeyRate keyRate = new KeyRate();
|
KeyRate keyRate = new KeyRate();
|
||||||
keyRate.setEndDate(req.getEndDate());
|
keyRate.setEndDate(req.getEndDate());
|
||||||
keyRate.setDocument(req.getDocument());
|
keyRate.setDocument(req.getDocument());
|
||||||
|
|
@ -60,7 +60,7 @@ public class KeyRateService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
||||||
private void updateKeyRate(BaseRequest<KeyRateUpdateRequest> userRequest) {
|
private void updateKeyRate(BaseRequest<KeyRateUpdateRequest> userRequest) {
|
||||||
KeyRateUpdateRequest req = userRequest.getRequestPayload();
|
KeyRateUpdateRequest req = userRequest.getRequestPayload();
|
||||||
log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId());
|
log.debug("KeyRateUpdateRequest received id = {}", req.getId());
|
||||||
KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId());
|
KeyRate keyRate = keyRateMap.getSingleObjectByID(req.getId());
|
||||||
keyRate.setEndDate(req.getEndDate());
|
keyRate.setEndDate(req.getEndDate());
|
||||||
keyRate.setDocument(req.getDocument());
|
keyRate.setDocument(req.getDocument());
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
|
||||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
|
||||||
public enum TransactionStatus implements IEnumKey {
|
public enum TransactionStatus implements IEnumKey {
|
||||||
stld("STLD"), notSent("NSNT"), cher("CHER"), sent("SENT");
|
stld("STLD"), notSent("NSNT"), cher("CHER"), sent("SENT"), ok("OK"), fail("FAIL");
|
||||||
|
|
||||||
private final String key;
|
private final String key;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -50,6 +50,12 @@ public class QueueConsumer implements AutoCloseable {
|
||||||
this.supportStartOffsetTimeWindow = false;
|
this.supportStartOffsetTimeWindow = false;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* используем этот конструктор, если хотим класть
|
||||||
|
* в кафку "ответ" - информацию о статусе обработки команд
|
||||||
|
* @param kafkaQueue
|
||||||
|
* @param kafkaResponseQueue
|
||||||
|
*/
|
||||||
public QueueConsumer(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue) {
|
public QueueConsumer(Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue) {
|
||||||
this(kafkaQueue);
|
this(kafkaQueue);
|
||||||
this.producer = kafkaResponseQueue;
|
this.producer = kafkaResponseQueue;
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue