Compare commits
4 commits
10_01_2024
...
kill-sessi
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a91ef94a44 | ||
|
|
efe5744e16 | ||
|
|
ec9586eef5 | ||
|
|
e88eef031c |
15 changed files with 350 additions and 21 deletions
|
|
@ -46,6 +46,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl;
|
|||
|
||||
import java.util.function.BiFunction;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
@Configuration
|
||||
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")
|
||||
public UserRoleVerification userRoleVerification(ImdgProvider imdgProvider, IMessageResolver msgs) {
|
||||
return new UserRoleVerification(imdgProvider, msgs, ClearingError.UserVerifyDenial);
|
||||
|
|
|
|||
|
|
@ -54,6 +54,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
private final BalanceRevise balanceRevise;
|
||||
private final Sdf05Sender sdf05Sender;
|
||||
private final StatementServiceV2 statementService;
|
||||
private final SessionTerminator sessionTerminator;
|
||||
private final PaymentInstructionOutboundService pmtOutboundService;
|
||||
|
||||
@Autowired
|
||||
|
|
@ -65,7 +66,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
||||
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
|
||||
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);
|
||||
this.errorResolver = errorResolver;
|
||||
this.clearingService = clearingService;
|
||||
|
|
@ -83,6 +84,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
this.balanceRevise = balanceRevise;
|
||||
this.sdf05Sender = sdf05Sender;
|
||||
this.statementService = statementService;
|
||||
this.sessionTerminator = sessionTerminator;
|
||||
this.pmtOutboundService = pmtOutboundService;
|
||||
}
|
||||
|
||||
|
|
@ -184,6 +186,11 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
callback(LauncherCommandRequest.class)
|
||||
.setConsumer(task -> sdf05Sender.sendSdf05("9"))
|
||||
.forDestination(Task.sdf05WithCode9Final.topic(), callbacks::put);
|
||||
|
||||
callback(Object.class)
|
||||
.setFunction(sessionTerminator::stopCurrentSession)
|
||||
.forDestination(Consts.KILL_SESSION, callbacks::put);
|
||||
|
||||
init();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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.SingleAssetResponse;
|
||||
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.registry.AssetTBFProcessing;
|
||||
import ru.spcex.clearing.service.registry.DmiService;
|
||||
import ru.spcex.clearing.service.schedule.TradingTimeService;
|
||||
import ru.spcex.clearing.service.validation.Sdf06NewValidationRule;
|
||||
|
|
@ -66,6 +68,7 @@ public class Sdf06Executor {
|
|||
private final KafkaSender kafkaSender;
|
||||
private final FilenameObtainer filenameObtainer;
|
||||
private final DmiService dmiService;
|
||||
private final AssetTBFProcessing assets;
|
||||
|
||||
private final static BigDecimal successResult = BigDecimal.ZERO;
|
||||
//идет сессия (не возвращаем такую ошибку)
|
||||
|
|
@ -87,7 +90,7 @@ public class Sdf06Executor {
|
|||
public Sdf06Executor(ImdgProvider imdgProvider,
|
||||
IMessageResolver messageResolver,
|
||||
@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.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
|
||||
this.tradingTimeService = tradingTimeService;
|
||||
|
|
@ -102,6 +105,7 @@ public class Sdf06Executor {
|
|||
this.sdf07Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf07, SDf07.class);
|
||||
this.filenameObtainer = filenameObtainer;
|
||||
this.dmiService = dmiService;
|
||||
this.assets = assets;
|
||||
}
|
||||
|
||||
public void execute(BaseRequest<StatementRequest> systemRequest) {
|
||||
|
|
@ -160,6 +164,8 @@ public class Sdf06Executor {
|
|||
CurrencyCode.RUB.getKey(),
|
||||
safeBD(sDf06.getSum()),
|
||||
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;
|
||||
}
|
||||
|
|
@ -266,11 +272,16 @@ public class Sdf06Executor {
|
|||
Instant updatedTime = Instant.now();
|
||||
if (gatewayMsg.isApproved()) {
|
||||
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(),
|
||||
safeBD(sdf06.getSum()),
|
||||
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 {
|
||||
log.trace("statement.id={}, sdf07.id={}, sdf06.id={} rejected (by gateway answer)",
|
||||
statementId,
|
||||
|
|
|
|||
|
|
@ -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.importexport.SwtExporterRequest;
|
||||
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.schedule.TradingTimeService;
|
||||
import ru.spcex.clearing.service.validation.ValidationStored;
|
||||
|
|
@ -68,6 +69,7 @@ public class Sdf10Executor {
|
|||
private final TradingTimeService tradingTimeService;
|
||||
private final KafkaSender kafkaSender;
|
||||
private final SecuritySelector<Security> scrtSlct;
|
||||
private final AssetTBFProcessing assets;
|
||||
|
||||
private final static String OK = "OK";
|
||||
private final static String SYNTAX_ERROR = "Синтаксическая ошибка (файл сформирован неверно)";
|
||||
|
|
@ -88,7 +90,7 @@ public class Sdf10Executor {
|
|||
public Sdf10Executor(ImdgProvider imdgProvider,
|
||||
RegistryManager rgsMng, IMessageResolver messageResolver,
|
||||
@Qualifier("sdf10Validator") Function<SDf10, IValidator> sDf10Validator,
|
||||
TradingTimeService tradingTimeService, KafkaSender kafkaSender) {
|
||||
TradingTimeService tradingTimeService, KafkaSender kafkaSender, AssetTBFProcessing assets) {
|
||||
this.imdgProvider = imdgProvider;
|
||||
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.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.plannerAllTodayImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_PlannerAllToday, PlannerAllToday.class);
|
||||
this.scrtSlct = new SecuritySelector<>(imdgProvider, Security.class);
|
||||
this.assets = assets;
|
||||
}
|
||||
|
||||
public void execute(BaseRequest<StatementRequest> systemRequest) {
|
||||
|
|
@ -166,7 +169,11 @@ public class Sdf10Executor {
|
|||
SDf11 sDf11 = processedApproved(stmt, sDf10, now, sdf11GroupId);
|
||||
BigDecimal amount = InOutDirection.in.equalsByKey(stmt.getInOutDirection()) ?
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
@ -190,6 +197,7 @@ public class Sdf10Executor {
|
|||
return;
|
||||
}
|
||||
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) {
|
||||
|
|
@ -203,6 +211,7 @@ public class Sdf10Executor {
|
|||
return;
|
||||
}
|
||||
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) {
|
||||
|
|
|
|||
|
|
@ -195,7 +195,7 @@ public class Sdf20Executor {
|
|||
Long securityId) {
|
||||
Statement stmt = new Statement();
|
||||
stmt.setAddresseeId(companyId);
|
||||
stmt.setSenderId(Sender.Prc.getId());
|
||||
stmt.setSenderId(Sender.Rdc.getId());
|
||||
stmt.setStatementType(StatementType.incr.getKey());
|
||||
//fixme stmt.setComment(sdf10.());
|
||||
stmt.setAccountId(account.getId());
|
||||
|
|
|
|||
|
|
@ -260,7 +260,7 @@ public class Sdf21Executor extends AbstractExecutor<SDf21> {
|
|||
private Statement create(SDf21 sdf21, Company cmp, Account acc) {
|
||||
Statement statement = new Statement();
|
||||
statement.setAddresseeId(cmp.getId());
|
||||
statement.setSenderId(Sender.Prc.getId());
|
||||
statement.setSenderId(Sender.Rdc.getId());
|
||||
statement.setStatementType(StatementType.incr.getKey());
|
||||
//fixme statement.setComment(sdf21.());
|
||||
statement.setAccountId(acc.getId());
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
String securitySymbol,
|
||||
Long securityId,
|
||||
|
|
@ -244,6 +263,23 @@ public class AssetTBFProcessing {
|
|||
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) {
|
||||
return safeBD(registry.getBalance()).compareTo(safeBD(amount)) < 0;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -12,9 +12,7 @@ import ru.spcex.clearing.error.ClearingError;
|
|||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.service.AssetTrio;
|
||||
import ru.spcex.clearing.service.registry.RegistryManager;
|
||||
import ru.spcex.platform.enumeration.AccountType;
|
||||
import ru.spcex.platform.enumeration.NoRef;
|
||||
import ru.spcex.platform.enumeration.OperationCode;
|
||||
import ru.spcex.platform.enumeration.*;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
||||
|
|
@ -67,6 +65,12 @@ public enum Sdf20ValidationRule implements IValidationRule<ImdgValidationContext
|
|||
if (acc == null) {
|
||||
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);
|
||||
return empty();
|
||||
}
|
||||
|
|
@ -86,6 +90,9 @@ public enum Sdf20ValidationRule implements IValidationRule<ImdgValidationContext
|
|||
if (company == null) {
|
||||
return of(ClearingError.CompanyNotFoundB, acc.getCompanyId());
|
||||
}
|
||||
if (!WorkflowStatus.Active.equalsByKey(company.getWorkflowStatus())) {
|
||||
return of(ClearingError.CompanyNotActive, company.getId());
|
||||
}
|
||||
context.storeObject(ValidationStored.Sdf20Company, company);
|
||||
return empty();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
@ -32,5 +32,7 @@ public enum ValidationStored {
|
|||
|
||||
SecurityBySecurityCode,
|
||||
|
||||
SessionTerminate,
|
||||
|
||||
RegistrysByContract, SplitDepositMaxNumber
|
||||
}
|
||||
|
|
|
|||
|
|
@ -105,6 +105,11 @@ public abstract class AbstractSession {
|
|||
|
||||
protected void continueRunning(TaskType t) {
|
||||
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.currSession.setSessionStatus(t.getKey());
|
||||
this.currSession.setUpdated(Instant.now());
|
||||
|
|
|
|||
|
|
@ -66,17 +66,7 @@ public class SessionManager {
|
|||
return;
|
||||
}
|
||||
BaseRequest<?> baseRequest = new BaseRequest<>();
|
||||
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;
|
||||
}
|
||||
|
||||
AbstractSession session = sessionByType(sessionType);
|
||||
if (session != null) {
|
||||
//checkActive
|
||||
checkAllowSessionStart();
|
||||
|
|
@ -99,4 +89,30 @@ public class SessionManager {
|
|||
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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -104,6 +104,9 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
|
|||
public final static RegistryTradingParams AM__;
|
||||
public final static RegistryTradingParams AS__;
|
||||
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 DMAU;
|
||||
public final static RegistryTradingParams DMAI;
|
||||
|
|
@ -240,6 +243,18 @@ public record RegistryTradingParams(RegistryDesignation registryDesignation,
|
|||
null,
|
||||
null,
|
||||
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);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -139,6 +139,7 @@ public interface Consts {
|
|||
String SDF56_PROCESS = "sdf56-process";
|
||||
String CONTINUE_SESSION_BN_FIRST_PART = "sdf57-process";
|
||||
String CONTINUE_SESSION_BN_SECOND_PART = "sdf13-process";
|
||||
String KILL_SESSION = "kill-session";
|
||||
String SWT_EXPORTER = "swt-exporter";
|
||||
String REVISE_PROCESS = "revise-process";
|
||||
String EXPORT_PROCESS = "export-process";
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue