start sessions based on request
This commit is contained in:
parent
d51cbf809a
commit
84bda2fa12
3 changed files with 26 additions and 3 deletions
|
|
@ -1,8 +1,11 @@
|
||||||
package ru.spcex.clearing.service;
|
package ru.spcex.clearing.service;
|
||||||
|
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.InitializingBean;
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRequest;
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.CreateRegistryRequest;
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
|
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.Sdf04Request;
|
||||||
|
|
@ -12,26 +15,32 @@ import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImported
|
||||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session;
|
import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session;
|
||||||
import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
|
import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
|
||||||
|
import ru.spcex.clearing.session.stage.SecondaryAuctionT0Session;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
|
import ru.spcex.platform.utils.enumeration.IEnumKey;
|
||||||
|
|
||||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
private final ClearingService clearingService;
|
private final ClearingService clearingService;
|
||||||
private final RegistryService registryService;
|
private final RegistryService registryService;
|
||||||
private final PrimaryAuctionBnSession primaryAuctionBnSession;
|
private final PrimaryAuctionBnSession primaryAuctionBnSession;
|
||||||
|
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
||||||
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
||||||
|
|
||||||
public EventsReceiver(Consumer<String, Object> kafkaQueue,
|
public EventsReceiver(Consumer<String, Object> kafkaQueue,
|
||||||
ClearingService clearingService,
|
ClearingService clearingService,
|
||||||
RegistryService registryService,
|
RegistryService registryService,
|
||||||
PrimaryAuctionBnSession primaryAuctionBnSession,
|
PrimaryAuctionBnSession primaryAuctionBnSession,
|
||||||
|
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
||||||
PrimaryAuctionB0Session primaryAuctionB0Session) {
|
PrimaryAuctionB0Session primaryAuctionB0Session) {
|
||||||
super(kafkaQueue);
|
super(kafkaQueue);
|
||||||
this.clearingService = clearingService;
|
this.clearingService = clearingService;
|
||||||
this.registryService = registryService;
|
this.registryService = registryService;
|
||||||
this.primaryAuctionBnSession = primaryAuctionBnSession;
|
this.primaryAuctionBnSession = primaryAuctionBnSession;
|
||||||
|
this.secondaryAuctionT0Session = secondaryAuctionT0Session;
|
||||||
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -61,8 +70,8 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
// .forDestination(Task.startOfB0.topic(), callbacks::put);
|
// .forDestination(Task.startOfB0.topic(), callbacks::put);
|
||||||
|
|
||||||
|
|
||||||
callback(Object.class)
|
callback(LauncherCommandRequest.class)
|
||||||
.setConsumer(primaryAuctionBnSession::runSession)
|
.setConsumer(this::startSession)
|
||||||
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
||||||
callback(Object.class)
|
callback(Object.class)
|
||||||
.setConsumer(primaryAuctionBnSession::continueSession)
|
.setConsumer(primaryAuctionBnSession::continueSession)
|
||||||
|
|
@ -76,4 +85,17 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
.forDestination(Consts.REGISTRY_NEW, callbacks::put);
|
.forDestination(Consts.REGISTRY_NEW, callbacks::put);
|
||||||
init();
|
init();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void startSession(BaseRequest<LauncherCommandRequest> request) {
|
||||||
|
LauncherCommandRequest payload = request.getRequestPayload();
|
||||||
|
Task specificSession = IEnumKey.getEnumByKey(Task.class, payload.getTaskName());
|
||||||
|
if (specificSession == null) {
|
||||||
|
log.error("unknown session: {}", payload.getTaskName());
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
switch (specificSession) {
|
||||||
|
case startOfClearing -> primaryAuctionBnSession.runSession(request);
|
||||||
|
case startOfT0 -> secondaryAuctionT0Session.runSession(request);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -96,7 +96,7 @@ public class Sdf57Executor extends AbstractExecutor<SDf57> {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public boolean isNeedToSendCommand() {
|
public boolean isNeedToSendCommand() {
|
||||||
return true;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
|
||||||
|
|
@ -11,6 +11,7 @@ public enum Task implements IEnumKey {
|
||||||
@Deprecated /* todo GBLD удаляется по CLS-267, CLS-275 */ getBalance("GBLD"),// Поступление средств
|
@Deprecated /* todo GBLD удаляется по CLS-267, CLS-275 */ getBalance("GBLD"),// Поступление средств
|
||||||
startOfClearing("SCLR"),// Запуск клиринговой сессии
|
startOfClearing("SCLR"),// Запуск клиринговой сессии
|
||||||
startOfB0("IPO0"),// Запуск клиринговой сессии
|
startOfB0("IPO0"),// Запуск клиринговой сессии
|
||||||
|
startOfT0("TRDT"),// Запуск клиринговой сессии
|
||||||
startOfPreClearing("SPRC"),// Запуск преклиринга
|
startOfPreClearing("SPRC"),// Запуск преклиринга
|
||||||
startPostClearing("SPOC"),// Запуск постклиринга
|
startPostClearing("SPOC"),// Запуск постклиринга
|
||||||
createOrder("CORD"),
|
createOrder("CORD"),
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue