Compare commits

...
Sign in to create a new pull request.

4 commits

Author SHA1 Message Date
ialbert
a91ef94a44 http://jira.mfd.msk:8088/browse/CLS-479 2023-11-09 14:43:15 +03:00
ialbert
efe5744e16 http://jira.mfd.msk:8088/browse/CLS-388 2023-11-09 14:42:38 +03:00
ialbert
ec9586eef5 http://jira.mfd.msk:8088/browse/CLS-584 2023-11-08 19:06:46 +03:00
ialbert
e88eef031c kill session command 2023-11-03 18:57:41 +03:00
15 changed files with 350 additions and 21 deletions

View file

@ -46,6 +46,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl;
import java.util.function.BiFunction; import java.util.function.BiFunction;
import java.util.function.Function; import java.util.function.Function;
import java.util.function.Supplier;
@Configuration @Configuration
public class ValidationConfig { public class ValidationConfig {
@ -490,6 +491,19 @@ public class ValidationConfig {
}; };
} }
@Bean("terminateSessionValidator")
public Supplier<IValidator> terminateSessionValidator() {
return () -> {
ImdgValidationContext<Session> ctx = new ImdgValidationContext<>();
ctx.setLogPrefix(LogPrefixId.INSTANCE);
ctx.addImdg(IMDGDistributedNames.Map_Session, imdgSession);
return new ValidatorImpl<>(ctx,
SessionTerminationValidationRule.SessionPresent,
SessionTerminationValidationRule.StatusCheck
);
};
}
@Bean("userRoleVerification") @Bean("userRoleVerification")
public UserRoleVerification userRoleVerification(ImdgProvider imdgProvider, IMessageResolver msgs) { public UserRoleVerification userRoleVerification(ImdgProvider imdgProvider, IMessageResolver msgs) {
return new UserRoleVerification(imdgProvider, msgs, ClearingError.UserVerifyDenial); return new UserRoleVerification(imdgProvider, msgs, ClearingError.UserVerifyDenial);

View file

@ -54,6 +54,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
private final BalanceRevise balanceRevise; private final BalanceRevise balanceRevise;
private final Sdf05Sender sdf05Sender; private final Sdf05Sender sdf05Sender;
private final StatementServiceV2 statementService; private final StatementServiceV2 statementService;
private final SessionTerminator sessionTerminator;
private final PaymentInstructionOutboundService pmtOutboundService; private final PaymentInstructionOutboundService pmtOutboundService;
@Autowired @Autowired
@ -65,7 +66,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
SecondaryAuctionT0Session secondaryAuctionT0Session, SecondaryAuctionT0Session secondaryAuctionT0Session,
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager, PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
Sdf06Executor sdf06Executor, Sdf06Executor sdf06Executor,
Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, PaymentInstructionOutboundService pmtOutboundService) { Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, SessionTerminator sessionTerminator, PaymentInstructionOutboundService pmtOutboundService) {
super(kafkaQueue, kafkaResponseQueue); super(kafkaQueue, kafkaResponseQueue);
this.errorResolver = errorResolver; this.errorResolver = errorResolver;
this.clearingService = clearingService; this.clearingService = clearingService;
@ -83,6 +84,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
this.balanceRevise = balanceRevise; this.balanceRevise = balanceRevise;
this.sdf05Sender = sdf05Sender; this.sdf05Sender = sdf05Sender;
this.statementService = statementService; this.statementService = statementService;
this.sessionTerminator = sessionTerminator;
this.pmtOutboundService = pmtOutboundService; this.pmtOutboundService = pmtOutboundService;
} }
@ -184,6 +186,11 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(LauncherCommandRequest.class) callback(LauncherCommandRequest.class)
.setConsumer(task -> sdf05Sender.sendSdf05("9")) .setConsumer(task -> sdf05Sender.sendSdf05("9"))
.forDestination(Task.sdf05WithCode9Final.topic(), callbacks::put); .forDestination(Task.sdf05WithCode9Final.topic(), callbacks::put);
callback(Object.class)
.setFunction(sessionTerminator::stopCurrentSession)
.forDestination(Consts.KILL_SESSION, callbacks::put);
init(); init();
} }

View file

