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 06ad2680b..9fbeffca1 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 @@ -15,6 +15,7 @@ import org.springframework.context.annotation.Scope; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.misc.Currency; import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.sdf.SDf51; import ru.clearing.classes.statics.data.sdf.SDf56; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.error.ClearingError; @@ -57,6 +58,7 @@ public class BalanceRevise implements ISessionStage { private Imdg statementImdg; private Imdg registryImdg; private Imdg currencyImdg; + private Imdg sDf51Imdg; private Imdg sDf56Imdg; private KafkaSender kafkaSender; @@ -67,6 +69,7 @@ public class BalanceRevise implements ISessionStage { this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class); this.currencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Currency, Currency.class); this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.sDf51Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf51, SDf51.class); this.sDf56Imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_SDf56, SDf56.class); this.kafkaSender = kafkaSender; } @@ -74,12 +77,15 @@ public class BalanceRevise implements ISessionStage { @Override public StageResult submit(Task task) { switch (task.getTaskType()) { - case StartRevise, FormingPaymentInstruction -> { + case StartRevise -> { + return sendSdfs51(); + } + case FormingPaymentInstruction -> { LocalTime fromTime = null; if (task.getData() instanceof LocalTime) { fromTime = (LocalTime) task.getData(); } - return sendSdfs(fromTime); + return sendSdfs56(fromTime); } case ContinueRevise -> { return revise(); // основная сверка @@ -95,26 +101,31 @@ public class BalanceRevise implements ISessionStage { } - private StageResult sendSdfs(LocalTime fromTime) { - Collection currencies = currencyImdg.projectSingleAttribute("id"); + private StageResult sendSdfs56(LocalTime fromTime) { Long millis = null; if (fromTime != null) { - log.debug("sendSdf, parameter fromTime={}", fromTime); + log.debug("sendSdf 56, parameter fromTime={}", fromTime); Instant fromInstant = TimeUtil.localDateTimeToInstant(LocalDateTime.of(LocalDate.now(), fromTime)); millis = fromInstant.toEpochMilli(); } else { + Collection currencies = currencyImdg.projectSingleAttribute("id"); Statement statement = statementImdg.aggregateByMax("created", StatementRevisePredicate.get(statementImdg, currencies)); if (statement != null) { - log.debug("sendSdf, parameter fromTime not set, use max(statement.created)={}", statement.getCreated()); + log.debug("sendSdf 56, parameter fromTime not set, use max(statement.created)={}", statement.getCreated()); millis = statement.getCreated().toEpochMilli(); } else { - log.debug("sendSdf, parameter fromTime not set, and statement's not found."); + log.debug("sendSdf 56, parameter fromTime not set, and statement's not found."); } } newSDf56(millis); return new StageResult<>(null, true); } + private StageResult sendSdfs51() { + newSDf51(); + return new StageResult<>(null, true); + } + private StageResult revise() { String statementSQL = String.format("inOutSDfType = %s and operationStatus = %s and statementType = %s", InOutSDfType.type1.getKey(), OperationStatus.Pending.getKey(), StatementType.full.getKey()); @@ -220,4 +231,24 @@ public class BalanceRevise implements ISessionStage { kafkaSender.sendRequestToQueue(Consts.SDF56_PROCESS, requestForExporter); log.debug("successfully processed, new id {}", sDf56.getId()); } + + private void newSDf51() { + log.debug("creating sdf51"); + SDf51 sDf51 = new SDf51(); + Instant now = Instant.now(); + Long numberId = idGenerator.nextId(); // из генератора + sDf51.setNumber(numberId.toString()); + if (sDf51.getNumber().length() > 10) { + log.warn("SDf51 number='{}' too large that 10 symbols", sDf51.getNumber()); + } + String datetime = String.valueOf(now.minus(1, ChronoUnit.DAYS).toEpochMilli()); + sDf51.setDatetime(datetime); // Дата и время сообщения + sDf51.setGenerationTime(now); + sDf51.setGenerationId(numberId); + sDf51Imdg.insert(sDf51); + SdfClearingRequest requestForExporter = new SdfClearingRequest(); + requestForExporter.setGroupId(sDf51.getGenerationId()); + kafkaSender.sendRequestToQueue(Consts.SDF51_PROCESS, requestForExporter); + log.debug("successfully processed, new id {}", sDf51.getId()); + } }