From 79c1a0171e84208cd2db042fd46c3c1385e22046 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=93=D1=80=D0=B8=D0=B3=D0=BE=D1=80=D0=B8=D0=B9=20=D0=A1?= =?UTF-8?q?=D0=B0=D1=84=D1=80=D0=BE=D0=BD=D0=BE=D0=B2?= Date: Mon, 5 Oct 2026 21:27:17 +0000 Subject: [PATCH] Update internal/config/dynamic_config.go --- internal/config/dynamic_config.go | 147 +++++++++++++++++------------- 1 file changed, 85 insertions(+), 62 deletions(-) diff --git a/internal/config/dynamic_config.go b/internal/config/dynamic_config.go index 6ada71e..45991f4 100644 --- a/internal/config/dynamic_config.go +++ b/internal/config/dynamic_config.go @@ -9,7 +9,7 @@ */ // Файл: internal/config/dynamic_config.go -// Назначение: Централизованное управление конфигурацией через Raft +// Назначение: Централизованное управление конфигурацией через Raft. package config @@ -32,8 +32,8 @@ import ( type atomicBool struct{ v int32 } -func (b *atomicBool) Load() bool { return atomic.LoadInt32(&b.v) != 0 } -func (b *atomicBool) Store(v bool) { atomic.StoreInt32(&b.v, boolToInt32(v)) } +func (b *atomicBool) Load() bool { return atomic.LoadInt32(&b.v) != 0 } +func (b *atomicBool) Store(v bool) { atomic.StoreInt32(&b.v, boolToInt32(v)) } func (b *atomicBool) CAS(old, new bool) bool { return atomic.CompareAndSwapInt32(&b.v, boolToInt32(old), boolToInt32(new)) } @@ -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,22 +65,22 @@ type DynamicConfigChange struct { Changes map[string]interface{} `json:"changes"` } -// DynamicConfigSnapshot представляет снэпшот конфигурации +// DynamicConfigSnapshot представляет снэпшот конфигурации. type DynamicConfigSnapshot struct { - Version uint64 `json:"version"` - UpdatedAt int64 `json:"updated_at"` - Config *Config `json:"config"` - ChangeHistory []*DynamicConfigChange `json:"change_history"` + Version uint64 `json:"version"` + UpdatedAt int64 `json:"updated_at"` + Config *Config `json:"config"` + 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,29 +97,28 @@ type DynamicConfigManager struct { isLeader atomicBool } -// ConfigFSM реализует конечный автомат для конфигурации +// ConfigFSM реализует конечный автомат для конфигурации. type ConfigFSM struct { config *DynamicConfigSnapshot mu sync.RWMutex logger *log.Logger // onChange вызывается после успешного применения изменения. // Может быть установлен менеджером для уведомления слушателей. - onChange func(change *DynamicConfigChange) + onChange func(change *DynamicConfigChange) 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()