@ -22,7 +22,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRe
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.AssetTrio;
import ru.spcex.clearing.service.integration.GatewayRequestCreator; import ru.spcex.clearing.service.integration.GatewayRequestCreator;
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
import ru.spcex.clearing.service.registry.DmiService; import ru.spcex.clearing.service.registry.DmiService;
import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.schedule.TradingTimeService;
import ru.spcex.clearing.service.validation.Sdf06NewValidationRule; import ru.spcex.clearing.service.validation.Sdf06NewValidationRule;
@ -66,6 +68,7 @@ public class Sdf06Executor {
private final KafkaSender kafkaSender; private final KafkaSender kafkaSender;
private final FilenameObtainer filenameObtainer; private final FilenameObtainer filenameObtainer;
private final DmiService dmiService; private final DmiService dmiService;
private final AssetTBFProcessing assets;
private final static BigDecimal successResult = BigDecimal.ZERO; private final static BigDecimal successResult = BigDecimal.ZERO;
//идет сессия (не возвращаем такую ошибку) //идет сессия (не возвращаем такую ошибку)
@ -87,7 +90,7 @@ public class Sdf06Executor {
public Sdf06Executor(ImdgProvider imdgProvider, public Sdf06Executor(ImdgProvider imdgProvider,
IMessageResolver messageResolver, IMessageResolver messageResolver,
@Qualifier("sdf06ValidatorNew") Function<SDf06, IValidator> sDf06Validator, @Qualifier("sdf06ValidatorNew") Function<SDf06, IValidator> sDf06Validator,
TradingTimeService tradingTimeService, KafkaSender kafkaSender, FilenameObtainer filenameObtainer, DmiService dmiService) { TradingTimeService tradingTimeService, KafkaSender kafkaSender, FilenameObtainer filenameObtainer, DmiService dmiService, AssetTBFProcessing assets) {
this.imdgProvider = imdgProvider; this.imdgProvider = imdgProvider;
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.tradingTimeService = tradingTimeService; this.tradingTimeService = tradingTimeService;
@ -102,6 +105,7 @@ public class Sdf06Executor {
this.sdf07Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf07, SDf07.class); this.sdf07Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf07, SDf07.class);
this.filenameObtainer = filenameObtainer; this.filenameObtainer = filenameObtainer;
this.dmiService = dmiService; this.dmiService = dmiService;
this.assets = assets;
} }
public void execute(BaseRequest<StatementRequest> systemRequest) { public void execute(BaseRequest<StatementRequest> systemRequest) {
@ -160,6 +164,8 @@ public class Sdf06Executor {
CurrencyCode.RUB.getKey(), CurrencyCode.RUB.getKey(),
safeBD(sDf06.getSum()), safeBD(sDf06.getSum()),
sDf06.getNumber().toString()); sDf06.getNumber().toString());
Optional<AssetTrio> asts = assets.searchMoneyByAccAndCompany(company.getId(), tcr.getMoneyAccountId());
asts.ifPresent(ast -> assets.process(ast.a__b(), ast.a__t(), ast.a__f(), BigDecimal.ZERO));
} }
sdf07WasCreated = true; sdf07WasCreated = true;
} }
@ -266,11 +272,16 @@ public class Sdf06Executor {
Instant updatedTime = Instant.now(); Instant updatedTime = Instant.now();
if (gatewayMsg.isApproved()) { if (gatewayMsg.isApproved()) {
processedApproved(stmt, sdf06, updatedTime); processedApproved(stmt, sdf06, updatedTime);
dmiService.setProcContract(searchTcrOnGatewayResponse(sdf06).map(SpcexObjectBase::getId).orElse(null), Optional<TradingClearingRegistry> tcr = searchTcrOnGatewayResponse(sdf06);
dmiService.setProcContract(tcr.map(SpcexObjectBase::getId).orElse(null),
CurrencyCode.RUB.getKey(), CurrencyCode.RUB.getKey(),
safeBD(sdf06.getSum()), safeBD(sdf06.getSum()),
sdf06.getNumber().toString() sdf06.getNumber().toString()
); );
if (tcr.isPresent()) {
Optional<AssetTrio> asts = assets.searchMoneyByAccAndCompany(tcr.get().getCompanyId(), tcr.get().getMoneyAccountId());
asts.ifPresent(a -> assets.process(a.a__b(), a.a__t(), a.a__f(), BigDecimal.ZERO));
}
} else { } else {
log.trace("statement.id={}, sdf07.id={}, sdf06.id={} rejected (by gateway answer)", log.trace("statement.id={}, sdf07.id={}, sdf06.id={} rejected (by gateway answer)",
statementId, statementId,

View file

@ -25,6 +25,7 @@ import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApp
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest; import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
import ru.spcex.clearing.service.registry.RegistryManager; import ru.spcex.clearing.service.registry.RegistryManager;
import ru.spcex.clearing.service.schedule.TradingTimeService; import ru.spcex.clearing.service.schedule.TradingTimeService;
import ru.spcex.clearing.service.validation.ValidationStored; import ru.spcex.clearing.service.validation.ValidationStored;
@ -68,6 +69,7 @@ public class Sdf10Executor {
private final TradingTimeService tradingTimeService; private final TradingTimeService tradingTimeService;
private final KafkaSender kafkaSender; private final KafkaSender kafkaSender;
private final SecuritySelector<Security> scrtSlct; private final SecuritySelector<Security> scrtSlct;
private final AssetTBFProcessing assets;
private final static String OK = "OK"; private final static String OK = "OK";
private final static String SYNTAX_ERROR = "Синтаксическая ошибка (файл сформирован неверно)"; private final static String SYNTAX_ERROR = "Синтаксическая ошибка (файл сформирован неверно)";
@ -88,7 +90,7 @@ public class Sdf10Executor {
public Sdf10Executor(ImdgProvider imdgProvider, public Sdf10Executor(ImdgProvider imdgProvider,
RegistryManager rgsMng, IMessageResolver messageResolver, RegistryManager rgsMng, IMessageResolver messageResolver,
@Qualifier("sdf10Validator") Function<SDf10, IValidator> sDf10Validator, @Qualifier("sdf10Validator") Function<SDf10, IValidator> sDf10Validator,
TradingTimeService tradingTimeService, KafkaSender kafkaSender) { TradingTimeService tradingTimeService, KafkaSender kafkaSender, AssetTBFProcessing assets) {
this.imdgProvider = imdgProvider; this.imdgProvider = imdgProvider;
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
@ -102,6 +104,7 @@ public class Sdf10Executor {
this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class); this.sdf11Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf11, SDf11.class);
this.plannerAllTodayImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerAllToday, PlannerAllToday.class); this.plannerAllTodayImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerAllToday, PlannerAllToday.class);
this.scrtSlct = new SecuritySelector<>(imdgProvider, Security.class); this.scrtSlct = new SecuritySelector<>(imdgProvider, Security.class);
this.assets = assets;
} }
public void execute(BaseRequest<StatementRequest> systemRequest) { public void execute(BaseRequest<StatementRequest> systemRequest) {
@ -166,7 +169,11 @@ public class Sdf10Executor {
SDf11 sDf11 = processedApproved(stmt, sDf10, now, sdf11GroupId); SDf11 sDf11 = processedApproved(stmt, sDf10, now, sdf11GroupId);
BigDecimal amount = InOutDirection.in.equalsByKey(stmt.getInOutDirection()) ? BigDecimal amount = InOutDirection.in.equalsByKey(stmt.getInOutDirection()) ?
stmt.getAmount() : safeBD(stmt.getAmount()).negate(); stmt.getAmount() : safeBD(stmt.getAmount()).negate();
createDs_iWithoutGateway(tcr.getId(), company.getId(), security.getSecuritySymbol(), amount, sDf10.getOutDocument()); createDs_iWithoutGateway(tcr.getId(),
company.getId(),
security.getSecuritySymbol(),
amount,
sDf10.getOutDocument());
sdf11WasCreated = true; sdf11WasCreated = true;
} }
} }
@ -190,6 +197,7 @@ public class Sdf10Executor {
return; return;
} }
createDs_i(as_t.get(), summ, outDocument); createDs_i(as_t.get(), summ, outDocument);
assets.processByA__t(as_t.get(), BigDecimal.ZERO);
} }
private void createDs_iWithoutGateway(Long tcrId, Long companyId, String securitySymobl, BigDecimal summ, String outDocument) { private void createDs_iWithoutGateway(Long tcrId, Long companyId, String securitySymobl, BigDecimal summ, String outDocument) {
@ -203,6 +211,7 @@ public class Sdf10Executor {
return; return;
} }
createDs_i(as_t.get(), summ, outDocument); createDs_i(as_t.get(), summ, outDocument);
assets.processByA__t(as_t.get(), BigDecimal.ZERO);
} }
private void createDs_i(Registry as_t, BigDecimal summ, String outDocument) { private void createDs_i(Registry as_t, BigDecimal summ, String outDocument) {

View file

@ -195,7 +195,7 @@ public class Sdf20Executor {
Long securityId) { Long securityId) {
Statement stmt = new Statement(); Statement stmt = new Statement();
stmt.setAddresseeId(companyId); stmt.setAddresseeId(companyId);
stmt.setSenderId(Sender.Prc.getId()); stmt.setSenderId(Sender.Rdc.getId());
stmt.setStatementType(StatementType.incr.getKey()); stmt.setStatementType(StatementType.incr.getKey());
//fixme stmt.setComment(sdf10.()); //fixme stmt.setComment(sdf10.());
stmt.setAccountId(account.getId()); stmt.setAccountId(account.getId());

View file

@ -260,7 +260,7 @@ public class Sdf21Executor extends AbstractExecutor<SDf21> {
private Statement create(SDf21 sdf21, Company cmp, Account acc) { private Statement create(SDf21 sdf21, Company cmp, Account acc) {
Statement statement = new Statement(); Statement statement = new Statement();
statement.setAddresseeId(cmp.getId()); statement.setAddresseeId(cmp.getId());
statement.setSenderId(Sender.Prc.getId()); statement.setSenderId(Sender.Rdc.getId());
statement.setStatementType(StatementType.incr.getKey()); statement.setStatementType(StatementType.incr.getKey());
//fixme statement.setComment(sdf21.()); //fixme statement.setComment(sdf21.());
statement.setAccountId(acc.getId()); statement.setAccountId(acc.getId());

View file

@ -134,6 +134,25 @@ public class AssetTBFProcessing {
} }
} }
public Optional<AssetTrio> processByA__t(Registry a__t, BigDecimal sum) {
Optional<Registry> a__b = rgsMng.findRegByUnit(a__t.getCompanyId(),
a__t.getAccountId(),
a__t.getContract(),
a__t,
RegistryUnit.B);
Optional<Registry> a__f = rgsMng.findRegByUnit(a__t.getCompanyId(),
a__t.getAccountId(),
a__t.getContract(),
a__t,
RegistryUnit.F);
if (a__b.isPresent() && a__f.isPresent()) {
process(a__b.get(), a__t, a__f.get(), sum);
return Optional.of(new AssetTrio(a__t, a__b.get(), a__f.get()));
} else {
return Optional.empty();
}
}
public AssetTrio createSAssets(Account acc, public AssetTrio createSAssets(Account acc,
String securitySymbol, String securitySymbol,
Long securityId, Long securityId,
@ -244,6 +263,23 @@ public class AssetTBFProcessing {
return Optional.empty(); return Optional.empty();
} }
public Optional<AssetTrio> searchSecurityByAccAndCompany(Long companyId, Long accountId, String securitySymbol) {
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
Function<RegistryTradingParams, ImdgPredicate> prdct = rgsCode -> pb.and(
pb.sql(RegistryCodeSqlBuilder.getInstance(rgsCode).build()),
pb.equals("companyId", companyId),
pb.equals("accountId", accountId),
pb.equals("securitySymbol", securitySymbol)
);
Registry amf = rgsImdg.getFirstObjectByPredicate(prdct.apply(AS_F));
Registry amb = rgsImdg.getFirstObjectByPredicate(prdct.apply(AS_B));
Registry amt = rgsImdg.getFirstObjectByPredicate(prdct.apply(AS_T));
if (amf != null && amb != null && amt != null) {
return Optional.of(new AssetTrio(amf, amb, amt));
}
return Optional.empty();
}
public boolean insufficientBalance(Registry registry, BigDecimal amount) { public boolean insufficientBalance(Registry registry, BigDecimal amount) {
return safeBD(registry.getBalance()).compareTo(safeBD(amount)) < 0; return safeBD(registry.getBalance()).compareTo(safeBD(amount)) < 0;
} }

View file

@ -12,9 +12,7 @@ import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.AssetTrio; import ru.spcex.clearing.service.AssetTrio;
import ru.spcex.clearing.service.registry.RegistryManager; import ru.spcex.clearing.service.registry.RegistryManager;
import ru.spcex.platform.enumeration.AccountType; import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.enumeration.NoRef;
import ru.spcex.platform.enumeration.OperationCode;
import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
@ -67,6 +65,12 @@ public enum Sdf20ValidationRule implements IValidationRule<ImdgValidationContext
if (acc == null) { if (acc == null) {
return of(ClearingError.AccountNotFoundB, depoCodeCl); return of(ClearingError.AccountNotFoundB, depoCodeCl);
} }
ServiceStatus status = IEnumKey.getEnumByKey(ServiceStatus.class, acc.getStatus());
if (!ServiceStatus.Active.equals(status)
&& !ServiceStatus.Reopened.equals(status)
&& !ServiceStatus.Appl.equals(status)) {
return of(ClearingError.AccountNotActive, acc.getId());
}
context.storeObject(ValidationStored.Sdf20Account, acc); context.storeObject(ValidationStored.Sdf20Account, acc);
return empty(); return empty();
} }
@ -86,6 +90,9 @@ public enum Sdf20ValidationRule implements IValidationRule<ImdgValidationContext
if (company == null) { if (company == null) {
return of(ClearingError.CompanyNotFoundB, acc.getCompanyId()); return of(ClearingError.CompanyNotFoundB, acc.getCompanyId());
} }
if (!WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
return of(ClearingError.CompanyNotActive, company.getId());
}
context.storeObject(ValidationStored.Sdf20Company, company); context.storeObject(ValidationStored.Sdf20Company, company);
return empty(); return empty();
} }

