SDF56/51, UNIT session, registry#section

This commit is contained in:
ialbert 2024-06-07 18:09:05 +03:00
parent bebb114bea
commit a52f135ee2
20 changed files with 640 additions and 37 deletions

View file

@ -45,6 +45,7 @@ import ru.spcex.clearing.session.stage.SecondaryAuctionT0Session;
import ru.spcex.clearing.session.stage.SessionManager;
import ru.spcex.clearing.session.stage.SessionTerminator;
import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.clearing.session.stage.UnitedSession;
import ru.spcex.clearing.session.stage.impl.BalanceRevise;
import ru.spcex.clearing.statement.StatementServiceV2;
import ru.spcex.platform.enumeration.Task;
@ -65,6 +66,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
private final PrimaryAuctionB0Session primaryAuctionB0Session;
private final IntermediateMkrSession intermediateMkrSession;
private final FinalMkrSession finalMkrSession;
private final UnitedSession unitedSession;
private final ReturnDepositSession returnDepositSession;
private final SessionManager sessionManager;
private final Sdf06Executor sdf06Executor;
@ -83,7 +85,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
PrimaryAuctionBnSession primaryAuctionBnSession,
SecondaryAuctionT0Session secondaryAuctionT0Session,
CurrencySession currencySession,
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
PrimaryAuctionB0Session primaryAuctionB0Session, PrimaryAuctionT0Session primaryAuctionT0Session, IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, UnitedSession unitedSession, ReturnDepositSession returnDepositSession, SessionManager sessionManager,
Sdf06Executor sdf06Executor,
Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, SessionTerminator sessionTerminator, PaymentInstructionOutboundService pmtOutboundService) {
super(kafkaQueue, kafkaResponseQueue);
@ -98,6 +100,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
this.primaryAuctionT0Session = primaryAuctionT0Session;
this.intermediateMkrSession = intermediateMkrSession;
this.finalMkrSession = finalMkrSession;
this.unitedSession = unitedSession;
this.returnDepositSession = returnDepositSession;
this.sessionManager = sessionManager;
this.sdf06Executor = sdf06Executor;
@ -153,6 +156,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
primaryAuctionB0Session.continueSession(req);
intermediateMkrSession.continueSession(req);
finalMkrSession.continueSession(req);
unitedSession.continueSession(req);
returnDepositSession.continueSession(req);
})
.forDestination(Consts.CONTINUE_SESSION_BN_FIRST_PART, callbacks::put);
@ -207,7 +211,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
callback(LauncherCommandRequest.class)
.setConsumer(task -> {
balanceRevise.submit(new ru.spcex.clearing.session.stage.Task<>(TaskType.StartRevise, null));
balanceRevise.submit(new ru.spcex.clearing.session.stage.Task<>(TaskType.SDF51, null));
})
.forDestination(Task.getVerification.topic(), callbacks::put);

View file

