From 9bf616b61c38f6da6cb381de614585c22fd58a79 Mon Sep 17 00:00:00 2001 From: akulikov Date: Tue, 23 Apr 2024 15:01:17 +0300 Subject: [PATCH] force stop session feature --- .../clearing/service/EventsReceiver.java | 17 ++++- .../session/stage/SessionTerminator.java | 70 +++++++++++++++++-- .../platform/messaging/domain/Consts.java | 1 + 3 files changed, 79 insertions(+), 9 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index d8f2282d5..6e882a5bf 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -1,6 +1,8 @@ package ru.spcex.clearing.service; import java.time.LocalTime; +import java.util.Arrays; +import java.util.List; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.producer.Producer; import org.slf4j.Logger; @@ -8,9 +10,11 @@ import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.core.env.Environment; import org.springframework.stereotype.Service; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; +import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRequest; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request; @@ -46,11 +50,11 @@ import ru.spcex.clearing.statement.StatementServiceV2; import ru.spcex.platform.enumeration.Task; import ru.spcex.platform.utils.enumeration.IMessageResolver; import ru.spcex.platform.utils.error.ValidationException; -import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED; @Service public class EventsReceiver extends QueueConsumer implements InitializingBean { private final Logger log = LoggerFactory.getLogger(getClass()); + private final Environment environment; private final IMessageResolver errorResolver; private final ClearingService clearingService; private final RegistryService registryService; @@ -72,7 +76,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { @Autowired public EventsReceiver(@Qualifier("kafkaConsumer") Consumer kafkaQueue, Producer kafkaResponseQueue, - IMessageResolver errorResolver, + Environment environment, IMessageResolver errorResolver, ClearingService clearingService, RegistryService registryService, PrimaryAuctionBnSession primaryAuctionBnSession, @@ -81,6 +85,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { Sdf06Executor sdf06Executor, Sdf10Executor sdf10Executor, BalanceRevise balanceRevise, Sdf05Sender sdf05Sender, StatementServiceV2 statementService, SessionTerminator sessionTerminator, PaymentInstructionOutboundService pmtOutboundService) { super(kafkaQueue, kafkaResponseQueue); + this.environment = environment; this.errorResolver = errorResolver; this.clearingService = clearingService; this.registryService = registryService; @@ -210,6 +215,14 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { .setFunction(sessionTerminator::stopCurrentSession) .forDestination(Consts.KILL_SESSION, callbacks::put); + String[] profilesRaw = environment.getActiveProfiles(); + List profiles = Arrays.stream(profilesRaw).map(s -> s == null ? "" : s.trim().toUpperCase()).toList(); + if (profiles.contains("DEV")) { + callback(LauncherCommandRequest.class) + .setFunction(sessionTerminator::forceStopSession) + .forDestination(Consts.FORCE_STOP_SESSION, callbacks::put); + } + init(); } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java index 1b244a856..af65bd289 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/SessionTerminator.java @@ -1,5 +1,14 @@ package ru.spcex.clearing.session.stage; +import java.math.BigDecimal; +import java.time.Instant; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Optional; +import java.util.function.Consumer; +import java.util.function.Supplier; +import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; @@ -12,7 +21,9 @@ 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.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; import ru.spcex.clearing.service.registry.AssetTBFProcessing; import ru.spcex.clearing.service.validation.ValidationStored; import ru.spcex.clearing.util.services.RequestHelper; @@ -21,28 +32,24 @@ 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.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.imdg.api.predicate.specific.RegistryCodeSqlBuilder; import ru.spcex.platform.utils.enumeration.EnumMessage; +import ru.spcex.platform.utils.enumeration.IEnumKey; 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 execDpstImdg; private final Imdg execFondImdg; private final Imdg rgsImdg; + private final Imdg sessionImdg; private final NotificationSender notification; private final AssetTBFProcessing assets; @@ -63,6 +70,7 @@ public class SessionTerminator { 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.sessionImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Session, Session.class); this.reqInfo = reqInfo; this.validation = validation; this.sessionMng = sessionMng; @@ -83,6 +91,53 @@ public class SessionTerminator { return reqInfo.success(session.getId()); } + public RequestInfoUpdate forceStopSession(BaseRequest req) { + LauncherCommandRequest requestPayload = req.getRequestPayload(); + String sessionTypeStr = requestPayload.getSessionType(); + String sectionStr = requestPayload.getSection(); + SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, sessionTypeStr); + Section section = IEnumKey.getEnumByKey(Section.class, sectionStr); + if (sessionType == null || section == null) { + String msg = "Nothing force stop session: session-type or section is empty"; + log.warn(msg); + RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate(); + requestInfoUpdate.setMessage(msg); + requestInfoUpdate.setStatus(Status.Success); + return requestInfoUpdate; + } + ImdgPredicateBuilder pb = sessionImdg.predicateBuilder(); + Collection sessions = sessionImdg.getCollectionObjectsByPredicate( + pb.and( + pb.equals("sessionType", sessionTypeStr), + pb.equals("section", sectionStr) + ) + ); + List cantComplete = new ArrayList<>(); + for (Session session : sessions) { + try { + rollbackExecutions(session); + rollbackAssetBalances(session.getId()); + deleteLiabilitiesAndClaims(session.getId()); + sessionMng.endSession(session); + } catch (Exception e) { + log.warn("Can't force stop session id=%s".formatted(session.getId()), e); + cantComplete.add(session.getId()); + } + } + String msg; + if (cantComplete.isEmpty()) { + msg = sessions.isEmpty() ? "Can't find sessions" : "All sessions completed success"; + } else { + msg = "For some sessions can't force stop, id=[%s]".formatted( + cantComplete.stream().map(String::valueOf).collect(Collectors.joining(",")) + ); + } + RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate(); + requestInfoUpdate.setMessage(msg); + requestInfoUpdate.setStatus(Status.Success); + return requestInfoUpdate; + } + private void rollbackExecutions(Session session) { Instant now = Instant.now(); @@ -100,6 +155,7 @@ public class SessionTerminator { }; log.debug("session.id={} section {}", session.getId(), session.getSection()); + // todo check Section.CURR if (Section.FOND.equalsByKey(session.getSection())) { nullifier.accept(upcastImdg(execFondImdg)); } else if (Section.MKR.equalsByKey(session.getSection())) { diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 48cdf4822..cdc01e358 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -141,6 +141,7 @@ public interface Consts { String CONTINUE_SESSION_BN_FIRST_PART = "sdf57-process"; String CONTINUE_SESSION_BN_SECOND_PART = "sdf13-process"; String KILL_SESSION = "kill-session"; + String FORCE_STOP_SESSION = "force-stop-session"; String SWT_EXPORTER = "swt-exporter"; String REVISE_PROCESS = "revise-process"; String EXPORT_PROCESS = "export-process";