View file

@ -0,0 +1,55 @@
package ru.spcex.clearing.service.validation;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.clearing.classes.statics.data.misc.Session;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.platform.enumeration.SessionStatus;
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;
public enum SessionTerminationValidationRule implements IValidationRule<ImdgValidationContext<Session>> {
SessionPresent() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Session> context) {
Imdg<Session> sessionImdg = context.obtainMap(IMDGDistributedNames.Map_Session, Session.class);
Session session = sessionImdg.getFirstObjectByFieldValues(Map.of(
"workflowStatus", SessionStatus.ACTV.getKey()
));
if (session == null) {
log.warn("Cannot stop session, failed to find workflowStatus 'ACTV'");
return of(ClearingError.RecordNotFound, "с активной клиринговой сессией для обнуления");
}
context.storeObject(ValidationStored.SessionTerminate, session);
context.setValidatedObject(session);
return empty();
}
},
StatusCheck() {
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<Session> context) {
Session session = context.getStoredObject(ValidationStored.SessionTerminate);
TaskType status = IEnumKey.getEnumByKey(TaskType.class, session.getSessionStatus());
if (status == null || status.ordinal() >= TaskType.FormingPaymentInstruction.ordinal()) {
log.warn("Cannot stop session.id={}, unknown sessionStatus '{}'",
session.getId(), session.getSessionStatus());
return of(ClearingError.IncorrectValue, "статуса для обнуления клиринговой сессии");
}
return empty();
}
}
;
private final static Logger log = LoggerFactory.getLogger(SessionTerminationValidationRule.class);
@Override
public String ruleName() {
return "SessionTerminationValidationRule." + name();
}
}

