This commit is contained in:
parent
0cba4f2dbb
commit
87ea31c507
5 changed files with 41 additions and 31 deletions
|
|
@ -192,6 +192,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
registryF.setRegistryUnit(RegistryUnit.F.getKey());
|
registryF.setRegistryUnit(RegistryUnit.F.getKey());
|
||||||
registryF.setRegistryCode(RegistryUtil.clearingCode(registryF));
|
registryF.setRegistryCode(RegistryUtil.clearingCode(registryF));
|
||||||
registryF.setBalance(registry.getBalance().subtract(registryB.getBalance()));
|
registryF.setBalance(registry.getBalance().subtract(registryB.getBalance()));
|
||||||
|
registryF.setBalance(BigDecimal.ZERO);
|
||||||
registryF.setId(imdgProvider.getImdgIdGenerator().nextId());
|
registryF.setId(imdgProvider.getImdgIdGenerator().nextId());
|
||||||
|
|
||||||
registryImdg.insert(registry);
|
registryImdg.insert(registry);
|
||||||
|
|
|
||||||
|
|
@ -171,6 +171,7 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali
|
||||||
{
|
{
|
||||||
EndStageNotificationPayload payload = new EndStageNotificationPayload();
|
EndStageNotificationPayload payload = new EndStageNotificationPayload();
|
||||||
payload.setSection(currSession.getSection());
|
payload.setSection(currSession.getSection());
|
||||||
|
payload.setSessionId(currSession.getId());
|
||||||
runStage(TaskType.EndStageNotification, payload, endStageNotification);
|
runStage(TaskType.EndStageNotification, payload, endStageNotification);
|
||||||
}
|
}
|
||||||
} catch (StageException e) {
|
} catch (StageException e) {
|
||||||
|
|
|
||||||
|
|
@ -6,6 +6,7 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
||||||
import org.springframework.context.annotation.Scope;
|
import org.springframework.context.annotation.Scope;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.clearing.classes.statics.data.misc.Session;
|
||||||
import ru.clearing.classes.statics.data.registry.Registry;
|
import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
|
@ -17,11 +18,13 @@ import ru.spcex.clearing.session.stage.Task;
|
||||||
import ru.spcex.clearing.session.stage.task.EndStageNotificationPayload;
|
import ru.spcex.clearing.session.stage.task.EndStageNotificationPayload;
|
||||||
import ru.spcex.platform.enumeration.RegistryStatus;
|
import ru.spcex.platform.enumeration.RegistryStatus;
|
||||||
import ru.spcex.platform.enumeration.Section;
|
import ru.spcex.platform.enumeration.Section;
|
||||||
|
import ru.spcex.platform.enumeration.SessionStatus;
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
||||||
|
import java.time.Instant;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Objects;
|
import java.util.Objects;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
|
|
@ -34,12 +37,14 @@ public class EndStageNotification implements ISessionStage {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
private final Imdg<Registry> registryImdg;
|
private final Imdg<Registry> registryImdg;
|
||||||
|
private final Imdg<Session> sessionImdg;
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
private final IMessageResolver msgResolver;
|
private final IMessageResolver msgResolver;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public EndStageNotification(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver) {
|
public EndStageNotification(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver) {
|
||||||
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
||||||
|
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
this.msgResolver = msgResolver;
|
this.msgResolver = msgResolver;
|
||||||
}
|
}
|
||||||
|
|
@ -49,7 +54,7 @@ public class EndStageNotification implements ISessionStage {
|
||||||
EndStageNotificationPayload payload = (EndStageNotificationPayload) task.getData();
|
EndStageNotificationPayload payload = (EndStageNotificationPayload) task.getData();
|
||||||
switch (task.getTaskType()) {
|
switch (task.getTaskType()) {
|
||||||
case EndStageNotification -> {
|
case EndStageNotification -> {
|
||||||
return endStageNotification(payload.getSection());
|
return endStageNotification(payload.getSection(), payload.getSessionId());
|
||||||
}
|
}
|
||||||
default -> {
|
default -> {
|
||||||
throw new IllegalStateException("Unknown task type: " + task.getTaskType());
|
throw new IllegalStateException("Unknown task type: " + task.getTaskType());
|
||||||
|
|
@ -64,7 +69,19 @@ public class EndStageNotification implements ISessionStage {
|
||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
protected StageResult<?> endStageNotification(String section) {
|
protected StageResult<?> endStageNotification(String section, Long sessionId) {
|
||||||
|
// 1. Изменить статус
|
||||||
|
Session theSession = sessionImdg.getSingleObjectByID(sessionId);
|
||||||
|
if (theSession == null) {
|
||||||
|
log.warn("Session {} not found", sessionId);
|
||||||
|
} else {
|
||||||
|
log.info("Finish status for session {}", sessionId);
|
||||||
|
theSession.setUpdated(Instant.now());
|
||||||
|
theSession.setSessionStatus(SessionStatus.CLOS.getKey());
|
||||||
|
sessionImdg.update(theSession);
|
||||||
|
log.trace("Session {} was updated", theSession.getId());
|
||||||
|
}
|
||||||
|
|
||||||
Collection<Registry> forRegistries = selectRegistry();
|
Collection<Registry> forRegistries = selectRegistry();
|
||||||
Collection<Long> groups = forRegistries.stream()
|
Collection<Long> groups = forRegistries.stream()
|
||||||
.map(Registry::getGroupId)
|
.map(Registry::getGroupId)
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,6 @@ import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
||||||
import org.springframework.context.annotation.Scope;
|
import org.springframework.context.annotation.Scope;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import ru.clearing.classes.statics.data.misc.Session;
|
|
||||||
import ru.clearing.classes.statics.data.registry.Registry;
|
import ru.clearing.classes.statics.data.registry.Registry;
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
|
@ -19,7 +18,6 @@ import ru.spcex.clearing.session.stage.Task;
|
||||||
import ru.spcex.clearing.session.stage.task.FinishingSessionPayload;
|
import ru.spcex.clearing.session.stage.task.FinishingSessionPayload;
|
||||||
import ru.spcex.platform.enumeration.RegistryDesignation;
|
import ru.spcex.platform.enumeration.RegistryDesignation;
|
||||||
import ru.spcex.platform.enumeration.RegistryTradingParams;
|
import ru.spcex.platform.enumeration.RegistryTradingParams;
|
||||||
import ru.spcex.platform.enumeration.SessionStatus;
|
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
import ru.spcex.platform.imdg.api.predicate.ImdgPredicate;
|
||||||
|
|
@ -28,7 +26,6 @@ import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
|
||||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||||
|
|
||||||
import java.time.Instant;
|
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
|
|
||||||
import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
|
import static ru.spcex.clearing.error.ClearingErrorInternal.SessionGeneralError;
|
||||||
|
|
@ -39,14 +36,12 @@ public class FinishingSession implements ISessionStage {
|
||||||
private final Logger log = LoggerFactory.getLogger(getClass());
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
||||||
private final Imdg<Registry> registryImdg;
|
private final Imdg<Registry> registryImdg;
|
||||||
private final Imdg<Session> sessionImdg;
|
|
||||||
private final KafkaSender kafkaSender;
|
private final KafkaSender kafkaSender;
|
||||||
private final IMessageResolver msgResolver;
|
private final IMessageResolver msgResolver;
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
public FinishingSession(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver) {
|
public FinishingSession(ImdgProvider imdgProvider, KafkaSender kafkaSender, IMessageResolver msgResolver) {
|
||||||
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
|
||||||
this.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class);
|
|
||||||
this.kafkaSender = kafkaSender;
|
this.kafkaSender = kafkaSender;
|
||||||
this.msgResolver = msgResolver;
|
this.msgResolver = msgResolver;
|
||||||
}
|
}
|
||||||
|
|
@ -64,31 +59,7 @@ public class FinishingSession implements ISessionStage {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
protected Collection<Registry> selectRegistry(Long sessionId) {
|
|
||||||
RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(
|
|
||||||
new RegistryTradingParams(RegistryDesignation.A, null, null, null)
|
|
||||||
);
|
|
||||||
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
|
|
||||||
ImdgPredicate condition = registryCodeSqlBuilder.buildPredicate(pb);
|
|
||||||
condition = pb.and(condition, pb.equals("sessionId", sessionId));
|
|
||||||
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(condition);
|
|
||||||
log.trace("Selected {} registry's by sql: {}", result.size(), condition);
|
|
||||||
return result;
|
|
||||||
}
|
|
||||||
|
|
||||||
protected StageResult<?> finishingSession(Long sessionId) {
|
protected StageResult<?> finishingSession(Long sessionId) {
|
||||||
// 1. Изменить статус
|
|
||||||
Session theSession = sessionImdg.getSingleObjectByID(sessionId);
|
|
||||||
if (theSession == null) {
|
|
||||||
log.warn("Session {} not found", sessionId);
|
|
||||||
} else {
|
|
||||||
log.info("Finish status for session {}", sessionId);
|
|
||||||
theSession.setUpdated(Instant.now());
|
|
||||||
theSession.setSessionStatus(SessionStatus.CLOS.getKey());
|
|
||||||
sessionImdg.update(theSession);
|
|
||||||
log.trace("Session {} was updated", theSession.getId());
|
|
||||||
}
|
|
||||||
|
|
||||||
// 2. Отправка сообщений
|
// 2. Отправка сообщений
|
||||||
Collection<Registry> forRegistries = selectRegistry(sessionId);
|
Collection<Registry> forRegistries = selectRegistry(sessionId);
|
||||||
log.debug("Found {} registries for sessionId={}", forRegistries.size(), sessionId);
|
log.debug("Found {} registries for sessionId={}", forRegistries.size(), sessionId);
|
||||||
|
|
@ -142,4 +113,15 @@ public class FinishingSession implements ISessionStage {
|
||||||
return new StageResult<>(null, true);
|
return new StageResult<>(null, true);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
protected Collection<Registry> selectRegistry(Long sessionId) {
|
||||||
|
RegistryCodeSqlBuilder registryCodeSqlBuilder = RegistryCodeSqlBuilder.getInstance(
|
||||||
|
new RegistryTradingParams(RegistryDesignation.A, null, null, null)
|
||||||
|
);
|
||||||
|
ImdgPredicateBuilder pb = registryImdg.predicateBuilder();
|
||||||
|
ImdgPredicate condition = registryCodeSqlBuilder.buildPredicate(pb);
|
||||||
|
condition = pb.and(condition, pb.equals("sessionId", sessionId));
|
||||||
|
Collection<Registry> result = registryImdg.getCollectionObjectsByPredicate(condition);
|
||||||
|
log.trace("Selected {} registry's by sql: {}", result.size(), condition);
|
||||||
|
return result;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ package ru.spcex.clearing.session.stage.task;
|
||||||
|
|
||||||
public class EndStageNotificationPayload {
|
public class EndStageNotificationPayload {
|
||||||
private String section;
|
private String section;
|
||||||
|
private Long sessionId;
|
||||||
|
|
||||||
public String getSection() {
|
public String getSection() {
|
||||||
return section;
|
return section;
|
||||||
|
|
@ -10,4 +11,12 @@ public class EndStageNotificationPayload {
|
||||||
public void setSection(String section) {
|
public void setSection(String section) {
|
||||||
this.section = section;
|
this.section = section;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public Long getSessionId() {
|
||||||
|
return sessionId;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setSessionId(Long sessionId) {
|
||||||
|
this.sessionId = sessionId;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue