Update internal/config/dynamic_config.go
This commit is contained in:
1 parent
ba4f7bfaa7
commit
79c1a0171e
1 file changed
+78
-55
@@ -9,7 +9,7 @@
|
||||
*/
|
||||
|
||||
// Файл: internal/config/dynamic_config.go
|
||||
// Назначение: Централизованное управление конфигурацией через Raft
|
||||
// Назначение: Централизованное управление конфигурацией через Raft.
|
||||
|
||||
package config
|
||||
|
||||
@@ -55,7 +55,7 @@ func (a *atomicUint64) Add(d uint64) uint64 { return atomic.AddUint64(&a.v, d) }
|
||||
// ТИПЫ ДЛЯ ДИНАМИЧЕСКОЙ КОНФИГУРАЦИИ
|
||||
// =============================================================================
|
||||
|
||||
// DynamicConfigChange представляет изменение конфигурации
|
||||
// DynamicConfigChange представляет изменение конфигурации.
|
||||
type DynamicConfigChange struct {
|
||||
ID string `json:"id"`
|
||||
Version uint64 `json:"version"`
|
||||
@@ -65,7 +65,7 @@ type DynamicConfigChange struct {
|
||||
Changes map[string]interface{} `json:"changes"`
|
||||
}
|
||||
|
||||
// DynamicConfigSnapshot представляет снэпшот конфигурации
|
||||
// DynamicConfigSnapshot представляет снэпшот конфигурации.
|
||||
type DynamicConfigSnapshot struct {
|
||||
Version uint64 `json:"version"`
|
||||
UpdatedAt int64 `json:"updated_at"`
|
||||
@@ -73,14 +73,14 @@ type DynamicConfigSnapshot struct {
|
||||
ChangeHistory []*DynamicConfigChange `json:"change_history"`
|
||||
}
|
||||
|
||||
// ConfigChangeListener представляет слушатель изменений конфигурации
|
||||
// ConfigChangeListener представляет слушатель изменений конфигурации.
|
||||
type ConfigChangeListener func(change *DynamicConfigChange)
|
||||
|
||||
// =============================================================================
|
||||
// DYNAMIC CONFIG MANAGER
|
||||
// =============================================================================
|
||||
|
||||
// DynamicConfigManager управляет динамической конфигурацией через Raft
|
||||
// DynamicConfigManager управляет динамической конфигурацией через Raft.
|
||||
type DynamicConfigManager struct {
|
||||
config *Config
|
||||
logger *log.Logger
|
||||
@@ -97,7 +97,7 @@ type DynamicConfigManager struct {
|
||||
isLeader atomicBool
|
||||
}
|
||||
|
||||
// ConfigFSM реализует конечный автомат для конфигурации
|
||||
// ConfigFSM реализует конечный автомат для конфигурации.
|
||||
type ConfigFSM struct {
|
||||
config *DynamicConfigSnapshot
|
||||
mu sync.RWMutex
|
||||
@@ -108,18 +108,17 @@ type ConfigFSM struct {
|
||||
onChangeMu sync.RWMutex
|
||||
}
|
||||
|
||||
// ConfigSnapshot реализует raft.FSMSnapshot
|
||||
// ConfigSnapshot реализует raft.FSMSnapshot.
|
||||
type ConfigSnapshot struct {
|
||||
config *DynamicConfigSnapshot
|
||||
}
|
||||
|
||||
// NewDynamicConfigManager создаёт новый менеджер динамической конфигурации
|
||||
// NewDynamicConfigManager создаёт новый менеджер динамической конфигурации.
|
||||
func NewDynamicConfigManager(cfg *Config, logger *log.Logger) *DynamicConfigManager {
|
||||
if cfg == nil {
|
||||
cfg = &Config{}
|
||||
}
|
||||
|
||||
// Создаём начальный снэпшот
|
||||
initialSnapshot := &DynamicConfigSnapshot{
|
||||
Version: 1,
|
||||
UpdatedAt: time.Now().UnixMilli(),
|
||||
@@ -142,13 +141,11 @@ func NewDynamicConfigManager(cfg *Config, logger *log.Logger) *DynamicConfigMana
|
||||
}
|
||||
dm.version.Store(1)
|
||||
|
||||
// Подключаем FSM к менеджеру для уведомления слушателей
|
||||
fsm.onChangeMu.Lock()
|
||||
fsm.onChange = func(change *DynamicConfigChange) {
|
||||
dm.version.Store(change.Version)
|
||||
dm.mu.Lock()
|
||||
dm.changeHistory = append(dm.changeHistory, change)
|
||||
// Ограничиваем историю
|
||||
if len(dm.changeHistory) > 1000 {
|
||||
dm.changeHistory = dm.changeHistory[len(dm.changeHistory)-1000:]
|
||||
}
|
||||
@@ -160,19 +157,19 @@ func NewDynamicConfigManager(cfg *Config, logger *log.Logger) *DynamicConfigMana
|
||||
return dm
|
||||
}
|
||||
|
||||
// SetRaft устанавливает Raft для управления конфигурацией
|
||||
// SetRaft устанавливает Raft для управления конфигурацией.
|
||||
func (dm *DynamicConfigManager) SetRaft(r *raft.Raft) {
|
||||
dm.mu.Lock()
|
||||
defer dm.mu.Unlock()
|
||||
dm.raft = r
|
||||
}
|
||||
|
||||
// GetFSM возвращает FSM для конфигурации
|
||||
// GetFSM возвращает FSM для конфигурации.
|
||||
func (dm *DynamicConfigManager) GetFSM() *ConfigFSM {
|
||||
return dm.fsm
|
||||
}
|
||||
|
||||
// Start запускает менеджер динамической конфигурации
|
||||
// Start запускает менеджер динамической конфигурации.
|
||||
func (dm *DynamicConfigManager) Start() {
|
||||
if dm.logger != nil {
|
||||
dm.logger.Info("Dynamic config manager started")
|
||||
@@ -203,9 +200,8 @@ func (dm *DynamicConfigManager) GetConfig() *Config {
|
||||
return deepCopyConfig(dm.fsm.config.Config)
|
||||
}
|
||||
|
||||
// GetVersion возвращает текущую версию конфигурации
|
||||
// GetVersion возвращает текущую версию конфигурации.
|
||||
func (dm *DynamicConfigManager) GetVersion() uint64 {
|
||||
// Синхронизируем с FSM, чтобы значение было актуальным
|
||||
dm.fsm.mu.RLock()
|
||||
defer dm.fsm.mu.RUnlock()
|
||||
if dm.fsm.config != nil {
|
||||
@@ -214,7 +210,7 @@ func (dm *DynamicConfigManager) GetVersion() uint64 {
|
||||
return dm.version.Load()
|
||||
}
|
||||
|
||||
// GetChangeHistory возвращает историю изменений
|
||||
// GetChangeHistory возвращает историю изменений.
|
||||
func (dm *DynamicConfigManager) GetChangeHistory() []*DynamicConfigChange {
|
||||
dm.mu.RLock()
|
||||
defer dm.mu.RUnlock()
|
||||
@@ -224,7 +220,7 @@ func (dm *DynamicConfigManager) GetChangeHistory() []*DynamicConfigChange {
|
||||
return history
|
||||
}
|
||||
|
||||
// ApplyChange применяет изменение конфигурации через Raft
|
||||
// ApplyChange применяет изменение конфигурации через Raft.
|
||||
func (dm *DynamicConfigManager) ApplyChange(changes map[string]interface{}, description string, changedBy string) error {
|
||||
dm.mu.RLock()
|
||||
r := dm.raft
|
||||
@@ -234,14 +230,12 @@ func (dm *DynamicConfigManager) ApplyChange(changes map[string]interface{}, desc
|
||||
return fmt.Errorf("raft not initialized")
|
||||
}
|
||||
|
||||
// Проверяем, что мы лидер — через сам Raft, а не через внешний флаг
|
||||
if r.State() != raft.Leader {
|
||||
return fmt.Errorf("not the leader (current state: %s)", r.State())
|
||||
}
|
||||
|
||||
currentVersion := dm.GetVersion()
|
||||
|
||||
// Создаём команду изменения
|
||||
change := &DynamicConfigChange{
|
||||
ID: fmt.Sprintf("config_%d_%d", currentVersion+1, time.Now().UnixNano()),
|
||||
Version: currentVersion + 1,
|
||||
@@ -251,7 +245,6 @@ func (dm *DynamicConfigManager) ApplyChange(changes map[string]interface{}, desc
|
||||
Changes: changes,
|
||||
}
|
||||
|
||||
// Сериализуем и применяем через Raft
|
||||
data, err := json.Marshal(change)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to marshal config change: %v", err)
|
||||
@@ -262,7 +255,6 @@ func (dm *DynamicConfigManager) ApplyChange(changes map[string]interface{}, desc
|
||||
return fmt.Errorf("raft apply failed: %v", err)
|
||||
}
|
||||
|
||||
// Проверяем ошибку из FSM.Apply
|
||||
if resp := future.Response(); resp != nil {
|
||||
if err, ok := resp.(error); ok {
|
||||
return fmt.Errorf("fsm apply failed: %v", err)
|
||||
@@ -272,16 +264,15 @@ func (dm *DynamicConfigManager) ApplyChange(changes map[string]interface{}, desc
|
||||
if dm.logger != nil {
|
||||
dm.logger.Info(fmt.Sprintf("Applied config change: %s (version: %d)", description, change.Version))
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// UpdateConfig обновляет конфигурацию
|
||||
// UpdateConfig обновляет конфигурацию.
|
||||
func (dm *DynamicConfigManager) UpdateConfig(updates map[string]interface{}, description string, changedBy string) error {
|
||||
return dm.ApplyChange(updates, description, changedBy)
|
||||
}
|
||||
|
||||
// ApplyConfigChange применяет изменение к FSM (вызывается из Raft)
|
||||
// ApplyConfigChange применяет изменение к FSM (вызывается из Raft).
|
||||
func (fsm *ConfigFSM) ApplyConfigChange(change *DynamicConfigChange) error {
|
||||
if change == nil {
|
||||
return fmt.Errorf("nil config change")
|
||||
@@ -303,7 +294,6 @@ func (fsm *ConfigFSM) ApplyConfigChange(change *DynamicConfigChange) error {
|
||||
|
||||
config := fsm.config.Config
|
||||
|
||||
// Применяем изменения к конфигурации
|
||||
for path, value := range change.Changes {
|
||||
if err := setConfigField(config, path, value); err != nil {
|
||||
fsm.mu.Unlock()
|
||||
@@ -311,11 +301,9 @@ func (fsm *ConfigFSM) ApplyConfigChange(change *DynamicConfigChange) error {
|
||||
}
|
||||
}
|
||||
|
||||
// Обновляем метаданные
|
||||
fsm.config.Version = change.Version
|
||||
fsm.config.UpdatedAt = time.Now().UnixMilli()
|
||||
fsm.config.ChangeHistory = append(fsm.config.ChangeHistory, change)
|
||||
// Ограничиваем историю
|
||||
if len(fsm.config.ChangeHistory) > 1000 {
|
||||
fsm.config.ChangeHistory = fsm.config.ChangeHistory[len(fsm.config.ChangeHistory)-1000:]
|
||||
}
|
||||
@@ -326,32 +314,28 @@ func (fsm *ConfigFSM) ApplyConfigChange(change *DynamicConfigChange) error {
|
||||
fsm.logger.Debug(fmt.Sprintf("Applied config change to FSM: %s", change.Description))
|
||||
}
|
||||
|
||||
// Уведомляем слушателей (через callback, установленный менеджером)
|
||||
fsm.onChangeMu.RLock()
|
||||
cb := fsm.onChange
|
||||
fsm.onChangeMu.RUnlock()
|
||||
if cb != nil {
|
||||
cb(change)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Apply применяет команду к FSM (реализация raft.FSM)
|
||||
// Apply применяет команду к FSM (реализация raft.FSM).
|
||||
func (fsm *ConfigFSM) Apply(log *raft.Log) interface{} {
|
||||
var change DynamicConfigChange
|
||||
if err := json.Unmarshal(log.Data, &change); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal config change: %v", err)
|
||||
}
|
||||
|
||||
if err := fsm.ApplyConfigChange(&change); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Snapshot создаёт снэпшот состояния FSM
|
||||
// Snapshot создаёт снэпшот состояния FSM.
|
||||
func (fsm *ConfigFSM) Snapshot() (raft.FSMSnapshot, error) {
|
||||
fsm.mu.RLock()
|
||||
defer fsm.mu.RUnlock()
|
||||
@@ -365,7 +349,6 @@ func (fsm *ConfigFSM) Snapshot() (raft.FSMSnapshot, error) {
|
||||
}}, nil
|
||||
}
|
||||
|
||||
// Глубокая копия через JSON, чтобы снэпшот не зависел от мутаций FSM
|
||||
data, err := json.Marshal(fsm.config)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to marshal snapshot: %v", err)
|
||||
@@ -384,7 +367,7 @@ func (fsm *ConfigFSM) Snapshot() (raft.FSMSnapshot, error) {
|
||||
return &ConfigSnapshot{config: &snapshotData}, nil
|
||||
}
|
||||
|
||||
// Restore восстанавливает состояние FSM из снэпшота
|
||||
// Restore восстанавливает состояние FSM из снэпшота.
|
||||
func (fsm *ConfigFSM) Restore(snapshot io.ReadCloser) error {
|
||||
defer snapshot.Close()
|
||||
|
||||
@@ -407,37 +390,34 @@ func (fsm *ConfigFSM) Restore(snapshot io.ReadCloser) error {
|
||||
fsm.mu.Lock()
|
||||
fsm.config = &data
|
||||
fsm.mu.Unlock()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Persist сохраняет снэпшот
|
||||
// Persist сохраняет снэпшот.
|
||||
func (s *ConfigSnapshot) Persist(sink raft.SnapshotSink) error {
|
||||
data, err := json.Marshal(s.config)
|
||||
if err != nil {
|
||||
sink.Cancel()
|
||||
return err
|
||||
}
|
||||
|
||||
if _, err := sink.Write(data); err != nil {
|
||||
sink.Cancel()
|
||||
return err
|
||||
}
|
||||
|
||||
return sink.Close()
|
||||
}
|
||||
|
||||
// Release освобождает ресурсы снэпшота
|
||||
// Release освобождает ресурсы снэпшота.
|
||||
func (s *ConfigSnapshot) Release() {}
|
||||
|
||||
// RegisterListener регистрирует слушатель изменений конфигурации
|
||||
// RegisterListener регистрирует слушатель изменений конфигурации.
|
||||
func (dm *DynamicConfigManager) RegisterListener(name string, listener ConfigChangeListener) {
|
||||
dm.listenerMu.Lock()
|
||||
defer dm.listenerMu.Unlock()
|
||||
dm.listeners[name] = listener
|
||||
}
|
||||
|
||||
// UnregisterListener удаляет слушатель
|
||||
// UnregisterListener удаляет слушатель.
|
||||
func (dm *DynamicConfigManager) UnregisterListener(name string) {
|
||||
dm.listenerMu.Lock()
|
||||
defer dm.listenerMu.Unlock()
|
||||
@@ -445,8 +425,6 @@ func (dm *DynamicConfigManager) UnregisterListener(name string) {
|
||||
}
|
||||
|
||||
// notifyListeners уведомляет слушателей об изменении.
|
||||
// Каждый слушатель вызывается в отдельной горутине; паника в слушателе
|
||||
// не должна ронять процесс.
|
||||
func (dm *DynamicConfigManager) notifyListeners(change *DynamicConfigChange) {
|
||||
dm.listenerMu.RLock()
|
||||
defer dm.listenerMu.RUnlock()
|
||||
@@ -468,28 +446,22 @@ func (dm *DynamicConfigManager) notifyListeners(change *DynamicConfigChange) {
|
||||
// =============================================================================
|
||||
|
||||
// deepCopyConfig создаёт глубокую копию конфигурации через JSON.
|
||||
// Если src == nil, возвращает &Config{} (не nil), чтобы избежать паник.
|
||||
func deepCopyConfig(src *Config) *Config {
|
||||
if src == nil {
|
||||
return &Config{}
|
||||
}
|
||||
|
||||
data, err := json.Marshal(src)
|
||||
if err != nil {
|
||||
// Возвращаем пустую копию, а не nil, чтобы не сломать вызывающий код
|
||||
return &Config{}
|
||||
}
|
||||
|
||||
var dst Config
|
||||
if err := json.Unmarshal(data, &dst); err != nil {
|
||||
return &Config{}
|
||||
}
|
||||
|
||||
return &dst
|
||||
}
|
||||
|
||||
// coerceInt пытается привести значение к int.
|
||||
// Поддерживает int, int32, int64, uint, uint32, uint64, float32, float64, json.Number.
|
||||
func coerceInt(value interface{}) (int, bool) {
|
||||
switch v := value.(type) {
|
||||
case int:
|
||||
@@ -627,8 +599,7 @@ func coerceString(value interface{}) (string, bool) {
|
||||
}
|
||||
|
||||
// setConfigField устанавливает значение поля конфигурации по пути.
|
||||
// Поддерживает основные секции: cluster, storage, replication, wal, mvcc,
|
||||
// saga, backpressure, api, log, plugins, autoscaling, monitoring, transactions.
|
||||
// Используется FSM при применении изменений из Raft.
|
||||
func setConfigField(config *Config, path string, value interface{}) error {
|
||||
if path == "" {
|
||||
return fmt.Errorf("empty config path")
|
||||
@@ -1087,12 +1058,64 @@ func setConfigField(config *Config, path string, value interface{}) error {
|
||||
}
|
||||
return fmt.Errorf("transactions.max_savepoints_per_tx: expected int, got %T", value)
|
||||
|
||||
// ===== compression =====
|
||||
case "compression.enabled":
|
||||
if v, ok := coerceBool(value); ok {
|
||||
config.Compression.Enabled = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("compression.enabled: expected bool, got %T", value)
|
||||
case "compression.algorithm":
|
||||
if v, ok := coerceString(value); ok {
|
||||
config.Compression.Algorithm = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("compression.algorithm: expected string, got %T", value)
|
||||
case "compression.level":
|
||||
if v, ok := coerceInt(value); ok {
|
||||
config.Compression.Level = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("compression.level: expected int, got %T", value)
|
||||
case "compression.min_size":
|
||||
if v, ok := coerceInt(value); ok {
|
||||
config.Compression.MinSize = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("compression.min_size: expected int, got %T", value)
|
||||
|
||||
// ===== performance =====
|
||||
case "performance.enable_pipeline":
|
||||
if v, ok := coerceBool(value); ok {
|
||||
config.Performance.EnablePipeline = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("performance.enable_pipeline: expected bool, got %T", value)
|
||||
case "performance.batch_size":
|
||||
if v, ok := coerceInt(value); ok {
|
||||
config.Performance.BatchSize = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("performance.batch_size: expected int, got %T", value)
|
||||
case "performance.max_connections":
|
||||
if v, ok := coerceInt(value); ok {
|
||||
config.Performance.MaxConnections = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("performance.max_connections: expected int, got %T", value)
|
||||
case "performance.read_from_follower":
|
||||
if v, ok := coerceBool(value); ok {
|
||||
config.Performance.ReadFromFollower = v
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("performance.read_from_follower: expected bool, got %T", value)
|
||||
|
||||
default:
|
||||
return fmt.Errorf("unknown config path: %s", path)
|
||||
}
|
||||
}
|
||||
|
||||
// GetLeaderStatus возвращает статус лидера
|
||||
// GetLeaderStatus возвращает статус лидера.
|
||||
func (dm *DynamicConfigManager) GetLeaderStatus() bool {
|
||||
dm.mu.RLock()
|
||||
r := dm.raft
|
||||
@@ -1110,7 +1133,7 @@ func (dm *DynamicConfigManager) SetLeaderStatus(isLeader bool) {
|
||||
dm.isLeader.Store(isLeader)
|
||||
}
|
||||
|
||||
// GetListeners возвращает список зарегистрированных слушателей
|
||||
// GetListeners возвращает список зарегистрированных слушателей.
|
||||
func (dm *DynamicConfigManager) GetListeners() []string {
|
||||
dm.listenerMu.RLock()
|
||||
defer dm.listenerMu.RUnlock()
|
||||
|
||||
Reference in new issue
Block a user