This commit is contained in:
AKurakin 2023-06-09 20:52:10 +03:00
parent c94fe7d4db
commit f1cb711326
6 changed files with 350 additions and 5 deletions

View file

@ -74,7 +74,7 @@ public class EquitySecurityService extends QueueConsumer implements Initializing
init();
}
private RequestInfoUpdate newEquity(BaseRequest<EquitySecurityNewRequest> userRequest) {
protected synchronized RequestInfoUpdate newEquity(BaseRequest<EquitySecurityNewRequest> userRequest) {
ImdgTransaction transaction = imdgProvider.newTransaction();
EquitySecurityNewRequest req = userRequest.getRequestPayload();
@ -123,7 +123,7 @@ public class EquitySecurityService extends QueueConsumer implements Initializing
return null; // default success
}
private RequestInfoUpdate updateEquity(BaseRequest<EquitySecurityUpdateRequest> userRequest) {
protected synchronized RequestInfoUpdate updateEquity(BaseRequest<EquitySecurityUpdateRequest> userRequest) {
EquitySecurityUpdateRequest req = userRequest.getRequestPayload();
log.debug("MoneyMarketSecurityUpdateRequest received id = {}", req.getId());

View file

@ -73,7 +73,7 @@ public class FixedIncomeSecurityService extends QueueConsumer implements Initial
init();
}
private RequestInfoUpdate newFixedIncome(BaseRequest<FixedIncomeSecurityNewRequest> userRequest) {
protected synchronized RequestInfoUpdate newFixedIncome(BaseRequest<FixedIncomeSecurityNewRequest> userRequest) {
ImdgTransaction transaction = imdgProvider.newTransaction();
FixedIncomeSecurityNewRequest req = userRequest.getRequestPayload();
@ -127,7 +127,7 @@ public class FixedIncomeSecurityService extends QueueConsumer implements Initial
return null; // default success
}
private RequestInfoUpdate updateFixedIncome(BaseRequest<FixedIncomeSecurityUpdateRequest> userRequest) {
protected synchronized RequestInfoUpdate updateFixedIncome(BaseRequest<FixedIncomeSecurityUpdateRequest> userRequest) {
FixedIncomeSecurityUpdateRequest req = userRequest.getRequestPayload();
log.debug("updateFixedIncome received id = {}", req.getId());

View file

@ -0,0 +1,277 @@
package ru.spcex.clearing.securities.service.cud;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import ru.clearing.classes.statics.data.instrument.issue.EquitySecurity;
import ru.clearing.classes.statics.data.instrument.issue.FixedIncomeSecurity;
import ru.clearing.classes.statics.data.misc.Listing;
import ru.clearing.classes.statics.data.security.Security;
import ru.spcex.clearing.imdg.IMDGDistributedNames;
import ru.spcex.clearing.platform.messaging.domain.ActionType;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonDeleteRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.securitites.*;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.securities.validation.ValidationProvider;
import ru.spcex.clearing.util.security.UserRoleVerification;
import ru.spcex.clearing.validation.common.ValidationHelper;
import ru.spcex.platform.enumeration.UserRole;
import ru.spcex.platform.enumeration.WorkflowStatus;
import ru.spcex.platform.imdg.api.Imdg;
import ru.spcex.platform.imdg.api.ImdgId;
import ru.spcex.platform.imdg.api.ImdgProvider;
import ru.spcex.platform.imdg.api.ImdgTransaction;
import java.time.Instant;
import java.util.Collection;
import java.util.Map;
@Service
public class MultiSecurityService extends QueueConsumer implements InitializingBean {
private final Logger log = LoggerFactory.getLogger(getClass());
private final Imdg<FixedIncomeSecurity> fixedIncomeSecurityImdg;
private final Imdg<EquitySecurity> equitySecurityImdg;
private final Imdg<Listing> listingImdg;
private final ImdgProvider imdgProvider;
private final ImdgId idGenerator;
FixedIncomeSecurityService fixedIncomeSecurityService;
EquitySecurityService equitySecurityService;
ListingService listingService;
FixedIncomeCashFlowService fixedIncomeCashFlowService;
private final ValidationProvider validation;
private final UserRoleVerification userRoleVerification;
private final ValidationHelper validationHelper;
@Autowired
public MultiSecurityService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
ImdgProvider imdgProvider,
ValidationProvider validation,
UserRoleVerification userRoleVerification,
ValidationHelper validationHelper) {
super(kafkaQueue, kafkaProducer);
this.validation = validation;
this.fixedIncomeSecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_FixedIncomeSecurity, FixedIncomeSecurity.class);
this.equitySecurityImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_EquitySecurity, EquitySecurity.class);
this.listingImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Listing, Listing.class);
this.imdgProvider = imdgProvider;
this.idGenerator = imdgProvider.getImdgIdGenerator();
this.userRoleVerification = userRoleVerification.setRoleForVerification(UserRole.Admin);
this.validationHelper = validationHelper;
}
@Override
public void afterPropertiesSet() {
callback(MultiSecurityRequest.class)
.setFunction(this::newSecurities)
.forDestination(Consts.DESTINATION_SECUIRTY_MULTIREQUEST, callbacks::put);
imdgProvider.waitAvailable();
init();
}
<T extends Security> T findSecurityBySecuritySymbol(Imdg<T> imdg, String securitySymbol) {
Collection<T> securities = imdg.getCollectionObjectsByFieldValues(
Map.of("securitySymbol", securitySymbol)
);
if (securities.isEmpty()) {
log.trace("{} not found by securitySymbol={}", imdg.getMapName(), securitySymbol);
return null;
} else {
if (securities.size() > 1)
log.warn("Found {} {} by securitySymbol={}", securities.size(), imdg.getMapName(), securitySymbol);
T security = securities.iterator().next();
log.trace("Found {}[{}] with securitySymbol={}", imdg.getMapName(), security.getId(), securitySymbol);
return security;
}
}
private RequestInfoUpdate newSecurities(BaseRequest<MultiSecurityRequest> userRequest) {
log.debug("MultiSecurityRequest id={} received", userRequest.getId());
MultiSecurityRequest req = userRequest.getRequestPayload();
if (req.getFixedIncomeSecurityNewRequest() == null && req.getEquitySecurityNewRequest() == null) {
log.error("Required one of sub-request: FixedIncomeSecurityNewRequest or FixedIncomeSecurityNewRequest");
} else if (req.getFixedIncomeSecurityNewRequest() != null && req.getEquitySecurityNewRequest() != null) {
log.error("Filling both one of sub-request: FixedIncomeSecurityNewRequest or getEquitySecurityNewRequest");
}
;
final Long companyId = req.getCompanyId();
// ImdgTransaction transaction = imdgProvider.newTransaction(); todo transaction
Security security = null;
if (req.getFixedIncomeSecurityNewRequest() != null) {
log.trace("FixedIncomeSecurityNewRequest - BOND");
String securitySymbol = req.getFixedIncomeSecurityNewRequest().getSecuritySymbol();
security = findSecurityBySecuritySymbol(fixedIncomeSecurityImdg, securitySymbol);
if (security == null) {
RequestInfoUpdate resp = fixedIncomeSecurityService.newFixedIncome(wrapRequest(userRequest, req.getFixedIncomeSecurityNewRequest(), null));
validateReply(companyId, "newFixedIncome", resp);
security = findSecurityBySecuritySymbol(fixedIncomeSecurityImdg, securitySymbol);
} else {
FixedIncomeSecurityUpdateRequest updateRequest = toUpdateRequest(security, req.getFixedIncomeSecurityNewRequest());
RequestInfoUpdate resp = fixedIncomeSecurityService.updateFixedIncome(wrapRequest(userRequest, updateRequest, ActionType.UPDATE));
validateReply(companyId, "newFixedIncome", resp);
}
}
if (req.getEquitySecurityNewRequest() != null) {
log.trace("EquitySecurityNewRequest - FOND");
String securitySymbol = req.getFixedIncomeSecurityNewRequest().getSecuritySymbol();
security = findSecurityBySecuritySymbol(equitySecurityImdg, securitySymbol);
if (security == null) {
RequestInfoUpdate resp = equitySecurityService.newEquity(wrapRequest(userRequest, req.getEquitySecurityNewRequest(), null));
validateReply(companyId, "newFixedIncome", resp);
security = findSecurityBySecuritySymbol(equitySecurityImdg, securitySymbol);
} else {
EquitySecurityUpdateRequest updateRequest = toUpdateRequest(security, req.getEquitySecurityNewRequest());
RequestInfoUpdate resp = equitySecurityService.updateEquity(wrapRequest(userRequest, updateRequest, ActionType.UPDATE));
validateReply(companyId, "newFixedIncome", resp);
}
}
log.error("TODO listing,couponSchedule,nominal");
/*
todo оставшиеся части:
for (IncomeListing listing : security.getListingList()) {
if (listing.isInvalidData()) continue;
if (!security.isAlreadyExist() || !listing.isAlreadyExist()) {
ListingNewRequest listingNewRequest = new ListingNewRequest();
listingNewRequest.setSecurityId(securityId);
listingNewRequest.setMarket(listing.getCode());
listingNewRequest.setLotSize(listing.getLotSize());
listingNewRequest.setTradingCurrency(listing.getTradingCurrency());
listingNewRequest.setWorkflowStatus(listing.getWorkflowStatus());
kafkaSender.sendRequestToQueue(Consts.LISTING_NEW, listingNewRequest);
} else {
ListingUpdateRequest listingUpdateRequest = new ListingUpdateRequest();
listingUpdateRequest.setId(listing.getMapId());
listingUpdateRequest.setSecurityId(securityId);
listingUpdateRequest.setMarket(listing.getCode());
listingUpdateRequest.setLotSize(listing.getLotSize());
listingUpdateRequest.setTradingCurrency(listing.getTradingCurrency());
listingUpdateRequest.setWorkflowStatus(listing.getWorkflowStatus());
kafkaSender.sendRequestToQueue(Consts.LISTING_UPDATE, listingUpdateRequest);
}
}
for (CouponSchedule couponSchedule : security.getCouponScheduleList()) {
if (couponSchedule.isInvalidData()) continue;
if (!security.isAlreadyExist() || !couponSchedule.isAlreadyExist()) {
CouponPeriodNewRequest couponPeriodNewRequest = new CouponPeriodNewRequest();
couponPeriodNewRequest.setSecurityId(securityId);
couponPeriodNewRequest.setCouponRate(couponSchedule.getCouponRate());
couponPeriodNewRequest.setNumber(couponSchedule.getCouponNumber());
couponPeriodNewRequest.setPeriodEndDate(couponSchedule.getPeriodEndDate());
couponPeriodNewRequest.setPeriodStartDate(couponSchedule.getPeriodStartDate());
kafkaSender.sendRequestToQueue(Consts.DESTINATION_COUPON_PERIOD_NEW, couponPeriodNewRequest);
} else {
CouponPeriodUpdateRequest couponPeriodUpdateRequest = new CouponPeriodUpdateRequest();
couponPeriodUpdateRequest.setId(couponSchedule.getMapId());
couponPeriodUpdateRequest.setSecurityId(securityId);
couponPeriodUpdateRequest.setCouponRate(couponSchedule.getCouponRate());
couponPeriodUpdateRequest.setNumber(couponSchedule.getCouponNumber());
couponPeriodUpdateRequest.setPeriodEndDate(couponSchedule.getPeriodEndDate());
couponPeriodUpdateRequest.setPeriodStartDate(couponSchedule.getPeriodStartDate());
kafkaSender.sendRequestToQueue(Consts.DESTINATION_COUPON_PERIOD_UPDATE, couponPeriodUpdateRequest);
}
}
for (Nominal nominal : security.getNominalList()) {
if (nominal.isInvalidData()) continue;
if (!security.isAlreadyExist() || !nominal.isAlreadyExist()) {
FixedIncomeCashFlowNewRequest fixedIncomeCashFlowNewRequest = new FixedIncomeCashFlowNewRequest();
fixedIncomeCashFlowNewRequest.setSecuritySymbol(security.getSecuritySymbol());
fixedIncomeCashFlowNewRequest.setAccruedCoupon(nominal.getAccruedCoupon());
fixedIncomeCashFlowNewRequest.setNominalValue(nominal.getNominal());
fixedIncomeCashFlowNewRequest.setNumber(nominal.getCouponNumber());
kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_CASH_FLOW_NEW, fixedIncomeCashFlowNewRequest);
} else {
FixedIncomeCashFlowUpdateRequest fixedIncomeCashFlowUpdateRequest = new FixedIncomeCashFlowUpdateRequest();
fixedIncomeCashFlowUpdateRequest.setId(nominal.getMapId());
fixedIncomeCashFlowUpdateRequest.setSecuritySymbol(security.getSecuritySymbol());
fixedIncomeCashFlowUpdateRequest.setAccruedCoupon(nominal.getAccruedCoupon());
fixedIncomeCashFlowUpdateRequest.setNominalValue(nominal.getNominal());
fixedIncomeCashFlowUpdateRequest.setNumber(nominal.getCouponNumber());
kafkaSender.sendRequestToQueue(Consts.DESTINATION_FIXED_INCOME_CASH_FLOW_UPDATE, fixedIncomeCashFlowUpdateRequest);
}
}
*/
log.debug("successfully processed, request id {}", userRequest.getId());
return null; // default success
}
FixedIncomeSecurityUpdateRequest toUpdateRequest(Security security, FixedIncomeSecurityNewRequest newRequest) {
FixedIncomeSecurityUpdateRequest updateRequest = new FixedIncomeSecurityUpdateRequest();
updateRequest.setId(security.getId());
updateRequest.setSecuritySymbol(newRequest.getSecuritySymbol());
updateRequest.setShortName(newRequest.getShortName());
updateRequest.setFullName(newRequest.getFullName());
updateRequest.setIsin(newRequest.getIsin());
updateRequest.setBondType(newRequest.getBondType());
updateRequest.setLotSize(newRequest.getLotSize());
updateRequest.setNominalValue(newRequest.getNominalValue());
updateRequest.setNominalCurrency(newRequest.getNominalCurrency());
updateRequest.setMaturityDate(newRequest.getMaturityDate());
updateRequest.setCoupon(newRequest.getCoupon());
updateRequest.setCouponFrequency(newRequest.getCouponFrequency());
updateRequest.setIssuerId(newRequest.getIssuerId());
updateRequest.setShortNameEng(newRequest.getShortNameEng());
updateRequest.setFullNameEng(newRequest.getFullNameEng());
updateRequest.setWorkflowStatus(newRequest.getWorkflowStatus());
updateRequest.setInstrumentType(newRequest.getInstrumentType());
return updateRequest;
}
EquitySecurityUpdateRequest toUpdateRequest(Security security, EquitySecurityNewRequest newRequest) {
EquitySecurityUpdateRequest updateRequest = new EquitySecurityUpdateRequest();
updateRequest.setId(security.getId());
updateRequest.setSecuritySymbol(newRequest.getSecuritySymbol());
updateRequest.setShortName(newRequest.getShortName());
updateRequest.setFullName(newRequest.getFullName());
updateRequest.setIsin(newRequest.getIsin());
updateRequest.setShareType(newRequest.getShareType());
updateRequest.setLotSize(newRequest.getLotSize());
updateRequest.setIssuerId(newRequest.getIssuerId());
updateRequest.setShortNameEng(newRequest.getShortNameEng());
updateRequest.setFullNameEng(newRequest.getFullNameEng());
updateRequest.setWorkflowStatus(newRequest.getWorkflowStatus());
updateRequest.setInstrumentType(newRequest.getInstrumentType());
return updateRequest;
}
private <T> BaseRequest<T> wrapRequest(BaseRequest<?> template, T payload, ActionType action) {
BaseRequest<T> r = new BaseRequest<>();
r.setId(template.getId());
if (action == null) {
r.setActionType(template.getActionType());
} else {
r.setActionType(action);
}
r.setUserId(template.getUserId());
r.setCorrelationId(template.getCorrelationId());
r.setRequestPayload(payload);
return r;
}
private void validateReply(Long companyId, String process, RequestInfoUpdate replyI) {
if (replyI != null && replyI.getMessage() != null) {
log.warn("Error process {} for companyId={}: {}", process, companyId, replyI.getMessage());
}
}
}

View file

@ -13,6 +13,8 @@ public interface Consts {
String DESTINATION_FIXED_INCOME_SECURITY_UPDATE = "fixed-income-security-update";
String DESTINATION_FIXED_INCOME_SECURITY_DELETE = "fixed-income-security-delete";
String DESTINATION_SECUIRTY_MULTIREQUEST = "secuirty-multirequest-new";
String DESTINATION_FIXED_INCOME_CASH_FLOW_NEW = "fixed-income-cash-flow-new";
String DESTINATION_FIXED_INCOME_CASH_FLOW_UPDATE = "fixed-income-cash-flow-update";

View file

@ -8,7 +8,7 @@ import java.util.List;
public class MultiCompanyRequest {
@JsonProperty
String uuid; // для лога
String uuid;
@JsonProperty
CompanyNewRequest company;

View file

@ -0,0 +1,66 @@
package ru.spcex.clearing.platform.messaging.domain.cud.securitites;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.util.ArrayList;
import java.util.List;
public class MultiSecurityRequest {
@JsonProperty
Long companyId;
FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest; // BOND
EquitySecurityNewRequest equitySecurityNewRequest; // FOND
List<ListingNewRequest> listingNewRequests = new ArrayList<>();
List<CouponPeriodNewRequest> couponPeriodNewRequests = new ArrayList<>();
List<FixedIncomeCashFlowNewRequest> fixedIncomeCashFlowNewRequests = new ArrayList<>();
public Long getCompanyId() {
return companyId;
}
public void setCompanyId(Long companyId) {
this.companyId = companyId;
}
public FixedIncomeSecurityNewRequest getFixedIncomeSecurityNewRequest() {
return fixedIncomeSecurityNewRequest;
}
public void setFixedIncomeSecurityNewRequest(FixedIncomeSecurityNewRequest fixedIncomeSecurityNewRequest) {
this.fixedIncomeSecurityNewRequest = fixedIncomeSecurityNewRequest;
}
public EquitySecurityNewRequest getEquitySecurityNewRequest() {
return equitySecurityNewRequest;
}
public void setEquitySecurityNewRequest(EquitySecurityNewRequest equitySecurityNewRequest) {
this.equitySecurityNewRequest = equitySecurityNewRequest;
}
public List<ListingNewRequest> getListingNewRequests() {
return listingNewRequests;
}
public void setListingNewRequests(List<ListingNewRequest> listingNewRequests) {
this.listingNewRequests = listingNewRequests;
}
public List<CouponPeriodNewRequest> getCouponPeriodNewRequests() {
return couponPeriodNewRequests;
}
public void setCouponPeriodNewRequests(List<CouponPeriodNewRequest> couponPeriodNewRequests) {
this.couponPeriodNewRequests = couponPeriodNewRequests;
}
public List<FixedIncomeCashFlowNewRequest> getFixedIncomeCashFlowNewRequests() {
return fixedIncomeCashFlowNewRequests;
}
public void setFixedIncomeCashFlowNewRequests(List<FixedIncomeCashFlowNewRequest> fixedIncomeCashFlowNewRequests) {
this.fixedIncomeCashFlowNewRequests = fixedIncomeCashFlowNewRequests;
}
}