@ -180,7 +180,7 @@ public class ExecutionCurrencyComponent {
eCurrency.setTradingDate(sTrades.getTradeDate());
eCurrency.setClearingDate(TimeUtil.toLocalDate(now));
eCurrency.setExchangeExecutionId(sTrades.getTradeNum());
eCurrency.setExchangeExecutionTime(sTrades.getTradeDateTime());
eCurrency.setExchangeExecutionTime(TimeUtil.instPlusMicros(sTrades.getTradeDateTime(), sTrades.getTradeTimeMs()));
eCurrency.setExchangeExecutionMicroseconds(sTrades.getTradeDateTime()); // todo проверить это дата+время или нет
eCurrency.setPartyTradingClearingRegistryId(rgstr.getId()); // setTradingClearingRegistryId
eCurrency.setPartyTradingClearingRegistry(rgstr.getCode()); // todo уточнить rgstr.code/strades.account?

View file

@ -183,7 +183,7 @@ public class ExecutionDepositComponent {
eDeposit.setTradingDate(sTrades.getTradeDate());
eDeposit.setClearingDate(TimeUtil.toLocalDate(now));
eDeposit.setExchangeExecutionId(sTrades.getTradeNum());
eDeposit.setExchangeExecutionTime(sTrades.getTradeDateTime());
eDeposit.setExchangeExecutionTime(TimeUtil.instPlusMicros(sTrades.getTradeDateTime(), sTrades.getTradeTimeMs()));
eDeposit.setTradingClearingRegistryId(rgstr.getId());
eDeposit.setMarket(sTrades.getClassCode());
eDeposit.setPrice(sTrades.getPrice());

View file

@ -204,7 +204,7 @@ public class ExecutionFondComponent {
eFond.setTradingDate(sTrades.getTradeDate());
eFond.setClearingDate(TimeUtil.toLocalDate(now));
eFond.setExchangeExecutionId(sTrades.getTradeNum());
eFond.setExchangeExecutionTime(sTrades.getTradeDateTime());
eFond.setExchangeExecutionTime(TimeUtil.instPlusMicros(sTrades.getTradeDateTime(), sTrades.getTradeTimeMs()));
eFond.setTradingClearingRegistryId(rgstr.getId());
//todo можем ли просто переложить sTrades.getClassCode() или все такие искать, одно и тоже же

View file

@ -38,6 +38,7 @@ public class SessionManager {
private final IntermediateMkrSession intermediateMkrSession;
private final FinalMkrSession finalMkrSession;
private final ReturnDepositSession returnDepositSession;
private final UnitedSession unitedSession;
private final TradingTimeService time;
public SessionManager(ImdgProvider imdgProvider,
@ -45,7 +46,7 @@ public class SessionManager {
PrimaryAuctionBnSession primaryAuctionBnSession,
PrimaryAuctionB0Session primaryAuctionB0Session,
SecondaryAuctionT0Session secondaryAuctionT0Session, CurrencySession currencySession,
IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, TradingTimeService time) {
IntermediateMkrSession intermediateMkrSession, FinalMkrSession finalMkrSession, ReturnDepositSession returnDepositSession, UnitedSession unitedSession, TradingTimeService time) {
this.notification = notification;
this.msgs = msgs;
this.primaryAuctionT0Session = primaryAuctionT0Session;
@ -58,6 +59,7 @@ public class SessionManager {
this.returnDepositSession = returnDepositSession;
sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
this.unitedSession = unitedSession;
this.time = time;
}
@ -118,6 +120,7 @@ public class SessionManager {
case MEDM -> session = intermediateMkrSession;
case FINL -> session = finalMkrSession;
case XDEP -> session = returnDepositSession;
case UNIT -> session = unitedSession;
}
return session;
}

View file

@ -9,6 +9,7 @@ public enum TaskType implements IEnumKey {
StartRevise("CLR0"),
ContinueRevise("CLR 0_0"),
StartRevisePart1("С0_1"),
SDF51("SDF51"),
/**
* step 1

View file

@ -0,0 +1,331 @@
package ru.spcex.clearing.session.stage;
import java.time.Instant;
import java.time.LocalDate;
import java.util.List;
import java.util.function.Supplier;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.ObjectFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.execution.ExecutionFond;
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.BalanceRevise;
import ru.spcex.clearing.session.stage.impl.DealsPrepare;
import ru.spcex.clearing.session.stage.impl.EndStageNotification;
import ru.spcex.clearing.session.stage.impl.FinishingSession;
import ru.spcex.clearing.session.stage.impl.FormingPaymentInstructionAssets;
import ru.spcex.clearing.session.stage.impl.FormingRegistersOnOS;
import ru.spcex.clearing.session.stage.impl.InclusionObligations;
import ru.spcex.clearing.session.stage.impl.InspectionObligations;
import ru.spcex.clearing.session.stage.impl.InspectionObligationsDepositReturn;
import ru.spcex.clearing.session.stage.impl.ObligationAdmission;
import ru.spcex.clearing.session.stage.impl.RequirementsAndObligationCreationCompound;
import ru.spcex.clearing.session.stage.impl.UnlockResources;
import ru.spcex.clearing.session.stage.impl.compound.CompoundStageDealsPrepare;
import ru.spcex.clearing.session.stage.monitor.SessionMonitor;
import ru.spcex.clearing.session.stage.monitor.SessionMonitorFactory;
import ru.spcex.clearing.session.stage.task.DealsPreparePayload;
import ru.spcex.clearing.session.stage.task.EndStageNotificationPayload;
import ru.spcex.clearing.session.stage.task.FinishingSessionPayload;
import ru.spcex.clearing.session.stage.task.FormingRegistersOnOSPayload;
import ru.spcex.clearing.session.stage.task.InclusionToPoolPayload;
import ru.spcex.clearing.session.stage.task.InspectionPoolPayload;
import ru.spcex.clearing.session.stage.task.RequirementsAndObligationCreationCompoundPayload;
import ru.spcex.clearing.session.stage.task.result.DealsPrepareCompoundResult;
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
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.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
@Service
//todo эта сессия
public class UnitedSession extends AbstractSession implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final BalanceRevise balanceRevise;
private final DealsPrepare dealsPrepareCurrency;
private final DealsPrepare dealsPrepareTRDT;
private final DealsPrepare dealsPrepareFinal;
private final RequirementsAndObligationCreationCompound requirementsAndObligationCreation;
private final ObligationAdmission obligationsAdmission;
private final InclusionObligations inclusionObligations;
private final InspectionObligationsDepositReturn inspectionObligationsReturn;
private final InspectionObligations inspectionObligations;
private final FormingRegistersOnOS formingRegistersOnOS;
private final FormingPaymentInstructionAssets formingPaymentInstructionAssets;
private final UnlockResources unlockResources;
private final FinishingSession finishingSession;
private final EndStageNotification endStageNotification;
private final Imdg<ExecutionFond> executionFondImdg;
private final Supplier<List<String>> marketCodesTrdt;
private final Supplier<List<String>> marketCodesCurr;
private SessionMonitor firstReviseMonitor;
private SessionMonitor afterPaymentsSdf4And13Monitor;
private SessionMonitor afterPaymentsReviseMonitor;
private SessionMonitor afterReviseErrorMonitor;
public UnitedSession(
ImdgProvider imdgProvider,
BalanceRevise balanceRevise,
ObjectFactory<DealsPrepare> deals,
RequirementsAndObligationCreationCompound requirementsAndObligationCreation,
ObligationAdmission obligationsAdmission,
InclusionObligations inclusionObligations, InspectionObligationsDepositReturn inspectionObligationsReturn,
FormingRegistersOnOS formingRegistersOnOS,
FormingPaymentInstructionAssets formingPaymentInstructionAssets,
UnlockResources unlockResources,
FinishingSession finishingSession,
EndStageNotification endStageNotification,
IMessageResolver messageResolver,
InspectionObligations inspectionObligations,
@Qualifier("marketCodesForCurr") Supplier<List<String>> marketCodesCurr,
@Qualifier("marketCodesForT0") Supplier<List<String>> marketCodesTrdt) {
super(imdgProvider, messageResolver);
this.balanceRevise = balanceRevise;
this.dealsPrepareCurrency = deals.getObject();
this.dealsPrepareTRDT = deals.getObject();
this.dealsPrepareFinal = deals.getObject();
this.requirementsAndObligationCreation = requirementsAndObligationCreation;
this.obligationsAdmission = obligationsAdmission;
this.inclusionObligations = inclusionObligations;
this.inspectionObligationsReturn = inspectionObligationsReturn;
this.formingRegistersOnOS = formingRegistersOnOS;
this.formingPaymentInstructionAssets = formingPaymentInstructionAssets;
this.unlockResources = unlockResources;
this.finishingSession = finishingSession;
this.endStageNotification = endStageNotification;
this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
this.inspectionObligations = inspectionObligations;
this.marketCodesTrdt = marketCodesTrdt;
this.marketCodesCurr = marketCodesCurr;
}
@Override
public void afterPropertiesSet() throws Exception {
{
dealsPrepareCurrency.searchForExecutions(ExecutionType.ExecutionCurrency);
ImdgPredicateBuilder pb = executionFondImdg.predicateBuilder();
dealsPrepareCurrency.addExecutionFondCondition(pb.regex("settlementCode", "^T0.*$"));
dealsPrepareCurrency.addExecutionFondCondition(pb.in("market", marketCodesCurr.get().toArray(new String[0])));
dealsPrepareTRDT.searchForExecutions(ExecutionType.ExecutionFond);
ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder();
dealsPrepareTRDT.addExecutionFondCondition(execFondPb.regex("settlementCode", "^T0.*$"));
dealsPrepareTRDT.addExecutionFondCondition(execFondPb.in("market", marketCodesTrdt.get().toArray(new String[0])));
dealsPrepareFinal.searchForExecutions(ExecutionType.ExecutionDeposit);
dealsPrepareFinal.addExecutionDepositCondition(pb.regex("firstLegSettlementCode", "^T0.*$"));
}
//1 шаг deals prepare вызываем 3 раза с разными параметрами объединяемых сессий
inclusionObligations.setSessionType(sessionType());
inspectionObligationsReturn.setSessionType(sessionType());
inspectionObligations.setSection(section());
inspectionObligations.setSessionType(sessionType());
finishingSession.setSection(section());
finishingSession.setSessionType(sessionType());
imdgProvider.waitAvailable();
initSessionIfPresent();
}
@Override
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");
this.firstReviseMonitor = SessionMonitorFactory.waitRevise();
}
}
public void continueSession(BaseRequest<?> req) {
try {
if (!isRunning()) {
return;
}
log.info("session is running, stage {}, monitors: {}",
currStage.get(),
logMonitors(firstReviseMonitor, afterPaymentsSdf4And13Monitor, afterPaymentsReviseMonitor));
if (firstReviseMonitor != null && isMonitorPassed(firstReviseMonitor, req.getRequestPayload())) {
firstReviseMonitor = null;
firstPart();
return;
}
if (afterPaymentsSdf4And13Monitor != null && isMonitorPassed(afterPaymentsSdf4And13Monitor, req.getRequestPayload())) {
afterPaymentsSdf4And13Monitor = null;
checkStageAndThrow(TaskType.FormingPaymentInstruction);
finishPart();
return;
}
if (afterReviseErrorMonitor != null && isMonitorPassed(afterReviseErrorMonitor, req.getRequestPayload())) {
afterReviseErrorMonitor = null;
finishPart();
return;
}
//if (afterPaymentsReviseMonitor != null && isMonitorPassed(afterPaymentsReviseMonitor, req.getRequestPayload())) {
// afterPaymentsReviseMonitor = null;
// finishPart();
//}
} catch (StageException e) {
//already logged
}
}
private void firstPart() {
try {
if (!checkStage(TaskType.StartRevise)) {
log.error("cannot continue session, current stage is {}", currStage.get());
throw new StageException();
}
runStage(TaskType.StartRevisePart1, currSession.getId(), balanceRevise);
//stage 1
StageResult<DealsPrepareCompoundResult> dealsPreparationResult;
{
DealsPreparePayload payload = new DealsPreparePayload();
payload.setSessionId(currSession.getId());
dealsPreparationResult = runStage(TaskType.DealsPrepare, payload, new CompoundStageDealsPrepare(
dealsPrepareTRDT, dealsPrepareCurrency, dealsPrepareFinal
));
}
//stage 2
{
RequirementsAndObligationCreationCompoundPayload p;
p = RequirementsAndObligationCreationCompoundPayload.create(dealsPreparationResult.getStageResult());
runStage(TaskType.RequirementsAndObligationsCreate, p, requirementsAndObligationCreation);
}
//stage 3
runStage(TaskType.ObligationsAdmission, currSession.getId(), obligationsAdmission);
//stage 4
{
InclusionToPoolPayload inclusionToPoolPayload = new InclusionToPoolPayload();
inclusionToPoolPayload.setSessionType(currSession.getSessionType());
inclusionToPoolPayload.setSessionId(currSession.getId());
runStage(TaskType.InclusionToPool, inclusionToPoolPayload, inclusionObligations);
}
{
InspectionPoolPayload companyIdPayload = new InspectionPoolPayload();
companyIdPayload.setSessionId(currSession.getId());
companyIdPayload.setProcessedCompanyId(currSession.getCompanyId());
//fixme checks inside the stages
runStage(TaskType.InspectionObligations, companyIdPayload, inspectionObligationsReturn);
runStage(TaskType.InspectionObligations, companyIdPayload, inspectionObligations);
}
//stage 6
{
FormingRegistersOnOSPayload payload = new FormingRegistersOnOSPayload();
payload.setSessionId(currSession.getId());
runStage(TaskType.FormingRegistersOnOS, payload, formingRegistersOnOS); //returns Collection<Registry>
}
//stage 7
//fixmeStageResult<PaymentInfo> paymentResult = null;
//fixme {
//fixme FormingPaymentInstructionPayload payload = new FormingPaymentInstructionPayload();
//fixme payload.setSessionId(currSession.getId());
//fixme paymentResult = runStage(TaskType.FormingPaymentInstruction, payload, formingPaymentInstructionAssets);
//fixme }
//fixme if (paymentResult.getStageResult().getPaymentInstructions().isEmpty()) {
//fixme log.info("no payment instructions were created");
//fixme finishPart();
//fixme } else {
//fixme this.afterPaymentsSdf4And13Monitor = SessionMonitorFactory.waitStep7(section());
//fixme log.info("created {} PaymentInstructions, waiting for {}",
//fixme paymentResult.getStageResult().getPaymentInstructions().size(),
//fixme this.afterPaymentsSdf4And13Monitor.allConditions());
//fixme }
} catch (StageException e) {
//already logged
}
}
public void finishPart() {
try {
//stage 9 continue revision
{
StageResult<Object> reviseRes = runStage(TaskType.AgainRevise, currSession.getId(), balanceRevise, false);
if (!reviseRes.isSuccess()) {
this.afterReviseErrorMonitor = SessionMonitorFactory.waitReviseError(section());
log.warn("{} stage error, created monitor for {}",
TaskType.AgainRevise,
afterReviseErrorMonitor.allConditions());
return;
}
}
//stage 10
{
FinishingSessionPayload payload = new FinishingSessionPayload();
payload.setSessionId(currSession.getId());
payload.setPr("1");
runStage(TaskType.FinishingSession, payload, finishingSession);
}
//stage 11
{
EndStageNotificationPayload payload = new EndStageNotificationPayload();
payload.setSection(currSession.getSection());
payload.setSessionId(currSession.getId());
runStage(TaskType.EndStageNotification, payload, endStageNotification);
endSession();
}
} catch (StageException e) {
//already logged
}
}
private void sendSdf56() {
runStage(TaskType.FormingPaymentInstruction, balanceRevise);
this.afterPaymentsReviseMonitor = SessionMonitorFactory.waitRevise();
}
private boolean startSession() {
synchronized (this.currStage) {
if (this.currStage.get() != null) {
log.info("already running session.id={}", this.currSession.getId());
return false;
} else {
TaskType startStatus = TaskType.StartRevise;
Session newSession = new Session();
newSession.setSection(section().getKey());
newSession.setSessionType(sessionType().getKey());
newSession.setSessionStatus(startStatus.getKey());
newSession.setWorkflowStatus(SessionStatus.ACTV.getKey());
newSession.setClearingDate(LocalDate.now());
newSession.setCreated(Instant.now());
//todo companyId/securityId/userId передается из сообщения очереди
sessionImdg.insert(newSession);
currSession = newSession;
log.info("started new session.id={}", this.currSession.getId());
currStage.set(TaskType.StartRevise);
return true;
}
}
}
@Override
protected Section section() {
return Section.FOND;
}
@Override
protected SessionType sessionType() {
return SessionType.UNIT;
}
}

View file

@ -77,6 +77,9 @@ public class BalanceRevise implements ISessionStage {
@Override
public StageResult<?> submit(Task<?> task) {
switch (task.getTaskType()) {
case SDF51 -> {
return sendSdfs51();
}
case StartRevise, FormingPaymentInstruction -> {
LocalTime fromTime = null;
if (task.getData() instanceof LocalTime) {

View file

@ -123,7 +123,7 @@ public class InspectionObligations implements ISessionStage {
boolean isUncovered = false;
List<CheckResult> checkResults = new ArrayList<>();
for (Registry obligation : obligationsInGroup) {
if ((SessionType.FINL.equals(sessionType) || SessionType.MEDM.equals(sessionType)) && RegistryInstrumentType.S.equalsByKey(obligation.getRegistryInstrumentType())) {
if ((SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType) || SessionType.MEDM.equals(sessionType)) && RegistryInstrumentType.S.equalsByKey(obligation.getRegistryInstrumentType())) {
checkResults.add(new CheckResult(obligation, false));
continue;
}

View file

@ -1,5 +1,12 @@
package ru.spcex.clearing.session.stage.impl;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.function.BiConsumer;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
@ -18,21 +25,19 @@ 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.InspectionPoolPayload;
import ru.spcex.platform.enumeration.*;
import ru.spcex.platform.enumeration.ClearingCategory;
import ru.spcex.platform.enumeration.RegistryCapacity;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.RegistryTradingParams;
import static ru.spcex.platform.enumeration.RegistryTradingParams.*;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.enumeration.SessionType;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import java.math.BigDecimal;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.function.BiConsumer;
import java.util.stream.Collectors;
import static ru.spcex.platform.enumeration.RegistryTradingParams.*;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
@ -205,7 +210,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage {
}
}
if (SessionType.FINL.equals(sessionType)
if ((SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType))
&& amfBalance.compareTo(omtBalance) >= 0
&& initiatorV(omtRgs) && gatewayIfNeeded(omtRgs)) {
log.debug("OM*T.id={} -> no DM*X/DM*T(INFO/CLRN) registry found. " +
@ -262,7 +267,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage {
}
private RegistryStatus registryStatusFailed() {
if (SessionType.FINL.equals(sessionType)) {
if (SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType)) {
return RegistryStatus.FAIL;
} else {
return RegistryStatus.MNG;
@ -270,7 +275,7 @@ public class InspectionObligationsDepositReturn implements ISessionStage {
}
private RegistryStatus registryStatusFailed(Registry registry) {
if (SessionType.FINL.equals(sessionType)) {
if (SessionType.FINL.equals(sessionType) || SessionType.UNIT.equals(sessionType)) {
if (equalByRgs(OM_T, registry)) {
return RegistryStatus.UNCV;
} else {

View file

@ -33,6 +33,7 @@ import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
@ -195,6 +196,7 @@ public class RegistryCurrBuilder implements IRegistryBuilder {
reg.setGroupId(groupId());
reg.setSessionId(exec.getSessionId());
reg.setSessionType(sessionType());
reg.setSection(Section.CURR.getKey());
return reg;
}

View file

@ -1,5 +1,8 @@
package ru.spcex.clearing.session.stage.impl;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Map;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClearingAccount;
import ru.clearing.classes.statics.data.account.DepoAccount;
@ -15,14 +18,20 @@ 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.enumeration.AccountType;
import ru.spcex.platform.enumeration.BalanceDimension;
import ru.spcex.platform.enumeration.ISide;
import ru.spcex.platform.enumeration.MoneyFlowSide;
import ru.spcex.platform.enumeration.RegistryCapacity;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Map;
public class RegistryDepoBuilder implements IRegistryBuilder {
private Imdg<Company> companyImdg;
private Imdg<TradingClearingRegistry> tradingClearingRegistryImdg;
@ -148,6 +157,7 @@ public class RegistryDepoBuilder implements IRegistryBuilder {
reg.setGroupId(groupId());
reg.setSessionId(exec.getSessionId());
reg.setSessionType(sessionType());
reg.setSection(Section.MKR.getKey());
return reg;
}

View file

@ -1,5 +1,10 @@
package ru.spcex.clearing.session.stage.impl;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Map;
import ru.clearing.classes.statics.data.account.Account;
import ru.clearing.classes.statics.data.account.ClearingAccount;
import ru.clearing.classes.statics.data.account.DepoAccount;
@ -16,17 +21,21 @@ import ru.clearing.classes.statics.data.security.Security;
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.enumeration.AccountType;
import ru.spcex.platform.enumeration.BalanceDimension;
import ru.spcex.platform.enumeration.ISide;
import ru.spcex.platform.enumeration.InstrumentType;
import ru.spcex.platform.enumeration.MoneyFlowSide;
import ru.spcex.platform.enumeration.RegistryCapacity;
import ru.spcex.platform.enumeration.RegistryDesignation;
import ru.spcex.platform.enumeration.RegistryInstrumentType;
import ru.spcex.platform.enumeration.RegistryStatus;
import ru.spcex.platform.enumeration.RegistryUnit;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.predicate.specific.SecuritySelector;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.Instant;
import java.time.LocalDate;
import java.util.Map;
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
public class RegistryFondBuilder implements IRegistryBuilder {
@ -162,6 +171,7 @@ public class RegistryFondBuilder implements IRegistryBuilder {
reg.setGroupId(groupId());
reg.setSessionId(exec.getSessionId());
reg.setSessionType(sessionType());
reg.setSection(Section.FOND.getKey());
return reg;
}

View file

@ -65,7 +65,7 @@ public class RequirementsAndObligationCreation implements ISessionStage {
throw new IllegalStateException("Unknown task type: " + task.getTaskType());
}
private StageResult<?> createRegisters(List<ExecutionCommon> data) {
protected StageResult<?> createRegisters(List<ExecutionCommon> data) {
data.sort(Comparator.comparing(ExecutionCommon::getExchangeExecutionId));
//см. описание к #matchExecutions
for (int i = 0; i < data.size(); ) {
@ -219,7 +219,7 @@ public class RequirementsAndObligationCreation implements ISessionStage {
* если это не так, передалать на коллекцию Long executionId в правильном порядке
* и для каждого искать мэтч отдельно в Imdg, с сохранением уже обработанных для избежания дублирования
*/
private Pair<ExecutionCommon, ExecutionCommon> matchExecutions(List<ExecutionCommon> data, int i) {
protected Pair<ExecutionCommon, ExecutionCommon> matchExecutions(List<ExecutionCommon> data, int i) {
//if the last execution, then no match
if (i >= data.size() - 1) {
return null;

View file

@ -0,0 +1,101 @@
package ru.spcex.clearing.session.stage.impl;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Comparator;
import java.util.List;
import java.util.Objects;
import java.util.stream.Stream;
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.execution.ExecutionCommon;
import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task;
import ru.spcex.clearing.session.stage.TaskType;
import ru.spcex.clearing.session.stage.task.RequirementsAndObligationCreationCompoundPayload;
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.utils.collection.Pair;
/**
* LiabilitiesAndClaims
*/
@Service
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class RequirementsAndObligationCreationCompound extends RequirementsAndObligationCreation {
private final Logger log = LoggerFactory.getLogger(getClass());
private static final Comparator<ExecutionCommon> tradeTimeComparator
= Comparator.comparing(ExecutionCommon::getExchangeExecutionTime);
private static final Comparator<ExecutionCommon> execIdComparator
= Comparator.comparing(ExecutionCommon::getExchangeExecutionId);
@Autowired
public RequirementsAndObligationCreationCompound(ImdgProvider imdgProvider) {
super(imdgProvider);
}
@SuppressWarnings("unchecked")
@Override
public StageResult<?> submit(Task<?> task) {
if (task.getTaskType() == TaskType.RequirementsAndObligationsCreate) {
return mergeReqsAndClaims((RequirementsAndObligationCreationCompoundPayload) task.getData());
}
throw new IllegalStateException("Unknown task type: " + task.getTaskType());
}
private StageResult<?> mergeReqsAndClaims(RequirementsAndObligationCreationCompoundPayload data) {
List<ExecutionCommon> tmpList = new ArrayList<>(Arrays.asList(new ExecutionCommon[2]));
List<ExecutionCommon> executionFond = data.getExecutionsTrdt();
List<ExecutionCommon> executionCurrency = data.getExecutionsCurrency();
List<ExecutionCommon> executionDeposit = data.getExecutionsFinal();
Stream.of(executionFond, executionCurrency, executionDeposit).forEach(c -> c.sort(execIdComparator));
int efs = executionFond.size();
int eds = executionDeposit.size();
int ecs = executionCurrency.size();
int[] ief= {0}, ied= {0}, iec= {0};
while (ief[0] < efs || ied[0] < eds || iec[0] < ecs) {
Pair<ExecutionCommon, ExecutionCommon> pair = getMinTimeExecutions(
safeMatch(executionFond, ief),
safeMatch(executionDeposit, ied),
safeMatch(executionCurrency, iec)
);
if (pair == null) continue;
if (pair.getFirst().type().equals(ExecutionType.ExecutionFond)) ief[0] += 2;
if (pair.getFirst().type().equals(ExecutionType.ExecutionDeposit)) ied[0] += 2;
if (pair.getFirst().type().equals(ExecutionType.ExecutionCurrency)) iec[0] += 2;
pair.map(f -> tmpList.set(0, f), s -> tmpList.set(1, s));
createRegisters(tmpList);
}
return new StageResult<>(null, true);
}
private Pair<ExecutionCommon, ExecutionCommon> safeMatch(List<ExecutionCommon> l, int[] i) {
Pair<ExecutionCommon, ExecutionCommon> pair = matchExecutions(l, i[0]);
if (pair == null) {
i[0]++;
return null;
}
// i[0] += 2;
return pair;
}
private Pair<ExecutionCommon, ExecutionCommon> getMinTimeExecutions(
Pair<ExecutionCommon, ExecutionCommon> fond,
Pair<ExecutionCommon, ExecutionCommon> deposit,
Pair<ExecutionCommon, ExecutionCommon> currency) {
return Stream.of(fond, deposit, currency)
.filter(Objects::nonNull)
.max((p1, p2) -> tradeTimeComparator.compare(p1.getFirst(), p2.getFirst())).orElse(null);
}
}

View file

@ -0,0 +1,34 @@
package ru.spcex.clearing.session.stage.impl.compound;
import java.util.ArrayList;
import java.util.List;
import ru.spcex.clearing.session.stage.ISessionStage;
import ru.spcex.clearing.session.stage.StageResult;
import ru.spcex.clearing.session.stage.Task;
public abstract class CompoundStage<T> implements ISessionStage {
private final List<ISessionStage> stages;
public CompoundStage(ISessionStage... stages) {
this.stages = new ArrayList<>();
this.stages.addAll(List.of(stages));
}
protected abstract Object mapper(List<T> results);
@Override
public StageResult<?> submit(Task<?> task) {
List<Object> innerResults = new ArrayList<>();
for (ISessionStage stage : stages) {
StageResult<?> submit = stage.submit(task);
if (!submit.isSuccess()) {
return submit;
}
innerResults.add(submit.getStageResult());
}
StageResult<Object> result = new StageResult<>(null, true);
result.setStageResult(mapper((List<T>) innerResults));
return result;
}
}

View file

@ -0,0 +1,22 @@
package ru.spcex.clearing.session.stage.impl.compound;
import java.util.List;
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
import ru.spcex.clearing.session.stage.impl.DealsPrepare;
import ru.spcex.clearing.session.stage.task.result.DealsPrepareCompoundResult;
public class CompoundStageDealsPrepare extends CompoundStage<List<ExecutionCommon>> {
public CompoundStageDealsPrepare(DealsPrepare... stages) {
super(stages);
}
@Override
protected Object mapper(List<List<ExecutionCommon>> results) {
DealsPrepareCompoundResult res = new DealsPrepareCompoundResult();
res.setExecutionsTrdt(results.get(0));
res.setExecutionsCurrency(results.get(1));
res.setExecutionsFinal(results.get(2));
return res;
}
}

View file

@ -0,0 +1,43 @@
package ru.spcex.clearing.session.stage.task;
import java.util.List;
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
import ru.spcex.clearing.session.stage.task.result.DealsPrepareCompoundResult;
public class RequirementsAndObligationCreationCompoundPayload {
private List<ExecutionCommon> executionsTrdt;
private List<ExecutionCommon> executionsCurrency;
private List<ExecutionCommon> executionsFinal;
public static RequirementsAndObligationCreationCompoundPayload create(DealsPrepareCompoundResult stageResult) {
RequirementsAndObligationCreationCompoundPayload p = new RequirementsAndObligationCreationCompoundPayload();
p.setExecutionsTrdt(stageResult.getExecutionsTrdt());
p.setExecutionsCurrency(stageResult.getExecutionsCurrency());
p.setExecutionsFinal(stageResult.getExecutionsFinal());
return p;
}
public List<ExecutionCommon> getExecutionsTrdt() {
return executionsTrdt;
}
public void setExecutionsTrdt(List<ExecutionCommon> executionsTrdt) {
this.executionsTrdt = executionsTrdt;
}
public List<ExecutionCommon> getExecutionsCurrency() {
return executionsCurrency;
}
public void setExecutionsCurrency(List<ExecutionCommon> executionsCurrency) {
this.executionsCurrency = executionsCurrency;
}
public List<ExecutionCommon> getExecutionsFinal() {
return executionsFinal;
}
public void setExecutionsFinal(List<ExecutionCommon> executionsFinal) {
this.executionsFinal = executionsFinal;
}
}

View file

@ -0,0 +1,34 @@
package ru.spcex.clearing.session.stage.task.result;
import java.util.List;
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
public class DealsPrepareCompoundResult {
private List<ExecutionCommon> executionsTrdt;
private List<ExecutionCommon> executionsCurrency;
private List<ExecutionCommon> executionsFinal;
public List<ExecutionCommon> getExecutionsTrdt() {
return executionsTrdt;
}
public void setExecutionsTrdt(List<ExecutionCommon> executionsTrdt) {
this.executionsTrdt = executionsTrdt;
}
public List<ExecutionCommon> getExecutionsCurrency() {
return executionsCurrency;
}
public void setExecutionsCurrency(List<ExecutionCommon> executionsCurrency) {
this.executionsCurrency = executionsCurrency;
}
public List<ExecutionCommon> getExecutionsFinal() {
return executionsFinal;
}
public void setExecutionsFinal(List<ExecutionCommon> executionsFinal) {
this.executionsFinal = executionsFinal;
}
}

View file

@ -5,7 +5,7 @@ import ru.spcex.platform.utils.enumeration.IEnumKey;
public enum SessionType implements IEnumKey {
IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"),
MEDM("MEDM"), FINL("FINL"), XDEP("XDEP"), LIQU("LIQU"),
CURR("CURR")
CURR("CURR"), UNIT("UNIT"),
;
SessionType(String key) {