ialbert 2026-07-03 16:00:22 +03:00
parent 56287f6a8e
commit fb57c44f9c
9 changed files with 290 additions and 188 deletions

View file

@ -50,6 +50,7 @@ public enum ClearingError implements IErrorEnumId {
WrongField(5004L),
GatewayTimeout(6003L),
GatewayNotApproved(6004L),
CounterStradeNotFound(-1L)
;
private final Long id;

View file

@ -0,0 +1,69 @@
package ru.spcex.clearing.service.execution;
import java.util.Collection;
import org.slf4j.Logger;
import org.slf4j.event.Level;
import ru.spcex.clearing.notification.NotificationSender;
import ru.spcex.platform.enumeration.ObjectType;
import ru.spcex.platform.enumeration.Priority;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
public abstract class ExecutionAbstractComponent implements IExecutionUploadComponent {
protected static final int EXTRA_ITERATIONS_WARN = 500;
protected static final int EXTRA_ITERATIONS_ERROR = 1000;
protected final NotificationSender notifications;
private final IMessageResolver msgResolver;
protected ExecutionAbstractComponent(NotificationSender notifications, IMessageResolver msgResolver) {
this.notifications = notifications;
this.msgResolver = msgResolver;
}
protected abstract Logger log();
protected int iterationsWarn(Collection<?> items) {
return items.size() + EXTRA_ITERATIONS_WARN;
}
protected int iterationsError(Collection<?> items) {
return items.size() + EXTRA_ITERATIONS_ERROR;
}
protected boolean cycleExceptionalInterrupt(int i, int iterationWarn, int iterationErr) {
if (i == iterationWarn) {
logNotifyE("Высокое число сделок без пары");
}
if (i > iterationErr) {
logNotifyE("Превышено число сделок без пары");
log().error("FATAL: too many unpaired STrades. breaking the cycle");
return true;
}
return false;
}
protected void logError(IEnumId subject, Object... args) {
String err = msgResolver.resolve(new EnumMessage(subject, args));
log().info("{}", err);
// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
protected void logNotifyW(String msg) {
logNotify(msg, Level.WARN);
}
protected void logNotifyE(String msg) {
logNotify(msg, Level.ERROR);
}
protected void logNotify(String msg, Level lvl) {
switch (lvl) {
case WARN -> log().warn(msg);
case ERROR -> log().error(msg);
default -> log().error("lvl not supported {}. {}", lvl, msg);
}
notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH);
}
}

View file

@ -6,15 +6,14 @@ import java.time.Instant;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import static java.lang.String.format;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.event.Level;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.scheduling.annotation.EnableScheduling;
@ -42,8 +41,9 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.service.validation.strades.CompanyByTradingCodeCashingValidationRule;
import ru.spcex.clearing.service.validation.strades.CounterStradePresentValidationRule;
import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentValidationRule;
import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule;
import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRuleV2;
import ru.spcex.clearing.validation.common.rules.FieldNotBlankRequiredRule;
import ru.spcex.platform.enumeration.MarketType;
import ru.spcex.platform.enumeration.ObjectType;
@ -60,8 +60,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.imdg.validation.LogPrefixId;
import ru.spcex.platform.utils.collection.CollectionUtil;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.log.ExceptionUtils;
@ -76,7 +76,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl;
*/
@Component
@EnableScheduling
public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
public class ExecutionCurrencyComponent extends ExecutionAbstractComponent {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<STrades> sTradeImdg;
@ -109,6 +109,7 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
IMessageResolver msgResolver,
NotificationSender notifications,
ClearingServiceSettings serviceSettings) {
super(notifications, msgResolver);
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
this.executionCurrencyImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionCurrency, ExecutionCurrency.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
@ -133,7 +134,12 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
ExecutionUploadCashUtil.loadExecutions(executionCurrencyImdg, cash);
}
// /**
@Override
protected Logger log() {
return log;
}
// /**
// * Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?"
// */
// @Scheduled(cron = "0 1 0 1 * ?")
@ -158,12 +164,14 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
LocalDate today = LocalDate.now();
log.debug("Start check new S_TRADE at {}, fromId={}", today, fromId);
//выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades;
List<STrades> sTrades;
if (cleanLoad) {
sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
sTrades = CollectionUtil.castOrCopyToModifyableList(
sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
"tradeDate", LocalDate.now(),
"section", Section.CURR.getKey()
));
))
);
} else {
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate filter = pb.and(
@ -171,7 +179,9 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
pb.equals("section", Section.CURR.getKey()),
pb.greatEqual("id", fromId)
);
sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter);
sTrades = CollectionUtil.castOrCopyToModifyableList(
sTradeImdg.getCollectionObjectsByPredicate(filter)
);
}
log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId);
@ -202,14 +212,29 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
}
log.info("{} strades left after already-added filtering", sTrades.size());
//защита от бесконтрольно разрастающегося цикла
int iterationsWarn = iterationsWarn(sTrades);
int iterationsMax = iterationsError(sTrades);
Map<Long, ExecutionCurrency> execsToInsert = new HashMap<>();
for (STrades sTrd : sTrades) {
for (int i = 0; i < sTrades.size(); i++) {
if (cycleExceptionalInterrupt(i, iterationsWarn, iterationsMax)) {
return;
}
STrades sTrd = sTrades.get(i);
log.trace("S_TRADE[{}] new", sTrd.getId());
//проверки
IValidator validator = valFor(sTrd);
Optional<EnumMessage> error = validator.tillFirstError();
STrades counter = validator.getStored(ValidationStored.STradesCounterTrade);
if (counter != null && fromId != null && counter.getId() < fromId) {
sTrades.add(counter);
}
if (error.isPresent()) {
if (error.get().getSubject().equals(ClearingError.CounterStradeNotFound)) {
log.info("STrade № {} id {} not processed. should be paired in next cycles.", sTrd.getTradeNum(), sTrd.getId());
continue;
}
String err = format("Сделка № %d не смогла быть обработана. %s", sTrd.getTradeNum(), msgResolver.resolve(error.get()));
logNotifyE(err);
continue;
@ -286,6 +311,8 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
Company company = validator.getStored(ValidationStored.STradesCompany);
Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany);
TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry);
STrades counterSTrades = validator.getStored(ValidationStored.STradesCounterTrade);
TradingClearingRegistry counterTCR = validator.getStored(ValidationStored.STradesCounterTradingClearingRegistry);
Listing listing = listingImdg.getFirstObjectByFieldValues(
Map.of(
"securityId", security.getId(),
@ -373,20 +400,8 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
eCurrency.setSecurityId(security.getId());
eCurrency.setCounterPartyId(counterCompany.getId());
{
ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder();
STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and(
strPb.equals("tradeNum", sTrades.getTradeNum()),
strPb.equals("section", sTrades.getSection()),
strPb.equals("classCode", sTrades.getClassCode()),
strPb.not(strPb.equals("operation", sTrades.getOperation()))));
if (counterSTrades != null) {
eCurrency.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount());
searchTcrByStrades(counterSTrades)
.ifPresent(tcr -> eCurrency.setCounterPartyTradingClearingRegistryId(tcr.getId()));
} else {
log.trace("Counter STRade for {tradeNum={}, section={}, operation={}} not found.",
sTrades.getTradeNum(), sTrades.getSection(), sTrades.getOperation());
}
eCurrency.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount());
eCurrency.setCounterPartyTradingClearingRegistryId(counterTCR.getId());
}
return eCurrency;
}
@ -416,45 +431,6 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
}
}
private void logError(IEnumId subject, Object... args) {
String err = msgResolver.resolve(new EnumMessage(subject, args));
log.info("{}", err);
// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
private void logError(Long sTradeNum, EnumMessage msg) {
String err = format("Сделка № %d не смогла быть обработана. %s", sTradeNum, msgResolver.resolve(msg));
log.warn(err);
notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
private void logNotifyW(String msg) {
logNotify(msg, Level.WARN);
}
private void logNotifyE(String msg) {
logNotify(msg, Level.ERROR);
}
private void logNotify(String msg, Level lvl) {
switch (lvl) {
case WARN -> log.warn(msg);
case ERROR -> log.error(msg);
default -> log.error("lvl not supported {}. {}", lvl, msg);
}
notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH);
}
private Optional<TradingClearingRegistry> searchTcrByStrades(STrades sTrades) {
IValidator counterValidator = valFor(sTrades);
Optional<EnumMessage> err = counterValidator.tillFirstError();
if (err.isPresent()) {
log.error("couldn't extract setCounterPartyTradingClearingRegistryId from strades.id={}", sTrades.getId());
return Optional.empty();
}
return Optional.ofNullable(counterValidator.getStored(ValidationStored.STradesTradingClearingRegistry));
}
public IValidator valFor(STrades sTrades) {
ImdgValidationContext<STrades> context = new ImdgValidationContext<>();
context.setValidatedObject(sTrades);
@ -465,6 +441,7 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
context.addImdg(IMDGDistributedNames.Map_Company, cmpImdg);
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, tcrImdg);
context.addImdg(IMDGDistributedNames.Map_DigitalCertificateSecurity, imdgDigitCertSecurity);
context.addImdg(IMDGDistributedNames.Map_STrades, sTradeImdg);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
new STradesSecurityPresentValidationRule(valCash.fixIncSecCash,
@ -478,11 +455,19 @@ public class ExecutionCurrencyComponent implements IExecutionUploadComponent {
new CompanyByTradingCodeCashingValidationRule<>(
STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCounterCompany
),
new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash),
new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash,
ImdgValidationContext::getValidatedObject,
ValidationStored.STradesCompany,
ValidationStored.STradesTradingClearingRegistry),
FieldNotBlankRequiredRule.instance(
"classCode",
STrades::getClassCode,
ClearingError.RequiredFieldEmpty, true)
ClearingError.RequiredFieldEmpty, true),
new CounterStradePresentValidationRule(sTrades),
new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash,
ctx -> ctx.getStoredObject(ValidationStored.STradesCounterTrade),
ValidationStored.STradesCounterCompany,
ValidationStored.STradesCounterTradingClearingRegistry)
);
}

View file

@ -5,15 +5,14 @@ import java.time.Instant;
import java.time.LocalDate;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import static java.lang.String.format;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.event.Level;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.scheduling.annotation.EnableScheduling;
@ -36,11 +35,10 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.service.validation.strades.CompanyByTradingCodeCashingValidationRule;
import ru.spcex.clearing.service.validation.strades.CounterStradePresentValidationRule;
import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentDepositValidationRule;
import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule;
import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRuleV2;
import ru.spcex.platform.enumeration.MoneyFlowSide;
import ru.spcex.platform.enumeration.ObjectType;
import ru.spcex.platform.enumeration.Priority;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Side;
import ru.spcex.platform.imdg.api.Imdg;
@ -53,8 +51,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.imdg.validation.LogPrefixId;
import ru.spcex.platform.utils.collection.CollectionUtil;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.log.ExceptionUtils;
@ -68,7 +66,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl;
*/
@Component
@EnableScheduling
public class ExecutionDepositComponent implements IExecutionUploadComponent {
public class ExecutionDepositComponent extends ExecutionAbstractComponent {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<STrades> sTradeImdg;
@ -97,6 +95,7 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
NotificationSender notifications,
@Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender,
IMessageResolver msgResolver) {
super(notifications, msgResolver);
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
this.executionDepositImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionDeposit, ExecutionDeposit.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
@ -111,6 +110,11 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
resetTradingDay();
}
@Override
protected Logger log() {
return log;
}
/**
* Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?"
*/
@ -135,12 +139,14 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
}
log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId);
//выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades;
List<STrades> sTrades;
if (cleanLoad) {
sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
sTrades = CollectionUtil.castOrCopyToModifyableList(
sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
"tradeDate", LocalDate.now(),
"section", Section.MKR.getKey()
));
))
);
} else {
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate filter = pb.and(
@ -148,7 +154,9 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
pb.equals("section", Section.MKR.getKey()),
pb.greatEqual("id", fromId)
);
sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter);
sTrades = CollectionUtil.castOrCopyToModifyableList(
sTradeImdg.getCollectionObjectsByPredicate(filter)
);
}
log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId);
@ -180,14 +188,28 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
}
log.info("{} strades left after already-added filtering", sTrades.size());
//защита от бесконтрольно разрастающегося цикла
int iterationsWarn = iterationsWarn(sTrades);
int iterationsMax = iterationsError(sTrades);
Map<Long, ExecutionDeposit> execsToInsert = new HashMap<>();
for (STrades sTrd : sTrades) {
for (int i = 0; i < sTrades.size(); i++) {
if (cycleExceptionalInterrupt(i, iterationsWarn, iterationsMax)) {
return;
}
STrades sTrd = sTrades.get(i);
log.trace("S_TRADE[{}] new", sTrd.getId());
//проверки
IValidator validator = valFor(sTrd);
Optional<EnumMessage> error = validator.tillFirstError();
STrades counter = validator.getStored(ValidationStored.STradesCounterTrade);
if (counter != null && fromId != null && counter.getId() < fromId) {
sTrades.add(counter);
}
if (error.isPresent()) {
if (error.get().getSubject().equals(ClearingError.CounterStradeNotFound)) {
log.info("STrade № {} id {} not processed. should be paired in next cycles.", sTrd.getTradeNum(), sTrd.getId());
continue;
}
String err = format("Сделка № %d не смогла быть обработана. %s", sTrd.getTradeNum(), msgResolver.resolve(error.get()));
logNotifyE(err);
continue;
@ -253,6 +275,8 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
Company company = validator.getStored(ValidationStored.STradesCompany);
Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany);
TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry);
STrades counterSTrades = validator.getStored(ValidationStored.STradesCounterTrade);
TradingClearingRegistry counterTCR = validator.getStored(ValidationStored.STradesCounterTradingClearingRegistry);
Listing listing = listingImdg.getFirstObjectByFieldValues(Map.of("securityId", security.getId()));
ExecutionDeposit eDeposit = new ExecutionDeposit();
@ -320,17 +344,8 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
eDeposit.setCounterPartyId(counterCompany.getId());
eDeposit.setSecurityCode(sTrades.getSecCode());
{
ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder();
STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and(
strPb.equals("tradeNum", sTrades.getTradeNum()),
strPb.equals("section", sTrades.getSection()),
strPb.equals("classCode", sTrades.getClassCode()),
strPb.not(strPb.equals("operation", sTrades.getOperation()))));
if (counterSTrades != null) {
eDeposit.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount());
searchTcrByStrades(counterSTrades)
.ifPresent(tcr -> eDeposit.setCounterPartyTradingClearingRegistryId(tcr.getId()));
}
eDeposit.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount());
eDeposit.setCounterPartyTradingClearingRegistryId(counterTCR.getId());
}
return eDeposit;
}
@ -351,46 +366,6 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
}
}
private void logError(IEnumId subject, Object... args) {
String err = msgResolver.resolve(new EnumMessage(subject, args));
log.info("{}", err);
// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
private void logError(Long sTradeNum, EnumMessage msg) {
String err = format("Сделка № %d не смогла быть обработана. %s", sTradeNum, msgResolver.resolve(msg));
log.warn(err);
notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
private void logNotifyW(String msg) {
logNotify(msg, Level.WARN);
}
private void logNotifyE(String msg) {
logNotify(msg, Level.ERROR);
}
private void logNotify(String msg, Level lvl) {
switch (lvl) {
case WARN -> log.warn(msg);
case ERROR -> log.error(msg);
default -> log.error("lvl not supported {}. {}", lvl, msg);
}
notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH);
}
private Optional<TradingClearingRegistry> searchTcrByStrades(STrades sTrades) {
IValidator counterValidator = valFor(sTrades);
Optional<EnumMessage> err = counterValidator.tillFirstError();
if (err.isPresent()) {
log.error("couldn't extract setCounterPartyTradingClearingRegistryId from strades.id={}", sTrades.getId());
return Optional.empty();
}
return Optional.ofNullable(counterValidator.getStored(ValidationStored.STradesTradingClearingRegistry));
}
@Override
public void initExecCash() {
ExecutionUploadCashUtil.loadExecutions(executionDepositImdg, cash);
@ -403,6 +378,7 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
context.addImdg(IMDGDistributedNames.Map_MoneyMarketSecurity, imdgMoneyMarketSecurity);
context.addImdg(IMDGDistributedNames.Map_Company, imdgCompany);
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
context.addImdg(IMDGDistributedNames.Map_STrades, sTradeImdg);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
new STradesSecurityPresentDepositValidationRule(valCash.mmsCash),
@ -412,7 +388,15 @@ public class ExecutionDepositComponent implements IExecutionUploadComponent {
new CompanyByTradingCodeCashingValidationRule<>(
STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCounterCompany
),
new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash)
new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash,
ImdgValidationContext::getValidatedObject,
ValidationStored.STradesCompany,
ValidationStored.STradesTradingClearingRegistry),
new CounterStradePresentValidationRule(sTrades),
new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash,
ctx -> ctx.getStoredObject(ValidationStored.STradesCounterTrade),
ValidationStored.STradesCounterCompany,
ValidationStored.STradesCounterTradingClearingRegistry)
);
}

