From f64910154cb407a26ea6d221c2914a8b967c63a0 Mon Sep 17 00:00:00 2001 From: ialbert Date: Fri, 9 Feb 2024 18:20:40 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-647 --- .../clearing/service/EventsReceiver.java | 27 ++++++++-- .../session/stage/impl/BalanceRevise.java | 54 +++++++++++++------ .../spcex/platform/utils/time/TimeUtil.java | 11 +++- 3 files changed, 72 insertions(+), 20 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 4527e3f19..d8f2282d5 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,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); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java index c55337788..ba757fb8d 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/BalanceRevise.java @@ -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 currencies = currencyImdg.projectSingleAttribute("id"); - Statement statement = statementImdg.aggregateByMax("created", StatementRevisePredicate.get(statementImdg, currencies)); - newSDf56(statement); + 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)); + 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 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); diff --git a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java index 43624f183..3da417437 100644 --- a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java +++ b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/time/TimeUtil.java @@ -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);