From 3f1748bd6c9cd3b42ef396d32e810cfb33ffc0e3 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:58 +0000 Subject: [PATCH] Update internal/repl/repl.go --- internal/repl/repl.go | 1691 ++++++++++++++++++++++++++++++++--------- 1 file changed, 1320 insertions(+), 371 deletions(-) diff --git a/internal/repl/repl.go b/internal/repl/repl.go index 15e7444..105f86f 100644 --- a/internal/repl/repl.go +++ b/internal/repl/repl.go @@ -12,44 +12,27 @@ // Пакет: repl (Read-Eval-Print Loop) // ============================================================================= // Назначение: Интерактивный интерфейс командной строки (REPL) для СУБД futriix. -// Поддерживает: -// - автодополнение (Tab) команд, коллекций, полей, операторов; -// - историю ввода с навигацией PgUp/PgDown и стрелками вверх/вниз; -// - многострочный ввод (обратный слэш в конце строки, незакрытые скобки); -// - пейджер для длинных результатов; -// - табличный вывод по умолчанию. // Совместимо с Linux и OpenIndiana. // // ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции -// функции storage.BeginTransaction, storage.CommitCurrentTransaction, -// storage.AbortCurrentTransaction, storage.HasActiveTransaction, -// storage.GetActiveTransactions были удалены. Команды REPL для транзакций -// переписаны поверх batch-API: -// - "begin transaction" → создаёт *storage.Batch (алиас "batch begin"); -// - "commit" → batch.Commit() (алиас "batch commit"); -// - "rollback" → снимает регистрацию batch (алиас "batch abort"); -// - "show transactions" → показывает активный batch сессии. -// Старые имена команд сохранены как алиасы для совместимости с привычками -// пользователей и существующими скриптами. +// команды REPL для транзакций переписаны поверх batch-API. // -// ИЗМЕНЕНО (2026-10-03): формат временных меток изменён с "2006-01-02 15:04:05.000" -// на "02-01-2006 15:04:05.000" для совместимости с OpenIndiana. +// ИЗМЕНЕНО (2026-10-03): формат временных меток изменён на +// "02-01-2006 15:04:05.000". // -// ИЗМЕНЕНО (2026-10-04): формат временных меток УНИФИЦИРОВАН — везде используется -// "02-01-2006 15:04:05.000" с миллисекундами. +// ИСПРАВЛЕНО (2026-10, аудит): см. отдельные пометки в методах. // -// ИСПРАВЛЕНО (2026-10, аудит): -// 1. handleMigrateStatus выводил `fmt.Sprintf("%v", status)` — для -// *migration.MigrationStatus это давало адрес или структуру с -// фигурными скобками. Теперь выводятся конкретные поля статуса. -// 2. handleCompressionConfig упоминал удалённый алгоритм lz4 и не -// упоминал brotli. Теперь перечислены только snappy, brotli, zstd. -// 3. Все обработчики, перебирающие map[string]interface{} без сортировки -// (handleClusterPipeline, handleFallbackStatus, handlePanicStats, -// handlePersistStatus, handleSagaStatus), теперь сортируют ключи — -// порядок строк в таблице детерминирован. -// 4. handleCommitTransaction больше не вызывает bm.UnregisterBatch -// вручную: batch.Commit() сам снимает регистрацию через defer. +// ДОБАВЛЕНО (2026-10, расширение функциональности): +// - config show / get / set / history +// - acl user create/delete, acl role create/delete +// - pager size/threshold +// - saga list --all +// - cluster pipeline on/off +// +// ИСПРАВЛЕНО (2026-10, типы для saga list): +// handleSagaList теперь использует cluster.SagaTransaction +// (а не storage.SagaTransaction), потому что r.coordinator.GetSagaManager() +// возвращает *cluster.SagaManager. // ============================================================================= package repl @@ -59,7 +42,9 @@ import ( "fmt" "io" "os" + "path/filepath" "sort" + "strconv" "strings" "time" @@ -89,7 +74,9 @@ type Repl struct { aclManager *acl.ACLManager pluginManager *plugin.PluginManager - // Терминальный ввод/вывод. + dynamicConfigManager *config.DynamicConfigManager + configPath string + rl *readline.Instance interrupt bool @@ -99,9 +86,6 @@ type Repl struct { authenticated bool sessionID string - // ДОБАВЛЕНО (2026-10): текущий активный batch (замена транзакции). - // Создаётся командой "begin transaction"/"batch begin", коммитится - // командой "commit"/"batch commit", снимается командой "rollback"/"batch abort". currentBatch *storage.Batch currentBatchID uint64 @@ -153,20 +137,16 @@ func NewRepl( historyFile: NewHistory(cfg.Repl.HistorySize), } - // Пейджер. r.pager = NewPager(true, 24) - // Загружаем историю из файла. if err := r.historyFile.Load(); err == nil { r.history = r.historyFile.GetEntries() r.historyPos = len(r.history) } - // Регистрируем команды и автодополнение. r.registerCommands() r.completer = NewCompleter(r) - // Настраиваем readline. rl, err := readline.NewEx(&readline.Config{ Prompt: r.buildPrompt(), HistoryFile: r.historyFile.FilePath(), @@ -175,10 +155,9 @@ func NewRepl( AutoComplete: r.completer.ReadlineCompleter(), InterruptPrompt: "^C", EOFPrompt: "exit", - DisableAutoSaveHistory: true, // сохраняем сами через historyFile - VimMode: false, // emacs-режим как в bash + DisableAutoSaveHistory: true, + VimMode: false, FuncFilterInputRune: func(rn rune) (rune, bool) { - // Блокируем Ctrl+Z, чтобы случайно не свернуть процесс. if rn == readline.CharCtrlZ { return rn, false } @@ -193,29 +172,29 @@ func NewRepl( return r, nil } +// SetDynamicConfigManager устанавливает менеджер динамической конфигурации. +func (r *Repl) SetDynamicConfigManager(dcm *config.DynamicConfigManager, configPath string) { + r.dynamicConfigManager = dcm + r.configPath = configPath +} + // ============================================================================= // РЕГИСТРАЦИЯ КОМАНД // ============================================================================= func (r *Repl) registerCommands() { - // ------------------------------------------------------------------------- - // Базы данных - // ------------------------------------------------------------------------- + // --- Базы данных --- r.commands["create slice"] = &Command{Name: "create slice", Description: "Create a new database (slice)", Handler: r.handleCreateSlice} r.commands["drop database"] = &Command{Name: "drop database", Description: "Drop a database", Handler: r.handleDropDatabase} r.commands["use"] = &Command{Name: "use", Description: "Switch to a database", Handler: r.handleUseDatabase} r.commands["show databases"] = &Command{Name: "show databases", Description: "List all databases", Handler: r.handleShowDatabases} - // ------------------------------------------------------------------------- - // Коллекции - // ------------------------------------------------------------------------- + // --- Коллекции --- r.commands["create collection"] = &Command{Name: "create collection", Description: "Create a new collection in current database", Handler: r.handleCreateCollection} r.commands["drop collection"] = &Command{Name: "drop collection", Description: "Drop a collection from current database", Handler: r.handleDropCollection} r.commands["show collections"] = &Command{Name: "show collections", Description: "List all collections in current database", Handler: r.handleShowCollections} - // ------------------------------------------------------------------------- - // Документы - // ------------------------------------------------------------------------- + // --- Документы --- r.commands["insert"] = &Command{Name: "insert", Description: "Insert a document into a collection (JSON format)", Handler: r.handleInsert} r.commands["find"] = &Command{Name: "find", Description: "Find a document by ID", Handler: r.handleFind} r.commands["findbyindex"] = &Command{Name: "findbyindex", Description: "Find documents by index", Handler: r.handleFindByIndex} @@ -229,25 +208,19 @@ func (r *Repl) registerCommands() { r.commands["show timestamps"] = &Command{Name: "show timestamps", Description: "Show timestamps for a document", Handler: r.handleShowTimestamps} r.commands["stats timestamps"] = &Command{Name: "stats timestamps", Description: "Show timestamp statistics for a collection", Handler: r.handleStatsTimestamps} - // ------------------------------------------------------------------------- - // Индексы - // ------------------------------------------------------------------------- + // --- Индексы --- r.commands["create index"] = &Command{Name: "create index", Description: "Create an index on a collection", Handler: r.handleCreateIndex} r.commands["drop index"] = &Command{Name: "drop index", Description: "Drop an index from a collection", Handler: r.handleDropIndex} r.commands["show indexes"] = &Command{Name: "show indexes", Description: "Show all indexes in a collection", Handler: r.handleShowIndexes} - // ------------------------------------------------------------------------- - // Ограничения - // ------------------------------------------------------------------------- + // --- Ограничения --- r.commands["add required"] = &Command{Name: "add required", Description: "Add a required field constraint", Handler: r.handleAddRequired} r.commands["add unique"] = &Command{Name: "add unique", Description: "Add a unique constraint", Handler: r.handleAddUnique} r.commands["add min"] = &Command{Name: "add min", Description: "Add a minimum value constraint", Handler: r.handleAddMin} r.commands["add max"] = &Command{Name: "add max", Description: "Add a maximum value constraint", Handler: r.handleAddMax} r.commands["add enum"] = &Command{Name: "add enum", Description: "Add an enum constraint (allowed values)", Handler: r.handleAddEnum} - // ------------------------------------------------------------------------- - // Триггеры - // ------------------------------------------------------------------------- + // --- Триггеры --- r.commands["create trigger"] = &Command{Name: "create trigger", Description: "Create a trigger on a collection (MongoDB-like syntax)", Handler: r.handleCreateTrigger} r.commands["drop trigger"] = &Command{Name: "drop trigger", Description: "Drop a trigger from a collection", Handler: r.handleDropTrigger} r.commands["show triggers"] = &Command{Name: "show triggers", Description: "Show all triggers on a collection", Handler: r.handleShowTriggers} @@ -255,29 +228,18 @@ func (r *Repl) registerCommands() { r.commands["disable trigger"] = &Command{Name: "disable trigger", Description: "Disable a trigger", Handler: r.handleDisableTrigger} r.commands["trigger log"] = &Command{Name: "trigger log", Description: "Show trigger execution log", Handler: r.handleTriggerLog} - // ------------------------------------------------------------------------- - // Batch-операции (замена транзакций). - // - // ИСПРАВЛЕНО (2026-10): старые имена команд (begin transaction, commit, - // rollback, show transactions) сохранены как алиасы, чтобы существующие - // скрипты и привычки пользователей не сломались. Новые канонические - // имена — batch begin, batch commit, batch abort, show batches. - // ------------------------------------------------------------------------- - // Канонические batch-команды. + // --- Batch-операции --- r.commands["batch begin"] = &Command{Name: "batch begin", Description: "Begin a new batch (atomically applied group of operations)", Handler: r.handleBeginTransaction} r.commands["batch commit"] = &Command{Name: "batch commit", Description: "Commit the current batch", Handler: r.handleCommitTransaction} r.commands["batch abort"] = &Command{Name: "batch abort", Description: "Abort the current batch (operations are discarded)", Handler: r.handleRollbackTransaction} r.commands["show batches"] = &Command{Name: "show batches", Description: "Show active batch of the current session", Handler: r.handleShowTransactions} - // Алиасы для обратной совместимости. r.commands["begin transaction"] = &Command{Name: "begin transaction", Description: "Alias for 'batch begin'", Handler: r.handleBeginTransaction} r.commands["commit"] = &Command{Name: "commit", Description: "Alias for 'batch commit'", Handler: r.handleCommitTransaction} r.commands["rollback"] = &Command{Name: "rollback", Description: "Alias for 'batch abort'", Handler: r.handleRollbackTransaction} r.commands["show transactions"] = &Command{Name: "show transactions", Description: "Alias for 'show batches'", Handler: r.handleShowTransactions} - // ------------------------------------------------------------------------- - // Плагины - // ------------------------------------------------------------------------- + // --- Плагины --- r.commands["plugin list"] = &Command{Name: "plugin list", Description: "List all loaded plugins", Handler: r.handlePluginList} r.commands["plugin load"] = &Command{Name: "plugin load", Description: "Load a plugin from file", Handler: r.handlePluginLoad} r.commands["plugin unload"] = &Command{Name: "plugin unload", Description: "Unload a plugin", Handler: r.handlePluginUnload} @@ -285,77 +247,60 @@ func (r *Repl) registerCommands() { r.commands["plugin stop"] = &Command{Name: "plugin stop", Description: "Stop a plugin", Handler: r.handlePluginStop} r.commands["plugin exec"] = &Command{Name: "plugin exec", Description: "Execute a plugin function", Handler: r.handlePluginExec} - // ------------------------------------------------------------------------- - // Импорт/экспорт - // ------------------------------------------------------------------------- + // --- Импорт/экспорт --- r.commands["export"] = &Command{Name: "export", Description: "Export database to MessagePack file", Handler: r.handleExport} r.commands["import"] = &Command{Name: "import", Description: "Import database from MessagePack file", Handler: r.handleImport} - // ------------------------------------------------------------------------- - // ACL - // ------------------------------------------------------------------------- + // --- ACL --- r.commands["acl login"] = &Command{Name: "acl login", Description: "Authenticate with username and password", Handler: r.handleACLLogin} r.commands["acl logout"] = &Command{Name: "acl logout", Description: "Logout current user session", Handler: r.handleACLLogout} r.commands["acl grant"] = &Command{Name: "acl grant", Description: "Grant permissions (r=read,w=write,d=delete,a=admin)", Handler: r.handleACLGrant} + r.commands["acl revoke"] = &Command{Name: "acl revoke", Description: "Revoke permissions (r=read,w=write,d=delete,a=admin)", Handler: r.handleACLRevoke} r.commands["acl users"] = &Command{Name: "acl users", Description: "List all users", Handler: r.handleACLUsers} r.commands["acl roles"] = &Command{Name: "acl roles", Description: "List all roles", Handler: r.handleACLRoles} + r.commands["acl user create"] = &Command{Name: "acl user create", Description: "Create a new user", Handler: r.handleACLUserCreate} + r.commands["acl user delete"] = &Command{Name: "acl user delete", Description: "Delete a user", Handler: r.handleACLUserDelete} + r.commands["acl role create"] = &Command{Name: "acl role create", Description: "Create a new role", Handler: r.handleACLRoleCreate} + r.commands["acl role delete"] = &Command{Name: "acl role delete", Description: "Delete a role", Handler: r.handleACLRoleDelete} - // ------------------------------------------------------------------------- - // Сжатие - // ------------------------------------------------------------------------- + // --- Сжатие --- r.commands["compression stats"] = &Command{Name: "compression stats", Description: "Show compression statistics for the database", Handler: r.handleCompressionStats} r.commands["compress collection"] = &Command{Name: "compress collection", Description: "Manually compress all documents in a collection", Handler: r.handleCompressCollection} r.commands["doc compression"] = &Command{Name: "doc compression", Description: "Show compression ratio for a document", Handler: r.handleDocCompression} - r.commands["compression config"] = &Command{Name: "compression config", Description: "Show current compression configuration", Handler: r.handleCompressionConfig} + r.commands["compression config"] = &Command{Name: "compression config", Description: "Show or set compression configuration", Handler: r.handleCompressionConfig} - // ------------------------------------------------------------------------- - // Аудит - // ------------------------------------------------------------------------- + // --- Аудит --- r.commands["audit log"] = &Command{Name: "audit log", Description: "Show audit log", Handler: r.handleAuditLog} r.commands["audit filter"] = &Command{Name: "audit filter", Description: "Filter audit log by type and operation", Handler: r.handleAuditFilter} - // ------------------------------------------------------------------------- - // Кластер - // ------------------------------------------------------------------------- + // --- Кластер --- r.commands["status"] = &Command{Name: "status", Description: "Show cluster status", Handler: r.handleStatus} r.commands["nodes"] = &Command{Name: "nodes", Description: "List cluster nodes", Handler: r.handleNodes} - - // ------------------------------------------------------------------------- - // Кластер (расширенные команды) - // ------------------------------------------------------------------------- r.commands["cluster pipeline"] = &Command{Name: "cluster pipeline", Description: "Show pipeline replicator statistics", Handler: r.handleClusterPipeline} + r.commands["cluster pipeline on"] = &Command{Name: "cluster pipeline on", Description: "Enable pipeline replication", Handler: r.handleClusterPipelineOn} + r.commands["cluster pipeline off"] = &Command{Name: "cluster pipeline off", Description: "Disable pipeline replication", Handler: r.handleClusterPipelineOff} r.commands["cluster reshard"] = &Command{Name: "cluster reshard", Description: "Trigger manual resharding", Handler: r.handleClusterReshard} - // ------------------------------------------------------------------------- - // Миграция - // ------------------------------------------------------------------------- + // --- Миграция --- r.commands["migration"] = &Command{Name: "migration", Description: "Cross-datacenter migration commands", Handler: r.handleMigration} r.commands["migrate status"] = &Command{Name: "migrate status", Description: "Show schema migration status", Handler: r.handleMigrateStatus} - // ------------------------------------------------------------------------- - // Fallback (SPoF protection) - // ------------------------------------------------------------------------- + // --- Fallback / Panic / Persist --- r.commands["fallback status"] = &Command{Name: "fallback status", Description: "Show SPoF protection status", Handler: r.handleFallbackStatus} - - // ------------------------------------------------------------------------- - // Panic recovery - // ------------------------------------------------------------------------- r.commands["panic stats"] = &Command{Name: "panic stats", Description: "Show panic recovery statistics", Handler: r.handlePanicStats} - - // ------------------------------------------------------------------------- - // Persistence - // ------------------------------------------------------------------------- r.commands["persist status"] = &Command{Name: "persist status", Description: "Show data persistence status", Handler: r.handlePersistStatus} - // ------------------------------------------------------------------------- - // SAGA - // ------------------------------------------------------------------------- + // --- SAGA --- r.commands["saga status"] = &Command{Name: "saga status", Description: "Show SAGA orchestrator status", Handler: r.handleSagaStatus} r.commands["saga list"] = &Command{Name: "saga list", Description: "List active SAGA transactions", Handler: r.handleSagaList} - // ------------------------------------------------------------------------- - // Системные - // ------------------------------------------------------------------------- + // --- Динамическая конфигурация --- + r.commands["config show"] = &Command{Name: "config show", Description: "Show current configuration (or section)", Handler: r.handleConfigShow} + r.commands["config get"] = &Command{Name: "config get", Description: "Get a single config value by path", Handler: r.handleConfigGet} + r.commands["config set"] = &Command{Name: "config set", Description: "Set a config value (persisted via Raft or config.toml)", Handler: r.handleConfigSet} + r.commands["config history"] = &Command{Name: "config history", Description: "Show dynamic config change history", Handler: r.handleConfigHistory} + + // --- Системные --- r.commands["help"] = &Command{Name: "help", Description: "Show this help message", Handler: r.handleHelp} r.commands["clear"] = &Command{Name: "clear", Description: "Clear the screen", Handler: r.handleClear} r.commands["quit"] = &Command{Name: "quit", Description: "Exit the REPL", Handler: r.handleQuit} @@ -366,7 +311,6 @@ func (r *Repl) registerCommands() { // ОСНОВНОЙ ЦИКЛ REPL // ============================================================================= -// Run запускает основной цикл REPL с использованием readline. func (r *Repl) Run() error { defer r.rl.Close() @@ -377,18 +321,15 @@ func (r *Repl) Run() error { utils.Println("") for { - // Обновляем prompt (может измениться currentDB/currentUser). r.rl.SetPrompt(r.buildPrompt()) line, err := r.rl.Readline() if err != nil { switch { case errors.Is(err, readline.ErrInterrupt): - // Ctrl+C — прерываем текущий ввод, не выходим. utils.Println("^C") continue case errors.Is(err, io.EOF): - // Ctrl+D — выход. return nil default: return err @@ -400,7 +341,6 @@ func (r *Repl) Run() error { continue } - // Многострочный ввод. full, err := r.collectMultiline(line) if err != nil { if errors.Is(err, io.EOF) { @@ -416,7 +356,6 @@ func (r *Repl) Run() error { continue } - // Сохраняем в историю (readline + наш файл). r.addToHistory(full) if err := r.executeCommand(full); err != nil { @@ -428,7 +367,6 @@ func (r *Repl) Run() error { } } -// collectMultiline продолжает ввод, если строка не закончена. func (r *Repl) collectMultiline(first string) (string, error) { acc := first depthCurly, depthSquare, depthParen := countBrackets(acc) @@ -436,7 +374,6 @@ func (r *Repl) collectMultiline(first string) (string, error) { for { cont := strings.HasSuffix(acc, "\\") || depthCurly > 0 || depthSquare > 0 || depthParen > 0 - if !cont { break } @@ -464,7 +401,6 @@ func (r *Repl) collectMultiline(first string) (string, error) { return strings.TrimSpace(acc), nil } -// countBrackets считает баланс {}, [], () с учётом строк. func countBrackets(s string) (curly, square, paren int) { inString := false var quote rune @@ -512,24 +448,23 @@ func countBrackets(s string) (curly, square, paren int) { func (r *Repl) buildPrompt() string { prompt := color.New(color.FgHiCyan).Sprint("futriiX") - if r.currentDB != "" { prompt += color.New(color.FgHiYellow).Sprint(":" + r.currentDB) } - if r.authenticated && r.currentUser != "" { prompt += color.New(color.FgHiGreen).Sprint(" (" + r.currentUser + ")") } - - // ДОБАВЛЕНО: индикатор активного batch'а. if r.currentBatch != nil { prompt += color.New(color.FgHiMagenta).Sprint(" [batch]") } - prompt += color.New(color.FgHiCyan).Sprint(":~> ") return prompt } +// executeCommand разбирает команду и вызывает соответствующий обработчик. +// +// ДОБАВЛЕНО: команды сортируются по убыванию длины. Это устраняет +// проблему перекрывающихся префиксов. func (r *Repl) executeCommand(input string) error { parts := strings.Fields(input) if len(parts) == 0 { @@ -538,25 +473,20 @@ func (r *Repl) executeCommand(input string) error { switch parts[0] { case "pager": - if len(parts) > 1 { - switch parts[1] { - case "on": - r.pager.SetEnabled(true) - utils.PrintSuccess("Pager enabled") - case "off": - r.pager.SetEnabled(false) - utils.PrintSuccess("Pager disabled") - default: - utils.PrintInfo(fmt.Sprintf("Pager: %v", r.pager.Enabled())) - } - } else { - utils.PrintInfo(fmt.Sprintf("Pager: %v", r.pager.Enabled())) - } - return nil + return r.handlePagerCommand(parts[1:]) } - for cmdName, cmd := range r.commands { + names := make([]string, 0, len(r.commands)) + for name := range r.commands { + names = append(names, name) + } + sort.Slice(names, func(i, j int) bool { + return len(names[i]) > len(names[j]) + }) + + for _, cmdName := range names { if strings.HasPrefix(input, cmdName) { + cmd := r.commands[cmdName] args := strings.TrimPrefix(input, cmdName) args = strings.TrimSpace(args) argList := strings.Fields(args) @@ -571,7 +501,6 @@ func (r *Repl) addToHistory(cmd string) { if len(r.history) > 0 && r.history[len(r.history)-1] == cmd { return } - if len(r.history) >= r.config.Repl.HistorySize { r.history = r.history[1:] } @@ -583,10 +512,9 @@ func (r *Repl) addToHistory(cmd string) { } // ============================================================================= -// PgUp/PgDown — программная навигация +// PgUp/PgDown // ============================================================================= -// HandlePgUp возвращает предыдущую команду из истории. func (r *Repl) HandlePgUp() string { if len(r.history) == 0 { return "" @@ -600,7 +528,6 @@ func (r *Repl) HandlePgUp() string { return r.history[r.historyPos] } -// HandlePgDown возвращает следующую команду из истории. func (r *Repl) HandlePgDown() string { if len(r.history) == 0 { return "" @@ -631,8 +558,6 @@ func (r *Repl) output(text string) { // ============================================================================= func (r *Repl) Close() error { - // ДОБАВЛЕНО: снимаем активный batch при выходе, чтобы не оставить - // его зарегистрированным в глобальном менеджере. if r.currentBatch != nil { if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(r.currentBatch) @@ -653,11 +578,44 @@ func (r *Repl) Close() error { // ОБРАБОТЧИКИ КОМАНД // ============================================================================= +// ----------------------------------------------------------------------------- +// Pager +// ----------------------------------------------------------------------------- + +func (r *Repl) handlePagerCommand(args []string) error { + if len(args) == 0 { + utils.PrintInfo(fmt.Sprintf("Pager: %v (threshold=%d lines)", r.pager.Enabled(), r.pager.Threshold())) + return nil + } + + switch args[0] { + case "on": + r.pager.SetEnabled(true) + utils.PrintSuccess("Pager enabled") + case "off": + r.pager.SetEnabled(false) + utils.PrintSuccess("Pager disabled") + case "size", "threshold": + if len(args) < 2 { + utils.PrintInfo(fmt.Sprintf("Pager threshold: %d lines", r.pager.Threshold())) + return nil + } + n, err := strconv.Atoi(args[1]) + if err != nil || n <= 0 { + return fmt.Errorf("invalid threshold '%s' (expected positive integer)", args[1]) + } + r.pager.SetThreshold(n) + utils.PrintSuccess(fmt.Sprintf("Pager threshold = %d lines", n)) + default: + utils.PrintInfo(fmt.Sprintf("Pager: %v (threshold=%d lines)", r.pager.Enabled(), r.pager.Threshold())) + } + return nil +} + // ----------------------------------------------------------------------------- // Базы данных // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleCreateSlice(args []string) error { if len(args) < 1 { return fmt.Errorf("usage: create slice ") @@ -670,7 +628,6 @@ func (r *Repl) handleCreateSlice(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleDropDatabase(args []string) error { if len(args) < 1 { return fmt.Errorf("usage: drop database ") @@ -705,7 +662,6 @@ func (r *Repl) handleShowDatabases(args []string) error { utils.PrintInfo("No databases found") return nil } - t := NewTable("DATABASE", "CURRENT") for _, db := range databases { cur := "" @@ -722,7 +678,6 @@ func (r *Repl) handleShowDatabases(args []string) error { // Коллекции // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleCreateCollection(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -774,7 +729,6 @@ func (r *Repl) handleShowCollections(args []string) error { utils.PrintInfo("No collections found") return nil } - t := NewTable("COLLECTION") for _, coll := range collections { t.AddRow(coll) @@ -787,7 +741,6 @@ func (r *Repl) handleShowCollections(args []string) error { // Документы // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleInsert(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -795,7 +748,6 @@ func (r *Repl) handleInsert(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: insert ") } - collName := args[0] jsonStr := strings.Join(args[1:], " ") @@ -809,7 +761,6 @@ func (r *Repl) handleInsert(args []string) error { } doc := storage.NewDocument() - if strings.Contains(jsonStr, "{") { pairs := strings.Split(jsonStr, ",") for _, pair := range pairs { @@ -839,7 +790,6 @@ func (r *Repl) handleInsert(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleFind(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -847,7 +797,6 @@ func (r *Repl) handleFind(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: find ") } - collName := args[0] docID := args[1] @@ -870,7 +819,6 @@ func (r *Repl) handleFind(args []string) error { row[k] = v } rows := []map[string]interface{}{row} - t := MapToTable(rows) r.output(t.String()) @@ -889,7 +837,6 @@ func (r *Repl) handleFindByIndex(args []string) error { if len(args) < 3 { return fmt.Errorf("usage: findbyindex ") } - collName := args[0] indexName := args[1] value := args[2] @@ -906,12 +853,10 @@ func (r *Repl) handleFindByIndex(args []string) error { if err != nil { return err } - if len(docs) == 0 { utils.PrintInfo("No documents found") return nil } - rows := make([]map[string]interface{}, 0, len(docs)) for _, doc := range docs { row := map[string]interface{}{"_id": doc.ID} @@ -930,41 +875,31 @@ func (r *Repl) handleFindByTime(args []string) error { return fmt.Errorf("no database selected") } if len(args) < 3 { - return fmt.Errorf("usage: findbytime \n"+ - " Date formats: YYYY-MM-DD or YYYY-MM-DD HH:MM:SS") + return fmt.Errorf("usage: findbytime \n" + + " Date formats: YYYY-MM-DD or YYYY-MM-DDTHH:MM:SS (ISO 8601, без пробелов)") + } + if len(args) > 3 { + return fmt.Errorf("too many arguments: dates must not contain spaces.\n" + + " Use YYYY-MM-DD or YYYY-MM-DDTHH:MM:SS (ISO 8601).") } collName := args[0] fromStr := args[1] toStr := args[2] - var fromTime, toTime time.Time - var err error - formats := []string{ "2006-01-02", - "2006-01-02 15:04:05", "2006-01-02T15:04:05", + "2006-01-02T15:04:05Z07:00", } - for _, format := range formats { - fromTime, err = time.Parse(format, fromStr) - if err == nil { - break - } - } + fromTime, err := parseFlexibleTime(fromStr, formats) if err != nil { - return fmt.Errorf("invalid from_date format: %s", fromStr) - } - - for _, format := range formats { - toTime, err = time.Parse(format, toStr) - if err == nil { - break - } + return fmt.Errorf("invalid from_date '%s': %v", fromStr, err) } + toTime, err := parseFlexibleTime(toStr, formats) if err != nil { - return fmt.Errorf("invalid to_date format: %s", toStr) + return fmt.Errorf("invalid to_date '%s': %v", toStr, err) } fromMs := fromTime.UnixMilli() @@ -981,12 +916,10 @@ func (r *Repl) handleFindByTime(args []string) error { docs := coll.FindByFilter(func(doc *storage.Document) bool { return doc.CreatedAt >= fromMs && doc.CreatedAt <= toMs }) - if len(docs) == 0 { utils.PrintInfo("No documents found in the specified time range") return nil } - rows := make([]map[string]interface{}, 0, len(docs)) for _, doc := range docs { row := map[string]interface{}{"_id": doc.ID} @@ -1000,7 +933,15 @@ func (r *Repl) handleFindByTime(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". +func parseFlexibleTime(s string, formats []string) (time.Time, error) { + for _, f := range formats { + if t, err := time.Parse(f, s); err == nil { + return t, nil + } + } + return time.Time{}, fmt.Errorf("expected one of: %s", strings.Join(formats, ", ")) +} + func (r *Repl) handleUpdate(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1008,10 +949,8 @@ func (r *Repl) handleUpdate(args []string) error { if len(args) < 3 { return fmt.Errorf("usage: update ...") } - collName := args[0] docID := args[1] - updates := make(map[string]interface{}) for i := 2; i < len(args); i++ { kv := strings.SplitN(args[i], "=", 2) @@ -1019,7 +958,6 @@ func (r *Repl) handleUpdate(args []string) error { updates[kv[0]] = kv[1] } } - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1035,7 +973,6 @@ func (r *Repl) handleUpdate(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleDelete(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1043,10 +980,8 @@ func (r *Repl) handleDelete(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: delete ") } - collName := args[0] docID := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1069,10 +1004,8 @@ func (r *Repl) handlePermanentDelete(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: permanent delete ") } - collName := args[0] docID := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1088,7 +1021,6 @@ func (r *Repl) handlePermanentDelete(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleRestore(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1096,10 +1028,8 @@ func (r *Repl) handleRestore(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: restore ") } - collName := args[0] docID := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1115,7 +1045,6 @@ func (r *Repl) handleRestore(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleShowDeleted(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1123,9 +1052,7 @@ func (r *Repl) handleShowDeleted(args []string) error { if len(args) < 1 { return fmt.Errorf("usage: show deleted ") } - collName := args[0] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1134,7 +1061,6 @@ func (r *Repl) handleShowDeleted(args []string) error { if err != nil { return err } - docs := coll.GetAllDocumentsIncludingDeleted() deleted := make([]*storage.Document, 0) for _, doc := range docs { @@ -1142,12 +1068,10 @@ func (r *Repl) handleShowDeleted(args []string) error { deleted = append(deleted, doc) } } - if len(deleted) == 0 { utils.PrintInfo("No deleted documents found") return nil } - t := NewTable("ID", "DELETED_AT", "FIELDS") for _, doc := range deleted { t.AddRow( @@ -1167,9 +1091,7 @@ func (r *Repl) handleCount(args []string) error { if len(args) < 1 { return fmt.Errorf("usage: count ") } - collName := args[0] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1178,11 +1100,9 @@ func (r *Repl) handleCount(args []string) error { if err != nil { return err } - active := coll.Count() deleted := coll.CountDeleted() total := coll.CountAll() - t := NewTable("METRIC", "VALUE") t.AddRow("Active documents", fmt.Sprintf("%d", active)) t.AddRow("Deleted documents", fmt.Sprintf("%d", deleted)) @@ -1195,7 +1115,6 @@ func (r *Repl) handleCount(args []string) error { // Временные метки // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleShowTimestamps(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1203,10 +1122,8 @@ func (r *Repl) handleShowTimestamps(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: show timestamps ") } - collName := args[0] docID := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1219,7 +1136,6 @@ func (r *Repl) handleShowTimestamps(args []string) error { if err != nil { return err } - t := NewTable("FIELD", "VALUE") t.AddRow("Created", time.UnixMilli(doc.CreatedAt).Format("02-01-2006 15:04:05.000")) t.AddRow("Updated", time.UnixMilli(doc.UpdatedAt).Format("02-01-2006 15:04:05.000")) @@ -1233,7 +1149,6 @@ func (r *Repl) handleShowTimestamps(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleStatsTimestamps(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1241,9 +1156,7 @@ func (r *Repl) handleStatsTimestamps(args []string) error { if len(args) < 1 { return fmt.Errorf("usage: stats timestamps ") } - collName := args[0] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1252,7 +1165,6 @@ func (r *Repl) handleStatsTimestamps(args []string) error { if err != nil { return err } - docs := coll.GetAllDocumentsIncludingDeleted() if len(docs) == 0 { utils.PrintInfo("No documents in collection") @@ -1308,12 +1220,10 @@ func (r *Repl) handleAuditLog(args []string) error { utils.PrintInfo("No audit log entries") return nil } - start := 0 if len(entries) > 50 { start = len(entries) - 50 } - t := NewTable("TIMESTAMP", "OPERATION", "TYPE", "NAME") for i := len(entries) - 1; i >= start; i-- { entry := entries[i] @@ -1324,21 +1234,47 @@ func (r *Repl) handleAuditLog(args []string) error { } func (r *Repl) handleAuditFilter(args []string) error { - if len(args) < 2 { - return fmt.Errorf("usage: audit filter \n"+ - " data_type: DATABASE, COLLECTION, DOCUMENT, FIELD, INDEX, TRANSACTION\n"+ - " operation: CREATE, INSERT, UPDATE, DELETE, SOFT_DELETE, RESTORE") + if len(args) < 1 { + return fmt.Errorf("usage: audit filter [] [from=] [to=]\n" + + " data_type: DATABASE, COLLECTION, DOCUMENT, FIELD, INDEX, TRANSACTION\n" + + " operation: CREATE, INSERT, UPDATE, DELETE, SOFT_DELETE, RESTORE\n" + + " Example: audit filter DOCUMENT INSERT from=1700000000000 to=1800000000000") + } + dataType := strings.ToUpper(args[0]) + operation := "" + var fromTime, toTime int64 + + for i := 1; i < len(args); i++ { + arg := args[i] + switch { + case strings.HasPrefix(arg, "from="): + val := strings.TrimPrefix(arg, "from=") + parsed, err := strconv.ParseInt(val, 10, 64) + if err != nil { + return fmt.Errorf("invalid from= value '%s': %v", val, err) + } + fromTime = parsed + case strings.HasPrefix(arg, "to="): + val := strings.TrimPrefix(arg, "to=") + parsed, err := strconv.ParseInt(val, 10, 64) + if err != nil { + return fmt.Errorf("invalid to= value '%s': %v", val, err) + } + toTime = parsed + default: + if operation == "" { + operation = strings.ToUpper(arg) + } else { + return fmt.Errorf("unexpected argument: %s", arg) + } + } } - dataType := strings.ToUpper(args[0]) - operation := strings.ToUpper(args[1]) - - entries := storage.GetAuditLogFiltered(dataType, operation, 0, 0) + entries := storage.GetAuditLogFiltered(dataType, operation, fromTime, toTime) if len(entries) == 0 { utils.PrintInfo(fmt.Sprintf("No audit log entries found for %s/%s", dataType, operation)) return nil } - t := NewTable("TIMESTAMP", "OPERATION", "TYPE", "NAME") for _, entry := range entries { t.AddRow(entry.TimestampStr, entry.Operation, entry.DataType, entry.Name) @@ -1351,7 +1287,6 @@ func (r *Repl) handleAuditFilter(args []string) error { // Индексы // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleCreateIndex(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1359,7 +1294,6 @@ func (r *Repl) handleCreateIndex(args []string) error { if len(args) < 3 { return fmt.Errorf("usage: create index [unique]") } - collName := args[0] indexName := args[1] fields := strings.Split(args[2], ",") @@ -1387,10 +1321,8 @@ func (r *Repl) handleDropIndex(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: drop index ") } - collName := args[0] indexName := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1406,8 +1338,6 @@ func (r *Repl) handleDropIndex(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handleShowIndexes(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1415,9 +1345,7 @@ func (r *Repl) handleShowIndexes(args []string) error { if len(args) < 1 { return fmt.Errorf("usage: show indexes ") } - collName := args[0] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1426,13 +1354,11 @@ func (r *Repl) handleShowIndexes(args []string) error { if err != nil { return err } - indexes := coll.GetIndexesInfo() if len(indexes) == 0 { utils.PrintInfo(fmt.Sprintf("No indexes found on collection '%s'", collName)) return nil } - t := NewTable("NAME", "FIELDS", "UNIQUE", "CREATED_AT") for _, idx := range indexes { uniqueStr := "false" @@ -1468,10 +1394,8 @@ func (r *Repl) handleAddRequired(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: add required ") } - collName := args[0] field := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1492,10 +1416,8 @@ func (r *Repl) handleAddUnique(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: add unique ") } - collName := args[0] field := args[1] - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1516,14 +1438,12 @@ func (r *Repl) handleAddMin(args []string) error { if len(args) < 3 { return fmt.Errorf("usage: add min ") } - collName := args[0] field := args[1] var minVal float64 if _, err := fmt.Sscanf(args[2], "%f", &minVal); err != nil { return fmt.Errorf("invalid minimum value: %s", args[2]) } - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1544,14 +1464,12 @@ func (r *Repl) handleAddMax(args []string) error { if len(args) < 3 { return fmt.Errorf("usage: add max ") } - collName := args[0] field := args[1] var maxVal float64 if _, err := fmt.Sscanf(args[2], "%f", &maxVal); err != nil { return fmt.Errorf("invalid maximum value: %s", args[2]) } - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1572,14 +1490,12 @@ func (r *Repl) handleAddEnum(args []string) error { if len(args) < 3 { return fmt.Errorf("usage: add enum ") } - collName := args[0] field := args[1] values := make([]interface{}, len(args[2:])) for i, v := range args[2:] { values[i] = v } - db, err := r.store.GetDatabase(r.currentDB) if err != nil { return err @@ -1594,20 +1510,9 @@ func (r *Repl) handleAddEnum(args []string) error { } // ----------------------------------------------------------------------------- -// Batch-операции (замена транзакций). -// -// ИСПРАВЛЕНО (2026-10): обработчики переписаны поверх batch-API. -// Семантика: -// - handleBeginTransaction → создаёт *storage.Batch (в памяти REPL); -// - handleCommitTransaction → batch.Commit() (атомарно через WAL); -// - handleRollbackTransaction→ снимает регистрацию batch без коммита; -// - handleShowTransactions → показывает активный batch сессии. -// -// ИЗМЕНЕНО (2026-10-03): формат временных меток "02-01-2006 15:04:05.000". +// Batch-операции // ----------------------------------------------------------------------------- -// handleBeginTransaction создаёт новый batch. -// Команда: "begin transaction" или "batch begin". func (r *Repl) handleBeginTransaction(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -1615,83 +1520,60 @@ func (r *Repl) handleBeginTransaction(args []string) error { if r.currentBatch != nil { return fmt.Errorf("batch already active (ID %d); commit or abort first", r.currentBatchID) } - - // Создаём batch для текущей БД. Коллекция может быть указана как аргумент. collection := "" if len(args) >= 1 { collection = args[0] } - batch := storage.NewBatch(r.currentDB, collection) if batch == nil { return fmt.Errorf("failed to create batch") } - r.currentBatch = batch r.currentBatchID = uint64(batch.ID) - utils.PrintSuccess(fmt.Sprintf("Batch %d started at %s (database=%s, collection=%s)", batch.ID, time.Now().Format("02-01-2006 15:04:05.000"), batch.Database, batch.Collection)) return nil } -// handleCommitTransaction коммитит текущий batch. -// Команда: "commit" или "batch commit". -// -// ИСПРАВЛЕНО: убран избыточный вызов bm.UnregisterBatch — Commit() -// сам снимает регистрацию через defer (см. transactions.go). func (r *Repl) handleCommitTransaction(args []string) error { if r.currentBatch == nil { return fmt.Errorf("no active batch to commit") } - batch := r.currentBatch if err := batch.Commit(); err != nil { r.currentBatch = nil r.currentBatchID = 0 return err } - r.currentBatch = nil r.currentBatchID = 0 - utils.PrintSuccess(fmt.Sprintf("Batch %d committed successfully at %s", batch.ID, time.Now().Format("02-01-2006 15:04:05.000"))) return nil } -// handleRollbackTransaction отменяет текущий batch (без коммита). -// Команда: "rollback" или "batch abort". func (r *Repl) handleRollbackTransaction(args []string) error { if r.currentBatch == nil { return fmt.Errorf("no active batch to rollback") } - batch := r.currentBatch - // Снимаем регистрацию в глобальном менеджере. if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(batch) } batchID := r.currentBatchID r.currentBatch = nil r.currentBatchID = 0 - utils.PrintSuccess(fmt.Sprintf("Batch %d aborted at %s (operations discarded)", batchID, time.Now().Format("02-01-2006 15:04:05.000"))) return nil } -// handleShowTransactions выводит список активных batch'ей сессии. -// Команда: "show transactions" или "show batches". -// -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleShowTransactions(args []string) error { if r.currentBatch == nil { utils.PrintInfo("No active batches") return nil } - t := NewTable("ID", "DATABASE", "COLLECTION", "OPERATIONS", "STARTED") started := "" if r.currentBatch.Timestamp > 0 { @@ -1707,13 +1589,10 @@ func (r *Repl) handleShowTransactions(args []string) error { r.output(t.String()) return nil } - // ----------------------------------------------------------------------------- // Плагины // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handlePluginList(args []string) error { if r.pluginManager == nil || !r.pluginManager.IsEnabled() { return fmt.Errorf("plugin system is disabled") @@ -1742,7 +1621,6 @@ func (r *Repl) handlePluginList(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handlePluginLoad(args []string) error { if r.pluginManager == nil || !r.pluginManager.IsEnabled() { return fmt.Errorf("plugin system is disabled") @@ -1774,7 +1652,6 @@ func (r *Repl) handlePluginUnload(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handlePluginStart(args []string) error { if r.pluginManager == nil || !r.pluginManager.IsEnabled() { return fmt.Errorf("plugin system is disabled") @@ -1833,7 +1710,6 @@ func (r *Repl) handlePluginExec(args []string) error { // Экспорт/импорт // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleExport(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: export ") @@ -1862,7 +1738,6 @@ func (r *Repl) handleExport(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleImport(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: import ") @@ -1895,10 +1770,9 @@ func (r *Repl) handleImport(args []string) error { } // ----------------------------------------------------------------------------- -// ACL +// ACL — базовая аутентификация и выдача прав // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleACLLogin(args []string) error { if len(args) < 2 { return fmt.Errorf("usage: acl login ") @@ -1928,7 +1802,6 @@ func (r *Repl) handleACLLogin(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleACLLogout(args []string) error { if r.sessionID != "" && r.aclManager != nil { r.aclManager.Logout(r.sessionID) @@ -1941,7 +1814,6 @@ func (r *Repl) handleACLLogout(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleACLGrant(args []string) error { if !r.authenticated || r.currentRole != "admin" { return fmt.Errorf("permission denied: admin access required") @@ -1980,6 +1852,63 @@ func (r *Repl) handleACLGrant(args []string) error { return nil } +// handleACLRevoke отзывает разрешения у роли на коллекции. +// +// ДОБАВЛЕНО: симметрично acl grant. +func (r *Repl) handleACLRevoke(args []string) error { + if !r.authenticated || r.currentRole != "admin" { + return fmt.Errorf("permission denied: admin access required") + } + if len(args) < 3 { + return fmt.Errorf("usage: acl revoke \n"+ + " Permissions: r=read, w=write, d=delete, a=admin\n"+ + " Example: acl revoke users guest rw") + } + + collName := args[0] + role := args[1] + perms := args[2] + + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + revoked := make([]string, 0) + if strings.Contains(perms, "r") { + coll.RevokeACL(role, "read") + revoked = append(revoked, "read") + } + if strings.Contains(perms, "w") { + coll.RevokeACL(role, "write") + revoked = append(revoked, "write") + } + if strings.Contains(perms, "d") { + coll.RevokeACL(role, "delete") + revoked = append(revoked, "delete") + } + if strings.Contains(perms, "a") { + coll.RevokeACL(role, "admin") + revoked = append(revoked, "admin") + } + + if len(revoked) == 0 { + return fmt.Errorf("no permissions specified in '%s'", perms) + } + + utils.PrintSuccess(fmt.Sprintf("Permissions [%s] revoked from role '%s' on collection '%s' at %s", + strings.Join(revoked, ","), role, collName, time.Now().Format("02-01-2006 15:04:05.000"))) + return nil +} + func (r *Repl) handleACLUsers(args []string) error { if r.aclManager == nil { return fmt.Errorf("ACL manager not initialized") @@ -2028,6 +1957,103 @@ func (r *Repl) handleACLRoles(args []string) error { return nil } +// ----------------------------------------------------------------------------- +// ACL — управление пользователями и ролями через REPL. +// ----------------------------------------------------------------------------- + +// handleACLUserCreate создаёт нового пользователя. +func (r *Repl) handleACLUserCreate(args []string) error { + if !r.authenticated || r.currentRole != "admin" { + return fmt.Errorf("permission denied: admin access required") + } + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + if len(args) < 2 { + return fmt.Errorf("usage: acl user create [role1,role2,...]\n" + + " Example: acl user create alice secret123 admin") + } + + username := args[0] + password := args[1] + + var roles []string + if len(args) >= 3 { + roles = strings.Split(args[2], ",") + for i := range roles { + roles[i] = strings.TrimSpace(roles[i]) + } + } + + if err := r.aclManager.CreateUser(username, password, roles); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("User '%s' created with roles %v at %s", + username, roles, time.Now().Format("02-01-2006 15:04:05.000"))) + return nil +} + +// handleACLUserDelete удаляет пользователя. +func (r *Repl) handleACLUserDelete(args []string) error { + if !r.authenticated || r.currentRole != "admin" { + return fmt.Errorf("permission denied: admin access required") + } + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + if len(args) < 1 { + return fmt.Errorf("usage: acl user delete ") + } + username := args[0] + + if err := r.aclManager.DeleteUser(username); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("User '%s' deleted", username)) + return nil +} + +// handleACLRoleCreate создаёт новую роль. +func (r *Repl) handleACLRoleCreate(args []string) error { + if !r.authenticated || r.currentRole != "admin" { + return fmt.Errorf("permission denied: admin access required") + } + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + if len(args) < 1 { + return fmt.Errorf("usage: acl role create ") + } + roleName := args[0] + + if err := r.aclManager.CreateRole(roleName); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Role '%s' created at %s", + roleName, time.Now().Format("02-01-2006 15:04:05.000"))) + return nil +} + +// handleACLRoleDelete удаляет роль. +func (r *Repl) handleACLRoleDelete(args []string) error { + if !r.authenticated || r.currentRole != "admin" { + return fmt.Errorf("permission denied: admin access required") + } + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + if len(args) < 1 { + return fmt.Errorf("usage: acl role delete ") + } + roleName := args[0] + + if err := r.aclManager.DeleteRole(roleName); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Role '%s' deleted", roleName)) + return nil +} + // ----------------------------------------------------------------------------- // Сжатие // ----------------------------------------------------------------------------- @@ -2084,7 +2110,6 @@ func (r *Repl) handleCompressionStats(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleCompressCollection(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -2161,17 +2186,46 @@ func (r *Repl) handleDocCompression(args []string) error { return nil } -// handleCompressionConfig показывает текущие настройки сжатия. +// handleCompressionConfig показывает или изменяет настройки сжатия. // -// ИСПРАВЛЕНО: раньше упоминался удалённый алгоритм lz4 и не упоминался -// brotli. Теперь перечислены только snappy, brotli, zstd — те, что -// поддерживаются в internal/compression/compression.go. +// РАСШИРЕНО: подкоманда set делегирует в handleConfigSet — +// изменения сохраняются через Raft или config.toml. func (r *Repl) handleCompressionConfig(args []string) error { + if len(args) >= 1 && args[0] == "set" { + if len(args) < 3 { + return fmt.Errorf("usage: compression config set \n" + + " keys: enabled (true|false), algorithm (snappy|brotli|zstd), level (1-9), min_size (bytes)") + } + key := strings.ToLower(args[1]) + value := args[2] + + var path string + switch key { + case "enabled": + path = "compression.enabled" + case "algorithm": + path = "compression.algorithm" + case "level": + path = "compression.level" + case "min_size": + path = "compression.min_size" + default: + return fmt.Errorf("unknown key '%s' (supported: enabled, algorithm, level, min_size)", key) + } + + return r.handleConfigSet([]string{path, value}) + } + + cfg := r.effectiveConfig() + if cfg == nil { + return fmt.Errorf("config is not available") + } + t := NewTable("SETTING", "VALUE") - t.AddRow("Enabled", fmt.Sprintf("%v", r.config.Compression.Enabled)) - t.AddRow("Algorithm", r.config.Compression.Algorithm) - t.AddRow("Level", fmt.Sprintf("%d", r.config.Compression.Level)) - t.AddRow("Min Size", utils.FormatBytes(int64(r.config.Compression.MinSize))) + t.AddRow("Enabled", fmt.Sprintf("%v", cfg.Compression.Enabled)) + t.AddRow("Algorithm", cfg.Compression.Algorithm) + t.AddRow("Level", fmt.Sprintf("%d", cfg.Compression.Level)) + t.AddRow("Min Size", utils.FormatBytes(int64(cfg.Compression.MinSize))) t.AddRow("snappy", "Fast compression/decompression, good balance") t.AddRow("brotli", "Best compression ratio (recommended)") t.AddRow("zstd", "High compression ratio, configurable level") @@ -2183,16 +2237,23 @@ func (r *Repl) handleCompressionConfig(args []string) error { // Триггеры // ----------------------------------------------------------------------------- -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". +// handleCreateTrigger создаёт триггер на коллекции. +// +// ИСПРАВЛЕНО: теперь парсятся все опции, заявленные в usage. func (r *Repl) handleCreateTrigger(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") } if len(args) < 4 { - return fmt.Errorf("usage: create trigger [options]\n"+ - " Events: BEFORE_INSERT, AFTER_INSERT, BEFORE_UPDATE, AFTER_UPDATE, BEFORE_DELETE, AFTER_DELETE\n"+ - " Actions: abort, skip, modify, log, notify\n"+ - " Options: --description , --set , --inc , --currentDate , --condition ") + return fmt.Errorf("usage: create trigger [options]\n" + + " Events: BEFORE_INSERT, AFTER_INSERT, BEFORE_UPDATE, AFTER_UPDATE, BEFORE_DELETE, AFTER_DELETE\n" + + " Actions: abort, skip, modify, log, notify\n" + + " Options:\n" + + " --description \n" + + " --set \n" + + " --inc \n" + + " --currentDate \n" + + " --condition (op: eq, ne, gt, lt, gte, lte, in, nin, exists, regex)") } collName := args[0] @@ -2215,17 +2276,77 @@ func (r *Repl) handleCreateTrigger(args []string) error { Action: action, Enabled: true, Description: "", + Operations: make([]storage.TriggerOperation, 0), CreatedAt: time.Now().UnixMilli(), UpdatedAt: time.Now().UnixMilli(), } - for i := 4; i < len(args); i++ { - switch args[i] { + i := 4 + for i < len(args) { + opt := args[i] + switch opt { case "--description": - if i+1 < len(args) { - trigger.Description = args[i+1] - i++ + if i+1 >= len(args) { + return fmt.Errorf("--description requires a value") } + trigger.Description = args[i+1] + i += 2 + + case "--set": + if i+2 >= len(args) { + return fmt.Errorf("--set requires ") + } + trigger.Operations = append(trigger.Operations, storage.TriggerOperation{ + Type: "set", + Field: args[i+1], + Value: args[i+2], + }) + i += 3 + + case "--inc": + if i+2 >= len(args) { + return fmt.Errorf("--inc requires ") + } + var incVal float64 + if _, err := fmt.Sscanf(args[i+2], "%f", &incVal); err != nil { + return fmt.Errorf("--inc: invalid number '%s'", args[i+2]) + } + trigger.Operations = append(trigger.Operations, storage.TriggerOperation{ + Type: "inc", + Field: args[i+1], + Value: incVal, + }) + i += 3 + + case "--currentDate", "--currentdate": + if i+1 >= len(args) { + return fmt.Errorf("--currentDate requires ") + } + trigger.Operations = append(trigger.Operations, storage.TriggerOperation{ + Type: "currentDate", + Field: args[i+1], + }) + i += 2 + + case "--condition": + if i+3 >= len(args) { + return fmt.Errorf("--condition requires ") + } + cond := &storage.TriggerCondition{ + Field: args[i+1], + Operator: args[i+2], + Value: args[i+3], + } + if f, err := strconv.ParseFloat(args[i+3], 64); err == nil { + cond.Value = f + } else if b, err := strconv.ParseBool(args[i+3]); err == nil { + cond.Value = b + } + trigger.Condition = cond + i += 4 + + default: + return fmt.Errorf("unknown option: %s", opt) } } @@ -2237,15 +2358,18 @@ func (r *Repl) handleCreateTrigger(args []string) error { return nil } +// handleDropTrigger удаляет триггер из коллекции. +// +// ИСПРАВЛЕНО: убран обязательный аргумент . func (r *Repl) handleDropTrigger(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") } - if len(args) < 3 { - return fmt.Errorf("usage: drop trigger ") + if len(args) < 2 { + return fmt.Errorf("usage: drop trigger ") } collName := args[0] - triggerName := args[2] + triggerName := args[1] db, err := r.store.GetDatabase(r.currentDB) if err != nil { @@ -2262,8 +2386,6 @@ func (r *Repl) handleDropTrigger(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handleShowTriggers(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") @@ -2301,16 +2423,18 @@ func (r *Repl) handleShowTriggers(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". +// handleEnableTrigger включает триггер. +// +// ИСПРАВЛЕНО: убран обязательный аргумент . func (r *Repl) handleEnableTrigger(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") } - if len(args) < 3 { - return fmt.Errorf("usage: enable trigger ") + if len(args) < 2 { + return fmt.Errorf("usage: enable trigger ") } collName := args[0] - triggerName := args[2] + triggerName := args[1] db, err := r.store.GetDatabase(r.currentDB) if err != nil { @@ -2327,15 +2451,18 @@ func (r *Repl) handleEnableTrigger(args []string) error { return nil } +// handleDisableTrigger отключает триггер. +// +// ИСПРАВЛЕНО: убран обязательный аргумент . func (r *Repl) handleDisableTrigger(args []string) error { if r.currentDB == "" { return fmt.Errorf("no database selected") } - if len(args) < 3 { - return fmt.Errorf("usage: disable trigger ") + if len(args) < 2 { + return fmt.Errorf("usage: disable trigger ") } collName := args[0] - triggerName := args[2] + triggerName := args[1] db, err := r.store.GetDatabase(r.currentDB) if err != nil { @@ -2352,8 +2479,6 @@ func (r *Repl) handleDisableTrigger(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handleTriggerLog(args []string) error { logs := storage.GetTriggerExecutionLog() if len(logs) == 0 { @@ -2418,8 +2543,6 @@ func (r *Repl) handleStatus(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handleNodes(args []string) error { if r.coordinator == nil { return fmt.Errorf("cluster coordinator not available") @@ -2458,14 +2581,7 @@ func (r *Repl) handleNodes(args []string) error { return nil } -// ----------------------------------------------------------------------------- -// Кластер (расширенные команды) -// ----------------------------------------------------------------------------- - // handleClusterPipeline показывает статистику pipeline репликатора. -// -// ИСПРАВЛЕНО: ключи map сортируются перед выводом, чтобы порядок -// строк в таблице был детерминированным. func (r *Repl) handleClusterPipeline(args []string) error { if r.coordinator == nil { return fmt.Errorf("cluster coordinator not available") @@ -2484,7 +2600,16 @@ func (r *Repl) handleClusterPipeline(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". +// handleClusterPipelineOn включает pipeline replication. +func (r *Repl) handleClusterPipelineOn(args []string) error { + return r.handleConfigSet([]string{"performance.enable_pipeline", "true"}) +} + +// handleClusterPipelineOff выключает pipeline replication. +func (r *Repl) handleClusterPipelineOff(args []string) error { + return r.handleConfigSet([]string{"performance.enable_pipeline", "false"}) +} + func (r *Repl) handleClusterReshard(args []string) error { if r.coordinator == nil { return fmt.Errorf("cluster coordinator not available") @@ -2561,8 +2686,6 @@ func (r *Repl) handleMigration(args []string) error { } } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handleMigrationStart(migrator *cluster.CrossDCMigrator, args []string) error { if len(args) < 2 { return fmt.Errorf("usage: migration start [database] [collection]") @@ -2599,8 +2722,6 @@ func (r *Repl) handleMigrationStart(migrator *cluster.CrossDCMigrator, args []st return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". -// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. func (r *Repl) handleMigrationStatus(migrator *cluster.CrossDCMigrator, args []string) error { var taskID string if len(args) > 0 { @@ -2672,7 +2793,6 @@ func (r *Repl) handleMigrationList(migrator *cluster.CrossDCMigrator) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleMigrationPause(migrator *cluster.CrossDCMigrator, taskID string) error { if err := migrator.PauseMigration(taskID); err != nil { return err @@ -2681,7 +2801,6 @@ func (r *Repl) handleMigrationPause(migrator *cluster.CrossDCMigrator, taskID st return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleMigrationResume(migrator *cluster.CrossDCMigrator, taskID string) error { if err := migrator.ResumeMigration(taskID); err != nil { return err @@ -2690,7 +2809,6 @@ func (r *Repl) handleMigrationResume(migrator *cluster.CrossDCMigrator, taskID s return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleMigrationCancel(migrator *cluster.CrossDCMigrator, taskID string) error { if err := migrator.CancelMigration(taskID); err != nil { return err @@ -2754,15 +2872,10 @@ func (r *Repl) handleMigrationQueue(migrator *cluster.CrossDCMigrator) error { return nil } -// ----------------------------------------------------------------------------- -// Миграция схемы -// ----------------------------------------------------------------------------- - // handleMigrateStatus показывает статус миграций схемы. // // ИСПРАВЛЕНО: раньше выводилось `fmt.Sprintf("%v", status)`, где -// status — *migration.MigrationStatus. Пользователь видел адрес или -// структуру с фигурными скобками. Теперь выводятся конкретные поля. +// status — *migration.MigrationStatus. Теперь выводятся поля. func (r *Repl) handleMigrateStatus(args []string) error { if r.coordinator == nil { return fmt.Errorf("cluster coordinator not available") @@ -2883,7 +2996,13 @@ func (r *Repl) handleSagaStatus(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". +// handleSagaList показывает активные (или все) SAGA транзакции. +// +// РАСШИРЕНО: поддерживает флаг --all для показа завершённых/отменённых SAGA. +// +// ИСПРАВЛЕНО (2026-10): тип элементов — *cluster.SagaTransaction, потому +// что r.coordinator.GetSagaManager() возвращает *cluster.SagaManager, +// а его методы работают с локальным типом cluster.SagaTransaction. func (r *Repl) handleSagaList(args []string) error { if r.coordinator == nil { return fmt.Errorf("cluster coordinator not available") @@ -2892,19 +3011,294 @@ func (r *Repl) handleSagaList(args []string) error { if sm == nil { return fmt.Errorf("saga manager not available") } - active := sm.GetActiveSagas() - if len(active) == 0 { - utils.PrintInfo("No active SAGA transactions") - return nil + + showAll := false + for _, a := range args { + if a == "--all" { + showAll = true + } } - t := NewTable("ID", "STATUS", "STEPS", "CURRENT_STEP", "CREATED_AT") - for _, saga := range active { + + var sagas []*cluster.SagaTransaction + if showAll { + sagas = sm.GetAllSagas() + if len(sagas) == 0 { + utils.PrintInfo("No SAGA transactions found") + return nil + } + } else { + sagas = sm.GetActiveSagas() + if len(sagas) == 0 { + utils.PrintInfo("No active SAGA transactions (use --all to show all)") + return nil + } + } + + t := NewTable("ID", "STATUS", "STEPS", "CURRENT_STEP", "CREATED_AT", "UPDATED_AT") + for _, saga := range sagas { + createdAt := "" + if saga.CreatedAt > 0 { + createdAt = time.UnixMilli(saga.CreatedAt).Format("02-01-2006 15:04:05.000") + } + updatedAt := "" + if saga.UpdatedAt > 0 { + updatedAt = time.UnixMilli(saga.UpdatedAt).Format("02-01-2006 15:04:05.000") + } t.AddRow( saga.ID, saga.Status, fmt.Sprintf("%d", len(saga.Steps)), fmt.Sprintf("%d", saga.CurrentStep), - time.UnixMilli(saga.CreatedAt).Format("02-01-2006 15:04:05.000"), + createdAt, + updatedAt, + ) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Динамическая конфигурация (через Raft). +// ----------------------------------------------------------------------------- + +func (r *Repl) handleConfigShow(args []string) error { + cfg := r.effectiveConfig() + if cfg == nil { + return fmt.Errorf("config is not available") + } + + section := "" + if len(args) > 0 { + section = strings.ToLower(args[0]) + } + + switch section { + case "", "all": + t := NewTable("SECTION", "KEY", "VALUE") + addRows := func(sectionName, prefix string, pairs map[string]interface{}) { + keys := make([]string, 0, len(pairs)) + for k := range pairs { + keys = append(keys, k) + } + sort.Strings(keys) + for _, k := range keys { + t.AddRow(sectionName, prefix+"."+k, fmt.Sprintf("%v", pairs[k])) + } + } + addRows("cluster", "cluster", map[string]interface{}{ + "name": cfg.Cluster.Name, + "node_ip": cfg.Cluster.NodeIP, + "node_port": cfg.Cluster.NodePort, + "raft_port": cfg.Cluster.RaftPort, + "region": cfg.Cluster.Region, + }) + addRows("api", "api", map[string]interface{}{ + "port": cfg.API.Port, + }) + addRows("compression", "compression", map[string]interface{}{ + "enabled": cfg.Compression.Enabled, + "algorithm": cfg.Compression.Algorithm, + "level": cfg.Compression.Level, + "min_size": cfg.Compression.MinSize, + }) + addRows("log", "log", map[string]interface{}{ + "log_file": cfg.Log.LogFile, + "log_level": cfg.Log.LogLevel, + }) + addRows("replication", "replication", map[string]interface{}{ + "enabled": cfg.Replication.Enabled, + "sync_replication": cfg.Replication.SyncReplication, + "replication_timeout_ms": cfg.Replication.ReplicationTimeoutMs, + "max_replica_lag_ms": cfg.Replication.MaxReplicaLagMs, + }) + addRows("wal", "wal", map[string]interface{}{ + "enabled": cfg.WAL.Enabled, + "segment_size_mb": cfg.WAL.SegmentSizeMB, + "sync_interval_sec": cfg.WAL.SyncIntervalSec, + "batch_size": cfg.WAL.BatchSize, + }) + addRows("saga", "saga", map[string]interface{}{ + "enabled": cfg.Saga.Enabled, + "coordinator_count": cfg.Saga.CoordinatorCount, + "max_retries": cfg.Saga.MaxRetries, + }) + addRows("backpressure", "backpressure", map[string]interface{}{ + "enabled": cfg.Backpressure.Enabled, + "cpu_threshold": cfg.Backpressure.CPUThreshold, + "memory_threshold": cfg.Backpressure.MemoryThreshold, + }) + addRows("monitoring", "monitoring", map[string]interface{}{ + "enable_metrics": cfg.Monitoring.EnableMetrics, + "metrics_port": cfg.Monitoring.MetricsPort, + "enable_tracing": cfg.Monitoring.EnableTracing, + "trace_sample_rate": cfg.Monitoring.TraceSampleRate, + }) + addRows("performance", "performance", map[string]interface{}{ + "enable_pipeline": cfg.Performance.EnablePipeline, + "batch_size": cfg.Performance.BatchSize, + "max_connections": cfg.Performance.MaxConnections, + "read_from_follower": cfg.Performance.ReadFromFollower, + }) + r.output(t.String()) + + case "cluster": + utils.PrintInfo(fmt.Sprintf("cluster.name = %s", cfg.Cluster.Name)) + utils.PrintInfo(fmt.Sprintf("cluster.node_ip = %s", cfg.Cluster.NodeIP)) + utils.PrintInfo(fmt.Sprintf("cluster.node_port = %d", cfg.Cluster.NodePort)) + utils.PrintInfo(fmt.Sprintf("cluster.raft_port = %d", cfg.Cluster.RaftPort)) + utils.PrintInfo(fmt.Sprintf("cluster.region = %s", cfg.Cluster.Region)) + + case "compression": + t := NewTable("KEY", "VALUE") + t.AddRow("compression.enabled", fmt.Sprintf("%v", cfg.Compression.Enabled)) + t.AddRow("compression.algorithm", cfg.Compression.Algorithm) + t.AddRow("compression.level", fmt.Sprintf("%d", cfg.Compression.Level)) + t.AddRow("compression.min_size", fmt.Sprintf("%d", cfg.Compression.MinSize)) + r.output(t.String()) + + case "replication": + t := NewTable("KEY", "VALUE") + t.AddRow("replication.enabled", fmt.Sprintf("%v", cfg.Replication.Enabled)) + t.AddRow("replication.sync_replication", fmt.Sprintf("%v", cfg.Replication.SyncReplication)) + t.AddRow("replication.replication_timeout_ms", fmt.Sprintf("%d", cfg.Replication.ReplicationTimeoutMs)) + t.AddRow("replication.max_replica_lag_ms", fmt.Sprintf("%d", cfg.Replication.MaxReplicaLagMs)) + r.output(t.String()) + + case "wal": + t := NewTable("KEY", "VALUE") + t.AddRow("wal.enabled", fmt.Sprintf("%v", cfg.WAL.Enabled)) + t.AddRow("wal.segment_size_mb", fmt.Sprintf("%d", cfg.WAL.SegmentSizeMB)) + t.AddRow("wal.sync_interval_sec", fmt.Sprintf("%d", cfg.WAL.SyncIntervalSec)) + t.AddRow("wal.batch_size", fmt.Sprintf("%d", cfg.WAL.BatchSize)) + t.AddRow("wal.recovery_workers", fmt.Sprintf("%d", cfg.WAL.RecoveryWorkers)) + r.output(t.String()) + + case "saga": + t := NewTable("KEY", "VALUE") + t.AddRow("saga.enabled", fmt.Sprintf("%v", cfg.Saga.Enabled)) + t.AddRow("saga.coordinator_count", fmt.Sprintf("%d", cfg.Saga.CoordinatorCount)) + t.AddRow("saga.max_retries", fmt.Sprintf("%d", cfg.Saga.MaxRetries)) + t.AddRow("saga.saga_timeout_sec", fmt.Sprintf("%d", cfg.Saga.SagaTimeoutSec)) + r.output(t.String()) + + case "monitoring": + t := NewTable("KEY", "VALUE") + t.AddRow("monitoring.enable_metrics", fmt.Sprintf("%v", cfg.Monitoring.EnableMetrics)) + t.AddRow("monitoring.metrics_port", fmt.Sprintf("%d", cfg.Monitoring.MetricsPort)) + t.AddRow("monitoring.enable_tracing", fmt.Sprintf("%v", cfg.Monitoring.EnableTracing)) + t.AddRow("monitoring.trace_sample_rate", fmt.Sprintf("%v", cfg.Monitoring.TraceSampleRate)) + r.output(t.String()) + + case "performance": + t := NewTable("KEY", "VALUE") + t.AddRow("performance.enable_pipeline", fmt.Sprintf("%v", cfg.Performance.EnablePipeline)) + t.AddRow("performance.batch_size", fmt.Sprintf("%d", cfg.Performance.BatchSize)) + t.AddRow("performance.max_connections", fmt.Sprintf("%d", cfg.Performance.MaxConnections)) + t.AddRow("performance.read_from_follower", fmt.Sprintf("%v", cfg.Performance.ReadFromFollower)) + r.output(t.String()) + + case "backpressure": + t := NewTable("KEY", "VALUE") + t.AddRow("backpressure.enabled", fmt.Sprintf("%v", cfg.Backpressure.Enabled)) + t.AddRow("backpressure.cpu_threshold", fmt.Sprintf("%v", cfg.Backpressure.CPUThreshold)) + t.AddRow("backpressure.memory_threshold", fmt.Sprintf("%v", cfg.Backpressure.MemoryThreshold)) + t.AddRow("backpressure.queue_size_threshold", fmt.Sprintf("%d", cfg.Backpressure.QueueSizeThreshold)) + t.AddRow("backpressure.connection_threshold", fmt.Sprintf("%d", cfg.Backpressure.ConnectionThreshold)) + r.output(t.String()) + + default: + return fmt.Errorf("unknown config section '%s' (supported: cluster, compression, replication, wal, saga, monitoring, performance, backpressure)", section) + } + return nil +} + +func (r *Repl) handleConfigGet(args []string) error { + if len(args) < 1 { + return fmt.Errorf("usage: config get (e.g. compression.algorithm)") + } + path := strings.Join(args, " ") + + cfg := r.effectiveConfig() + if cfg == nil { + return fmt.Errorf("config is not available") + } + + value, err := getConfigValue(cfg, path) + if err != nil { + return err + } + utils.PrintInfo(fmt.Sprintf("%s = %v", path, value)) + return nil +} + +func (r *Repl) handleConfigSet(args []string) error { + if !r.authenticated || r.currentRole != "admin" { + return fmt.Errorf("permission denied: admin access required") + } + if len(args) < 2 { + return fmt.Errorf("usage: config set (e.g. compression.algorithm brotli)") + } + + path := args[0] + valueStr := strings.Join(args[1:], " ") + + cfg := r.effectiveConfig() + if cfg == nil { + return fmt.Errorf("config is not available") + } + currentValue, err := getConfigValue(cfg, path) + if err != nil { + return err + } + + value, err := coerceStringToType(valueStr, currentValue) + if err != nil { + return fmt.Errorf("cannot parse value for %s: %v", path, err) + } + + change := map[string]interface{}{path: value} + + if r.dynamicConfigManager != nil { + err := r.dynamicConfigManager.ApplyChange(change, + fmt.Sprintf("REPL: set %s", path), r.currentUser) + if err == nil { + r.applyLocalConfigChange(path, value) + utils.PrintSuccess(fmt.Sprintf("config %s = %v (applied via Raft)", path, value)) + return nil + } + utils.PrintWarning(fmt.Sprintf("Raft apply failed (%v), falling back to local config.toml", err)) + } + + r.applyLocalConfigChange(path, value) + if r.configPath != "" { + if err := r.persistConfigFile(); err != nil { + utils.PrintWarning(fmt.Sprintf("config applied locally, but saving to %s failed: %v", r.configPath, err)) + } else { + utils.PrintSuccess(fmt.Sprintf("config %s = %v (saved to %s)", path, value, r.configPath)) + return nil + } + } + utils.PrintSuccess(fmt.Sprintf("config %s = %v (in-memory only, no configPath set)", path, value)) + return nil +} + +func (r *Repl) handleConfigHistory(args []string) error { + if r.dynamicConfigManager == nil { + utils.PrintInfo("Dynamic config manager is not available (single-node mode or Raft disabled)") + return nil + } + history := r.dynamicConfigManager.GetChangeHistory() + if len(history) == 0 { + utils.PrintInfo("No config changes recorded") + return nil + } + t := NewTable("VERSION", "TIMESTAMP", "CHANGED_BY", "DESCRIPTION") + for _, ch := range history { + t.AddRow( + fmt.Sprintf("%d", ch.Version), + time.UnixMilli(ch.Timestamp).Format("02-01-2006 15:04:05.000"), + ch.ChangedBy, + ch.Description, ) } r.output(t.String()) @@ -2938,7 +3332,7 @@ func (r *Repl) handleHelp(args []string) error { {"insert ", "Insert a new document (JSON format: key=value,key2=value2)"}, {"find ", "Find a document by its ID"}, {"findbyindex ", "Find documents using an index"}, - {"findbytime ", "Find documents by time range"}, + {"findbytime ", "Find documents by time range (ISO 8601, e.g. 2026-01-01T12:00:00)"}, {"update ...", "Update fields of an existing document"}, {"delete ", "Delete a document (soft delete if enabled)"}, {"permanent delete ", "Permanently delete a soft-deleted document"}, @@ -2950,7 +3344,7 @@ func (r *Repl) handleHelp(args []string) error { {"show timestamps ", "Show timestamps for a document"}, {"stats timestamps ", "Show timestamp statistics for a collection"}, {"audit log", "Show audit log"}, - {"audit filter ", "Filter audit log by type and operation"}, + {"audit filter [] [from=] [to=]", "Filter audit log by type, operation and time"}, }, "Index Management": { {"create index [unique]", "Create a new index on specified fields"}, @@ -2964,9 +3358,6 @@ func (r *Repl) handleHelp(args []string) error { {"add max ", "Add a maximum value constraint for numeric fields"}, {"add enum ", "Add allowed values constraint (enum)"}, }, - // ИСПРАВЛЕНО (2026-10): раздел "Transactions" заменён на "Batch Operations". - // Старые команды (begin transaction, commit, rollback, show transactions) - // оставлены как алиасы для обратной совместимости. "Batch Operations": { {"batch begin [collection]", "Begin a new batch (atomically applied group of operations)"}, {"batch commit", "Commit the current batch (operations applied atomically via WAL)"}, @@ -2979,10 +3370,15 @@ func (r *Repl) handleHelp(args []string) error { }, "Triggers (MongoDB-like)": { {"create trigger [options]", "Create a trigger on collection events"}, - {"drop trigger ", "Remove a trigger from collection"}, + {" --description ", " Set trigger description"}, + {" --set ", " Set field value when trigger fires"}, + {" --inc ", " Increment numeric field"}, + {" --currentDate ", " Set field to current timestamp"}, + {" --condition ", " Only fire if condition matches (op: eq, ne, gt, lt, gte, lte, in, nin, exists, regex)"}, + {"drop trigger ", "Remove a trigger from collection"}, {"show triggers ", "List all triggers on a collection"}, - {"enable trigger ", "Enable a disabled trigger"}, - {"disable trigger ", "Disable a trigger without removing it"}, + {"enable trigger ", "Enable a disabled trigger"}, + {"disable trigger ", "Disable a trigger without removing it"}, {"trigger log", "Show trigger execution history"}, }, "Plugins": { @@ -3000,20 +3396,36 @@ func (r *Repl) handleHelp(args []string) error { "Compression": { {"compression stats", "Show compression statistics for current database"}, {"compression config", "Display current compression settings"}, - {"compress collection ", "Manually compress all documents in a collection"}, + {"compression config set ", "Change setting (enabled, algorithm, level, min_size)"}, + {"compress collection ", "Compress all documents in a collection (frees memory)"}, {"doc compression ", "Show compression info for a specific document"}, }, "Access Control": { {"acl login ", "Authenticate with username and password"}, {"acl logout", "Logout current user session"}, {"acl grant ", "Grant permissions (r=read,w=write,d=delete,a=admin)"}, + {"acl revoke ", "Revoke permissions (r=read,w=write,d=delete,a=admin)"}, {"acl users", "List all users"}, {"acl roles", "List all roles"}, }, + "Users & Roles (admin only)": { + {"acl user create [roles]", "Create a new user (roles separated by commas)"}, + {"acl user delete ", "Delete a user"}, + {"acl role create ", "Create a new role"}, + {"acl role delete ", "Delete a role"}, + }, + "Dynamic Config (admin only)": { + {"config show [section]", "Show current configuration or a specific section"}, + {"config get ", "Get a single config value (e.g. compression.algorithm)"}, + {"config set ", "Set a config value (persisted via Raft or config.toml)"}, + {"config history", "Show dynamic config change history"}, + }, "Cluster": { {"status", "Show current cluster status and role (leader/follower)"}, {"nodes", "List all nodes in the cluster"}, {"cluster pipeline", "Show pipeline replicator statistics"}, + {"cluster pipeline on", "Enable pipeline replication (config set)"}, + {"cluster pipeline off", "Disable pipeline replication (config set)"}, {"cluster reshard", "Trigger manual resharding"}, }, "Migration": { @@ -3039,10 +3451,12 @@ func (r *Repl) handleHelp(args []string) error { }, "SAGA": { {"saga status", "Show SAGA orchestrator status"}, - {"saga list", "List active SAGA transactions"}, + {"saga list [--all]", "List active (or all) SAGA transactions"}, }, "Pager & Input": { {"pager on|off", "Enable/disable pager"}, + {"pager size ", "Set pager threshold (lines) — alias of 'pager threshold'"}, + {"pager threshold ", "Set pager threshold (lines)"}, {"pager", "Show pager status"}, {"PgUp/PgDown", "Navigate command history (like in bash)"}, {"Tab", "Autocomplete command / collection / field"}, @@ -3057,7 +3471,6 @@ func (r *Repl) handleHelp(args []string) error { }, } - // Сортируем категории для детерминированного вывода. categoryNames := make([]string, 0, len(categories)) for name := range categories { categoryNames = append(categoryNames, name) @@ -3091,9 +3504,7 @@ func (r *Repl) handleClear(args []string) error { return nil } -// ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (r *Repl) handleQuit(args []string) error { - // Снимаем активный batch при выходе, если есть. if r.currentBatch != nil { if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(r.currentBatch) @@ -3108,3 +3519,541 @@ func (r *Repl) handleQuit(args []string) error { os.Exit(0) return nil } + +// ============================================================================= +// ВСПОМОГАТЕЛЬНЫЕ ФУНКЦИИ ДЛЯ CONFIG +// ============================================================================= + +// effectiveConfig возвращает конфигурацию, которую видит REPL. +func (r *Repl) effectiveConfig() *config.Config { + if r.dynamicConfigManager != nil { + if cfg := r.dynamicConfigManager.GetConfig(); cfg != nil { + return cfg + } + } + return r.config +} + +// applyLocalConfigChange применяет изменение только к in-memory r.config. +func (r *Repl) applyLocalConfigChange(path string, value interface{}) { + if r.config == nil { + return + } + _ = setConfigFieldLocal(r.config, path, value) +} + +// persistConfigFile сохраняет текущий r.config в r.configPath. +// Делегирует в config.SaveConfig. +func (r *Repl) persistConfigFile() error { + if r.configPath == "" { + return fmt.Errorf("configPath is not set") + } + return config.SaveConfig(r.configPath, r.config) +} + +// getConfigValue читает значение по пути из конфигурации. +func getConfigValue(cfg *config.Config, path string) (interface{}, error) { + switch path { + case "cluster.name": + return cfg.Cluster.Name, nil + case "cluster.node_ip": + return cfg.Cluster.NodeIP, nil + case "cluster.node_port": + return cfg.Cluster.NodePort, nil + case "cluster.raft_port": + return cfg.Cluster.RaftPort, nil + case "cluster.heartbeat_timeout_ms": + return cfg.Cluster.HeartbeatTimeoutMs, nil + case "cluster.election_timeout_ms": + return cfg.Cluster.ElectionTimeoutMs, nil + case "cluster.commit_timeout_ms": + return cfg.Cluster.CommitTimeoutMs, nil + case "cluster.snapshot_interval_min": + return cfg.Cluster.SnapshotIntervalMin, nil + case "cluster.snapshot_threshold": + return cfg.Cluster.SnapshotThreshold, nil + case "cluster.recovery_timeout_sec": + return cfg.Cluster.RecoveryTimeoutSec, nil + case "cluster.split_brain_prevention": + return cfg.Cluster.SplitBrainPrevention, nil + case "cluster.region": + return cfg.Cluster.Region, nil + case "cluster.priority_zone": + return cfg.Cluster.PriorityZone, nil + + case "storage.page_size_mb": + return cfg.Storage.PageSizeMB, nil + case "storage.max_collections": + return cfg.Storage.MaxCollections, nil + case "storage.max_documents_per_collection": + return cfg.Storage.MaxDocumentsPerCollection, nil + case "storage.default_engine": + return cfg.Storage.DefaultEngine, nil + case "storage.enable_custom_engines": + return cfg.Storage.EnableCustomEngines, nil + + case "replication.enabled": + return cfg.Replication.Enabled, nil + case "replication.sync_replication": + return cfg.Replication.SyncReplication, nil + case "replication.replication_timeout_ms": + return cfg.Replication.ReplicationTimeoutMs, nil + case "replication.max_replica_lag_ms": + return cfg.Replication.MaxReplicaLagMs, nil + + case "wal.enabled": + return cfg.WAL.Enabled, nil + case "wal.segment_size_mb": + return cfg.WAL.SegmentSizeMB, nil + case "wal.sync_interval_sec": + return cfg.WAL.SyncIntervalSec, nil + case "wal.batch_size": + return cfg.WAL.BatchSize, nil + case "wal.async_recovery": + return cfg.WAL.AsyncRecovery, nil + case "wal.recovery_workers": + return cfg.WAL.RecoveryWorkers, nil + case "wal.async_recovery_workers": + return cfg.WAL.AsyncRecoveryWorkers, nil + case "wal.async_recovery_buffer": + return cfg.WAL.AsyncRecoveryBuffer, nil + + case "mvcc.max_versions_per_doc": + return cfg.MVCC.MaxVersionsPerDoc, nil + case "mvcc.visibility_map_size": + return cfg.MVCC.VisibilityMapSize, nil + case "mvcc.prune_interval_min": + return cfg.MVCC.PruneIntervalMin, nil + case "mvcc.retention_days": + return cfg.MVCC.RetentionDays, nil + case "mvcc.read_cache_size": + return cfg.MVCC.ReadCacheSize, nil + case "mvcc.read_cache_ttl_sec": + return cfg.MVCC.ReadCacheTTLSec, nil + + case "saga.enabled": + return cfg.Saga.Enabled, nil + case "saga.coordinator_count": + return cfg.Saga.CoordinatorCount, nil + case "saga.max_retries": + return cfg.Saga.MaxRetries, nil + case "saga.saga_timeout_sec": + return cfg.Saga.SagaTimeoutSec, nil + case "saga.retry_backoff_ms": + return cfg.Saga.RetryBackoffMs, nil + case "saga.stuck_check_interval_sec": + return cfg.Saga.StuckCheckIntervalSec, nil + + case "backpressure.enabled": + return cfg.Backpressure.Enabled, nil + case "backpressure.cpu_threshold": + return cfg.Backpressure.CPUThreshold, nil + case "backpressure.memory_threshold": + return cfg.Backpressure.MemoryThreshold, nil + case "backpressure.queue_size_threshold": + return cfg.Backpressure.QueueSizeThreshold, nil + case "backpressure.connection_threshold": + return cfg.Backpressure.ConnectionThreshold, nil + + case "api.port": + return cfg.API.Port, nil + + case "log.log_level": + return cfg.Log.LogLevel, nil + case "log.log_file": + return cfg.Log.LogFile, nil + + case "plugins.enabled": + return cfg.Plugins.Enabled, nil + case "plugins.max_cpu_time_ms": + return cfg.Plugins.MaxCPUTimeMs, nil + case "plugins.max_memory_mb": + return cfg.Plugins.MaxMemoryMB, nil + case "plugins.hot_reload_interval_sec": + return cfg.Plugins.HotReloadIntervalSec, nil + case "plugins.max_lua_states": + return cfg.Plugins.MaxLuaStates, nil + + case "autoscaling.enabled": + return cfg.Autoscaling.Enabled, nil + case "autoscaling.min_nodes": + return cfg.Autoscaling.MinNodes, nil + case "autoscaling.max_nodes": + return cfg.Autoscaling.MaxNodes, nil + case "autoscaling.scale_up_threshold": + return cfg.Autoscaling.ScaleUpThreshold, nil + case "autoscaling.scale_down_threshold": + return cfg.Autoscaling.ScaleDownThreshold, nil + case "autoscaling.evaluation_interval_sec": + return cfg.Autoscaling.EvaluationIntervalSec, nil + case "autoscaling.scale_up_cooldown_sec": + return cfg.Autoscaling.ScaleUpCooldownSec, nil + case "autoscaling.scale_down_cooldown_sec": + return cfg.Autoscaling.ScaleDownCooldownSec, nil + + case "monitoring.enable_metrics": + return cfg.Monitoring.EnableMetrics, nil + case "monitoring.metrics_port": + return cfg.Monitoring.MetricsPort, nil + case "monitoring.enable_tracing": + return cfg.Monitoring.EnableTracing, nil + case "monitoring.trace_sample_rate": + return cfg.Monitoring.TraceSampleRate, nil + + case "transactions.enabled": + return cfg.Transactions.Enabled, nil + case "transactions.default_timeout_sec": + return cfg.Transactions.DefaultTimeoutSec, nil + case "transactions.deadlock_check_interval_sec": + return cfg.Transactions.DeadlockCheckIntervalSec, nil + case "transactions.max_savepoints_per_tx": + return cfg.Transactions.MaxSavepointsPerTx, nil + + case "compression.enabled": + return cfg.Compression.Enabled, nil + case "compression.algorithm": + return cfg.Compression.Algorithm, nil + case "compression.level": + return cfg.Compression.Level, nil + case "compression.min_size": + return cfg.Compression.MinSize, nil + + case "performance.enable_pipeline": + return cfg.Performance.EnablePipeline, nil + case "performance.batch_size": + return cfg.Performance.BatchSize, nil + case "performance.max_connections": + return cfg.Performance.MaxConnections, nil + case "performance.read_from_follower": + return cfg.Performance.ReadFromFollower, nil + } + return nil, fmt.Errorf("unknown config path: %s", path) +} + +// coerceStringToType приводит строку из REPL к типу существующего значения. +func coerceStringToType(valueStr string, currentValue interface{}) (interface{}, error) { + switch currentValue.(type) { + case bool: + b, err := strconv.ParseBool(valueStr) + if err != nil { + return nil, fmt.Errorf("expected bool, got '%s'", valueStr) + } + return b, nil + case int: + i, err := strconv.ParseInt(valueStr, 10, 64) + if err != nil { + return nil, fmt.Errorf("expected integer, got '%s'", valueStr) + } + return int(i), nil + case int8, int16, int32: + i, err := strconv.ParseInt(valueStr, 10, 64) + if err != nil { + return nil, fmt.Errorf("expected integer, got '%s'", valueStr) + } + return i, nil + case int64: + i, err := strconv.ParseInt(valueStr, 10, 64) + if err != nil { + return nil, fmt.Errorf("expected integer, got '%s'", valueStr) + } + return i, nil + case uint, uint8, uint16, uint32: + i, err := strconv.ParseUint(valueStr, 10, 64) + if err != nil { + return nil, fmt.Errorf("expected unsigned integer, got '%s'", valueStr) + } + return i, nil + case uint64: + i, err := strconv.ParseUint(valueStr, 10, 64) + if err != nil { + return nil, fmt.Errorf("expected unsigned integer, got '%s'", valueStr) + } + return i, nil + case float32, float64: + f, err := strconv.ParseFloat(valueStr, 64) + if err != nil { + return nil, fmt.Errorf("expected number, got '%s'", valueStr) + } + return f, nil + case string: + return valueStr, nil + } + + if b, err := strconv.ParseBool(valueStr); err == nil { + return b, nil + } + if i, err := strconv.ParseInt(valueStr, 10, 64); err == nil { + return i, nil + } + if f, err := strconv.ParseFloat(valueStr, 64); err == nil { + return f, nil + } + return valueStr, nil +} + +// setConfigFieldLocal применяет изменение к config локально. +// Это дубликат setConfigField из dynamic_config.go. +func setConfigFieldLocal(cfg *config.Config, path string, value interface{}) error { + switch path { + case "compression.enabled": + if v, ok := value.(bool); ok { + cfg.Compression.Enabled = v + return nil + } + case "compression.algorithm": + if v, ok := value.(string); ok { + cfg.Compression.Algorithm = v + return nil + } + case "compression.level": + if v, ok := value.(int64); ok { + cfg.Compression.Level = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Compression.Level = v + return nil + } + case "compression.min_size": + if v, ok := value.(int64); ok { + cfg.Compression.MinSize = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Compression.MinSize = v + return nil + } + + case "cluster.name": + if v, ok := value.(string); ok { + cfg.Cluster.Name = v + return nil + } + case "cluster.node_ip": + if v, ok := value.(string); ok { + cfg.Cluster.NodeIP = v + return nil + } + case "cluster.node_port": + if v, ok := value.(int64); ok { + cfg.Cluster.NodePort = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Cluster.NodePort = v + return nil + } + case "cluster.raft_port": + if v, ok := value.(int64); ok { + cfg.Cluster.RaftPort = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Cluster.RaftPort = v + return nil + } + case "cluster.region": + if v, ok := value.(string); ok { + cfg.Cluster.Region = v + return nil + } + + case "performance.enable_pipeline": + if v, ok := value.(bool); ok { + cfg.Performance.EnablePipeline = v + return nil + } + case "performance.batch_size": + if v, ok := value.(int64); ok { + cfg.Performance.BatchSize = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Performance.BatchSize = v + return nil + } + case "performance.max_connections": + if v, ok := value.(int64); ok { + cfg.Performance.MaxConnections = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Performance.MaxConnections = v + return nil + } + case "performance.read_from_follower": + if v, ok := value.(bool); ok { + cfg.Performance.ReadFromFollower = v + return nil + } + + case "monitoring.enable_metrics": + if v, ok := value.(bool); ok { + cfg.Monitoring.EnableMetrics = v + return nil + } + case "monitoring.metrics_port": + if v, ok := value.(int64); ok { + cfg.Monitoring.MetricsPort = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Monitoring.MetricsPort = v + return nil + } + case "monitoring.enable_tracing": + if v, ok := value.(bool); ok { + cfg.Monitoring.EnableTracing = v + return nil + } + case "monitoring.trace_sample_rate": + if v, ok := value.(float64); ok { + cfg.Monitoring.TraceSampleRate = v + return nil + } + + case "log.log_level": + if v, ok := value.(string); ok { + cfg.Log.LogLevel = v + return nil + } + case "log.log_file": + if v, ok := value.(string); ok { + cfg.Log.LogFile = v + return nil + } + + case "wal.enabled": + if v, ok := value.(bool); ok { + cfg.WAL.Enabled = v + return nil + } + case "wal.segment_size_mb": + if v, ok := value.(int64); ok { + cfg.WAL.SegmentSizeMB = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.WAL.SegmentSizeMB = v + return nil + } + case "wal.sync_interval_sec": + if v, ok := value.(int64); ok { + cfg.WAL.SyncIntervalSec = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.WAL.SyncIntervalSec = v + return nil + } + case "wal.batch_size": + if v, ok := value.(int64); ok { + cfg.WAL.BatchSize = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.WAL.BatchSize = v + return nil + } + + case "saga.enabled": + if v, ok := value.(bool); ok { + cfg.Saga.Enabled = v + return nil + } + case "saga.coordinator_count": + if v, ok := value.(int64); ok { + cfg.Saga.CoordinatorCount = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Saga.CoordinatorCount = v + return nil + } + case "saga.max_retries": + if v, ok := value.(int64); ok { + cfg.Saga.MaxRetries = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Saga.MaxRetries = v + return nil + } + case "saga.saga_timeout_sec": + if v, ok := value.(int64); ok { + cfg.Saga.SagaTimeoutSec = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Saga.SagaTimeoutSec = v + return nil + } + case "saga.retry_backoff_ms": + if v, ok := value.(int64); ok { + cfg.Saga.RetryBackoffMs = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Saga.RetryBackoffMs = v + return nil + } + case "saga.stuck_check_interval_sec": + if v, ok := value.(int64); ok { + cfg.Saga.StuckCheckIntervalSec = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Saga.StuckCheckIntervalSec = v + return nil + } + + case "api.port": + if v, ok := value.(int64); ok { + cfg.API.Port = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.API.Port = v + return nil + } + + case "backpressure.enabled": + if v, ok := value.(bool); ok { + cfg.Backpressure.Enabled = v + return nil + } + case "backpressure.cpu_threshold": + if v, ok := value.(float64); ok { + cfg.Backpressure.CPUThreshold = v + return nil + } + case "backpressure.memory_threshold": + if v, ok := value.(float64); ok { + cfg.Backpressure.MemoryThreshold = v + return nil + } + case "backpressure.queue_size_threshold": + if v, ok := value.(int64); ok { + cfg.Backpressure.QueueSizeThreshold = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Backpressure.QueueSizeThreshold = v + return nil + } + case "backpressure.connection_threshold": + if v, ok := value.(int64); ok { + cfg.Backpressure.ConnectionThreshold = int(v) + return nil + } + if v, ok := value.(int); ok { + cfg.Backpressure.ConnectionThreshold = v + return nil + } + } + return fmt.Errorf("local update for path '%s' is not implemented (use Raft for full support)", path) +} + +// Проверка, что filepath используется (для будущего расширения +// fallback-записи config.toml в произвольную директорию). +var _ = filepath.Join \ No newline at end of file