http://jira.mfd.msk:8088/browse/CLS-1060 IPOB + inclusion to pool action V2
This commit is contained in:
parent
6f416f6eda
commit
b8a4b114fa
4 changed files with 283 additions and 7 deletions
|
|
@ -27,7 +27,7 @@ import ru.spcex.clearing.session.state.action.DealsPrepareAction;
|
|||
import ru.spcex.clearing.session.state.action.DiscardRegistriesAction;
|
||||
import ru.spcex.clearing.session.state.action.FinishingSessionAction;
|
||||
import ru.spcex.clearing.session.state.action.FormingPaymentInstructionSecuritiesAction;
|
||||
import ru.spcex.clearing.session.state.action.InclusionToPoolAction;
|
||||
import ru.spcex.clearing.session.state.action.InclusionToPoolActionV2;
|
||||
import ru.spcex.clearing.session.state.action.InspectionObligationsV2Action;
|
||||
import ru.spcex.clearing.session.state.action.ObligationAdmissionAction;
|
||||
import ru.spcex.clearing.session.state.action.RequirementAndObligationCreationAction;
|
||||
|
|
@ -46,6 +46,7 @@ import ru.spcex.clearing.session.state.guard.StashedRegistriesPresentGuard;
|
|||
import ru.spcex.clearing.session.state.listener.MachineMonitoringListener;
|
||||
import static ru.spcex.clearing.util.StateMachineUtil.*;
|
||||
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||
import ru.spcex.platform.enumeration.RegistryStatus;
|
||||
import ru.spcex.platform.enumeration.Section;
|
||||
import ru.spcex.platform.enumeration.SessionType;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
|
|
@ -63,7 +64,7 @@ public class PrimaryAuctionBnSessionStateMachineConfig
|
|||
private final DealsPrepareAction dealsPrepareAction;
|
||||
private final RequirementAndObligationCreationAction reqAndOblAction;
|
||||
private final ObligationAdmissionAction obligationAdmissionAction;
|
||||
private final InclusionToPoolAction inclusionToPoolAction;
|
||||
private final InclusionToPoolActionV2 inclusionToPoolAction;
|
||||
private final InspectionObligationsV2Action inspectionObligationsV2Action;
|
||||
private final FormingPaymentInstructionSecuritiesAction formingPaymentInstructionSecuritiesAction;
|
||||
private final FinishingSessionAction finishingSessionAction;
|
||||
|
|
@ -96,7 +97,7 @@ public class PrimaryAuctionBnSessionStateMachineConfig
|
|||
@Qualifier("requirementAndObligationCreationAction")
|
||||
RequirementAndObligationCreationAction reqAndOblAction,
|
||||
ObligationAdmissionAction obligationAdmissionAction,
|
||||
InclusionToPoolAction inclusionToPoolAction,
|
||||
InclusionToPoolActionV2 inclusionToPoolAction,
|
||||
InspectionObligationsV2Action inspectionObligationsV2Action,
|
||||
FormingPaymentInstructionSecuritiesAction formingPaymentInstructionSecuritiesAction,
|
||||
FinishingSessionAction finishingSessionAction,
|
||||
|
|
@ -129,10 +130,18 @@ public class PrimaryAuctionBnSessionStateMachineConfig
|
|||
|
||||
//stages settings:
|
||||
this.dealsPrepareAction.searchForExecutions(ExecutionType.ExecutionFond);
|
||||
ImdgPredicateBuilder execFondPb = executionFondImdg.predicateBuilder();
|
||||
dealsPrepareAction.addExecutionFondCondition(execFondPb.regex("settlementCode", "^B\\d{2}$"));
|
||||
dealsPrepareAction.addExecutionFondCondition(execFondPb.in("market", marketCodes.get().toArray(new String[0])));
|
||||
ImdgPredicateBuilder pb = executionFondImdg.predicateBuilder();
|
||||
dealsPrepareAction.addExecutionFondCondition(pb.regex("settlementCode", "^B\\d{2}$"));
|
||||
dealsPrepareAction.addExecutionFondCondition(pb.in("market", marketCodes.get().toArray(new String[0])));
|
||||
this.obligationAdmissionContinueAction.setDataEnum(DataEnum.obligationAdmissionStashedRgs);
|
||||
|
||||
this.inclusionToPoolAction.addRgsConditionSessionTypeIn(
|
||||
SessionType.IPOB, SessionType.IPOM
|
||||
);
|
||||
this.inclusionToPoolAction.addRgsConditionStatusIn(
|
||||
RegistryStatus.PROC, RegistryStatus.MNG
|
||||
);
|
||||
this.inclusionToPoolAction.setExecType(ExecutionType.ExecutionFond);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -49,6 +49,22 @@ import ru.spcex.platform.utils.enumeration.IEnumKey;
|
|||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||
import ru.spcex.platform.utils.text.TextUtil;
|
||||
|
||||
/**
|
||||
* при изменениях в этом шаге необходимо перевести сессию на <code>{@link InclusionToPoolActionV2}</code> <br/>
|
||||
* все настроечные методы для "повторения" логики из этого кода сделаны в новом классе. <br/>
|
||||
* (можно добавить setExecTypeBySection для закрытия некоторых кейсов
|
||||
* определения imdg execution'ов из updateExecutions)
|
||||
* @see InclusionToPoolActionV2#addRegistryCondition
|
||||
* @see InclusionToPoolActionV2#addRgsConditionSessionTypeIn
|
||||
* @see InclusionToPoolActionV2#addRgsConditionSessionTypeEquals
|
||||
* @see InclusionToPoolActionV2#addRgsConditionDynamic
|
||||
* @see InclusionToPoolActionV2#addRgsConditionStatusIn
|
||||
* @see InclusionToPoolActionV2#addRgsConditionStatusEquals
|
||||
* @see InclusionToPoolActionV2#skipExecutionLoad
|
||||
* @see InclusionToPoolActionV2#setExecType
|
||||
*
|
||||
*/
|
||||
@Deprecated
|
||||
@Service
|
||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||
public class InclusionToPoolAction extends AbstractSessionActionForOkErrorHandling {
|
||||
|
|
|
|||
|
|
@ -0,0 +1,251 @@
|
|||
package ru.spcex.clearing.session.state.action;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.Comparator;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
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.statemachine.ExtendedState;
|
||||
import org.springframework.statemachine.StateContext;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.util.StringUtils;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionCommon;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionCurrency;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionDeposit;
|
||||
import ru.clearing.classes.statics.data.execution.ExecutionFond;
|
||||
import ru.clearing.classes.statics.data.registry.Registry;
|
||||
import ru.spcex.clearing.component.GroupRgsKey;
|
||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.session.stage.TaskType;
|
||||
import ru.spcex.clearing.session.state.DataEnum;
|
||||
import ru.spcex.clearing.session.state.SsnEvent;
|
||||
import ru.spcex.clearing.util.ClearingUtil;
|
||||
import ru.spcex.platform.classes.base.interfaces.ExecutionType;
|
||||
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.SessionType;
|
||||
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.utils.text.TextUtil;
|
||||
|
||||
/**
|
||||
* рефактор старого депрекейтед класса <code>{@link InclusionToPoolAction}</code>
|
||||
*/
|
||||
@Service
|
||||
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||
public class InclusionToPoolActionV2 extends AbstractSessionActionForOkErrorHandling {
|
||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||
private final Imdg<Registry> registryImdg;
|
||||
private final Imdg<ExecutionDeposit> executionDepositImdg;
|
||||
private final Imdg<ExecutionFond> executionFondImdg;
|
||||
private final Imdg<ExecutionCurrency> executionCurrImdg;
|
||||
private SessionType sessionType;
|
||||
private final ImdgPredicateBuilder rgsPb;
|
||||
|
||||
private boolean loadExecs = true;
|
||||
private ExecutionType execType; //
|
||||
private final List<ImdgPredicate> registryConditionsAndV2 = new ArrayList<>();
|
||||
private final List<Function<ExtendedState, ImdgPredicate>> dynamicConditions = new ArrayList<>();
|
||||
|
||||
@Autowired
|
||||
public InclusionToPoolActionV2(ImdgProvider imdgProvider) {
|
||||
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
||||
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
|
||||
this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
|
||||
this.executionCurrImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class);
|
||||
this.rgsPb = registryImdg.predicateBuilder();
|
||||
}
|
||||
|
||||
public InclusionToPoolActionV2 addRegistryCondition(ImdgPredicate imdgPredicate) {
|
||||
this.registryConditionsAndV2.add(imdgPredicate);
|
||||
return this;
|
||||
}
|
||||
|
||||
public void addRgsConditionSessionTypeIn(SessionType... sessionTypes) {
|
||||
ImdgPredicate sessionTypePredicate = rgsPb.or(
|
||||
Arrays.stream(sessionTypes)
|
||||
.map(sessionType -> rgsPb.equals("sessionType", sessionType.getKey()))
|
||||
.toArray(ImdgPredicate[]::new)
|
||||
);
|
||||
this.registryConditionsAndV2.add(sessionTypePredicate);
|
||||
}
|
||||
|
||||
/**
|
||||
* можно использовать метод sessionInCondition, но не хочется лишнего or в предикате хазелкаста.
|
||||
*/
|
||||
public void addRgsConditionSessionTypeEquals(SessionType sessionType) {
|
||||
ImdgPredicate sessionTypePredicate = rgsPb.equals(
|
||||
"sessionType", sessionType.getKey()
|
||||
);
|
||||
this.registryConditionsAndV2.add(sessionTypePredicate);
|
||||
}
|
||||
|
||||
public void addRgsConditionDynamic(Function<ExtendedState, ImdgPredicate> dynamicCondition) {
|
||||
this.dynamicConditions.add(dynamicCondition);
|
||||
}
|
||||
|
||||
public void addRgsConditionStatusIn(RegistryStatus... registryStatuses) {
|
||||
ImdgPredicate registryStatusesCondition = rgsPb.or(
|
||||
Arrays.stream(registryStatuses)
|
||||
.map(registryStatus -> rgsPb.equals("registryStatus", registryStatus.getKey()))
|
||||
.toArray(ImdgPredicate[]::new)
|
||||
);
|
||||
this.registryConditionsAndV2.add(registryStatusesCondition);
|
||||
}
|
||||
|
||||
public void addRgsConditionStatusEquals(RegistryStatus registryStatus) {
|
||||
ImdgPredicate registryStatusCondition = rgsPb.equals("registryStatus", registryStatus.getKey());
|
||||
this.registryConditionsAndV2.add(registryStatusCondition);
|
||||
}
|
||||
|
||||
public void skipExecutionLoad() {
|
||||
this.loadExecs = false;
|
||||
}
|
||||
|
||||
public void setExecType(ExecutionType execType) {
|
||||
this.execType = execType;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void actualExecute(StateContext<TaskType, SsnEvent> ctx) {
|
||||
sessionType = ctx.getExtendedState().get(DataEnum.sessionType, SessionType.class);
|
||||
Long sessionId = ctx.getExtendedState().get(DataEnum.sessionId, Long.class);
|
||||
Long counterPartyId = ctx.getExtendedState().get(DataEnum.counterPartyId, Long.class);
|
||||
SessionType sessionType = ctx.getExtendedState().get(DataEnum.sessionType, SessionType.class);
|
||||
String sqlCondition = String.format("registryDesignation in ('%s', '%s') and " +
|
||||
"registryInstrumentType in ('%s', '%s') and " +
|
||||
"registryUnit = '%s'",
|
||||
RegistryDesignation.O.getKey(), RegistryDesignation.T.getKey(),
|
||||
RegistryInstrumentType.S.getKey(), RegistryInstrumentType.M.getKey(),
|
||||
RegistryUnit.T.getKey());
|
||||
ImdgPredicateBuilder prdctBuilder = registryImdg.predicateBuilder();
|
||||
ImdgPredicate defaultCondition = prdctBuilder.sql(sqlCondition);
|
||||
|
||||
ImdgPredicate registryPredicate;
|
||||
ImdgPredicate[] conditions = Stream.of(
|
||||
dynamicConditions.stream().map(d -> d.apply(ctx.getExtendedState())),
|
||||
registryConditionsAndV2.stream(),
|
||||
Stream.of(defaultCondition))
|
||||
.flatMap(Function.identity())
|
||||
.toArray(ImdgPredicate[]::new);
|
||||
registryPredicate = prdctBuilder.and(conditions);
|
||||
|
||||
log.info("Inclusion to pool with predicate: {}", registryPredicate.toString());
|
||||
Collection<Registry> obligations = registryImdg.getCollectionObjectsByPredicate(registryPredicate);
|
||||
Map<GroupRgsKey, List<Registry>> registryByGroupId = obligations.stream()
|
||||
.filter(registry -> registry.getSettlementDate().isEqual(LocalDate.now()))
|
||||
.collect(Collectors.groupingBy(registry -> new GroupRgsKey(registry.getGroupId(), registry.getMarket())));
|
||||
|
||||
Optional<SessionType> ssnType = Optional.of(sessionType);
|
||||
Map<Long, Registry> rgsToUpdate = new HashMap<>();
|
||||
List<ExecutionCommon> execsToUpdate = new ArrayList<>();
|
||||
for (Map.Entry<GroupRgsKey, List<Registry>> entrySet : registryByGroupId.entrySet()) {
|
||||
log.debug("Processing set of registry with groupId/market: {}/{}", entrySet.getKey().groupId(), entrySet.getKey().market());
|
||||
String rgsSection = null;
|
||||
for (Registry registry : entrySet.getValue()) {
|
||||
registry.setRegistryStatus(RegistryStatus.POOL.getKey());
|
||||
registry.setClearingDate(LocalDate.now());
|
||||
registry.setSessionId(sessionId);
|
||||
ssnType.ifPresent(st -> registry.setSessionType(st.getKey()));
|
||||
rgsToUpdate.put(registry.getId(), registry);
|
||||
if (TextUtil.isEmpty(rgsSection)) {
|
||||
rgsSection = registry.getSection();
|
||||
}
|
||||
}
|
||||
if (loadExecs) {
|
||||
execsToUpdate.addAll(updateExecutionsV2(sessionId, entrySet.getKey(), rgsSection));
|
||||
}
|
||||
}
|
||||
registryImdg.putAll(rgsToUpdate);
|
||||
if (loadExecs) {
|
||||
updateExecsBatchV2(execsToUpdate);
|
||||
}
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private Collection<ExecutionCommon> updateExecutionsV2(Long sessionId, GroupRgsKey groupRgsKey, String rgsSection) {
|
||||
Instant now = Instant.now();
|
||||
if (sessionId == null) {
|
||||
log.error("cannot update executions#sessionId - session id is null");
|
||||
return Collections.emptyList();
|
||||
}
|
||||
Imdg<? extends ExecutionCommon> execImdg = null;
|
||||
if (execType == null) {
|
||||
throw new IllegalStateException("execution type not set");
|
||||
}
|
||||
switch (execType) {
|
||||
case ExecutionDeposit -> execImdg = executionDepositImdg;
|
||||
case ExecutionFond -> execImdg = executionFondImdg;
|
||||
case ExecutionCurrency -> execImdg = executionCurrImdg;
|
||||
}
|
||||
if (execImdg == null) {
|
||||
log.warn("couldn't define Execution Type for session {}. Will not update executions#sessionId",
|
||||
this.sessionType);
|
||||
return Collections.emptyList();
|
||||
}
|
||||
ImdgPredicateBuilder pb = execImdg.predicateBuilder();
|
||||
ImdgPredicate prdct = pb.and(
|
||||
pb.equals("exchangeExecutionId", groupRgsKey.groupId()),
|
||||
StringUtils.hasText(groupRgsKey.market()) ? pb.equals("market", groupRgsKey.market()) : pb.alwaysTrue()
|
||||
);
|
||||
Collection<ExecutionCommon> execs = (Collection<ExecutionCommon>) execImdg.getCollectionObjectsByPredicate(prdct);
|
||||
execs.forEach(e -> {
|
||||
e.setUpdated(now);
|
||||
e.setSessionId(sessionId);
|
||||
});
|
||||
return execs;
|
||||
}
|
||||
|
||||
private void updateExecsBatchV2(List<ExecutionCommon> execs) {
|
||||
log.debug("updating {} executions", execs.size());
|
||||
execs.sort(Comparator.comparing(ExecutionCommon::type));
|
||||
log.debug("sorted executions by type");
|
||||
|
||||
Map<Long, ExecutionCommon> m = new HashMap<>();
|
||||
for (int i = 0; i < execs.size(); i++) {
|
||||
ExecutionCommon exec = execs.get(i);
|
||||
ExecutionType execType = exec.type();
|
||||
log.debug("batch from {} position type {}", i, execType);
|
||||
m.put(exec.getId(), exec);
|
||||
int j = i + 1;
|
||||
while (j < execs.size() && j < i + 100) {
|
||||
ExecutionCommon other = execs.get(j);
|
||||
if (!Objects.equals(other.type(), execType)) {
|
||||
break;
|
||||
}
|
||||
m.put(other.getId(), other);
|
||||
j++;
|
||||
}
|
||||
i = j;
|
||||
log.debug("batch size {} type {}", m.size(), execType);
|
||||
if (!m.isEmpty()) {
|
||||
switch (execType) {
|
||||
case ExecutionDeposit -> executionDepositImdg.putAll(ClearingUtil.castMap(m));
|
||||
case ExecutionFond -> executionFondImdg.putAll(ClearingUtil.castMap(m));
|
||||
case ExecutionCurrency -> executionCurrImdg.putAll(ClearingUtil.castMap(m));
|
||||
}
|
||||
}
|
||||
log.debug("batch size {} type {} done", m.size(), execType);
|
||||
m.clear();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,7 @@ package ru.spcex.platform.enumeration;
|
|||
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||
|
||||
public enum SessionType implements IEnumKey {
|
||||
IPOB("IPOB"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"),
|
||||
IPOB("IPOB"), IPOM("IPOM"), IPO0("IPO0"), IPOT("IPOT"), TRDT("TRDT"),
|
||||
MEDM("MEDM"), FINL("FINL"), XDEP("XDEP"), LIQU("LIQU"),
|
||||
CURR("CURR"), UNIT("UNIT"),
|
||||
PREP("PREP"), PREC("PREC"), PAYM("PAYM"), RSLT("RSLT"),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue