SDF01 statement#clearingDate,securityId fix
This commit is contained in:
parent
5f121a09b0
commit
547c993c53
4 changed files with 14 additions and 23 deletions
|
|
@ -4,6 +4,7 @@ import org.apache.kafka.clients.consumer.Consumer;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.InitializingBean;
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
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;
|
||||||
|
|
@ -15,6 +16,7 @@ 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.clearing.session.stage.SecondaryAuctionT0Session;
|
||||||
|
import ru.spcex.clearing.session.stage.SessionManager;
|
||||||
import ru.spcex.platform.enumeration.Task;
|
import ru.spcex.platform.enumeration.Task;
|
||||||
|
|
||||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.S_TRADES_IMPORTED;
|
||||||
|
|
@ -27,19 +29,22 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
private final PrimaryAuctionBnSession primaryAuctionBnSession;
|
private final PrimaryAuctionBnSession primaryAuctionBnSession;
|
||||||
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
private final SecondaryAuctionT0Session secondaryAuctionT0Session;
|
||||||
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
private final PrimaryAuctionB0Session primaryAuctionB0Session;
|
||||||
|
private final SessionManager sessionManager;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
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,
|
SecondaryAuctionT0Session secondaryAuctionT0Session,
|
||||||
PrimaryAuctionB0Session primaryAuctionB0Session) {
|
PrimaryAuctionB0Session primaryAuctionB0Session, SessionManager sessionManager) {
|
||||||
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.secondaryAuctionT0Session = secondaryAuctionT0Session;
|
||||||
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
||||||
|
this.sessionManager = sessionManager;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -60,22 +65,10 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
||||||
.setConsumer(clearingService::continueClearing)
|
.setConsumer(clearingService::continueClearing)
|
||||||
.forDestination(Consts.CONTINUE_CLEARING, callbacks::put);
|
.forDestination(Consts.CONTINUE_CLEARING, callbacks::put);
|
||||||
|
|
||||||
callback(Object.class)
|
|
||||||
.setConsumer(primaryAuctionB0Session::runSession)
|
|
||||||
.forDestination(Task.startOfB0.topic(), callbacks::put);
|
|
||||||
// callback(Object.class)
|
|
||||||
// .setConsumer(primaryAuctionB0Session::continueSession)
|
|
||||||
// .forDestination(Task.startOfB0.topic(), callbacks::put);
|
|
||||||
|
|
||||||
|
|
||||||
callback(LauncherCommandRequest.class)
|
callback(LauncherCommandRequest.class)
|
||||||
.setConsumer(primaryAuctionBnSession::runSession)
|
.setConsumer(sessionManager::defineAndStartSession)
|
||||||
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
.forDestination(Task.startOfClearing.topic(), callbacks::put);
|
||||||
|
|
||||||
callback(LauncherCommandRequest.class)
|
|
||||||
.setConsumer(secondaryAuctionT0Session::runSession)
|
|
||||||
.forDestination(Task.startOfT0.topic(), callbacks::put);
|
|
||||||
|
|
||||||
callback(Object.class)
|
callback(Object.class)
|
||||||
.setConsumer(req -> {
|
.setConsumer(req -> {
|
||||||
primaryAuctionBnSession.continueSession(req);
|
primaryAuctionBnSession.continueSession(req);
|
||||||
|
|
|
||||||
|
|
@ -110,8 +110,7 @@ public class PrimaryAuctionBnSession extends AbstractSession implements Initiali
|
||||||
} else if (checkStage(TaskType.FormingPaymentInstruction)) {
|
} else if (checkStage(TaskType.FormingPaymentInstruction)) {
|
||||||
finishPart(req);
|
finishPart(req);
|
||||||
} else {
|
} else {
|
||||||
log.error("cannot continue session, current stage is {}", currStage.get());
|
log.info("will not continue session, current stage is {}", currStage.get());
|
||||||
throw new StageException();
|
|
||||||
}
|
}
|
||||||
} catch (StageException e) {
|
} catch (StageException e) {
|
||||||
//already logged
|
//already logged
|
||||||
|
|
|
||||||
|
|
@ -104,14 +104,12 @@ public class SecondaryAuctionT0Session extends AbstractSession implements Initia
|
||||||
if (!isRunning()) {
|
if (!isRunning()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (!checkStage(TaskType.StartRevise)) {
|
|
||||||
log.error("cannot continue session, current stage is {}", currStage.get());
|
|
||||||
throw new StageException();
|
|
||||||
}
|
|
||||||
if (checkStage(TaskType.StartRevise)) {
|
if (checkStage(TaskType.StartRevise)) {
|
||||||
firstPart(req);
|
firstPart(req);
|
||||||
} else if (checkStage(TaskType.FormingPaymentInstruction)) {
|
} else if (checkStage(TaskType.FormingPaymentInstruction)) {
|
||||||
finishPart(req);
|
finishPart(req);
|
||||||
|
} else {
|
||||||
|
log.info("will not continue session, current stage is {}", currStage.get());
|
||||||
}
|
}
|
||||||
} catch (StageException e) {
|
} catch (StageException e) {
|
||||||
//already logged
|
//already logged
|
||||||
|
|
|
||||||
|
|
@ -24,10 +24,11 @@ public class SessionManager {
|
||||||
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
this.primaryAuctionB0Session = primaryAuctionB0Session;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void defineAndStartSession(LauncherCommandRequest commandRequest) {
|
public void defineAndStartSession(BaseRequest<LauncherCommandRequest> r) {
|
||||||
|
LauncherCommandRequest payload = r.getRequestPayload();
|
||||||
// Long sessionId = commandRequest.getSessionId();
|
// Long sessionId = commandRequest.getSessionId();
|
||||||
SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, commandRequest.getSessionType());
|
SessionType sessionType = IEnumKey.getEnumByKey(SessionType.class, payload.getSessionType());
|
||||||
Section section = IEnumKey.getEnumByKey(Section.class, commandRequest.getSection());
|
Section section = IEnumKey.getEnumByKey(Section.class, payload.getSection());
|
||||||
if (sessionType == null || section == null) {
|
if (sessionType == null || section == null) {
|
||||||
log.warn("Unknown sessionType: {} or section: {}", section, sessionType);
|
log.warn("Unknown sessionType: {} or section: {}", section, sessionType);
|
||||||
return;
|
return;
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue