context) {
+ SDf57 sdf57 = context.getValidatedObject();
+ if (!ru.spcex.platform.enumeration.CurrencyCode.RUR.equalsByKey(sdf57.getPay_val())) {
+ return of(BalanceError.CurrencyNotFound, sdf57.getPay_val());
+ }
+ return empty();
+ }
+ };
+
+ @Override
+ public String ruleName() {
+ return "Sdf57ValidationRule." + name();
+ }
+}
diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/ValidationStored.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/ValidationStored.java
index 25eaddc69..b6074f057 100644
--- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/ValidationStored.java
+++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/validation/ValidationStored.java
@@ -1,5 +1,9 @@
package ru.spcex.clearing.balance.validation;
public enum ValidationStored {
- Account, Company;
+ Account, Company,
+
+ Sdf57CompanyDeb, Sdf57CompanyCred, Sdf57AccountDeb, Sdf57AccountCred
+ ;
+
}
diff --git a/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/register/MoneyBalanceRegister.java b/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/register/MoneyBalanceRegister.java
index 90dd23ad5..dff2d1116 100644
--- a/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/register/MoneyBalanceRegister.java
+++ b/clearing-parent/classes/src/main/java/ru/clearing/classes/statics/data/register/MoneyBalanceRegister.java
@@ -7,7 +7,7 @@ import java.math.BigDecimal;
/**
* Реестр остатков денежных средств
- *
+ *
* DB table: MONEY_BALANCE_REGISTER
**/
public class MoneyBalanceRegister extends BusinessObject {
@@ -24,74 +24,108 @@ public class MoneyBalanceRegister extends BusinessObject {
private String companyFullName;
private Long companyId;
+ public static MoneyBalanceRegister makeMoneyBalanceRegister(String setHouseName,
+ String account,
+ String infoAccount,
+ BigDecimal remainderSum,
+ BigDecimal blockedSum,
+ BigDecimal unblockedSum,
+ String inn,
+ Long sessionId,
+ String companyFullName,
+ Long companyId) {
+ MoneyBalanceRegister moneyBalanceRegister = new MoneyBalanceRegister();
+ moneyBalanceRegister.setSetHouseName(setHouseName);
+ moneyBalanceRegister.setAccount(account);
+ moneyBalanceRegister.setInfoAccount(infoAccount);
+ moneyBalanceRegister.setRemainderSum(remainderSum);
+ moneyBalanceRegister.setBlockedSum(blockedSum);
+ moneyBalanceRegister.setUnblockedSum(unblockedSum);
+ moneyBalanceRegister.setInn(inn);
+ moneyBalanceRegister.setSessionId(sessionId);
+ moneyBalanceRegister.setCompanyFullName(companyFullName);
+ moneyBalanceRegister.setCompanyId(companyId);
+ return moneyBalanceRegister;
+ }
+
public String getSetHouseName() {
return setHouseName;
}
+
public void setSetHouseName(String value) {
- this.setHouseName=value;
+ this.setHouseName = value;
}
public String getAccount() {
return account;
}
+
public void setAccount(String value) {
- this.account=value;
+ this.account = value;
}
public String getInfoAccount() {
return infoAccount;
}
+
public void setInfoAccount(String value) {
- this.infoAccount=value;
+ this.infoAccount = value;
}
public BigDecimal getRemainderSum() {
return remainderSum;
}
+
public void setRemainderSum(BigDecimal value) {
- this.remainderSum=value;
+ this.remainderSum = value;
}
public BigDecimal getBlockedSum() {
return blockedSum;
}
+
public void setBlockedSum(BigDecimal value) {
- this.blockedSum=value;
+ this.blockedSum = value;
}
public BigDecimal getUnblockedSum() {
return unblockedSum;
}
+
public void setUnblockedSum(BigDecimal value) {
- this.unblockedSum=value;
+ this.unblockedSum = value;
}
public String getInn() {
return inn;
}
+
public void setInn(String value) {
- this.inn=value;
+ this.inn = value;
}
public Long getSessionId() {
return sessionId;
}
+
public void setSessionId(Long value) {
- this.sessionId=value;
+ this.sessionId = value;
}
public String getCompanyFullName() {
return companyFullName;
}
+
public void setCompanyFullName(String value) {
- this.companyFullName=value;
+ this.companyFullName = value;
}
public Long getCompanyId() {
return companyId;
}
+
public void setCompanyId(Long value) {
- this.companyId=value;
+ this.companyId = value;
}
}
\ No newline at end of file
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java
index f4677be37..95f979e72 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java
@@ -8,11 +8,13 @@ import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.company.relation.Relation;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.misc.STrades;
+import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.validation.ClrngValidationStored;
import ru.spcex.clearing.service.validation.ExecutionDepositValidationRule;
+import ru.spcex.clearing.service.validation.RegistryStep3ValidationRule;
import ru.spcex.clearing.service.validation.STradesValidationRule;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.ClearingCategory;
@@ -86,4 +88,23 @@ public class ValidationConfig {
STradesValidationRule.TradingClearingRegistryPresent);
};
}
+
+ @Bean("obligationAndRequirementsAdmissionValidator")
+ public Function registryValidator() {
+ return rgs -> {
+ ImdgValidationContext context = new ImdgValidationContext<>();
+ context.setValidatedObject(rgs);
+ addImdg.accept(context, IMDGDistributedNames.Map_Relation);
+ addImdg.accept(context, IMDGDistributedNames.Map_Session);
+ addImdg.accept(context, IMDGDistributedNames.Map_SectionDictionary);
+ addImdg.accept(context, IMDGDistributedNames.Map_Account);
+ addImdg.accept(context, IMDGDistributedNames.Map_Company);
+ addImdg.accept(context, IMDGDistributedNames.Map_TradingClearingRegistry);
+ return new ValidatorImpl<>(context,
+ RegistryStep3ValidationRule.ClearingAvailable,
+ RegistryStep3ValidationRule.AccountActive,
+ RegistryStep3ValidationRule.CompanyActive,
+ RegistryStep3ValidationRule.TradingClearingRegistryActive);
+ };
+ }
}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java
index 7b8cb5e01..9dd2a2595 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/error/ClearingError.java
@@ -5,12 +5,15 @@ import ru.spcex.platform.utils.enumeration.IErrorEnumId;
public enum ClearingError implements IErrorEnumId {
GeneralError(5400L),
RecordNotFound(5406L),
+ CompanyNotActive(5411L),
CompanyCreditCheck(5412L),
CompanyDebitCheck(5413L),
CompanyNotFound(5410L),
+ AccountNotActive(5415L),
SecurityNotFound(5416L),
TradingClearingRegistryNotFound(5418L),
TradingClearingRegistryNotActive(5419L),
+ ClearingUnavailableForCompany(5421L),
InsecurityObligation(5422L),
NewDealsNotFound(5423L),
;
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/models/RegistryInfo.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/models/RegistryInfo.java
deleted file mode 100644
index 1c618492f..000000000
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/models/RegistryInfo.java
+++ /dev/null
@@ -1,16 +0,0 @@
-package ru.spcex.clearing.models;
-
-import ru.spcex.platform.enumeration.RegistryCapacity;
-import ru.spcex.platform.enumeration.RegistryDesignation;
-import ru.spcex.platform.enumeration.RegistryInstrumentType;
-import ru.spcex.platform.enumeration.RegistryUnit;
-
-public record RegistryInfo(RegistryDesignation registryDesignation,
- RegistryInstrumentType registryInstrumentType,
- RegistryCapacity registryCapacity,
- RegistryUnit registryUnit) {
-
- public String buildSqlCondition(){
- return "";
- }
-}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java
index 19774c888..758b9be7b 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java
@@ -9,19 +9,23 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
+import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
import ru.spcex.platform.enumeration.Task;
@Service
public class EventsReceiver extends QueueConsumer implements InitializingBean {
private final ClearingService clearingService;
private final RegistryService registryService;
+ private final PrimaryAuctionBnSession primaryAuctionBnSession;
public EventsReceiver(Consumer kafkaQueue,
ClearingService clearingService,
- RegistryService registryService) {
+ RegistryService registryService,
+ PrimaryAuctionBnSession primaryAuctionBnSession) {
super(kafkaQueue);
this.clearingService = clearingService;
this.registryService = registryService;
+ this.primaryAuctionBnSession = primaryAuctionBnSession;
}
@Override
@@ -38,12 +42,15 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(LauncherCommandRequest.class)
.setConsumer(event -> clearingService.executeVerification())
.forDestination(Task.getVerification.topic(), callbacks::put);
- callback(Object.class) //todo check Object suitable
- .setConsumer(event -> clearingService.executeClearing())
+ callback(Object.class)
+ .setConsumer(primaryAuctionBnSession::runSession)
.forDestination(Task.startOfClearing.topic(), callbacks::put);
callback(CommonIdRequest.class)
.setConsumer(clearingService::continueClearing)
.forDestination(Consts.CONTINUE_CLEARING, callbacks::put);
+ callback(Object.class)
+ .setConsumer(primaryAuctionBnSession::continueSession)
+ .forDestination(Consts.SDF57_PROCESS, callbacks::put);
callback(Object.class)
.setConsumer(event -> clearingService.executeSTrade())
.forDestination(Task.getOfTrades.topic(), callbacks::put);
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java
index 81cddeb36..d7756f23a 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/ExecutionDepositComponent.java
@@ -54,7 +54,7 @@ public class ExecutionDepositComponent {
//fixme ждать ТЗ
Long tradeNum;
Instant tradingDay;
- private final IMessageResolver msgResolver = new SimpleMessageResolver();
+ private final IMessageResolver msgResolver;
private final Function stradesValidator;
private final KafkaSender kafkaSender;
@@ -62,12 +62,14 @@ public class ExecutionDepositComponent {
@Autowired
public ExecutionDepositComponent(ImdgProvider imdgProvider, Producer kafka,
@Qualifier("sTradesValidator") Function stradesValidator,
- @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender) {
+ @Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender,
+ IMessageResolver msgResolver) {
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
this.stradesValidator = stradesValidator;
this.kafkaSender = kafkaSender;
+ this.msgResolver = msgResolver;
resetTradingDay();
}
@@ -127,7 +129,7 @@ public class ExecutionDepositComponent {
}
ExecutionDeposit newED;
try {
- newED = createExecutionDeposit(sTrd, validator);
+ newED = createExecutionDeposit(sTrd, validator);
executionDepositImdg.insert(newED);
sendNotification(newED);
log.debug("New executionDeposit.id={} was created.", newED.getId());
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java
index f6ab3915f..d70722db4 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java
@@ -6,7 +6,6 @@ import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
-import ru.spcex.clearing.models.RegistryInfo;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRequest;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.imdg.api.Imdg;
@@ -42,7 +41,7 @@ public class RegistryService {
log.trace("Found tradingClearingRegistry with id: {}", tcrByCompanyId.getId());
Registry registry = new Registry();
tcrByCompanyId.getTradingClearingRegistryType();
- RegistryInfo registryInfo = createRegistryInfoByType(tcrByCompanyId.getTradingClearingRegistryType());
+ RegistryTradingParams registryInfo = createRegistryInfoByType(tcrByCompanyId.getTradingClearingRegistryType());
registry.setRegistryDesignation(registryInfo.registryDesignation().getKey());
registry.setRegistryInstrumentType(registryInfo.registryInstrumentType().getKey());
registry.setRegistryCapacity(registryInfo.registryCapacity().getKey());
@@ -55,20 +54,20 @@ public class RegistryService {
}
- private RegistryInfo createRegistryInfoByType(String type){
+ private RegistryTradingParams createRegistryInfoByType(String type){
TradingClearingRegistryType registryType = IEnumKey.getEnumByKey(TradingClearingRegistryType.class, type);
- RegistryInfo registryInfo = null;
+ RegistryTradingParams registryInfo = null;
//todo спросить про дубль C и T
switch (registryType){
- case Owner_A -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.M, RegistryCapacity.A, RegistryUnit.T);
- case Client_B -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.M, RegistryCapacity.B, RegistryUnit.T);
- case Trustee_C, DepoOfShares_T -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.T);
- case TrusteeManager_D -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.A);
- case Issuer_E -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.E, RegistryUnit.R);
- case ToThePlacementOrRedeem_Z -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.Z, RegistryUnit.R);
- case DepoBond_H -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.R);
- case DepoOfTrustedBond_N -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.C, RegistryUnit.R);
- case DepoOfTrustedShares_S -> registryInfo = new RegistryInfo(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.C, RegistryUnit.T);
+ case Owner_A -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.M, RegistryCapacity.A, RegistryUnit.T);
+ case Client_B -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.M, RegistryCapacity.B, RegistryUnit.T);
+ case Trustee_C, DepoOfShares_T -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.T);
+ case TrusteeManager_D -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.A);
+ case Issuer_E -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.E, RegistryUnit.R);
+ case ToThePlacementOrRedeem_Z -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.Z, RegistryUnit.R);
+ case DepoBond_H -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.R);
+ case DepoOfTrustedBond_N -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.C, RegistryUnit.R);
+ case DepoOfTrustedShares_S -> registryInfo = new RegistryTradingParams(RegistryDesignation.A, RegistryInstrumentType.S, RegistryCapacity.C, RegistryUnit.T);
}
return registryInfo;
}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java
new file mode 100644
index 000000000..495574f9e
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryStep3ValidationRule.java
@@ -0,0 +1,104 @@
+package ru.spcex.clearing.service.validation;
+
+import ru.clearing.classes.statics.data.account.Account;
+import ru.clearing.classes.statics.data.company.Company;
+import ru.clearing.classes.statics.data.company.relation.Relation;
+import ru.clearing.classes.statics.data.misc.Session;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
+import ru.clearing.platform.dictionary.SectionDictionary;
+import ru.spcex.clearing.error.ClearingError;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.platform.enumeration.ServiceStatus;
+import ru.spcex.platform.enumeration.WorkflowStatus;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.validation.ImdgValidationContext;
+import ru.spcex.platform.utils.enumeration.EnumMessage;
+import ru.spcex.platform.utils.enumeration.IEnumKey;
+import ru.spcex.platform.utils.validation.IValidationRule;
+
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.Supplier;
+
+public enum RegistryStep3ValidationRule implements IValidationRule> {
+ ClearingAvailable() {
+ @Override
+ public Optional validate(ImdgValidationContext context) {
+ Registry validatedObject = context.getValidatedObject();
+ Supplier sectionFind = () -> {
+ Imdg sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class);
+ Imdg sectionDictionaryImdg = context.obtainMap(IMDGDistributedNames.Map_SectionDictionary, SectionDictionary.class);
+ Session session = sessionImdg.getSingleObjectByID(validatedObject.getSessionId());
+ if (session == null) {
+ return "null";
+ } else {
+ SectionDictionary section = sectionDictionaryImdg.getSingleObjectBySQL("code ='" + session.getSection() + "'");
+ if (section == null) {
+ return "null";
+ } else {
+ return section.getName();
+ }
+ }
+ };
+ if (validatedObject.getCompanyId() == null) {
+ return of(ClearingError.ClearingUnavailableForCompany, sectionFind.get());
+ }
+ Imdg relationImdg = context.obtainMap(IMDGDistributedNames.Map_Relation, Relation.class);
+ Relation relation = relationImdg.getSingleObjectByFieldValues(Map.of("consumerId", validatedObject.getCompanyId()));
+ if (relation == null || (!ServiceStatus.Active.equalsByKey(relation.getServiceStatus()) && !ServiceStatus.Reopened.equalsByKey(relation.getServiceStatus()))) {
+ return of(ClearingError.ClearingUnavailableForCompany, sectionFind.get(), validatedObject.getCompanyId());
+ }
+ context.storeObject(RegistryValidationStored.Relation, relation);
+ return empty();
+ }
+ },
+ AccountActive() {
+ @Override
+ public Optional validate(ImdgValidationContext context) {
+ Registry validatedObject = context.getValidatedObject();
+ Relation relation = context.getStoredObject(RegistryValidationStored.Relation);
+ Imdg accountImdg = context.obtainMap(IMDGDistributedNames.Map_Account, Account.class);
+ Account account = accountImdg.getSingleObjectBySQL("id = %d and relationId = %d".formatted(validatedObject.getAccountId(), relation.getId()));
+ if (account == null || (!IEnumKey.contains(account.getStatus(), ServiceStatus.Active, ServiceStatus.Reopened))) {
+ return of(ClearingError.AccountNotActive, validatedObject.getAccountId());
+ }
+ return empty();
+ }
+ },
+ CompanyActive() {
+ @Override
+ public Optional validate(ImdgValidationContext context) {
+ Registry validatedObject = context.getValidatedObject();
+ Imdg companyImdg = context.obtainMap(IMDGDistributedNames.Map_Company, Company.class);
+ if (validatedObject.getCompanyId() == null) {
+ return of(ClearingError.CompanyNotActive, validatedObject.getCompanyId());
+ }
+ Company company = companyImdg.getSingleObjectBySQL("id = %d".formatted(validatedObject.getCompanyId()));
+ if (company == null || !WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
+ return of(ClearingError.CompanyNotActive, validatedObject.getCompanyId());
+ }
+ return empty();
+ }
+ },
+ TradingClearingRegistryActive() {
+ @Override
+ public Optional validate(ImdgValidationContext context) {
+ Registry validatedObject = context.getValidatedObject();
+ if (validatedObject.getTradingClearingRegistryId() == null) {
+ return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId());
+ }
+ Imdg tradingClearingRegistryImdg = context.obtainMap(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
+ TradingClearingRegistry tcr = tradingClearingRegistryImdg.getSingleObjectBySQL("id = " + validatedObject.getTradingClearingRegistryId());
+ if (tcr == null || (ServiceStatus.Active.equalsByKey(tcr.getStatus()) || ServiceStatus.Reopened.equalsByKey(tcr.getStatus())))
+ return of(ClearingError.TradingClearingRegistryNotActive, validatedObject.getTradingClearingRegistryId());
+ return empty();
+ }
+ },
+ ;
+
+ @Override
+ public String ruleName() {
+ return "STradesValidationRule." + name();
+ }
+}
\ No newline at end of file
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryValidationStored.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryValidationStored.java
new file mode 100644
index 000000000..0d7bba6b7
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/RegistryValidationStored.java
@@ -0,0 +1,5 @@
+package ru.spcex.clearing.service.validation;
+
+public enum RegistryValidationStored {
+ Relation
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java
new file mode 100644
index 000000000..3d7a35f1d
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/PrimaryAuctionBnSession.java
@@ -0,0 +1,200 @@
+package ru.spcex.clearing.session.stage;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.execution.ExecutionCommon;
+import ru.clearing.classes.statics.data.misc.Session;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
+import ru.spcex.clearing.session.stage.impl.*;
+import ru.spcex.clearing.session.stage.task.*;
+import ru.spcex.platform.enumeration.Section;
+import ru.spcex.platform.enumeration.SessionStatus;
+import ru.spcex.platform.enumeration.SessionType;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.utils.enumeration.IMessageResolver;
+
+import java.util.List;
+import java.util.concurrent.atomic.AtomicReference;
+
+@Service
+public class PrimaryAuctionBnSession {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+ private final BalanceRevise balanceRevise;
+ private final DealsPrepare dealsPrepare;
+ private final RequirementsAndObligationCreation requirementsAndObligationCreation;
+ private final ObligationAdmission obligationsAdmission;
+
+ private final InclusionObligations inclusionObligations;
+ private final FormingRegistersOnOS formingRegistersOnOS;
+ private final FormingPaymentInstruction formingPaymentInstruction;
+ private final UnlockResources unlockResources;
+ private final FinishingSession finishingSession;
+ private final EndStageNotification endStageNotification;
+
+ private Imdg sessionImdg;
+ private final IMessageResolver messageResolver;
+ private final AtomicReference currStage = new AtomicReference<>();
+ private Session currSession;
+
+ public PrimaryAuctionBnSession(
+ ImdgProvider imdgProvider,
+ BalanceRevise balanceRevise,
+ DealsPrepare dealsPrepare,
+ RequirementsAndObligationCreation requirementsAndObligationCreation,
+ ObligationAdmission obligationsAdmission,
+ InclusionObligations inclusionObligations,
+ FormingRegistersOnOS formingRegistersOnOS,
+ FormingPaymentInstruction formingPaymentInstruction,
+ UnlockResources unlockResources,
+ FinishingSession finishingSession, EndStageNotification endStageNotification, IMessageResolver messageResolver) {
+ this.balanceRevise = balanceRevise;
+ this.dealsPrepare = dealsPrepare;
+ this.requirementsAndObligationCreation = requirementsAndObligationCreation;
+ this.obligationsAdmission = obligationsAdmission;
+ this.inclusionObligations = inclusionObligations;
+ this.formingRegistersOnOS = formingRegistersOnOS;
+ this.formingPaymentInstruction = formingPaymentInstruction;
+ this.unlockResources = unlockResources;
+ this.finishingSession = finishingSession;
+ this.endStageNotification = endStageNotification;
+
+ this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
+ this.messageResolver = messageResolver;
+ }
+
+ public void runSession(BaseRequest> req) {
+ if (!startSession()) {
+ return;
+ }
+ StageResult> submit = balanceRevise.submit(new Task<>(TaskType.StartRevise, null));
+ if (!submit.success) {
+ log.error("stage {} error={}", balanceRevise.getClass().getSimpleName(), messageResolver.resolve(submit.error));
+ endSession();
+ } else {
+ log.info("stage BalanceRevise success, waiting for a response from kafka");
+ }
+ }
+
+ public void continueSession(BaseRequest> req) {
+ try {
+ if (!checkStage(TaskType.StartRevise)) {
+ log.error("cannot continue session, current stage is {}", currStage.get());
+ throw new StageException();
+ }
+ //stage 0
+ runStage(TaskType.ContinueRevise, balanceRevise);
+ //stage 1
+ StageResult> dealsPreparationResult;
+ {
+ DealsPreparePayload payload = new DealsPreparePayload();
+ payload.setSessionId(currSession.getId());
+ dealsPreparationResult = runStage(TaskType.DealsPrepare , payload, dealsPrepare);
+ }
+ //stage 2
+ runStage(TaskType.RequirementsAndObligationsCreate, dealsPreparationResult.getStageResult(), requirementsAndObligationCreation);
+ //stage 3
+ runStage(TaskType.ObligationsAdmission, currSession.getId(), obligationsAdmission);
+ //stage 4
+ {
+ InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload();
+ inclusionToPoolPayload.setSessionType(currSession.getSessionType());
+ runStage(TaskType.InclusionToPool, inclusionToPoolPayload, inclusionObligations);
+ }
+ //stage 5
+ {
+ InspectionPoolPayload companyIdPayload = new InspectionPoolPayload();
+ companyIdPayload.setProcessedCompanyId(currSession.getCompanyId());
+ runStage(TaskType.InspectionObligations, companyIdPayload, inclusionObligations);
+ }
+ //stage 6
+ runStage(TaskType.FormingRegistersOnOS, formingRegistersOnOS); //returns Collection
+ //stage 7
+ runStage(TaskType.FormingPaymentInstruction, formingPaymentInstruction);
+ //stage 8
+ {
+ UnlockResourcesPayload unlockResourcesPayload = new UnlockResourcesPayload();
+ //todo set arguments
+ runStage(TaskType.UnlockResources, unlockResourcesPayload, unlockResources); //returns Collection
+ }
+ //stage 9
+ {
+ FinishingSessionPayload payload = new FinishingSessionPayload();
+ payload.setSessionId(currSession.getId());
+ runStage(TaskType.FinishingSession, payload, finishingSession);
+ }
+ {
+ EndStageNotificationPayload payload = new EndStageNotificationPayload();
+ payload.setSection(currSession.getSection());
+ runStage(TaskType.EndStageNotification, payload, endStageNotification);
+ }
+ } catch (StageException e) {
+ //already logged
+ }
+ }
+
+ private StageResult runStage(TaskType type, ISessionStage stage) {
+ return runStage(type, null, stage);
+ }
+
+ @SuppressWarnings("unchecked")
+ private StageResult runStage(TaskType type, T payload, ISessionStage stage) {
+ continueRunning(type);
+ log.info("session.id={} step {} started", currSession.getId(), currStage.get());
+ Task t = new Task<>(type, payload);
+ StageResult> stgRes = stage.submit(t);
+ log.info("session.id={} step {} result: {} ",
+ currSession.getId(),
+ currStage.get(),
+ stgRes.success ? "success" : messageResolver.resolve(stgRes.error));
+ if (!stgRes.success) {
+ endSession();
+ throw new StageException();
+ }
+ return (StageResult) stgRes;
+ }
+
+ private boolean startSession() {
+ synchronized (this.currStage) {
+ if (this.currStage.get() != null) {
+ log.info("already running session.id={}", this.currSession.getId());
+ return false;
+ } else {
+ Session newSession = new Session();
+ newSession.setSection(Section.FOND.getKey());
+ newSession.setSessionType(SessionType.IPOB.getKey());
+ newSession.setSessionStatus(SessionStatus.CLRN.getKey());
+ //todo companyId/securityId/userId передается из сообщения очереди
+ sessionImdg.insert(newSession);
+ currSession = newSession;
+ log.info("started new session.id={}", this.currSession.getId());
+ currStage.set(TaskType.StartRevise);
+ return true;
+ }
+ }
+ }
+
+ private void endSession() {
+ synchronized (this.currStage) {
+ this.currStage.set(null);
+ this.currSession = null;
+ }
+ }
+
+ private void continueRunning(TaskType t) {
+ synchronized (this.currStage) {
+ this.currStage.set(t);
+ }
+ }
+
+ private boolean checkStage(TaskType t) {
+ synchronized (this.currStage) {
+ return this.currStage.get().equals(t);
+ }
+ }
+
+ private static class StageException extends RuntimeException {
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java
index e25a6d371..9ec11e771 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/TaskType.java
@@ -15,19 +15,42 @@ public enum TaskType {
* step 2
*/
RequirementsAndObligationsCreate,
+ /**
+ * step 3
+ */
+ ObligationsAdmission,
/**
* step 4
*/
InclusionToPool,
-
+ /**
+ * step 5
+ */
+ InspectionObligations,
/**
* Step 6: forming Registry on Obligations and Settlement requirements
*/
FormingRegistersOnOS,
+ /**
+ * step 7
+ */
+ FormingPaymentInstruction,
/**
* step 8
*/
UnlockResources,
+
+
+
+ /**
+ * Step 10
+ */
+ FinishingSession,
+
+ /**
+ * Step 11
+ */
+ EndStageNotification;
}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java
index 4614f2276..d962c1355 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java
@@ -3,13 +3,14 @@ package ru.spcex.clearing.session.stage.impl;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.config.ConfigurableBeanFactory;
+import org.springframework.context.annotation.Scope;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.misc.Currency;
import ru.clearing.classes.statics.data.registry.Registry;
+import ru.clearing.classes.statics.data.sdf.SDf56;
import ru.clearing.classes.statics.data.statement.Statement;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
-import ru.spcex.clearing.platform.messaging.domain.Consts;
-import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf56And51Request;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult;
@@ -19,19 +20,15 @@ import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.specific.StatementRevisePredicate;
-import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumKey;
-import ru.spcex.platform.utils.time.TimeUtil;
import java.math.BigDecimal;
import java.time.Instant;
-import java.time.LocalDate;
-import java.time.Month;
+import java.time.temporal.ChronoUnit;
import java.util.Collection;
-import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
-
@Service
+@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class BalanceRevise implements ISessionStage {
private final Logger log = LoggerFactory.getLogger(getClass());
//todo remove (set all in single method setImdg(provider -> setImdg1();setIdGenerator();...)
@@ -40,15 +37,18 @@ public class BalanceRevise implements ISessionStage {
private Imdg statementImdg;
private Imdg registryImdg;
private Imdg currencyImdg;
+ private Imdg sDf56Imdg;
private KafkaSender kafkaSender;
@Autowired
- public BalanceRevise(ImdgProvider imdgProvider) {
+ public BalanceRevise(ImdgProvider imdgProvider, KafkaSender kafkaSender) {
this.imdgProvider = imdgProvider;
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.currencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Currency, Currency.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ this.sDf56Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf56, SDf56.class);
+ this.kafkaSender = kafkaSender;
}
@Override
@@ -61,28 +61,29 @@ public class BalanceRevise implements ISessionStage {
cashFlow(); //((SdfClearingRequest) task.getData()).getGroupId() if needed
return revise();
}
+ default -> throw new IllegalStateException("unknown task " + task.getTaskType());
}
- return null;
}
private StageResult> sendSdfs() {
Collection currencies = currencyImdg.projectSingleAttribute("id");
Statement statement = statementImdg.aggregateByMax("created", StatementRevisePredicate.get(statementImdg, currencies));
-
- Instant now = Instant.now();
- Sdf56And51Request sdf56Request = new Sdf56And51Request();
- sdf56Request.setDf56Number(idGenerator.nextId()); //fixme day scope id generator
- sdf56Request.setsDateTime(statement != null ?
- statement.getCreated() : TimeUtil.localDateToInstant(LocalDate.of(2023, Month.JANUARY, 1))); //fixme default sDate
- sdf56Request.seteDateTime(now);
- sdf56Request.setDf51Number(idGenerator.nextId()); //fixme day scope id generator
- sdf56Request.setDateTime(now);
- Long msgKey = kafkaSender.sendRequestToQueue(Consts.SDF56_PROCESS, sdf56Request);
- if (msgKey == null) {
- log.error("failed to put SDF56 request to kafka queue");
- return new StageResult<>(new EnumMessage(SessionGeneralError), false);
- }
+ newSDf56(statement);
+ //fixme теперь не отправляем команду в модуль dbf-export, он получит ее из другого места
+// Instant now = Instant.now();
+// Sdf56And51Request sdf56Request = new Sdf56And51Request();
+// sdf56Request.setDf56Number(idGenerator.nextId()); //fixme day scope id generator
+// sdf56Request.setsDateTime(statement != null ?
+// statement.getCreated() : TimeUtil.localDateToInstant(LocalDate.of(2023, Month.JANUARY, 1))); //fixme default sDate
+// sdf56Request.seteDateTime(now);
+// sdf56Request.setDf51Number(idGenerator.nextId()); //fixme day scope id generator
+// sdf56Request.setDateTime(now);
+// Long msgKey = kafkaSender.sendRequestToQueue(Consts.REVISE_PROCESS, sdf56Request);
+// if (msgKey == null) {
+// log.error("failed to put SDF56 request to kafka queue");
+// return new StageResult<>(new EnumMessage(SessionGeneralError), false);
+// }
return new StageResult<>(null, true);
}
@@ -132,7 +133,9 @@ public class BalanceRevise implements ISessionStage {
rgsAMT.setBalance(safeBD(rgsAMT.getBalance()).add(safeBD(stmt.getAmount())));
rgsAMT.setCredit(safeBD(rgsAMT.getCredit()).add(safeBD(stmt.getAmount())));
}
- default -> {throw new IllegalStateException("null direction");}
+ default -> {
+ throw new IllegalStateException("null direction");
+ }
}
rgsAMF.setBalance(safeBD(rgsAMT.getBalance()).subtract(safeBD(rgsAMB.getBalance())));
registryImdg.update(rgsAMT);
@@ -170,4 +173,21 @@ public class BalanceRevise implements ISessionStage {
private BigDecimal safeBD(BigDecimal value) {
return value != null ? value : BigDecimal.ZERO;
}
+
+ private void newSDf56(Statement statement) {
+ log.debug("GALB request received; creating sdf56");
+ SDf56 sDf56 = new SDf56();
+ sDf56.setNumber(idGenerator.nextId().toString());
+ Instant now = Instant.now();
+ String startTime = String.valueOf(statement != null ?
+ statement.getCreated().toEpochMilli() : now.minus(1, ChronoUnit.DAYS).toEpochMilli());
+ sDf56.setStart_datetime(startTime);
+ sDf56.setEnd_datetime(String.valueOf(now.toEpochMilli()));
+ sDf56.setAccount("ТБС");
+ sDf56.setDeal("КОДУ");
+ sDf56.setGenerationTime(now);
+ sDf56.setGenerationId(idGenerator.nextId());
+ sDf56Imdg.insert(sDf56);
+ log.debug("successfully processed, new id {}", sDf56.getId());
+ }
}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java
index ab25bb192..d0613654c 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/DealsPrepare.java
@@ -4,6 +4,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.execution.ExecutionCommon;
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
import ru.clearing.classes.statics.data.execution.ExecutionFond;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
@@ -11,7 +12,6 @@ import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task;
import ru.spcex.clearing.session.stage.task.DealsPreparePayload;
-import ru.spcex.platform.classes.base.interfaces.IExecution;
import ru.spcex.platform.classes.base.interfaces.WithExchangeExecutionId;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@@ -74,11 +74,11 @@ public class DealsPrepare implements ISessionStage {
ImdgPredicate excFondPrct = prdComposer.apply(execFondPredicates, executionFondImdg);
Collection excDpsts = executionDepositImdg.getCollectionObjectsByPredicate(excDepPrct);
Collection excFonds = executionFondImdg.getCollectionObjectsByPredicate(excFondPrct);
- List excs = Stream.concat(excDpsts.stream().map(execToInterface()),
+ List excs = Stream.concat(excDpsts.stream().map(execToInterface()),
excFonds.stream().map(execToInterface()))
.sorted(Comparator.comparing(WithExchangeExecutionId::getExchangeExecutionId))
.toList();
- for (IExecution exc : excs) {
+ for (ExecutionCommon exc : excs) {
exc.setSessionId(sessionId);
if (exc instanceof ExecutionDeposit) {
executionDepositImdg.update((ExecutionDeposit) exc);
@@ -86,7 +86,7 @@ public class DealsPrepare implements ISessionStage {
executionFondImdg.update((ExecutionFond) exc);
}
}
- StageResult> res = new StageResult<>(null, true);
+ StageResult> res = new StageResult<>(null, true);
res.setStageResult(excs);
return res;
}
@@ -103,7 +103,7 @@ public class DealsPrepare implements ISessionStage {
- private static Function execToInterface() {
+ private static Function execToInterface() {
return (e) -> e;
}
}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java
new file mode 100644
index 000000000..ad6481639
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/EndStageNotification.java
@@ -0,0 +1,112 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.registry.Registry;
+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.session.stage.ISessionStage;
+import ru.spcex.clearing.session.stage.StageResult;
+import ru.spcex.clearing.session.stage.Task;
+import ru.spcex.clearing.session.stage.task.EndStageNotificationPayload;
+import ru.spcex.platform.enumeration.*;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.utils.enumeration.EnumMessage;
+import ru.spcex.platform.utils.enumeration.IMessageResolver;
+
+import java.util.Collection;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
+
+@Service
+public class EndStageNotification implements ISessionStage {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+
+ private final Imdg registryImdg;
+ private final KafkaSender kafkaSender;
+ private final IMessageResolver msgResolver;
+
+ @Autowired
+ public EndStageNotification(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver) {
+ this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ this.kafkaSender = kafkaSender;
+ this.msgResolver = msgResolver;
+ }
+
+ @Override
+ public StageResult> submit(Task> task) {
+ EndStageNotificationPayload payload = (EndStageNotificationPayload) task.getData();
+ switch (task.getTaskType()) {
+ case EndStageNotification -> {
+ return endStageNotification(payload.getSection());
+ }
+ default -> {
+ throw new IllegalStateException("Unknown task type: " + task.getTaskType());
+ }
+ }
+ }
+
+ protected Collection selectRegistry() {
+ String registrySQL = "registryStatus='" + RegistryStatus.OK.getKey() + "'";
+ Collection result = registryImdg.getCollectionObjectsBySQL(registrySQL);
+ log.trace("Selected {} registry's by sql: {}", result.size(), registrySQL);
+ return result;
+ }
+
+ protected StageResult> endStageNotification(String section) {
+ Collection forRegistries = selectRegistry();
+ Collection groups = forRegistries.stream()
+ .map(Registry::getGroupId)
+ .filter(Objects::nonNull)
+ .distinct().collect(Collectors.toList());
+ log.debug("Sending notifications for {} groups ({} registers) on section {}",
+ groups, forRegistries.size(), section);
+ for (Long groupId : groups) {
+ log.trace("For registry group {} send notification",
+ groupId);
+ StageResult sResult;
+ if (Section.FOND.equalsByKey(section)) {
+ sResult = notificationDF14(groupId);
+ } else if (Section.MKR.equalsByKey(section)) {
+ sResult = notificationDF05(groupId);
+ } else {
+ throw new IllegalArgumentException("Unsupported section " + section);
+ }
+ if (sResult.getError() != null) {
+ log.warn("When sending groupId={} has error: {}", groupId, msgResolver.resolve(sResult.getError()));
+ }
+ }
+ StageResult> res = new StageResult<>(null, true);
+ return res;
+ }
+
+ protected StageResult notificationDF14(Long groupId) {
+ SdfClearingRequest sdf14Request = new SdfClearingRequest();
+ sdf14Request.setGroupId(groupId);
+ Long msgKey = kafkaSender.sendRequestToQueue(Consts.SDF14_PROCESS, sdf14Request);
+ if (msgKey == null) {
+ log.error("failed to put SDF14 request to kafka queue");
+ return new StageResult<>(new EnumMessage(SessionGeneralError), false);
+ }
+ return new StageResult<>(null, true);
+ }
+
+ protected StageResult notificationDF05(Long groupId) {
+ SdfClearingRequest sdf05Request = new SdfClearingRequest();
+ sdf05Request.setGroupId(groupId);
+ Long msgKey = kafkaSender.sendRequestToQueue(Consts.SDF05_PROCESS, sdf05Request);
+ if (msgKey == null) {
+ log.error("failed to put SDF05 request to kafka queue");
+ return new StageResult<>(new EnumMessage(SessionGeneralError), false);
+ }
+ return new StageResult<>(null, true);
+ }
+
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java
new file mode 100644
index 000000000..051a05eca
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FinishingSession.java
@@ -0,0 +1,142 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.misc.Session;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.platform.messaging.domain.Consts;
+import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequestWithRegistryId;
+import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequestWithSessionId;
+import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
+import ru.spcex.clearing.session.stage.ISessionStage;
+import ru.spcex.clearing.session.stage.StageResult;
+import ru.spcex.clearing.session.stage.Task;
+import ru.spcex.clearing.session.stage.task.FinishingSessionPayload;
+import ru.spcex.platform.enumeration.RegistryDesignation;
+import ru.spcex.platform.enumeration.RegistryTradingParams;
+import ru.spcex.platform.enumeration.SessionStatus;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
+import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
+import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
+import ru.spcex.platform.utils.enumeration.EnumMessage;
+import ru.spcex.platform.utils.enumeration.IMessageResolver;
+
+import java.time.Instant;
+import java.util.Collection;
+
+import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
+
+@Service
+public class FinishingSession implements ISessionStage {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+
+ private final Imdg registryImdg;
+ private final Imdg sessionImdg;
+ private final KafkaSender kafkaSender;
+ private final IMessageResolver msgResolver;
+
+ @Autowired
+ public FinishingSession(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver) {
+ this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
+ this.kafkaSender = kafkaSender;
+ this.msgResolver = msgResolver;
+ }
+
+ @Override
+ public StageResult> submit(Task> task) {
+ FinishingSessionPayload payload = (FinishingSessionPayload) task.getData();
+ switch (task.getTaskType()) {
+ case FinishingSession -> {
+ return finishingSession(payload.getSessionId());
+ }
+ default -> {
+ throw new IllegalStateException("Unknown task type: " + task.getTaskType());
+ }
+ }
+ }
+
+ protected Collection selectRegistry(Long sessionId) {
+ RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(
+ new RegistryTradingParams(RegistryDesignation.A, null, null, null)
+ );
+ ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
+ ImdgPredicate condition = registryCodeSqlBuilder.buildPredicate(pb);
+ condition = pb.and(condition, pb.equals("sessionId", sessionId));
+ Collection result = registryImdg.getCollectionObjectsByPredicate(condition);
+ log.trace("Selected {} registry's by sql: {}", result.size(), condition);
+ return result;
+ }
+
+ protected StageResult> finishingSession(Long sessionId) {
+ // 1. Изменить статус
+ Session theSession = sessionImdg.getSingleObjectByID(sessionId);
+ if (theSession == null) {
+ log.warn("Session {} not found", sessionId);
+ } else {
+ log.info("Finish status for session {}", sessionId);
+ theSession.setUpdated(Instant.now());
+ theSession.setSessionStatus(SessionStatus.CLOS.getKey());
+ sessionImdg.update(theSession);
+ log.trace("Session {} was updated", theSession.getId());
+ }
+
+ // 2. Отправка сообщений
+ Collection forRegistries = selectRegistry(sessionId);
+ log.debug("Found {} registries for sessionId={}", forRegistries.size(), sessionId);
+
+ StageResult sResult = toReportSession(sessionId);
+ if (sResult.getError() != null) {
+ log.error("When sending sessionId={} has error: {}", sessionId, msgResolver.resolve(sResult.getError()));
+ return sResult;
+ }
+ for (Registry registryA : forRegistries) {
+ if (!RegistryDesignation.A.equalsByKey(registryA.getRegistryUnit())) {
+ log.warn("For registry {}.RegistryDesignation is not A.", registryA.getId());
+ continue;
+ }
+ sResult = toReportMoney(sessionId, registryA);
+ if (sResult.getError() != null) {
+ log.warn("When sending sessionId={}, registry.id={} has error: {}", sessionId, registryA.getId(), msgResolver.resolve(sResult.getError()));
+ return sResult;
+ }
+ }
+ StageResult> res = new StageResult<>(null, true);
+ return res;
+ }
+
+ /**
+ * формирование операционного отчета об обязательствах
+ **/
+ protected StageResult toReportSession(Long sessionId) {
+ ReportRequestWithSessionId reportSessionRequest = new ReportRequestWithSessionId();
+ reportSessionRequest.setSessionId(sessionId);
+ Long msgKey = kafkaSender.sendRequestToQueue(Consts.CREATE_REPORT_FOR_SESSION_ID, reportSessionRequest);
+ if (msgKey == null) {
+ log.error("failed to put Report Session request to kafka queue");
+ return new StageResult<>(new EnumMessage(SessionGeneralError), false);
+ }
+ return new StageResult<>(null, true);
+ }
+
+ /**
+ * формирование операционного отчета о денежных средствах
+ **/
+ protected StageResult toReportMoney(Long sessionId, Registry registry) {
+ ReportRequestWithRegistryId reportMoneyRequest = new ReportRequestWithRegistryId();
+ reportMoneyRequest.setSessionId(sessionId);
+ reportMoneyRequest.setRegistryId(registry.getId());
+ Long msgKey = kafkaSender.sendRequestToQueue(Consts.CREATE_REPORT_FOR_REGISTRY, reportMoneyRequest);
+ if (msgKey == null) {
+ log.error("failed to put Report Registry request to kafka queue");
+ return new StageResult<>(new EnumMessage(SessionGeneralError), false);
+ }
+ return new StageResult<>(null, true);
+ }
+
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstruction.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstruction.java
new file mode 100644
index 000000000..3c99589e9
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingPaymentInstruction.java
@@ -0,0 +1,232 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.payment.PaymentInstruction;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
+import ru.spcex.clearing.session.stage.ISessionStage;
+import ru.spcex.clearing.session.stage.StageResult;
+import ru.spcex.clearing.session.stage.Task;
+import ru.spcex.platform.enumeration.*;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgId;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
+import ru.spcex.platform.utils.enumeration.IEnumKey;
+import ru.spcex.platform.utils.enumeration.IMessageResolver;
+import ru.spcex.platform.utils.enumeration.SimpleMessageResolver;
+
+import java.math.BigDecimal;
+import java.util.Collection;
+import java.util.List;
+
+@Service
+public class FormingPaymentInstruction implements ISessionStage {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+ //todo remove (set all in single method setImdg(provider -> setImdg1();setIdGenerator();...)
+ private ImdgProvider imdgProvider;
+ private ImdgId idGenerator;
+ private Imdg registryImdg;
+ private Imdg paymentInstructionImdg;
+ private KafkaSender kafkaSender;
+ private final IMessageResolver msgResolver = new SimpleMessageResolver();
+
+ @Autowired
+ public FormingPaymentInstruction(ImdgProvider imdgProvider) {
+ this.imdgProvider = imdgProvider;
+ this.idGenerator = imdgProvider.getImdgIdGenerator();
+ this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ this.paymentInstructionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PaymentInstruction, PaymentInstruction.class);
+ }
+
+ @Override
+ public StageResult submit(Task> task) {
+ switch (task.getTaskType()) {
+ case FormingPaymentInstruction -> {
+ return formingPaymentInstruction();
+ }
+ default -> {
+ throw new IllegalStateException("Unknown task type: " + task.getTaskType());
+ }
+ }
+ }
+
+
+ private StageResult formingPaymentInstruction() {
+ RegistryTradingParams registryTradingParamsL = new RegistryTradingParams(RegistryDesignation.L,
+ RegistryInstrumentType.S, null, RegistryUnit.T);
+ RegistryTradingParams registryTradingParamsC = new RegistryTradingParams(RegistryDesignation.C,
+ RegistryInstrumentType.M, null, null);
+ RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParamsL, registryTradingParamsC);
+ String registryCodeCondition = registryCodeSqlBuilder.build();
+
+ Collection obligations = registryImdg.getCollectionObjectsBySQL(registryCodeCondition);
+
+ List obligationsByMoney = obligations.stream().filter(registry ->
+ new RegistryTradingParams(RegistryDesignation.L, RegistryInstrumentType.M, null, RegistryUnit.T).equalByRegistry(
+ IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, registry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, registry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, registry.getRegistryUnit())))
+ .toList();
+
+ for (Registry obligationByMoney : obligationsByMoney) {
+ String sql = String.format("tradingClearingRegistryId = %s and companyId = %s",
+ obligationByMoney.getTradingClearingRegistryId(), obligationByMoney.getCompanyId());
+ Collection relatedRegistries = registryImdg.getCollectionObjectsBySQL(sql);
+
+ RegistryTradingParams registryTradingParamsF = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.M, null, RegistryUnit.F);
+ RegistryTradingParams registryTradingParamsT = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.M, null, RegistryUnit.T);
+ RegistryTradingParams registryTradingParamsB = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.M, null, RegistryUnit.B);
+
+ for (Registry relatedRegistry : relatedRegistries) {
+ if (registryTradingParamsF.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setBalance(relatedRegistry.getBalance().subtract(obligationByMoney.getBalance()));
+ } else if (registryTradingParamsT.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setSettledDebit(relatedRegistry.getSettledDebit().add(obligationByMoney.getBalance()));
+ } else if (registryTradingParamsB.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setBalance(relatedRegistry.getBalance().add(obligationByMoney.getBalance()));
+ }
+ registryImdg.update(relatedRegistry);
+ }
+ }
+
+ List requirementsByIssue = obligations.stream().filter(registry ->
+ new RegistryTradingParams(RegistryDesignation.C, RegistryInstrumentType.S, null, RegistryUnit.T).equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, registry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, registry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, registry.getRegistryUnit())))
+ .toList();
+
+ for (Registry requirementByIssue : requirementsByIssue) {
+ String sql = String.format("tradingClearingRegistryId = %s and companyId = %s",
+ requirementByIssue.getTradingClearingRegistryId(), requirementByIssue.getCompanyId());
+ Collection relatedRegistries = registryImdg.getCollectionObjectsBySQL(sql);
+
+ RegistryTradingParams registryTradingParamsA = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.S, null, RegistryUnit.T);
+
+ for (Registry relatedRegistry : relatedRegistries) {
+ if (registryTradingParamsA.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setSettledDebit(relatedRegistry.getSettledDebit().add(requirementByIssue.getBalance()));
+ }
+ registryImdg.update(relatedRegistry);
+ }
+ }
+
+ List requirementsByMoney = obligations.stream().filter(registry ->
+ new RegistryTradingParams(RegistryDesignation.C, RegistryInstrumentType.M, null, RegistryUnit.T).equalByRegistry(
+ IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, registry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, registry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, registry.getRegistryUnit())))
+ .toList();
+
+ for (Registry requirementByMoney : requirementsByMoney) {
+ String sql = String.format("tradingClearingRegistryId = %s and companyId = %s",
+ requirementByMoney.getTradingClearingRegistryId(), requirementByMoney.getCompanyId());
+ Collection relatedRegistries = registryImdg.getCollectionObjectsBySQL(sql);
+
+ RegistryTradingParams registryTradingParamsA = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.M, null, RegistryUnit.T);
+
+ for (Registry relatedRegistry : relatedRegistries) {
+ if (registryTradingParamsA.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setSettledDebit(relatedRegistry.getSettledDebit().add(requirementByMoney.getBalance()));
+ }
+ registryImdg.update(relatedRegistry);
+ }
+ }
+
+ List obligationsByIssue = obligations.stream().filter(registry ->
+ new RegistryTradingParams(RegistryDesignation.L, RegistryInstrumentType.S, null, RegistryUnit.T).equalByRegistry(
+ IEnumKey.getEnumByKey(RegistryDesignation.class, registry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, registry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, registry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, registry.getRegistryUnit())))
+ .toList();
+
+ for (Registry obligationByIssue : obligationsByIssue) {
+ String sql = String.format("tradingClearingRegistryId = %s and companyId = %s",
+ obligationByIssue.getTradingClearingRegistryId(), obligationByIssue.getCompanyId());
+ Collection relatedRegistries = registryImdg.getCollectionObjectsBySQL(sql);
+
+ RegistryTradingParams registryTradingParamsF = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.S, null, RegistryUnit.F);
+ RegistryTradingParams registryTradingParamsT = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.S, null, RegistryUnit.T);
+ RegistryTradingParams registryTradingParamsB = new RegistryTradingParams(RegistryDesignation.A,
+ RegistryInstrumentType.S, null, RegistryUnit.B);
+
+ for (Registry relatedRegistry : relatedRegistries) {
+ if (registryTradingParamsF.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setBalance(relatedRegistry.getBalance().subtract(obligationByIssue.getBalance()));
+ } else if (registryTradingParamsT.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setSettledDebit(relatedRegistry.getSettledDebit().add(obligationByIssue.getBalance()));
+ } else if (registryTradingParamsB.equalByRegistry(IEnumKey.getEnumByKey(RegistryDesignation.class, relatedRegistry.getRegistryDesignation()),
+ IEnumKey.getEnumByKey(RegistryInstrumentType.class, relatedRegistry.getRegistryInstrumentType()),
+ IEnumKey.getEnumByKey(RegistryCapacity.class, relatedRegistry.getRegistryCapacity()),
+ IEnumKey.getEnumByKey(RegistryUnit.class, relatedRegistry.getRegistryUnit()))) {
+ relatedRegistry.setBalance(relatedRegistry.getBalance().add(obligationByIssue.getBalance()));
+ }
+ registryImdg.update(relatedRegistry);
+ }
+ }
+
+ for (Registry obligation : obligations) {
+ RegistryTradingParams registryTradingParamsLT = new RegistryTradingParams(RegistryDesignation.L,
+ null, null, RegistryUnit.T);
+ registryCodeCondition = RegistryCodeSqlBuilder.getInstance(registryTradingParamsLT).build();
+ String sqlCondition = String.format("(%s) and securityId = %s and tradingClearingRegistryId = %s and companyId = %s and counterPartyId = %s",
+ registryCodeCondition, obligation.getSecurityId(), obligation.getTradingClearingRegistryId(), obligation.getCompanyId(), obligation.getCounterPartyId());
+ Collection registries = registryImdg.getCollectionObjectsBySQL(sqlCondition);
+
+ BigDecimal sumBalance = registries.stream().map(Registry::getBalance).reduce(BigDecimal.ZERO, BigDecimal::add);
+ PaymentInstruction paymentInstruction = createPaymentInstruction(obligation, sumBalance);
+ paymentInstructionImdg.insert(paymentInstruction);
+ obligation.setPaymentId(paymentInstruction.getId());
+ registryImdg.update(obligation);
+ }
+
+
+ //todo send message to queue for forming sDf03 and sDf12?
+
+ return new StageResult(null, true);
+ }
+
+ private PaymentInstruction createPaymentInstruction(Registry registry, BigDecimal balance) {
+ PaymentInstruction paymentInstruction = new PaymentInstruction();
+ paymentInstruction.setSenderId(registry.getCompanyId());
+ paymentInstruction.setAddresseeId(registry.getCounterPartyId());
+ //todo add builder for paymentInstruction
+ return paymentInstruction;
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingRegistersOnOS.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingRegistersOnOS.java
index 5204abe62..54b16f861 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingRegistersOnOS.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/FormingRegistersOnOS.java
@@ -9,7 +9,7 @@ import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task;
-import ru.spcex.clearing.session.stage.task.RegistryOnObligationsAndSettlementRequirementsPayload;
+import ru.spcex.clearing.session.stage.task.FormingRegistersOnOSPayload;
import ru.spcex.clearing.session.stage.util.RegistryUtil;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
@@ -39,10 +39,10 @@ public class FormingRegistersOnOS implements ISessionStage {
@Override
public StageResult> submit(Task> task) {
- RegistryOnObligationsAndSettlementRequirementsPayload payload = (RegistryOnObligationsAndSettlementRequirementsPayload) task.getData();
+ FormingRegistersOnOSPayload payload = (FormingRegistersOnOSPayload) task.getData();
switch (task.getTaskType()) {
case FormingRegistersOnOS -> {
- return createRegistryOnObligationsAndSettlementRequirements(payload.getSessionId());
+ return createRegistryOnObligationsAndSettlementRequirements();
}
default -> {
throw new IllegalStateException("Unknown task type: " + task.getTaskType());
@@ -53,12 +53,10 @@ public class FormingRegistersOnOS implements ISessionStage {
/**
* Select Registry by: registryCode = [O/T][S/M][*][T] & registryStatus=OK
*
- * @param sessionId
* @return
*/
- protected Collection selectRegistry(Long sessionId) {
- String registrySQL = "sessionId = " + sessionId;
- registrySQL += " and (registryDesignation=" + RegistryDesignation.O.getKey() + " or registryDesignation=" + RegistryDesignation.T.getKey() + ")";
+ protected Collection selectRegistry() {
+ String registrySQL = "(registryDesignation=" + RegistryDesignation.O.getKey() + " or registryDesignation=" + RegistryDesignation.T.getKey() + ")";
registrySQL += " and (registryInstrumentType=" + RegistryInstrumentType.S.getKey() + " or registryInstrumentType=" + RegistryInstrumentType.M.getKey() + ")";
registrySQL += " and (registryUnit=" + RegistryUnit.T.getKey() + ")";
registrySQL += " and registryStatus=" + RegistryStatus.OK.getKey() + ")";
@@ -67,8 +65,8 @@ public class FormingRegistersOnOS implements ISessionStage {
return result;
}
- protected StageResult> createRegistryOnObligationsAndSettlementRequirements(Long sessionId) {
- Collection forRegistries = selectRegistry(sessionId);
+ protected StageResult> createRegistryOnObligationsAndSettlementRequirements() {
+ Collection forRegistries = selectRegistry();
ArrayList newRegistries = new ArrayList<>();
//todo oreder by обрабатываться группами по полю registry.groupId
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java
index eea004140..b3aae427f 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InclusionObligations.java
@@ -67,7 +67,6 @@ public class InclusionObligations implements ISessionStage {
sessionType);
Collection obligations = registryImdg.getCollectionObjectsBySQL(sqlCondition);
Map> registryByGroupId = obligations.stream().
- filter(registry -> registry.getCompanyId().equals(counterPartyId)).
collect(Collectors.groupingBy(Registry::getGroupId));
for (Map.Entry> entrySet : registryByGroupId.entrySet()) {
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java
index 47b770347..5f60447d6 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/InspectionObligations.java
@@ -51,7 +51,7 @@ public class InspectionObligations implements ISessionStage {
public StageResult submit(Task> task) {
InspectionPoolPayload payload = (InspectionPoolPayload) task.getData();
switch (task.getTaskType()) {
- case InclusionToPool -> {
+ case InspectionObligations -> {
return inspectionPool(payload.getProcessedCompanyId());
}
default -> {
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java
new file mode 100644
index 000000000..f62cdd211
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/ObligationAdmission.java
@@ -0,0 +1,79 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.session.stage.ISessionStage;
+import ru.spcex.clearing.session.stage.StageResult;
+import ru.spcex.clearing.session.stage.Task;
+import ru.spcex.clearing.session.stage.TaskType;
+import ru.spcex.platform.enumeration.RegistryStatus;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.utils.enumeration.EnumMessage;
+import ru.spcex.platform.utils.enumeration.IMessageResolver;
+import ru.spcex.platform.utils.validation.IValidator;
+
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+@Service
+public class ObligationAdmission implements ISessionStage {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+ private final Imdg registryImdg;
+ private final ImdgProvider imdgProvider;
+ private final Function validatorFactory;
+ private final IMessageResolver messageResolver;
+
+ @Autowired
+ public ObligationAdmission(ImdgProvider imdgProvider,
+ @Qualifier("obligationAndRequirementsAdmissionValidator") Function validatorFactory,
+ IMessageResolver messageResolver) {
+ this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ this.imdgProvider = imdgProvider;
+ this.validatorFactory = validatorFactory;
+ this.messageResolver = messageResolver;
+ }
+
+ @SuppressWarnings("unchecked")
+ @Override
+ public StageResult> submit(Task> task) {
+ if (task.getTaskType() == TaskType.ObligationsAdmission) {
+ return obligationAdmission((Long) task.getData());
+ }
+ throw new IllegalStateException("Unknown task type: " + task.getTaskType());
+ }
+
+ private StageResult> obligationAdmission(Long sessionId) {
+ Collection registries = registryImdg.getCollectionObjectsBySQL("sessionId = " + sessionId);
+ Map> byGroups = registries.stream().collect(Collectors.groupingBy(Registry::getGroupId));
+ for (Map.Entry> grpEntry : byGroups.entrySet()) {
+ EnumMessage groupError = null;
+ List rgsGroup = grpEntry.getValue();
+ for (Registry rgs : rgsGroup) {
+ IValidator validator = validatorFactory.apply(rgs);
+ Optional error = validator.tillFirstError();
+ if (error.isPresent()) {
+ groupError = error.get();
+ break;
+ }
+ }
+ if (groupError != null) {
+ log.warn("error {} for registries groupId = {}", messageResolver.resolve(groupError), grpEntry.getKey());
+ for (Registry rgs : rgsGroup) {
+ rgs.setRegistryStatus(RegistryStatus.NACK.getKey());
+ registryImdg.update(rgs);
+ }
+ }
+ }
+ return new StageResult<>(null, true);
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RegistryBuilder.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RegistryBuilder.java
new file mode 100644
index 000000000..29e40efae
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RegistryBuilder.java
@@ -0,0 +1,151 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import ru.clearing.classes.statics.data.account.Account;
+import ru.clearing.classes.statics.data.company.Company;
+import ru.clearing.classes.statics.data.execution.ExecutionCommon;
+import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
+import ru.clearing.classes.statics.data.execution.ExecutionFond;
+import ru.clearing.classes.statics.data.misc.Market;
+import ru.clearing.classes.statics.data.misc.Session;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.session.stage.util.RegistryUtil;
+import ru.spcex.platform.classes.base.interfaces.ExecutionType;
+import ru.spcex.platform.enumeration.*;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+
+import java.time.LocalDate;
+import java.time.format.DateTimeFormatter;
+
+public class RegistryBuilder {
+ private Imdg companyImdg;
+ private Imdg tradingClearingRegistryImdg;
+ private Imdg accountImdg;
+ private Imdg marketImdg;
+ private Imdg sessionImdg;
+ private ExecutionCommon exec;
+ private RegistryDesignation regDsgn;
+
+ private RegistryBuilder() {
+ }
+
+ public static RegistryBuilder builder() {
+ return new RegistryBuilder();
+ }
+
+ public RegistryBuilder imdg(ImdgProvider imdgProvider) {
+ this.companyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class);
+ this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class);
+ this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class);
+ this.marketImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Market, Market.class);
+ this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
+ return this;
+ }
+
+ public RegistryBuilder exec(ExecutionCommon exec) {
+ this.exec = exec;
+ return this;
+ }
+
+ public RegistryBuilder registryDesignation(RegistryDesignation registryDesignation) {
+ this.regDsgn = registryDesignation;
+ return this;
+ }
+
+ public Registry build() {
+ ISide side = getSide(exec);
+ Registry reg = new Registry();
+ reg.setCompanyId(exec.getCompanyId());
+ Company company = searchCompany();
+ reg.setTradingCode(company.getTradingCode());
+ reg.setClearingCode(company.getClearingCode());
+ reg.setShortName(company.getShortName());
+ reg.setFullName(company.getFullName());
+ TradingClearingRegistry tcr = searchTradingClearingRegistry();
+ if ((regDsgn.equals(RegistryDesignation.O) && side.isBuy()) || (regDsgn.equals(RegistryDesignation.T) && side.isSell())) {
+ reg.setAccountId(tcr.getMoneyAccountId());
+ reg.setRegistryInstrumentType(RegistryInstrumentType.M.getKey());
+ reg.setBalanceDimension(BalanceDimension.MONY.getKey()); //fixme add second leg code branch
+ } else if ((regDsgn.equals(RegistryDesignation.T) && side.isBuy()) || (regDsgn.equals(RegistryDesignation.O) && side.isSell())) {
+ reg.setAccountId(tcr.getDepoAccountId());
+ reg.setRegistryInstrumentType(RegistryInstrumentType.S.getKey());
+ reg.setBalanceDimension(BalanceDimension.PICS.getKey()); //fixme add second leg code branch
+ }
+ Account account = accountImdg.getSingleObjectByID(reg.getAccountId());
+ reg.setAccountType(account.getAccountType());
+ reg.setAccount(account.getAccount());
+ reg.setRegistryDesignation(regDsgn.getKey());
+ reg.setRegistryCapacity(account.getAccountType());
+ reg.setRegistryUnit(RegistryUnit.T.getKey());
+ reg.setRegistryCode(RegistryUtil.clearingCode(reg));
+ reg.setTradingClearingRegistryId(exec.getTradingClearingRegistryId());
+ reg.setTradingClearingRegistry(tcr.getCode());
+ reg.setRegistryStatus(RegistryStatus.PROC.getKey());
+ reg.setSecurityId(exec.getSecurityId());
+ reg.setSecuritySymbol(exec.getSecuritySymbol());
+ switch (exec.type()) {
+ case ExecutionDeposit -> {
+ ExecutionDeposit execDep = (ExecutionDeposit) this.exec;
+ reg.setBalance((execDep).getFirstLegAmount());
+ if (regDsgn.equals(RegistryDesignation.T)) {
+ reg.setSettledCredit((execDep).getFirstLegAmount());
+ } else if (regDsgn.equals(RegistryDesignation.O)) {
+ reg.setSettledDebit((execDep).getFirstLegAmount());
+ }
+ reg.setSettlementDate(execDep.getFirstLegSettlementDate()); //fixme add second leg code branch
+ reg.setSettlementCode(execDep.getFirstLegSettlementCode()); //fixme add second leg code branch
+ reg.setRefundDate(execDep.getSecondLegSettlementDate()); //fixme add second leg code branch??
+ reg.setValueDate(execDep.getFirstLegSettlementDate());
+ reg.setContract(execDep.getContract());
+ }
+ case ExecutionFond -> {
+ ExecutionFond execFond = (ExecutionFond) this.exec;
+ reg.setBalance(execFond.getSettlementAmount());
+ reg.setSettlementDate(execFond.getSettlementDate());
+ reg.setSettlementCode(execFond.getSettlementCode());
+ }
+ }
+ reg.setTradingDate(exec.getTradingDate());
+ reg.setClearingDate(LocalDate.now());
+ reg.setPrice(exec.getPrice());
+ reg.setCounterPartyId(exec.getCounterPartyId());
+ reg.setGroupId(groupId());
+ reg.setSessionId(exec.getSessionId());
+ reg.setSessionType(sessionType());
+ return reg;
+ }
+
+ private static DateTimeFormatter yyyyMMdd = DateTimeFormatter.ofPattern("yyyyMMdd");
+
+ private Long groupId() {
+ LocalDate now = LocalDate.now();
+ Market market = marketImdg.getSingleObjectBySQL("code = '" + exec.getMarket() + "'");
+ return Long.valueOf(now.format(yyyyMMdd) + exec.getExchangeExecutionId() + market.getId());
+ }
+
+ //можно передать из стейджа
+ private String sessionType() {
+ Session session = sessionImdg.getSingleObjectByID(exec.getSessionId());
+ return session.getSessionType();
+ }
+
+ private Company searchCompany() {
+ return companyImdg.getSingleObjectBySQL("id = " + exec.getCompanyId());
+ }
+
+ private TradingClearingRegistry searchTradingClearingRegistry() {
+ return tradingClearingRegistryImdg.getSingleObjectBySQL("tradingClearingRegistryId = " + exec.getTradingClearingRegistryId());
+ }
+
+ private static ISide getSide(ExecutionCommon exec) {
+ if (exec.type().equals(ExecutionType.ExecutionDeposit)) {
+ return ISide.parse(MoneyFlowSide.class, exec.getSide());
+ } else if (exec.type().equals(ExecutionType.ExecutionFond)) {
+ return ISide.parse(Side.class, exec.getSide());
+ } else {
+ throw new RuntimeException("Unknown Execution type: " + exec.type());
+ }
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RegistryUpdater.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RegistryUpdater.java
new file mode 100644
index 000000000..a5c08b03e
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RegistryUpdater.java
@@ -0,0 +1,55 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import ru.clearing.classes.statics.data.execution.ExecutionCommon;
+import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
+import ru.clearing.classes.statics.data.execution.ExecutionFond;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.spcex.platform.classes.base.interfaces.ExecutionType;
+import ru.spcex.platform.enumeration.RegistryDesignation;
+
+import java.math.BigDecimal;
+import java.time.Instant;
+
+public class RegistryUpdater {
+ public static RegistryUpdater updater() {
+ return new RegistryUpdater();
+ }
+
+ private RegistryUpdater() {
+ }
+
+ private ExecutionCommon exec;
+ private Registry rgs;
+
+ public RegistryUpdater registry(Registry rgs) {
+ this.rgs = rgs;
+ return this;
+ }
+
+ public RegistryUpdater exec(ExecutionCommon exec) {
+ this.exec = exec;
+ return this;
+ }
+
+ public void update() {
+ if (exec.type().equals(ExecutionType.ExecutionDeposit)) {
+ rgs.setBalance(safe(rgs.getBalance()).add(safe(((ExecutionDeposit) exec).getFirstLegAmount())));
+ if (rgs.getRegistryDesignation().equals(RegistryDesignation.T.getKey())) {
+ rgs.setSettledCredit(safe(rgs.getSettledCredit()).add(safe(((ExecutionDeposit) exec).getFirstLegAmount())));
+ }
+ if (rgs.getRegistryDesignation().equals(RegistryDesignation.O.getKey())) {
+ rgs.setSettledDebit(safe(rgs.getSettledDebit()).add(safe(((ExecutionDeposit) exec).getFirstLegAmount())));
+ }
+ rgs.setValueDate(((ExecutionDeposit) exec).getFirstLegSettlementDate());
+
+ } else if (exec.type().equals(ExecutionType.ExecutionFond)) {
+ rgs.setBalance(safe(rgs.getBalance()).add(safe(((ExecutionFond) exec).getSettlementAmount())));
+ }
+ rgs.setUpdated(Instant.now());
+
+ }
+
+ private BigDecimal safe(BigDecimal value) {
+ return value == null ? BigDecimal.ZERO : value;
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RequirementsAndObligationCreation.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RequirementsAndObligationCreation.java
new file mode 100644
index 000000000..150456493
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/RequirementsAndObligationCreation.java
@@ -0,0 +1,223 @@
+package ru.spcex.clearing.session.stage.impl;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import ru.clearing.classes.statics.data.execution.ExecutionCommon;
+import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
+import ru.clearing.classes.statics.data.execution.ExecutionFond;
+import ru.clearing.classes.statics.data.registry.Registry;
+import ru.spcex.clearing.imdg.IMDGDistributedNames;
+import ru.spcex.clearing.session.stage.ISessionStage;
+import ru.spcex.clearing.session.stage.StageResult;
+import ru.spcex.clearing.session.stage.Task;
+import ru.spcex.clearing.session.stage.TaskType;
+import ru.spcex.platform.classes.base.interfaces.ExecutionType;
+import ru.spcex.platform.enumeration.*;
+import ru.spcex.platform.imdg.api.Imdg;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
+import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
+import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
+import ru.spcex.platform.utils.collection.Pair;
+
+import java.time.LocalDate;
+import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.BiConsumer;
+import java.util.function.BiFunction;
+
+@Service
+public class RequirementsAndObligationCreation implements ISessionStage {
+ private final Logger log = LoggerFactory.getLogger(getClass());
+ private final Imdg registryImdg;
+ private final ImdgProvider imdgProvider;
+
+
+ @Autowired
+ public RequirementsAndObligationCreation(ImdgProvider imdgProvider) {
+ this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ this.imdgProvider = imdgProvider;
+ }
+
+ @SuppressWarnings("unchecked")
+ @Override
+ public StageResult> submit(Task> task) {
+ if (task.getTaskType() == TaskType.RequirementsAndObligationsCreate) {
+ return createRegisters((List) task.getData());
+ }
+ throw new IllegalStateException("Unknown task type: " + task.getTaskType());
+ }
+
+ private StageResult> createRegisters(List data) {
+ if (data.size() % 2 != 0) {
+ throw new IllegalStateException("Data size must be even (executions must match)");
+ }
+ //см. описание к #matchExecutions
+ for (int i = 0; i < data.size(); ) {
+ Pair matched = matchExecutions(data, i);
+ if (matched == null) {
+ log.error("couldn't find match execution for id = {}", data.get(i).getId());
+ i += 1;
+ continue;
+ }
+ i += 2;
+ ExecutionCommon partyExec = matched.getFirst();
+ ExecutionCommon counterExec = matched.getSecond();
+ //бывает поиск обычный, бывает по ExecutionDeposit.secondLegSettlementDt
+ AtomicReference>> regSearchMethod = new AtomicReference<>();
+ //обычный поиск
+ regSearchMethod.set(this::findReg);
+ BiConsumer findAndUpdateOrCreate = (exec, regDsgn) ->
+ regSearchMethod.get().apply(exec, regDsgn)
+ .ifPresentOrElse(registry -> {
+ RegistryUpdater.updater()
+ .exec(partyExec)
+ .registry(registry)
+ .update();
+ registryImdg.update(registry);
+ }, () -> {
+ Registry newRegister = RegistryBuilder.builder()
+ .imdg(imdgProvider)
+ .exec(partyExec)
+ .registryDesignation(regDsgn)
+ .build();
+ registryImdg.insert(newRegister);
+ });
+ findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.O);
+ findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.T);
+ findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.O);
+ findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.T);
+ if (partyExec.type().equals(ExecutionType.ExecutionDeposit) && ((ExecutionDeposit) partyExec).getSecondLegSettlementDate() != null) {
+ //поиск по дате расчетов ExecutionDeposit.secondLegSettlementDt
+ regSearchMethod.set(this::findRegBySettlementDt);
+ findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.T);
+ findAndUpdateOrCreate.accept(partyExec, RegistryDesignation.O);
+ findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.T);
+ findAndUpdateOrCreate.accept(counterExec, RegistryDesignation.O);
+ }
+ }
+ return new StageResult<>(null, true);
+ }
+
+
+ private Optional findReg(ExecutionCommon exec, RegistryDesignation des) {
+ ISide side = getSide(exec);
+ RegistryTradingParams p;
+ if (side.isBuy() && des.equals(RegistryDesignation.O)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.O, RegistryInstrumentType.M, null, RegistryUnit.T
+ );
+ } else if (side.isBuy() && des.equals(RegistryDesignation.T)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.T, RegistryInstrumentType.S, null, RegistryUnit.T
+ );
+ } else if (side.isSell() && des.equals(RegistryDesignation.O)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.O, RegistryInstrumentType.S, null, RegistryUnit.T
+ );
+ } else if (side.isSell() && des.equals(RegistryDesignation.T)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.T, RegistryInstrumentType.M, null, RegistryUnit.T
+ );
+ } else {
+ throw new IllegalStateException("cannot construct for " + des + " " + side);
+ }
+ return imdgRegSearch(p, exec.getTradingClearingRegistryId(), exec.getCompanyId(), settlementDt(exec));
+ }
+
+ private Optional findRegBySettlementDt(ExecutionCommon exec, RegistryDesignation des) {
+ if (!exec.type().equals(ExecutionType.ExecutionDeposit) || ((ExecutionDeposit) exec).getSecondLegSettlementDate() == null)
+ throw new IllegalStateException("cannot searchSettlementDt for " + exec.type());
+ ISide side = getSide(exec);
+ RegistryTradingParams p;
+ if (side.isBuy() && des.equals(RegistryDesignation.T)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.T, RegistryInstrumentType.M, RegistryCapacity.A, RegistryUnit.T
+ );
+ } else if (side.isBuy() && des.equals(RegistryDesignation.O)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.O, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.T
+ );
+ } else if (side.isSell() && des.equals(RegistryDesignation.O)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.O, RegistryInstrumentType.M, RegistryCapacity.A, RegistryUnit.T
+ );
+ } else if (side.isSell() && des.equals(RegistryDesignation.T)) {
+ p = new RegistryTradingParams(
+ RegistryDesignation.T, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.T
+ );
+ } else {
+ throw new IllegalStateException("cannot construct by settleDt for " + des + " " + side);
+ }
+ return imdgRegSearch(p, exec.getTradingClearingRegistryId(), exec.getCompanyId(), ((ExecutionDeposit) exec).getSecondLegSettlementDate());
+ }
+
+ public Optional imdgRegSearch(RegistryTradingParams p, Long tcrId, Long companyId, LocalDate settlementDt) {
+ String sql = RegistryCodeSqlBuilder.getInstance(p).build();
+ ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
+ ImdgPredicate rgstrPredicate = pb.and(pb.sql(sql),
+ pb.equals("tradingClearingRegistryId", tcrId),
+ pb.equals("companyId", companyId),
+ pb.equals("settlementDate", settlementDt)
+ );
+ return Optional.ofNullable(registryImdg.getSingleObjectByPredicate(rgstrPredicate));
+ }
+
+ private LocalDate settlementDt(ExecutionCommon exec) {
+ if (exec.type().equals(ExecutionType.ExecutionDeposit)) {
+ return ((ExecutionDeposit) exec).getFirstLegSettlementDate();
+ } else if (exec.type().equals(ExecutionType.ExecutionFond)) {
+ return ((ExecutionFond) exec).getSettlementDate();
+ } else {
+ throw new RuntimeException("Unknown Execution type: " + exec.type());
+ }
+ }
+
+
+ /**
+ * входные данные отсортированы по exchangeExecutionId.
+ * сортировка выполнена на предыдущем шаге. (такое описание шагов в ТЗ)
+ * по идее, Execution с одинаковым exchangeExecutionId должно быть всего два.
+ * если это не так, передалать на коллекцию Long executionId в правильном порядке
+ * и для каждого искать мэтч отдельно в Imdg, с сохранением уже обработанных для избежания дублирования
+ */
+ private Pair matchExecutions(List data, int i) {
+ //if the last execution, then no match
+ if (i >= data.size() - 1) {
+ return null;
+ }
+ ExecutionCommon exec1 = data.get(i);
+ ExecutionCommon exec2 = data.get(i + 1);
+ //Встречная сделка контрагента выбирается из executionDeposit/Fond по условию:
+ //exchangeExecutionId=currentExecutionDeposit/Fond.exchangeExecutionId
+ //и [side=SELL (если currentExecutionDeposit/Fond.side=BUY) или side=BUY (если currentExecutionDeposit/Fond.side=SELL) по справочнику moneyFlowSide или справочнику side в зависимости от секции обрабатываемой сделки]
+ //и companyId=currentExecutionDeposit/Fond.counterPartyId
+ boolean valid = true;
+ if (!Objects.equals(exec1.getExchangeExecutionId(), exec2.getExchangeExecutionId())) {
+ valid = false;
+ } else if (Objects.equals(getSide(exec1), getSide(exec2))) {
+ valid = false;
+ } else if (!Objects.equals(exec1.getCounterPartyId(), exec2.getCounterPartyId())) {
+ valid = false;
+ }
+ if (!valid) {
+ log.error("FATAL skipping ExchangeExecutionId: {}", exec1.getId());
+ }
+ return new Pair<>(exec1, exec2);
+ }
+
+ private static ISide getSide(ExecutionCommon exec) {
+ if (exec.type().equals(ExecutionType.ExecutionDeposit)) {
+ return ISide.parse(MoneyFlowSide.class, exec.getSide());
+ } else if (exec.type().equals(ExecutionType.ExecutionFond)) {
+ return ISide.parse(Side.class, exec.getSide());
+ } else {
+ throw new RuntimeException("Unknown Execution type: " + exec.type());
+ }
+ }
+
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/UnlockResources.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/UnlockResources.java
index c563c84b5..91c0a60f2 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/UnlockResources.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/UnlockResources.java
@@ -15,6 +15,7 @@ import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.imdg.api.ImdgTransaction;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
@@ -26,10 +27,12 @@ import java.util.Collection;
public class UnlockResources implements ISessionStage {
private final Logger log = LoggerFactory.getLogger(getClass());
+ protected final ImdgProvider imdgProvider;
private final Imdg registryImdg;
@Autowired
public UnlockResources(ImdgProvider imdgProvider) {
+ this.imdgProvider = imdgProvider;
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
}
@@ -39,7 +42,7 @@ public class UnlockResources implements ISessionStage {
switch (task.getTaskType()) {
case UnlockResources -> {
UnlockResourcesPayload payload = (UnlockResourcesPayload) task.getData();
- return unlockResources(payload.getSdfMode(), payload.getSessionId(),
+ return unlockResources(payload.getSdfMode(),
payload.getAccount(),
payload.getSecurityId(),
payload.getFullNames(),
@@ -52,26 +55,18 @@ public class UnlockResources implements ISessionStage {
}
}
- /**
- * Select Registry by: account; fullName* / securityId
- */
- protected Collection selectRegistry(Long sessionId, String account, Long securityId, Collection fullNames) {
+ protected Collection selectRegistryForSDF04(String account, Collection fullNames) {
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
- ImdgPredicate queryPart;
- if (securityId != null) {
- queryPart = pb.in("securityId", securityId);
- } else if (fullNames != null && !fullNames.isEmpty()) {
- queryPart = pb.in("fullName", fullNames.toArray(new String[fullNames.size()]));
- } else {
- throw new IllegalArgumentException("Required securityId or fullNames");
- }
ImdgPredicate query = pb.and(
pb.and(
- pb.equals("sessionId", sessionId),
- pb.equals("registryDesignation", RegistryDesignation.A.getKey()) // не все нужны, только с этим кодом отфильтруем.
+ pb.equals("registryDesignation", RegistryDesignation.A.getKey()),
+ pb.equals("registryInstrumentType", RegistryInstrumentType.M.getKey()),
+ // pb.equals("registryCapacity", *),
+ pb.or(pb.equals("registryUnit", RegistryUnit.F.getKey()),
+ pb.equals("registryUnit", RegistryUnit.B.getKey()))
),
pb.equals("account", account),
- queryPart
+ pb.in("fullName", fullNames.toArray(new String[fullNames.size()]))
);
Collection result = registryImdg.getCollectionObjectsByPredicate(query);
@@ -79,50 +74,74 @@ public class UnlockResources implements ISessionStage {
return result;
}
+ /**
+ * Select Registry by: account; fullName* / securityId
+ */
+ protected Collection selectRegistryForSDF12(String account, Long securityId) {
+ ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
+ ImdgPredicate query = pb.and(
+ pb.and(
+ pb.equals("registryDesignation", RegistryDesignation.A.getKey()),
+ pb.equals("registryInstrumentType", RegistryInstrumentType.S.getKey()),
+ // pb.equals("registryCapacity", *),
+ pb.or(pb.equals("registryUnit", RegistryUnit.F.getKey()),
+ pb.equals("registryUnit", RegistryUnit.B.getKey()))
+ ),
+ pb.equals("account", account),
+ pb.in("securityId", securityId)
+ );
+
+ Collection result = registryImdg.getCollectionObjectsByPredicate(query);
+ log.trace("Selected {} registry's by sql: {}", result.size(), query);
+ return result;
+ }
+
+
/**
* @param sdfMode UnlockResourcesPayload.sdfMode
- * @param sessionId
* @param account
* @param securityId
* @param fullNames
* @return
*/
- protected StageResult> unlockResources(String sdfMode, Long sessionId, String account, Long securityId, Collection fullNames, BigDecimal value) {
- Collection forRegistries = selectRegistry(sessionId, account, securityId, fullNames);
+ protected StageResult> unlockResources(String sdfMode, String account, Long securityId, Collection fullNames, BigDecimal value) {
+ Collection forRegistries;
- //todo при перезапуске после незапланированного завершения стадии: надо ли проверять уже созданные регистры и не создавать дубликаты?
if (UnlockResourcesPayload.MODE_SDF04.equals(sdfMode)) {
- forRegistries = forRegistries.stream().filter(
- (Registry r) ->
- RegistryDesignation.A.equalsByKey(r.getRegistryDesignation())
- &&
- RegistryInstrumentType.M.equalsByKey(r.getRegistryInstrumentType())
- &&
- (RegistryUnit.F.equalsByKey(r.getRegistryUnit()) || RegistryUnit.B.equalsByKey(r.getRegistryUnit()))
- ).toList();
+ forRegistries = selectRegistryForSDF04(account, fullNames);
} else if (UnlockResourcesPayload.MODE_SDF12.equals(sdfMode)) {
- forRegistries = forRegistries.stream().filter(
- (Registry r) ->
- RegistryDesignation.A.equalsByKey(r.getRegistryDesignation())
- &&
- RegistryInstrumentType.S.equalsByKey(r.getRegistryInstrumentType())
- &&
- (RegistryUnit.F.equalsByKey(r.getRegistryUnit()) || RegistryUnit.B.equalsByKey(r.getRegistryUnit()))
- ).toList();
+ forRegistries = selectRegistryForSDF12(account, securityId);
} else {
throw new IllegalArgumentException("Mode not support: " + sdfMode);
}
- for (Registry registry : forRegistries) {
- boolean modified = unlockRegistry(registry, value);
- if (modified) {
- registry.setUpdated(Instant.now());
- registryImdg.update(registry);
- log.trace("Registry {} changed; value +- registry",
- registry.getId(), registry.getRegistryCode(), value);
+ ImdgTransaction tx = imdgProvider.newTransaction();
+ boolean txOk = false;
+ try {
+ log.debug("Processing transaction {}, input {} registers.", tx, forRegistries.size());
+ Imdg registryTxImdg = tx.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
+ int nUpdates = 0;
+ for (Registry registry : forRegistries) {
+ boolean modified = unlockRegistry(registry, value);
+ if (modified) {
+ registry.setUpdated(Instant.now());
+ registryTxImdg.update(registry);
+ nUpdates++;
+ log.trace("Registry {} (registryCode={}) changed; value +- {}",
+ registry.getId(), registry.getRegistryCode(), value);
+ }
+ }
+ log.info("Updated {} registry's (under transaction {}", nUpdates, tx);
+ txOk = true;
+ } finally {
+ if (txOk) {
+ log.debug("Commit transaction {}.", tx);
+ tx.commitTransaction();
+ } else {
+ log.info("Rollback transaction {}", tx);
+ tx.rollbackTransaction();
}
}
- //todo рекомендуется делать транзакцией.
StageResult> res = new StageResult<>(null, true);
res.setStageResult(forRegistries);
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/EndStageNotificationPayload.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/EndStageNotificationPayload.java
new file mode 100644
index 000000000..ba2b3836c
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/EndStageNotificationPayload.java
@@ -0,0 +1,13 @@
+package ru.spcex.clearing.session.stage.task;
+
+public class EndStageNotificationPayload {
+ private String section;
+
+ public String getSection() {
+ return section;
+ }
+
+ public void setSection(String section) {
+ this.section = section;
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/FinishingSessionPayload.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/FinishingSessionPayload.java
new file mode 100644
index 000000000..9842cdc30
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/FinishingSessionPayload.java
@@ -0,0 +1,13 @@
+package ru.spcex.clearing.session.stage.task;
+
+public class FinishingSessionPayload {
+ private Long sessionId;
+
+ public Long getSessionId() {
+ return sessionId;
+ }
+
+ public void setSessionId(Long sessionId) {
+ this.sessionId = sessionId;
+ }
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/FormingRegistersOnOSPayload.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/FormingRegistersOnOSPayload.java
new file mode 100644
index 000000000..c09807ebb
--- /dev/null
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/FormingRegistersOnOSPayload.java
@@ -0,0 +1,5 @@
+package ru.spcex.clearing.session.stage.task;
+
+public class FormingRegistersOnOSPayload {
+
+}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/RegistryOnObligationsAndSettlementRequirementsPayload.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/RegistryOnObligationsAndSettlementRequirementsPayload.java
deleted file mode 100644
index e0ddc15da..000000000
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/RegistryOnObligationsAndSettlementRequirementsPayload.java
+++ /dev/null
@@ -1,31 +0,0 @@
-package ru.spcex.clearing.session.stage.task;
-
-public class RegistryOnObligationsAndSettlementRequirementsPayload {
- private Long sessionId;
-// private Long companyId;
-// private Long securityId;
-
- public Long getSessionId() {
- return sessionId;
- }
-
- public void setSessionId(Long sessionId) {
- this.sessionId = sessionId;
- }
-
-// public Long getCompanyId() {
-// return companyId;
-// }
-//
-// public void setCompanyId(Long companyId) {
-// this.companyId = companyId;
-// }
-//
-// public Long getSecurityId() {
-// return securityId;
-// }
-//
-// public void setSecurityId(Long securityId) {
-// this.securityId = securityId;
-// }
-}
diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/UnlockResourcesPayload.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/UnlockResourcesPayload.java
index 1b9525f23..fc16b5d2b 100644
--- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/UnlockResourcesPayload.java
+++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/task/UnlockResourcesPayload.java
@@ -12,8 +12,6 @@ public class UnlockResourcesPayload {
*/
private String sdfMode;
- private Long sessionId;
-
/**
* sDf04.c_acc_cred
* sDf12.depoCodeSender
@@ -46,13 +44,6 @@ public class UnlockResourcesPayload {
}
- public Long getSessionId() {
- return sessionId;
- }
-
- public void setSessionId(Long sessionId) {
- this.sessionId = sessionId;
- }
public String getAccount() {
return account;
diff --git a/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/service/builder/sql/RegistryCodeSqlBuilderTest.java b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/service/builder/sql/RegistryCodeSqlBuilderTest.java
new file mode 100644
index 000000000..ca1c7b68c
--- /dev/null
+++ b/clearing-parent/clearing-service/src/test/java/ru/spcex/clearing/service/builder/sql/RegistryCodeSqlBuilderTest.java
@@ -0,0 +1,54 @@
+package ru.spcex.clearing.service.builder.sql;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import ru.spcex.clearing.models.RegistryTradingParams;
+import ru.spcex.platform.enumeration.RegistryCapacity;
+import ru.spcex.platform.enumeration.RegistryDesignation;
+import ru.spcex.platform.enumeration.RegistryInstrumentType;
+import ru.spcex.platform.enumeration.RegistryUnit;
+
+public class RegistryCodeSqlBuilderTest {
+
+ @Test
+ public void testBuildByOneObject() {
+ RegistryTradingParams registryTradingParams = new RegistryTradingParams(RegistryDesignation.C, RegistryInstrumentType.M, RegistryCapacity.B, RegistryUnit.R);
+ RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParams);
+ String sql = registryCodeSqlBuilder.build();
+ Assertions.assertEquals("(registryDesignation = 'C' and registryInstrumentType = 'M' and registryCapacity = 'B' and registryUnit = 'R')", sql);
+
+ registryTradingParams = new RegistryTradingParams(null, RegistryInstrumentType.M, RegistryCapacity.B, RegistryUnit.R);
+ registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParams);
+ sql = registryCodeSqlBuilder.build();
+ Assertions.assertEquals("(registryInstrumentType = 'M' and registryCapacity = 'B' and registryUnit = 'R')", sql);
+
+ registryTradingParams = new RegistryTradingParams(RegistryDesignation.C, null, null, RegistryUnit.R);
+ registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParams);
+ sql = registryCodeSqlBuilder.build();
+ Assertions.assertEquals("(registryDesignation = 'C' and registryUnit = 'R')", sql);
+ }
+
+ @Test
+ public void testBuildByFewObjects() {
+ RegistryTradingParams registryTradingParams_first = new RegistryTradingParams(RegistryDesignation.C, RegistryInstrumentType.M, RegistryCapacity.B, RegistryUnit.R);
+ RegistryTradingParams registryTradingParams_second = new RegistryTradingParams(RegistryDesignation.O, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.F);
+ RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParams_first, registryTradingParams_second);
+ String sql = registryCodeSqlBuilder.build();
+ Assertions.assertEquals("(registryDesignation = 'C' and registryInstrumentType = 'M' and registryCapacity = 'B' and registryUnit = 'R')" +
+ " or (registryDesignation = 'O' and registryInstrumentType = 'S' and registryCapacity = 'A' and registryUnit = 'F')", sql);
+
+ registryTradingParams_first = new RegistryTradingParams(null, RegistryInstrumentType.M, RegistryCapacity.B, RegistryUnit.R);
+ registryTradingParams_second = new RegistryTradingParams(null, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.F);
+ registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParams_first, registryTradingParams_second);
+ sql = registryCodeSqlBuilder.build();
+ Assertions.assertEquals("(registryInstrumentType = 'M' and registryCapacity = 'B' and registryUnit = 'R') or " +
+ "(registryInstrumentType = 'S' and registryCapacity = 'A' and registryUnit = 'F')", sql);
+
+ registryTradingParams_first = new RegistryTradingParams(RegistryDesignation.C, null, null, RegistryUnit.R);
+ registryTradingParams_second = new RegistryTradingParams(null, RegistryInstrumentType.S, RegistryCapacity.A, RegistryUnit.F);
+ registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(registryTradingParams_first, registryTradingParams_second);
+ sql = registryCodeSqlBuilder.build();
+ Assertions.assertEquals("(registryDesignation = 'C' and registryUnit = 'R') or " +
+ "(registryInstrumentType = 'S' and registryCapacity = 'A' and registryUnit = 'F')", sql);
+ }
+}
diff --git a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/validation/common/rules/FieldRequiredRule.java b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/validation/common/rules/FieldRequiredRule.java
index 43c6e0e0a..f73172da4 100644
--- a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/validation/common/rules/FieldRequiredRule.java
+++ b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/validation/common/rules/FieldRequiredRule.java
@@ -65,6 +65,7 @@ public record FieldRequiredRule(
}
@Override
+ // todo replace for error fieldName to fieldValue
public Optional validate(ImdgValidationContext context) {
R validatedObject = context.getValidatedObject();
V value = getter.apply(validatedObject);
diff --git a/clearing-parent/dbf-exporter/pom.xml b/clearing-parent/dbf-exporter/pom.xml
index fc2108169..12f8fbc23 100644
--- a/clearing-parent/dbf-exporter/pom.xml
+++ b/clearing-parent/dbf-exporter/pom.xml
@@ -23,11 +23,15 @@
org.springframework.boot
- spring-boot-starter-web
+ spring-boot-autoconfigure
- org.springframework.boot
- spring-boot-autoconfigure
+ org.springframework.integration
+ spring-integration-sftp
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
@@ -49,6 +53,21 @@
ru.spcex.clearing
classes
+
+ ru.spcex.clearing
+ classes
+
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+ ru.spcex.clearing
+ test-clearing
+ test
+
diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExporterConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExporterConfig.java
index 10261ad81..7994454f4 100644
--- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExporterConfig.java
+++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/DBFExporterConfig.java
@@ -1,80 +1,11 @@
package ru.spcex.clearing.dbf.exporter.config;
-import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
-import org.springframework.context.ApplicationContext;
-import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
-import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
-import ru.spcex.clearing.dbf.exporter.config.settings.ExportDBFServiceSettings;
-import ru.spcex.clearing.dbf.exporter.logic.stages.ExportFromHazelcast;
-import ru.spcex.clearing.dbf.exporter.logic.stages.PrepareDBFFile;
-import ru.spcex.clearing.dbf.exporter.logic.stages.Stage;
-import ru.spcex.platform.imdg.api.ImdgProvider;
-import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
-import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
-
-import java.util.LinkedList;
-import java.util.List;
@Configuration
@EnableConfigurationProperties
@ComponentScan(basePackages = {"ru.spcex.clearing.dbf.exporter"})
public class DBFExporterConfig {
-
- @Bean("taskExecutorHazelcastClientInitializer")
- public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
- return createThreadPoolTaskExecutor(1, true);
- }
-
- @Bean("taskExecutorIdGeneratorAwaiter")
- public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
- return createThreadPoolTaskExecutor(1, false);
- }
-
- @Bean("imdgProvider")
- public ImdgProvider imdgProvider(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
- @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
- ExportDBFServiceSettings settings) {
- HazelcastClientParams params = new HazelcastClientParams();
- params.setClusterMembers(settings.getHazelcast().getClusterMembers());
- params.setLogin(settings.getHazelcast().getLogin());
- params.setPassword(settings.getHazelcast().getPassword());
- return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params);
- }
-
- @Bean("pipeline")
- public List pipeline(ApplicationContext context) {
- List pipeline = new LinkedList<>();
-
- pipeline.add(context.getBean(PrepareDBFFile.class));
- pipeline.add(context.getBean(ExportFromHazelcast.class));
-
- return pipeline;
- }
-
- @Bean("executor")
- public ThreadPoolTaskExecutor executor(ExportDBFServiceSettings settings) {
- ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
- executor.setMaxPoolSize(settings.getCommon().getThreadsCount());
- executor.setCorePoolSize(settings.getCommon().getThreadsCount());
- executor.setThreadNamePrefix("dbf-exporter");
- executor.setWaitForTasksToCompleteOnShutdown(true);
- executor.setAwaitTerminationSeconds(300);
- executor.initialize();
- return executor;
- }
-
- private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
- ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
- if (maxPoolSz > 2) {
- pool.setKeepAliveSeconds(60);
- pool.setAllowCoreThreadTimeOut(true);
- }
- pool.setCorePoolSize(maxPoolSz);
- pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
- return pool;
- }
-
}
diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/ImdgConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/ImdgConfig.java
new file mode 100644
index 000000000..bc37887f9
--- /dev/null
+++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/ImdgConfig.java
@@ -0,0 +1,63 @@
+package ru.spcex.clearing.dbf.exporter.config;
+
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.context.properties.EnableConfigurationProperties;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.ComponentScan;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
+import ru.spcex.clearing.dbf.exporter.config.settings.ExportDBFServiceSettings;
+import ru.spcex.platform.imdg.api.ImdgProvider;
+import ru.spcex.platform.imdg.iml.hazelcast.config.HazelcastClientParams;
+import ru.spcex.platform.imdg.iml.hazelcast.service.HazelcastService;
+
+@Configuration
+@EnableConfigurationProperties
+@ComponentScan(basePackages = {"ru.spcex.clearing.dbf.exporter"})
+public class ImdgConfig {
+
+ @Bean("taskExecutorHazelcastClientInitializer")
+ public ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer() {
+ return createThreadPoolTaskExecutor(1, true);
+ }
+
+ @Bean("taskExecutorIdGeneratorAwaiter")
+ public ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter() {
+ return createThreadPoolTaskExecutor(1, false);
+ }
+
+ @Bean("imdgProvider")
+ public ImdgProvider imdgProvider(@Qualifier("taskExecutorHazelcastClientInitializer") ThreadPoolTaskExecutor taskExecutorHazelcastClientInitializer,
+ @Qualifier("taskExecutorIdGeneratorAwaiter") ThreadPoolTaskExecutor taskExecutorIdGeneratorAwaiter,
+ ExportDBFServiceSettings settings) {
+ HazelcastClientParams params = new HazelcastClientParams();
+ params.setClusterMembers(settings.getHazelcast().getClusterMembers());
+ params.setLogin(settings.getHazelcast().getLogin());
+ params.setPassword(settings.getHazelcast().getPassword());
+ return new HazelcastService(taskExecutorHazelcastClientInitializer, taskExecutorIdGeneratorAwaiter, params);
+ }
+
+ @Bean("executor")
+ public ThreadPoolTaskExecutor executor(ExportDBFServiceSettings settings) {
+ ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
+ executor.setMaxPoolSize(settings.getCommon().getThreadsCount());
+ executor.setCorePoolSize(settings.getCommon().getThreadsCount());
+ executor.setThreadNamePrefix("dbf-exporter");
+ executor.setWaitForTasksToCompleteOnShutdown(true);
+ executor.setAwaitTerminationSeconds(300);
+ executor.initialize();
+ return executor;
+ }
+
+ private static ThreadPoolTaskExecutor createThreadPoolTaskExecutor(int maxPoolSz, boolean waitForCompletion) {
+ ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
+ if (maxPoolSz > 2) {
+ pool.setKeepAliveSeconds(60);
+ pool.setAllowCoreThreadTimeOut(true);
+ }
+ pool.setCorePoolSize(maxPoolSz);
+ pool.setWaitForTasksToCompleteOnShutdown(waitForCompletion);
+ return pool;
+ }
+
+}
diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java
index 5f0a59e40..bd4994aca 100644
--- a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java
+++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/KafkaSenderConfig.java
@@ -17,6 +17,8 @@ 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;
+
//отдельный конфиг для sender чтобы сделать required false
@Configuration
public class KafkaSenderConfig {
@@ -38,7 +40,6 @@ public class KafkaSenderConfig {
return KafkaProducerFactory.producerFactory(kafkaSettings);
}
- @Autowired(required = false)
@Bean("kafkaTemplate")
public KafkaTemplate kafkaTemplate(ProducerFactory pf) {
if (pf == null) {
@@ -47,16 +48,11 @@ public class KafkaSenderConfig {
return new KafkaTemplate<>(pf);
}
- @Autowired(required = false)
@Bean
- public KafkaSender kafkaSender(KafkaTemplate kafkaTemplate,
- ImdgProvider imdgProvider) {
- if (kafkaTemplate == null) {
- log.info("Can not create KafkaSender: no kafka-producer settings");
- return null;
- }
+ public Supplier kafkaSenderSupplier(KafkaTemplate kafkaTemplate,
+ ImdgProvider imdgProvider) {
ImdgId imdgIdGenerator = imdgProvider.getImdgIdGenerator();
- return KafkaSender
+ return () -> KafkaSender
.setup()
.setKafkaTemplate(kafkaTemplate)
.idGenerator(imdgIdGenerator::nextId)
diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/PipelineConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/PipelineConfig.java
new file mode 100644
index 000000000..bccca7609
--- /dev/null
+++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/PipelineConfig.java
@@ -0,0 +1,26 @@
+package ru.spcex.clearing.dbf.exporter.config;
+
+import org.springframework.context.ApplicationContext;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import ru.spcex.clearing.dbf.exporter.logic.stages.ExportFromHazelcast;
+import ru.spcex.clearing.dbf.exporter.logic.stages.Journal;
+import ru.spcex.clearing.dbf.exporter.logic.stages.PrepareDBFFile;
+import ru.spcex.clearing.dbf.exporter.logic.stages.Stage;
+
+import java.util.LinkedList;
+import java.util.List;
+
+@Configuration
+public class PipelineConfig {
+ @Bean("pipeline")
+ public List pipeline(ApplicationContext context) {
+ List pipeline = new LinkedList<>();
+
+ pipeline.add(context.getBean(PrepareDBFFile.class));
+ pipeline.add(context.getBean(ExportFromHazelcast.class));
+ pipeline.add(context.getBean(Journal.class));
+
+ return pipeline;
+ }
+}
diff --git a/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/SFTPConfig.java b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/SFTPConfig.java
new file mode 100644
index 000000000..5f07a6eac
--- /dev/null
+++ b/clearing-parent/dbf-exporter/src/main/java/ru/spcex/clearing/dbf/exporter/config/SFTPConfig.java
@@ -0,0 +1,95 @@
+package ru.spcex.clearing.dbf.exporter.config;
+
+import com.jcraft.jsch.ChannelSftp;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.expression.common.LiteralExpression;
+import org.springframework.integration.annotation.Gateway;
+import org.springframework.integration.annotation.MessagingGateway;
+import org.springframework.integration.annotation.ServiceActivator;
+import org.springframework.integration.channel.DirectChannel;
+import org.springframework.integration.dsl.IntegrationFlow;
+import org.springframework.integration.dsl.IntegrationFlows;
+import org.springframework.integration.file.remote.session.CachingSessionFactory;
+import org.springframework.integration.file.remote.session.SessionFactory;
+import org.springframework.integration.sftp.gateway.SftpOutboundGateway;
+import org.springframework.integration.sftp.outbound.SftpMessageHandler;
+import org.springframework.integration.sftp.session.DefaultSftpSessionFactory;
+import org.springframework.integration.sftp.session.SftpFileInfo;
+import org.springframework.messaging.MessageChannel;
+import org.springframework.messaging.MessageHandler;
+import ru.spcex.clearing.dbf.exporter.config.settings.ExportDBFServiceSettings;
+
+import java.io.File;
+import java.util.List;
+
+import static org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.Command.LS;
+
+@Configuration
+public class SFTPConfig {
+
+ @Bean
+ public SessionFactory sftpSessionFactory(ExportDBFServiceSettings settings) {
+ DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true);
+ factory.setHost(settings.getStore().getServerIp());
+ factory.setPort(settings.getStore().getServerPort());
+ factory.setUser(settings.getStore().getUser());
+ factory.setPassword(settings.getStore().getPassword());
+ factory.setAllowUnknownKeys(true);
+ return new CachingSessionFactory<>(factory);
+ }
+
+ @Bean
+ @ServiceActivator(inputChannel = "toSftpChannel")
+ public MessageHandler handler(SessionFactory sessionFactory, ExportDBFServiceSettings settings) {
+ SftpMessageHandler handler = new SftpMessageHandler(sessionFactory);
+ handler.setRemoteDirectoryExpression(new LiteralExpression(settings.getStore().getOutDir()));
+ handler.setAutoCreateDirectory(true);
+ handler.setFileNameGenerator(message -> {
+ if (message.getPayload() instanceof File) {
+ return ((File) message.getPayload()).getName();
+ }else {
+ throw new IllegalArgumentException("File must expected as payload.");
+ }
+ });
+ return handler;
+ }
+
+ @MessagingGateway
+ public interface DbfGateway {
+ @Gateway(requestChannel = "toSftpChannel")
+ void sendToSftp(File file);
+
+ @Gateway(requestChannel = "listSftpChannel")
+ List listFiles(String dir);
+ }
+
+ @Bean
+ public MessageChannel listSftpChannel(SessionFactory sessionFactory, ExportDBFServiceSettings settings) {
+ DirectChannel dc = new DirectChannel();
+ dc.subscribe(handlerList(sessionFactory, settings));
+ return dc;
+ }
+
+ @Bean
+ public MessageChannel toSftpChannel(SessionFactory sessionFactory, ExportDBFServiceSettings settings) {
+ DirectChannel dc = new DirectChannel();
+ dc.subscribe(handler(sessionFactory, settings));
+ return dc;
+ }
+
+ @Bean
+ @ServiceActivator(inputChannel = "listSftpChannel")
+ public MessageHandler handlerList(SessionFactory sessionFactory, ExportDBFServiceSettings settings) {
+ String expression = "'/%s'".formatted(settings.getStore().getOutDir());
+ SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sessionFactory, LS.getCommand(), expression);
+ return sftpOutboundGateway;
+ }
+
+ @Bean
+ public IntegrationFlow sftpOutboundListFlow(SessionFactory