account-service utility-service http://jira.mfd.msk:8088/browse/CLS-527 DF52 ожидание акцепта пользователя через notification'sы
This commit is contained in:
parent
eef8e8499a
commit
9a97006dcb
5 changed files with 240 additions and 7 deletions
|
|
@ -15,8 +15,8 @@ import ru.clearing.classes.statics.data.account.ClearingAccount;
|
||||||
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
|
import ru.clearing.classes.statics.data.company.ClearingMemberCategory;
|
||||||
import ru.clearing.classes.statics.data.company.Company;
|
import ru.clearing.classes.statics.data.company.Company;
|
||||||
import ru.clearing.classes.statics.data.company.relation.Relation;
|
import ru.clearing.classes.statics.data.company.relation.Relation;
|
||||||
|
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf52;
|
import ru.clearing.classes.statics.data.sdf.SDf52;
|
||||||
import ru.clearing.platform.dictionary.AbstractDictionary;
|
|
||||||
import ru.clearing.platform.dictionary.ClearingCategoryDictionary;
|
import ru.clearing.platform.dictionary.ClearingCategoryDictionary;
|
||||||
import ru.spcex.clearing.account.errors.AccountError;
|
import ru.spcex.clearing.account.errors.AccountError;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
|
@ -28,6 +28,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf0
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
|
import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountSdfToStatementRequestPart;
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountSdfToStatementRequestPart;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationFeedbackRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
|
|
@ -42,7 +45,6 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||||
import ru.spcex.platform.utils.collection.Pair;
|
import ru.spcex.platform.utils.collection.Pair;
|
||||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
|
||||||
import ru.spcex.platform.utils.enumeration.IErrorEnumId;
|
import ru.spcex.platform.utils.enumeration.IErrorEnumId;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
import ru.spcex.platform.utils.validation.IValidator;
|
import ru.spcex.platform.utils.validation.IValidator;
|
||||||
|
|
@ -50,7 +52,6 @@ import ru.spcex.platform.utils.validation.IValidator;
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.util.*;
|
import java.util.*;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
import java.util.stream.Collectors;
|
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class ClearingAccountService extends QueueConsumer implements InitializingBean {
|
public class ClearingAccountService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
@ -70,6 +71,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
private final Imdg<Relation> relationImdg;
|
private final Imdg<Relation> relationImdg;
|
||||||
private final Imdg<ClearingMemberCategory> clearingMemberCategoryImdg;
|
private final Imdg<ClearingMemberCategory> clearingMemberCategoryImdg;
|
||||||
private final Imdg<ClearingCategoryDictionary> clearingCategoryImdg;
|
private final Imdg<ClearingCategoryDictionary> clearingCategoryImdg;
|
||||||
|
private final Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public ClearingAccountService(Consumer<String, Object> kafkaQueue,
|
public ClearingAccountService(Consumer<String, Object> kafkaQueue,
|
||||||
|
|
@ -101,6 +103,7 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
this.relationImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class);
|
this.relationImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Relation, Relation.class);
|
||||||
this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
|
this.clearingMemberCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingMemberCategory, ClearingMemberCategory.class);
|
||||||
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCategoryDictionary, ClearingCategoryDictionary.class);
|
this.clearingCategoryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingCategoryDictionary, ClearingCategoryDictionary.class);
|
||||||
|
this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -114,9 +117,13 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
callback(AccountSdf01Request.class)
|
callback(AccountSdf01Request.class)
|
||||||
.setFunction(this::accountNewSdf01)
|
.setFunction(this::accountNewSdf01)
|
||||||
.forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put);
|
.forDestination(Consts.ACCOUNT_NEW_SDF01, callbacks::put);
|
||||||
|
|
||||||
callback(StatementRequest.class)
|
callback(StatementRequest.class)
|
||||||
.setFunction(this::accountUpdateSdf52)
|
.setFunction(this::accountUpdateSdf52)
|
||||||
.forDestination(Consts.ACCOUNT_PROCESS_SDF52, callbacks::put);
|
.forDestination(Consts.ACCOUNT_PROCESS_SDF52, callbacks::put);
|
||||||
|
callback(NotificationFeedbackRequest.class)
|
||||||
|
.setFunction(this::accountUpdateSdf52_part2Notification)
|
||||||
|
.forDestination(Consts.ACCOUNT_NOTIFICATION_FEEDBACK, callbacks::put);
|
||||||
|
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
@ -359,7 +366,79 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
log.debug("Selected to update {} account's", toUpdate.size());
|
log.debug("Selected to update {} account's", toUpdate.size());
|
||||||
|
|
||||||
|
// 4. to notification
|
||||||
|
List<Long> notificationIds = new ArrayList<>();
|
||||||
|
boolean needWait = false;
|
||||||
|
{
|
||||||
|
for (Pair<SDf52, Account> item : toUpdate) {
|
||||||
|
SDf52 sdf = item.getFirst();
|
||||||
|
Account account = item.getSecond();
|
||||||
|
AccountStatus newStatus = sdfProcessService.parseSdf52Status(sdf.getStatus());
|
||||||
|
if (newStatus == null) {
|
||||||
|
throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus());
|
||||||
|
}
|
||||||
|
if (newStatus.equalsByKey(account.getStatus())) {
|
||||||
|
// одинаковых обычно не бывает.
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (AccountStatus.BLOCKED == newStatus || AccountStatus.CLOSE == newStatus) { // статус 0/2
|
||||||
|
notificationIds.add(sendNotificationRequest(ObjectType.account_block, account, newStatus));
|
||||||
|
needWait = true;
|
||||||
|
}
|
||||||
|
if (AccountStatus.ACTIVE == newStatus) { // статус 1/3
|
||||||
|
notificationIds.add(sendNotificationRequest(ObjectType.account_active, account, newStatus));
|
||||||
|
needWait = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (needWait) {
|
||||||
|
SDFProcessService.SDF52WaitingData data = new SDFProcessService.SDF52WaitingData(systemRequest, toProcessSDF53, toUpdate, notificationIds);
|
||||||
|
sdfProcessService.putNotificationWaiting(data);
|
||||||
|
log.info("Account's need wait accept over notification. Waiting user. GroupId={}", groupId);
|
||||||
|
} else {
|
||||||
|
log.info("No one account for wait notification reply. Do immediatly stage 2. groupId={}", groupId);
|
||||||
|
accountUpdateSdf52_part2(systemRequest, toProcessSDF53, toUpdate, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
log.debug("successfully processed, grouping id={} with {} accounts.",
|
||||||
|
groupId, toUpdate.size());
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
protected RequestInfoUpdate accountUpdateSdf52_part2Notification(BaseRequest<NotificationFeedbackRequest> secondSystemRequest2) {
|
||||||
|
log.debug("accountUpdateSdf52 NotificationFeedbackRequest received, id={}", secondSystemRequest2.getId());
|
||||||
|
SDFProcessService.SDF52WaitingData trigger = sdfProcessService.onNotificationResponse(secondSystemRequest2.getRequestPayload());
|
||||||
|
if (trigger == null) {
|
||||||
|
log.trace("Not triggered, continue waiting");
|
||||||
|
return null;
|
||||||
|
} else {
|
||||||
|
log.debug("Triggered, groupId={}, notification's status: {}, last notification user id={}",
|
||||||
|
trigger.getGroupId(), trigger.globalStatus, secondSystemRequest2.getUserId());
|
||||||
|
if (NotificationStatus.ACPT.equalsByKey(trigger.globalStatus)) {
|
||||||
|
return accountUpdateSdf52_part2(trigger.systemRequest, trigger.toProcessSDF53, trigger.toUpdate, secondSystemRequest2);
|
||||||
|
} else if (NotificationStatus.CNCL.equalsByKey(trigger.globalStatus)) {
|
||||||
|
log.debug("groupId={} rejected by user {}.", trigger.getGroupId(), secondSystemRequest2.getUserId());
|
||||||
|
return null;
|
||||||
|
} else {
|
||||||
|
log.warn("Unknown waiting group status \"{}\"", trigger.globalStatus);
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
protected RequestInfoUpdate accountUpdateSdf52_part2(BaseRequest<StatementRequest> systemRequest1,
|
||||||
|
List<Triple<SDf52, Account, String>> toProcessSDF53,
|
||||||
|
List<Pair<SDf52, Account>> toUpdate,
|
||||||
|
BaseRequest<NotificationFeedbackRequest> secondSystemRequest2) {
|
||||||
|
// 5. from notification:
|
||||||
|
StatementRequest req = systemRequest1.getRequestPayload();
|
||||||
|
Long groupId = req.getGroupId();
|
||||||
|
log.info("Continue Sdf52 groupId={}, first request id={}, second request id={}",
|
||||||
|
groupId, systemRequest1 == null ? null : systemRequest1.getId(), secondSystemRequest2 == null ? null : secondSystemRequest2.getId());
|
||||||
|
|
||||||
int countOfUpdated = 0;
|
int countOfUpdated = 0;
|
||||||
|
List<TradingClearingRegistryUpdateRequest> toTCRRequests = new ArrayList<>();
|
||||||
ImdgTransaction imdgTransaction = imdgProvider.newTransaction();
|
ImdgTransaction imdgTransaction = imdgProvider.newTransaction();
|
||||||
boolean txOk = false;
|
boolean txOk = false;
|
||||||
imdgTransaction.beginTransaction();
|
imdgTransaction.beginTransaction();
|
||||||
|
|
@ -374,14 +453,28 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
if (newStatus == null) {
|
if (newStatus == null) {
|
||||||
throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus());
|
throw new IllegalArgumentException("Can not parse sdf status " + sdf.getStatus());
|
||||||
}
|
}
|
||||||
|
// За долгое время ожидания пользователя счёт мог быть обновлён, прочитать его ещё раз
|
||||||
|
account = accountImdg.getSingleObjectByID(account.getId());
|
||||||
|
item.setSecond(account); // В toProcessSDF53 есть тоже account, но он на следующих этапах не используется
|
||||||
toProcessSDF53.add(new MutableTriple<>(sdf, account, SDFProcessService.SDF_STATUS_OK));
|
toProcessSDF53.add(new MutableTriple<>(sdf, account, SDFProcessService.SDF_STATUS_OK));
|
||||||
if (!newStatus.equalsByKey(account.getStatus())) {
|
if (!newStatus.equalsByKey(account.getStatus())) {
|
||||||
// Обновление счёта
|
// Обновление счёта
|
||||||
account.setStatus(newStatus.getKey());
|
account.setStatus(newStatus.getKey());
|
||||||
account.setUpdated(now);
|
account.setUpdated(now);
|
||||||
accountImdg.update(account);
|
accountImdg.update(account);
|
||||||
log.trace("S_DF52[{}] do update status to {} for account[{}]", sdf.getId(), newStatus.getKey(), account.getId());
|
log.trace("transaction {}, S_DF52[{}] do update status to {} for account[{}]",
|
||||||
|
imdgTransaction, sdf.getId(), newStatus.getKey(), account.getId());
|
||||||
|
if (AccountStatus.ACTIVE == newStatus) {
|
||||||
|
TradingClearingRegistryUpdateRequest tcrReq = new TradingClearingRegistryUpdateRequest();
|
||||||
|
tcrReq.setId(account.getId());
|
||||||
|
tcrReq.setCompanyId(account.getCompanyId());
|
||||||
|
// tcrReq.setId(); // select TCR; - выполнить поиск вне этой транзакции, чтобы
|
||||||
|
tcrReq.setStatus(account.getStatus());
|
||||||
|
toTCRRequests.add(tcrReq);
|
||||||
|
}
|
||||||
countOfUpdated++;
|
countOfUpdated++;
|
||||||
|
} else {
|
||||||
|
log.info("Account[{}] \"{}\" do not updated - same status \"{}\"", account.getId(), account.getAccount(), newStatus);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
txOk = true;
|
txOk = true;
|
||||||
|
|
@ -389,12 +482,28 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
if (txOk) {
|
if (txOk) {
|
||||||
imdgTransaction.commitTransaction();
|
imdgTransaction.commitTransaction();
|
||||||
} else {
|
} else {
|
||||||
log.debug("failed update, clearing accounts, rollback transaction. Request id={}; sdf52 groupId={}", systemRequest.getId(), groupId);
|
log.debug("failed update, clearing accounts, rollback transaction {}. Request id={}; sdf52 groupId={}",
|
||||||
|
imdgTransaction, systemRequest1 == null ? null : systemRequest1.getId(), groupId);
|
||||||
imdgTransaction.rollbackTransaction();
|
imdgTransaction.rollbackTransaction();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (txOk) {
|
if (txOk) {
|
||||||
|
for (TradingClearingRegistryUpdateRequest tcrReq : toTCRRequests) {
|
||||||
|
TradingClearingRegistry tcr = tradingClearingRegistryImdg.getFirstObjectByFieldValues(Map.of(
|
||||||
|
"moneyAccountId", tcrReq.getMoneyAccountId(),
|
||||||
|
"companyId", tcrReq.getCompanyId()));
|
||||||
|
if (tcr == null) {
|
||||||
|
log.warn("TradingClearingRegistry by moneyAccountId={} and companyId={} not found, do not send update TCR.",
|
||||||
|
tcrReq.getMoneyAccountId(), tcrReq.getCompanyId());
|
||||||
|
} else {
|
||||||
|
tcrReq.setId(tcr.getId());
|
||||||
|
String destination = Consts.DESTINATION_TRADING_CLEARING_REGISTRY_UPDATE;
|
||||||
|
log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(tcrReq));
|
||||||
|
kafkaSender.sendRequestToQueue(destination, tcrReq);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
sdfProcessService.process(req, toProcessSDF53);
|
sdfProcessService.process(req, toProcessSDF53);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -422,4 +531,34 @@ public class ClearingAccountService extends QueueConsumer implements Initializin
|
||||||
log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request));
|
log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request));
|
||||||
kafkaSender.sendRequestToQueue(destination, request);
|
kafkaSender.sendRequestToQueue(destination, request);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public void sendTCRUpdateRequest(Long groupingSdf01Id, Long groupSdf02Id, List<AccountSdfToStatementRequestPart> results) {
|
||||||
|
final String destination = Consts.STATEMENT_PROCESS;
|
||||||
|
StatementRequest request = new StatementRequest();
|
||||||
|
request.setGroupId(groupingSdf01Id);
|
||||||
|
request.setChildGenerationId(groupSdf02Id);
|
||||||
|
request.setAccountCreationResults(results);
|
||||||
|
request.setContinueSdf(true);
|
||||||
|
request.setTable(SdfTable.SDF_01); // по нему запрос получили
|
||||||
|
log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request));
|
||||||
|
kafkaSender.sendRequestToQueue(destination, request);
|
||||||
|
}
|
||||||
|
|
||||||
|
public Long sendNotificationRequest(ObjectType objectType, Account account, AccountStatus requestToStatus) {
|
||||||
|
String statusText = AccountStatus.BLOCKED == requestToStatus ? "Заблокирован"
|
||||||
|
: AccountStatus.CLOSE == requestToStatus ? "Закрыт"
|
||||||
|
: AccountStatus.ACTIVE == requestToStatus ? "Разблокирован" :
|
||||||
|
requestToStatus.getKey();
|
||||||
|
String message = String.format("Для счета «%s» будет изменен статус на «%s»", account.getAccount(), statusText);
|
||||||
|
final String destination = Consts.NOTIFICATION_NEW;
|
||||||
|
NotificationNewRequest request = new NotificationNewRequest();
|
||||||
|
request.setObjectType(objectType.getKey());
|
||||||
|
request.setObjectId(account.getId());
|
||||||
|
request.setPriority(Priority.HIGH.getKey());
|
||||||
|
request.setComment(message);
|
||||||
|
log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request));
|
||||||
|
Long rKey = kafkaSender.sendRequestToQueue(destination, request);
|
||||||
|
log.trace("About account.id={} send notification request id={}", account.getId(), rKey);
|
||||||
|
return rKey;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -9,18 +9,23 @@ import ru.clearing.classes.statics.data.account.Account;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf52;
|
import ru.clearing.classes.statics.data.sdf.SDf52;
|
||||||
import ru.clearing.classes.statics.data.sdf.SDf53;
|
import ru.clearing.classes.statics.data.sdf.SDf53;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationFeedbackRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
import ru.spcex.platform.enumeration.AccountStatus;
|
import ru.spcex.platform.enumeration.AccountStatus;
|
||||||
|
import ru.spcex.platform.enumeration.NotificationStatus;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgId;
|
import ru.spcex.platform.imdg.api.ImdgId;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.imdg.api.ImdgTransaction;
|
import ru.spcex.platform.imdg.api.ImdgTransaction;
|
||||||
|
import ru.spcex.platform.utils.collection.Pair;
|
||||||
|
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
import java.util.*;
|
import java.util.*;
|
||||||
|
import java.util.concurrent.CopyOnWriteArrayList;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class SDFProcessService {
|
public class SDFProcessService {
|
||||||
|
|
@ -144,4 +149,89 @@ public class SDFProcessService {
|
||||||
log.debug("Send ExportToFileRequest({}, {}) message id={} to kafka \"{}\"",
|
log.debug("Send ExportToFileRequest({}, {}) message id={} to kafka \"{}\"",
|
||||||
groupId, request.getNameOfTable(), msgId, destination);
|
groupId, request.getNameOfTable(), msgId, destination);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --------- notification apply system -----------
|
||||||
|
public static class SDF52WaitingData {
|
||||||
|
public BaseRequest<StatementRequest> systemRequest;
|
||||||
|
public List<Triple<SDf52, Account, String>> toProcessSDF53;
|
||||||
|
public List<Pair<SDf52, Account>> toUpdate;
|
||||||
|
public HashSet<Long> notAnsweredNotifications = new HashSet<>();
|
||||||
|
public String globalStatus;
|
||||||
|
|
||||||
|
public SDF52WaitingData(BaseRequest<StatementRequest> systemRequest, List<Triple<SDf52, Account, String>> toProcessSDF53, List<Pair<SDf52, Account>> toUpdate,
|
||||||
|
Collection<Long> notAnsweredNotifications) {
|
||||||
|
this.systemRequest = systemRequest;
|
||||||
|
this.toProcessSDF53 = toProcessSDF53;
|
||||||
|
this.toUpdate = toUpdate;
|
||||||
|
this.notAnsweredNotifications = new HashSet<>(notAnsweredNotifications);
|
||||||
|
}
|
||||||
|
|
||||||
|
public Long getGroupId() {
|
||||||
|
if (systemRequest != null && systemRequest.getRequestPayload() != null)
|
||||||
|
return systemRequest.getRequestPayload().getGroupId();
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
protected CopyOnWriteArrayList<SDF52WaitingData> waitingList = new CopyOnWriteArrayList<>();
|
||||||
|
|
||||||
|
public boolean putNotificationWaiting(SDF52WaitingData data) {
|
||||||
|
for (SDF52WaitingData i : waitingList) {
|
||||||
|
if (Objects.equals(i.getGroupId(), data.getGroupId())) {
|
||||||
|
log.warn("groupId={} already waiting", data.getGroupId());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
data.globalStatus = null;
|
||||||
|
waitingList.add(data);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @return triggered SDF52WaitingData or null
|
||||||
|
*/
|
||||||
|
public SDF52WaitingData onNotificationResponse(NotificationFeedbackRequest onNotification) {
|
||||||
|
final Long nId = onNotification.getNotificationId();
|
||||||
|
if (nId == null)
|
||||||
|
return null;
|
||||||
|
SDF52WaitingData inWaitingLst = null;
|
||||||
|
for (SDF52WaitingData i : waitingList)
|
||||||
|
synchronized (i) {
|
||||||
|
if (i.notAnsweredNotifications.contains(nId)) {
|
||||||
|
inWaitingLst = i;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (inWaitingLst == null) {
|
||||||
|
log.debug("SDF52 waiting list not found for notificationId={}", nId);
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
boolean returnTrigger = false;
|
||||||
|
synchronized (inWaitingLst) {
|
||||||
|
boolean removedId = inWaitingLst.notAnsweredNotifications.remove(nId);
|
||||||
|
if (!removedId) { // never
|
||||||
|
log.warn("Can not remove notificationId={} from waiting list", nId);
|
||||||
|
}
|
||||||
|
if (NotificationStatus.ACPT.equalsByKey(onNotification.getNotificationStatus())) {
|
||||||
|
if (inWaitingLst.notAnsweredNotifications.isEmpty()) {
|
||||||
|
log.trace("Waiting list is empty, by notificationId={}", nId);
|
||||||
|
inWaitingLst.globalStatus = onNotification.getNotificationStatus();
|
||||||
|
returnTrigger = true;
|
||||||
|
}
|
||||||
|
} else if (NotificationStatus.CNCL.equalsByKey(onNotification.getNotificationStatus())) {
|
||||||
|
log.trace("Waiting list cancel, by notificationId={}", nId);
|
||||||
|
inWaitingLst.globalStatus = onNotification.getNotificationStatus();
|
||||||
|
returnTrigger = true;
|
||||||
|
} else {
|
||||||
|
log.warn("Unknown notification[{}] status={}", onNotification.getNotificationId(), onNotification.getNotificationStatus());
|
||||||
|
}
|
||||||
|
if (returnTrigger || inWaitingLst.notAnsweredNotifications.isEmpty()) {
|
||||||
|
log.debug("Remove waiting list by notificationId={} (list groupId={})", nId, inWaitingLst.getGroupId());
|
||||||
|
waitingList.remove(inWaitingLst);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (returnTrigger)
|
||||||
|
return inWaitingLst;
|
||||||
|
else
|
||||||
|
return null;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,7 @@ import java.time.LocalDate;
|
||||||
|
|
||||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.*;
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.*;
|
||||||
import static ru.spcex.platform.enumeration.NotificationStatus.PEND;
|
import static ru.spcex.platform.enumeration.NotificationStatus.PEND;
|
||||||
import static ru.spcex.platform.enumeration.ObjectType.statement;
|
import static ru.spcex.platform.enumeration.ObjectType.*;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class NotificationService extends QueueConsumer implements InitializingBean {
|
public class NotificationService extends QueueConsumer implements InitializingBean {
|
||||||
|
|
@ -94,6 +94,8 @@ public class NotificationService extends QueueConsumer implements InitializingBe
|
||||||
String objectType = notification.getObjectType();
|
String objectType = notification.getObjectType();
|
||||||
if (objectType.equalsIgnoreCase(statement.getKey())) {
|
if (objectType.equalsIgnoreCase(statement.getKey())) {
|
||||||
kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(notification.getId(), notification.getNotificationStatus()));
|
kafkaSender.sendRequestToQueue(CLEARING_NOTIFICATION_FEEDBACK, buildFeedbackRequest(notification.getId(), notification.getNotificationStatus()));
|
||||||
|
} else if (objectType.equalsIgnoreCase(account_active.getKey()) || objectType.equalsIgnoreCase(account_block.getKey())) {
|
||||||
|
kafkaSender.sendRequestToQueue(ACCOUNT_NOTIFICATION_FEEDBACK, buildFeedbackRequest(notification.getId(), notification.getNotificationStatus()));
|
||||||
} else {
|
} else {
|
||||||
log.warn("Unsupported notification type {}", objectType);
|
log.warn("Unsupported notification type {}", objectType);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,8 @@ package ru.spcex.platform.enumeration;
|
||||||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
|
||||||
public enum ObjectType implements IEnumKey {
|
public enum ObjectType implements IEnumKey {
|
||||||
statement("STMT"), vfrs("VFRS"), rgst("RGST"), gateway("GTWY"), session("SESN");
|
statement("STMT"), vfrs("VFRS"), rgst("RGST"), gateway("GTWY"), session("SESN"),
|
||||||
|
account_block("ACCB"), account_active("ACCA");
|
||||||
|
|
||||||
private final String key;
|
private final String key;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -111,6 +111,7 @@ public interface Consts {
|
||||||
String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block";
|
String DESTINATION_TRADING_CLEARING_REGISTRY_BLOCK = "trading-clearing-registry-block";
|
||||||
|
|
||||||
String CLEARING_NOTIFICATION_FEEDBACK = "clearing-notification-feedback";
|
String CLEARING_NOTIFICATION_FEEDBACK = "clearing-notification-feedback";
|
||||||
|
String ACCOUNT_NOTIFICATION_FEEDBACK = "account-notification-feedback";
|
||||||
String NOTIFICATION_NEW = "notification-new";
|
String NOTIFICATION_NEW = "notification-new";
|
||||||
String NOTIFICATION_UPDATE = "notification-update";
|
String NOTIFICATION_UPDATE = "notification-update";
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue