From 79c56c457180d1e8ee0a88bd1df0c08a6ee69c5e Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 12:14:22 +0300 Subject: [PATCH 1/6] =?UTF-8?q?=D0=91=D0=BE=D0=BB=D1=8C=D1=88=D0=B5=20?= =?UTF-8?q?=D0=BB=D0=BE=D0=B3=D0=BE=D0=B2=20SDF08...?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../java/ru/spcex/clearing/service/StatementService.java | 6 +++++- .../clearing/swt/importer/logic/stages/ImportToDB.java | 3 +++ 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index 2189b0065..d8ec4bde6 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -68,9 +68,9 @@ public class StatementService extends QueueConsumer implements InitializingBean } private void process(BaseRequest systemRequest) { - log.debug("Receiving StatementRequest id={}", systemRequest.getId()); StatementRequest statementRequest = systemRequest.getRequestPayload(); SdfTable table = statementRequest.getTable(); + log.debug("Receiving StatementRequest id={}; table {}", systemRequest.getId(), table); Optional completePairKey = saveRequest(statementRequest); if (completePairKey.isPresent()) { @@ -133,6 +133,8 @@ public class StatementService extends QueueConsumer implements InitializingBean AbstractExecutor service = executorsMap.get(SdfTable.SDF_04); if (service != null) { service.execute(sdfGroup, statementRequest); + } else { + log.warn("Executor for SDF_04 not set"); } } @@ -143,6 +145,8 @@ public class StatementService extends QueueConsumer implements InitializingBean AbstractExecutor service = executorsMap.get(SdfTable.SDF_08); if (service != null) { service.execute(sdfGroup, statementRequest); + } else { + log.warn("Executor for SDF_08 not set"); } } diff --git a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java index 2b1d634d8..3e42efde8 100644 --- a/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java +++ b/clearing-parent/swt-importer/src/main/java/ru/spcex/clearing/swt/importer/logic/stages/ImportToDB.java @@ -54,6 +54,7 @@ public class ImportToDB extends Stage { Long fileId = hazelcastService.getImdgIdGenerator().nextId(); table.setFileId(fileId); int counter = 0; + int storedCounter = 0; SWTHeaderData swtHeaderData = swtReader.getSWTHeader(); do { counter++; @@ -63,7 +64,9 @@ public class ImportToDB extends Stage { continue; } table.injectEntity(table.getEntity(swtHeaderData, entity)); + storedCounter++; } while (swtReader.hasNextRecord()); + log.debug("{} record read; {} entity stored.", counter, storedCounter); kafkaMessenger.notifySystemIfNeeded(currTable, fileId); } catch (IOException exception) { log.warn(exception.getMessage()); From 5cffd474639b432b99ce6893fdd929bf9f0244f8 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 13:44:10 +0300 Subject: [PATCH 2/6] =?UTF-8?q?company-service=20http://jira.mfd.msk:8088/?= =?UTF-8?q?browse/CLS-347=20MultiCompanyService=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=B8=D0=BB=20=D0=BE=D0=B1=D0=BD=D0=BE=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=D0=B8=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../company/service/MultiCompanyService.java | 155 +++++++++++++-- .../service/ProfileDocumentService.java | 2 +- .../service/MultiCompanyServiceTest.java | 184 ++++++++++++++++++ .../clearing/test/TestObjectCreator.java | 14 ++ 4 files changed, 336 insertions(+), 19 deletions(-) create mode 100644 clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/MultiCompanyServiceTest.java diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/MultiCompanyService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/MultiCompanyService.java index 360fa37f4..0511f22f9 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/MultiCompanyService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/MultiCompanyService.java @@ -10,7 +10,10 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.company.CompanySymbols; +import ru.clearing.classes.statics.data.profile.Contact; +import ru.clearing.classes.statics.data.profile.ProfileDocument; 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.account.ClientCodeNewRequest; @@ -21,7 +24,6 @@ import ru.spcex.clearing.util.security.UserRoleVerification; import ru.spcex.clearing.util.services.RequestHelper; import ru.spcex.clearing.validation.common.ValidationHelper; import ru.spcex.platform.enumeration.CompanySymbol; -import ru.spcex.platform.enumeration.UserRole; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.utils.error.ValidationException; @@ -29,6 +31,7 @@ import ru.spcex.platform.utils.log.ExceptionUtils; import java.util.ArrayList; import java.util.Collection; +import java.util.HashMap; import java.util.Map; /** @@ -55,6 +58,8 @@ public class MultiCompanyService final ImdgProvider imdgProvider; final Imdg companyIMap; final Imdg companySymbolsImdg; + final Imdg profileDocumentImdg; + final Imdg contactImdg; @Autowired public MultiCompanyService(Consumer kafkaQueue, Producer kafkaProducer, @@ -88,6 +93,8 @@ public class MultiCompanyService companyIMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Company, Company.class); companySymbolsImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); + profileDocumentImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ProfileDocument, ProfileDocument.class); + contactImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Contact, Contact.class); } @Override @@ -106,35 +113,62 @@ public class MultiCompanyService if (requestInfoUpdate != null) return requestInfoUpdate; } - /* - todo: - 1) добавить апдейт кейсы - 2) сделать транзакции. - - */ Long companyId = null; - boolean txOk = false; + boolean txOk = false; // todo IMDG transaction synchronized (companyService) { try { companyId = getCompanyIdForCompanySymbols(req.getUuid(), req.getCompany()); log.debug("For request {} (uuid {}) company {}.", baseRequest.getId(), req.getUuid(), companyId == null ? "not found" : ("found, id=" + companyId)); RequestInfoUpdate replyI; - Company company = companyService.createOrUpdateCompany(wrapRequest(baseRequest, req.getCompany()), companyId); + Company company = companyService.createOrUpdateCompany(wrapRequest(baseRequest, req.getCompany(), null), companyId); companyId = company.getId(); log.debug("The companyId={}", companyId); fillCompanyId(req, companyId); - //todo update case: - replyI = companyInfoService.companyInfoUpdate(wrapRequest(baseRequest, req.getCompanyInfo())); + if (req.getCompanyInfo() != null) { + log.trace("For company {} do update CompanyInfo", companyId); + replyI = companyInfoService.companyInfoUpdate(wrapRequest(baseRequest, req.getCompanyInfo(), null)); + validateReply(companyId, "CompanyInfo", replyI); + } for (ProfileDocumentNewRequest partRequest : req.getProfileDocuments()) { - replyI = profileDocumentService.profileDocumentNew(wrapRequest(baseRequest, partRequest)); + ProfileDocument existDocument = findCompanyDocument(partRequest); + if (existDocument == null) { + log.trace("For company[{}] do new profileDocument", companyId); + replyI = profileDocumentService.profileDocumentNew(wrapRequest(baseRequest, partRequest, ActionType.NEW)); + validateReply(companyId, "profileDocumentNew", replyI); + } else { + log.trace("For company[{}] do update profileDocument[{}]", companyId, existDocument.getId()); + ProfileDocumentUpdateRequest partUpdateRequest = createUpdateRequest(existDocument, partRequest); + replyI = profileDocumentService.profileDocumentUpdate(wrapRequest(baseRequest, partUpdateRequest, ActionType.UPDATE)); + validateReply(companyId, "profileDocumentUpdate", replyI); + } } for (CompanySymbolNewRequest partRequest : req.getCompanySymbols()) { - replyI = companySymbolService.companySymbolNew(wrapRequest(baseRequest, partRequest)); + CompanySymbols existSymbol = findCompanySymbols(partRequest); + if (existSymbol == null) { + log.trace("For company[{}] do new companySymbol", companyId); + replyI = companySymbolService.companySymbolNew(wrapRequest(baseRequest, partRequest, ActionType.NEW)); + validateReply(companyId, "companySymbolNew", replyI); + } else { + log.trace("For company[{}] do update companySymbol[{}]", companyId, existSymbol.getId()); + CompanySymbolUpdateRequest partUpdateRequest = createUpdateRequest(existSymbol, partRequest); + replyI = companySymbolService.companySymbolUpdate(wrapRequest(baseRequest, partUpdateRequest, ActionType.UPDATE)); + validateReply(companyId, "companySymbolUpdate", replyI); + } } for (ContactNewRequest partRequest : req.getContacts()) { - replyI = contactService.contactNew(wrapRequest(baseRequest, partRequest)); + Contact existContact = findContact(partRequest); + if (existContact == null) { + log.trace("For company[{}] do new contact", companyId); + replyI = contactService.contactNew(wrapRequest(baseRequest, partRequest, ActionType.NEW)); + validateReply(companyId, "contactNew", replyI); + } else { + log.trace("For company[{}] do update contact[{}]", companyId, existContact.getId()); + ContactUpdateRequest partUpdateRequest = createUpdateRequest(existContact, partRequest); + replyI = contactService.contactUpdate(wrapRequest(baseRequest, partUpdateRequest, ActionType.UPDATE)); + validateReply(companyId, "contactUpdate", replyI); + } } txOk = true; } catch (ValidationException vex) { @@ -151,6 +185,83 @@ public class MultiCompanyService return null; } + void validateReply(Long companyId, String process, RequestInfoUpdate replyI) { + if (replyI != null && replyI.getMessage() != null) { + log.warn("Error process {} for companyId={}: {}", process, companyId, replyI.getMessage()); + } + } + + ProfileDocument findCompanyDocument(ProfileDocumentNewRequest partRequest) { + Map> query = new HashMap<>(); + query.put("companyId", partRequest.getCompanyId()); + query.put("documentType", partRequest.getDocumentType()); + query.put("issueDate", partRequest.getIssueDate()); // может быть null + Collection allDoc = profileDocumentImdg.getCollectionObjectsByFieldValues(query); + if (allDoc.isEmpty()) + return null; + if (allDoc.size() > 1) + log.warn("Fount {} ProfileDocument by: {}", allDoc.size(), query); + return allDoc.iterator().next(); + } + + CompanySymbols findCompanySymbols(CompanySymbolNewRequest partRequest) { + Map> query = new HashMap<>(); + query.put("companyId", partRequest.getCompanyId()); + query.put("companySymbol", partRequest.getCompanySymbol()); + Collection allCS = companySymbolsImdg.getCollectionObjectsByFieldValues(query); + if (allCS.isEmpty()) + return null; + if (allCS.size() > 1) + log.warn("Fount {} CompanySymbols by: {}", allCS.size(), query); + return allCS.iterator().next(); + } + + Contact findContact(ContactNewRequest partRequest) { + Map> query = new HashMap<>(); + query.put("companyId", partRequest.getCompanyId()); + query.put("contactType", partRequest.getContactType()); + Collection allCS = contactImdg.getCollectionObjectsByFieldValues(query); + if (allCS.isEmpty()) + return null; + if (allCS.size() > 1) + log.warn("Fount {} Contact by: {}", allCS.size(), query); + return allCS.iterator().next(); + } + + ProfileDocumentUpdateRequest createUpdateRequest(ProfileDocument existDocument, ProfileDocumentNewRequest partRequest) { + ProfileDocumentUpdateRequest r = new ProfileDocumentUpdateRequest(); + r.setId(existDocument.getId()); + r.setCompanyId(partRequest.getCompanyId()); + r.setDocumentType(partRequest.getDocumentType()); + r.setIssueDate(partRequest.getIssueDate()); + r.setIssuePlace(partRequest.getIssuePlace()); + r.setIssuer(partRequest.getIssuer()); + r.setIssuerCode(partRequest.getIssuerCode()); + r.setName(partRequest.getName()); + r.setNumber(partRequest.getNumber()); + r.setPlace(partRequest.getPlace()); + r.setValidFromDate(partRequest.getValidFromDate()); + r.setValidToDate(partRequest.getValidToDate()); + r.setLink(partRequest.getLink()); + return r; + } + + CompanySymbolUpdateRequest createUpdateRequest(CompanySymbols existSymbol, CompanySymbolNewRequest partRequest) { + CompanySymbolUpdateRequest r = new CompanySymbolUpdateRequest(); + r.setId(existSymbol.getId()); + r.setCompanyId(partRequest.getCompanyId()); + r.setCompanySymbol(partRequest.getCompanySymbol()); + r.setCompanySymbolValue(partRequest.getCompanySymbolValue()); + return r; + } + + ContactUpdateRequest createUpdateRequest(Contact existContact, ContactNewRequest partRequest) { + ContactUpdateRequest r = new ContactUpdateRequest(); + r.setId(existContact.getId()); + r.setContactType(partRequest.getContactType()); + r.setContactValue(partRequest.getContactValue()); + return r; + } Long getCompanyIdForCompanySymbols(String uuid, CompanyNewRequest cnr) { CompanySymbols companySymbol = null; @@ -190,10 +301,14 @@ public class MultiCompanyService return null; } - private BaseRequest wrapRequest(BaseRequest template, T payload) { + private BaseRequest wrapRequest(BaseRequest template, T payload, ActionType action) { BaseRequest r = new BaseRequest<>(); r.setId(template.getId()); - r.setActionType(template.getActionType()); // todo? + if (action == null) { + r.setActionType(template.getActionType()); + } else { + r.setActionType(action); + } r.setUserId(template.getUserId()); r.setCorrelationId(template.getCorrelationId()); r.setRequestPayload(payload); @@ -201,8 +316,12 @@ public class MultiCompanyService } void fillCompanyId(MultiCompanyRequest req, Long companyId) { - req.getCompany().setId(companyId); - req.getCompanyInfo().setId(companyId); + if (req.getCompany() != null) { + req.getCompany().setId(companyId); + } + if (req.getCompanyInfo() != null) { + req.getCompanyInfo().setId(companyId); + } if (req.getCompanySymbols() == null) req.setCompanySymbols(new ArrayList<>()); if (req.getClientCodes() == null) req.setClientCodes(new ArrayList<>()); if (req.getProfileDocuments() == null) req.setProfileDocuments(new ArrayList<>()); diff --git a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ProfileDocumentService.java b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ProfileDocumentService.java index af4c24586..8a548c097 100644 --- a/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ProfileDocumentService.java +++ b/clearing-parent/company-service/src/main/java/ru/spcex/clearing/company/service/ProfileDocumentService.java @@ -168,7 +168,7 @@ public class ProfileDocumentService extends QueueConsumer implements Initializin } - private RequestInfoUpdate profileDocumentUpdate(BaseRequest profileDocumentUpdateRequestBaseRequest) { + protected RequestInfoUpdate profileDocumentUpdate(BaseRequest profileDocumentUpdateRequestBaseRequest) { log.trace("Start processing ProfileDocumentUpdateRequest!"); RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(profileDocumentUpdateRequestBaseRequest); diff --git a/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/MultiCompanyServiceTest.java b/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/MultiCompanyServiceTest.java new file mode 100644 index 000000000..9535bc82c --- /dev/null +++ b/clearing-parent/company-service/src/test/java/ru/spcex/clearing/company/service/MultiCompanyServiceTest.java @@ -0,0 +1,184 @@ +package ru.spcex.clearing.company.service; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit.jupiter.SpringExtension; +import ru.clearing.classes.statics.data.company.Company; +import ru.clearing.classes.statics.data.company.CompanySymbols; +import ru.clearing.classes.statics.data.profile.Contact; +import ru.clearing.classes.statics.data.profile.ProfileDocument; +import ru.clearing.platform.dictionary.CompanySymbolDictionary; +import ru.clearing.platform.dictionary.WorkflowStatusDictionary; +import ru.spcex.clearing.company.config.BeanConfiguration; +import ru.spcex.clearing.company.config.validation.*; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.cud.company.*; +import ru.spcex.clearing.test.config.ImdgTestConfig; +import ru.spcex.clearing.test.config.KafkaTestConfig; +import ru.spcex.platform.enumeration.CompanySymbol; +import ru.spcex.platform.enumeration.ContactTypes; +import ru.spcex.platform.enumeration.DocumentTypes; +import ru.spcex.platform.enumeration.WorkflowStatus; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import javax.annotation.PostConstruct; + +import java.time.LocalDate; + +import static org.junit.jupiter.api.Assertions.*; +import static ru.spcex.clearing.test.config.ImdgTestConfig.waitAvailableImdgProviderAndAddAdminWithDefaultId; + +@ExtendWith(SpringExtension.class) +@ContextConfiguration(classes = { + MultiCompanyService.class, + + CompanyService.class, + CompanyValidationConfig.class, + ValidationConfig.class, + CompanySymbolService.class, + AccountNotificationHelper.class, + RelationService.class, RelationValidationConfig.class, + + CompanySymbolService.class, + CompanySymbolValidationConfig.class, + CompanyInfoService.class, + ProfileDocumentService.class, + ProfileDocumentValidationConfig.class, + ContactService.class, + ContactValidationConfig.class, + + KafkaTestConfig.class, + ImdgTestConfig.class, + BeanConfiguration.class}) +class MultiCompanyServiceTest { + final Long COMPANY_ID = 123L; + + @Autowired + protected MultiCompanyService multiCompanyService; + + @Autowired + @Qualifier("hazelcastServiceTest") + private ImdgProvider hazelcastServiceTest; + + @PostConstruct + private void init() { + waitAvailableImdgProviderAndAddAdminWithDefaultId(); + Imdg companyImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Company, Company.class); + Imdg companySymbolsImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_CompanySymbols, CompanySymbols.class); + + Company company = new Company(); + company.setId(COMPANY_ID); + company.setWorkflowStatus(WorkflowStatus.Active.getKey()); + company.setFullName("Test company prime"); + companyImdg.insert(company); + + + Imdg companySymbolDictionaryImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_CompanySymbolDictionary, CompanySymbolDictionary.class); + { + CompanySymbolDictionary cSymbol = new CompanySymbolDictionary(); + cSymbol.setId(11L); + cSymbol.setCode(CompanySymbol.CLRC.getKey()); + cSymbol.setName(CompanySymbol.CLRC.getKey()); + cSymbol.setShortname(" CLRC key"); + companySymbolDictionaryImdg.insert(cSymbol); + CompanySymbolDictionary cioSymbol = new CompanySymbolDictionary(); + cioSymbol.setId(22L); + cioSymbol.setCode(CompanySymbol.CIO.getKey()); + cioSymbol.setName(CompanySymbol.CIO.getKey()); + companySymbolDictionaryImdg.insert(cioSymbol); + } + + Imdg contactImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_Contact, Contact.class); + Contact contact = new Contact(); + contact.setId(100L); + contact.setCompanyId(company.getId()); + contact.setContactType(ContactTypes.adrs.getKey()); + contact.setContactValue("UAR, st.Uarus, anystreet st., house 1"); + contactImdg.insert(contact); + + CompanySymbols companySymbol = new CompanySymbols(); + companySymbol.setId(222L); + companySymbol.setCompanyId(company.getId()); + companySymbol.setCompanySymbol(CompanySymbol.CLRC.getKey()); + companySymbol.setCompanySymbolValue("CL-VALUE"); + companySymbolsImdg.insert(companySymbol); + + Imdg profileDocumentImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_ProfileDocument, ProfileDocument.class); + ProfileDocument doc = new ProfileDocument(); + doc.setId(10002L); + doc.setCompanyId(COMPANY_ID); + doc.setDocumentType(DocumentTypes.form.getKey()); + doc.setIssueDate(LocalDate.now()); // тест может показать ложное срабатывание в 00:00:00.001 + profileDocumentImdg.insert(doc); + } + + @Test + void findCompanyDocument() { + ProfileDocumentNewRequest req = new ProfileDocumentNewRequest(); + req.setCompanyId(COMPANY_ID); + req.setDocumentType(DocumentTypes.form.getKey()); + req.setIssueDate(LocalDate.now()); + ProfileDocument doc = multiCompanyService.findCompanyDocument(req); + assertNotNull(doc); + } + + @Test + void findCompanySymbols() { + CompanySymbolNewRequest req = new CompanySymbolNewRequest(); + req.setCompanyId(COMPANY_ID); + req.setCompanySymbol(CompanySymbol.CLRC.getKey()); + req.setCompanySymbolValue("CL-VALUE-2"); + CompanySymbols symbol = multiCompanyService.findCompanySymbols(req); + assertNotNull(symbol); + assertEquals("CL-VALUE", symbol.getCompanySymbolValue()); + } + + @Test + void findContact() { + ContactNewRequest req = new ContactNewRequest(); + req.setCompanyId(COMPANY_ID); + req.setContactType(ContactTypes.adrs.getKey()); + req.setContactValue("UAR, st.Uarus, anystreet st., house 2"); + Contact contact = multiCompanyService.findContact(req); + assertNotNull(contact); + assertNotEquals(req.getContactValue(), contact.getContactValue()); + } + + @Test + void getCompanyIdForCompanySymbols() { + CompanyNewRequest req = new CompanyNewRequest(); + req.setCompanySymbol(CompanySymbol.CLRC.getKey()); + req.setCompanySymbolValue("CL-VALUE"); + Long theId = multiCompanyService.getCompanyIdForCompanySymbols("SAMPLE-NO-uuid", req); + assertEquals(COMPANY_ID, theId); + } + + @Test + void fillCompanyId() { + { + MultiCompanyRequest testReq = new MultiCompanyRequest(); + testReq.setCompany(new CompanyNewRequest()); + testReq.setProfileDocuments(null); + testReq.setCompanySymbols(null); + testReq.setContacts(null); + testReq.setClientCodes(null); + multiCompanyService.fillCompanyId(testReq, COMPANY_ID); + assertEquals(COMPANY_ID, testReq.getCompany().getId()); + assertNotNull(testReq.getProfileDocuments()); + assertNotNull(testReq.getCompanySymbols()); + assertNotNull(testReq.getContacts()); + assertNotNull(testReq.getClientCodes()); + } + { + MultiCompanyRequest testReq = new MultiCompanyRequest(); + testReq.setCompany(new CompanyNewRequest()); + testReq.setCompanyInfo(new CompanyInfoUpdateRequest()); + multiCompanyService.fillCompanyId(testReq, COMPANY_ID); + assertEquals(COMPANY_ID, testReq.getCompany().getId()); + } + } +} \ No newline at end of file diff --git a/clearing-parent/test-clearing/src/main/java/ru/spcex/clearing/test/TestObjectCreator.java b/clearing-parent/test-clearing/src/main/java/ru/spcex/clearing/test/TestObjectCreator.java index 84468145f..1cb6d557f 100644 --- a/clearing-parent/test-clearing/src/main/java/ru/spcex/clearing/test/TestObjectCreator.java +++ b/clearing-parent/test-clearing/src/main/java/ru/spcex/clearing/test/TestObjectCreator.java @@ -38,4 +38,18 @@ public class TestObjectCreator { userRoleSession.setStatus(WorkflowStatus.Active.getKey()); userRoleSessionImdg.insert(userRoleSession); } + +// public void fillWorkflowStatusDictionary() { +// Imdg workflowStatusDictionaryImdg = hazelcastServiceTest.getImdg(IMDGDistributedNames.Map_WorkflowStatusDictionary, WorkflowStatusDictionary.class); +// WorkflowStatusDictionary ws1 = new WorkflowStatusDictionary(); +// ws1.setId(1L); +// ws1.setCode(WorkflowStatus.Active.getKey()); +// ws1.setName("Active"); +// workflowStatusDictionaryImdg.insert(ws1); +// WorkflowStatusDictionary ws2 = new WorkflowStatusDictionary(); +// ws2.setId(2L); +// ws2.setCode(WorkflowStatus.Blocked.getKey()); +// ws2.setName("Blocked"); +// workflowStatusDictionaryImdg.insert(ws2); +// } } From f9fa84ad64f7dd875143b0205d7428ed1422cad3 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 16:50:12 +0300 Subject: [PATCH 3/6] =?UTF-8?q?clearing-service=20sdf08=20=D0=BE=D0=B1?= =?UTF-8?q?=D1=80=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B0=D0=B2=D1=82?= =?UTF-8?q?=D0=BE=D1=81=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20=D1=81?= =?UTF-8?q?=D1=87=D0=B5=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../spcex/clearing/service/StatementService.java | 16 ++++++++++++---- .../logic/stages/DbfImportKafkaMessenger.java | 9 +++++++-- 2 files changed, 19 insertions(+), 6 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index d8ec4bde6..c2a748287 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -123,7 +123,8 @@ public class StatementService extends QueueConsumer implements InitializingBean Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_57); - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } private void processSdf04(StatementRequest statementRequest) { @@ -132,7 +133,8 @@ public class StatementService extends QueueConsumer implements InitializingBean "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_04); if (service != null) { - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } else { log.warn("Executor for SDF_04 not set"); } @@ -144,7 +146,8 @@ public class StatementService extends QueueConsumer implements InitializingBean "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_08); if (service != null) { - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } else { log.warn("Executor for SDF_08 not set"); } @@ -155,7 +158,8 @@ public class StatementService extends QueueConsumer implements InitializingBean Collection sdfGroup = sdfImdg.getCollectionObjectsByFieldValues(Map.of( "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_13); - service.execute(sdfGroup, statementRequest); + Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); } private void processSdf01(StatementRequest statementRequest) { @@ -174,6 +178,10 @@ public class StatementService extends QueueConsumer implements InitializingBean } AbstractExecutor service = executorsMap.get(SdfTable.SDF_01); Result res = service.execute(sdfGroup, statementRequest); + finishSendCommand(res, service, statementRequest); + } + + void finishSendCommand(Result res, AbstractExecutor service, StatementRequest statementRequest) { if (res.getAccountRequests().size() != 0) { kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); } else if (service.isNeedToSendCommand()) { diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java index ae14aecb4..eed6b0b63 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/logic/stages/DbfImportKafkaMessenger.java @@ -1,5 +1,7 @@ package ru.spcex.clearing.dbf.importer.logic.stages; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.stereotype.Component; import ru.spcex.clearing.dbf.importer.logic.data.enums.ETable; @@ -16,6 +18,7 @@ import java.util.function.Supplier; @Component public class DbfImportKafkaMessenger implements InitializingBean { + final Logger log = LoggerFactory.getLogger(getClass()); private final Supplier kafka; private final Map> messengers; @@ -30,7 +33,7 @@ public class DbfImportKafkaMessenger implements InitializingBean { messengers.put(ETable.DF_09, groupId -> messageBalance(groupId, SdfTable.SDF_09)); messengers.put(ETable.DF_16, groupId -> messageBalance(groupId, SdfTable.SDF_16)); messengers.put(ETable.DF_57, groupId -> messageBalance(groupId, SdfTable.SDF_57)); - messengers.put(ETable.DF_04, groupId -> messageBalance(groupId, SdfTable.SDF_04)); + messengers.put(ETable.DF_04, groupId -> messageBalance(groupId, SdfTable.SDF_04)); } /** @@ -48,7 +51,9 @@ public class DbfImportKafkaMessenger implements InitializingBean { StatementRequest statementRequest = new StatementRequest(); statementRequest.setGroupId(groupId); statementRequest.setTable(table); - kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + Long msgId = kafka.get().sendRequestToQueue(Consts.STATEMENT_PROCESS, statementRequest); + log.debug("Send StatementRequest({}, {}) message id={} to kafka \"{}\"", + groupId, table, msgId, Consts.STATEMENT_PROCESS); } private void messageDf04(Long groupId) { From ad09a31171f6d892361f67537f620456d32efb74 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 17:15:50 +0300 Subject: [PATCH 4/6] =?UTF-8?q?clearing-service=20sdf08=20=D0=BE=D0=B1?= =?UTF-8?q?=D1=80=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B0=D0=B2=D1=82?= =?UTF-8?q?=D0=BE=D1=81=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20=D1=81?= =?UTF-8?q?=D1=87=D0=B5=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../java/ru/spcex/clearing/service/StatementService.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index c2a748287..6a9f10790 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -124,7 +124,7 @@ public class StatementService extends QueueConsumer implements InitializingBean "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_57); Result res = service.execute(sdfGroup, statementRequest); - finishSendCommand(res, service, statementRequest); +// finishSendCommand(res, service, statementRequest); } private void processSdf04(StatementRequest statementRequest) { @@ -134,7 +134,7 @@ public class StatementService extends QueueConsumer implements InitializingBean AbstractExecutor service = executorsMap.get(SdfTable.SDF_04); if (service != null) { Result res = service.execute(sdfGroup, statementRequest); - finishSendCommand(res, service, statementRequest); +// finishSendCommand(res, service, statementRequest); } else { log.warn("Executor for SDF_04 not set"); } @@ -159,7 +159,7 @@ public class StatementService extends QueueConsumer implements InitializingBean "generationId", statementRequest.getGroupId())); AbstractExecutor service = executorsMap.get(SdfTable.SDF_13); Result res = service.execute(sdfGroup, statementRequest); - finishSendCommand(res, service, statementRequest); +// finishSendCommand(res, service, statementRequest); } private void processSdf01(StatementRequest statementRequest) { From 872b66984dbee31b79c200c59237817ce8f448bf Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 17:35:10 +0300 Subject: [PATCH 5/6] =?UTF-8?q?clearing-service=20sdf08=20=D0=BE=D0=B1?= =?UTF-8?q?=D1=80=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B0=D0=B2=D1=82?= =?UTF-8?q?=D0=BE=D1=81=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20=D1=81?= =?UTF-8?q?=D1=87=D0=B5=D1=82=D0=BE=D0=B2.=20=D0=9E=D0=BA=D0=B0=D0=B7?= =?UTF-8?q?=D0=B0=D0=BB=D0=BE=D1=81=D1=8C=20=D1=87=D1=82=D0=BE=20DEPO=20?= =?UTF-8?q?=D1=83=D0=B6=D0=BD=D0=BE=20=D0=B4=D0=B5=D0=BB=D0=B0=D1=82=D1=8C?= =?UTF-8?q?=20(=D1=87=D0=B0=D1=81=D1=82=201/2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../ru/spcex/clearing/service/StatementService.java | 13 ++++++------- .../clearing/platform/messaging/domain/Consts.java | 1 + 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java index 6a9f10790..c8bd2fe62 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/StatementService.java @@ -147,7 +147,11 @@ public class StatementService extends QueueConsumer implements InitializingBean AbstractExecutor service = executorsMap.get(SdfTable.SDF_08); if (service != null) { Result res = service.execute(sdfGroup, statementRequest); - finishSendCommand(res, service, statementRequest); + if (res.getAccountRequests().size() != 0) { + kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF08, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); + } else if (service.isNeedToSendCommand()) { + service.sendCommand(kafkaSender, res); + } } else { log.warn("Executor for SDF_08 not set"); } @@ -177,12 +181,7 @@ public class StatementService extends QueueConsumer implements InitializingBean .collect(Collectors.toList()); } AbstractExecutor service = executorsMap.get(SdfTable.SDF_01); - Result res = service.execute(sdfGroup, statementRequest); - finishSendCommand(res, service, statementRequest); - } - - void finishSendCommand(Result res, AbstractExecutor service, StatementRequest statementRequest) { - if (res.getAccountRequests().size() != 0) { + Result res = service.execute(sdfGroup, statementRequest); if (res.getAccountRequests().size() != 0) { kafkaSender.sendRequestToQueue(Consts.ACCOUNT_NEW_SDF01, createAccountsRequest(statementRequest.getGroupId(), res.getAccountRequests())); } else if (service.isNeedToSendCommand()) { service.sendCommand(kafkaSender, res); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 49e1f3772..4103f6812 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -86,6 +86,7 @@ public interface Consts { String CREATE_REPORT_FOR_REGISTRY = "create-report-for-registry"; String ACCOUNT_NEW_SDF01 = "account-new-sdf01"; + String ACCOUNT_NEW_SDF08 = "account-new-sdf08"; String DESTINATION_RELATION_NEW = "relation-new"; String DESTINATION_RELATION_UPDATE = "relation-update"; From dd151e80d2682d8b2af1a0263865bf8234a66059 Mon Sep 17 00:00:00 2001 From: AKurakin Date: Fri, 9 Jun 2023 18:23:47 +0300 Subject: [PATCH 6/6] =?UTF-8?q?clearing-service=20sdf08=20=D0=BE=D0=B1?= =?UTF-8?q?=D1=80=D0=B0=D0=B1=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B0=D0=B2=D1=82?= =?UTF-8?q?=D0=BE=D1=81=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F=20=D1=81?= =?UTF-8?q?=D1=87=D0=B5=D1=82=D0=BE=D0=B2.=20=D0=9E=D0=BA=D0=B0=D0=B7?= =?UTF-8?q?=D0=B0=D0=BB=D0=BE=D1=81=D1=8C=20=D1=87=D1=82=D0=BE=20DEPO=20?= =?UTF-8?q?=D1=83=D0=B6=D0=BD=D0=BE=20=D0=B4=D0=B5=D0=BB=D0=B0=D1=82=D1=8C?= =?UTF-8?q?=20(=D1=87=D0=B0=D1=81=D1=82=202/2).=20=D0=9F=D0=BE=D0=BF=D1=80?= =?UTF-8?q?=D0=B0=D0=B2=D0=B8=D0=BB=20=D1=82=D1=80=D0=B0=D0=BD=D0=B7=D0=B0?= =?UTF-8?q?=D0=BA=D1=86=D0=B8=D0=B8.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../account/service/BankAccountService.java | 37 ++---- .../service/ClearingAccountService.java | 58 +++++---- .../account/service/DepoAccountService.java | 123 +++++++++++++++++- .../service/InformationAccountService.java | 19 +-- .../account/sdf01/AccountSdf01Request.java | 3 + 5 files changed, 180 insertions(+), 60 deletions(-) diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java index 44d37f90d..95e422f71 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/BankAccountService.java @@ -55,11 +55,11 @@ public class BankAccountService extends QueueConsumer implements InitializingBea ValidationHelper validationHelper, AccountService accountService, @Qualifier("bankAccountNewRequestValidator") - Function bankAccountNewRequestValidator, + Function bankAccountNewRequestValidator, @Qualifier("bankAccountUpdateRequestValidator") - Function bankAccountUpdateRequestValidator, + Function bankAccountUpdateRequestValidator, @Qualifier("bankAccountBlockRequestValidator") - Function bankAccountBlockRequestValidator) { + Function bankAccountBlockRequestValidator) { super(kafkaQueue, kafkaProducer); this.bankAccountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_BankAccount, BankAccount.class); this.accountMap = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); @@ -115,6 +115,8 @@ public class BankAccountService extends QueueConsumer implements InitializingBea boolean txOk = false; imdgTransaction.beginTransaction(); try { + Imdg bankAccountMap = imdgTransaction.getImdg(IMDGDistributedNames.Map_BankAccount, BankAccount.class); + Imdg accountMap = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); accountId = accountMap.insert(account); BankAccount bankAccount = new BankAccount(); @@ -175,6 +177,8 @@ public class BankAccountService extends QueueConsumer implements InitializingBea boolean txOk = false; imdgTransaction.beginTransaction(); try { + Imdg bankAccountMap = imdgTransaction.getImdg(IMDGDistributedNames.Map_BankAccount, BankAccount.class); + Imdg accountMap = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); accountMap.update(account); bankAccountMap.update(bankAccount); txOk = true; @@ -189,7 +193,7 @@ public class BankAccountService extends QueueConsumer implements InitializingBea bankAccount.getId(), account.getId()); imdgTransaction.rollbackTransaction(); - } + } } log.debug("successfully update, existing bankAccount with id {}", bankAccount.getId()); return null; @@ -210,27 +214,10 @@ public class BankAccountService extends QueueConsumer implements InitializingBea Account account = accountMap.getSingleObjectByID(bankAccount.getAccountId()); account.setStatus(ServiceStatus.Blocked.getKey()); account.setUpdated(Instant.now()); - - ImdgTransaction imdgTransaction = imdgProvider.newTransaction(); - boolean txOk = false; - imdgTransaction.beginTransaction(); - try { - accountMap.update(account); - //без изменений bankAccountMap.update(bankAccount); - txOk = true; - } finally { - if (txOk) { - imdgTransaction.commitTransaction(); - log.debug("successfully processed, new bank account id {}, account id {}", - bankAccount.getId(), - account.getId()); - } else { - log.debug("failed block, bank account id {}, new account id {}", - bankAccount.getId(), - account.getId()); - imdgTransaction.rollbackTransaction(); - } - } + + accountMap.update(account); + //без изменений bankAccountMap.update(bankAccount); + log.debug("successfully block, existing bankAccount with id {}", bankAccount.getId()); return null; } diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java index de3bd3411..113c2cad0 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/ClearingAccountService.java @@ -52,8 +52,6 @@ public class ClearingAccountService extends QueueConsumer implements Initializin private final AccountService accountService; private final ValidationHelper validationHelper; private final ImdgProvider imdgProvider; - private final Imdg accountImdg; - private final Imdg clearingAccountImdg; private final IMessageResolver messageResolver; private final RequestHelper requestHelper; private final Function clearingAccountNewRequestValidator; @@ -77,8 +75,6 @@ public class ClearingAccountService extends QueueConsumer implements Initializin this.accountService = accountService; this.validationHelper = validationHelper; this.imdgProvider = imdgProvider; - this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); - this.clearingAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class); this.messageResolver = messageResolver; this.requestHelper = requestHelper; this.clearingAccountNewRequestValidator = clearingAccountNewRequestValidator; @@ -126,6 +122,8 @@ public class ClearingAccountService extends QueueConsumer implements Initializin imdgTransaction.beginTransaction(); ClearingAccount clearingAccount = null; try { + Imdg accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Imdg clearingAccountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class); accountId = accountImdg.insert(account); clearingAccount = new ClearingAccount(); @@ -161,28 +159,41 @@ public class ClearingAccountService extends QueueConsumer implements Initializin if (requestInfoUpdate != null) return requestInfoUpdate; ClearingAccountUpdateRequest req = userRequest.getRequestPayload(); - ImdgPredicateBuilder pb = accountImdg.predicateBuilder(); - ImdgPredicate accountValuePredicate = pb.equals("account", req.getAccount()); - ImdgPredicate accountStatusPredicate = pb.equals("status", ServiceStatus.Active.getKey()); - ImdgPredicate accountTypePredicate = pb.equals("accountType", AccountType.Clrn.getKey()); - ImdgPredicate finalPredicate = pb.and(accountValuePredicate, - accountStatusPredicate, - accountTypePredicate); - Account account = accountImdg.getCollectionObjectsByPredicate(finalPredicate).iterator().next(); - Long accountId = account.getId(); + ImdgTransaction imdgTransaction = imdgProvider.newTransaction(); + boolean txOk = false; + imdgTransaction.beginTransaction(); + try { + Imdg accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Imdg clearingAccountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class); - ClearingAccount clearingAccount = clearingAccountImdg.getSingleObjectByFieldValues(Map.of("accountId", accountId)); - if (clearingAccount == null) - return requestHelper.makeErrorResponse(userRequest, AccountError.AccountNotFound, account.getAccount()); + ImdgPredicateBuilder pb = accountImdg.predicateBuilder(); + ImdgPredicate accountValuePredicate = pb.equals("account", req.getAccount()); + ImdgPredicate accountStatusPredicate = pb.equals("status", ServiceStatus.Active.getKey()); + ImdgPredicate accountTypePredicate = pb.equals("accountType", AccountType.Clrn.getKey()); + ImdgPredicate finalPredicate = pb.and(accountValuePredicate, + accountStatusPredicate, + accountTypePredicate); - Integer statusValue = req.getStatus(); - if (statusValue == 0) account.setStatus(ServiceStatus.Blocked.getKey()); - else if (statusValue == 1) account.setStatus(ServiceStatus.Active.getKey()); - else if (statusValue == 2) account.setStatus(ServiceStatus.Closed.getKey()); - account.setUpdated(Instant.now()); + Account account = accountImdg.getCollectionObjectsByPredicate(finalPredicate).iterator().next(); + Long accountId = account.getId(); - accountImdg.update(account); + ClearingAccount clearingAccount = clearingAccountImdg.getSingleObjectByFieldValues(Map.of("accountId", accountId)); + if (clearingAccount == null) + return requestHelper.makeErrorResponse(userRequest, AccountError.AccountNotFound, account.getAccount()); + + Integer statusValue = req.getStatus(); + if (statusValue == 0) account.setStatus(ServiceStatus.Blocked.getKey()); + else if (statusValue == 1) account.setStatus(ServiceStatus.Active.getKey()); + else if (statusValue == 2) account.setStatus(ServiceStatus.Closed.getKey()); + account.setUpdated(Instant.now()); + + accountImdg.update(account); + txOk = true; + } finally { + if (txOk) imdgTransaction.commitTransaction(); + else imdgTransaction.rollbackTransaction(); + } return null; } @@ -199,6 +210,9 @@ public class ClearingAccountService extends QueueConsumer implements Initializin boolean txOk = false; imdgTransaction.beginTransaction(); try { + Imdg accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Imdg clearingAccountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_ClearingAccount, ClearingAccount.class); + accountsLoop: for (AccountSdfRequestPart accountReq : req.getAccounts()) { Instant now = Instant.now(); diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java index e9c050283..596136ebd 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/DepoAccountService.java @@ -9,10 +9,15 @@ import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.account.DepoAccount; +import ru.spcex.clearing.account.errors.AccountError; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.DepoAccountNewRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdf01Request; +import ru.spcex.clearing.platform.messaging.domain.cud.account.sdf01.AccountSdfRequestPart; +import ru.spcex.clearing.platform.messaging.domain.cud.balance.AccountSdfToStatementRequestPart; +import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest; import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; @@ -20,6 +25,7 @@ import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.validation.common.ValidationHelper; import ru.spcex.platform.enumeration.AccountType; +import ru.spcex.platform.enumeration.SdfTable; import ru.spcex.platform.enumeration.ServiceStatus; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @@ -27,6 +33,8 @@ import ru.spcex.platform.imdg.api.ImdgTransaction; import ru.spcex.platform.utils.validation.IValidator; import java.time.Instant; +import java.util.ArrayList; +import java.util.List; import java.util.function.Function; @Service @@ -38,8 +46,6 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea private final ImdgProvider imdgProvider; private final AccountService accountService; private final Function depoAccountNewRequestValidator; - private final Imdg accountImdg; - private final Imdg depoAccountImdg; public DepoAccountService(Consumer kafkaQueue, Producer kafkaProducer, @@ -55,8 +61,6 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea this.imdgProvider = imdgProvider; this.accountService = accountService; this.depoAccountNewRequestValidator = depoAccountNewRequestValidator; - this.accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); - this.depoAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class); } @Override @@ -64,6 +68,10 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea callback(DepoAccountNewRequest.class) .setFunction(this::depoAccountNew) .forDestination(Consts.DESTINATION_DEPO_ACCOUNT_NEW, callbacks::put); + callback(AccountSdf01Request.class) + .setFunction(this::accountNewSdf08) + .forDestination(Consts.ACCOUNT_NEW_SDF08, callbacks::put); + init(); } @@ -94,6 +102,8 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea imdgTransaction.beginTransaction(); DepoAccount depoAccount = null; try { + Imdg accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Imdg depoAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class); accountId = accountImdg.insert(account); depoAccount = new DepoAccount(); @@ -120,4 +130,109 @@ public class DepoAccountService extends QueueConsumer implements InitializingBea } return null; } + + public RequestInfoUpdate accountNewSdf08(BaseRequest userRequest) { + log.debug("AccountSdf08Request received, id={}", userRequest.getId()); + + AccountSdf01Request req = userRequest.getRequestPayload(); + List accountToStatement = new ArrayList<>(); + + ImdgTransaction imdgTransaction = imdgProvider.newTransaction(); + List toTCRRequests = new ArrayList<>(); + boolean txOk = false; + imdgTransaction.beginTransaction(); + try { + Imdg accountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Account, Account.class); + Imdg depoAccountImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_DepoAccount, DepoAccount.class); + accountsLoop: + for (AccountSdfRequestPart accountReq : req.getAccounts()) { + +// DepoAccountNewRequest req = userRequest.getRequestPayload(); + + Instant now = Instant.now(); + Account account = new Account(); + account.setAccount(accountReq.getAccount()); + account.setAccountType(AccountType.Depo.getKey()); + account.setStatus(ServiceStatus.Active.getKey()); + account.setCompanyId(accountReq.getCompanyId()); + account.setCreated(now); + account.setUpdated(now); + RequestInfoUpdate requestInfoUpdate = accountService.fillAccountFromRelation(account, userRequest.getId(), true); + if (requestInfoUpdate != null) { + log.warn("Error fill new account from relation. {}", /*account.getId(),*/ requestInfoUpdate.getMessage()); + { + AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart(); + responsePart.setSdfId(accountReq.getSdfId()); + responsePart.setErrorCode(AccountError.ClearingCategoryNotFound.getId()); // see accountService.fillAccountFromRelation + responsePart.setErrorText(requestInfoUpdate.getMessage()); + accountToStatement.add(responsePart); + continue accountsLoop; + } + } + + Long depoAccountId = -1L; + Long accountId = -1L; + + accountId = accountImdg.insert(account); + + DepoAccount depoAccount = new DepoAccount(); + depoAccount.setCompanyId(accountReq.getCompanyId()); + depoAccount.setAccountId(accountId); + depoAccount.setDepoAccountType(accountReq.getAccountType()); + depoAccountId = depoAccountImdg.insert(depoAccount); + + + + log.debug("New account {}, depoAccount {} was created.", accountId, depoAccountId); + + { + TradingClearingRegistryNewRequest request = new TradingClearingRegistryNewRequest(); + request.setDepoAccountId(accountId); + request.setCompanyId(depoAccount.getCompanyId()); +// request.setTradingClearingRegistryType(cl); + toTCRRequests.add(request); + } + { + AccountSdfToStatementRequestPart responsePart = new AccountSdfToStatementRequestPart(); + responsePart.setSdfId(accountReq.getSdfId()); + responsePart.setErrorCode(null); + responsePart.setErrorText(null); + accountToStatement.add(responsePart); + } + } + txOk = true; + } finally { + if (txOk) { + imdgTransaction.commitTransaction(); + } else { + log.debug("failed insert, new clearing accounts. Request id={}", userRequest.getId()); + imdgTransaction.rollbackTransaction(); + } + } + + if (txOk) { + log.debug("Sending {} messages of TradingClearingRegistryNewRequest", toTCRRequests.size()); + for (TradingClearingRegistryNewRequest tcrReq : toTCRRequests) { + log.debug("Send message to kafka \"{}\": {}", Consts.DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW, + LogFormatter.toStringWrapper(tcrReq)); + Long kafkaId = kafkaSender.sendRequestToQueue(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_AUTO_NEW, tcrReq); + log.trace("successfully send request {} to kafka: new clearing account MoneyAccountId {}, DepoAccountId {}", + kafkaId, tcrReq.getMoneyAccountId(), tcrReq.getDepoAccountId()); + } + } + + sendStatementRequestBack(req.getGroupingSdf01Id(), accountToStatement); + log.debug("successfully processed, grouping id={}, processed number={}", req.getGroupingSdf01Id(), accountToStatement.size()); + + return null; + } + + public void sendStatementRequestBack(Long groupingSdf01Id, List results) { + StatementRequest request = new StatementRequest(); + request.setGroupId(groupingSdf01Id); + request.setAccountCreationResults(results); + request.setTable(SdfTable.SDF_08); // по нему запрос получили + log.debug("Send message to kafka \"{}\": {}", Consts.STATEMENT_PROCESS, LogFormatter.toStringWrapper(request)); + kafkaSender.sendRequestToQueue(Consts.STATEMENT_PROCESS, request); + } } diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/InformationAccountService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/InformationAccountService.java index 7aead7923..88eac050b 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/InformationAccountService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/InformationAccountService.java @@ -17,7 +17,7 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.account.InformationAccountNewRequest; import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryNewRequest; -import ru.spcex.clearing.platform.messaging.domain.cud.reports.ReportRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.reports.NotificationRequest; import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; @@ -25,7 +25,6 @@ import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.clearing.util.services.RequestHelper; import ru.spcex.clearing.validation.common.ValidationHelper; import ru.spcex.platform.enumeration.AccountType; -import ru.spcex.platform.enumeration.ReportType; import ru.spcex.platform.enumeration.ServiceStatus; import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgProvider; @@ -142,6 +141,9 @@ public class InformationAccountService extends QueueConsumer implements Initiali Long accountId = -1L; InformationAccount informationAccount = null; try { + Imdg informationAccountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_InformationAccount, InformationAccount.class); + Imdg accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); + accountId = accountImdg.insert(account); informationAccount = new InformationAccount(); @@ -239,6 +241,8 @@ public class InformationAccountService extends QueueConsumer implements Initiali Long accountId = -1L; InformationAccount informationAccount = null; try { + Imdg informationAccountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_InformationAccount, InformationAccount.class); + Imdg accountImdg = imdgTransaction.getImdg(IMDGDistributedNames.Map_Account, Account.class); accountId = accountImdg.insert(account); informationAccount = new InformationAccount(); @@ -284,13 +288,10 @@ public class InformationAccountService extends QueueConsumer implements Initiali * Формирование уведолмения о регистрации УК */ protected void sendNotificationToReport(InformationAccount informationAccount, Account account) { - //todo актуализировать ТЗ или ReportRequest -// ReportRequest request = new ReportRequest(); -// request.setReportId(ReportType.REGISTRACTION_UK/NEW_INFO_ACCOUNT); -// request.setCompanyId(informationAccount.getCompanyId()); -// log.debug("Send message to kafka \"{}\": {}", Consts.CREATE_REPORT_WITH_COMPANY_ID, LogFormatter.toStringWrapper(request)); -// kafkaSender.sendRequestToQueue(Consts.CREATE_REPORT_WITH_COMPANY_ID, request); -// // см. в report-service: ROOT_ACTV_NotificationBuilder, ru.spcex.clearing.reports.services.ReportService + NotificationRequest request = new NotificationRequest(); + request.setConsumerId(account.getCompanyId()); + log.debug("Send message to kafka \"{}\": {}", Consts.CREATE_NOTIFICATION_NTCR, LogFormatter.toStringWrapper(request)); + kafkaSender.sendRequestToQueue(Consts.CREATE_NOTIFICATION_NCMP, request); } public String generateInfoAccount(Long id) { diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/account/sdf01/AccountSdf01Request.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/account/sdf01/AccountSdf01Request.java index 2e9cb13a2..52c1a5c2e 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/account/sdf01/AccountSdf01Request.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/account/sdf01/AccountSdf01Request.java @@ -4,6 +4,9 @@ import com.fasterxml.jackson.annotation.JsonProperty; import java.util.List; +/** + * SDF01 / SDF08 + */ public class AccountSdf01Request { private Long groupingSdf01Id;