View file

@ -32,5 +32,7 @@ public enum ValidationStored {
SecurityBySecurityCode, SecurityBySecurityCode,
SessionTerminate,
RegistrysByContract, SplitDepositMaxNumber RegistrysByContract, SplitDepositMaxNumber
} }

View file

@ -105,6 +105,11 @@ public abstract class AbstractSession {
protected void continueRunning(TaskType t) { protected void continueRunning(TaskType t) {
synchronized (this.currStage) { synchronized (this.currStage) {
//if (WorkflowStatus.Blocked.equalsByKey(this.currSession.getWorkflowStatus())) {
// log.error("cannot run task {}.{}. session.id={} workflowStatus is {}",
// t, t.getKey(), this.currSession.getId(), this.currSession.getWorkflowStatus());
// throw new StageException();
//}
this.currStage.set(t); this.currStage.set(t);
this.currSession.setSessionStatus(t.getKey()); this.currSession.setSessionStatus(t.getKey());
this.currSession.setUpdated(Instant.now()); this.currSession.setUpdated(Instant.now());

View file

@ -66,17 +66,7 @@ public class SessionManager {
return; return;
} }
BaseRequest<?> baseRequest = new BaseRequest<>(); BaseRequest<?> baseRequest = new BaseRequest<>();
AbstractSession session = null; AbstractSession session = sessionByType(sessionType);
switch (sessionType){
case IPOB -> session = primaryAuctionBnSession;
case IPOT -> session = primaryAuctionT0Session;
case IPO0 -> session = primaryAuctionB0Session;
case TRDT -> session = secondaryAuctionT0Session;
case MEDM -> session = intermediateMkrSession;
case FINL -> session = finalMkrSession;
case XDEP -> session = returnDepositSession;
}
if (session != null) { if (session != null) {
//checkActive //checkActive
checkAllowSessionStart(); checkAllowSessionStart();
@ -99,4 +89,30 @@ public class SessionManager {
throw new ValidationException(err); throw new ValidationException(err);
} }
} }
public void endSession(Session session) {
SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, session.getSessionType());
String err = "failed to stop session.id=%d: unknown sessionType: %s"
.formatted(session.getId(), session.getSessionType());
if (sessionType == null) {
throw new IllegalStateException(err);
}
AbstractSession sessionFlow = sessionByType(sessionType);
if (sessionFlow == null) throw new IllegalStateException(err);
sessionFlow.endSession();
}
private AbstractSession sessionByType(SessionType sessionType) {
AbstractSession session = null;
switch (sessionType){
case IPOB -> session = primaryAuctionBnSession;
case IPOT -> session = primaryAuctionT0Session;
case IPO0 -> session = primaryAuctionB0Session;
case TRDT -> session = secondaryAuctionT0Session;
case MEDM -> session = intermediateMkrSession;
case FINL -> session = finalMkrSession;
case XDEP -> session = returnDepositSession;
}
return session;
}
} }

