response = client.send(request, HttpResponse.BodyHandlers.ofString());
-
- if (response.statusCode() == 200) {
- log.info("ConfigMap '{}' успешно обновлён, тенантов: {}", CONFIGMAP_NAME, tenants.size());
- return true;
- } else {
- log.error("Не удалось обновить ConfigMap: HTTP {} — {}", response.statusCode(), response.body());
- return false;
- }
-
- } catch (Exception e) {
- log.error("Ошибка при обновлении ConfigMap: {}", e.getMessage());
- return false;
- }
- }
-
- /**
- * Создаёт HttpClient, который доверяет self-signed сертификатам K8s API.
- */
- private HttpClient createInsecureClient() {
- try {
- TrustManager[] trustAll = new TrustManager[]{
- new X509TrustManager() {
- public X509Certificate[] getAcceptedIssuers() { return new X509Certificate[0]; }
- public void checkClientTrusted(X509Certificate[] certs, String authType) {}
- public void checkServerTrusted(X509Certificate[] certs, String authType) {}
- }
- };
-
- SSLContext sslContext = SSLContext.getInstance("TLS");
- sslContext.init(null, trustAll, new SecureRandom());
-
- return HttpClient.newBuilder()
- .sslContext(sslContext)
- .build();
- } catch (Exception e) {
- log.warn("Не удалось создать клиент без проверки сертификата, используем стандартный: {}", e.getMessage());
- return HttpClient.newHttpClient();
- }
- }
-}
diff --git a/backend/src/main/java/com/magistr/app/config/tenant/KubernetesTenantSecretUpdater.java b/backend/src/main/java/com/magistr/app/config/tenant/KubernetesTenantSecretUpdater.java
new file mode 100644
index 0000000..0a6cf61
--- /dev/null
+++ b/backend/src/main/java/com/magistr/app/config/tenant/KubernetesTenantSecretUpdater.java
@@ -0,0 +1,178 @@
+package com.magistr.app.config.tenant;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.stereotype.Service;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.SSLParameters;
+import javax.net.ssl.TrustManagerFactory;
+import java.io.InputStream;
+import java.net.URI;
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.security.KeyStore;
+import java.security.cert.Certificate;
+import java.security.cert.CertificateFactory;
+import java.util.Base64;
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Обновляет Kubernetes Secret с конфигурацией тенантов через Kubernetes REST API.
+ *
+ * Работает только внутри Kubernetes pod и использует service-account token и
+ * service-account CA. В локальной разработке персистенция пропускается.
+ */
+@Service
+public class KubernetesTenantSecretUpdater {
+
+ private static final Logger log = LoggerFactory.getLogger(KubernetesTenantSecretUpdater.class);
+
+ private static final Path DEFAULT_TOKEN_PATH =
+ Path.of("/var/run/secrets/kubernetes.io/serviceaccount/token");
+ private static final Path DEFAULT_NAMESPACE_PATH =
+ Path.of("/var/run/secrets/kubernetes.io/serviceaccount/namespace");
+ private static final Path DEFAULT_CA_PATH =
+ Path.of("/var/run/secrets/kubernetes.io/serviceaccount/ca.crt");
+ private static final String DEFAULT_API_BASE = "https://kubernetes.default.svc";
+ private static final String DEFAULT_SECRET_NAME = "tenants-secret";
+
+ private final ObjectMapper objectMapper;
+ private final Path tokenPath;
+ private final Path namespacePath;
+ private final Path caPath;
+ private final String apiBase;
+ private final String secretName;
+ private final boolean runningInKubernetes;
+
+ public KubernetesTenantSecretUpdater() {
+ this(
+ DEFAULT_TOKEN_PATH,
+ DEFAULT_NAMESPACE_PATH,
+ DEFAULT_CA_PATH,
+ DEFAULT_API_BASE,
+ DEFAULT_SECRET_NAME,
+ new ObjectMapper()
+ );
+ }
+
+ KubernetesTenantSecretUpdater(Path tokenPath,
+ Path namespacePath,
+ Path caPath,
+ String apiBase,
+ String secretName,
+ ObjectMapper objectMapper) {
+ this.tokenPath = tokenPath;
+ this.namespacePath = namespacePath;
+ this.caPath = caPath;
+ this.apiBase = apiBase.endsWith("/") ? apiBase.substring(0, apiBase.length() - 1) : apiBase;
+ this.secretName = secretName;
+ this.objectMapper = objectMapper;
+ this.runningInKubernetes = Files.exists(tokenPath);
+
+ if (!runningInKubernetes) {
+ log.info("Приложение запущено вне Kubernetes — обновление tenant Secret будет пропущено");
+ }
+ }
+
+ /**
+ * Сохраняет полный список тенантов в ключе {@code tenants.json} Kubernetes Secret.
+ *
+ * @return {@code true}, если Secret обновлён или приложение запущено вне Kubernetes
+ */
+ public boolean updateTenantsConfig(List tenants) {
+ if (!runningInKubernetes) {
+ log.warn("Приложение запущено вне Kubernetes, персистенция tenant Secret пропущена");
+ return true;
+ }
+
+ try {
+ String token = Files.readString(tokenPath).trim();
+ String namespace = Files.readString(namespacePath).trim();
+ if (token.isBlank() || namespace.isBlank()) {
+ throw new IllegalStateException("ServiceAccount token или namespace не настроены");
+ }
+
+ String tenantsJson = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(tenants);
+ String encodedTenants = Base64.getEncoder()
+ .encodeToString(tenantsJson.getBytes(StandardCharsets.UTF_8));
+ String patchBody = objectMapper.writeValueAsString(Map.of(
+ "data", Map.of("tenants.json", encodedTenants)
+ ));
+
+ URI uri = URI.create(String.format(
+ "%s/api/v1/namespaces/%s/secrets/%s",
+ apiBase,
+ namespace,
+ secretName
+ ));
+ HttpRequest request = HttpRequest.newBuilder()
+ .uri(uri)
+ .header("Authorization", "Bearer " + token)
+ .header("Content-Type", "application/strategic-merge-patch+json")
+ .method("PATCH", HttpRequest.BodyPublishers.ofString(patchBody))
+ .build();
+
+ HttpResponse response = createSecureClient(caPath)
+ .send(request, HttpResponse.BodyHandlers.discarding());
+ if (response.statusCode() == 200) {
+ log.info("Tenant Secret успешно обновлён: tenantCount={}", tenants.size());
+ return true;
+ }
+
+ log.error("Kubernetes API отклонил обновление tenant Secret: httpStatus={}", response.statusCode());
+ return false;
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ log.error("Обновление tenant Secret прервано");
+ log.debug("Технические детали прерывания tenant Secret", e);
+ return false;
+ } catch (Exception e) {
+ log.error("Не удалось безопасно обновить tenant Secret: errorType={}",
+ e.getClass().getSimpleName());
+ log.debug("Технические детали ошибки tenant Secret", e);
+ return false;
+ }
+ }
+
+ HttpClient createSecureClient(Path serviceAccountCaPath) throws Exception {
+ CertificateFactory certificateFactory = CertificateFactory.getInstance("X.509");
+ Collection extends Certificate> certificates;
+ try (InputStream input = Files.newInputStream(serviceAccountCaPath)) {
+ certificates = certificateFactory.generateCertificates(input);
+ }
+ if (certificates.isEmpty()) {
+ throw new IllegalStateException("ServiceAccount CA не содержит сертификатов");
+ }
+
+ KeyStore trustStore = KeyStore.getInstance(KeyStore.getDefaultType());
+ trustStore.load(null, null);
+ int index = 0;
+ for (Certificate certificate : certificates) {
+ trustStore.setCertificateEntry("kubernetes-ca-" + index++, certificate);
+ }
+
+ TrustManagerFactory trustManagerFactory = TrustManagerFactory.getInstance(
+ TrustManagerFactory.getDefaultAlgorithm()
+ );
+ trustManagerFactory.init(trustStore);
+
+ SSLContext sslContext = SSLContext.getInstance("TLS");
+ sslContext.init(null, trustManagerFactory.getTrustManagers(), null);
+
+ SSLParameters sslParameters = new SSLParameters();
+ sslParameters.setEndpointIdentificationAlgorithm("HTTPS");
+
+ return HttpClient.newBuilder()
+ .sslContext(sslContext)
+ .sslParameters(sslParameters)
+ .build();
+ }
+}
diff --git a/backend/src/main/java/com/magistr/app/config/tenant/TenantConfigWatcher.java b/backend/src/main/java/com/magistr/app/config/tenant/TenantConfigWatcher.java
index e27fb62..a22539d 100755
--- a/backend/src/main/java/com/magistr/app/config/tenant/TenantConfigWatcher.java
+++ b/backend/src/main/java/com/magistr/app/config/tenant/TenantConfigWatcher.java
@@ -2,25 +2,23 @@ package com.magistr.app.config.tenant;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.magistr.app.service.TenantLifecycleService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
-import org.springframework.core.io.ClassPathResource;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
-import javax.sql.DataSource;
import java.io.File;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.HexFormat;
-import java.util.*;
-import java.util.stream.Collectors;
+import java.util.List;
/**
- * Периодически перечитывает tenants.json (mounted ConfigMap).
- * Если ConfigMap был обновлён через K8s API, этот компонент
+ * Периодически перечитывает tenants.json (mounted Secret).
+ * Если Secret был обновлён через K8s API, этот компонент
* подхватит изменения и синхронизирует in-memory datasource'ы.
*
* Также отвечает за инициализацию БД (init.sql) для новых тенантов.
@@ -30,8 +28,7 @@ public class TenantConfigWatcher {
private static final Logger log = LoggerFactory.getLogger(TenantConfigWatcher.class);
- private final TenantRoutingDataSource routingDataSource;
- private final DataSource dataSource;
+ private final TenantLifecycleService tenantLifecycleService;
private final ObjectMapper objectMapper = new ObjectMapper();
@Value("${app.tenants.config-path:tenants.json}")
@@ -40,9 +37,8 @@ public class TenantConfigWatcher {
// Хеш последнего прочитанного конфига — чтобы не перезагружать зря
private String lastConfigHash = "";
- public TenantConfigWatcher(TenantRoutingDataSource routingDataSource, DataSource dataSource) {
- this.routingDataSource = routingDataSource;
- this.dataSource = dataSource;
+ public TenantConfigWatcher(TenantLifecycleService tenantLifecycleService) {
+ this.tenantLifecycleService = tenantLifecycleService;
}
/**
@@ -50,6 +46,10 @@ public class TenantConfigWatcher {
*/
@Scheduled(fixedDelay = 30_000, initialDelay = 30_000)
public void watchForChanges() {
+ tenantLifecycleService.executeSerialized(this::watchForChangesSerialized);
+ }
+
+ private void watchForChangesSerialized() {
try {
File file = new File(tenantsConfigPath);
if (!file.exists()) return;
@@ -62,20 +62,24 @@ public class TenantConfigWatcher {
}
log.info("Обнаружено изменение tenants.json (хеш: {} -> {}), перечитываем конфиг", lastConfigHash, hash);
- lastConfigHash = hash;
-
List newTenants = objectMapper.readValue(content, new TypeReference<>() {});
syncTenants(newTenants);
+ lastConfigHash = hash;
} catch (Exception e) {
- log.error("Ошибка при проверке конфига тенантов: {}", e.getMessage());
+ log.error("Ошибка при проверке конфига тенантов: errorType={}", e.getClass().getSimpleName());
+ log.debug("Технические детали синхронизации конфига тенантов", e);
}
}
/**
- * Обновляет хеш конфига (вызывается после ручного обновления ConfigMap с этого же пода).
+ * Обновляет хеш конфига после ручного обновления Secret с этого же пода.
*/
public void refreshHash() {
+ tenantLifecycleService.executeSerialized(this::refreshHashSerialized);
+ }
+
+ private void refreshHashSerialized() {
try {
File file = new File(tenantsConfigPath);
if (file.exists()) {
@@ -83,7 +87,9 @@ public class TenantConfigWatcher {
lastConfigHash = configHash(content);
}
} catch (Exception e) {
- log.warn("Не удалось обновить хеш конфига тенантов: {}", e.getMessage());
+ log.warn("Не удалось обновить хеш конфига тенантов: errorType={}",
+ e.getClass().getSimpleName());
+ log.debug("Технические детали обновления хеша тенантов", e);
}
}
@@ -100,63 +106,7 @@ public class TenantConfigWatcher {
* Синхронизирует in-memory тенантов с конфигом из файла.
*/
private void syncTenants(List newTenants) {
- Map current = routingDataSource.getTenantConfigs();
- Set newDomains = newTenants.stream()
- .map(t -> t.getDomain().toLowerCase())
- .collect(Collectors.toSet());
-
- // Добавить новые тенанты
- for (TenantConfig tenant : newTenants) {
- String domain = tenant.getDomain().toLowerCase();
- if (!current.containsKey(domain)) {
- log.info("Добавляем нового тенанта '{}' из обновлённого ConfigMap", domain);
- routingDataSource.addTenant(tenant);
- // Инициализируем БД для нового тенанта
- initDatabaseForTenant(tenant);
- }
- }
-
- // Удалить тенанты, которых больше нет в конфиге
- for (String existingDomain : new ArrayList<>(current.keySet())) {
- if (!newDomains.contains(existingDomain)) {
- log.info("Удаляем тенанта '{}' — его больше нет в ConfigMap", existingDomain);
- routingDataSource.removeTenant(existingDomain);
- }
- }
+ tenantLifecycleService.synchronizeFromPersistedConfig(newTenants);
}
- /**
- * Выполняет миграции Flyway для конкретного тенанта пи подключении.
- * Если БД уже существует, но история Flyway пуста —
- * делает baseline (считает V1_init.sql уже выполненным).
- */
- public void initDatabaseForTenant(TenantConfig tenant) {
- String domain = tenant.getDomain();
- try {
- TenantContext.setCurrentTenant(domain);
-
- log.info("[{}] Запускаем миграции Flyway", domain);
-
- // Получаем DataSource конкретно для этого тенанта
- javax.sql.DataSource tenantDs = routingDataSource.getResolvedDataSources().get(domain);
- if (tenantDs == null) {
- // Если ещё не resolve'нулся (первый запуск), берём обёртку
- tenantDs = dataSource;
- }
-
- org.flywaydb.core.Flyway flyway = org.flywaydb.core.Flyway.configure()
- .dataSource(tenantDs)
- .baselineOnMigrate(true)
- .baselineVersion("1")
- .load();
-
- flyway.migrate();
- log.info("[{}] Миграции Flyway успешно выполнены", domain);
-
- } catch (Exception e) {
- log.error("[{}] Ошибка миграции Flyway: {}", domain, e.getMessage());
- } finally {
- TenantContext.clear();
- }
- }
}
diff --git a/backend/src/main/java/com/magistr/app/config/tenant/TenantDataSourceConfig.java b/backend/src/main/java/com/magistr/app/config/tenant/TenantDataSourceConfig.java
index 16d9ff6..134cc03 100755
--- a/backend/src/main/java/com/magistr/app/config/tenant/TenantDataSourceConfig.java
+++ b/backend/src/main/java/com/magistr/app/config/tenant/TenantDataSourceConfig.java
@@ -2,6 +2,8 @@ package com.magistr.app.config.tenant;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.magistr.app.service.TenantDatabaseMigrationService;
+import com.zaxxer.hikari.HikariDataSource;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
@@ -17,6 +19,7 @@ import jakarta.persistence.EntityManagerFactory;
import javax.sql.DataSource;
import java.io.File;
import java.io.IOException;
+import java.sql.Connection;
import java.util.*;
/**
@@ -45,7 +48,7 @@ public class TenantDataSourceConfig {
@Bean
@Primary
- public DataSource dataSource() {
+ public DataSource dataSource(TenantDatabaseMigrationService migrationService) {
TenantRoutingDataSource routingDataSource = new TenantRoutingDataSource();
// Загружаем тенантов из JSON (read-only ConfigMap mount)
@@ -57,16 +60,12 @@ public class TenantDataSourceConfig {
"Default", "default", defaultDbUrl, defaultDbUsername, defaultDbPassword
);
tenants.add(defaultTenant);
- log.info("No tenants config found, using default datasource: {}", defaultDbUrl);
+ log.info("Конфигурация тенантов отсутствует, используется default DataSource");
}
// Регистрируем тенантов
for (TenantConfig tenant : tenants) {
- try {
- routingDataSource.addTenant(tenant);
- } catch (Exception e) {
- log.error("Не удалось добавить тенанта '{}': {}", tenant.getDomain(), e.getMessage());
- }
+ registerPreparedTenant(routingDataSource, migrationService, tenant);
}
// Если всё ещё нет ни одного тенанта — H2 in-memory заглушка
@@ -80,7 +79,9 @@ public class TenantDataSourceConfig {
"jdbc:h2:mem:placeholder;DB_CLOSE_DELAY=-1",
"sa", ""
);
- routingDataSource.addTenant(h2Fallback);
+ if (!registerPreparedTenant(routingDataSource, migrationService, h2Fallback)) {
+ throw new IllegalStateException("Не удалось создать резервный H2 DataSource");
+ }
}
return routingDataSource;
@@ -120,18 +121,64 @@ public class TenantDataSourceConfig {
private List loadTenantsFromFile() {
File file = new File(tenantsConfigPath);
if (!file.exists()) {
- log.info("Tenants config file not found: {}", tenantsConfigPath);
+ log.info("Файл конфигурации тенантов не найден");
return new ArrayList<>();
}
try {
ObjectMapper mapper = new ObjectMapper();
List list = mapper.readValue(file, new TypeReference<>() {});
- log.info("Loaded {} tenant(s) from {}", list.size(), tenantsConfigPath);
+ log.info("Загружено конфигураций тенантов: {}", list.size());
return list;
} catch (IOException e) {
- log.error("Не удалось прочитать конфиг тенантов: {}", e.getMessage());
+ log.error("Не удалось прочитать конфигурацию тенантов: errorType={}",
+ e.getClass().getSimpleName());
+ log.debug("Технические детали чтения конфигурации тенантов", e);
return new ArrayList<>();
}
}
+
+ private boolean registerPreparedTenant(TenantRoutingDataSource routingDataSource,
+ TenantDatabaseMigrationService migrationService,
+ TenantConfig tenant) {
+ HikariDataSource candidate = null;
+ boolean activated = false;
+ try {
+ candidate = routingDataSource.prepareTenantDataSource(tenant);
+ try (Connection connection = candidate.getConnection()) {
+ if (!connection.isValid(5)) {
+ throw new IllegalStateException("База данных не подтвердила готовность подключения");
+ }
+ }
+ if (!tenant.getUrl().trim().startsWith("jdbc:h2:")) {
+ migrationService.migrate(candidate);
+ }
+ TenantRoutingDataSource.TenantState previous = routingDataSource.swapTenant(tenant, candidate);
+ activated = true;
+ closeDataSource(previous == null ? null : previous.dataSource());
+ log.info("Tenant-БД '{}' проверена и активирована", tenant.getDomain());
+ return true;
+ } catch (Exception startupFailure) {
+ log.error("Не удалось безопасно активировать tenant-БД '{}': errorType={}",
+ tenant.getDomain(), startupFailure.getClass().getSimpleName());
+ log.debug("Технические детали startup lifecycle tenant-БД", startupFailure);
+ return false;
+ } finally {
+ if (!activated) {
+ closeDataSource(candidate);
+ }
+ }
+ }
+
+ private void closeDataSource(DataSource dataSource) {
+ if (dataSource instanceof HikariDataSource hikariDataSource) {
+ try {
+ hikariDataSource.close();
+ } catch (RuntimeException closeFailure) {
+ log.warn("Не удалось штатно закрыть startup Hikari pool: errorType={}",
+ closeFailure.getClass().getSimpleName());
+ log.debug("Технические детали закрытия startup Hikari pool", closeFailure);
+ }
+ }
+ }
}
diff --git a/backend/src/main/java/com/magistr/app/config/tenant/TenantRoutingDataSource.java b/backend/src/main/java/com/magistr/app/config/tenant/TenantRoutingDataSource.java
old mode 100755
new mode 100644
index ea152b0..d06ae19
--- a/backend/src/main/java/com/magistr/app/config/tenant/TenantRoutingDataSource.java
+++ b/backend/src/main/java/com/magistr/app/config/tenant/TenantRoutingDataSource.java
@@ -8,143 +8,252 @@ import org.springframework.jdbc.datasource.lookup.AbstractRoutingDataSource;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.SQLException;
-import java.util.HashMap;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.Locale;
import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
+import java.util.Objects;
+import java.util.Optional;
/**
* DataSource, который переключается между БД разных тенантов.
- * На каждый запрос determineCurrentLookupKey() возвращает текущий тенант из TenantContext.
+ * Runtime-состояние публикуется одним неизменяемым snapshot, поэтому запрос
+ * никогда не видит конфигурацию и DataSource из разных версий.
*/
public class TenantRoutingDataSource extends AbstractRoutingDataSource {
private static final Logger log = LoggerFactory.getLogger(TenantRoutingDataSource.class);
- private final Map tenantConfigs = new ConcurrentHashMap<>();
- private final Map