Files
futriix/internal/repl/repl.go
T

4060 lines
126 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/*
* 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
*/
// // Файл: internal/repl/repl.go
// =============================================================================
// Пакет: repl (Read-Eval-Print Loop)
// =============================================================================
// Назначение: Интерактивный интерфейс командной строки (REPL) для СУБД futriix.
// Совместимо с Linux и OpenIndiana.
//
// ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции
// команды REPL для транзакций переписаны поверх batch-API.
//
// ИЗМЕНЕНО (2026-10-03): формат временных меток изменён на
// "02-01-2006 15:04:05.000".
//
// ИСПРАВЛЕНО (2026-10, аудит): см. отдельные пометки в методах.
//
// ДОБАВЛЕНО (2026-10, расширение функциональности):
// - config show / get / set / history
// - acl user create/delete, acl role create/delete
// - pager size/threshold <N>
// - saga list --all
// - cluster pipeline on/off
//
// ИСПРАВЛЕНО (2026-10, типы для saga list):
// handleSagaList теперь использует cluster.SagaTransaction
// (а не storage.SagaTransaction), потому что r.coordinator.GetSagaManager()
// возвращает *cluster.SagaManager.
// =============================================================================
package repl
import (
"errors"
"fmt"
"io"
"os"
"path/filepath"
"sort"
"strconv"
"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
dynamicConfigManager *config.DynamicConfigManager
configPath string
rl *readline.Instance
interrupt bool
currentDB string
currentUser string
currentRole string
authenticated bool
sessionID string
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)
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,
VimMode: false,
FuncFilterInputRune: func(rn rune) (rune, bool) {
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
}
// SetDynamicConfigManager устанавливает менеджер динамической конфигурации.
func (r *Repl) SetDynamicConfigManager(dcm *config.DynamicConfigManager, configPath string) {
r.dynamicConfigManager = dcm
r.configPath = configPath
}
// =============================================================================
// РЕГИСТРАЦИЯ КОМАНД
// =============================================================================
func (r *Repl) registerCommands() {
// --- Базы данных ---
r.commands["create slice"] = &Command{Name: "create slice", Description: "Create a new database (slice)", Handler: r.handleCreateSlice}
r.commands["drop database"] = &Command{Name: "drop database", Description: "Drop a database", Handler: r.handleDropDatabase}
r.commands["use"] = &Command{Name: "use", Description: "Switch to a database", Handler: r.handleUseDatabase}
r.commands["show databases"] = &Command{Name: "show databases", Description: "List all databases", Handler: r.handleShowDatabases}
// --- Коллекции ---
r.commands["create collection"] = &Command{Name: "create collection", Description: "Create a new collection in current database", Handler: r.handleCreateCollection}
r.commands["drop collection"] = &Command{Name: "drop collection", Description: "Drop a collection from current database", Handler: r.handleDropCollection}
r.commands["show collections"] = &Command{Name: "show collections", Description: "List all collections in current database", Handler: r.handleShowCollections}
// --- Документы ---
r.commands["insert"] = &Command{Name: "insert", Description: "Insert a document into a collection (JSON format)", Handler: r.handleInsert}
r.commands["find"] = &Command{Name: "find", Description: "Find a document by ID", Handler: r.handleFind}
r.commands["findbyindex"] = &Command{Name: "findbyindex", Description: "Find documents by index", Handler: r.handleFindByIndex}
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-операции ---
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 revoke"] = &Command{Name: "acl revoke", Description: "Revoke permissions (r=read,w=write,d=delete,a=admin)", Handler: r.handleACLRevoke}
r.commands["acl users"] = &Command{Name: "acl users", Description: "List all users", Handler: r.handleACLUsers}
r.commands["acl roles"] = &Command{Name: "acl roles", Description: "List all roles", Handler: r.handleACLRoles}
r.commands["acl user create"] = &Command{Name: "acl user create", Description: "Create a new user", Handler: r.handleACLUserCreate}
r.commands["acl user delete"] = &Command{Name: "acl user delete", Description: "Delete a user", Handler: r.handleACLUserDelete}
r.commands["acl role create"] = &Command{Name: "acl role create", Description: "Create a new role", Handler: r.handleACLRoleCreate}
r.commands["acl role delete"] = &Command{Name: "acl role delete", Description: "Delete a role", Handler: r.handleACLRoleDelete}
// --- Сжатие ---
r.commands["compression stats"] = &Command{Name: "compression stats", Description: "Show compression statistics for the database", Handler: r.handleCompressionStats}
r.commands["compress collection"] = &Command{Name: "compress collection", Description: "Manually compress all documents in a collection", Handler: r.handleCompressCollection}
r.commands["doc compression"] = &Command{Name: "doc compression", Description: "Show compression ratio for a document", Handler: r.handleDocCompression}
r.commands["compression config"] = &Command{Name: "compression config", Description: "Show or set compression configuration", Handler: r.handleCompressionConfig}
// --- Аудит ---
r.commands["audit log"] = &Command{Name: "audit log", Description: "Show audit log", Handler: r.handleAuditLog}
r.commands["audit filter"] = &Command{Name: "audit filter", Description: "Filter audit log by type and operation", Handler: r.handleAuditFilter}
// --- Кластер ---
r.commands["status"] = &Command{Name: "status", Description: "Show cluster status", Handler: r.handleStatus}
r.commands["nodes"] = &Command{Name: "nodes", Description: "List cluster nodes", Handler: r.handleNodes}
r.commands["cluster pipeline"] = &Command{Name: "cluster pipeline", Description: "Show pipeline replicator statistics", Handler: r.handleClusterPipeline}
r.commands["cluster pipeline on"] = &Command{Name: "cluster pipeline on", Description: "Enable pipeline replication", Handler: r.handleClusterPipelineOn}
r.commands["cluster pipeline off"] = &Command{Name: "cluster pipeline off", Description: "Disable pipeline replication", Handler: r.handleClusterPipelineOff}
r.commands["cluster reshard"] = &Command{Name: "cluster reshard", Description: "Trigger manual resharding", Handler: r.handleClusterReshard}
// --- Миграция ---
r.commands["migration"] = &Command{Name: "migration", Description: "Cross-datacenter migration commands", Handler: r.handleMigration}
r.commands["migrate status"] = &Command{Name: "migrate status", Description: "Show schema migration status", Handler: r.handleMigrateStatus}
// --- Fallback / Panic / Persist ---
r.commands["fallback status"] = &Command{Name: "fallback status", Description: "Show SPoF protection status", Handler: r.handleFallbackStatus}
r.commands["panic stats"] = &Command{Name: "panic stats", Description: "Show panic recovery statistics", Handler: r.handlePanicStats}
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["config show"] = &Command{Name: "config show", Description: "Show current configuration (or section)", Handler: r.handleConfigShow}
r.commands["config get"] = &Command{Name: "config get", Description: "Get a single config value by path", Handler: r.handleConfigGet}
r.commands["config set"] = &Command{Name: "config set", Description: "Set a config value (persisted via Raft or config.toml)", Handler: r.handleConfigSet}
r.commands["config history"] = &Command{Name: "config history", Description: "Show dynamic config change history", Handler: r.handleConfigHistory}
// --- Системные ---
r.commands["help"] = &Command{Name: "help", Description: "Show this help message", Handler: r.handleHelp}
r.commands["clear"] = &Command{Name: "clear", Description: "Clear the screen", Handler: r.handleClear}
r.commands["quit"] = &Command{Name: "quit", Description: "Exit the REPL", Handler: r.handleQuit}
r.commands["exit"] = &Command{Name: "exit", Description: "Exit the REPL", Handler: r.handleQuit}
}
// =============================================================================
// ОСНОВНОЙ ЦИКЛ REPL
// =============================================================================
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 {
r.rl.SetPrompt(r.buildPrompt())
line, err := r.rl.Readline()
if err != nil {
switch {
case errors.Is(err, readline.ErrInterrupt):
utils.Println("^C")
continue
case errors.Is(err, io.EOF):
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
}
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())
}
}
}
}
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
}
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 + ")")
}
if r.currentBatch != nil {
prompt += color.New(color.FgHiMagenta).Sprint(" [batch]")
}
prompt += color.New(color.FgHiCyan).Sprint(":~> ")
return prompt
}
// executeCommand разбирает команду и вызывает соответствующий обработчик.
//
// ДОБАВЛЕНО: команды сортируются по убыванию длины. Это устраняет
// проблему перекрывающихся префиксов.
func (r *Repl) executeCommand(input string) error {
parts := strings.Fields(input)
if len(parts) == 0 {
return nil
}
switch parts[0] {
case "pager":
return r.handlePagerCommand(parts[1:])
}
names := make([]string, 0, len(r.commands))
for name := range r.commands {
names = append(names, name)
}
sort.Slice(names, func(i, j int) bool {
return len(names[i]) > len(names[j])
})
for _, cmdName := range names {
if strings.HasPrefix(input, cmdName) {
cmd := r.commands[cmdName]
args := strings.TrimPrefix(input, cmdName)
args = strings.TrimSpace(args)
argList := strings.Fields(args)
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
// =============================================================================
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]
}
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 {
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
}
// =============================================================================
// ОБРАБОТЧИКИ КОМАНД
// =============================================================================
// -----------------------------------------------------------------------------
// Pager
// -----------------------------------------------------------------------------
func (r *Repl) handlePagerCommand(args []string) error {
if len(args) == 0 {
utils.PrintInfo(fmt.Sprintf("Pager: %v (threshold=%d lines)", r.pager.Enabled(), r.pager.Threshold()))
return nil
}
switch args[0] {
case "on":
r.pager.SetEnabled(true)
utils.PrintSuccess("Pager enabled")
case "off":
r.pager.SetEnabled(false)
utils.PrintSuccess("Pager disabled")
case "size", "threshold":
if len(args) < 2 {
utils.PrintInfo(fmt.Sprintf("Pager threshold: %d lines", r.pager.Threshold()))
return nil
}
n, err := strconv.Atoi(args[1])
if err != nil || n <= 0 {
return fmt.Errorf("invalid threshold '%s' (expected positive integer)", args[1])
}
r.pager.SetThreshold(n)
utils.PrintSuccess(fmt.Sprintf("Pager threshold = %d lines", n))
default:
utils.PrintInfo(fmt.Sprintf("Pager: %v (threshold=%d lines)", r.pager.Enabled(), r.pager.Threshold()))
}
return nil
}
// -----------------------------------------------------------------------------
// Базы данных
// -----------------------------------------------------------------------------
func (r *Repl) handleCreateSlice(args []string) error {
if len(args) < 1 {
return fmt.Errorf("usage: create slice <name>")
}
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-01-2006 15:04:05.000")))
return nil
}
func (r *Repl) handleDropDatabase(args []string) error {
if len(args) < 1 {
return fmt.Errorf("usage: drop database <name>")
}
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-01-2006 15:04:05.000")))
return nil
}
func (r *Repl) handleUseDatabase(args []string) error {
if len(args) < 1 {
return fmt.Errorf("usage: use <database>")
}
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
}
// -----------------------------------------------------------------------------
// Коллекции
// -----------------------------------------------------------------------------
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>")
}
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-01-2006 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>")
}
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
}
// -----------------------------------------------------------------------------
// Документы
// -----------------------------------------------------------------------------
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 <collection> <json>")
}
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-01-2006 15:04:05.000")))
return nil
}
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 <collection> <id>")
}
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-01-2006 15:04:05.000")))
utils.PrintInfo(fmt.Sprintf("updated_at: %s", time.UnixMilli(doc.UpdatedAt).Format("02-01-2006 15:04:05.000")))
if doc.DeletedAt > 0 {
utils.PrintWarning(fmt.Sprintf("deleted_at: %s", time.UnixMilli(doc.DeletedAt).Format("02-01-2006 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 <collection> <index> <value>")
}
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 <collection> <from_date> <to_date>\n" +
" Date formats: YYYY-MM-DD or YYYY-MM-DDTHH:MM:SS (ISO 8601, без пробелов)")
}
if len(args) > 3 {
return fmt.Errorf("too many arguments: dates must not contain spaces.\n" +
" Use YYYY-MM-DD or YYYY-MM-DDTHH:MM:SS (ISO 8601).")
}
collName := args[0]
fromStr := args[1]
toStr := args[2]
formats := []string{
"2006-01-02",
"2006-01-02T15:04:05",
"2006-01-02T15:04:05Z07:00",
}
fromTime, err := parseFlexibleTime(fromStr, formats)
if err != nil {
return fmt.Errorf("invalid from_date '%s': %v", fromStr, err)
}
toTime, err := parseFlexibleTime(toStr, formats)
if err != nil {
return fmt.Errorf("invalid to_date '%s': %v", toStr, err)
}
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
}
func parseFlexibleTime(s string, formats []string) (time.Time, error) {
for _, f := range formats {
if t, err := time.Parse(f, s); err == nil {
return t, nil
}
}
return time.Time{}, fmt.Errorf("expected one of: %s", strings.Join(formats, ", "))
}
func (r *Repl) handleUpdate(args []string) error {
if r.currentDB == "" {
return fmt.Errorf("no database selected")
}
if len(args) < 3 {
return fmt.Errorf("usage: update <collection> <id> <field=value>...")
}
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-01-2006 15:04:05.000")))
return nil
}
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 <collection> <id>")
}
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-01-2006 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 <collection> <id>")
}
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
}
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 <collection> <id>")
}
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-01-2006 15:04:05.000")))
return nil
}
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 <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.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-01-2006 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 <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
}
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
}
// -----------------------------------------------------------------------------
// Временные метки
// -----------------------------------------------------------------------------
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 <collection> <id>")
}
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-01-2006 15:04:05.000"))
t.AddRow("Updated", time.UnixMilli(doc.UpdatedAt).Format("02-01-2006 15:04:05.000"))
if doc.DeletedAt > 0 {
t.AddRow("Deleted", time.UnixMilli(doc.DeletedAt).Format("02-01-2006 15:04:05.000"))
} else {
t.AddRow("Deleted", "not deleted")
}
t.AddRow("Version", fmt.Sprintf("%d", doc.Version))
r.output(t.String())
return nil
}
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 <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.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-01-2006 15:04:05.000"),
time.UnixMilli(minUpdated).Format("02-01-2006 15:04:05.000"))
t.AddRow("Latest",
time.UnixMilli(maxCreated).Format("02-01-2006 15:04:05.000"),
time.UnixMilli(maxUpdated).Format("02-01-2006 15:04:05.000"))
t.AddRow("Average",
time.UnixMilli(avgCreated).Format("02-01-2006 15:04:05.000"),
time.UnixMilli(avgUpdated).Format("02-01-2006 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) < 1 {
return fmt.Errorf("usage: audit filter <data_type> [<operation>] [from=<ms>] [to=<ms>]\n" +
" data_type: DATABASE, COLLECTION, DOCUMENT, FIELD, INDEX, TRANSACTION\n" +
" operation: CREATE, INSERT, UPDATE, DELETE, SOFT_DELETE, RESTORE\n" +
" Example: audit filter DOCUMENT INSERT from=1700000000000 to=1800000000000")
}
dataType := strings.ToUpper(args[0])
operation := ""
var fromTime, toTime int64
for i := 1; i < len(args); i++ {
arg := args[i]
switch {
case strings.HasPrefix(arg, "from="):
val := strings.TrimPrefix(arg, "from=")
parsed, err := strconv.ParseInt(val, 10, 64)
if err != nil {
return fmt.Errorf("invalid from= value '%s': %v", val, err)
}
fromTime = parsed
case strings.HasPrefix(arg, "to="):
val := strings.TrimPrefix(arg, "to=")
parsed, err := strconv.ParseInt(val, 10, 64)
if err != nil {
return fmt.Errorf("invalid to= value '%s': %v", val, err)
}
toTime = parsed
default:
if operation == "" {
operation = strings.ToUpper(arg)
} else {
return fmt.Errorf("unexpected argument: %s", arg)
}
}
}
entries := storage.GetAuditLogFiltered(dataType, operation, fromTime, toTime)
if len(entries) == 0 {
utils.PrintInfo(fmt.Sprintf("No audit log entries found for %s/%s", dataType, operation))
return nil
}
t := NewTable("TIMESTAMP", "OPERATION", "TYPE", "NAME")
for _, entry := range entries {
t.AddRow(entry.TimestampStr, entry.Operation, entry.DataType, entry.Name)
}
r.output(t.String())
return nil
}
// -----------------------------------------------------------------------------
// Индексы
// -----------------------------------------------------------------------------
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 <collection> <name> <fields> [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-01-2006 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 <collection> <name>")
}
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
}
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 <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
}
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 unique, ok := idx["unique"].(bool); ok && unique {
uniqueStr = "true"
}
var createdAt string
switch v := idx["created_at"].(type) {
case int64:
createdAt = time.UnixMilli(v).Format("02-01-2006 15:04:05.000")
case int:
createdAt = time.UnixMilli(int64(v)).Format("02-01-2006 15:04:05.000")
case float64:
createdAt = time.UnixMilli(int64(v)).Format("02-01-2006 15:04:05.000")
default:
createdAt = "unknown"
}
name, _ := idx["name"].(string)
t.AddRow(name, 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 <collection> <field>")
}
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 <collection> <field>")
}
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 <collection> <field> <value>")
}
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 <collection> <field> <value>")
}
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 <collection> <field> <values...>")
}
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-операции
// -----------------------------------------------------------------------------
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)
}
collection := ""
if len(args) >= 1 {
collection = args[0]
}
batch := storage.NewBatch(r.currentDB, collection)
if batch == nil {
return fmt.Errorf("failed to create batch")
}
r.currentBatch = batch
r.currentBatchID = uint64(batch.ID)
utils.PrintSuccess(fmt.Sprintf("Batch %d started at %s (database=%s, collection=%s)",
batch.ID, time.Now().Format("02-01-2006 15:04:05.000"),
batch.Database, batch.Collection))
return nil
}
func (r *Repl) handleCommitTransaction(args []string) error {
if r.currentBatch == nil {
return fmt.Errorf("no active batch to commit")
}
batch := r.currentBatch
if err := batch.Commit(); err != nil {
r.currentBatch = nil
r.currentBatchID = 0
return err
}
r.currentBatch = nil
r.currentBatchID = 0
utils.PrintSuccess(fmt.Sprintf("Batch %d committed successfully at %s",
batch.ID, time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
func (r *Repl) handleRollbackTransaction(args []string) error {
if r.currentBatch == nil {
return fmt.Errorf("no active batch to rollback")
}
batch := r.currentBatch
if bm := storage.GetBatchManager(); bm != nil {
bm.UnregisterBatch(batch)
}
batchID := r.currentBatchID
r.currentBatch = nil
r.currentBatchID = 0
utils.PrintSuccess(fmt.Sprintf("Batch %d aborted at %s (operations discarded)",
batchID, time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
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-01-2006 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
}
// -----------------------------------------------------------------------------
// Плагины
// -----------------------------------------------------------------------------
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-01-2006 15:04:05.000"))
}
r.output(t.String())
return nil
}
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> <filepath>")
}
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-01-2006 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>")
}
name := args[0]
if err := r.pluginManager.UnloadPlugin(name); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Plugin '%s' unloaded", name))
return nil
}
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>")
}
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-01-2006 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>")
}
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 <plugin> <function> [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
}
// -----------------------------------------------------------------------------
// Экспорт/импорт
// -----------------------------------------------------------------------------
func (r *Repl) handleExport(args []string) error {
if len(args) < 2 {
return fmt.Errorf("usage: export <database> <filename>")
}
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-01-2006 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
}
func (r *Repl) handleImport(args []string) error {
if len(args) < 2 {
return fmt.Errorf("usage: import <database> <filename>")
}
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-01-2006 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 — базовая аутентификация и выдача прав
// -----------------------------------------------------------------------------
func (r *Repl) handleACLLogin(args []string) error {
if len(args) < 2 {
return fmt.Errorf("usage: acl login <username> <password>")
}
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-01-2006 15:04:05.000")))
return nil
}
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-01-2006 15:04:05.000")))
return nil
}
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 <collection> <role> <permissions>\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-01-2006 15:04:05.000")))
return nil
}
// handleACLRevoke отзывает разрешения у роли на коллекции.
//
// ДОБАВЛЕНО: симметрично acl grant.
func (r *Repl) handleACLRevoke(args []string) error {
if !r.authenticated || r.currentRole != "admin" {
return fmt.Errorf("permission denied: admin access required")
}
if len(args) < 3 {
return fmt.Errorf("usage: acl revoke <collection> <role> <permissions>\n"+
" Permissions: r=read, w=write, d=delete, a=admin\n"+
" Example: acl revoke users guest rw")
}
collName := args[0]
role := args[1]
perms := args[2]
if r.currentDB == "" {
return fmt.Errorf("no database selected")
}
db, err := r.store.GetDatabase(r.currentDB)
if err != nil {
return err
}
coll, err := db.GetCollection(collName)
if err != nil {
return err
}
revoked := make([]string, 0)
if strings.Contains(perms, "r") {
coll.RevokeACL(role, "read")
revoked = append(revoked, "read")
}
if strings.Contains(perms, "w") {
coll.RevokeACL(role, "write")
revoked = append(revoked, "write")
}
if strings.Contains(perms, "d") {
coll.RevokeACL(role, "delete")
revoked = append(revoked, "delete")
}
if strings.Contains(perms, "a") {
coll.RevokeACL(role, "admin")
revoked = append(revoked, "admin")
}
if len(revoked) == 0 {
return fmt.Errorf("no permissions specified in '%s'", perms)
}
utils.PrintSuccess(fmt.Sprintf("Permissions [%s] revoked from role '%s' on collection '%s' at %s",
strings.Join(revoked, ","), role, collName, time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
func (r *Repl) handleACLUsers(args []string) error {
if r.aclManager == nil {
return fmt.Errorf("ACL manager not initialized")
}
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
}
// -----------------------------------------------------------------------------
// ACL — управление пользователями и ролями через REPL.
// -----------------------------------------------------------------------------
// handleACLUserCreate создаёт нового пользователя.
func (r *Repl) handleACLUserCreate(args []string) error {
if !r.authenticated || r.currentRole != "admin" {
return fmt.Errorf("permission denied: admin access required")
}
if r.aclManager == nil {
return fmt.Errorf("ACL manager not initialized")
}
if len(args) < 2 {
return fmt.Errorf("usage: acl user create <username> <password> [role1,role2,...]\n" +
" Example: acl user create alice secret123 admin")
}
username := args[0]
password := args[1]
var roles []string
if len(args) >= 3 {
roles = strings.Split(args[2], ",")
for i := range roles {
roles[i] = strings.TrimSpace(roles[i])
}
}
if err := r.aclManager.CreateUser(username, password, roles); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("User '%s' created with roles %v at %s",
username, roles, time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
// handleACLUserDelete удаляет пользователя.
func (r *Repl) handleACLUserDelete(args []string) error {
if !r.authenticated || r.currentRole != "admin" {
return fmt.Errorf("permission denied: admin access required")
}
if r.aclManager == nil {
return fmt.Errorf("ACL manager not initialized")
}
if len(args) < 1 {
return fmt.Errorf("usage: acl user delete <username>")
}
username := args[0]
if err := r.aclManager.DeleteUser(username); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("User '%s' deleted", username))
return nil
}
// handleACLRoleCreate создаёт новую роль.
func (r *Repl) handleACLRoleCreate(args []string) error {
if !r.authenticated || r.currentRole != "admin" {
return fmt.Errorf("permission denied: admin access required")
}
if r.aclManager == nil {
return fmt.Errorf("ACL manager not initialized")
}
if len(args) < 1 {
return fmt.Errorf("usage: acl role create <name>")
}
roleName := args[0]
if err := r.aclManager.CreateRole(roleName); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Role '%s' created at %s",
roleName, time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
// handleACLRoleDelete удаляет роль.
func (r *Repl) handleACLRoleDelete(args []string) error {
if !r.authenticated || r.currentRole != "admin" {
return fmt.Errorf("permission denied: admin access required")
}
if r.aclManager == nil {
return fmt.Errorf("ACL manager not initialized")
}
if len(args) < 1 {
return fmt.Errorf("usage: acl role delete <name>")
}
roleName := args[0]
if err := r.aclManager.DeleteRole(roleName); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Role '%s' deleted", roleName))
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
}
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 <name>")
}
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-01-2006 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 <collection> <id>")
}
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
}
// handleCompressionConfig показывает или изменяет настройки сжатия.
//
// РАСШИРЕНО: подкоманда set делегирует в handleConfigSet —
// изменения сохраняются через Raft или config.toml.
func (r *Repl) handleCompressionConfig(args []string) error {
if len(args) >= 1 && args[0] == "set" {
if len(args) < 3 {
return fmt.Errorf("usage: compression config set <key> <value>\n" +
" keys: enabled (true|false), algorithm (snappy|brotli|zstd), level (1-9), min_size (bytes)")
}
key := strings.ToLower(args[1])
value := args[2]
var path string
switch key {
case "enabled":
path = "compression.enabled"
case "algorithm":
path = "compression.algorithm"
case "level":
path = "compression.level"
case "min_size":
path = "compression.min_size"
default:
return fmt.Errorf("unknown key '%s' (supported: enabled, algorithm, level, min_size)", key)
}
return r.handleConfigSet([]string{path, value})
}
cfg := r.effectiveConfig()
if cfg == nil {
return fmt.Errorf("config is not available")
}
t := NewTable("SETTING", "VALUE")
t.AddRow("Enabled", fmt.Sprintf("%v", cfg.Compression.Enabled))
t.AddRow("Algorithm", cfg.Compression.Algorithm)
t.AddRow("Level", fmt.Sprintf("%d", cfg.Compression.Level))
t.AddRow("Min Size", utils.FormatBytes(int64(cfg.Compression.MinSize)))
t.AddRow("snappy", "Fast compression/decompression, good balance")
t.AddRow("brotli", "Best compression ratio (recommended)")
t.AddRow("zstd", "High compression ratio, configurable level")
r.output(t.String())
return nil
}
// -----------------------------------------------------------------------------
// Триггеры
// -----------------------------------------------------------------------------
// handleCreateTrigger создаёт триггер на коллекции.
//
// ИСПРАВЛЕНО: теперь парсятся все опции, заявленные в usage.
func (r *Repl) handleCreateTrigger(args []string) error {
if r.currentDB == "" {
return fmt.Errorf("no database selected")
}
if len(args) < 4 {
return fmt.Errorf("usage: create trigger <collection> <name> <event> <action> [options]\n" +
" Events: BEFORE_INSERT, AFTER_INSERT, BEFORE_UPDATE, AFTER_UPDATE, BEFORE_DELETE, AFTER_DELETE\n" +
" Actions: abort, skip, modify, log, notify\n" +
" Options:\n" +
" --description <text>\n" +
" --set <field> <value>\n" +
" --inc <field> <number>\n" +
" --currentDate <field>\n" +
" --condition <field> <op> <value> (op: eq, ne, gt, lt, gte, lte, in, nin, exists, regex)")
}
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: "",
Operations: make([]storage.TriggerOperation, 0),
CreatedAt: time.Now().UnixMilli(),
UpdatedAt: time.Now().UnixMilli(),
}
i := 4
for i < len(args) {
opt := args[i]
switch opt {
case "--description":
if i+1 >= len(args) {
return fmt.Errorf("--description requires a value")
}
trigger.Description = args[i+1]
i += 2
case "--set":
if i+2 >= len(args) {
return fmt.Errorf("--set requires <field> <value>")
}
trigger.Operations = append(trigger.Operations, storage.TriggerOperation{
Type: "set",
Field: args[i+1],
Value: args[i+2],
})
i += 3
case "--inc":
if i+2 >= len(args) {
return fmt.Errorf("--inc requires <field> <number>")
}
var incVal float64
if _, err := fmt.Sscanf(args[i+2], "%f", &incVal); err != nil {
return fmt.Errorf("--inc: invalid number '%s'", args[i+2])
}
trigger.Operations = append(trigger.Operations, storage.TriggerOperation{
Type: "inc",
Field: args[i+1],
Value: incVal,
})
i += 3
case "--currentDate", "--currentdate":
if i+1 >= len(args) {
return fmt.Errorf("--currentDate requires <field>")
}
trigger.Operations = append(trigger.Operations, storage.TriggerOperation{
Type: "currentDate",
Field: args[i+1],
})
i += 2
case "--condition":
if i+3 >= len(args) {
return fmt.Errorf("--condition requires <field> <op> <value>")
}
cond := &storage.TriggerCondition{
Field: args[i+1],
Operator: args[i+2],
Value: args[i+3],
}
if f, err := strconv.ParseFloat(args[i+3], 64); err == nil {
cond.Value = f
} else if b, err := strconv.ParseBool(args[i+3]); err == nil {
cond.Value = b
}
trigger.Condition = cond
i += 4
default:
return fmt.Errorf("unknown option: %s", opt)
}
}
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-01-2006 15:04:05.000")))
return nil
}
// handleDropTrigger удаляет триггер из коллекции.
//
// ИСПРАВЛЕНО: убран обязательный аргумент <event>.
func (r *Repl) handleDropTrigger(args []string) error {
if r.currentDB == "" {
return fmt.Errorf("no database selected")
}
if len(args) < 2 {
return fmt.Errorf("usage: drop trigger <collection> <name>")
}
collName := args[0]
triggerName := 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.DropTrigger(triggerName); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Trigger '%s' dropped from collection '%s'", triggerName, collName))
return nil
}
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 <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
}
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-01-2006 15:04:05.000"), tr.Description)
}
r.output(t.String())
return nil
}
// handleEnableTrigger включает триггер.
//
// ИСПРАВЛЕНО: убран обязательный аргумент <event>.
func (r *Repl) handleEnableTrigger(args []string) error {
if r.currentDB == "" {
return fmt.Errorf("no database selected")
}
if len(args) < 2 {
return fmt.Errorf("usage: enable trigger <collection> <name>")
}
collName := args[0]
triggerName := 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.EnableTrigger(triggerName); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Trigger '%s' enabled at %s", triggerName, time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
// handleDisableTrigger отключает триггер.
//
// ИСПРАВЛЕНО: убран обязательный аргумент <event>.
func (r *Repl) handleDisableTrigger(args []string) error {
if r.currentDB == "" {
return fmt.Errorf("no database selected")
}
if len(args) < 2 {
return fmt.Errorf("usage: disable trigger <collection> <name>")
}
collName := args[0]
triggerName := 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.DisableTrigger(triggerName); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Trigger '%s' disabled", triggerName))
return nil
}
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-01-2006 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
}
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-01-2006 15:04:05.000")
t.AddRow(prefix, node.ID, fmt.Sprintf("%s:%d", nodeIP, nodePort), node.Status, lastSeen)
}
r.output(t.String())
return nil
}
// handleClusterPipeline показывает статистику pipeline репликатора.
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")
keys := make([]string, 0, len(stats))
for k := range stats {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
t.AddRow(k, fmt.Sprintf("%v", stats[k]))
}
r.output(t.String())
return nil
}
// handleClusterPipelineOn включает pipeline replication.
func (r *Repl) handleClusterPipelineOn(args []string) error {
return r.handleConfigSet([]string{"performance.enable_pipeline", "true"})
}
// handleClusterPipelineOff выключает pipeline replication.
func (r *Repl) handleClusterPipelineOff(args []string) error {
return r.handleConfigSet([]string{"performance.enable_pipeline", "false"})
}
func (r *Repl) handleClusterReshard(args []string) error {
if r.coordinator == nil {
return fmt.Errorf("cluster coordinator not available")
}
if err := r.coordinator.TriggerResharding("manual"); err != nil {
return err
}
utils.PrintSuccess(fmt.Sprintf("Resharding triggered at %s", time.Now().Format("02-01-2006 15:04:05.000")))
return nil
}
// -----------------------------------------------------------------------------
// Миграция
// -----------------------------------------------------------------------------
func (r *Repl) handleMigration(args []string) error {
if len(args) == 0 {
return fmt.Errorf("usage: migration <subcommand>\n"+
" Subcommands:\n"+
" start <source_dc> <target_dc> [database] [collection] - Start a new migration\n"+
" status [task_id] - Show migration status\n"+
" list - List all migration tasks\n"+
" pause <task_id> - Pause a running migration\n"+
" resume <task_id> - Resume a paused migration\n"+
" cancel <task_id> - 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 <source_dc> <target_dc> [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 <task_id>")
}
return r.handleMigrationPause(migrator, args[1])
case "resume":
if len(args) < 2 {
return fmt.Errorf("usage: migration resume <task_id>")
}
return r.handleMigrationResume(migrator, args[1])
case "cancel":
if len(args) < 2 {
return fmt.Errorf("usage: migration cancel <task_id>")
}
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)
}
}
func (r *Repl) handleMigrationStart(migrator *cluster.CrossDCMigrator, args []string) error {
if len(args) < 2 {
return fmt.Errorf("usage: migration start <source_dc> <target_dc> [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-01-2006 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
}
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-01-2006 15:04:05.000"))
}
if task.EndTime > 0 {
t.AddRow("Completed", time.UnixMilli(task.EndTime).Format("02-01-2006 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
}
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
}
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
}
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
}
// handleMigrateStatus показывает статус миграций схемы.
//
// ИСПРАВЛЕНО: раньше выводилось `fmt.Sprintf("%v", status)`, где
// status — *migration.MigrationStatus. Теперь выводятся поля.
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("Current Version", status.CurrentVersion)
t.AddRow("Total Migrations", fmt.Sprintf("%d", status.TotalMigrations))
t.AddRow("Applied", fmt.Sprintf("%d", status.AppliedMigrations))
t.AddRow("Pending", fmt.Sprintf("%d", status.PendingMigrations))
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")
keys := make([]string, 0, len(stats))
for k := range stats {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
t.AddRow(k, fmt.Sprintf("%v", stats[k]))
}
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")
keys := make([]string, 0, len(stats))
for k := range stats {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
t.AddRow(k, fmt.Sprintf("%v", stats[k]))
}
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")
keys := make([]string, 0, len(info))
for k := range info {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
t.AddRow(k, fmt.Sprintf("%v", info[k]))
}
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")
keys := make([]string, 0, len(metrics))
for k := range metrics {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
t.AddRow(k, fmt.Sprintf("%v", metrics[k]))
}
r.output(t.String())
return nil
}
// handleSagaList показывает активные (или все) SAGA транзакции.
//
// РАСШИРЕНО: поддерживает флаг --all для показа завершённых/отменённых SAGA.
//
// ИСПРАВЛЕНО (2026-10): тип элементов — *cluster.SagaTransaction, потому
// что r.coordinator.GetSagaManager() возвращает *cluster.SagaManager,
// а его методы работают с локальным типом cluster.SagaTransaction.
func (r *Repl) handleSagaList(args []string) error {
if r.coordinator == nil {
return fmt.Errorf("cluster coordinator not available")
}
sm := r.coordinator.GetSagaManager()
if sm == nil {
return fmt.Errorf("saga manager not available")
}
showAll := false
for _, a := range args {
if a == "--all" {
showAll = true
}
}
var sagas []*cluster.SagaTransaction
if showAll {
sagas = sm.GetAllSagas()
if len(sagas) == 0 {
utils.PrintInfo("No SAGA transactions found")
return nil
}
} else {
sagas = sm.GetActiveSagas()
if len(sagas) == 0 {
utils.PrintInfo("No active SAGA transactions (use --all to show all)")
return nil
}
}
t := NewTable("ID", "STATUS", "STEPS", "CURRENT_STEP", "CREATED_AT", "UPDATED_AT")
for _, saga := range sagas {
createdAt := ""
if saga.CreatedAt > 0 {
createdAt = time.UnixMilli(saga.CreatedAt).Format("02-01-2006 15:04:05.000")
}
updatedAt := ""
if saga.UpdatedAt > 0 {
updatedAt = time.UnixMilli(saga.UpdatedAt).Format("02-01-2006 15:04:05.000")
}
t.AddRow(
saga.ID,
saga.Status,
fmt.Sprintf("%d", len(saga.Steps)),
fmt.Sprintf("%d", saga.CurrentStep),
createdAt,
updatedAt,
)
}
r.output(t.String())
return nil
}
// -----------------------------------------------------------------------------
// Динамическая конфигурация (через Raft).
// -----------------------------------------------------------------------------
func (r *Repl) handleConfigShow(args []string) error {
cfg := r.effectiveConfig()
if cfg == nil {
return fmt.Errorf("config is not available")
}
section := ""
if len(args) > 0 {
section = strings.ToLower(args[0])
}
switch section {
case "", "all":
t := NewTable("SECTION", "KEY", "VALUE")
addRows := func(sectionName, prefix string, pairs map[string]interface{}) {
keys := make([]string, 0, len(pairs))
for k := range pairs {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
t.AddRow(sectionName, prefix+"."+k, fmt.Sprintf("%v", pairs[k]))
}
}
addRows("cluster", "cluster", map[string]interface{}{
"name": cfg.Cluster.Name,
"node_ip": cfg.Cluster.NodeIP,
"node_port": cfg.Cluster.NodePort,
"raft_port": cfg.Cluster.RaftPort,
"region": cfg.Cluster.Region,
})
addRows("api", "api", map[string]interface{}{
"port": cfg.API.Port,
})
addRows("compression", "compression", map[string]interface{}{
"enabled": cfg.Compression.Enabled,
"algorithm": cfg.Compression.Algorithm,
"level": cfg.Compression.Level,
"min_size": cfg.Compression.MinSize,
})
addRows("log", "log", map[string]interface{}{
"log_file": cfg.Log.LogFile,
"log_level": cfg.Log.LogLevel,
})
addRows("replication", "replication", map[string]interface{}{
"enabled": cfg.Replication.Enabled,
"sync_replication": cfg.Replication.SyncReplication,
"replication_timeout_ms": cfg.Replication.ReplicationTimeoutMs,
"max_replica_lag_ms": cfg.Replication.MaxReplicaLagMs,
})
addRows("wal", "wal", map[string]interface{}{
"enabled": cfg.WAL.Enabled,
"segment_size_mb": cfg.WAL.SegmentSizeMB,
"sync_interval_sec": cfg.WAL.SyncIntervalSec,
"batch_size": cfg.WAL.BatchSize,
})
addRows("saga", "saga", map[string]interface{}{
"enabled": cfg.Saga.Enabled,
"coordinator_count": cfg.Saga.CoordinatorCount,
"max_retries": cfg.Saga.MaxRetries,
})
addRows("backpressure", "backpressure", map[string]interface{}{
"enabled": cfg.Backpressure.Enabled,
"cpu_threshold": cfg.Backpressure.CPUThreshold,
"memory_threshold": cfg.Backpressure.MemoryThreshold,
})
addRows("monitoring", "monitoring", map[string]interface{}{
"enable_metrics": cfg.Monitoring.EnableMetrics,
"metrics_port": cfg.Monitoring.MetricsPort,
"enable_tracing": cfg.Monitoring.EnableTracing,
"trace_sample_rate": cfg.Monitoring.TraceSampleRate,
})
addRows("performance", "performance", map[string]interface{}{
"enable_pipeline": cfg.Performance.EnablePipeline,
"batch_size": cfg.Performance.BatchSize,
"max_connections": cfg.Performance.MaxConnections,
"read_from_follower": cfg.Performance.ReadFromFollower,
})
r.output(t.String())
case "cluster":
utils.PrintInfo(fmt.Sprintf("cluster.name = %s", cfg.Cluster.Name))
utils.PrintInfo(fmt.Sprintf("cluster.node_ip = %s", cfg.Cluster.NodeIP))
utils.PrintInfo(fmt.Sprintf("cluster.node_port = %d", cfg.Cluster.NodePort))
utils.PrintInfo(fmt.Sprintf("cluster.raft_port = %d", cfg.Cluster.RaftPort))
utils.PrintInfo(fmt.Sprintf("cluster.region = %s", cfg.Cluster.Region))
case "compression":
t := NewTable("KEY", "VALUE")
t.AddRow("compression.enabled", fmt.Sprintf("%v", cfg.Compression.Enabled))
t.AddRow("compression.algorithm", cfg.Compression.Algorithm)
t.AddRow("compression.level", fmt.Sprintf("%d", cfg.Compression.Level))
t.AddRow("compression.min_size", fmt.Sprintf("%d", cfg.Compression.MinSize))
r.output(t.String())
case "replication":
t := NewTable("KEY", "VALUE")
t.AddRow("replication.enabled", fmt.Sprintf("%v", cfg.Replication.Enabled))
t.AddRow("replication.sync_replication", fmt.Sprintf("%v", cfg.Replication.SyncReplication))
t.AddRow("replication.replication_timeout_ms", fmt.Sprintf("%d", cfg.Replication.ReplicationTimeoutMs))
t.AddRow("replication.max_replica_lag_ms", fmt.Sprintf("%d", cfg.Replication.MaxReplicaLagMs))
r.output(t.String())
case "wal":
t := NewTable("KEY", "VALUE")
t.AddRow("wal.enabled", fmt.Sprintf("%v", cfg.WAL.Enabled))
t.AddRow("wal.segment_size_mb", fmt.Sprintf("%d", cfg.WAL.SegmentSizeMB))
t.AddRow("wal.sync_interval_sec", fmt.Sprintf("%d", cfg.WAL.SyncIntervalSec))
t.AddRow("wal.batch_size", fmt.Sprintf("%d", cfg.WAL.BatchSize))
t.AddRow("wal.recovery_workers", fmt.Sprintf("%d", cfg.WAL.RecoveryWorkers))
r.output(t.String())
case "saga":
t := NewTable("KEY", "VALUE")
t.AddRow("saga.enabled", fmt.Sprintf("%v", cfg.Saga.Enabled))
t.AddRow("saga.coordinator_count", fmt.Sprintf("%d", cfg.Saga.CoordinatorCount))
t.AddRow("saga.max_retries", fmt.Sprintf("%d", cfg.Saga.MaxRetries))
t.AddRow("saga.saga_timeout_sec", fmt.Sprintf("%d", cfg.Saga.SagaTimeoutSec))
r.output(t.String())
case "monitoring":
t := NewTable("KEY", "VALUE")
t.AddRow("monitoring.enable_metrics", fmt.Sprintf("%v", cfg.Monitoring.EnableMetrics))
t.AddRow("monitoring.metrics_port", fmt.Sprintf("%d", cfg.Monitoring.MetricsPort))
t.AddRow("monitoring.enable_tracing", fmt.Sprintf("%v", cfg.Monitoring.EnableTracing))
t.AddRow("monitoring.trace_sample_rate", fmt.Sprintf("%v", cfg.Monitoring.TraceSampleRate))
r.output(t.String())
case "performance":
t := NewTable("KEY", "VALUE")
t.AddRow("performance.enable_pipeline", fmt.Sprintf("%v", cfg.Performance.EnablePipeline))
t.AddRow("performance.batch_size", fmt.Sprintf("%d", cfg.Performance.BatchSize))
t.AddRow("performance.max_connections", fmt.Sprintf("%d", cfg.Performance.MaxConnections))
t.AddRow("performance.read_from_follower", fmt.Sprintf("%v", cfg.Performance.ReadFromFollower))
r.output(t.String())
case "backpressure":
t := NewTable("KEY", "VALUE")
t.AddRow("backpressure.enabled", fmt.Sprintf("%v", cfg.Backpressure.Enabled))
t.AddRow("backpressure.cpu_threshold", fmt.Sprintf("%v", cfg.Backpressure.CPUThreshold))
t.AddRow("backpressure.memory_threshold", fmt.Sprintf("%v", cfg.Backpressure.MemoryThreshold))
t.AddRow("backpressure.queue_size_threshold", fmt.Sprintf("%d", cfg.Backpressure.QueueSizeThreshold))
t.AddRow("backpressure.connection_threshold", fmt.Sprintf("%d", cfg.Backpressure.ConnectionThreshold))
r.output(t.String())
default:
return fmt.Errorf("unknown config section '%s' (supported: cluster, compression, replication, wal, saga, monitoring, performance, backpressure)", section)
}
return nil
}
func (r *Repl) handleConfigGet(args []string) error {
if len(args) < 1 {
return fmt.Errorf("usage: config get <path> (e.g. compression.algorithm)")
}
path := strings.Join(args, " ")
cfg := r.effectiveConfig()
if cfg == nil {
return fmt.Errorf("config is not available")
}
value, err := getConfigValue(cfg, path)
if err != nil {
return err
}
utils.PrintInfo(fmt.Sprintf("%s = %v", path, value))
return nil
}
func (r *Repl) handleConfigSet(args []string) error {
if !r.authenticated || r.currentRole != "admin" {
return fmt.Errorf("permission denied: admin access required")
}
if len(args) < 2 {
return fmt.Errorf("usage: config set <path> <value> (e.g. compression.algorithm brotli)")
}
path := args[0]
valueStr := strings.Join(args[1:], " ")
cfg := r.effectiveConfig()
if cfg == nil {
return fmt.Errorf("config is not available")
}
currentValue, err := getConfigValue(cfg, path)
if err != nil {
return err
}
value, err := coerceStringToType(valueStr, currentValue)
if err != nil {
return fmt.Errorf("cannot parse value for %s: %v", path, err)
}
change := map[string]interface{}{path: value}
if r.dynamicConfigManager != nil {
err := r.dynamicConfigManager.ApplyChange(change,
fmt.Sprintf("REPL: set %s", path), r.currentUser)
if err == nil {
r.applyLocalConfigChange(path, value)
utils.PrintSuccess(fmt.Sprintf("config %s = %v (applied via Raft)", path, value))
return nil
}
utils.PrintWarning(fmt.Sprintf("Raft apply failed (%v), falling back to local config.toml", err))
}
r.applyLocalConfigChange(path, value)
if r.configPath != "" {
if err := r.persistConfigFile(); err != nil {
utils.PrintWarning(fmt.Sprintf("config applied locally, but saving to %s failed: %v", r.configPath, err))
} else {
utils.PrintSuccess(fmt.Sprintf("config %s = %v (saved to %s)", path, value, r.configPath))
return nil
}
}
utils.PrintSuccess(fmt.Sprintf("config %s = %v (in-memory only, no configPath set)", path, value))
return nil
}
func (r *Repl) handleConfigHistory(args []string) error {
if r.dynamicConfigManager == nil {
utils.PrintInfo("Dynamic config manager is not available (single-node mode or Raft disabled)")
return nil
}
history := r.dynamicConfigManager.GetChangeHistory()
if len(history) == 0 {
utils.PrintInfo("No config changes recorded")
return nil
}
t := NewTable("VERSION", "TIMESTAMP", "CHANGED_BY", "DESCRIPTION")
for _, ch := range history {
t.AddRow(
fmt.Sprintf("%d", ch.Version),
time.UnixMilli(ch.Timestamp).Format("02-01-2006 15:04:05.000"),
ch.ChangedBy,
ch.Description,
)
}
r.output(t.String())
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 <name>", "Create a new database (slice)"},
{"drop database <name>", "Delete an existing database"},
{"use <database>", "Switch to a specific database"},
{"show databases", "List all available databases"},
},
"Collection Management": {
{"create collection <name>", "Create a new collection in current database"},
{"drop collection <name>", "Delete a collection from current database"},
{"show collections", "List all collections in current database"},
},
"Document Operations": {
{"insert <collection> <json>", "Insert a new document (JSON format: key=value,key2=value2)"},
{"find <collection> <id>", "Find a document by its ID"},
{"findbyindex <collection> <index> <value>", "Find documents using an index"},
{"findbytime <collection> <from_date> <to_date>", "Find documents by time range (ISO 8601, e.g. 2026-01-01T12:00:00)"},
{"update <collection> <id> <field=value>...", "Update fields of an existing document"},
{"delete <collection> <id>", "Delete a document (soft delete if enabled)"},
{"permanent delete <collection> <id>", "Permanently delete a soft-deleted document"},
{"restore <collection> <id>", "Restore a soft-deleted document"},
{"show deleted <collection>", "Show soft-deleted documents in a collection"},
{"count <collection>", "Count total documents in a collection"},
},
"Timestamp Management": {
{"show timestamps <collection> <id>", "Show timestamps for a document"},
{"stats timestamps <collection>", "Show timestamp statistics for a collection"},
{"audit log", "Show audit log"},
{"audit filter <type> [<op>] [from=<ms>] [to=<ms>]", "Filter audit log by type, operation and time"},
},
"Index Management": {
{"create index <collection> <name> <fields> [unique]", "Create a new index on specified fields"},
{"drop index <collection> <name>", "Remove an existing index"},
{"show indexes <collection>", "List all indexes on a collection"},
},
"Constraints": {
{"add required <collection> <field>", "Add a required field constraint"},
{"add unique <collection> <field>", "Add a unique constraint on a field"},
{"add min <collection> <field> <value>", "Add a minimum value constraint for numeric fields"},
{"add max <collection> <field> <value>", "Add a maximum value constraint for numeric fields"},
{"add enum <collection> <field> <values...>", "Add allowed values constraint (enum)"},
},
"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 <collection> <name> <event> <action> [options]", "Create a trigger on collection events"},
{" --description <text>", " Set trigger description"},
{" --set <field> <value>", " Set field value when trigger fires"},
{" --inc <field> <number>", " Increment numeric field"},
{" --currentDate <field>", " Set field to current timestamp"},
{" --condition <field> <op> <value>", " Only fire if condition matches (op: eq, ne, gt, lt, gte, lte, in, nin, exists, regex)"},
{"drop trigger <collection> <name>", "Remove a trigger from collection"},
{"show triggers <collection>", "List all triggers on a collection"},
{"enable trigger <collection> <name>", "Enable a disabled trigger"},
{"disable trigger <collection> <name>", "Disable a trigger without removing it"},
{"trigger log", "Show trigger execution history"},
},
"Plugins": {
{"plugin list", "List all loaded plugins"},
{"plugin load <name> <filepath>", "Load a plugin from Lua file"},
{"plugin unload <name>", "Unload a plugin"},
{"plugin start <name>", "Start a loaded plugin"},
{"plugin stop <name>", "Stop a running plugin"},
{"plugin exec <plugin> <function> [args...]", "Execute a plugin function"},
},
"Import/Export": {
{"export <database> <filename>", "Export entire database to file"},
{"import <database> <filename>", "Import database from file"},
},
"Compression": {
{"compression stats", "Show compression statistics for current database"},
{"compression config", "Display current compression settings"},
{"compression config set <key> <value>", "Change setting (enabled, algorithm, level, min_size)"},
{"compress collection <name>", "Compress all documents in a collection (frees memory)"},
{"doc compression <collection> <id>", "Show compression info for a specific document"},
},
"Access Control": {
{"acl login <username> <password>", "Authenticate with username and password"},
{"acl logout", "Logout current user session"},
{"acl grant <collection> <role> <permissions>", "Grant permissions (r=read,w=write,d=delete,a=admin)"},
{"acl revoke <collection> <role> <permissions>", "Revoke permissions (r=read,w=write,d=delete,a=admin)"},
{"acl users", "List all users"},
{"acl roles", "List all roles"},
},
"Users & Roles (admin only)": {
{"acl user create <user> <pass> [roles]", "Create a new user (roles separated by commas)"},
{"acl user delete <user>", "Delete a user"},
{"acl role create <name>", "Create a new role"},
{"acl role delete <name>", "Delete a role"},
},
"Dynamic Config (admin only)": {
{"config show [section]", "Show current configuration or a specific section"},
{"config get <path>", "Get a single config value (e.g. compression.algorithm)"},
{"config set <path> <value>", "Set a config value (persisted via Raft or config.toml)"},
{"config history", "Show dynamic config change history"},
},
"Cluster": {
{"status", "Show current cluster status and role (leader/follower)"},
{"nodes", "List all nodes in the cluster"},
{"cluster pipeline", "Show pipeline replicator statistics"},
{"cluster pipeline on", "Enable pipeline replication (config set)"},
{"cluster pipeline off", "Disable pipeline replication (config set)"},
{"cluster reshard", "Trigger manual resharding"},
},
"Migration": {
{"migration start <source_dc> <target_dc> [db] [coll]", "Start cross-datacenter migration"},
{"migration status [task_id]", "Show migration status"},
{"migration list", "List all migration tasks"},
{"migration pause <task_id>", "Pause a running migration"},
{"migration resume <task_id>", "Resume a paused migration"},
{"migration cancel <task_id>", "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 [--all]", "List active (or all) SAGA transactions"},
},
"Pager & Input": {
{"pager on|off", "Enable/disable pager"},
{"pager size <N>", "Set pager threshold (lines) — alias of 'pager threshold'"},
{"pager threshold <N>", "Set pager threshold (lines)"},
{"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)"},
},
}
categoryNames := make([]string, 0, len(categories))
for name := range categories {
categoryNames = append(categoryNames, name)
}
sort.Strings(categoryNames)
for _, category := range categoryNames {
commands := categories[category]
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
}
func (r *Repl) handleQuit(args []string) error {
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-01-2006 15:04:05.000")))
if r.rl != nil {
_ = r.rl.Close()
}
os.Exit(0)
return nil
}
// =============================================================================
// ВСПОМОГАТЕЛЬНЫЕ ФУНКЦИИ ДЛЯ CONFIG
// =============================================================================
// effectiveConfig возвращает конфигурацию, которую видит REPL.
func (r *Repl) effectiveConfig() *config.Config {
if r.dynamicConfigManager != nil {
if cfg := r.dynamicConfigManager.GetConfig(); cfg != nil {
return cfg
}
}
return r.config
}
// applyLocalConfigChange применяет изменение только к in-memory r.config.
func (r *Repl) applyLocalConfigChange(path string, value interface{}) {
if r.config == nil {
return
}
_ = setConfigFieldLocal(r.config, path, value)
}
// persistConfigFile сохраняет текущий r.config в r.configPath.
// Делегирует в config.SaveConfig.
func (r *Repl) persistConfigFile() error {
if r.configPath == "" {
return fmt.Errorf("configPath is not set")
}
return config.SaveConfig(r.configPath, r.config)
}
// getConfigValue читает значение по пути из конфигурации.
func getConfigValue(cfg *config.Config, path string) (interface{}, error) {
switch path {
case "cluster.name":
return cfg.Cluster.Name, nil
case "cluster.node_ip":
return cfg.Cluster.NodeIP, nil
case "cluster.node_port":
return cfg.Cluster.NodePort, nil
case "cluster.raft_port":
return cfg.Cluster.RaftPort, nil
case "cluster.heartbeat_timeout_ms":
return cfg.Cluster.HeartbeatTimeoutMs, nil
case "cluster.election_timeout_ms":
return cfg.Cluster.ElectionTimeoutMs, nil
case "cluster.commit_timeout_ms":
return cfg.Cluster.CommitTimeoutMs, nil
case "cluster.snapshot_interval_min":
return cfg.Cluster.SnapshotIntervalMin, nil
case "cluster.snapshot_threshold":
return cfg.Cluster.SnapshotThreshold, nil
case "cluster.recovery_timeout_sec":
return cfg.Cluster.RecoveryTimeoutSec, nil
case "cluster.split_brain_prevention":
return cfg.Cluster.SplitBrainPrevention, nil
case "cluster.region":
return cfg.Cluster.Region, nil
case "cluster.priority_zone":
return cfg.Cluster.PriorityZone, nil
case "storage.page_size_mb":
return cfg.Storage.PageSizeMB, nil
case "storage.max_collections":
return cfg.Storage.MaxCollections, nil
case "storage.max_documents_per_collection":
return cfg.Storage.MaxDocumentsPerCollection, nil
case "storage.default_engine":
return cfg.Storage.DefaultEngine, nil
case "storage.enable_custom_engines":
return cfg.Storage.EnableCustomEngines, nil
case "replication.enabled":
return cfg.Replication.Enabled, nil
case "replication.sync_replication":
return cfg.Replication.SyncReplication, nil
case "replication.replication_timeout_ms":
return cfg.Replication.ReplicationTimeoutMs, nil
case "replication.max_replica_lag_ms":
return cfg.Replication.MaxReplicaLagMs, nil
case "wal.enabled":
return cfg.WAL.Enabled, nil
case "wal.segment_size_mb":
return cfg.WAL.SegmentSizeMB, nil
case "wal.sync_interval_sec":
return cfg.WAL.SyncIntervalSec, nil
case "wal.batch_size":
return cfg.WAL.BatchSize, nil
case "wal.async_recovery":
return cfg.WAL.AsyncRecovery, nil
case "wal.recovery_workers":
return cfg.WAL.RecoveryWorkers, nil
case "wal.async_recovery_workers":
return cfg.WAL.AsyncRecoveryWorkers, nil
case "wal.async_recovery_buffer":
return cfg.WAL.AsyncRecoveryBuffer, nil
case "mvcc.max_versions_per_doc":
return cfg.MVCC.MaxVersionsPerDoc, nil
case "mvcc.visibility_map_size":
return cfg.MVCC.VisibilityMapSize, nil
case "mvcc.prune_interval_min":
return cfg.MVCC.PruneIntervalMin, nil
case "mvcc.retention_days":
return cfg.MVCC.RetentionDays, nil
case "mvcc.read_cache_size":
return cfg.MVCC.ReadCacheSize, nil
case "mvcc.read_cache_ttl_sec":
return cfg.MVCC.ReadCacheTTLSec, nil
case "saga.enabled":
return cfg.Saga.Enabled, nil
case "saga.coordinator_count":
return cfg.Saga.CoordinatorCount, nil
case "saga.max_retries":
return cfg.Saga.MaxRetries, nil
case "saga.saga_timeout_sec":
return cfg.Saga.SagaTimeoutSec, nil
case "saga.retry_backoff_ms":
return cfg.Saga.RetryBackoffMs, nil
case "saga.stuck_check_interval_sec":
return cfg.Saga.StuckCheckIntervalSec, nil
case "backpressure.enabled":
return cfg.Backpressure.Enabled, nil
case "backpressure.cpu_threshold":
return cfg.Backpressure.CPUThreshold, nil
case "backpressure.memory_threshold":
return cfg.Backpressure.MemoryThreshold, nil
case "backpressure.queue_size_threshold":
return cfg.Backpressure.QueueSizeThreshold, nil
case "backpressure.connection_threshold":
return cfg.Backpressure.ConnectionThreshold, nil
case "api.port":
return cfg.API.Port, nil
case "log.log_level":
return cfg.Log.LogLevel, nil
case "log.log_file":
return cfg.Log.LogFile, nil
case "plugins.enabled":
return cfg.Plugins.Enabled, nil
case "plugins.max_cpu_time_ms":
return cfg.Plugins.MaxCPUTimeMs, nil
case "plugins.max_memory_mb":
return cfg.Plugins.MaxMemoryMB, nil
case "plugins.hot_reload_interval_sec":
return cfg.Plugins.HotReloadIntervalSec, nil
case "plugins.max_lua_states":
return cfg.Plugins.MaxLuaStates, nil
case "autoscaling.enabled":
return cfg.Autoscaling.Enabled, nil
case "autoscaling.min_nodes":
return cfg.Autoscaling.MinNodes, nil
case "autoscaling.max_nodes":
return cfg.Autoscaling.MaxNodes, nil
case "autoscaling.scale_up_threshold":
return cfg.Autoscaling.ScaleUpThreshold, nil
case "autoscaling.scale_down_threshold":
return cfg.Autoscaling.ScaleDownThreshold, nil
case "autoscaling.evaluation_interval_sec":
return cfg.Autoscaling.EvaluationIntervalSec, nil
case "autoscaling.scale_up_cooldown_sec":
return cfg.Autoscaling.ScaleUpCooldownSec, nil
case "autoscaling.scale_down_cooldown_sec":
return cfg.Autoscaling.ScaleDownCooldownSec, nil
case "monitoring.enable_metrics":
return cfg.Monitoring.EnableMetrics, nil
case "monitoring.metrics_port":
return cfg.Monitoring.MetricsPort, nil
case "monitoring.enable_tracing":
return cfg.Monitoring.EnableTracing, nil
case "monitoring.trace_sample_rate":
return cfg.Monitoring.TraceSampleRate, nil
case "transactions.enabled":
return cfg.Transactions.Enabled, nil
case "transactions.default_timeout_sec":
return cfg.Transactions.DefaultTimeoutSec, nil
case "transactions.deadlock_check_interval_sec":
return cfg.Transactions.DeadlockCheckIntervalSec, nil
case "transactions.max_savepoints_per_tx":
return cfg.Transactions.MaxSavepointsPerTx, nil
case "compression.enabled":
return cfg.Compression.Enabled, nil
case "compression.algorithm":
return cfg.Compression.Algorithm, nil
case "compression.level":
return cfg.Compression.Level, nil
case "compression.min_size":
return cfg.Compression.MinSize, nil
case "performance.enable_pipeline":
return cfg.Performance.EnablePipeline, nil
case "performance.batch_size":
return cfg.Performance.BatchSize, nil
case "performance.max_connections":
return cfg.Performance.MaxConnections, nil
case "performance.read_from_follower":
return cfg.Performance.ReadFromFollower, nil
}
return nil, fmt.Errorf("unknown config path: %s", path)
}
// coerceStringToType приводит строку из REPL к типу существующего значения.
func coerceStringToType(valueStr string, currentValue interface{}) (interface{}, error) {
switch currentValue.(type) {
case bool:
b, err := strconv.ParseBool(valueStr)
if err != nil {
return nil, fmt.Errorf("expected bool, got '%s'", valueStr)
}
return b, nil
case int:
i, err := strconv.ParseInt(valueStr, 10, 64)
if err != nil {
return nil, fmt.Errorf("expected integer, got '%s'", valueStr)
}
return int(i), nil
case int8, int16, int32:
i, err := strconv.ParseInt(valueStr, 10, 64)
if err != nil {
return nil, fmt.Errorf("expected integer, got '%s'", valueStr)
}
return i, nil
case int64:
i, err := strconv.ParseInt(valueStr, 10, 64)
if err != nil {
return nil, fmt.Errorf("expected integer, got '%s'", valueStr)
}
return i, nil
case uint, uint8, uint16, uint32:
i, err := strconv.ParseUint(valueStr, 10, 64)
if err != nil {
return nil, fmt.Errorf("expected unsigned integer, got '%s'", valueStr)
}
return i, nil
case uint64:
i, err := strconv.ParseUint(valueStr, 10, 64)
if err != nil {
return nil, fmt.Errorf("expected unsigned integer, got '%s'", valueStr)
}
return i, nil
case float32, float64:
f, err := strconv.ParseFloat(valueStr, 64)
if err != nil {
return nil, fmt.Errorf("expected number, got '%s'", valueStr)
}
return f, nil
case string:
return valueStr, nil
}
if b, err := strconv.ParseBool(valueStr); err == nil {
return b, nil
}
if i, err := strconv.ParseInt(valueStr, 10, 64); err == nil {
return i, nil
}
if f, err := strconv.ParseFloat(valueStr, 64); err == nil {
return f, nil
}
return valueStr, nil
}
// setConfigFieldLocal применяет изменение к config локально.
// Это дубликат setConfigField из dynamic_config.go.
func setConfigFieldLocal(cfg *config.Config, path string, value interface{}) error {
switch path {
case "compression.enabled":
if v, ok := value.(bool); ok {
cfg.Compression.Enabled = v
return nil
}
case "compression.algorithm":
if v, ok := value.(string); ok {
cfg.Compression.Algorithm = v
return nil
}
case "compression.level":
if v, ok := value.(int64); ok {
cfg.Compression.Level = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Compression.Level = v
return nil
}
case "compression.min_size":
if v, ok := value.(int64); ok {
cfg.Compression.MinSize = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Compression.MinSize = v
return nil
}
case "cluster.name":
if v, ok := value.(string); ok {
cfg.Cluster.Name = v
return nil
}
case "cluster.node_ip":
if v, ok := value.(string); ok {
cfg.Cluster.NodeIP = v
return nil
}
case "cluster.node_port":
if v, ok := value.(int64); ok {
cfg.Cluster.NodePort = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Cluster.NodePort = v
return nil
}
case "cluster.raft_port":
if v, ok := value.(int64); ok {
cfg.Cluster.RaftPort = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Cluster.RaftPort = v
return nil
}
case "cluster.region":
if v, ok := value.(string); ok {
cfg.Cluster.Region = v
return nil
}
case "performance.enable_pipeline":
if v, ok := value.(bool); ok {
cfg.Performance.EnablePipeline = v
return nil
}
case "performance.batch_size":
if v, ok := value.(int64); ok {
cfg.Performance.BatchSize = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Performance.BatchSize = v
return nil
}
case "performance.max_connections":
if v, ok := value.(int64); ok {
cfg.Performance.MaxConnections = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Performance.MaxConnections = v
return nil
}
case "performance.read_from_follower":
if v, ok := value.(bool); ok {
cfg.Performance.ReadFromFollower = v
return nil
}
case "monitoring.enable_metrics":
if v, ok := value.(bool); ok {
cfg.Monitoring.EnableMetrics = v
return nil
}
case "monitoring.metrics_port":
if v, ok := value.(int64); ok {
cfg.Monitoring.MetricsPort = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Monitoring.MetricsPort = v
return nil
}
case "monitoring.enable_tracing":
if v, ok := value.(bool); ok {
cfg.Monitoring.EnableTracing = v
return nil
}
case "monitoring.trace_sample_rate":
if v, ok := value.(float64); ok {
cfg.Monitoring.TraceSampleRate = v
return nil
}
case "log.log_level":
if v, ok := value.(string); ok {
cfg.Log.LogLevel = v
return nil
}
case "log.log_file":
if v, ok := value.(string); ok {
cfg.Log.LogFile = v
return nil
}
case "wal.enabled":
if v, ok := value.(bool); ok {
cfg.WAL.Enabled = v
return nil
}
case "wal.segment_size_mb":
if v, ok := value.(int64); ok {
cfg.WAL.SegmentSizeMB = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.WAL.SegmentSizeMB = v
return nil
}
case "wal.sync_interval_sec":
if v, ok := value.(int64); ok {
cfg.WAL.SyncIntervalSec = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.WAL.SyncIntervalSec = v
return nil
}
case "wal.batch_size":
if v, ok := value.(int64); ok {
cfg.WAL.BatchSize = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.WAL.BatchSize = v
return nil
}
case "saga.enabled":
if v, ok := value.(bool); ok {
cfg.Saga.Enabled = v
return nil
}
case "saga.coordinator_count":
if v, ok := value.(int64); ok {
cfg.Saga.CoordinatorCount = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Saga.CoordinatorCount = v
return nil
}
case "saga.max_retries":
if v, ok := value.(int64); ok {
cfg.Saga.MaxRetries = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Saga.MaxRetries = v
return nil
}
case "saga.saga_timeout_sec":
if v, ok := value.(int64); ok {
cfg.Saga.SagaTimeoutSec = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Saga.SagaTimeoutSec = v
return nil
}
case "saga.retry_backoff_ms":
if v, ok := value.(int64); ok {
cfg.Saga.RetryBackoffMs = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Saga.RetryBackoffMs = v
return nil
}
case "saga.stuck_check_interval_sec":
if v, ok := value.(int64); ok {
cfg.Saga.StuckCheckIntervalSec = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Saga.StuckCheckIntervalSec = v
return nil
}
case "api.port":
if v, ok := value.(int64); ok {
cfg.API.Port = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.API.Port = v
return nil
}
case "backpressure.enabled":
if v, ok := value.(bool); ok {
cfg.Backpressure.Enabled = v
return nil
}
case "backpressure.cpu_threshold":
if v, ok := value.(float64); ok {
cfg.Backpressure.CPUThreshold = v
return nil
}
case "backpressure.memory_threshold":
if v, ok := value.(float64); ok {
cfg.Backpressure.MemoryThreshold = v
return nil
}
case "backpressure.queue_size_threshold":
if v, ok := value.(int64); ok {
cfg.Backpressure.QueueSizeThreshold = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Backpressure.QueueSizeThreshold = v
return nil
}
case "backpressure.connection_threshold":
if v, ok := value.(int64); ok {
cfg.Backpressure.ConnectionThreshold = int(v)
return nil
}
if v, ok := value.(int); ok {
cfg.Backpressure.ConnectionThreshold = v
return nil
}
}
return fmt.Errorf("local update for path '%s' is not implemented (use Raft for full support)", path)
}
// Проверка, что filepath используется (для будущего расширения
// fallback-записи config.toml в произвольную директорию).
var _ = filepath.Join