View file

@ -0,0 +1,151 @@
package ru.spcex.clearing.session.stage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
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.misc.Session;
import ru.clearing.classes.statics.data.registry.Registry;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.notification.NotificationSender;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.service.registry.AssetTBFProcessing;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.util.services.RequestHelper;
import ru.spcex.platform.classes.base.SpcexObjectBase;
import ru.spcex.platform.enumeration.ObjectType;
import ru.spcex.platform.enumeration.Priority;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import ru.spcex.platform.enumeration.Section;
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 ru.spcex.platform.utils.validation.IValidator;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.Collection;
import java.util.Optional;
import java.util.function.Consumer;
import java.util.function.Supplier;
@Service
public class SessionTerminator {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<ExecutionDeposit> execDpstImdg;
private final Imdg<ExecutionFond> execFondImdg;
private final Imdg<Registry> rgsImdg;
private final NotificationSender notification;
private final AssetTBFProcessing assets;
private final IMessageResolver msgs;
private final RequestHelper reqInfo;
private final Supplier<IValidator> validation;
private final SessionManager sessionMng;
public SessionTerminator(ImdgProvider imdgProvider,
NotificationSender notification,
AssetTBFProcessing assets, IMessageResolver msgs,
RequestHelper reqInfo,
@Qualifier("terminateSessionValidator") Supplier<IValidator> validation, SessionManager sessionMng) {
this.notification = notification;
this.assets = assets;
this.msgs = msgs;
this.execDpstImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
this.execFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
this.rgsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
this.reqInfo = reqInfo;
this.validation = validation;
this.sessionMng = sessionMng;
}
public RequestInfoUpdate stopCurrentSession(BaseRequest<?> req) {
IValidator validator = validation.get();
Optional<EnumMessage> err = validator.tillFirstError();
if (err.isPresent()) {
notification.sendNotification(ObjectType.session, msgs.resolve(err.get()), Priority.HIGH);
return reqInfo.error(req.getId(), err.get());
}
Session session = validator.getStored(ValidationStored.SessionTerminate);
rollbackExecutions(session);
rollbackAssetBalances(session.getId());
deleteLiabilitiesAndClaims(session.getId());
sessionMng.endSession(session);
return reqInfo.success(session.getId());
}
private void rollbackExecutions(Session session) {
Instant now = Instant.now();
Consumer<Imdg<ExecutionCommon>> nullifier = imdg -> {
Collection<ExecutionCommon> execs = imdg.getCollectionObjectsBySQL(
"sessionId = " + session.getId()
);
for (ExecutionCommon exec : execs) {
exec.setSessionId(null);
exec.setCoverageStatus(null);
exec.setUpdated(now);
imdg.update(exec);
}
log.info("updated {} executions with sessionId={}", execs.size(), session.getId());
};
log.debug("session.id={} section {}", session.getId(), session.getSection());
if (Section.FOND.equalsByKey(session.getSection())) {
nullifier.accept(upcastImdg(execFondImdg));
} else if (Section.MKR.equalsByKey(session.getSection())) {
nullifier.accept(upcastImdg(execDpstImdg));
} else throw new IllegalStateException("unknown section.");
}
private void rollbackAssetBalances(Long sessionId) {
Instant now = Instant.now();
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.equals("sessionId", sessionId),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.A__B).build()),
pb.not(pb.equals("balance", BigDecimal.ZERO))
);
Collection<Registry> a__bs = rgsImdg.getCollectionObjectsByPredicate(prdct);
a__bs.forEach(rgs -> {
rgs.setBalance(BigDecimal.ZERO);
rgs.setUpdated(now);
rgsImdg.update(rgs);
assets.processByAm_b(rgs, BigDecimal.ZERO);
});
log.debug("set ZERO balance for {} A**B by predicate: {}", a__bs.size(), prdct);
}
private void deleteLiabilitiesAndClaims(Long sessionId) {
ImdgPredicateBuilder pb = rgsImdg.predicateBuilder();
ImdgPredicate prdct = pb.and(
pb.equals("sessionId", sessionId),
pb.or(
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.C__T).build()),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.L__T).build()),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.O__T).build()),
pb.sql(RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.T__T).build())
)
);
Collection<Registry> liabilitiesAndClaims = rgsImdg.getCollectionObjectsByPredicate(prdct);
liabilitiesAndClaims.forEach(rgsImdg::delete);
log.debug("deleted {} liabilities and claims by predicate: {}", liabilitiesAndClaims.size(), prdct);
}
private static <T extends SpcexObjectBase, U extends T> Imdg<T> upcastImdg(Imdg<U> imdg) {
@SuppressWarnings("unchecked")
Imdg<T> result = (Imdg<T>) imdg;
return result;
}
}

