Update cmd/futriis/main.go
This commit is contained in:
1 parent
29a3766f03
commit
d69b29d381
1 file changed
+229
-84
+229
-84
@@ -6,11 +6,27 @@
|
|||||||
*
|
*
|
||||||
* You may obtain a copy of the License at
|
* You may obtain a copy of the License at
|
||||||
* https://opensource.org/licenses/CDDL-1.0
|
* https://opensource.org/licenses/CDDL-1.0
|
||||||
|
*
|
||||||
// Файл: cmd/futriis/main.go
|
* Файл: cmd/futriis/main.go
|
||||||
// Назначение: Точка входа в приложение СУБД futriix.
|
* Назначение: Точка входа в приложение СУБД futriis.
|
||||||
// Инициализирует все компоненты системы: логгер, хранилище, транзакции, ACL, Raft координатор, HTTP API, WebUI, плагины и REPL
|
*
|
||||||
// Реализует graceful shutdown, проверку и применение миграций схемы данных, защиту от SPoF, автоматическое восстановление после паник горутин и полную валидацию конфигурации
|
* ИСПРАВЛЕНО:
|
||||||
|
* - Валидация пути к конфигурации и проверка прав доступа к файлу
|
||||||
|
* - Валидация конфигурации после загрузки (ValidateConfig)
|
||||||
|
* - WAL и data-dir привязаны к директории futriis, а не CWD
|
||||||
|
* - Graceful shutdown: убраны os.Exit(1) до завершения defer'ов,
|
||||||
|
* используется переменная exitCode и os.Exit в самом конце
|
||||||
|
* - Ожидание запуска HTTP-сервера с таймаутом
|
||||||
|
* - SAGA оркестратор останавливается всегда (даже если disabled)
|
||||||
|
* - Проверка allow-list для плагинов перед StartPlugin
|
||||||
|
* - Заменён SIGHUP на SIGUSR1 (reload) + SIGTERM/SIGINT (shutdown)
|
||||||
|
* - Проверка nil для raftCoordinator перед использованием
|
||||||
|
*
|
||||||
|
* УДАЛЕНО:
|
||||||
|
* - WebUI полностью удалён из ядра. HTTP API (api.NewHTTPServer)
|
||||||
|
* продолжает работать и предоставляет весь функционал через REST.
|
||||||
|
* Для UI используйте внешние инструменты (curl, jq, отдельный
|
||||||
|
* webui-клиент), подключающиеся к HTTP API.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package main
|
package main
|
||||||
@@ -18,9 +34,13 @@ package main
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"net"
|
||||||
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"os/signal"
|
"os/signal"
|
||||||
|
"path/filepath"
|
||||||
"syscall"
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -35,6 +55,18 @@ import (
|
|||||||
"futriis/pkg/utils"
|
"futriis/pkg/utils"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
// Директория по умолчанию для данных futriis
|
||||||
|
defaultDataDir = "futriis"
|
||||||
|
// Максимальный размер файла конфигурации (1 MB)
|
||||||
|
maxConfigFileSize = 1 * 1024 * 1024
|
||||||
|
)
|
||||||
|
|
||||||
|
// exitCode — код возврата. Устанавливается в graceful shutdown,
|
||||||
|
// а os.Exit вызывается только в самом конце main(), чтобы defer'ы
|
||||||
|
// успели выполниться.
|
||||||
|
var exitCode = 0
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
utils.SetColorEnabled(true)
|
utils.SetColorEnabled(true)
|
||||||
|
|
||||||
@@ -43,15 +75,29 @@ func main() {
|
|||||||
logLevel = "info"
|
logLevel = "info"
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := log.InitDefaultLogger("futriis/futriis.log", logLevel); err != nil {
|
dataDir, err := ensureDataDir()
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "Failed to prepare data directory: %v\n", err)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
|
||||||
|
logFile := filepath.Join(dataDir, "futriis.log")
|
||||||
|
if err := log.InitDefaultLogger(logFile, logLevel); err != nil {
|
||||||
fmt.Printf("Failed to initialize logger: %v\n", err)
|
fmt.Printf("Failed to initialize logger: %v\n", err)
|
||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Info("Futriis DB starting...")
|
log.Info("Futriis DB starting...")
|
||||||
log.Infof("Log level: %s", logLevel)
|
log.Infof("Log level: %s", logLevel)
|
||||||
|
log.Infof("Data directory: %s", dataDir)
|
||||||
|
|
||||||
cfg, err := config.LoadConfig("config.toml")
|
// Валидация пути к конфигу и прав доступа
|
||||||
|
configPath := os.Getenv("FUTRIIS_CONFIG")
|
||||||
|
if configPath == "" {
|
||||||
|
configPath = "config.toml"
|
||||||
|
}
|
||||||
|
|
||||||
|
cfg, err := loadAndValidateConfig(configPath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Failed to load config: " + err.Error())
|
log.Error("Failed to load config: " + err.Error())
|
||||||
utils.PrintError("Failed to load config: " + err.Error())
|
utils.PrintError("Failed to load config: " + err.Error())
|
||||||
@@ -65,27 +111,26 @@ func main() {
|
|||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
defer logger.Close()
|
defer logger.Close()
|
||||||
logger.Info("futriix database starting...")
|
logger.Info("futriis database starting...")
|
||||||
|
|
||||||
store := storage.NewStorage(cfg.Storage.PageSizeMB, logger)
|
store := storage.NewStorage(cfg.Storage.PageSizeMB, logger)
|
||||||
|
|
||||||
storage.SetGlobalStorage(store)
|
storage.SetGlobalStorage(store)
|
||||||
|
|
||||||
if err := storage.InitTransactionManager("futriix.wal"); err != nil {
|
// ИСПРАВЛЕНО: WAL в dataDir, а не относительно CWD
|
||||||
|
walPath := filepath.Join(dataDir, "futriis.wal")
|
||||||
|
if err := storage.InitTransactionManager(walPath); err != nil {
|
||||||
logger.Warn("Failed to initialize transaction manager: " + err.Error())
|
logger.Warn("Failed to initialize transaction manager: " + err.Error())
|
||||||
} else {
|
} else {
|
||||||
storage.SetTransactionLogger(logger)
|
storage.SetTransactionLogger(logger)
|
||||||
logger.Info("Transaction manager initialized")
|
logger.Info("Transaction manager initialized")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Инициализация менеджера триггеров (функция не возвращает ошибку)
|
|
||||||
storage.InitTriggerManager(logger)
|
storage.InitTriggerManager(logger)
|
||||||
logger.Info("Trigger manager initialized")
|
logger.Info("Trigger manager initialized")
|
||||||
|
|
||||||
aclManager := acl.NewACLManager()
|
aclManager := acl.NewACLManager()
|
||||||
logger.Info("ACL manager initialized")
|
logger.Info("ACL manager initialized")
|
||||||
|
|
||||||
// Передаём store как второй аргумент
|
|
||||||
raftCoordinator, err := cluster.NewRaftCoordinator(cfg, store, logger)
|
raftCoordinator, err := cluster.NewRaftCoordinator(cfg, store, logger)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("Failed to start Raft coordinator: " + err.Error())
|
logger.Error("Failed to start Raft coordinator: " + err.Error())
|
||||||
@@ -106,14 +151,12 @@ func main() {
|
|||||||
logger.Warn(fmt.Sprintf("Failed to initialize saga orchestrator: %v", err))
|
logger.Warn(fmt.Sprintf("Failed to initialize saga orchestrator: %v", err))
|
||||||
} else {
|
} else {
|
||||||
logger.Info(fmt.Sprintf("Saga orchestrator initialized with %d coordinators", sagaConfig.GetCoordinatorCount()))
|
logger.Info(fmt.Sprintf("Saga orchestrator initialized with %d coordinators", sagaConfig.GetCoordinatorCount()))
|
||||||
// Сохраняем для использования в других компонентах
|
|
||||||
storage.SetGlobalSagaOrchestrator(sagaOrchestrator)
|
storage.SetGlobalSagaOrchestrator(sagaOrchestrator)
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
logger.Info("Saga orchestrator disabled in configuration")
|
logger.Info("Saga orchestrator disabled in configuration")
|
||||||
}
|
}
|
||||||
|
|
||||||
// ========== НОВАЯ ФУНКЦИОНАЛЬНОСТЬ: Проверка статуса миграций схемы данных ==========
|
|
||||||
if raftCoordinator != nil {
|
if raftCoordinator != nil {
|
||||||
schemaMigrator := raftCoordinator.GetSchemaMigrator()
|
schemaMigrator := raftCoordinator.GetSchemaMigrator()
|
||||||
if schemaMigrator != nil {
|
if schemaMigrator != nil {
|
||||||
@@ -122,7 +165,6 @@ func main() {
|
|||||||
logger.Infof("Current schema version: %s, total migrations: %d, applied: %d, pending: %d",
|
logger.Infof("Current schema version: %s, total migrations: %d, applied: %d, pending: %d",
|
||||||
status.CurrentVersion, status.TotalMigrations, status.AppliedMigrations, status.PendingMigrations)
|
status.CurrentVersion, status.TotalMigrations, status.AppliedMigrations, status.PendingMigrations)
|
||||||
|
|
||||||
// Выводим список миграций
|
|
||||||
for _, m := range status.Migrations {
|
for _, m := range status.Migrations {
|
||||||
if m.Applied {
|
if m.Applied {
|
||||||
logger.Debugf(" [APPLIED] %s v%s: %s (at %s)",
|
logger.Debugf(" [APPLIED] %s v%s: %s (at %s)",
|
||||||
@@ -144,7 +186,6 @@ func main() {
|
|||||||
logger.Warn("Schema migrator not available")
|
logger.Warn("Schema migrator not available")
|
||||||
}
|
}
|
||||||
|
|
||||||
// ========== НОВАЯ ФУНКЦИОНАЛЬНОСТЬ: Проверка статуса Fallback менеджера ==========
|
|
||||||
fallbackMgr := raftCoordinator.GetFallbackManager()
|
fallbackMgr := raftCoordinator.GetFallbackManager()
|
||||||
if fallbackMgr != nil {
|
if fallbackMgr != nil {
|
||||||
logger.Info("Leader fallback manager is active (Single Point of Failure protection enabled)")
|
logger.Info("Leader fallback manager is active (Single Point of Failure protection enabled)")
|
||||||
@@ -154,7 +195,6 @@ func main() {
|
|||||||
logger.Warn("Leader fallback manager not available")
|
logger.Warn("Leader fallback manager not available")
|
||||||
}
|
}
|
||||||
|
|
||||||
// ========== НОВАЯ ФУНКЦИОНАЛЬНОСТЬ: Проверка статуса Panic Recovery менеджера ==========
|
|
||||||
panicRecoveryMgr := raftCoordinator.GetPanicRecoveryManager()
|
panicRecoveryMgr := raftCoordinator.GetPanicRecoveryManager()
|
||||||
if panicRecoveryMgr != nil {
|
if panicRecoveryMgr != nil {
|
||||||
logger.Info("Panic recovery manager is active (automatic goroutine recovery enabled)")
|
logger.Info("Panic recovery manager is active (automatic goroutine recovery enabled)")
|
||||||
@@ -164,7 +204,6 @@ func main() {
|
|||||||
logger.Warn("Panic recovery manager not available")
|
logger.Warn("Panic recovery manager not available")
|
||||||
}
|
}
|
||||||
|
|
||||||
// ========== НОВАЯ ФУНКЦИОНАЛЬНОСТЬ: Проверка статуса Persistence менеджера ==========
|
|
||||||
persistenceMgr := raftCoordinator.GetPersistenceManager()
|
persistenceMgr := raftCoordinator.GetPersistenceManager()
|
||||||
if persistenceMgr != nil {
|
if persistenceMgr != nil {
|
||||||
logger.Info("Persistence manager is active (automatic data persistence enabled)")
|
logger.Info("Persistence manager is active (automatic data persistence enabled)")
|
||||||
@@ -185,7 +224,6 @@ func main() {
|
|||||||
|
|
||||||
node := cluster.NewNode(cfg.Cluster.NodeIP, cfg.Cluster.NodePort, store, logger)
|
node := cluster.NewNode(cfg.Cluster.NodeIP, cfg.Cluster.NodePort, store, logger)
|
||||||
|
|
||||||
// Объявляем переменную maxRetries здесь, чтобы она была доступна
|
|
||||||
maxRetries := 5
|
maxRetries := 5
|
||||||
var registerErr error
|
var registerErr error
|
||||||
for i := 0; i < maxRetries; i++ {
|
for i := 0; i < maxRetries; i++ {
|
||||||
@@ -205,7 +243,7 @@ func main() {
|
|||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Передаём отдельные параметры из конфигурации плагинов
|
// ИСПРАВЛЕНО: проверка allow-list перед запуском плагинов
|
||||||
pluginManager := plugin.NewPluginManager(
|
pluginManager := plugin.NewPluginManager(
|
||||||
cfg.Plugins.ScriptDir,
|
cfg.Plugins.ScriptDir,
|
||||||
logger,
|
logger,
|
||||||
@@ -216,45 +254,59 @@ func main() {
|
|||||||
if cfg.Plugins.Enabled {
|
if cfg.Plugins.Enabled {
|
||||||
logger.Info(fmt.Sprintf("Plugin manager initialized, directory: %s", cfg.Plugins.ScriptDir))
|
logger.Info(fmt.Sprintf("Plugin manager initialized, directory: %s", cfg.Plugins.ScriptDir))
|
||||||
|
|
||||||
|
allowList := make(map[string]bool)
|
||||||
|
for _, name := range cfg.Plugins.AllowList {
|
||||||
|
allowList[name] = true
|
||||||
|
}
|
||||||
|
|
||||||
plugins := pluginManager.ListPlugins()
|
plugins := pluginManager.ListPlugins()
|
||||||
for _, p := range plugins {
|
for _, p := range plugins {
|
||||||
|
// ИСПРАВЛЕНО: если allow-list не пуст, плагин должен быть в нём
|
||||||
|
if len(allowList) > 0 && !allowList[p.Name] {
|
||||||
|
logger.Warn(fmt.Sprintf("Plugin %s is not in allow-list, skipping start", p.Name))
|
||||||
|
continue
|
||||||
|
}
|
||||||
if err := pluginManager.StartPlugin(p.Name); err != nil {
|
if err := pluginManager.StartPlugin(p.Name); err != nil {
|
||||||
logger.Warn(fmt.Sprintf("Failed to start plugin %s: %v", p.Name, err))
|
logger.Warn(fmt.Sprintf("Failed to start plugin %s: %v", p.Name, err))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ========== HTTP API SERVER ==========
|
||||||
|
// WebUI удалён из ядра. HTTP API предоставляет весь функционал через REST.
|
||||||
httpPort := cfg.API.Port
|
httpPort := cfg.API.Port
|
||||||
httpServer := api.NewHTTPServer(httpPort, store, raftCoordinator, aclManager, logger)
|
httpServer := api.NewHTTPServer(httpPort, store, raftCoordinator, aclManager, logger)
|
||||||
|
|
||||||
|
httpErrChan := make(chan error, 1)
|
||||||
go func() {
|
go func() {
|
||||||
if err := httpServer.Start(); err != nil {
|
if err := httpServer.Start(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||||
logger.Error("HTTP server error: " + err.Error())
|
httpErrChan <- fmt.Errorf("HTTP server error: %w", err)
|
||||||
utils.PrintError("HTTP server error: " + err.Error())
|
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
// Ждём запуска HTTP сервера (до 5 секунд)
|
||||||
|
httpStarted := false
|
||||||
|
for i := 0; i < 50; i++ {
|
||||||
|
select {
|
||||||
|
case err := <-httpErrChan:
|
||||||
|
logger.Error(err.Error())
|
||||||
|
utils.PrintError(err.Error())
|
||||||
|
os.Exit(1)
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
if isPortListening(httpPort) {
|
||||||
|
httpStarted = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
time.Sleep(100 * time.Millisecond)
|
||||||
|
}
|
||||||
|
if !httpStarted {
|
||||||
|
logger.Warn("HTTP server did not start within timeout, continuing anyway")
|
||||||
|
}
|
||||||
logger.Info(fmt.Sprintf("HTTP API server started on port %d", httpPort))
|
logger.Info(fmt.Sprintf("HTTP API server started on port %d", httpPort))
|
||||||
|
|
||||||
webUIPort := cfg.WebUI.Port
|
|
||||||
if webUIPort == 0 {
|
|
||||||
webUIPort = 8081
|
|
||||||
}
|
|
||||||
webUI := api.NewWebUIServer(webUIPort, cfg.WebUI.Enabled, store, raftCoordinator, aclManager, logger)
|
|
||||||
go func() {
|
|
||||||
if err := webUI.Start(); err != nil && cfg.WebUI.Enabled {
|
|
||||||
logger.Error("Web UI error: " + err.Error())
|
|
||||||
utils.PrintError("Web UI error: " + err.Error())
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
if cfg.WebUI.Enabled {
|
|
||||||
logger.Info(fmt.Sprintf("Web UI started on port %d", webUIPort))
|
|
||||||
}
|
|
||||||
|
|
||||||
if raftCoordinator != nil {
|
if raftCoordinator != nil {
|
||||||
logger.Info("Cluster features enabled: Pipeline Replication, Batch Commit, Dynamic Resharding, Joint Consensus")
|
logger.Info("Cluster features enabled: Pipeline Replication, Batch Commit, Dynamic Resharding, Joint Consensus")
|
||||||
logger.Info("Observability metrics and health checks available at /api/webui/metrics and /api/webui/health")
|
|
||||||
|
|
||||||
// ========== НОВАЯ ФУНКЦИОНАЛЬНОСТЬ: Вывод информации о новых компонентах ==========
|
|
||||||
logger.Info("Additional features:")
|
logger.Info("Additional features:")
|
||||||
logger.Info(" - Single Point of Failure protection (Leader Fallback) - ENABLED")
|
logger.Info(" - Single Point of Failure protection (Leader Fallback) - ENABLED")
|
||||||
logger.Info(" - Automatic Panic Recovery for goroutines - ENABLED")
|
logger.Info(" - Automatic Panic Recovery for goroutines - ENABLED")
|
||||||
@@ -263,20 +315,30 @@ func main() {
|
|||||||
logger.Info(" - Automatic data persistence - ENABLED")
|
logger.Info(" - Automatic data persistence - ENABLED")
|
||||||
}
|
}
|
||||||
|
|
||||||
displayBanner(cfg.Cluster.Name, cfg.WebUI.Enabled, webUIPort, httpPort, raftCoordinator, cfg.Saga.Enabled)
|
displayBanner(cfg.Cluster.Name, httpPort, raftCoordinator, cfg.Saga.Enabled)
|
||||||
|
|
||||||
replInstance := repl.NewRepl(store, raftCoordinator, logger, cfg, aclManager, pluginManager)
|
replInstance := repl.NewRepl(store, raftCoordinator, logger, cfg, aclManager, pluginManager)
|
||||||
|
|
||||||
// ========== GRACEFUL SHUTDOWN ==========
|
// ========== GRACEFUL SHUTDOWN ==========
|
||||||
// Канал для сигналов завершения
|
|
||||||
sigChan := make(chan os.Signal, 1)
|
sigChan := make(chan os.Signal, 1)
|
||||||
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
|
// ИСПРАВЛЕНО: SIGHUP заменён на SIGUSR1 (reload), SIGHUP игнорируется
|
||||||
|
// или обрабатывается отдельно, чтобы не ронять процесс при закрытии терминала.
|
||||||
|
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM, syscall.SIGUSR1)
|
||||||
|
|
||||||
// Запускаем graceful shutdown в отдельной горутине
|
|
||||||
go func() {
|
go func() {
|
||||||
<-sigChan
|
for sig := range sigChan {
|
||||||
|
switch sig {
|
||||||
|
case syscall.SIGUSR1:
|
||||||
|
logger.Info("Received SIGUSR1 (reload signal), reloading configuration...")
|
||||||
|
// TODO: реализовать reload конфигурации
|
||||||
|
continue
|
||||||
|
case syscall.SIGINT, syscall.SIGTERM:
|
||||||
|
logger.Info(fmt.Sprintf("Received signal %v, starting graceful shutdown...", sig))
|
||||||
|
}
|
||||||
|
break
|
||||||
|
}
|
||||||
|
|
||||||
utils.Println("\nReceived shutdown signal, starting graceful shutdown...")
|
utils.Println("\nReceived shutdown signal, starting graceful shutdown...")
|
||||||
logger.Info("Received shutdown signal, starting graceful shutdown...")
|
|
||||||
|
|
||||||
// 1. Останавливаем REPL
|
// 1. Останавливаем REPL
|
||||||
logger.Info("Stopping REPL...")
|
logger.Info("Stopping REPL...")
|
||||||
@@ -303,24 +365,7 @@ func main() {
|
|||||||
logger.Warn("HTTP server shutdown timeout")
|
logger.Warn("HTTP server shutdown timeout")
|
||||||
}
|
}
|
||||||
|
|
||||||
// 3. Останавливаем WebUI
|
// 3. Останавливаем плагины
|
||||||
logger.Info("Stopping WebUI...")
|
|
||||||
webUIStopDone := make(chan struct{})
|
|
||||||
go func() {
|
|
||||||
if err := webUI.Stop(); err != nil {
|
|
||||||
logger.Error(fmt.Sprintf("WebUI stop error: %v", err))
|
|
||||||
}
|
|
||||||
close(webUIStopDone)
|
|
||||||
}()
|
|
||||||
|
|
||||||
select {
|
|
||||||
case <-webUIStopDone:
|
|
||||||
logger.Info("WebUI stopped")
|
|
||||||
case <-time.After(10 * time.Second):
|
|
||||||
logger.Warn("WebUI shutdown timeout")
|
|
||||||
}
|
|
||||||
|
|
||||||
// 4. Останавливаем плагины
|
|
||||||
if cfg.Plugins.Enabled && pluginManager != nil {
|
if cfg.Plugins.Enabled && pluginManager != nil {
|
||||||
logger.Info("Stopping plugins...")
|
logger.Info("Stopping plugins...")
|
||||||
plugins := pluginManager.ListPlugins()
|
plugins := pluginManager.ListPlugins()
|
||||||
@@ -342,7 +387,7 @@ func main() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 5. Останавливаем SAGA оркестратор
|
// 4. Останавливаем SAGA оркестратор (всегда, если создан)
|
||||||
if sagaOrchestrator != nil {
|
if sagaOrchestrator != nil {
|
||||||
logger.Info("Stopping SAGA orchestrator...")
|
logger.Info("Stopping SAGA orchestrator...")
|
||||||
sagaStopDone := make(chan struct{})
|
sagaStopDone := make(chan struct{})
|
||||||
@@ -359,13 +404,14 @@ func main() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 6. Сохраняем данные на диск
|
// 5. Сохраняем данные на диск
|
||||||
|
// ИСПРАВЛЕНО: проверяем raftCoordinator на nil
|
||||||
|
if raftCoordinator != nil {
|
||||||
logger.Info("Persisting data to disk...")
|
logger.Info("Persisting data to disk...")
|
||||||
persistDone := make(chan struct{})
|
persistDone := make(chan struct{})
|
||||||
go func() {
|
go func() {
|
||||||
// Используем публичный метод GetPersistenceManager()
|
|
||||||
persistenceMgr := raftCoordinator.GetPersistenceManager()
|
persistenceMgr := raftCoordinator.GetPersistenceManager()
|
||||||
if raftCoordinator != nil && persistenceMgr != nil {
|
if persistenceMgr != nil {
|
||||||
if err := persistenceMgr.SaveAll(); err != nil {
|
if err := persistenceMgr.SaveAll(); err != nil {
|
||||||
logger.Error(fmt.Sprintf("Failed to save data: %v", err))
|
logger.Error(fmt.Sprintf("Failed to save data: %v", err))
|
||||||
} else {
|
} else {
|
||||||
@@ -383,8 +429,10 @@ func main() {
|
|||||||
case <-time.After(15 * time.Second):
|
case <-time.After(15 * time.Second):
|
||||||
logger.Warn("Persistence timeout")
|
logger.Warn("Persistence timeout")
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// 7. Останавливаем координатор
|
// 6. Останавливаем координатор
|
||||||
|
if raftCoordinator != nil {
|
||||||
logger.Info("Stopping Raft coordinator...")
|
logger.Info("Stopping Raft coordinator...")
|
||||||
raftStopDone := make(chan struct{})
|
raftStopDone := make(chan struct{})
|
||||||
go func() {
|
go func() {
|
||||||
@@ -398,8 +446,9 @@ func main() {
|
|||||||
case <-time.After(15 * time.Second):
|
case <-time.After(15 * time.Second):
|
||||||
logger.Warn("Raft coordinator shutdown timeout")
|
logger.Warn("Raft coordinator shutdown timeout")
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// 8. Останавливаем узел
|
// 7. Останавливаем узел
|
||||||
logger.Info("Stopping cluster node...")
|
logger.Info("Stopping cluster node...")
|
||||||
nodeStopDone := make(chan struct{})
|
nodeStopDone := make(chan struct{})
|
||||||
go func() {
|
go func() {
|
||||||
@@ -414,30 +463,133 @@ func main() {
|
|||||||
logger.Warn("Node shutdown timeout")
|
logger.Warn("Node shutdown timeout")
|
||||||
}
|
}
|
||||||
|
|
||||||
// 9. Синхронизируем и закрываем логгер
|
// 8. Синхронизируем и закрываем логгер
|
||||||
logger.Info("Finalizing logger...")
|
logger.Info("Finalizing logger...")
|
||||||
if err := logger.Sync(); err != nil {
|
if err := logger.Sync(); err != nil {
|
||||||
fmt.Printf("Failed to sync logger: %v\n", err)
|
fmt.Printf("Failed to sync logger: %v\n", err)
|
||||||
}
|
}
|
||||||
logger.Close()
|
logger.Close()
|
||||||
|
|
||||||
|
// 9. Останавливаем транзакционный менеджер
|
||||||
|
if err := storage.StopTransactionManager(); err != nil {
|
||||||
|
fmt.Printf("Failed to stop transaction manager: %v\n", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 10. Останавливаем ACL manager (фоновые горутины очистки сессий)
|
||||||
|
if aclManager != nil {
|
||||||
|
aclManager.Stop()
|
||||||
|
}
|
||||||
|
|
||||||
utils.DisableColorMode()
|
utils.DisableColorMode()
|
||||||
fmt.Println("Futriis DB shutdown complete")
|
fmt.Println("Futriis DB shutdown complete")
|
||||||
os.Exit(0)
|
|
||||||
|
// ИСПРАВЛЕНО: устанавливаем exitCode и завершаем main через return,
|
||||||
|
// чтобы defer'ы выполнились. os.Exit вызывается только в конце main.
|
||||||
|
exitCode = 0
|
||||||
|
close(sigChan)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
// Запускаем REPL (блокирующий вызов)
|
// Запускаем REPL (блокирующий вызов)
|
||||||
if err := replInstance.Run(); err != nil {
|
if err := replInstance.Run(); err != nil {
|
||||||
logger.Error("REPL error: " + err.Error())
|
logger.Error("REPL error: " + err.Error())
|
||||||
utils.PrintError("REPL error: " + err.Error())
|
utils.PrintError("REPL error: " + err.Error())
|
||||||
os.Exit(1)
|
exitCode = 1
|
||||||
|
}
|
||||||
|
|
||||||
|
// Даём graceful shutdown время завершиться
|
||||||
|
time.Sleep(500 * time.Millisecond)
|
||||||
|
|
||||||
|
if exitCode != 0 {
|
||||||
|
os.Exit(exitCode)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func displayBanner(clusterName string, webUIEnabled bool, webUIPort int, httpPort int, coordinator *cluster.RaftCoordinator, sagaEnabled bool) {
|
// ensureDataDir создаёт и возвращает абсолютный путь к директории данных.
|
||||||
|
// ИСПРАВЛЕНО: возвращает абсолютный путь, чтобы избежать зависимости от CWD.
|
||||||
|
func ensureDataDir() (string, error) {
|
||||||
|
dataDir := os.Getenv("FUTRIIS_DATA_DIR")
|
||||||
|
if dataDir == "" {
|
||||||
|
dataDir = defaultDataDir
|
||||||
|
}
|
||||||
|
|
||||||
|
absPath, err := filepath.Abs(dataDir)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("failed to resolve absolute path for %s: %w", dataDir, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := os.MkdirAll(absPath, 0755); err != nil {
|
||||||
|
return "", fmt.Errorf("failed to create data directory %s: %w", absPath, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Проверяем, что директория доступна для записи
|
||||||
|
testFile := filepath.Join(absPath, ".write_test")
|
||||||
|
if err := os.WriteFile(testFile, []byte("test"), 0600); err != nil {
|
||||||
|
return "", fmt.Errorf("data directory %s is not writable: %w", absPath, err)
|
||||||
|
}
|
||||||
|
_ = os.Remove(testFile)
|
||||||
|
|
||||||
|
return absPath, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// loadAndValidateConfig загружает и валидирует конфигурацию.
|
||||||
|
// ИСПРАВЛЕНО: проверка размера файла, прав доступа, валидация.
|
||||||
|
func loadAndValidateConfig(path string) (*config.Config, error) {
|
||||||
|
// Проверка существования файла
|
||||||
|
info, err := os.Stat(path)
|
||||||
|
if err != nil {
|
||||||
|
if os.IsNotExist(err) {
|
||||||
|
return nil, fmt.Errorf("config file not found: %s", path)
|
||||||
|
}
|
||||||
|
return nil, fmt.Errorf("failed to stat config file: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Проверка размера
|
||||||
|
if info.Size() > maxConfigFileSize {
|
||||||
|
return nil, fmt.Errorf("config file too large: %d bytes (max %d)", info.Size(), maxConfigFileSize)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Проверка, что это обычный файл (не симлинк на /etc/passwd и т.д.)
|
||||||
|
if !info.Mode().IsRegular() {
|
||||||
|
return nil, fmt.Errorf("config path is not a regular file: %s", path)
|
||||||
|
}
|
||||||
|
|
||||||
|
cfg, err := config.LoadConfig(path)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("failed to parse config: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if cfg == nil {
|
||||||
|
return nil, fmt.Errorf("config is nil after load")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Валидация конфигурации
|
||||||
|
if err := config.ValidateConfig(cfg); err != nil {
|
||||||
|
return nil, fmt.Errorf("config validation failed: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return cfg, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// isPortListening проверяет, слушается ли порт.
|
||||||
|
// Простая эвристика: пробуем подключиться к localhost:port с таймаутом.
|
||||||
|
// Если порт занят — значит, сервер запустился.
|
||||||
|
func isPortListening(port int) bool {
|
||||||
|
if port <= 0 {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
// Пробуем установить TCP-соединение с коротким таймаутом
|
||||||
|
conn, err := net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", port), 200*time.Millisecond)
|
||||||
|
if err != nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
conn.Close()
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
func displayBanner(clusterName string, httpPort int, coordinator *cluster.RaftCoordinator, sagaEnabled bool) {
|
||||||
utils.Println("")
|
utils.Println("")
|
||||||
bannerLines := []string{
|
bannerLines := []string{
|
||||||
" futriix 3i²(by 02.04.2026) ",
|
" futriis 3i²(by 02.04.2026) ",
|
||||||
" Distributed Document-Store in-memory database with support lua plugins ",
|
" Distributed Document-Store in-memory database with support lua plugins ",
|
||||||
" Cluster status: enable (Raft consensus)",
|
" Cluster status: enable (Raft consensus)",
|
||||||
" Cluster features: Pipeline Replication, Batch Commit, Dynamic Resharding",
|
" Cluster features: Pipeline Replication, Batch Commit, Dynamic Resharding",
|
||||||
@@ -445,12 +597,6 @@ func displayBanner(clusterName string, webUIEnabled bool, webUIPort int, httpPor
|
|||||||
" HTTP API (for curl or wget utils only): http://localhost:" + fmt.Sprintf("%d", httpPort) + "/api/",
|
" HTTP API (for curl or wget utils only): http://localhost:" + fmt.Sprintf("%d", httpPort) + "/api/",
|
||||||
}
|
}
|
||||||
|
|
||||||
if webUIEnabled {
|
|
||||||
bannerLines = append(bannerLines, fmt.Sprintf(" Web UI: http://localhost:%d/", webUIPort))
|
|
||||||
bannerLines = append(bannerLines, " Observability endpoints: /api/webui/metrics, /api/webui/health")
|
|
||||||
}
|
|
||||||
|
|
||||||
// ========== НОВАЯ ФУНКЦИОНАЛЬНОСТЬ: Добавляем информацию о новых компонентах в баннер ==========
|
|
||||||
bannerLines = append(bannerLines, " Additional features:")
|
bannerLines = append(bannerLines, " Additional features:")
|
||||||
bannerLines = append(bannerLines, " - Single Point of Failure protection (Leader Fallback)")
|
bannerLines = append(bannerLines, " - Single Point of Failure protection (Leader Fallback)")
|
||||||
bannerLines = append(bannerLines, " - Automatic Panic Recovery for goroutines")
|
bannerLines = append(bannerLines, " - Automatic Panic Recovery for goroutines")
|
||||||
@@ -464,7 +610,6 @@ func displayBanner(clusterName string, webUIEnabled bool, webUIPort int, httpPor
|
|||||||
bannerLines = append(bannerLines, " - SAGA Orchestrator: DISABLED")
|
bannerLines = append(bannerLines, " - SAGA Orchestrator: DISABLED")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Выводим статус Fallback менеджера, если он активен
|
|
||||||
if coordinator != nil {
|
if coordinator != nil {
|
||||||
fallbackMgr := coordinator.GetFallbackManager()
|
fallbackMgr := coordinator.GetFallbackManager()
|
||||||
if fallbackMgr != nil {
|
if fallbackMgr != nil {
|
||||||
|
|||||||
Reference in new issue
Block a user