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.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);

View file

@ -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();
}

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.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,

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.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) {

View file

@ -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());

View file

@ -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());

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,
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;
}

View file

@ -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();
}

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,
SessionTerminate,
RegistrysByContract, SplitDepositMaxNumber
}

View file

@ -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());

View file

@ -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;
}
}

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 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);
}
}

View file

@ -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";