This commit is contained in:
parent
5d40c9db34
commit
f64910154c
3 changed files with 72 additions and 20 deletions
|
|
@ -1,5 +1,6 @@
|
|||
package ru.spcex.clearing.service;
|
||||
|
||||
import java.time.LocalTime;
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.slf4j.Logger;
|
||||
|
|
@ -17,7 +18,11 @@ import ru.spcex.clearing.platform.messaging.domain.cud.clearing.SessionContinueE
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.payment.PIClearingOutbondActionNewRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.*;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.IdentificationFundsRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeRefundDateRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryChangeStatusExtractRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistryReturnDepositRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.registry.RegistrySplitDepositActionRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.STradesImportedRequest;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
|
|
@ -26,13 +31,21 @@ import ru.spcex.clearing.platform.messaging.service.Status;
|
|||
import ru.spcex.clearing.service.executors.Sdf06Executor;
|
||||
import ru.spcex.clearing.service.executors.Sdf10Executor;
|
||||
import ru.spcex.clearing.service.payment.PaymentInstructionOutboundService;
|
||||
import ru.spcex.clearing.session.stage.*;
|
||||
import ru.spcex.clearing.session.stage.FinalMkrSession;
|
||||
import ru.spcex.clearing.session.stage.IntermediateMkrSession;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionB0Session;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionBnSession;
|
||||
import ru.spcex.clearing.session.stage.PrimaryAuctionT0Session;
|
||||
import ru.spcex.clearing.session.stage.ReturnDepositSession;
|
||||
import ru.spcex.clearing.session.stage.SecondaryAuctionT0Session;
|
||||
import ru.spcex.clearing.session.stage.SessionManager;
|
||||
import ru.spcex.clearing.session.stage.SessionTerminator;
|
||||
import ru.spcex.clearing.session.stage.TaskType;
|
||||
import ru.spcex.clearing.session.stage.impl.BalanceRevise;
|
||||
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
|
||||
|
|
@ -173,7 +186,13 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean {
|
|||
|
||||
callback(LauncherCommandRequest.class)
|
||||
.setConsumer(task -> {
|
||||
balanceRevise.submit(new ru.spcex.clearing.session.stage.Task<>(TaskType.StartRevise, null));
|
||||
LocalTime fromTime = null;
|
||||
if (task != null
|
||||
&& task.getRequestPayload() != null
|
||||
&& task.getRequestPayload().getFromTime() != null) {
|
||||
fromTime = task.getRequestPayload().getFromTime();
|
||||
}
|
||||
balanceRevise.submit(new ru.spcex.clearing.session.stage.Task<>(TaskType.StartRevise, fromTime));
|
||||
})
|
||||
.forDestination(Task.getAllBalance.topic(), callbacks::put);
|
||||
|
||||
|
|
|
|||
|
|
@ -1,5 +1,12 @@
|
|||
package ru.spcex.clearing.session.stage.impl;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.LocalTime;
|
||||
import java.time.temporal.ChronoUnit;
|
||||
import java.util.Collection;
|
||||
import java.util.Objects;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
|
|
@ -19,7 +26,16 @@ import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
|||
import ru.spcex.clearing.session.stage.ISessionStage;
|
||||
import ru.spcex.clearing.session.stage.StageResult;
|
||||
import ru.spcex.clearing.session.stage.Task;
|
||||
import ru.spcex.platform.enumeration.*;
|
||||
import ru.spcex.platform.enumeration.AccountType;
|
||||
import ru.spcex.platform.enumeration.InOutSDfType;
|
||||
import ru.spcex.platform.enumeration.ObjectType;
|
||||
import ru.spcex.platform.enumeration.OperationStatus;
|
||||
import ru.spcex.platform.enumeration.Priority;
|
||||
import ru.spcex.platform.enumeration.RegistryDesignation;
|
||||
import ru.spcex.platform.enumeration.RegistryInstrumentType;
|
||||
import ru.spcex.platform.enumeration.RegistryTradingParams;
|
||||
import ru.spcex.platform.enumeration.RegistryUnit;
|
||||
import ru.spcex.platform.enumeration.StatementType;
|
||||
import ru.spcex.platform.imdg.api.Imdg;
|
||||
import ru.spcex.platform.imdg.api.ImdgId;
|
||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||
|
|
@ -28,12 +44,7 @@ import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
|
|||
import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder;
|
||||
import ru.spcex.platform.imdg.api.predicate.specific.StatementRevisePredicate;
|
||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.time.temporal.ChronoUnit;
|
||||
import java.util.Collection;
|
||||
import java.util.Objects;
|
||||
|
||||
import ru.spcex.platform.utils.time.TimeUtil;
|
||||
import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD;
|
||||
|
||||
@Service
|
||||
|
|
@ -64,7 +75,11 @@ public class BalanceRevise implements ISessionStage {
|
|||
public StageResult<?> submit(Task<?> task) {
|
||||
switch (task.getTaskType()) {
|
||||
case StartRevise, FormingPaymentInstruction -> {
|
||||
return sendSdfs();
|
||||
LocalTime fromTime = null;
|
||||
if (task.getData() instanceof LocalTime) {
|
||||
fromTime = (LocalTime) task.getData();
|
||||
}
|
||||
return sendSdfs(fromTime);
|
||||
}
|
||||
case ContinueRevise -> {
|
||||
return revise(); // основная сверка
|
||||
|
|
@ -80,10 +95,19 @@ public class BalanceRevise implements ISessionStage {
|
|||
}
|
||||
|
||||
|
||||
private StageResult<?> sendSdfs() {
|
||||
private StageResult<?> sendSdfs(LocalTime fromTime) {
|
||||
Collection<Long> currencies = currencyImdg.projectSingleAttribute("id");
|
||||
Long millis = null;
|
||||
if (fromTime != null) {
|
||||
Instant fromInstant = TimeUtil.localDateTimeToInstant(LocalDateTime.of(LocalDate.now(), fromTime));
|
||||
millis = fromInstant.toEpochMilli();
|
||||
} else {
|
||||
Statement statement = statementImdg.aggregateByMax("created", StatementRevisePredicate.get(statementImdg, currencies));
|
||||
newSDf56(statement);
|
||||
if (statement != null) {
|
||||
millis = statement.getCreated().toEpochMilli();
|
||||
}
|
||||
}
|
||||
newSDf56(millis);
|
||||
return new StageResult<>(null, true);
|
||||
}
|
||||
|
||||
|
|
@ -108,7 +132,7 @@ public class BalanceRevise implements ISessionStage {
|
|||
}
|
||||
|
||||
private StageResult<?> reviseStage1() {
|
||||
String sql = RegistryCodeSqlBuilder.getInstance(ru.spcex.platform.enumeration.RegistryTradingParams.A__T).build();
|
||||
String sql = RegistryCodeSqlBuilder.getInstance(RegistryTradingParams.A__T).build();
|
||||
Collection<Registry> regsAT = registryImdg.getCollectionObjectsBySQL(sql);
|
||||
log.trace("Select {} registry's by query \"{}\" for revision step 1", regsAT.size(), sql);
|
||||
|
||||
|
|
@ -175,13 +199,13 @@ public class BalanceRevise implements ISessionStage {
|
|||
|
||||
}
|
||||
|
||||
private void newSDf56(Statement statement) {
|
||||
private void newSDf56(Long millis) {
|
||||
log.debug("creating sdf56");
|
||||
SDf56 sDf56 = new SDf56();
|
||||
sDf56.setNumber(idGenerator.nextId().toString());
|
||||
Instant now = Instant.now();
|
||||
String startTime = String.valueOf(statement != null ?
|
||||
statement.getCreated().toEpochMilli() : now.minus(1, ChronoUnit.DAYS).toEpochMilli());
|
||||
String startTime = String.valueOf(millis != null ?
|
||||
millis : now.minus(1, ChronoUnit.DAYS).toEpochMilli());
|
||||
sDf56.setStart_datetime(startTime);
|
||||
sDf56.setEnd_datetime(String.valueOf(now.toEpochMilli()));
|
||||
sDf56.setGenerationTime(now);
|
||||
|
|
|
|||
|
|
@ -1,7 +1,12 @@
|
|||
package ru.spcex.platform.utils.time;
|
||||
|
||||
import java.sql.Time;
|
||||
import java.time.*;
|
||||
import java.time.Instant;
|
||||
import java.time.LocalDate;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.LocalTime;
|
||||
import java.time.ZoneId;
|
||||
import java.time.ZonedDateTime;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import java.time.temporal.ChronoUnit;
|
||||
import java.util.Date;
|
||||
|
|
@ -52,6 +57,10 @@ public class TimeUtil {
|
|||
return date.atStartOfDay(zone).toInstant();
|
||||
}
|
||||
|
||||
public static Instant localDateTimeToInstant(LocalDateTime date) {
|
||||
return date.atZone(zone).toInstant();
|
||||
}
|
||||
|
||||
public static LocalDate strToLocalDate(String date) {
|
||||
try {
|
||||
LocalDate parsed = LocalDate.parse(date, PROPERTY_DATE_FORMATTER);
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue