diff --git a/internal/repl/repl.go b/internal/repl/repl.go new file mode 100644 index 0000000..fe3efe1 --- /dev/null +++ b/internal/repl/repl.go @@ -0,0 +1,3038 @@ +/* + * Copyright 2026 Safronov Grigorii + * + * Licensed under the CDDL, Version 1.0 (the "License"); + * you may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * https://opensource.org/licenses/CDDL-1.0 + */ + +// ============================================================================= +// Пакет: 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 сессии. +// Старые имена команд сохранены как алиасы для совместимости с привычками +// пользователей и существующими скриптами. +// +// ИЗМЕНЕНО (2026-10-03): формат временных меток изменён с "2006-01-02 15:04:05.000" +// на "02-012006 15:04:05.000" для совместимости с OpenIndiana. +// +// ИЗМЕНЕНО (2026-10-04): формат временных меток УНИФИЦИРОВАН — везде используется +// "02-012006 15:04:05.000" с миллисекундами. Ранее в некоторых табличных выводах +// (индексы, плагины, триггеры, узлы кластера, статус миграции) формат был без +// миллисекунд — теперь это исправлено для консистентности. +// ============================================================================= + +package repl + +import ( + "errors" + "fmt" + "io" + "os" + "strings" + "time" + + "futriis/internal/acl" + "futriis/internal/cluster" + "futriis/internal/compression" + "futriis/internal/config" + "futriis/internal/log" + "futriis/internal/plugin" + "futriis/internal/storage" + "futriis/pkg/utils" + + "github.com/chzyer/readline" + "github.com/fatih/color" +) + +// ============================================================================= +// ТИПЫ ДАННЫХ +// ============================================================================= + +// Repl представляет основную структуру REPL. +type Repl struct { + store *storage.Storage + coordinator *cluster.RaftCoordinator + logger *log.Logger + config *config.Config + aclManager *acl.ACLManager + pluginManager *plugin.PluginManager + + // Терминальный ввод/вывод. + rl *readline.Instance + interrupt bool + + currentDB string + currentUser string + currentRole string + authenticated bool + sessionID string + + // ДОБАВЛЕНО (2026-10): текущий активный batch (замена транзакции). + // Создаётся командой "begin transaction"/"batch begin", коммитится + // командой "commit"/"batch commit", снимается командой "rollback"/"batch abort". + currentBatch *storage.Batch + currentBatchID uint64 + + commands map[string]*Command + history []string + historyPos int + historyFile *History + + pager *Pager + completer *Completer +} + +// Command представляет отдельную команду REPL. +type Command struct { + Name string + Description string + Handler func(args []string) error +} + +// ============================================================================= +// КОНСТРУКТОРЫ +// ============================================================================= + +// NewRepl создаёт новый экземпляр REPL с использованием chzyer/readline. +func NewRepl( + store *storage.Storage, + coordinator *cluster.RaftCoordinator, + logger *log.Logger, + cfg *config.Config, + aclManager *acl.ACLManager, + pluginManager *plugin.PluginManager, +) (*Repl, error) { + + r := &Repl{ + store: store, + coordinator: coordinator, + logger: logger, + config: cfg, + aclManager: aclManager, + pluginManager: pluginManager, + currentDB: "", + currentUser: "", + currentRole: "anonymous", + authenticated: false, + sessionID: "", + commands: make(map[string]*Command), + history: make([]string, 0, cfg.Repl.HistorySize), + historyPos: -1, + 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(), + HistoryLimit: cfg.Repl.HistorySize, + HistorySearchFold: true, + AutoComplete: r.completer.ReadlineCompleter(), + InterruptPrompt: "^C", + EOFPrompt: "exit", + DisableAutoSaveHistory: true, // сохраняем сами через historyFile + VimMode: false, // emacs-режим как в bash + FuncFilterInputRune: func(rn rune) (rune, bool) { + // Блокируем Ctrl+Z, чтобы случайно не свернуть процесс. + if rn == readline.CharCtrlZ { + return rn, false + } + return rn, true + }, + }) + if err != nil { + return nil, fmt.Errorf("failed to initialize readline: %w", err) + } + + r.rl = rl + return r, nil +} + +// ============================================================================= +// РЕГИСТРАЦИЯ КОМАНД +// ============================================================================= + +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} + r.commands["findbytime"] = &Command{Name: "findbytime", Description: "Find documents by time range (created_at)", Handler: r.handleFindByTime} + r.commands["update"] = &Command{Name: "update", Description: "Update a document", Handler: r.handleUpdate} + r.commands["delete"] = &Command{Name: "delete", Description: "Delete a document (soft delete if enabled)", Handler: r.handleDelete} + r.commands["permanent delete"] = &Command{Name: "permanent delete", Description: "Permanently delete a soft-deleted document", Handler: r.handlePermanentDelete} + r.commands["restore"] = &Command{Name: "restore", Description: "Restore a soft-deleted document", Handler: r.handleRestore} + r.commands["show deleted"] = &Command{Name: "show deleted", Description: "Show soft-deleted documents in a collection", Handler: r.handleShowDeleted} + r.commands["count"] = &Command{Name: "count", Description: "Count documents in a collection", Handler: r.handleCount} + 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} + r.commands["enable trigger"] = &Command{Name: "enable trigger", Description: "Enable a trigger", Handler: r.handleEnableTrigger} + 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-команды. + 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} + r.commands["plugin start"] = &Command{Name: "plugin start", Description: "Start a plugin", Handler: r.handlePluginStart} + 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 + // ------------------------------------------------------------------------- + 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 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["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["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 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) + // ------------------------------------------------------------------------- + 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 + // ------------------------------------------------------------------------- + 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["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} + r.commands["exit"] = &Command{Name: "exit", Description: "Exit the REPL", Handler: r.handleQuit} +} + +// ============================================================================= +// ОСНОВНОЙ ЦИКЛ REPL +// ============================================================================= + +// Run запускает основной цикл REPL с использованием readline. +func (r *Repl) Run() error { + defer r.rl.Close() + + utils.Println("") + utils.PrintInfo("Type 'help' for available commands") + utils.PrintInfo("Multiline: end line with '\\' to continue, or use '{' ... '}'") + utils.PrintInfo("PgUp/PgDown/↑/↓ — history, Tab — autocomplete, Ctrl+D — exit") + 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 + } + } + + line = strings.TrimSpace(line) + if line == "" { + continue + } + + // Многострочный ввод. + full, err := r.collectMultiline(line) + if err != nil { + if errors.Is(err, io.EOF) { + return nil + } + if errors.Is(err, readline.ErrInterrupt) { + utils.Println("^C") + continue + } + return err + } + if full == "" { + continue + } + + // Сохраняем в историю (readline + наш файл). + r.addToHistory(full) + + if err := r.executeCommand(full); err != nil { + utils.PrintError(err.Error()) + if r.logger != nil { + r.logger.Error("REPL command error: " + err.Error()) + } + } + } +} + +// collectMultiline продолжает ввод, если строка не закончена. +func (r *Repl) collectMultiline(first string) (string, error) { + acc := first + depthCurly, depthSquare, depthParen := countBrackets(acc) + + for { + cont := strings.HasSuffix(acc, "\\") || + depthCurly > 0 || depthSquare > 0 || depthParen > 0 + + if !cont { + break + } + + r.rl.SetPrompt("... ") + next, err := r.rl.Readline() + if err != nil { + return "", err + } + r.rl.SetPrompt(r.buildPrompt()) + + if strings.HasSuffix(acc, "\\") { + acc = strings.TrimSuffix(acc, "\\") + } + + next = strings.TrimSpace(next) + if next == "" && depthCurly <= 0 && depthSquare <= 0 && depthParen <= 0 { + break + } + + acc += " " + next + depthCurly, depthSquare, depthParen = countBrackets(acc) + } + + return strings.TrimSpace(acc), nil +} + +// countBrackets считает баланс {}, [], () с учётом строк. +func countBrackets(s string) (curly, square, paren int) { + inString := false + var quote rune + escaped := false + + for _, ch := range s { + if escaped { + escaped = false + continue + } + if ch == '\\' && inString { + escaped = true + continue + } + if inString { + if ch == quote { + inString = false + } + continue + } + switch ch { + case '"', '\'': + inString = true + quote = ch + case '{': + curly++ + case '}': + curly-- + case '[': + square++ + case ']': + square-- + case '(': + paren++ + case ')': + paren-- + } + } + return +} + +// ============================================================================= +// ВСПОМОГАТЕЛЬНЫЕ МЕТОДЫ +// ============================================================================= + +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 +} + +func (r *Repl) executeCommand(input string) error { + parts := strings.Fields(input) + if len(parts) == 0 { + return nil + } + + 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 + } + + for cmdName, cmd := range r.commands { + if strings.HasPrefix(input, cmdName) { + args := strings.TrimPrefix(input, cmdName) + args = strings.TrimSpace(args) + argList := strings.Fields(args) + return cmd.Handler(argList) + } + } + + return fmt.Errorf("unknown command: %s", parts[0]) +} + +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:] + } + r.history = append(r.history, cmd) + r.historyPos = len(r.history) + + _ = r.historyFile.Add(cmd) + _ = r.historyFile.Save() +} + +// ============================================================================= +// PgUp/PgDown — программная навигация +// ============================================================================= + +// HandlePgUp возвращает предыдущую команду из истории. +func (r *Repl) HandlePgUp() string { + if len(r.history) == 0 { + return "" + } + if r.historyPos > 0 { + r.historyPos-- + } + if r.historyPos < 0 { + r.historyPos = 0 + } + return r.history[r.historyPos] +} + +// HandlePgDown возвращает следующую команду из истории. +func (r *Repl) HandlePgDown() string { + if len(r.history) == 0 { + return "" + } + if r.historyPos < len(r.history)-1 { + r.historyPos++ + } else { + r.historyPos = len(r.history) + return "" + } + return r.history[r.historyPos] +} + +// ============================================================================= +// ВЫВОД ЧЕРЕЗ ПЕЙДЖЕР +// ============================================================================= + +func (r *Repl) output(text string) { + if r.pager != nil && r.pager.Enabled() { + _ = r.pager.Page(text) + return + } + fmt.Print(text) +} + +// ============================================================================= +// ЗАКРЫТИЕ +// ============================================================================= + +func (r *Repl) Close() error { + // ДОБАВЛЕНО: снимаем активный batch при выходе, чтобы не оставить + // его зарегистрированным в глобальном менеджере. + if r.currentBatch != nil { + if bm := storage.GetBatchManager(); bm != nil { + bm.UnregisterBatch(r.currentBatch) + } + r.currentBatch = nil + r.currentBatchID = 0 + } + if r.historyFile != nil { + _ = r.historyFile.Save() + } + if r.rl != nil { + return r.rl.Close() + } + return nil +} + +// ============================================================================= +// ОБРАБОТЧИКИ КОМАНД +// ============================================================================= + +// ----------------------------------------------------------------------------- +// Базы данных +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleCreateSlice(args []string) error { + if len(args) < 1 { + return fmt.Errorf("usage: create slice ") + } + name := args[0] + if err := r.store.CreateDatabase(name); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Slice '%s' created at %s", name, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleDropDatabase(args []string) error { + if len(args) < 1 { + return fmt.Errorf("usage: drop database ") + } + name := args[0] + if err := r.store.DropDatabase(name); err != nil { + return err + } + if r.currentDB == name { + r.currentDB = "" + } + utils.PrintSuccess(fmt.Sprintf("Database '%s' dropped at %s", name, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handleUseDatabase(args []string) error { + if len(args) < 1 { + return fmt.Errorf("usage: use ") + } + name := args[0] + if !r.store.ExistsDatabase(name) { + return fmt.Errorf("database '%s' does not exist", name) + } + r.currentDB = name + utils.PrintSuccess(fmt.Sprintf("Switched to database '%s'", name)) + return nil +} + +func (r *Repl) handleShowDatabases(args []string) error { + databases := r.store.ListDatabases() + if len(databases) == 0 { + utils.PrintInfo("No databases found") + return nil + } + + t := NewTable("DATABASE", "CURRENT") + for _, db := range databases { + cur := "" + if db == r.currentDB { + cur = "*" + } + t.AddRow(db, cur) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Коллекции +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleCreateCollection(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 1 { + return fmt.Errorf("usage: create collection ") + } + name := args[0] + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + if err := db.CreateCollection(name); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Collection '%s' created in database '%s' at %s", name, r.currentDB, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handleDropCollection(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 1 { + return fmt.Errorf("usage: drop collection ") + } + name := args[0] + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + if err := db.DropCollection(name); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Collection '%s' dropped from database '%s'", name, r.currentDB)) + return nil +} + +func (r *Repl) handleShowCollections(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + collections := db.ListCollections() + if len(collections) == 0 { + utils.PrintInfo("No collections found") + return nil + } + + t := NewTable("COLLECTION") + for _, coll := range collections { + t.AddRow(coll) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Документы +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleInsert(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 2 { + return fmt.Errorf("usage: insert ") + } + + collName := args[0] + jsonStr := strings.Join(args[1:], " ") + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + doc := storage.NewDocument() + + if strings.Contains(jsonStr, "{") { + pairs := strings.Split(jsonStr, ",") + for _, pair := range pairs { + pair = strings.TrimSpace(pair) + pair = strings.Trim(pair, "{}") + kv := strings.SplitN(pair, "=", 2) + if len(kv) == 2 { + doc.SetField(kv[0], kv[1]) + } + } + } else { + pairs := strings.Split(jsonStr, ",") + for _, pair := range pairs { + pair = strings.TrimSpace(pair) + kv := strings.SplitN(pair, "=", 2) + if len(kv) == 2 { + doc.SetField(kv[0], kv[1]) + } + } + } + + if err := coll.Insert(doc); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Document inserted with ID: %s (created at: %s)", + doc.ID, time.UnixMilli(doc.CreatedAt).Format("02-012006 15:04:05.000"))) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleFind(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 2 { + return fmt.Errorf("usage: find ") + } + + collName := args[0] + docID := args[1] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + doc, err := coll.Find(docID) + if err != nil { + return err + } + + fields := doc.GetFields() + row := map[string]interface{}{"_id": doc.ID} + for k, v := range fields { + row[k] = v + } + rows := []map[string]interface{}{row} + + t := MapToTable(rows) + r.output(t.String()) + + utils.PrintInfo(fmt.Sprintf("created_at: %s", time.UnixMilli(doc.CreatedAt).Format("02-012006 15:04:05.000"))) + utils.PrintInfo(fmt.Sprintf("updated_at: %s", time.UnixMilli(doc.UpdatedAt).Format("02-012006 15:04:05.000"))) + if doc.DeletedAt > 0 { + utils.PrintWarning(fmt.Sprintf("deleted_at: %s", time.UnixMilli(doc.DeletedAt).Format("02-012006 15:04:05.000"))) + } + return nil +} + +func (r *Repl) handleFindByIndex(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 3 { + return fmt.Errorf("usage: findbyindex ") + } + + collName := args[0] + indexName := args[1] + value := args[2] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + docs, err := coll.FindByIndex(indexName, value) + 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} + for k, v := range doc.GetFields() { + row[k] = v + } + rows = append(rows, row) + } + t := MapToTable(rows) + r.output(t.String()) + return nil +} + +func (r *Repl) handleFindByTime(args []string) error { + if r.currentDB == "" { + 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") + } + + 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", + } + + for _, format := range formats { + fromTime, err = time.Parse(format, fromStr) + if err == nil { + break + } + } + 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 + } + } + if err != nil { + return fmt.Errorf("invalid to_date format: %s", toStr) + } + + fromMs := fromTime.UnixMilli() + toMs := toTime.UnixMilli() + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + 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} + for k, v := range doc.GetFields() { + row[k] = v + } + rows = append(rows, row) + } + t := MapToTable(rows) + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleUpdate(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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) + if len(kv) == 2 { + updates[kv[0]] = kv[1] + } + } + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.Update(docID, updates); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Document '%s' updated at %s", docID, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleDelete(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.Delete(docID); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Document '%s' deleted at %s", docID, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handlePermanentDelete(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.PermanentDelete(docID); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Document '%s' permanently deleted", docID)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleRestore(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.RestoreDeleted(docID); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Document '%s' restored at %s", docID, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleShowDeleted(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + docs := coll.GetAllDocumentsIncludingDeleted() + deleted := make([]*storage.Document, 0) + for _, doc := range docs { + if doc.IsDeleted() { + 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( + doc.ID, + time.UnixMilli(doc.DeletedAt).Format("02-012006 15:04:05.000"), + utils.ColorizeTextAny(doc.GetFields()), + ) + } + r.output(t.String()) + return nil +} + +func (r *Repl) handleCount(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 1 { + return fmt.Errorf("usage: count ") + } + + collName := args[0] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + 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)) + t.AddRow("Total documents", fmt.Sprintf("%d", total)) + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Временные метки +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleShowTimestamps(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + doc, err := coll.FindIncludingDeleted(docID) + if err != nil { + return err + } + + t := NewTable("FIELD", "VALUE") + t.AddRow("Created", time.UnixMilli(doc.CreatedAt).Format("02-012006 15:04:05.000")) + t.AddRow("Updated", time.UnixMilli(doc.UpdatedAt).Format("02-012006 15:04:05.000")) + if doc.DeletedAt > 0 { + t.AddRow("Deleted", time.UnixMilli(doc.DeletedAt).Format("02-012006 15:04:05.000")) + } else { + t.AddRow("Deleted", "not deleted") + } + t.AddRow("Version", fmt.Sprintf("%d", doc.Version)) + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleStatsTimestamps(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + docs := coll.GetAllDocumentsIncludingDeleted() + if len(docs) == 0 { + utils.PrintInfo("No documents in collection") + return nil + } + + var minCreated, maxCreated int64 = 1<<63 - 1, 0 + var minUpdated, maxUpdated int64 = 1<<63 - 1, 0 + var totalCreated, totalUpdated int64 = 0, 0 + + for _, doc := range docs { + if doc.CreatedAt < minCreated { + minCreated = doc.CreatedAt + } + if doc.CreatedAt > maxCreated { + maxCreated = doc.CreatedAt + } + if doc.UpdatedAt < minUpdated { + minUpdated = doc.UpdatedAt + } + if doc.UpdatedAt > maxUpdated { + maxUpdated = doc.UpdatedAt + } + totalCreated += doc.CreatedAt + totalUpdated += doc.UpdatedAt + } + + avgCreated := totalCreated / int64(len(docs)) + avgUpdated := totalUpdated / int64(len(docs)) + + t := NewTable("METRIC", "CREATED", "UPDATED") + t.AddRow("Documents", fmt.Sprintf("%d", len(docs)), "") + t.AddRow("Earliest", + time.UnixMilli(minCreated).Format("02-012006 15:04:05.000"), + time.UnixMilli(minUpdated).Format("02-012006 15:04:05.000")) + t.AddRow("Latest", + time.UnixMilli(maxCreated).Format("02-012006 15:04:05.000"), + time.UnixMilli(maxUpdated).Format("02-012006 15:04:05.000")) + t.AddRow("Average", + time.UnixMilli(avgCreated).Format("02-012006 15:04:05.000"), + time.UnixMilli(avgUpdated).Format("02-012006 15:04:05.000")) + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Аудит +// ----------------------------------------------------------------------------- + +func (r *Repl) handleAuditLog(args []string) error { + entries := storage.GetAuditLog() + if len(entries) == 0 { + 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] + t.AddRow(entry.TimestampStr, entry.Operation, entry.DataType, entry.Name) + } + r.output(t.String()) + return nil +} + +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") + } + + dataType := strings.ToUpper(args[0]) + operation := strings.ToUpper(args[1]) + + entries := storage.GetAuditLogFiltered(dataType, operation, 0, 0) + 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) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Индексы +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleCreateIndex(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 3 { + return fmt.Errorf("usage: create index [unique]") + } + + collName := args[0] + indexName := args[1] + fields := strings.Split(args[2], ",") + unique := len(args) > 3 && args[3] == "unique" + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.CreateIndex(indexName, fields, unique); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Index '%s' created on collection '%s' at %s", indexName, collName, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handleDropIndex(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.DropIndex(indexName); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Index '%s' dropped from collection '%s'", indexName, collName)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. +func (r *Repl) handleShowIndexes(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + 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" + if idx["unique"].(bool) { + uniqueStr = "true" + } + createdAt := time.UnixMilli(idx["created_at"].(int64)).Format("02-012006 15:04:05.000") + t.AddRow(idx["name"].(string), fmt.Sprintf("%v", idx["fields"]), uniqueStr, createdAt) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Ограничения +// ----------------------------------------------------------------------------- + +func (r *Repl) handleAddRequired(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + coll.AddRequiredField(field) + utils.PrintSuccess(fmt.Sprintf("Required field '%s' added to collection '%s'", field, collName)) + return nil +} + +func (r *Repl) handleAddUnique(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + coll.AddUniqueConstraint(field) + utils.PrintSuccess(fmt.Sprintf("Unique constraint added for field '%s' on collection '%s'", field, collName)) + return nil +} + +func (r *Repl) handleAddMin(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + coll.AddMinConstraint(field, minVal) + utils.PrintSuccess(fmt.Sprintf("Min constraint added for field '%s' on collection '%s' (min: %.2f)", field, collName, minVal)) + return nil +} + +func (r *Repl) handleAddMax(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + coll.AddMaxConstraint(field, maxVal) + utils.PrintSuccess(fmt.Sprintf("Max constraint added for field '%s' on collection '%s' (max: %.2f)", field, collName, maxVal)) + return nil +} + +func (r *Repl) handleAddEnum(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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 + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + coll.AddEnumConstraint(field, values) + utils.PrintSuccess(fmt.Sprintf("Enum constraint added for field '%s' on collection '%s' (allowed: %v)", field, collName, values)) + return nil +} + +// ----------------------------------------------------------------------------- +// Batch-операции (замена транзакций). +// +// ИСПРАВЛЕНО (2026-10): обработчики переписаны поверх batch-API. +// Семантика: +// - handleBeginTransaction → создаёт *storage.Batch (в памяти REPL); +// - handleCommitTransaction → batch.Commit() (атомарно через WAL); +// - handleRollbackTransaction→ снимает регистрацию batch без коммита; +// - handleShowTransactions → показывает активный batch сессии. +// +// ИЗМЕНЕНО (2026-10-03): формат временных меток "02-012006 15:04:05.000". +// ----------------------------------------------------------------------------- + +// handleBeginTransaction создаёт новый batch. +// Команда: "begin transaction" или "batch begin". +func (r *Repl) handleBeginTransaction(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + 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-012006 15:04:05.000"), + batch.Database, batch.Collection)) + return nil +} + +// handleCommitTransaction коммитит текущий batch. +// Команда: "commit" или "batch commit". +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 { + // Снимаем регистрацию даже при ошибке. + if bm := storage.GetBatchManager(); bm != nil { + bm.UnregisterBatch(batch) + } + r.currentBatch = nil + r.currentBatchID = 0 + return err + } + + // Commit сам снимает регистрацию через defer. + r.currentBatch = nil + r.currentBatchID = 0 + + utils.PrintSuccess(fmt.Sprintf("Batch %d committed successfully at %s", + batch.ID, time.Now().Format("02-012006 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-012006 15:04:05.000"))) + return nil +} + +// handleShowTransactions выводит список активных batch'ей сессии. +// Команда: "show transactions" или "show batches". +// +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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 { + started = time.UnixMilli(r.currentBatch.Timestamp).Format("02-012006 15:04:05.000") + } + t.AddRow( + fmt.Sprintf("%d", r.currentBatch.ID), + r.currentBatch.Database, + r.currentBatch.Collection, + fmt.Sprintf("%d", len(r.currentBatch.Operations)), + started, + ) + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Плагины +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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") + } + plugins := r.pluginManager.ListPlugins() + if len(plugins) == 0 { + utils.PrintInfo("No plugins loaded") + utils.PrintInfo(fmt.Sprintf("Plugins directory: %s", r.pluginManager.GetPluginsDir())) + return nil + } + + t := NewTable("NAME", "VERSION", "STATUS", "AUTHOR", "LOADED_AT") + for _, p := range plugins { + status := "loaded" + switch p.Status.Load() { + case 1: + status = "running" + case 2: + status = "stopped" + case 3: + status = "error" + } + t.AddRow(p.Name, p.Version(), status, p.Author(), p.LoadedAt().Format("02-012006 15:04:05.000")) + } + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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") + } + if len(args) < 2 { + return fmt.Errorf("usage: plugin load ") + } + name := args[0] + filepath := args[1] + if err := r.pluginManager.LoadPlugin(name, filepath); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Plugin '%s' loaded from %s at %s", name, filepath, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handlePluginUnload(args []string) error { + if r.pluginManager == nil || !r.pluginManager.IsEnabled() { + return fmt.Errorf("plugin system is disabled") + } + if len(args) < 1 { + return fmt.Errorf("usage: plugin unload ") + } + name := args[0] + if err := r.pluginManager.UnloadPlugin(name); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Plugin '%s' unloaded", name)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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") + } + if len(args) < 1 { + return fmt.Errorf("usage: plugin start ") + } + name := args[0] + if err := r.pluginManager.StartPlugin(name); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Plugin '%s' started at %s", name, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handlePluginStop(args []string) error { + if r.pluginManager == nil || !r.pluginManager.IsEnabled() { + return fmt.Errorf("plugin system is disabled") + } + if len(args) < 1 { + return fmt.Errorf("usage: plugin stop ") + } + name := args[0] + if err := r.pluginManager.StopPlugin(name); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Plugin '%s' stopped", name)) + return nil +} + +func (r *Repl) handlePluginExec(args []string) error { + if r.pluginManager == nil || !r.pluginManager.IsEnabled() { + return fmt.Errorf("plugin system is disabled") + } + if len(args) < 2 { + return fmt.Errorf("usage: plugin exec [args...]") + } + pluginName := args[0] + funcName := args[1] + + var execArgs []interface{} + for _, arg := range args[2:] { + execArgs = append(execArgs, arg) + } + + result, err := r.pluginManager.ExecutePlugin(pluginName, funcName, execArgs...) + if err != nil { + return err + } + utils.PrintSuccess("Result:") + r.output(FormatCell(result)) + return nil +} + +// ----------------------------------------------------------------------------- +// Экспорт/импорт +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleExport(args []string) error { + if len(args) < 2 { + return fmt.Errorf("usage: export ") + } + dbName := args[0] + fileName := args[1] + + if !r.store.ExistsDatabase(dbName) { + return fmt.Errorf("database '%s' does not exist", dbName) + } + + utils.PrintInfo(fmt.Sprintf("Exporting database '%s' to %s at %s...", dbName, fileName, time.Now().Format("02-012006 15:04:05.000"))) + + db, err := r.store.GetDatabase(dbName) + if err != nil { + return err + } + data, err := db.SerializeDatabase() + if err != nil { + return fmt.Errorf("failed to serialize database: %v", err) + } + if err := os.WriteFile(fileName, data, 0644); err != nil { + return fmt.Errorf("failed to write export file: %v", err) + } + utils.PrintSuccess(fmt.Sprintf("Database '%s' exported to %s", dbName, fileName)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleImport(args []string) error { + if len(args) < 2 { + return fmt.Errorf("usage: import ") + } + dbName := args[0] + fileName := args[1] + + utils.PrintInfo(fmt.Sprintf("Importing data from %s to database '%s' at %s...", fileName, dbName, time.Now().Format("02-012006 15:04:05.000"))) + + data, err := os.ReadFile(fileName) + if err != nil { + return fmt.Errorf("failed to read import file: %v", err) + } + + if !r.store.ExistsDatabase(dbName) { + if err := r.store.CreateDatabase(dbName); err != nil { + return err + } + } + + db, err := r.store.GetDatabase(dbName) + if err != nil { + return err + } + if err := db.DeserializeDatabase(data); err != nil { + return fmt.Errorf("failed to deserialize database: %v", err) + } + utils.PrintSuccess(fmt.Sprintf("Data imported to database '%s' from %s", dbName, fileName)) + return nil +} + +// ----------------------------------------------------------------------------- +// ACL +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleACLLogin(args []string) error { + if len(args) < 2 { + return fmt.Errorf("usage: acl login ") + } + username := args[0] + password := args[1] + + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + + sessionID, err := r.aclManager.Authenticate(username, password) + if err != nil { + return err + } + + r.authenticated = true + r.currentUser = username + r.sessionID = sessionID + + roles := r.aclManager.GetUserRoles(sessionID) + if len(roles) > 0 { + r.currentRole = roles[0] + } + + utils.PrintSuccess(fmt.Sprintf("Logged in as '%s' with role '%s' at %s", username, r.currentRole, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleACLLogout(args []string) error { + if r.sessionID != "" && r.aclManager != nil { + r.aclManager.Logout(r.sessionID) + } + r.authenticated = false + r.currentUser = "" + r.currentRole = "anonymous" + r.sessionID = "" + utils.PrintSuccess(fmt.Sprintf("Logged out at %s", time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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") + } + if len(args) < 3 { + return fmt.Errorf("usage: acl grant \n"+ + " Permissions: r=read, w=write, d=delete, a=admin\n"+ + " Example: acl grant users admin rwa") + } + + 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 + } + + canRead := strings.Contains(perms, "r") + canWrite := strings.Contains(perms, "w") + canDelete := strings.Contains(perms, "d") + isAdmin := strings.Contains(perms, "a") + + coll.SetACL(role, canRead, canWrite, canDelete, isAdmin) + utils.PrintSuccess(fmt.Sprintf("Permissions '%s' granted to role '%s' on collection '%s' at %s", + perms, role, collName, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +func (r *Repl) handleACLUsers(args []string) error { + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + users := r.aclManager.ListUsers() + if len(users) == 0 { + utils.PrintInfo("No users found") + return nil + } + + t := NewTable("USERNAME", "ROLES", "STATUS") + for _, username := range users { + userInfo, err := r.aclManager.GetUserInfo(username) + if err != nil { + continue + } + status := "active" + if !userInfo.Active { + status = "disabled" + } + t.AddRow(username, fmt.Sprintf("%v", userInfo.Roles), status) + } + r.output(t.String()) + return nil +} + +func (r *Repl) handleACLRoles(args []string) error { + if r.aclManager == nil { + return fmt.Errorf("ACL manager not initialized") + } + roles := r.aclManager.ListRoles() + if len(roles) == 0 { + utils.PrintInfo("No roles found") + return nil + } + + t := NewTable("ROLE", "PERMISSIONS") + for _, roleName := range roles { + perms, err := r.aclManager.GetRolePermissions(roleName) + if err != nil { + continue + } + t.AddRow(roleName, fmt.Sprintf("%v", perms)) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Сжатие +// ----------------------------------------------------------------------------- + +func (r *Repl) handleCompressionStats(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + + collections := db.ListCollections() + totalDocs := int64(0) + compressedDocs := int64(0) + totalOriginalSize := int64(0) + totalCompressedSize := int64(0) + + for _, collName := range collections { + coll, err := db.GetCollection(collName) + if err != nil { + continue + } + docs := coll.GetAllDocuments() + for _, doc := range docs { + totalDocs++ + if doc.Compressed { + compressedDocs++ + totalOriginalSize += doc.OriginalSize + if data, err := doc.Serialize(); err == nil { + totalCompressedSize += int64(len(data)) + } + } + } + } + + t := NewTable("METRIC", "VALUE") + t.AddRow("Total Documents", fmt.Sprintf("%d", totalDocs)) + t.AddRow("Compressed Documents", fmt.Sprintf("%d", compressedDocs)) + if totalDocs > 0 { + t.AddRow("Compression Rate", fmt.Sprintf("%.2f%%", float64(compressedDocs)/float64(totalDocs)*100)) + } + if totalOriginalSize > 0 { + ratio := float64(totalCompressedSize) / float64(totalOriginalSize) + t.AddRow("Size Reduction", fmt.Sprintf("%.2f%%", (1-ratio)*100)) + t.AddRow("Original Size", utils.FormatBytes(totalOriginalSize)) + t.AddRow("Compressed Size", utils.FormatBytes(totalCompressedSize)) + } + t.AddRow("Algorithm", r.config.Compression.Algorithm) + t.AddRow("Compression Level", fmt.Sprintf("%d", r.config.Compression.Level)) + t.AddRow("Min Size Threshold", utils.FormatBytes(int64(r.config.Compression.MinSize))) + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleCompressCollection(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 1 { + return fmt.Errorf("usage: compress collection ") + } + collName := args[0] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + docs := coll.GetAllDocuments() + compressed := 0 + + utils.PrintInfo(fmt.Sprintf("Compressing collection '%s' at %s...", collName, time.Now().Format("02-012006 15:04:05.000"))) + + for _, doc := range docs { + if !doc.Compressed { + if err := doc.Compress(&compression.Config{ + Enabled: r.config.Compression.Enabled, + Algorithm: r.config.Compression.Algorithm, + Level: r.config.Compression.Level, + MinSize: r.config.Compression.MinSize, + }); err == nil { + compressed++ + } + } + } + utils.PrintSuccess(fmt.Sprintf("Compressed %d documents in collection '%s'", compressed, collName)) + return nil +} + +func (r *Repl) handleDocCompression(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 2 { + return fmt.Errorf("usage: doc compression ") + } + collName := args[0] + docID := args[1] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + doc, err := coll.Find(docID) + if err != nil { + return err + } + + t := NewTable("METRIC", "VALUE") + t.AddRow("Compressed", fmt.Sprintf("%v", doc.Compressed)) + if doc.Compressed { + ratio := doc.GetCompressionRatio() + t.AddRow("Ratio", fmt.Sprintf("%.2f%%", (1-ratio)*100)) + t.AddRow("Original Size", utils.FormatBytes(doc.OriginalSize)) + if data, err := doc.Serialize(); err == nil { + t.AddRow("Current Size", utils.FormatBytes(int64(len(data)))) + } + } + r.output(t.String()) + return nil +} + +func (r *Repl) handleCompressionConfig(args []string) error { + 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("snappy", "Fast compression/decompression, good balance (default)") + t.AddRow("lz4", "Extremely fast, lower compression ratio") + t.AddRow("zstd", "High compression ratio, slower") + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Триггеры +// ----------------------------------------------------------------------------- + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +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 ") + } + + collName := args[0] + triggerName := args[1] + event := args[2] + action := args[3] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + trigger := &storage.Trigger{ + Name: triggerName, + Event: event, + Action: action, + Enabled: true, + Description: "", + CreatedAt: time.Now().UnixMilli(), + UpdatedAt: time.Now().UnixMilli(), + } + + for i := 4; i < len(args); i++ { + switch args[i] { + case "--description": + if i+1 < len(args) { + trigger.Description = args[i+1] + i++ + } + } + } + + if err := coll.AddTrigger(trigger); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Trigger '%s' created on collection '%s' for event %s at %s", + triggerName, collName, event, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +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 ") + } + collName := args[0] + triggerName := args[2] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.DropTrigger(triggerName); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Trigger '%s' dropped from collection '%s'", triggerName, collName)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. +func (r *Repl) handleShowTriggers(args []string) error { + if r.currentDB == "" { + return fmt.Errorf("no database selected") + } + if len(args) < 1 { + return fmt.Errorf("usage: show triggers ") + } + collName := args[0] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + triggers := coll.ListTriggers() + if len(triggers) == 0 { + utils.PrintInfo(fmt.Sprintf("No triggers found on collection '%s'", collName)) + return nil + } + + t := NewTable("NAME", "EVENT", "ACTION", "STATUS", "CREATED_AT", "DESCRIPTION") + for _, tr := range triggers { + status := "enabled" + if !tr.Enabled { + status = "disabled" + } + t.AddRow(tr.Name, tr.Event, tr.Action, status, + time.UnixMilli(tr.CreatedAt).Format("02-012006 15:04:05.000"), tr.Description) + } + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +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 ") + } + collName := args[0] + triggerName := args[2] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.EnableTrigger(triggerName); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Trigger '%s' enabled at %s", triggerName, time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +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 ") + } + collName := args[0] + triggerName := args[2] + + db, err := r.store.GetDatabase(r.currentDB) + if err != nil { + return err + } + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + if err := coll.DisableTrigger(triggerName); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Trigger '%s' disabled", triggerName)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. +func (r *Repl) handleTriggerLog(args []string) error { + logs := storage.GetTriggerExecutionLog() + if len(logs) == 0 { + utils.PrintInfo("No trigger executions logged") + return nil + } + + t := NewTable("TIMESTAMP", "TRIGGER", "EVENT", "COLLECTION", "DOCUMENT") + for _, entry := range logs { + t.AddRow( + entry.Timestamp.Format("02-012006 15:04:05.000"), + entry.TriggerName, + entry.Event, + entry.Collection, + entry.DocumentID, + ) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Кластер +// ----------------------------------------------------------------------------- + +func (r *Repl) handleStatus(args []string) error { + if r.coordinator == nil { + utils.PrintWarning("Cluster coordinator not available") + return nil + } + + t := NewTable("METRIC", "VALUE") + isLeader := r.coordinator.IsLeader() + if isLeader { + t.AddRow("Role", "LEADER") + } else { + t.AddRow("Role", "FOLLOWER") + } + + leader := r.coordinator.GetLeader() + if leader != nil { + leaderIP := "unknown" + leaderPort := 0 + if leader.IP != "" { + leaderIP = leader.IP + } + if leader.Port > 0 { + leaderPort = leader.Port + } + t.AddRow("Leader", fmt.Sprintf("%s:%d", leaderIP, leaderPort)) + } else { + t.AddRow("Leader", "none") + } + + status := r.coordinator.GetClusterStatus() + t.AddRow("Cluster Name", status.Name) + t.AddRow("Total Nodes", fmt.Sprintf("%d", status.TotalNodes)) + t.AddRow("Active Nodes", fmt.Sprintf("%d", status.ActiveNodes)) + t.AddRow("Health", status.Health) + t.AddRow("Replication Factor", fmt.Sprintf("%d", status.ReplicationFactor)) + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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") + } + + nodes := r.coordinator.GetAllNodes() + if len(nodes) == 0 { + utils.PrintInfo("No nodes in cluster") + return nil + } + + leader := r.coordinator.GetLeader() + leaderID := "" + if leader != nil { + leaderID = leader.ID + } + + t := NewTable("", "ID", "ADDRESS", "STATUS", "LAST_SEEN") + for _, node := range nodes { + prefix := " " + if node.ID == leaderID { + prefix = " *" + } + nodeIP := "unknown" + nodePort := 0 + if node.IP != "" { + nodeIP = node.IP + } + if node.Port > 0 { + nodePort = node.Port + } + lastSeen := time.UnixMilli(node.LastSeen).Format("02-012006 15:04:05.000") + t.AddRow(prefix, node.ID, fmt.Sprintf("%s:%d", nodeIP, nodePort), node.Status, lastSeen) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Кластер (расширенные команды) +// ----------------------------------------------------------------------------- + +func (r *Repl) handleClusterPipeline(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + stats := r.coordinator.GetPipelineStats() + t := NewTable("METRIC", "VALUE") + for k, v := range stats { + t.AddRow(k, fmt.Sprintf("%v", v)) + } + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleClusterReshard(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + if err := r.coordinator.TriggerResharding("manual"); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Resharding triggered at %s", time.Now().Format("02-012006 15:04:05.000"))) + return nil +} + +// ----------------------------------------------------------------------------- +// Миграция +// ----------------------------------------------------------------------------- + +func (r *Repl) handleMigration(args []string) error { + if len(args) == 0 { + return fmt.Errorf("usage: migration \n"+ + " Subcommands:\n"+ + " start [database] [collection] - Start a new migration\n"+ + " status [task_id] - Show migration status\n"+ + " list - List all migration tasks\n"+ + " pause - Pause a running migration\n"+ + " resume - Resume a paused migration\n"+ + " cancel - Cancel a migration\n"+ + " stats - Show migration statistics\n"+ + " config - Show migration configuration\n"+ + " queue - Show change queue status") + } + + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + + migrator := r.coordinator.GetCrossDCMigrator() + if migrator == nil { + return fmt.Errorf("migration system not initialized (check config: migration.enabled = true)") + } + + subcommand := args[0] + switch subcommand { + case "start": + if len(args) < 3 { + return fmt.Errorf("usage: migration start [database] [collection]") + } + return r.handleMigrationStart(migrator, args[1:]) + case "status": + return r.handleMigrationStatus(migrator, args[1:]) + case "list": + return r.handleMigrationList(migrator) + case "pause": + if len(args) < 2 { + return fmt.Errorf("usage: migration pause ") + } + return r.handleMigrationPause(migrator, args[1]) + case "resume": + if len(args) < 2 { + return fmt.Errorf("usage: migration resume ") + } + return r.handleMigrationResume(migrator, args[1]) + case "cancel": + if len(args) < 2 { + return fmt.Errorf("usage: migration cancel ") + } + return r.handleMigrationCancel(migrator, args[1]) + case "stats": + return r.handleMigrationStats(migrator) + case "config": + return r.handleMigrationConfig(migrator) + case "queue": + return r.handleMigrationQueue(migrator) + default: + return fmt.Errorf("unknown migration subcommand: %s", subcommand) + } +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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]") + } + sourceDC := args[0] + targetDC := args[1] + + var databases []string + var collections map[string][]string + if len(args) > 2 { + databases = []string{args[2]} + if len(args) > 3 { + collections = map[string][]string{ + args[2]: {args[3]}, + } + } + } + + utils.PrintInfo(fmt.Sprintf("Starting migration from %s to %s at %s...", + sourceDC, targetDC, time.Now().Format("02-012006 15:04:05.000"))) + + task, err := migrator.StartMigration(sourceDC, targetDC, databases, collections) + if err != nil { + return fmt.Errorf("failed to start migration: %v", err) + } + + t := NewTable("FIELD", "VALUE") + t.AddRow("Task ID", task.ID) + t.AddRow("Source", task.SourceDC) + t.AddRow("Target", task.TargetDC) + t.AddRow("Total documents", fmt.Sprintf("%d", task.TotalDocuments)) + t.AddRow("Status", string(task.Status)) + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +// ИЗМЕНЕНО (2026-10-04): УНИФИЦИРОВАН — добавлены миллисекунды. +func (r *Repl) handleMigrationStatus(migrator *cluster.CrossDCMigrator, args []string) error { + var taskID string + if len(args) > 0 { + taskID = args[0] + } else { + taskID = migrator.GetCurrentTaskID() + if taskID == "" { + return fmt.Errorf("no active migration task") + } + } + + task, err := migrator.GetMigrationStatus(taskID) + if err != nil { + return err + } + + t := NewTable("FIELD", "VALUE") + t.AddRow("Task ID", task.ID) + t.AddRow("Source DC", task.SourceDC) + t.AddRow("Target DC", task.TargetDC) + t.AddRow("Status", string(task.Status)) + t.AddRow("Progress", fmt.Sprintf("%.1f%%", task.ProgressPercent)) + t.AddRow("Total Docs", fmt.Sprintf("%d", task.TotalDocuments)) + t.AddRow("Migrated", fmt.Sprintf("%d", task.MigratedDocs)) + t.AddRow("Failed", fmt.Sprintf("%d", task.FailedDocs)) + t.AddRow("Skipped", fmt.Sprintf("%d", task.SkippedDocs)) + + if task.StartTime > 0 { + t.AddRow("Started", time.UnixMilli(task.StartTime).Format("02-012006 15:04:05.000")) + } + if task.EndTime > 0 { + t.AddRow("Completed", time.UnixMilli(task.EndTime).Format("02-012006 15:04:05.000")) + } + if task.Error != "" { + t.AddRow("Error", task.Error) + } + if len(task.Databases) > 0 { + t.AddRow("Databases", strings.Join(task.Databases, ", ")) + } + if len(task.Collections) > 0 { + for db, colls := range task.Collections { + t.AddRow("Collections", fmt.Sprintf("%s: %s", db, strings.Join(colls, ", "))) + } + } + r.output(t.String()) + return nil +} + +func (r *Repl) handleMigrationList(migrator *cluster.CrossDCMigrator) error { + tasks := migrator.ListTasks() + if len(tasks) == 0 { + utils.PrintInfo("No migration tasks found") + return nil + } + + t := NewTable("ID", "STATUS", "PROGRESS", "SOURCE", "TARGET") + for _, task := range tasks { + id := task.ID + if len(id) > 36 { + id = id[:33] + "..." + } + progress := fmt.Sprintf("%.1f%%", task.ProgressPercent) + if task.Status == cluster.MigrationStatusCompleted { + progress = "100%" + } + t.AddRow(id, string(task.Status), progress, task.SourceDC, task.TargetDC) + } + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleMigrationPause(migrator *cluster.CrossDCMigrator, taskID string) error { + if err := migrator.PauseMigration(taskID); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Migration %s paused", taskID)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleMigrationResume(migrator *cluster.CrossDCMigrator, taskID string) error { + if err := migrator.ResumeMigration(taskID); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Migration %s resumed", taskID)) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleMigrationCancel(migrator *cluster.CrossDCMigrator, taskID string) error { + if err := migrator.CancelMigration(taskID); err != nil { + return err + } + utils.PrintSuccess(fmt.Sprintf("Migration %s cancelled", taskID)) + return nil +} + +func (r *Repl) handleMigrationStats(migrator *cluster.CrossDCMigrator) error { + stats := migrator.GetMigrationStats() + t := NewTable("METRIC", "VALUE") + t.AddRow("Total Changes", fmt.Sprintf("%d", stats.TotalChanges)) + t.AddRow("Applied Changes", fmt.Sprintf("%d", stats.AppliedChanges)) + t.AddRow("Failed Changes", fmt.Sprintf("%d", stats.FailedChanges)) + t.AddRow("Skipped Changes", fmt.Sprintf("%d", stats.SkippedChanges)) + t.AddRow("Avg Latency", fmt.Sprintf("%d ms", stats.LatencyAvg)) + t.AddRow("Max Latency", fmt.Sprintf("%d ms", stats.LatencyMax)) + t.AddRow("Throughput", fmt.Sprintf("%d docs/sec", stats.Throughput)) + r.output(t.String()) + return nil +} + +func (r *Repl) handleMigrationConfig(migrator *cluster.CrossDCMigrator) error { + cfg := migrator.GetConfig() + t := NewTable("SETTING", "VALUE") + t.AddRow("Enabled", fmt.Sprintf("%v", cfg.Enabled)) + t.AddRow("Mode", cfg.Mode) + if cfg.Source != nil { + t.AddRow("Source DC", fmt.Sprintf("%s (%s)", cfg.Source.Name, cfg.Source.Endpoint)) + } + if cfg.Target != nil { + t.AddRow("Target DC", fmt.Sprintf("%s (%s)", cfg.Target.Name, cfg.Target.Endpoint)) + } + if cfg.Settings != nil { + t.AddRow("Batch Size", fmt.Sprintf("%d", cfg.Settings.BatchSize)) + t.AddRow("Workers", fmt.Sprintf("%d", cfg.Settings.Workers)) + t.AddRow("Compression", cfg.Settings.Compression) + t.AddRow("Resume Enabled", fmt.Sprintf("%v", cfg.Settings.ResumeEnabled)) + t.AddRow("Max Retries", fmt.Sprintf("%d", cfg.Settings.MaxRetries)) + } + if cfg.Delta != nil { + t.AddRow("Delta Sync", fmt.Sprintf("%v", cfg.Delta.Enabled)) + t.AddRow("Delta Interval", fmt.Sprintf("%d sec", cfg.Delta.IntervalSec)) + } + if cfg.Validation != nil { + t.AddRow("Validation", fmt.Sprintf("%v", cfg.Validation.Enabled)) + t.AddRow("Sample Percent", fmt.Sprintf("%d%%", cfg.Validation.SamplePercent)) + } + r.output(t.String()) + return nil +} + +func (r *Repl) handleMigrationQueue(migrator *cluster.CrossDCMigrator) error { + stats := migrator.GetQueueStats() + t := NewTable("METRIC", "VALUE") + t.AddRow("Queue Size", fmt.Sprintf("%v", stats["size"])) + t.AddRow("Last LSN", fmt.Sprintf("%v", stats["last_lsn"])) + t.AddRow("Max Size", fmt.Sprintf("%v", stats["max_size"])) + t.AddRow("File Path", fmt.Sprintf("%v", stats["file_path"])) + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Миграция схемы +// ----------------------------------------------------------------------------- + +func (r *Repl) handleMigrateStatus(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + status := r.coordinator.GetMigrationStatus() + if status == nil { + utils.PrintInfo("No schema migration status available") + return nil + } + t := NewTable("METRIC", "VALUE") + t.AddRow("Status", fmt.Sprintf("%v", status)) + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Fallback (SPoF protection) +// ----------------------------------------------------------------------------- + +func (r *Repl) handleFallbackStatus(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + stats := r.coordinator.GetFallbackStats() + t := NewTable("METRIC", "VALUE") + for k, v := range stats { + t.AddRow(k, fmt.Sprintf("%v", v)) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Panic recovery +// ----------------------------------------------------------------------------- + +func (r *Repl) handlePanicStats(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + stats := r.coordinator.GetPanicRecoveryStats() + t := NewTable("METRIC", "VALUE") + for k, v := range stats { + t.AddRow(k, fmt.Sprintf("%v", v)) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Persistence +// ----------------------------------------------------------------------------- + +func (r *Repl) handlePersistStatus(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + pm := r.coordinator.GetPersistenceManager() + if pm == nil { + return fmt.Errorf("persistence manager not available") + } + info := pm.GetLastCheckpointInfo() + t := NewTable("METRIC", "VALUE") + for k, v := range info { + t.AddRow(k, fmt.Sprintf("%v", v)) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// SAGA +// ----------------------------------------------------------------------------- + +func (r *Repl) handleSagaStatus(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + sm := r.coordinator.GetSagaManager() + if sm == nil { + return fmt.Errorf("saga manager not available") + } + orchestrator := sm.GetOrchestrator() + if orchestrator == nil { + utils.PrintInfo("SAGA orchestrator not initialized") + return nil + } + metrics := orchestrator.GetMetrics() + t := NewTable("METRIC", "VALUE") + for k, v := range metrics { + t.AddRow(k, fmt.Sprintf("%v", v)) + } + r.output(t.String()) + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 15:04:05.000". +func (r *Repl) handleSagaList(args []string) error { + if r.coordinator == nil { + return fmt.Errorf("cluster coordinator not available") + } + sm := r.coordinator.GetSagaManager() + 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 + } + t := NewTable("ID", "STATUS", "STEPS", "CURRENT_STEP", "CREATED_AT") + for _, saga := range active { + t.AddRow( + saga.ID, + saga.Status, + fmt.Sprintf("%d", len(saga.Steps)), + fmt.Sprintf("%d", saga.CurrentStep), + time.UnixMilli(saga.CreatedAt).Format("02-012006 15:04:05.000"), + ) + } + r.output(t.String()) + return nil +} + +// ----------------------------------------------------------------------------- +// Системные команды +// ----------------------------------------------------------------------------- + +func (r *Repl) handleHelp(args []string) error { + var sb strings.Builder + sb.WriteString("\n=== Futriix available commands ===\n") + + categories := map[string][]struct { + cmd string + description string + }{ + "Database Management": { + {"create slice ", "Create a new database (slice)"}, + {"drop database ", "Delete an existing database"}, + {"use ", "Switch to a specific database"}, + {"show databases", "List all available databases"}, + }, + "Collection Management": { + {"create collection ", "Create a new collection in current database"}, + {"drop collection ", "Delete a collection from current database"}, + {"show collections", "List all collections in current database"}, + }, + "Document Operations": { + {"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"}, + {"update ...", "Update fields of an existing document"}, + {"delete ", "Delete a document (soft delete if enabled)"}, + {"permanent delete ", "Permanently delete a soft-deleted document"}, + {"restore ", "Restore a soft-deleted document"}, + {"show deleted ", "Show soft-deleted documents in a collection"}, + {"count ", "Count total documents in a collection"}, + }, + "Timestamp Management": { + {"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"}, + }, + "Index Management": { + {"create index [unique]", "Create a new index on specified fields"}, + {"drop index ", "Remove an existing index"}, + {"show indexes ", "List all indexes on a collection"}, + }, + "Constraints": { + {"add required ", "Add a required field constraint"}, + {"add unique ", "Add a unique constraint on a field"}, + {"add min ", "Add a minimum value constraint for numeric fields"}, + {"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)"}, + {"batch abort", "Abort the current batch (operations are discarded)"}, + {"show batches", "Show the active batch of the current session"}, + {"begin transaction", "Alias for 'batch begin'"}, + {"commit", "Alias for 'batch commit'"}, + {"rollback", "Alias for 'batch abort'"}, + {"show transactions", "Alias for 'show batches'"}, + }, + "Triggers (MongoDB-like)": { + {"create trigger [options]", "Create a trigger on collection events"}, + {"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"}, + {"trigger log", "Show trigger execution history"}, + }, + "Plugins": { + {"plugin list", "List all loaded plugins"}, + {"plugin load ", "Load a plugin from Lua file"}, + {"plugin unload ", "Unload a plugin"}, + {"plugin start ", "Start a loaded plugin"}, + {"plugin stop ", "Stop a running plugin"}, + {"plugin exec [args...]", "Execute a plugin function"}, + }, + "Import/Export": { + {"export ", "Export entire database to file"}, + {"import ", "Import database from file"}, + }, + "Compression": { + {"compression stats", "Show compression statistics for current database"}, + {"compression config", "Display current compression settings"}, + {"compress collection ", "Manually compress all documents in a collection"}, + {"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 users", "List all users"}, + {"acl roles", "List all roles"}, + }, + "Cluster": { + {"status", "Show current cluster status and role (leader/follower)"}, + {"nodes", "List all nodes in the cluster"}, + {"cluster pipeline", "Show pipeline replicator statistics"}, + {"cluster reshard", "Trigger manual resharding"}, + }, + "Migration": { + {"migration start [db] [coll]", "Start cross-datacenter migration"}, + {"migration status [task_id]", "Show migration status"}, + {"migration list", "List all migration tasks"}, + {"migration pause ", "Pause a running migration"}, + {"migration resume ", "Resume a paused migration"}, + {"migration cancel ", "Cancel a migration"}, + {"migration stats", "Show migration statistics"}, + {"migration config", "Show migration configuration"}, + {"migration queue", "Show change queue status"}, + {"migrate status", "Show schema migration status"}, + }, + "Fallback": { + {"fallback status", "Show SPoF protection status"}, + }, + "Panic Recovery": { + {"panic stats", "Show panic recovery statistics"}, + }, + "Persistence": { + {"persist status", "Show data persistence status"}, + }, + "SAGA": { + {"saga status", "Show SAGA orchestrator status"}, + {"saga list", "List active SAGA transactions"}, + }, + "Pager & Input": { + {"pager on|off", "Enable/disable pager"}, + {"pager", "Show pager status"}, + {"PgUp/PgDown", "Navigate command history (like in bash)"}, + {"Tab", "Autocomplete command / collection / field"}, + {"Ctrl+D", "Exit REPL"}, + {"Ctrl+C", "Interrupt current input"}, + }, + "System": { + {"help", "Display this help message with all available commands"}, + {"clear", "Clear the terminal screen"}, + {"quit", "Exit the futriix database REPL"}, + {"exit", "Exit the futriix database REPL (alias for quit)"}, + }, + } + + for category, commands := range categories { + sb.WriteString(fmt.Sprintf("\n%s:\n", category)) + for _, cmd := range commands { + if cmd.description != "" { + sb.WriteString(fmt.Sprintf(" %-50s %s\n", cmd.cmd, cmd.description)) + } else { + sb.WriteString(fmt.Sprintf(" %s\n", cmd.cmd)) + } + } + } + + sb.WriteString("\nMulti-line input:\n") + sb.WriteString(" End line with '\\' to continue\n") + sb.WriteString(" Or use '{' ... '}' — input continues until brackets are balanced\n") + sb.WriteString(" Empty line finishes multi-line input\n") + sb.WriteString(" \n") + + utils.PrintInfo(sb.String()) + return nil +} + +func (r *Repl) handleClear(args []string) error { + fmt.Print("\033[2J\033[H") + return nil +} + +// ИЗМЕНЕНО (2026-10-03): формат "02-012006 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) + } + r.currentBatch = nil + r.currentBatchID = 0 + } + utils.PrintInfo(fmt.Sprintf("Session ended at %s", time.Now().Format("02-012006 15:04:05.000"))) + if r.rl != nil { + _ = r.rl.Close() + } + os.Exit(0) + return nil +}