Compare commits
2 commits
1df96b7ef6
...
d3a6456d33
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d3a6456d33 | ||
|
|
95f7741111 |
17 changed files with 626 additions and 19 deletions
|
|
@ -43,6 +43,7 @@
|
||||||
<module>swt-importer</module>
|
<module>swt-importer</module>
|
||||||
<module>gateway-api</module>
|
<module>gateway-api</module>
|
||||||
<module>imdg-hist</module>
|
<module>imdg-hist</module>
|
||||||
|
<module>snapshot-maker</module>
|
||||||
</modules>
|
</modules>
|
||||||
|
|
||||||
<properties>
|
<properties>
|
||||||
|
|
|
||||||
82
clearing-parent/snapshot-maker/pom.xml
Normal file
82
clearing-parent/snapshot-maker/pom.xml
Normal file
|
|
@ -0,0 +1,82 @@
|
||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<project xmlns="http://maven.apache.org/POM/4.0.0"
|
||||||
|
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||||
|
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||||
|
<modelVersion>4.0.0</modelVersion>
|
||||||
|
|
||||||
|
<artifactId>snapshot-maker</artifactId>
|
||||||
|
<name>snapshot-maker</name>
|
||||||
|
<description>Database snapshot maker</description>
|
||||||
|
<version>SPCEX-3.12.5</version>
|
||||||
|
|
||||||
|
<parent>
|
||||||
|
<groupId>ru.spcex.clearing</groupId>
|
||||||
|
<artifactId>clearing-parent</artifactId>
|
||||||
|
<version>SPCEX-3.12.5</version>
|
||||||
|
</parent>
|
||||||
|
|
||||||
|
<properties>
|
||||||
|
<maven.compiler.source>17</maven.compiler.source>
|
||||||
|
<maven.compiler.target>17</maven.compiler.target>
|
||||||
|
</properties>
|
||||||
|
|
||||||
|
<dependencies>
|
||||||
|
<!-- Spring dependencies -->
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter</artifactId>
|
||||||
|
</dependency>
|
||||||
|
|
||||||
|
<!-- Kafka dependencies -->
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.platform</groupId>
|
||||||
|
<artifactId>platform-messaging</artifactId>
|
||||||
|
</dependency>
|
||||||
|
|
||||||
|
<!-- Misc dependencies -->
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.clearing</groupId>
|
||||||
|
<artifactId>classes</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>ru.spcex.platform</groupId>
|
||||||
|
<artifactId>platform-enum</artifactId>
|
||||||
|
</dependency>
|
||||||
|
|
||||||
|
|
||||||
|
<!-- Test dependencies -->
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-test</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
</dependencies>
|
||||||
|
|
||||||
|
<build>
|
||||||
|
<resources>
|
||||||
|
<resource>
|
||||||
|
<directory>src/main/resources</directory>
|
||||||
|
<excludes>
|
||||||
|
<exclude>application.properties</exclude>
|
||||||
|
</excludes>
|
||||||
|
<filtering>false</filtering>
|
||||||
|
</resource>
|
||||||
|
</resources>
|
||||||
|
<plugins>
|
||||||
|
<plugin>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-maven-plugin</artifactId>
|
||||||
|
<executions>
|
||||||
|
<execution>
|
||||||
|
<goals>
|
||||||
|
<goal>repackage</goal>
|
||||||
|
</goals>
|
||||||
|
</execution>
|
||||||
|
</executions>
|
||||||
|
<configuration>
|
||||||
|
<finalName>${project.artifactId}</finalName>
|
||||||
|
</configuration>
|
||||||
|
</plugin>
|
||||||
|
</plugins>
|
||||||
|
</build>
|
||||||
|
</project>
|
||||||
|
|
@ -0,0 +1,11 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker;
|
||||||
|
|
||||||
|
import org.springframework.boot.SpringApplication;
|
||||||
|
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||||
|
|
||||||
|
@SpringBootApplication
|
||||||
|
public class SnapshotMakerApplication {
|
||||||
|
public static void main(String[] args) {
|
||||||
|
SpringApplication.run(SnapshotMakerApplication.class, args);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,32 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker.config;
|
||||||
|
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.apache.kafka.clients.producer.Producer;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.context.annotation.Configuration;
|
||||||
|
import org.springframework.context.annotation.Scope;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory;
|
||||||
|
import ru.spcex.clearing.snapshot.maker.config.settings.SnapshotMakerSettings;
|
||||||
|
|
||||||
|
@Configuration
|
||||||
|
public class KafkaConfig {
|
||||||
|
@Autowired
|
||||||
|
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
|
||||||
|
@Bean
|
||||||
|
public Consumer<String, Object> createConsumer(SnapshotMakerSettings settings) {
|
||||||
|
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Bean
|
||||||
|
public Producer<String, Object> createProducer(SnapshotMakerSettings settings) {
|
||||||
|
if (settings.getKafkaProducer() != null) {
|
||||||
|
return KafkaProducerFactory.producer(settings.getKafkaProducer());
|
||||||
|
} else {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,31 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker.config.settings;
|
||||||
|
|
||||||
|
public class DatabaseSettings {
|
||||||
|
private String jdbcUrl;
|
||||||
|
private String username;
|
||||||
|
private String password;
|
||||||
|
|
||||||
|
public String getJdbcUrl() {
|
||||||
|
return jdbcUrl;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setJdbcUrl(String jdbcUrl) {
|
||||||
|
this.jdbcUrl = jdbcUrl;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getUsername() {
|
||||||
|
return username;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setUsername(String username) {
|
||||||
|
this.username = username;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getPassword() {
|
||||||
|
return password;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setPassword(String password) {
|
||||||
|
this.password = password;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,49 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker.config.settings;
|
||||||
|
|
||||||
|
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||||
|
import org.springframework.context.annotation.PropertySource;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings;
|
||||||
|
import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings;
|
||||||
|
|
||||||
|
@Component
|
||||||
|
@PropertySource("file:${spring.config.location}/application.properties")
|
||||||
|
@ConfigurationProperties("snapshot-maker")
|
||||||
|
public class SnapshotMakerSettings {
|
||||||
|
private DatabaseSettings database;
|
||||||
|
private KafkaConsumerSettings kafkaConsumer;
|
||||||
|
private KafkaProducerSettings kafkaProducer;
|
||||||
|
private String storeLocation;
|
||||||
|
|
||||||
|
public DatabaseSettings getDatabase() {
|
||||||
|
return database;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setDatabase(DatabaseSettings database) {
|
||||||
|
this.database = database;
|
||||||
|
}
|
||||||
|
|
||||||
|
public KafkaConsumerSettings getKafkaConsumer() {
|
||||||
|
return kafkaConsumer;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setKafkaConsumer(KafkaConsumerSettings kafkaConsumer) {
|
||||||
|
this.kafkaConsumer = kafkaConsumer;
|
||||||
|
}
|
||||||
|
|
||||||
|
public KafkaProducerSettings getKafkaProducer() {
|
||||||
|
return kafkaProducer;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setKafkaProducer(KafkaProducerSettings kafkaProducer) {
|
||||||
|
this.kafkaProducer = kafkaProducer;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getStoreLocation() {
|
||||||
|
return storeLocation;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setStoreLocation(String storeLocation) {
|
||||||
|
this.storeLocation = storeLocation;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,66 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker.service;
|
||||||
|
|
||||||
|
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.stereotype.Service;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.ExportedSDFFile;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||||
|
import ru.spcex.clearing.snapshot.maker.utils.FileStore;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class CommandService extends QueueConsumer implements InitializingBean {
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final SessionService sessionService;
|
||||||
|
private final FileStore fileStore;
|
||||||
|
|
||||||
|
public CommandService(Consumer<String, Object> kafkaQueue,
|
||||||
|
Producer<String, Object> kafkaResponseQueue,
|
||||||
|
SessionService sessionService,
|
||||||
|
FileStore fileStore) {
|
||||||
|
super(kafkaQueue, kafkaResponseQueue);
|
||||||
|
this.sessionService = sessionService;
|
||||||
|
this.fileStore = fileStore;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void afterPropertiesSet() throws Exception {
|
||||||
|
callback(LauncherCommandRequest.class)
|
||||||
|
.setFunction(this::process)
|
||||||
|
.forDestination(Consts.LAUNCHER_NEW, callbacks::put);
|
||||||
|
callback(ExportedSDFFile.class)
|
||||||
|
.setConsumer(r -> {
|
||||||
|
if (fileStore.isSessionStarted()) {
|
||||||
|
log.info("Adding {} to files list", r.getRequestPayload());
|
||||||
|
} else {
|
||||||
|
log.info("Session hasn't been started, ignoring {}", r.getRequestPayload());
|
||||||
|
}
|
||||||
|
fileStore.addFile(r.getRequestPayload());
|
||||||
|
})
|
||||||
|
.forDestination(Consts.FILE_CREATED, callbacks::put);
|
||||||
|
init();
|
||||||
|
}
|
||||||
|
|
||||||
|
private RequestInfoUpdate process(BaseRequest<LauncherCommandRequest> request) {
|
||||||
|
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
|
||||||
|
|
||||||
|
if (request.getRequestPayload().getTaskName().equals("STSS")) {
|
||||||
|
requestInfoUpdate = sessionService.startSession();
|
||||||
|
} else if (request.getRequestPayload().getTaskName().equals("SPSS")) {
|
||||||
|
requestInfoUpdate = sessionService.stopSession();
|
||||||
|
} else {
|
||||||
|
log.info("Unknown task: {}", request.getRequestPayload().getTaskName());
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Unknown task: " + request.getRequestPayload().getTaskName());
|
||||||
|
}
|
||||||
|
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,186 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker.service;
|
||||||
|
|
||||||
|
import java.io.BufferedReader;
|
||||||
|
import java.io.File;
|
||||||
|
import java.io.IOException;
|
||||||
|
import java.io.InputStreamReader;
|
||||||
|
import java.nio.file.Files;
|
||||||
|
import java.nio.file.Path;
|
||||||
|
import java.nio.file.Paths;
|
||||||
|
import java.time.LocalDateTime;
|
||||||
|
import java.time.format.DateTimeFormatter;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.regex.Matcher;
|
||||||
|
import java.util.regex.Pattern;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.ExportedSDFFile;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.FileType;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||||
|
import ru.spcex.clearing.snapshot.maker.config.settings.SnapshotMakerSettings;
|
||||||
|
import ru.spcex.clearing.snapshot.maker.utils.FileStore;
|
||||||
|
|
||||||
|
@Service
|
||||||
|
public class SessionService {
|
||||||
|
private final static Pattern JDBC_URL_PATTERN = Pattern.compile("(?<=//)(.+?)(?=:):(.+?)(?=/)/(.+?)(?=\\?|$)");
|
||||||
|
private final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
private final SnapshotMakerSettings settings;
|
||||||
|
private final DateTimeFormatter dtf = DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH-mm-ss");
|
||||||
|
private final FileStore fileStore;
|
||||||
|
private Path currentPath;
|
||||||
|
|
||||||
|
public SessionService(SnapshotMakerSettings settings,
|
||||||
|
FileStore fileStore) {
|
||||||
|
this.settings = settings;
|
||||||
|
this.fileStore = fileStore;
|
||||||
|
}
|
||||||
|
|
||||||
|
public RequestInfoUpdate startSession() {
|
||||||
|
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
|
||||||
|
|
||||||
|
if (currentPath != null || fileStore.isSessionStarted()) {
|
||||||
|
log.info("You need to stop session first!");
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("You need to stop session first!");
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
|
||||||
|
log.info("Starting making snapshot");
|
||||||
|
|
||||||
|
try {
|
||||||
|
this.currentPath = Files.createDirectories(Paths.get(settings.getStoreLocation()).resolve(LocalDateTime.now().format(dtf)));
|
||||||
|
log.info("Session folder created: {}", this.currentPath);
|
||||||
|
} catch (IOException e) {
|
||||||
|
log.info("Can't create directory, error={}", e.getMessage());
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Can't create directory, error = " + e.getMessage());
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
dumpDatabase(requestInfoUpdate);
|
||||||
|
|
||||||
|
if (requestInfoUpdate.getStatus() == Status.Error) {
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.fileStore.setSessionStarted(true);
|
||||||
|
|
||||||
|
requestInfoUpdate.setStatus(Status.Success);
|
||||||
|
requestInfoUpdate.setMessage("Snapshot created successfully in folder " + this.currentPath.toAbsolutePath());
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
|
||||||
|
public RequestInfoUpdate stopSession() {
|
||||||
|
RequestInfoUpdate requestInfoUpdate = new RequestInfoUpdate();
|
||||||
|
|
||||||
|
if (currentPath == null || !fileStore.isSessionStarted()) {
|
||||||
|
log.info("You need to start session first!");
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("You need to start session first!");
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
|
||||||
|
log.info("Stopping session");
|
||||||
|
|
||||||
|
this.fileStore.setSessionStarted(false);
|
||||||
|
copyDirectories(requestInfoUpdate);
|
||||||
|
this.currentPath = null;
|
||||||
|
this.fileStore.clearFiles();
|
||||||
|
|
||||||
|
if (requestInfoUpdate.getStatus() == Status.Error) {
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
|
||||||
|
requestInfoUpdate.setStatus(Status.Success);
|
||||||
|
requestInfoUpdate.setMessage("Session snapshot stopped successfully!");
|
||||||
|
return requestInfoUpdate;
|
||||||
|
}
|
||||||
|
|
||||||
|
private void dumpDatabase(RequestInfoUpdate requestInfoUpdate) {
|
||||||
|
Matcher matcher = JDBC_URL_PATTERN.matcher(settings.getDatabase().getJdbcUrl());
|
||||||
|
if (!matcher.find()) {
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Database URL is not valid!");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
Path dumpLocation = this.currentPath.resolve("database_dump.sql");
|
||||||
|
|
||||||
|
ProcessBuilder pb = new ProcessBuilder();
|
||||||
|
pb.environment().put("PGPASSWORD", settings.getDatabase().getPassword());
|
||||||
|
|
||||||
|
if (System.getProperty("os.name").contains("Windows")) {
|
||||||
|
pb.command("cmd", "/c", "pg_dump -U " + settings.getDatabase().getUsername() + " -h " + matcher.group(1) + " -p " + matcher.group(2) + " -d " + matcher.group(3) + " --column-inserts -f " + dumpLocation.toAbsolutePath());
|
||||||
|
} else {
|
||||||
|
pb.command("/bin/sh", "-c", "pg_dump -U " + settings.getDatabase().getUsername() + " -h " + matcher.group(1) + " -p " + matcher.group(2) + " -d " + matcher.group(3) + " --column-inserts -f " + dumpLocation.toAbsolutePath());
|
||||||
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
log.info("Executing command {}", pb.command());
|
||||||
|
Process process = pb.start();
|
||||||
|
StringBuilder out = new StringBuilder();
|
||||||
|
BufferedReader reader = new BufferedReader(new InputStreamReader(process.getInputStream()));
|
||||||
|
|
||||||
|
String line;
|
||||||
|
while ((line = reader.readLine()) != null) {
|
||||||
|
out.append(line).append(System.lineSeparator());
|
||||||
|
}
|
||||||
|
int statusCode = process.waitFor();
|
||||||
|
if (statusCode != 0) {
|
||||||
|
log.info("Can't dump database, status code = {}. Make sure that pg-dump is installed or PostgreSQL version and pg-dump version is the same.", statusCode);
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Can't dump database, status code = " + statusCode);
|
||||||
|
}
|
||||||
|
} catch (IOException | InterruptedException e) {
|
||||||
|
log.info("Can't dump database, error={}", e.getMessage());
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Can't dump database, error = " + e.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void copyDirectories(RequestInfoUpdate requestInfoUpdate) {
|
||||||
|
Path dbfFolder = this.currentPath.resolve("dbf");
|
||||||
|
Path xmlFolder = this.currentPath.resolve("xml");
|
||||||
|
Path swtFolder = this.currentPath.resolve("swt");
|
||||||
|
try {
|
||||||
|
Files.createDirectories(dbfFolder);
|
||||||
|
Files.createDirectories(xmlFolder);
|
||||||
|
Files.createDirectories(swtFolder);
|
||||||
|
} catch (IOException e) {
|
||||||
|
log.info("Can't create directory, error={}", e.getMessage());
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Can't create directory, error = " + e.getMessage());
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
List<ExportedSDFFile> exportedFiles = this.fileStore.getFiles().stream()
|
||||||
|
.filter(f -> f.getStatus().equals(Status.Success))
|
||||||
|
.toList();
|
||||||
|
|
||||||
|
for (ExportedSDFFile file : exportedFiles) {
|
||||||
|
try {
|
||||||
|
if (file.getFileType().equals(FileType.DBF) && new File(file.getPath()).exists()) {
|
||||||
|
log.info("Copying from {} to {}", file.getPath(), dbfFolder.resolve(file.getFilename()));
|
||||||
|
Files.copy(Paths.get(file.getPath()), dbfFolder.resolve(file.getFilename()));
|
||||||
|
} else if (file.getFileType().equals(FileType.XML) && new File(file.getPath()).exists()) {
|
||||||
|
log.info("Copying from {} to {}", file.getPath(), xmlFolder.resolve(file.getFilename()));
|
||||||
|
Files.copy(Paths.get(file.getPath()), xmlFolder.resolve(file.getFilename()));
|
||||||
|
} else if (file.getFileType().equals(FileType.SWT) && new File(file.getPath()).exists()) {
|
||||||
|
log.info("Copying from {} to {}", file.getPath(), swtFolder.resolve(file.getFilename()));
|
||||||
|
Files.copy(Paths.get(file.getPath()), swtFolder.resolve(file.getFilename()));
|
||||||
|
} else {
|
||||||
|
log.info("File {} have wrong type {}", file.getPath(), file.getFileType());
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("File " + file.getPath() + "have wrong type " + file.getFileType());
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
} catch (IOException e) {
|
||||||
|
log.info("Can't copy file {}, error={}", file.getPath(), e.getMessage());
|
||||||
|
requestInfoUpdate.setStatus(Status.Error);
|
||||||
|
requestInfoUpdate.setMessage("Can't copy file " + file.getPath() + ", error = " + e.getMessage());
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,39 @@
|
||||||
|
package ru.spcex.clearing.snapshot.maker.utils;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.List;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.ExportedSDFFile;
|
||||||
|
|
||||||
|
@Component
|
||||||
|
public class FileStore {
|
||||||
|
private final List<ExportedSDFFile> files;
|
||||||
|
private boolean isSessionStarted;
|
||||||
|
|
||||||
|
public FileStore() {
|
||||||
|
this.files = new ArrayList<>();
|
||||||
|
this.isSessionStarted = false;
|
||||||
|
}
|
||||||
|
|
||||||
|
public List<ExportedSDFFile> getFiles() {
|
||||||
|
return files;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void addFile(ExportedSDFFile file) {
|
||||||
|
if (isSessionStarted) {
|
||||||
|
files.add(file);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public void clearFiles() {
|
||||||
|
files.clear();
|
||||||
|
}
|
||||||
|
|
||||||
|
public boolean isSessionStarted() {
|
||||||
|
return isSessionStarted;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setSessionStarted(boolean sessionStarted) {
|
||||||
|
isSessionStarted = sessionStarted;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,24 @@
|
||||||
|
# Database settings
|
||||||
|
snapshot-maker.database.jdbcUrl=jdbc:postgresql://localhost:5433/clearing
|
||||||
|
snapshot-maker.database.username=clearing
|
||||||
|
snapshot-maker.database.password=Aa111111
|
||||||
|
|
||||||
|
# Locations setting
|
||||||
|
snapshot-maker.store-location=/mnt/c/Users/ivan/Desktop/snapshot-test/out
|
||||||
|
|
||||||
|
# Kafka producer settings
|
||||||
|
snapshot-maker.kafka-producer.bootstrap-servers=localhost:9092
|
||||||
|
snapshot-maker.kafka-producer.acks=all
|
||||||
|
snapshot-maker.kafka-producer.retries=0
|
||||||
|
snapshot-maker.kafka-producer.batch-size=16384
|
||||||
|
snapshot-maker.kafka-producer.linger-ms=1
|
||||||
|
snapshot-maker.kafka-producer.buffer-memory=33554432
|
||||||
|
|
||||||
|
# Kafka consumer settings
|
||||||
|
snapshot-maker.kafka-consumer.bootstrap-servers=localhost:9092
|
||||||
|
snapshot-maker.kafka-consumer.group-id=dev-group-clearing-service
|
||||||
|
snapshot-maker.kafka-consumer.enable-auto-commit=false
|
||||||
|
snapshot-maker.kafka-consumer.session-timeout-ms=30000
|
||||||
|
snapshot-maker.kafka-consumer.auto-offset-reset=latest
|
||||||
|
snapshot-maker.kafka-consumer.linger-ms=1
|
||||||
|
snapshot-maker.kafka-consumer.buffer-memory=33554432
|
||||||
|
|
@ -1,22 +1,10 @@
|
||||||
package ru.spcex.clearing.swt.exporter.services;
|
package ru.spcex.clearing.swt.exporter.services;
|
||||||
|
|
||||||
import org.slf4j.Logger;
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.FILE_CREATED;
|
||||||
import org.slf4j.LoggerFactory;
|
import static ru.spcex.clearing.platform.messaging.domain.Consts.JOURNAL_SERVICE;
|
||||||
import ru.clearing.classes.statics.data.misc.Session;
|
|
||||||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest;
|
|
||||||
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.JournalEventExportedRequest;
|
|
||||||
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
|
||||||
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
|
||||||
import ru.spcex.clearing.swt.exporter.util.ConvertionContext;
|
|
||||||
import ru.spcex.clearing.swt.exporter.util.SWTHeaderData;
|
|
||||||
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
|
||||||
import ru.spcex.platform.enumeration.ResultStatuses;
|
|
||||||
import ru.spcex.platform.enumeration.SwtTable;
|
|
||||||
import ru.spcex.platform.imdg.api.Imdg;
|
|
||||||
import ru.spcex.platform.imdg.api.ImdgProvider;
|
|
||||||
|
|
||||||
import java.io.ByteArrayOutputStream;
|
import java.io.ByteArrayOutputStream;
|
||||||
|
import java.io.File;
|
||||||
import java.io.OutputStream;
|
import java.io.OutputStream;
|
||||||
import java.io.PrintWriter;
|
import java.io.PrintWriter;
|
||||||
import java.nio.charset.Charset;
|
import java.nio.charset.Charset;
|
||||||
|
|
@ -25,8 +13,24 @@ import java.time.format.DateTimeFormatter;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
import java.util.Objects;
|
import java.util.Objects;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
|
import org.slf4j.Logger;
|
||||||
import static ru.spcex.clearing.platform.messaging.domain.Consts.JOURNAL_SERVICE;
|
import org.slf4j.LoggerFactory;
|
||||||
|
import ru.clearing.classes.statics.data.misc.Session;
|
||||||
|
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.importexport.SwtExporterRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.ExportedSDFFile;
|
||||||
|
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.JournalEventExportedRequest;
|
||||||
|
import ru.spcex.clearing.platform.messaging.serialization.LogFormatter;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.FileType;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
|
||||||
|
import ru.spcex.clearing.swt.exporter.util.ConvertionContext;
|
||||||
|
import ru.spcex.clearing.swt.exporter.util.SWTHeaderData;
|
||||||
|
import ru.spcex.platform.classes.base.SpcexObjectBase;
|
||||||
|
import ru.spcex.platform.enumeration.ResultStatuses;
|
||||||
|
import ru.spcex.platform.enumeration.SwtTable;
|
||||||
|
import ru.spcex.platform.imdg.api.Imdg;
|
||||||
|
import ru.spcex.platform.imdg.api.ImdgProvider;
|
||||||
|
|
||||||
public abstract class AbstractExporterService<T extends SpcexObjectBase> {
|
public abstract class AbstractExporterService<T extends SpcexObjectBase> {
|
||||||
protected final Logger log = LoggerFactory.getLogger(getClass());
|
protected final Logger log = LoggerFactory.getLogger(getClass());
|
||||||
|
|
@ -83,7 +87,8 @@ public abstract class AbstractExporterService<T extends SpcexObjectBase> {
|
||||||
makeSWTData(swtHeaderData, records, outBuffer);
|
makeSWTData(swtHeaderData, records, outBuffer);
|
||||||
data = outBuffer.toByteArray();
|
data = outBuffer.toByteArray();
|
||||||
}
|
}
|
||||||
fileStorage.saveFile(fileName, data);
|
File createdFile = fileStorage.saveFile(fileName, data);
|
||||||
|
sendFileCreatedNotification(createdFile);
|
||||||
} catch (Exception e) { // IOException, ...
|
} catch (Exception e) { // IOException, ...
|
||||||
log.error("Failed export {} file", fileName);
|
log.error("Failed export {} file", fileName);
|
||||||
sendSwtExportedNotification(exportAt, null, ResultStatuses.notSuccess);
|
sendSwtExportedNotification(exportAt, null, ResultStatuses.notSuccess);
|
||||||
|
|
@ -105,6 +110,16 @@ public abstract class AbstractExporterService<T extends SpcexObjectBase> {
|
||||||
|
|
||||||
protected abstract String getDocumentNameForJournal();
|
protected abstract String getDocumentNameForJournal();
|
||||||
|
|
||||||
|
void sendFileCreatedNotification(File file) {
|
||||||
|
ExportedSDFFile exportedSDFFile = new ExportedSDFFile();
|
||||||
|
exportedSDFFile.setStatus(Status.Success);
|
||||||
|
exportedSDFFile.setFileType(FileType.SWT);
|
||||||
|
exportedSDFFile.setFilename(file.getName());
|
||||||
|
exportedSDFFile.setPath(file.getAbsolutePath());
|
||||||
|
log.debug("Send message to kafka \"{}\": {}", FILE_CREATED, LogFormatter.toStringWrapper(exportedSDFFile));
|
||||||
|
kafkaSender.sendRequestToQueue(FILE_CREATED, exportedSDFFile);
|
||||||
|
}
|
||||||
|
|
||||||
void sendSwtExportedNotification(LocalDateTime registrationAt, Long registrationNumber, ResultStatuses resultStatus) {
|
void sendSwtExportedNotification(LocalDateTime registrationAt, Long registrationNumber, ResultStatuses resultStatus) {
|
||||||
JournalEventExportedRequest exportedRequest = new JournalEventExportedRequest();
|
JournalEventExportedRequest exportedRequest = new JournalEventExportedRequest();
|
||||||
exportedRequest.setRegistratoinDate(registrationAt.toLocalDate());
|
exportedRequest.setRegistratoinDate(registrationAt.toLocalDate());
|
||||||
|
|
|
||||||
|
|
@ -37,7 +37,7 @@ public class FileStorage {
|
||||||
this.gateway = gateway;
|
this.gateway = gateway;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void saveFile(String fileName, byte[] data) throws IOException {
|
public File saveFile(String fileName, byte[] data) throws IOException {
|
||||||
File toFile = new File(outPath, fileName);
|
File toFile = new File(outPath, fileName);
|
||||||
FileUtils.writeByteArrayToFile(toFile, data);
|
FileUtils.writeByteArrayToFile(toFile, data);
|
||||||
if (gateway == null) {
|
if (gateway == null) {
|
||||||
|
|
@ -46,5 +46,6 @@ public class FileStorage {
|
||||||
log.debug("sftp is enabled. sending {}", toFile.getAbsolutePath());
|
log.debug("sftp is enabled. sending {}", toFile.getAbsolutePath());
|
||||||
gateway.sendToSftp(toFile);
|
gateway.sendToSftp(toFile);
|
||||||
}
|
}
|
||||||
|
return toFile;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -59,6 +59,9 @@ public enum Task implements IEnumKey {
|
||||||
finishBadSessions("CCLR"),//Завершение неудачных клиринговых сессий
|
finishBadSessions("CCLR"),//Завершение неудачных клиринговых сессий
|
||||||
sendLim_LIMC("LIMC"), // Выгрузка в Торговую систему остатков по валюте (отправка lim)"
|
sendLim_LIMC("LIMC"), // Выгрузка в Торговую систему остатков по валюте (отправка lim)"
|
||||||
makeFiles_MTCR("MTCR"), // Формирование файлов с МТКР
|
makeFiles_MTCR("MTCR"), // Формирование файлов с МТКР
|
||||||
|
startSessionSnapshot("STSS"), // Старт формирования снэпшота сессии
|
||||||
|
stopSessionSnapshot("SPSS"), // Завершение формирования снэпшота сессии
|
||||||
|
fileCreated("FCRD"), // SDF файл создан
|
||||||
;
|
;
|
||||||
|
|
||||||
private final String key;
|
private final String key;
|
||||||
|
|
|
||||||
|
|
@ -161,6 +161,7 @@ public interface Consts {
|
||||||
String BALANCE_ACCOUNT_UPDATE = "balance-account-update";
|
String BALANCE_ACCOUNT_UPDATE = "balance-account-update";
|
||||||
String CONTINUE_CLEARING = "continue-clearing";
|
String CONTINUE_CLEARING = "continue-clearing";
|
||||||
String LAUNCHER_NEW = "launcher-new";
|
String LAUNCHER_NEW = "launcher-new";
|
||||||
|
String FILE_CREATED = "file-created";
|
||||||
|
|
||||||
String DESTINATION_DEPO_ACCOUNT_SYMBOLS_NEW = "depo-accounts-symbols-new";
|
String DESTINATION_DEPO_ACCOUNT_SYMBOLS_NEW = "depo-accounts-symbols-new";
|
||||||
String DESTINATION_DEPO_ACCOUNT_SYMBOLS_DELETE = "depo-accounts-symbols-delete";
|
String DESTINATION_DEPO_ACCOUNT_SYMBOLS_DELETE = "depo-accounts-symbols-delete";
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,58 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.domain.cud.utilities;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.FileType;
|
||||||
|
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||||
|
|
||||||
|
public class ExportedSDFFile {
|
||||||
|
@JsonProperty
|
||||||
|
private Status status;
|
||||||
|
@JsonProperty
|
||||||
|
private FileType fileType;
|
||||||
|
@JsonProperty
|
||||||
|
private String filename;
|
||||||
|
@JsonProperty
|
||||||
|
private String path;
|
||||||
|
|
||||||
|
public Status getStatus() {
|
||||||
|
return status;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setStatus(Status status) {
|
||||||
|
this.status = status;
|
||||||
|
}
|
||||||
|
|
||||||
|
public FileType getFileType() {
|
||||||
|
return fileType;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setFileType(FileType fileType) {
|
||||||
|
this.fileType = fileType;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getFilename() {
|
||||||
|
return filename;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setFilename(String filename) {
|
||||||
|
this.filename = filename;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String getPath() {
|
||||||
|
return path;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setPath(String path) {
|
||||||
|
this.path = path;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String toString() {
|
||||||
|
return "ExportedSDFFile{" +
|
||||||
|
"status=" + status +
|
||||||
|
", fileType=" + fileType +
|
||||||
|
", filename='" + filename + '\'' +
|
||||||
|
", path='" + path + '\'' +
|
||||||
|
'}';
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,7 @@
|
||||||
|
package ru.spcex.clearing.platform.messaging.service;
|
||||||
|
|
||||||
|
public enum FileType {
|
||||||
|
DBF,
|
||||||
|
SWT,
|
||||||
|
XML
|
||||||
|
}
|
||||||
1
pom.xml
1
pom.xml
|
|
@ -53,6 +53,7 @@
|
||||||
<folder_root_registry-service>${folder_root_clearing}/clearing-parent/registry-service</folder_root_registry-service>
|
<folder_root_registry-service>${folder_root_clearing}/clearing-parent/registry-service</folder_root_registry-service>
|
||||||
<folder_root_scheduler-service>${folder_root_clearing}/clearing-parent/scheduler-service</folder_root_scheduler-service>
|
<folder_root_scheduler-service>${folder_root_clearing}/clearing-parent/scheduler-service</folder_root_scheduler-service>
|
||||||
<folder_root_gateway-api>${folder_root_clearing}/clearing-parent/gateway-api</folder_root_gateway-api>
|
<folder_root_gateway-api>${folder_root_clearing}/clearing-parent/gateway-api</folder_root_gateway-api>
|
||||||
|
<folder_root_snapshot-maker>${folder_root_clearing}/clearing-parent/snapshot-maker</folder_root_snapshot-maker>
|
||||||
<!-- IMDG -->
|
<!-- IMDG -->
|
||||||
<external_libraries.hazelcast.version>3.12.4</external_libraries.hazelcast.version>
|
<external_libraries.hazelcast.version>3.12.4</external_libraries.hazelcast.version>
|
||||||
<external_libraries.slf4j.version>1.7.33</external_libraries.slf4j.version>
|
<external_libraries.slf4j.version>1.7.33</external_libraries.slf4j.version>
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue