force stop session feature

This commit is contained in:
akulikov 2024-04-23 15:01:17 +03:00
parent 4d37475cbc
commit 8a6f2abc8c
3 changed files with 79 additions and 9 deletions

View file

@ -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<String, Object> kafkaQueue, Producer<String, Object> 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<String> 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();
}

View file

@ -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<ExecutionDeposit> execDpstImdg;
private final Imdg<ExecutionFond> execFondImdg;
private final Imdg<Registry> rgsImdg;
private final Imdg<Session> 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<LauncherCommandRequest> 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<Session> sessions = sessionImdg.getCollectionObjectsByPredicate(
pb.and(
pb.equals("sessionType", sessionTypeStr),
pb.equals("section", sectionStr)
)
);
List<Long> 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())) {

View file

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