Update cmd/futriis/main.go
This commit is contained in:
1 parent
7804ce2338
commit
1607f4108d
1 file changed
+20
-101
+20
-101
@@ -10,58 +10,13 @@
|
|||||||
* Файл: cmd/futriis/main.go
|
* Файл: cmd/futriis/main.go
|
||||||
* Назначение: Точка входа в приложение СУБД futriis.
|
* Назначение: Точка входа в приложение СУБД futriis.
|
||||||
*
|
*
|
||||||
* ИСПРАВЛЕНО:
|
* ДОБАВЛЕНО (2026-10, динамическая конфигурация):
|
||||||
* - Валидация пути к конфигурации и проверка прав доступа к файлу
|
* После создания REPL передаём ему ссылку на DynamicConfigManager
|
||||||
* - Валидация конфигурации после загрузки (ValidateConfig)
|
* (через SetDynamicConfigManager), чтобы команды `config set`,
|
||||||
* - WAL и data-dir привязаны к директории futriis, а не CWD
|
* `compression config set`, `cluster pipeline on/off` могли
|
||||||
* - Graceful shutdown: убраны os.Exit(1) до завершения defer'ов,
|
* применять изменения через Raft и сохранять их на всех узлах.
|
||||||
* используется переменная exitCode и os.Exit в самом конце
|
* В fallback-режиме (Raft недоступен или узел не лидер) REPL сам
|
||||||
* - Ожидание запуска HTTP-сервера с таймаутом
|
* запишет изменения в config.toml через config.SaveConfig.
|
||||||
* - SAGA оркестратор останавливается всегда (даже если disabled)
|
|
||||||
* - Проверка allow-list для плагинов перед StartPlugin
|
|
||||||
* - Заменён SIGHUP на SIGUSR1 (reload) + SIGTERM/SIGINT (shutdown)
|
|
||||||
* - Проверка nil для raftCoordinator перед использованием
|
|
||||||
* - repl.NewRepl теперь возвращает (*Repl, error)
|
|
||||||
* - Интеграция с Prometheus: /metrics endpoint
|
|
||||||
* - ДОБАВЛЕНО: поддержка TLS на HTTP API (вариант A — TLS внутри futriiX).
|
|
||||||
* Теперь используется api.NewHTTPServerWithTLS(..., &cfg.Security).
|
|
||||||
* - ИСПРАВЛЕНО (UI): убраны лишние пустые строки перед баннером, чтобы
|
|
||||||
* строка "futriis 3i²(by 29.09.2026)" шла сразу после сообщения ACL
|
|
||||||
* о созданном admin-пользователе.
|
|
||||||
* - ИСПРАВЛЕНО (UI): текст баннера окрашивается в точный #00bfff
|
|
||||||
* (Deep Sky Blue) через utils.PrintlnDeepSkyBlueReset, чтобы совпадать
|
|
||||||
* с цветом надписи "futriis 3i²(by 29.09.2026)".
|
|
||||||
* - ИСПРАВЛЕНО: acl.NewACLManager() теперь возвращает (*ACLManager, error).
|
|
||||||
* Обрабатываем ошибку и завершаемся через os.Exit(1) с логом,
|
|
||||||
* без panic и stack trace.
|
|
||||||
*
|
|
||||||
* ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции
|
|
||||||
* функции storage.InitTransactionManager, storage.SetTransactionLogger,
|
|
||||||
* storage.StopTransactionManager были удалены. Заменены на:
|
|
||||||
* - storage.InitBatchManager(walPath) — инициализирует batch-менеджер,
|
|
||||||
* WAL и AOF (append-only file для мутаций вне batch);
|
|
||||||
* - storage.SetBatchLogger(logger) — устанавливает логгер;
|
|
||||||
* - storage.StopBatchManager() — останавливает batch-менеджер,
|
|
||||||
* WAL и AOF.
|
|
||||||
* Дополнительно: при shutdown логируется состояние batch-статистики
|
|
||||||
* (количество выполненных batch'ей, операций, ошибок) для диагностики.
|
|
||||||
*
|
|
||||||
* ДОБАВЛЕНО (2026-10, выборы лидера SAGA через Raft):
|
|
||||||
* - SAGA-оркестратор теперь создаётся через
|
|
||||||
* storage.NewSagaOrchestratorWithRaft и получает bridge из
|
|
||||||
* raftCoordinator (raftCoordinator.GetSagaLeaderBridge()).
|
|
||||||
* Это позволяет в кластерном режиме выбирать SAGA-лидера через Raft
|
|
||||||
* и выполнять компенсации только на одном узле.
|
|
||||||
* - В single-node режиме bridge также возвращает корректные значения
|
|
||||||
* (IsLeader() == true), так что поведение не меняется.
|
|
||||||
* - Раньше оркестратор создавался через NewSagaOrchestratorWithConfig
|
|
||||||
* с raftAccessor == nil, что в кластерном режиме приводило к тому,
|
|
||||||
* что каждый узел считал себя SAGA-лидером. Идемпотентность шагов
|
|
||||||
* через saga_executed_keys.json защищала от дублей операций, но не
|
|
||||||
* от одновременной компенсации одной SAGA с двух узлов.
|
|
||||||
* - Порядок shutdown безопасен: sagaOrchestrator.Stop() вызывается
|
|
||||||
* раньше raftCoordinator.Stop(), поэтому bridge ещё жив, когда
|
|
||||||
* оркестратор отписывается от него.
|
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package main
|
package main
|
||||||
@@ -146,13 +101,6 @@ func main() {
|
|||||||
store := storage.NewStorage(cfg.Storage.PageSizeMB, logger)
|
store := storage.NewStorage(cfg.Storage.PageSizeMB, logger)
|
||||||
storage.SetGlobalStorage(store)
|
storage.SetGlobalStorage(store)
|
||||||
|
|
||||||
// ========== BATCH MANAGER (замена TransactionManager) ==========
|
|
||||||
//
|
|
||||||
// ИСПРАВЛЕНО (2026-10): InitTransactionManager заменён на InitBatchManager.
|
|
||||||
// InitBatchManager инициализирует:
|
|
||||||
// - WAL (сегментированный журнал для batch-операций);
|
|
||||||
// - AOF (append-only file для мутаций вне batch).
|
|
||||||
// Директория AOF создаётся автоматически как filepath.Join(filepath.Dir(walPath), "aof").
|
|
||||||
walPath := filepath.Join(dataDir, "futriis.wal")
|
walPath := filepath.Join(dataDir, "futriis.wal")
|
||||||
if err := storage.InitBatchManager(walPath); err != nil {
|
if err := storage.InitBatchManager(walPath); err != nil {
|
||||||
logger.Warn("Failed to initialize batch manager: " + err.Error())
|
logger.Warn("Failed to initialize batch manager: " + err.Error())
|
||||||
@@ -166,8 +114,6 @@ func main() {
|
|||||||
storage.InitTriggerManager(logger)
|
storage.InitTriggerManager(logger)
|
||||||
logger.Info("Trigger manager initialized")
|
logger.Info("Trigger manager initialized")
|
||||||
|
|
||||||
// ИСПРАВЛЕНО: NewACLManager теперь возвращает (*ACLManager, error).
|
|
||||||
// Обрабатываем ошибку управляемо (без panic) и завершаемся с кодом 1.
|
|
||||||
aclManager, err := acl.NewACLManager()
|
aclManager, err := acl.NewACLManager()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("Failed to initialize ACL manager: " + err.Error())
|
logger.Error("Failed to initialize ACL manager: " + err.Error())
|
||||||
@@ -184,27 +130,11 @@ func main() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ========== SAGA ORCHESTRATOR ==========
|
// ========== SAGA ORCHESTRATOR ==========
|
||||||
//
|
|
||||||
// ДОБАВЛЕНО (2026-10): используем NewSagaOrchestratorWithRaft и
|
|
||||||
// передаём bridge из raftCoordinator. Это позволяет в кластерном
|
|
||||||
// режиме выбирать SAGA-лидера через Raft, чтобы компенсации
|
|
||||||
// выполнялись только на одном узле.
|
|
||||||
//
|
|
||||||
// В single-node режиме bridge возвращает IsLeader() == true, так что
|
|
||||||
// поведение не меняется. В кластерном режиме только лидер Raft
|
|
||||||
// будет IsSagaLeader() == true, и только он будет выполнять SAGA.
|
|
||||||
//
|
|
||||||
// Если raftCoordinator == nil (не должно случаться — выше os.Exit(1)
|
|
||||||
// при ошибке), создаём оркестратор без Raft (single-node fallback).
|
|
||||||
var sagaOrchestrator *storage.SagaOrchestrator
|
var sagaOrchestrator *storage.SagaOrchestrator
|
||||||
if cfg.Saga.Enabled {
|
if cfg.Saga.Enabled {
|
||||||
sagaConfig := &cfg.Saga
|
sagaConfig := &cfg.Saga
|
||||||
nodeID := cfg.Cluster.NodeIP + ":" + fmt.Sprint(cfg.Cluster.NodePort)
|
nodeID := cfg.Cluster.NodeIP + ":" + fmt.Sprint(cfg.Cluster.NodePort)
|
||||||
|
|
||||||
// Передаём bridge из Raft в качестве storage.SagaRaftAccessor.
|
|
||||||
// Тип bridge реализует интерфейс storage.SagaRaftAccessor
|
|
||||||
// (IsSagaLeader, GetSagaLeaderID, GetSagaCurrentTerm,
|
|
||||||
// RegisterSagaLeaderObserver).
|
|
||||||
var raftAccessor storage.SagaRaftAccessor
|
var raftAccessor storage.SagaRaftAccessor
|
||||||
if raftCoordinator != nil {
|
if raftCoordinator != nil {
|
||||||
raftAccessor = raftCoordinator.GetSagaLeaderBridge()
|
raftAccessor = raftCoordinator.GetSagaLeaderBridge()
|
||||||
@@ -311,8 +241,6 @@ func main() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ========== HTTP API SERVER ==========
|
// ========== HTTP API SERVER ==========
|
||||||
// ДОБАВЛЕНО: используем NewHTTPServerWithTLS и передаём cfg.Security.
|
|
||||||
// Если в config.toml [security].enable_tls = true, сервер поднимется по HTTPS.
|
|
||||||
httpPort := cfg.API.Port
|
httpPort := cfg.API.Port
|
||||||
httpServer := api.NewHTTPServerWithTLS(httpPort, store, raftCoordinator, aclManager, logger, &cfg.Security)
|
httpServer := api.NewHTTPServerWithTLS(httpPort, store, raftCoordinator, aclManager, logger, &cfg.Security)
|
||||||
if metricsCollector != nil {
|
if metricsCollector != nil {
|
||||||
@@ -362,10 +290,6 @@ func main() {
|
|||||||
logger.Info(" - Automatic data persistence - ENABLED")
|
logger.Info(" - Automatic data persistence - ENABLED")
|
||||||
}
|
}
|
||||||
|
|
||||||
// ИСПРАВЛЕНО (UI): убраны лишние пустые строки перед баннером.
|
|
||||||
// Внутри displayBanner уже есть utils.Println("") первой строкой,
|
|
||||||
// чтобы отделить баннер от ACL-сообщения. Строка "futriis 3i²(by
|
|
||||||
// 29.09.2026)" идёт сразу после сообщения ACL.
|
|
||||||
displayBanner(cfg.Cluster.Name, httpPort, raftCoordinator, cfg.Saga.Enabled, cfg.Metrics.Enabled)
|
displayBanner(cfg.Cluster.Name, httpPort, raftCoordinator, cfg.Saga.Enabled, cfg.Metrics.Enabled)
|
||||||
|
|
||||||
replInstance, err := repl.NewRepl(store, raftCoordinator, logger, cfg, aclManager, pluginManager)
|
replInstance, err := repl.NewRepl(store, raftCoordinator, logger, cfg, aclManager, pluginManager)
|
||||||
@@ -376,6 +300,19 @@ func main() {
|
|||||||
}
|
}
|
||||||
defer replInstance.Close()
|
defer replInstance.Close()
|
||||||
|
|
||||||
|
// ДОБАВЛЕНО (2026-10): передаём REPL ссылку на DynamicConfigManager,
|
||||||
|
// чтобы команды `config set`, `compression config set`,
|
||||||
|
// `cluster pipeline on/off` могли применять изменения через Raft.
|
||||||
|
// В fallback-режиме REPL пишет config.toml через config.SaveConfig.
|
||||||
|
if raftCoordinator != nil {
|
||||||
|
replInstance.SetDynamicConfigManager(
|
||||||
|
raftCoordinator.GetDynamicConfigManager(),
|
||||||
|
configPath,
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
replInstance.SetDynamicConfigManager(nil, configPath)
|
||||||
|
}
|
||||||
|
|
||||||
// ========== GRACEFUL SHUTDOWN ==========
|
// ========== GRACEFUL SHUTDOWN ==========
|
||||||
sigChan := make(chan os.Signal, 1)
|
sigChan := make(chan os.Signal, 1)
|
||||||
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM, syscall.SIGUSR1)
|
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM, syscall.SIGUSR1)
|
||||||
@@ -431,10 +368,6 @@ func main() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ВАЖНО: sagaOrchestrator останавливается ДО raftCoordinator.
|
|
||||||
// Это гарантирует, что оркестратор отпишется от sagaLeaderBridge
|
|
||||||
// (и остановит leaderElector) до того, как bridge будет остановлен
|
|
||||||
// вместе с координатором. Двойной остановки bridge не будет.
|
|
||||||
if sagaOrchestrator != nil {
|
if sagaOrchestrator != nil {
|
||||||
logger.Info("Stopping SAGA orchestrator...")
|
logger.Info("Stopping SAGA orchestrator...")
|
||||||
sagaStopDone := make(chan struct{})
|
sagaStopDone := make(chan struct{})
|
||||||
@@ -497,8 +430,6 @@ func main() {
|
|||||||
logger.Warn("Node shutdown timeout")
|
logger.Warn("Node shutdown timeout")
|
||||||
}
|
}
|
||||||
|
|
||||||
// ИСПРАВЛЕНО (2026-10): StopTransactionManager заменён на StopBatchManager.
|
|
||||||
// StopBatchManager синхронизирует WAL, закрывает его и закрывает AOF.
|
|
||||||
logger.Info("Stopping batch manager (WAL + AOF)...")
|
logger.Info("Stopping batch manager (WAL + AOF)...")
|
||||||
if stats := storage.GetBatchStats(); stats != nil {
|
if stats := storage.GetBatchStats(); stats != nil {
|
||||||
logger.Info(fmt.Sprintf(" Batch stats: total_batches=%v, total_operations=%v, total_failed=%v, active_batches=%v",
|
logger.Info(fmt.Sprintf(" Batch stats: total_batches=%v, total_operations=%v, total_failed=%v, active_batches=%v",
|
||||||
@@ -515,7 +446,6 @@ func main() {
|
|||||||
}()
|
}()
|
||||||
select {
|
select {
|
||||||
case <-batchStopDone:
|
case <-batchStopDone:
|
||||||
// уже залогировано внутри горутины
|
|
||||||
case <-time.After(10 * time.Second):
|
case <-time.After(10 * time.Second):
|
||||||
logger.Warn("Batch manager shutdown timeout")
|
logger.Warn("Batch manager shutdown timeout")
|
||||||
}
|
}
|
||||||
@@ -679,15 +609,6 @@ func isPortListening(port int) bool {
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
// displayBanner выводит стартовый баннер.
|
|
||||||
//
|
|
||||||
// ИСПРАВЛЕНО (UI): текст баннера окрашивается точным цветом
|
|
||||||
// #00bfff (Deep Sky Blue) через utils.PrintlnDeepSkyBlueReset, чтобы
|
|
||||||
// совпадать с цветом надписи "futriis 3i²(by 29.09.2026)".
|
|
||||||
//
|
|
||||||
// Внутри уже есть utils.Println("") первой строкой, чтобы отделить
|
|
||||||
// баннер от ACL-сообщения. Дополнительные пустые строки в main() не
|
|
||||||
// добавляются (см. комментарий перед вызовом displayBanner).
|
|
||||||
func displayBanner(clusterName string, httpPort int, coordinator *cluster.RaftCoordinator, sagaEnabled, metricsEnabled bool) {
|
func displayBanner(clusterName string, httpPort int, coordinator *cluster.RaftCoordinator, sagaEnabled, metricsEnabled bool) {
|
||||||
utils.Println("")
|
utils.Println("")
|
||||||
bannerLines := []string{
|
bannerLines := []string{
|
||||||
@@ -738,8 +659,6 @@ func displayBanner(clusterName string, httpPort int, coordinator *cluster.RaftCo
|
|||||||
}...)
|
}...)
|
||||||
|
|
||||||
for _, line := range bannerLines {
|
for _, line := range bannerLines {
|
||||||
// ИСПРАВЛЕНО: используем точный цвет #00bfff (Deep Sky Blue),
|
|
||||||
// совпадающий с цветом заголовка баннера.
|
|
||||||
utils.PrintlnDeepSkyBlueReset(line)
|
utils.PrintlnDeepSkyBlueReset(line)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in new issue
Block a user