View file

@ -5,14 +5,13 @@ import java.math.RoundingMode;
import java.time.Instant;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import static java.lang.String.format;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.event.Level;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.scheduling.annotation.EnableScheduling;
@ -45,9 +44,10 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.ExecutionType;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.clearing.service.validation.strades.CompanyByTradingCodeCashingValidationRule;
import ru.spcex.clearing.service.validation.strades.CounterStradePresentValidationRule;
import ru.spcex.clearing.service.validation.strades.STradesSecurityDGCTValidationRule;
import ru.spcex.clearing.service.validation.strades.STradesSecurityPresentValidationRule;
import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRule;
import ru.spcex.clearing.service.validation.strades.TcrByCodeAndCmpIsActiveCashingValidationRuleV2;
import ru.spcex.platform.enumeration.CurrencyCode;
import ru.spcex.platform.enumeration.InstrumentType;
import ru.spcex.platform.enumeration.ObjectType;
@ -65,8 +65,8 @@ import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByIdAndString;
import ru.spcex.platform.imdg.iml.hazelcast.adapter.CashV2ByString;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.imdg.validation.LogPrefixId;
import ru.spcex.platform.utils.collection.CollectionUtil;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.enumeration.IEnumId;
import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.enumeration.IMessageResolver;
import ru.spcex.platform.utils.log.ExceptionUtils;
@ -78,7 +78,7 @@ import ru.spcex.platform.utils.validation.ValidatorImpl;
@Component
@EnableScheduling
public class ExecutionFondComponent implements IExecutionUploadComponent {
public class ExecutionFondComponent extends ExecutionAbstractComponent {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<STrades> sTradeImdg;
@ -115,6 +115,7 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
@Qualifier("kafkaSenderWithoutRequestInfo") KafkaSender kafkaSender,
NotificationSender notifications,
IMessageResolver msgResolver) {
super(notifications, msgResolver);
this.sTradeImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_STrades, STrades.class);
this.executionFondImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ExecutionFond, ExecutionFond.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
@ -140,6 +141,11 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
resetTradingDay();
}
@Override
protected Logger log() {
return log;
}
/**
* Сбрасывать каждый день в 01:00:01 "0 1 0 1 * ?"
*/
@ -164,12 +170,14 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
}
log.debug("Start check new S_TRADE after {}, fromId={}", tradingDay, fromId);
//выбираем STrades на сегодня с правильным section
Collection<STrades> sTrades;
List<STrades> sTrades;
if (cleanLoad) {
sTrades = sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
sTrades = CollectionUtil.castOrCopyToModifyableList(
sTradeImdg.getCollectionObjectsByFieldValues(Map.of(
"tradeDate", LocalDate.now(),
"section", Section.FOND.getKey()
));
))
);
} else {
ImdgPredicateBuilder pb = sTradeImdg.predicateBuilder();
ImdgPredicate filter = pb.and(
@ -177,7 +185,9 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
pb.equals("section", Section.FOND.getKey()),
pb.greatEqual("id", fromId)
);
sTrades = sTradeImdg.getCollectionObjectsByPredicate(filter);
sTrades = CollectionUtil.castOrCopyToModifyableList(
sTradeImdg.getCollectionObjectsByPredicate(filter)
);
}
log.info("Found {} s_trade for today (fromId={})", sTrades.size(), fromId);
@ -209,14 +219,29 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
}
log.info("{} strades left after already-added filtering", sTrades.size());
//защита от бесконтрольно разрастающегося цикла
int iterationsWarn = iterationsWarn(sTrades);
int iterationsMax = iterationsError(sTrades);
Map<Long, ExecutionFond> execsToInsert = new HashMap<>();
for (STrades sTrd : sTrades) {
for (int i = 0; i < sTrades.size(); i++) {
if (cycleExceptionalInterrupt(i, iterationsWarn, iterationsMax)) {
return;
}
STrades sTrd = sTrades.get(i);
log.info("new S_TRADE[{}], valuation {}", sTrd.getId(), valuation);
//проверки
IValidator validator = valFor(sTrd);
Optional<EnumMessage> error = validator.tillFirstError();
STrades counter = validator.getStored(ValidationStored.STradesCounterTrade);
if (counter != null && fromId != null && counter.getId() < fromId) {
sTrades.add(counter);
}
if (error.isPresent()) {
if (error.get().getSubject().equals(ClearingError.CounterStradeNotFound)) {
log.info("STrade № {} id {} not processed. should be paired in next cycles.", sTrd.getTradeNum(), sTrd.getId());
continue;
}
String err = format("Сделка № %d не смогла быть обработана. %s", sTrd.getTradeNum(), msgResolver.resolve(error.get()));
logNotifyE(err);
continue;
@ -283,6 +308,7 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
Company company = validator.getStored(ValidationStored.STradesCompany);
Company counterCompany = validator.getStored(ValidationStored.STradesCounterCompany);
TradingClearingRegistry rgstr = validator.getStored(ValidationStored.STradesTradingClearingRegistry);
STrades counterSTrades = validator.getStored(ValidationStored.STradesCounterTrade);
ClientCode clientCode = clientCodeImdg.getFirstObjectByFieldValues(Map.of("code", sTrades.getClientCode()));
ExecutionFond eFond = new ExecutionFond();
@ -413,49 +439,10 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
eFond.setSettlementAmount(sTrades.getValue());
}
eFond.setCounterPartyId(counterCompany.getId());
{
ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder();
STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and(
strPb.equals("tradeNum", sTrades.getTradeNum()),
strPb.equals("section", sTrades.getSection()),
strPb.equals("classCode", sTrades.getClassCode()),
strPb.not(strPb.equals("operation", sTrades.getOperation()))));
if (counterSTrades != null) {
eFond.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount());
}
}
eFond.setCounterPartyTradingClearingRegistry(counterSTrades.getAccount());
return eFond;
}
private void logError(IEnumId subject, Object... args) {
String err = msgResolver.resolve(new EnumMessage(subject, args));
log.info("{}", err);
// notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
private void logError(Long sTradeId, EnumMessage msg) {
String err = format("Сделка № %d не смогла быть обработана. %s", sTradeId, msgResolver.resolve(msg));
log.warn(err);
notifications.sendNotification(ObjectType.vfrs, err, Priority.HIGH);
}
private void logNotifyW(String msg) {
logNotify(msg, Level.WARN);
}
private void logNotifyE(String msg) {
logNotify(msg, Level.ERROR);
}
private void logNotify(String msg, Level lvl) {
switch (lvl) {
case WARN -> log.warn(msg);
case ERROR -> log.error(msg);
default -> log.error("lvl not supported {}. {}", lvl, msg);
}
notifications.sendNotification(ObjectType.vfrs, msg, Priority.HIGH);
}
@Override
public void initExecCash() {
ExecutionUploadCashUtil.loadExecutions(executionFondImdg, cash);
@ -473,6 +460,7 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
context.addImdg(IMDGDistributedNames.Map_TradingClearingRegistry, imdgTradingClearingRegistry);
context.addImdg(IMDGDistributedNames.Map_DigitalCertificateSecurity, imdgDigitCertSecurity);
context.addImdg(IMDGDistributedNames.Map_UtilitarianDigitalRight, udrImdg);
context.addImdg(IMDGDistributedNames.Map_STrades, sTradeImdg);
context.setLogPrefix(LogPrefixId.INSTANCE);
return new ValidatorImpl<>(context,
new STradesSecurityPresentValidationRule(valCash.fixIncSecCash,
@ -486,9 +474,12 @@ public class ExecutionFondComponent implements IExecutionUploadComponent {
new CompanyByTradingCodeCashingValidationRule<>(
STrades::getCpFirmId, valCash.cmpCash, ValidationStored.STradesCounterCompany
),
new TcrByCodeAndCmpIsActiveCashingValidationRule(valCash.tcrCash),
new STradesSecurityDGCTValidationRule(valCash.udrCash)
);
new TcrByCodeAndCmpIsActiveCashingValidationRuleV2(valCash.tcrCash,
ImdgValidationContext::getValidatedObject,
ValidationStored.STradesCompany,
ValidationStored.STradesTradingClearingRegistry),
new STradesSecurityDGCTValidationRule(valCash.udrCash),
new CounterStradePresentValidationRule(sTrades));
}
private static class ExecFondCompCash extends CashCloser {

View file

@ -4,6 +4,8 @@ public enum ValidationStored {
STradesCompany, STradesCounterCompany, STradesSecurity, STradesTradingClearingRegistry,
Account, Company, STradesUtilitarianDigitalRight,
STradesCounterTrade, STradesCounterTradingClearingRegistry,
Sdf57CompanyDeb, Sdf57CompanyCred, Sdf57AccountDeb, Sdf57AccountCred,
Sdf01Company, Sdf01Account,

View file

@ -0,0 +1,44 @@
package ru.spcex.clearing.service.validation.strades;
import java.util.Optional;
import ru.clearing.classes.statics.data.misc.STrades;
import ru.spcex.clearing.error.ClearingError;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.service.validation.ValidationStored;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder;
import ru.spcex.platform.imdg.validation.ImdgValidationContext;
import ru.spcex.platform.utils.enumeration.EnumMessage;
import ru.spcex.platform.utils.validation.IValidationRule;
public class CounterStradePresentValidationRule implements IValidationRule<ImdgValidationContext<STrades>> {
private final STrades sTrades;
public CounterStradePresentValidationRule(STrades sTrades) {
this.sTrades = sTrades;
}
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
Imdg<STrades> sTradeImdg = context.obtainMap(IMDGDistributedNames.Map_STrades, STrades.class);
ImdgPredicateBuilder strPb = sTradeImdg.predicateBuilder();
STrades counterSTrades = sTradeImdg.getFirstObjectByPredicate(strPb.and(
strPb.equals("tradeNum", sTrades.getTradeNum()),
strPb.equals("section", sTrades.getSection()),
strPb.equals("classCode", sTrades.getClassCode()),
strPb.not(strPb.equals("operation", sTrades.getOperation()))));
if (counterSTrades != null) {
context.storeObject(ValidationStored.STradesCounterTrade, counterSTrades);
return empty();
}
return of(ClearingError.CounterStradeNotFound);
}
@Override
public String ruleName() {
return "counter-strades-present";
}
}

View file

@ -1,6 +1,7 @@
package ru.spcex.clearing.service.validation.strades;
import java.util.Optional;
import java.util.function.Function;
import ru.clearing.classes.statics.data.company.Company;
import ru.clearing.classes.statics.data.misc.STrades;
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
@ -18,20 +19,29 @@ import ru.spcex.platform.utils.enumeration.IEnumKey;
import ru.spcex.platform.utils.text.TextUtil;
import ru.spcex.platform.utils.validation.IValidationRule;
public class TcrByCodeAndCmpIsActiveCashingValidationRule implements IValidationRule<ImdgValidationContext<STrades>> {
public class TcrByCodeAndCmpIsActiveCashingValidationRuleV2 implements IValidationRule<ImdgValidationContext<STrades>> {
private final CashV2ByIdAndString<TradingClearingRegistry> tcrCash;
public TcrByCodeAndCmpIsActiveCashingValidationRule(CashV2ByIdAndString<TradingClearingRegistry> tcrCash) {
public TcrByCodeAndCmpIsActiveCashingValidationRuleV2(CashV2ByIdAndString<TradingClearingRegistry> tcrCash, Function<ImdgValidationContext<STrades>, STrades> tradeExtraction, Enum<ValidationStored> cmpEnum, Enum<ValidationStored> tcrEnum) {
this.tcrCash = tcrCash;
this.tradeExtraction = tradeExtraction;
this.cmpEnum = cmpEnum;
this.tcrEnum = tcrEnum;
}
private final Function<ImdgValidationContext<STrades>, STrades> tradeExtraction;
private final Enum<ValidationStored> cmpEnum;
private final Enum<ValidationStored> tcrEnum;
@Override
public Optional<EnumMessage> validate(ImdgValidationContext<STrades> context) {
STrades validatedObject = context.getValidatedObject();
STrades validatedObject = tradeExtraction.apply(context);
if (TextUtil.isEmpty(validatedObject.getAccount())) {
return of(ClearingError.TradingClearingRegistryNotFound, "", validatedObject.getAccount());
}
if (context.getStoredObject(ValidationStored.STradesCompany) == null) { //maybe unnecessary
if (context.getStoredObject(cmpEnum) == null) { //maybe unnecessary
return of(ClearingError.CompanyNotFound, validatedObject.getSecCode());
}
Imdg<TradingClearingRegistry> tcrImdg = context.obtainMap(

View file

@ -0,0 +1,16 @@
package ru.spcex.platform.utils.collection;
import java.util.ArrayList;
import java.util.Collection;
import java.util.LinkedList;
import java.util.List;
public class CollectionUtil {
public static <T> List<T> castOrCopyToModifyableList(Collection<T> c) {
if (c instanceof ArrayList<T> || c instanceof LinkedList<T>) {
return (List<T>) c;
}
return new ArrayList<>(c);
}
}