View file

@ -104,6 +104,9 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
public final static RegistryTradingParams AM__; public final static RegistryTradingParams AM__;
public final static RegistryTradingParams AS__; public final static RegistryTradingParams AS__;
public final static RegistryTradingParams CM__; public final static RegistryTradingParams CM__;
public final static RegistryTradingParams C__T;
public final static RegistryTradingParams O__T;
public final static RegistryTradingParams T__T;
public final static RegistryTradingParams LM__; public final static RegistryTradingParams LM__;
public final static RegistryTradingParams DMAU; public final static RegistryTradingParams DMAU;
public final static RegistryTradingParams DMAI; public final static RegistryTradingParams DMAI;
@ -240,6 +243,18 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
null, null,
null, null,
RegistryUnit.V); RegistryUnit.V);
C__T = new RegistryTradingParams(RegistryDesignation.C,
null,
null,
RegistryUnit.T);
O__T = new RegistryTradingParams(RegistryDesignation.O,
null,
null,
RegistryUnit.T);
T__T = new RegistryTradingParams(RegistryDesignation.T,
null,
null,
RegistryUnit.T);
} }
} }

View file

@ -139,6 +139,7 @@ public interface Consts {
String SDF56_PROCESS = "sdf56-process"; String SDF56_PROCESS = "sdf56-process";
String CONTINUE_SESSION_BN_FIRST_PART = "sdf57-process"; String CONTINUE_SESSION_BN_FIRST_PART = "sdf57-process";
String CONTINUE_SESSION_BN_SECOND_PART = "sdf13-process"; String CONTINUE_SESSION_BN_SECOND_PART = "sdf13-process";
String KILL_SESSION = "kill-session";
String SWT_EXPORTER = "swt-exporter"; String SWT_EXPORTER = "swt-exporter";
String REVISE_PROCESS = "revise-process"; String REVISE_PROCESS = "revise-process";
String EXPORT_PROCESS = "export-process"; String EXPORT_PROCESS = "export-process";