/* * 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/plugin/plugin.go // Назначение: Система плагинов на основе Lua для расширения функциональности СУБД. // // ОСНОВНЫЕ ФУНКЦИИ: // 1. Загрузка Lua-скриптов как плагинов с изолированным окружением (sandbox) // 2. Выполнение Lua-скриптов с ограничениями по CPU и памяти // 3. Взаимодействие с данными СУБД через API (базы данных, коллекции, документы) // 4. Логирование действий плагинов в общий лог-файл // 5. Поддержка зависимостей между плагинами (Plugin Dependencies) // 6. Горячая перезагрузка плагинов без остановки системы (Hot Reload) // 7. Пул Lua-состояний для эффективного использования ресурсов // 8. Поддержка кастомных движков хранения через Lua-плагины // 9. Плагин Marketplace для установки плагинов из репозитория // // АРХИТЕКТУРНЫЕ РЕШЕНИЯ: // - Каждый плагин выполняется в отдельном Lua-состоянии с ограничениями // - Используется пул состояний для предотвращения утечек ресурсов // - Плагины могут регистрировать собственные движки хранения данных // - Поддерживается горячая перезагрузка с сохранением состояния // // ОГРАНИЧЕНИЯ SANDBOX: // В gopher-lua (github.com/yuin/gopher-lua) отсутствует публичный API // lua_sethook/LUA_MASKCOUNT, поэтому лимит инструкций (MaxInstructions) // реализуется косвенно — через context.WithTimeout, который проверяется // виртуальной машиной gopher-lua на каждой инструкции при ненулевом ctx. // Это даёт реальное ограничение по времени выполнения, но не по точному // числу инструкций. Поле MaxInstructions сохранено для совместимости с // конфигурацией и возможного будущего использования. // // ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции // функции storage.BeginTransaction, storage.CommitCurrentTransaction, // storage.AbortCurrentTransaction, storage.HasActiveTransaction, // storage.GetCurrentTransactionID, storage.GetActiveTransactions и тип // storage.Transaction были удалены. Lua-API для транзакций переписан // поверх batch-операций. Сохранены старые имена Lua-функций // (begin_transaction, commit_transaction и т.д.) как обёртки над // batch-API, чтобы существующие Lua-плагины не сломались. // // Семантика batch отличается от транзакции: batch — это короткоживущий // объект для группы операций, применяемых атомарно. Он не поддерживает // savepoints, не имеет таймаутов, не отслеживается глобально как // долгоживущее состояние. Поэтому Lua-функция `begin_transaction` // создаёт batch и сохраняет его в реестре активных batch'ей плагина, // а `commit_transaction` применяет его. // // ИЗМЕНЕНО (2026-10-03): формат временных меток изменён с "2006-01-02 15:04:05.000" // на "02-01-2006 15:04:05.000" для совместимости с OpenIndiana. package plugin import ( "context" "encoding/json" "fmt" "math" "os" "path/filepath" "runtime" "strings" "sync" "sync/atomic" "time" "futriis/internal/storage" lua "github.com/yuin/gopher-lua" ) // ========== Константы ========== const ( // Лимиты песочницы по умолчанию DefaultCPULimit = 100 * time.Millisecond // Максимальное время CPU на выполнение DefaultMemoryLimit = 50 * 1024 * 1024 // Максимальная память (50 MB) DefaultExecutionTimeout = 5 * time.Second // Таймаут выполнения скрипта // Интервалы проверок HotReloadCheckInterval = 30 * time.Second // Проверка изменений плагинов MaxPluginLoadTime = 10 * time.Second // Максимальное время загрузки плагина // Размеры буферов DefaultMaxEventLogSize = 1000 // Максимальное количество записей в логе DefaultBatchSize = 100 // Размер пакета для операций DefaultWriteBufferSize = 64 * 1024 // Буфер записи (64 KB) // Константы для пула Lua-состояний MaxLuaStates = 100 // Максимальное количество состояний LuaStateTTL = 10 * time.Minute // Время жизни состояния в пуле LuaStatePoolSize = 50 // Размер пула состояний ) // PluginConfig представляет конфигурацию плагинов для main.go type PluginConfig struct { Enabled bool `toml:"enabled"` // Включена ли система плагинов ScriptDir string `toml:"script_dir"` // Директория с Lua-скриптами AllowList []string `toml:"allow_list"` // Список разрешённых плагинов MaxCPUMS int64 `toml:"max_cpu_ms"` // Максимальное время CPU (мс) MaxMemoryMB int64 `toml:"max_memory_mb"` // Максимальная память (MB) } // ============================================================================= // СОВМЕСТИМЫЕ АТОМАРНЫЕ ОБЁРТКИ (Go 1.13+) // ============================================================================= // Используются вместо atomic.Bool/Int32/Int64/Uint64/Float64 (Go 1.19+), // чтобы обеспечить сборку на старых версиях Go (актуально для OpenIndiana). type atomicBool struct{ v int32 } func (b *atomicBool) Load() bool { return atomic.LoadInt32(&b.v) != 0 } func (b *atomicBool) Store(v bool) { atomic.StoreInt32(&b.v, boolToInt32(v)) } func (b *atomicBool) CAS(old, new bool) bool { return atomic.CompareAndSwapInt32(&b.v, boolToInt32(old), boolToInt32(new)) } func boolToInt32(b bool) int32 { if b { return 1 } return 0 } type atomicInt32 struct{ v int32 } func (a *atomicInt32) Load() int32 { return atomic.LoadInt32(&a.v) } func (a *atomicInt32) Store(v int32) { atomic.StoreInt32(&a.v, v) } func (a *atomicInt32) Add(d int32) int32 { return atomic.AddInt32(&a.v, d) } func (a *atomicInt32) CAS(old, new int32) bool { return atomic.CompareAndSwapInt32(&a.v, old, new) } type atomicInt64 struct{ v int64 } func (a *atomicInt64) Load() int64 { return atomic.LoadInt64(&a.v) } func (a *atomicInt64) Store(v int64) { atomic.StoreInt64(&a.v, v) } func (a *atomicInt64) Add(d int64) int64 { return atomic.AddInt64(&a.v, d) } type atomicUint64 struct{ v uint64 } func (a *atomicUint64) Load() uint64 { return atomic.LoadUint64(&a.v) } func (a *atomicUint64) Store(v uint64) { atomic.StoreUint64(&a.v, v) } func (a *atomicUint64) Add(d uint64) uint64 { return atomic.AddUint64(&a.v, d) } type atomicFloat64 struct{ v uint64 } func (f *atomicFloat64) Load() float64 { return math.Float64frombits(atomic.LoadUint64(&f.v)) } func (f *atomicFloat64) Store(v float64) { atomic.StoreUint64(&f.v, math.Float64bits(v)) } // ========== Интерфейсы для кастомных движков ========== // CustomEngine определяет интерфейс для кастомного движка хранения. // Движки могут быть реализованы как на Go, так и на Lua через плагины. type CustomEngine interface { // Базовые операции Name() string Version() string // Операции с документами Insert(doc *storage.Document) error Find(id string) (*storage.Document, error) Update(id string, updates map[string]interface{}) error Delete(id string) error // Пакетные операции BatchInsert(docs []*storage.Document) error BatchUpdate(updates map[string]map[string]interface{}) error BatchDelete(ids []string) error // Поисковые операции FindByFilter(filter func(*storage.Document) bool) ([]*storage.Document, error) FindByIndex(indexName string, value interface{}) ([]*storage.Document, error) // Индексы CreateIndex(name string, fields []string, unique bool) error DropIndex(name string) error GetIndexes() []string // Статистика и метаданные Count() int64 Size() int64 GetStats() map[string]interface{} // Жизненный цикл Initialize(config map[string]interface{}) error Close() error // События (для интеграции с триггерами) OnDocumentInserted(doc *storage.Document) OnDocumentUpdated(oldDoc, newDoc *storage.Document) OnDocumentDeleted(doc *storage.Document) } // EngineRegistry управляет зарегистрированными кастомными движками. // Поддерживает регистрацию фабрик движков и создание экземпляров. type EngineRegistry struct { engines map[string]EngineFactory // Зарегистрированные фабрики activeEngines map[string]CustomEngine // Активные экземпляры движков mu sync.RWMutex // Защита доступа logger storage.LoggerInterface // Логгер pm *PluginManager // Менеджер плагинов } // EngineFactory создаёт экземпляр кастомного движка. type EngineFactory func(config map[string]interface{}, logger storage.LoggerInterface) (CustomEngine, error) // NewEngineRegistry создаёт новый реестр движков. func NewEngineRegistry(logger storage.LoggerInterface, pm *PluginManager) *EngineRegistry { return &EngineRegistry{ engines: make(map[string]EngineFactory), activeEngines: make(map[string]CustomEngine), logger: logger, pm: pm, } } // RegisterEngine регистрирует фабрику движка в реестре. func (er *EngineRegistry) RegisterEngine(name string, factory EngineFactory) { er.mu.Lock() defer er.mu.Unlock() er.engines[name] = factory if er.logger != nil { er.logger.Info(fmt.Sprintf("Engine registered: %s", name)) } } // CreateEngine создаёт экземпляр движка по имени. func (er *EngineRegistry) CreateEngine(name string, config map[string]interface{}) (CustomEngine, error) { er.mu.RLock() factory, ok := er.engines[name] er.mu.RUnlock() if !ok { return nil, fmt.Errorf("engine not found: %s", name) } engine, err := factory(config, er.logger) if err != nil { return nil, err } if err := engine.Initialize(config); err != nil { return nil, err } er.mu.Lock() er.activeEngines[engine.Name()] = engine er.mu.Unlock() return engine, nil } // GetEngine возвращает активный движок по имени. func (er *EngineRegistry) GetEngine(name string) CustomEngine { er.mu.RLock() defer er.mu.RUnlock() return er.activeEngines[name] } // ListEngines возвращает список зарегистрированных движков. func (er *EngineRegistry) ListEngines() []string { er.mu.RLock() defer er.mu.RUnlock() engines := make([]string, 0, len(er.engines)) for name := range er.engines { engines = append(engines, name) } return engines } // ListActiveEngines возвращает список активных движков. func (er *EngineRegistry) ListActiveEngines() []string { er.mu.RLock() defer er.mu.RUnlock() engines := make([]string, 0, len(er.activeEngines)) for name := range er.activeEngines { engines = append(engines, name) } return engines } // CloseAll закрывает все активные движки. func (er *EngineRegistry) CloseAll() error { er.mu.Lock() defer er.mu.Unlock() var lastErr error for name, engine := range er.activeEngines { if err := engine.Close(); err != nil { lastErr = err if er.logger != nil { er.logger.Error(fmt.Sprintf("Failed to close engine %s: %v", name, err)) } } } er.activeEngines = make(map[string]CustomEngine) return lastErr } // ========== Пул Lua-состояний ========== // LuaStatePool управляет пулом Lua-состояний для ограничения ресурсов. // // Close() не закрывает канал states, чтобы избежать паники // "send on closed channel" при одновременном Release из других горутин. // Вместо этого используется флаг closed (atomicBool). type LuaStatePool struct { states chan *InstrumentedLState active atomicInt32 maxSize int mu sync.RWMutex createdAt map[*InstrumentedLState]time.Time closed atomicBool } // Глобальный пул состояний (синглтон) var globalStatePool *LuaStatePool var statePoolOnce sync.Once // GetLuaStatePool возвращает глобальный пул Lua-состояний. func GetLuaStatePool() *LuaStatePool { statePoolOnce.Do(func() { globalStatePool = &LuaStatePool{ states: make(chan *InstrumentedLState, LuaStatePoolSize), maxSize: MaxLuaStates, createdAt: make(map[*InstrumentedLState]time.Time), } }) return globalStatePool } // Acquire получает состояние из пула или создаёт новое. func (p *LuaStatePool) Acquire(limits SandboxLimits) (*InstrumentedLState, error) { if p.closed.Load() { return nil, fmt.Errorf("state pool is closed") } // Пробуем получить готовое состояние select { case ls := <-p.states: if ls == nil { return p.Acquire(limits) } p.mu.RLock() createdAt, ok := p.createdAt[ls] p.mu.RUnlock() if !ok || time.Since(createdAt) > LuaStateTTL { // Состояние устарело - закрываем и создаём новое ls.Close() p.mu.Lock() delete(p.createdAt, ls) p.mu.Unlock() p.active.Add(-1) return p.Acquire(limits) } // Применяем новые лимиты (включая пересоздание контекста) ls.SetLimits(limits) return ls, nil default: } // Пул пуст — проверяем лимит if p.active.Load() >= int32(p.maxSize) { return nil, fmt.Errorf("max Lua states limit reached: %d", p.maxSize) } ls := NewInstrumentedLState(limits) p.active.Add(1) p.mu.Lock() p.createdAt[ls] = time.Now() p.mu.Unlock() return ls, nil } // Release возвращает состояние в пул для переиспользования. func (p *LuaStatePool) Release(ls *InstrumentedLState) { if ls == nil { return } if p.closed.Load() { ls.Close() p.active.Add(-1) p.mu.Lock() delete(p.createdAt, ls) p.mu.Unlock() return } // Сбрасываем счётчики для нового использования ls.Reset() select { case p.states <- ls: // Успешно вернули в пул default: // Пул заполнен - закрываем состояние ls.Close() p.active.Add(-1) p.mu.Lock() delete(p.createdAt, ls) p.mu.Unlock() } } // GetActiveCount возвращает количество активных состояний. func (p *LuaStatePool) GetActiveCount() int32 { return p.active.Load() } // GetPoolSize возвращает размер пула (количество свободных состояний). func (p *LuaStatePool) GetPoolSize() int { return len(p.states) } // Close закрывает пул и все состояния. // Безопасен для многократного вызова. НЕ закрывает канал states. func (p *LuaStatePool) Close() { if !p.closed.CAS(false, true) { return } // Дренируем канал без его закрытия for { select { case ls := <-p.states: if ls != nil { ls.Close() } default: goto drained } } drained: p.mu.Lock() for ls := range p.createdAt { if ls != nil { ls.Close() } } p.createdAt = make(map[*InstrumentedLState]time.Time) p.mu.Unlock() p.active.Store(0) } // ========== Sandbox с ограничениями ========== // SandboxLimits определяет ограничения для Lua-скрипта. // // ВАЖНО: MaxInstructions в gopher-lua не имеет прямого API для // подсчёта инструкций (нет аналога lua_sethook/LUA_MASKCOUNT). // Поле оставлено для совместимости с конфигом и может быть // реализовано через context.WithTimeout (косвенно, по времени). type SandboxLimits struct { MaxCPUTime time.Duration // Максимальное время CPU MaxMemory int64 // Максимальный размер памяти (в байтах) MaxExecutionTime time.Duration // Максимальное время выполнения скрипта MaxStackDepth int // Максимальная глубина стека MaxInstructions int64 // Максимальное количество инструкций (зарезервировано) } // InstrumentedLState расширяет lua.LState с мониторингом ресурсов. // // ОГРАНИЧЕНИЯ SANDBOX: // gopher-lua не предоставляет публичного API lua_sethook/LUA_MASKCOUNT // для подсчёта инструкций. Вместо этого используется L.SetContext(ctx) // с таймаутом — виртуальная машина gopher-lua проверяет ctx.Done() на // каждой инструкции, что даёт реальное прерывание по времени. // // НЕ используем runtime.SetFinalizer: финализатор держит ссылку на // объект, мешает сборке мусора и приводит к двойному закрытию. type InstrumentedLState struct { *lua.LState limits SandboxLimits instructions atomicInt64 // Зарезервировано; в gopher-lua не инкрементируется автоматически startTime time.Time ctx context.Context cancel context.CancelFunc mu sync.RWMutex closed bool } // NewInstrumentedLState создаёт Lua-состояние с ограничениями. // Открывает только безопасные библиотеки: base, string, table, math. func NewInstrumentedLState(limits SandboxLimits) *InstrumentedLState { ctx, cancel := context.WithCancel(context.Background()) ls := &InstrumentedLState{ LState: lua.NewState(lua.Options{SkipOpenLibs: true}), limits: limits, startTime: time.Now(), ctx: ctx, cancel: cancel, } // Открываем только безопасные библиотеки lua.OpenBase(ls.LState) lua.OpenString(ls.LState) lua.OpenTable(ls.LState) lua.OpenMath(ls.LState) // Устанавливаем контекст с таймаутом. // gopher-lua проверяет ctx.Done() на каждой инструкции — это даёт // реальное прерывание выполнения при превышении лимита времени, // что компенсирует отсутствие LUA_MASKCOUNT. ls.applyContext() return ls } // applyContext создаёт новый контекст с таймаутом на основе лимитов // и устанавливает его в LState. Вызывается при создании состояния // и при обновлении лимитов (SetLimits, Reset). // // Логика выбора таймаута: // - Если задан MaxExecutionTime > 0 — используем его. // - Иначе, если задан MaxCPUTime > 0 — используем его. // - Иначе — без таймаута (только context.Background с возможностью // внешней отмены через cancel()). func (ls *InstrumentedLState) applyContext() { ls.mu.Lock() defer ls.mu.Unlock() // Отменяем предыдущий контекст, если был if ls.cancel != nil { ls.cancel() } var ctx context.Context var cancel context.CancelFunc timeout := ls.limits.MaxExecutionTime if timeout <= 0 { timeout = ls.limits.MaxCPUTime } if timeout > 0 { ctx, cancel = context.WithTimeout(context.Background(), timeout) } else { ctx, cancel = context.WithCancel(context.Background()) } ls.ctx = ctx ls.cancel = cancel if ls.LState != nil { ls.LState.SetContext(ctx) } } // SetLimits обновляет лимиты для повторного использования состояния. // Пересоздаёт контекст с новым таймаутом. func (ls *InstrumentedLState) SetLimits(limits SandboxLimits) { ls.mu.Lock() ls.limits = limits ls.mu.Unlock() ls.applyContext() } // Reset сбрасывает счётчики для нового использования. // Пересоздаёт контекст с таймаутом, чтобы состояние можно было // безопасно переиспользовать после предыдущего выполнения. func (ls *InstrumentedLState) Reset() { ls.instructions.Store(0) ls.mu.Lock() ls.startTime = time.Now() ls.mu.Unlock() ls.applyContext() } // SetCPULimit устанавливает лимит CPU времени. func (ls *InstrumentedLState) SetCPULimit(limit time.Duration) { ls.mu.Lock() ls.limits.MaxCPUTime = limit ls.mu.Unlock() ls.applyContext() } // SetMemoryLimit устанавливает лимит памяти. func (ls *InstrumentedLState) SetMemoryLimit(limit int64) { ls.mu.Lock() defer ls.mu.Unlock() ls.limits.MaxMemory = limit } // CheckMemory проверяет текущее использование памяти. func (ls *InstrumentedLState) CheckMemory() bool { ls.mu.RLock() maxMem := ls.limits.MaxMemory ls.mu.RUnlock() if maxMem <= 0 { return true } var m runtime.MemStats runtime.ReadMemStats(&m) return int64(m.Alloc) <= maxMem } // GetSandboxStats возвращает статистику песочницы. func (ls *InstrumentedLState) GetSandboxStats() map[string]interface{} { var m runtime.MemStats runtime.ReadMemStats(&m) return map[string]interface{}{ "instructions": ls.instructions.Load(), "execution_ms": time.Since(ls.startTime).Milliseconds(), "memory_alloc": m.Alloc, "memory_sys": m.Sys, "num_gc": m.NumGC, "pause_total_ms": m.PauseTotalNs / 1000000, } } // IsClosed возвращает статус закрытия состояния. func (ls *InstrumentedLState) IsClosed() bool { ls.mu.RLock() defer ls.mu.RUnlock() return ls.closed } // Close закрывает Lua состояние и освобождает ресурсы. // Безопасно вызывать несколько раз. func (ls *InstrumentedLState) Close() { ls.mu.Lock() if ls.closed { ls.mu.Unlock() return } ls.closed = true ls.mu.Unlock() if ls.cancel != nil { ls.cancel() } if ls.LState != nil { ls.LState.Close() } } // ========== Plugin Dependencies ========== // PluginDependency определяет зависимость плагина. type PluginDependency struct { Name string `json:"name"` Version string `json:"version"` MinVersion string `json:"min_version,omitempty"` MaxVersion string `json:"max_version,omitempty"` Optional bool `json:"optional,omitempty"` } // PluginManifest описывает метаданные плагина. type PluginManifest struct { Name string `json:"name"` Version string `json:"version"` Author string `json:"author"` Description string `json:"description"` Dependencies []PluginDependency `json:"dependencies,omitempty"` APIVersion string `json:"api_version"` MinGoVersion string `json:"min_go_version,omitempty"` EntryPoint string `json:"entry_point,omitempty"` CreatedAt int64 `json:"created_at"` UpdatedAt int64 `json:"updated_at"` EngineType string `json:"engine_type,omitempty"` } // PluginWithDeps расширяет Plugin с поддержкой зависимостей. type PluginWithDeps struct { *Plugin Manifest *PluginManifest Dependencies map[string]*PluginWithDeps Dependents map[string]*PluginWithDeps LoadOrder int depMu sync.RWMutex } // DependencyResolver разрешает зависимости между плагинами. type DependencyResolver struct { plugins map[string]*PluginWithDeps mu sync.RWMutex } // NewDependencyResolver создаёт новый резолвер зависимостей. func NewDependencyResolver() *DependencyResolver { return &DependencyResolver{ plugins: make(map[string]*PluginWithDeps), } } // RegisterPlugin регистрирует плагин с его манифестом в резолвере. func (dr *DependencyResolver) RegisterPlugin(plugin *PluginWithDeps) error { dr.mu.Lock() defer dr.mu.Unlock() dr.plugins[plugin.Name] = plugin return nil } // ResolveOrder возвращает порядок загрузки плагинов (топологическая сортировка). func (dr *DependencyResolver) ResolveOrder() ([]*PluginWithDeps, error) { dr.mu.RLock() defer dr.mu.RUnlock() graph := make(map[string][]string) inDegree := make(map[string]int) for name, plugin := range dr.plugins { inDegree[name] = 0 for _, dep := range plugin.Manifest.Dependencies { if !dep.Optional { graph[dep.Name] = append(graph[dep.Name], name) inDegree[name]++ } } } queue := make([]string, 0) for name, degree := range inDegree { if degree == 0 { queue = append(queue, name) } } result := make([]*PluginWithDeps, 0) for len(queue) > 0 { current := queue[0] queue = queue[1:] result = append(result, dr.plugins[current]) for _, dependent := range graph[current] { inDegree[dependent]-- if inDegree[dependent] == 0 { queue = append(queue, dependent) } } } if len(result) != len(dr.plugins) { return nil, fmt.Errorf("circular dependency detected among plugins") } for i, plugin := range result { plugin.LoadOrder = i } return result, nil } // CheckVersionCompatibility проверяет совместимость версий. func (dr *DependencyResolver) CheckVersionCompatibility(dep PluginDependency, actualVersion string) bool { if dep.MinVersion != "" && actualVersion < dep.MinVersion { return false } if dep.MaxVersion != "" && actualVersion > dep.MaxVersion { return false } if dep.Version != "" && actualVersion != dep.Version { return false } return true } // ========== File Watcher (без fsnotify) ========== // FileWatcher самостоятельно отслеживает изменения файлов. type FileWatcher struct { files map[string]time.Time mu sync.RWMutex interval time.Duration onChange func(path string) error stopChan chan struct{} stopOnce sync.Once wg sync.WaitGroup logger storage.LoggerInterface } // NewFileWatcher создаёт новый файловый вотчер без fsnotify. func NewFileWatcher(interval time.Duration, onChange func(path string) error, logger storage.LoggerInterface) *FileWatcher { if interval <= 0 { interval = HotReloadCheckInterval } return &FileWatcher{ files: make(map[string]time.Time), interval: interval, onChange: onChange, stopChan: make(chan struct{}), logger: logger, } } // Add добавляет файл для отслеживания. func (fw *FileWatcher) Add(path string) error { info, err := os.Stat(path) if err != nil { return err } fw.mu.Lock() defer fw.mu.Unlock() fw.files[path] = info.ModTime() return nil } // Remove прекращает отслеживание файла. func (fw *FileWatcher) Remove(path string) { fw.mu.Lock() defer fw.mu.Unlock() delete(fw.files, path) } // Start запускает цикл проверки изменений. func (fw *FileWatcher) Start() { fw.wg.Add(1) go fw.watchLoop() } // Stop останавливает цикл проверки. Безопасен для многократного вызова. func (fw *FileWatcher) Stop() { fw.stopOnce.Do(func() { close(fw.stopChan) }) fw.wg.Wait() } // watchLoop периодически проверяет изменения файлов. func (fw *FileWatcher) watchLoop() { defer fw.wg.Done() ticker := time.NewTicker(fw.interval) defer ticker.Stop() for { select { case <-ticker.C: fw.checkChanges() case <-fw.stopChan: return } } } // checkChanges проверяет изменения всех отслеживаемых файлов. func (fw *FileWatcher) checkChanges() { fw.mu.RLock() files := make([]string, 0, len(fw.files)) for path := range fw.files { files = append(files, path) } fw.mu.RUnlock() for _, path := range files { info, err := os.Stat(path) if err != nil { if fw.logger != nil { fw.logger.Warn(fmt.Sprintf("Failed to stat file %s: %v", path, err)) } continue } fw.mu.RLock() lastMod := fw.files[path] fw.mu.RUnlock() if info.ModTime().After(lastMod) { if fw.logger != nil { fw.logger.Info(fmt.Sprintf("File changed: %s (was %v, now %v)", path, lastMod, info.ModTime())) } if err := fw.onChange(path); err != nil { if fw.logger != nil { fw.logger.Error(fmt.Sprintf("Failed to handle change for %s: %v", path, err)) } } fw.mu.Lock() fw.files[path] = info.ModTime() fw.mu.Unlock() } } } // ========== HotReloadManager управляет горячей перезагрузкой плагинов ========== // HotReloadManager управляет горячей перезагрузкой плагинов. type HotReloadManager struct { pluginManager *PluginManager watcher *FileWatcher reloadQueue chan string reloading map[string]bool mu sync.RWMutex logger storage.LoggerInterface stopOnce sync.Once stopped chan struct{} } // NewHotReloadManager создаёт менеджер горячей перезагрузки. func NewHotReloadManager(pm *PluginManager, logger storage.LoggerInterface) *HotReloadManager { hrm := &HotReloadManager{ pluginManager: pm, reloadQueue: make(chan string, 100), reloading: make(map[string]bool), logger: logger, stopped: make(chan struct{}), } hrm.watcher = NewFileWatcher(HotReloadCheckInterval, hrm.handleFileChange, logger) go hrm.processReloadQueue() return hrm } // Start запускает мониторинг изменений. func (hrm *HotReloadManager) Start() { hrm.watcher.Start() } // Stop останавливает мониторинг. Безопасен для многократного вызова. func (hrm *HotReloadManager) Stop() { hrm.stopOnce.Do(func() { hrm.watcher.Stop() close(hrm.stopped) }) } // WatchPlugin начинает отслеживать плагин. func (hrm *HotReloadManager) WatchPlugin(pluginPath string) error { return hrm.watcher.Add(pluginPath) } // handleFileChange обрабатывает изменение файла. func (hrm *HotReloadManager) handleFileChange(path string) error { baseName := filepath.Base(path) pluginName := strings.TrimSuffix(baseName, ".lua") hrm.mu.Lock() if hrm.reloading[pluginName] { hrm.mu.Unlock() return nil } hrm.reloading[pluginName] = true hrm.mu.Unlock() select { case hrm.reloadQueue <- pluginName: default: if hrm.logger != nil { hrm.logger.Warn(fmt.Sprintf("Reload queue full, skipping reload for %s", pluginName)) } } return nil } // processReloadQueue обрабатывает очередь перезагрузки. // Завершается при закрытии канала stopped. func (hrm *HotReloadManager) processReloadQueue() { for { select { case pluginName := <-hrm.reloadQueue: hrm.reloadPlugin(pluginName) hrm.mu.Lock() delete(hrm.reloading, pluginName) hrm.mu.Unlock() case <-hrm.stopped: return } } } // reloadPlugin выполняет горячую перезагрузку плагина. func (hrm *HotReloadManager) reloadPlugin(pluginName string) { if hrm.logger != nil { hrm.logger.Info(fmt.Sprintf("Hot reloading plugin: %s", pluginName)) } oldPlugin, err := hrm.pluginManager.GetPlugin(pluginName) if err != nil { if hrm.logger != nil { hrm.logger.Error(fmt.Sprintf("Plugin not found for reload: %s", pluginName)) } return } // Сохраняем состояние старого плагина var oldState map[string]interface{} if oldPlugin.LState != nil && !oldPlugin.IsClosed() { fn := oldPlugin.LState.GetGlobal("on_before_reload") if fn != lua.LNil { oldPlugin.LState.CallByParam(lua.P{ Fn: fn, NRet: 1, Protect: true, }) ret := oldPlugin.LState.Get(-1) if table, ok := ret.(*lua.LTable); ok { oldState = make(map[string]interface{}) table.ForEach(func(key, value lua.LValue) { oldState[key.String()] = value.String() }) } oldPlugin.LState.Pop(1) } } if err := hrm.pluginManager.StopPlugin(pluginName); err != nil { if hrm.logger != nil { hrm.logger.Warn(fmt.Sprintf("Failed to stop plugin before reload: %v", err)) } } if err := hrm.pluginManager.UnloadPlugin(pluginName); err != nil { if hrm.logger != nil { hrm.logger.Error(fmt.Sprintf("Failed to unload plugin for reload: %v", err)) } return } pluginPath := filepath.Join(hrm.pluginManager.GetPluginsDir(), pluginName+".lua") if err := hrm.pluginManager.LoadPluginWithSandbox(pluginName, pluginPath, SandboxLimits{ MaxCPUTime: DefaultCPULimit, MaxMemory: DefaultMemoryLimit, MaxExecutionTime: DefaultExecutionTimeout, MaxInstructions: 1000000, }); err != nil { if hrm.logger != nil { hrm.logger.Error(fmt.Sprintf("Failed to load plugin during reload: %v", err)) } return } newPlugin, _ := hrm.pluginManager.GetPlugin(pluginName) if newPlugin != nil && newPlugin.LState != nil && !newPlugin.IsClosed() && oldState != nil { fn := newPlugin.LState.GetGlobal("on_after_reload") if fn != lua.LNil { stateTable := newPlugin.LState.NewTable() for k, v := range oldState { stateTable.RawSetString(k, lua.LString(fmt.Sprintf("%v", v))) } newPlugin.LState.CallByParam(lua.P{ Fn: fn, NRet: 0, Protect: true, }, stateTable) } } if err := hrm.pluginManager.StartPlugin(pluginName); err != nil { if hrm.logger != nil { hrm.logger.Error(fmt.Sprintf("Failed to start plugin after reload: %v", err)) } } if hrm.logger != nil { hrm.logger.Info(fmt.Sprintf("Plugin hot reload completed: %s", pluginName)) } storage.LogAudit("PLUGIN_HOT_RELOAD", "PLUGIN", pluginName, map[string]interface{}{ "reload_time": time.Now().UnixMilli(), }) } // ========== Plugin Marketplace API ========== // MarketplacePlugin представляет плагин в маркетплейсе. type MarketplacePlugin struct { ID string `json:"id"` Name string `json:"name"` Version string `json:"version"` Author string `json:"author"` Description string `json:"description"` Downloads int64 `json:"downloads"` Rating float64 `json:"rating"` Tags []string `json:"tags"` Repository string `json:"repository"` CreatedAt int64 `json:"created_at"` UpdatedAt int64 `json:"updated_at"` } // PluginMarketplace управляет плагинами из маркетплейса. type PluginMarketplace struct { plugins map[string]*MarketplacePlugin installed map[string]bool apiEndpoint string mu sync.RWMutex logger storage.LoggerInterface pm *PluginManager } // NewPluginMarketplace создаёт новый маркетплейс. func NewPluginMarketplace(apiEndpoint string, pm *PluginManager, logger storage.LoggerInterface) *PluginMarketplace { return &PluginMarketplace{ plugins: make(map[string]*MarketplacePlugin), installed: make(map[string]bool), apiEndpoint: apiEndpoint, logger: logger, pm: pm, } } // ListAvailablePlugins возвращает список доступных плагинов. func (pmkt *PluginMarketplace) ListAvailablePlugins() []*MarketplacePlugin { pmkt.mu.RLock() defer pmkt.mu.RUnlock() result := make([]*MarketplacePlugin, 0, len(pmkt.plugins)) for _, p := range pmkt.plugins { result = append(result, p) } return result } // RegisterPlugin регистрирует плагин в маркетплейсе. func (pmkt *PluginMarketplace) RegisterPlugin(plugin *MarketplacePlugin) { pmkt.mu.Lock() defer pmkt.mu.Unlock() pmkt.plugins[plugin.ID] = plugin } // InstallPlugin устанавливает плагин из маркетплейса. func (pmkt *PluginMarketplace) InstallPlugin(pluginID, version string) error { pmkt.mu.RLock() plugin, ok := pmkt.plugins[pluginID] pmkt.mu.RUnlock() if !ok { return fmt.Errorf("plugin not found in marketplace: %s", pluginID) } pluginPath := filepath.Join(pmkt.pm.GetPluginsDir(), plugin.Name+".lua") if err := pmkt.pm.LoadPluginWithSandbox(plugin.Name, pluginPath, SandboxLimits{ MaxCPUTime: DefaultCPULimit, MaxMemory: DefaultMemoryLimit, MaxExecutionTime: DefaultExecutionTimeout, MaxInstructions: 1000000, }); err != nil { return fmt.Errorf("failed to load plugin: %v", err) } if err := pmkt.pm.StartPlugin(plugin.Name); err != nil { return fmt.Errorf("failed to start plugin: %v", err) } pmkt.mu.Lock() pmkt.installed[pluginID] = true pmkt.mu.Unlock() if pmkt.logger != nil { pmkt.logger.Info(fmt.Sprintf("Installed plugin from marketplace: %s v%s", plugin.Name, plugin.Version)) } return nil } // UninstallPlugin удаляет плагин из системы. func (pmkt *PluginMarketplace) UninstallPlugin(pluginID string) error { pmkt.mu.RLock() plugin, ok := pmkt.plugins[pluginID] pmkt.mu.RUnlock() if !ok { return fmt.Errorf("plugin not found: %s", pluginID) } if err := pmkt.pm.StopPlugin(plugin.Name); err != nil { return err } if err := pmkt.pm.UnloadPlugin(plugin.Name); err != nil { return err } pmkt.mu.Lock() delete(pmkt.installed, pluginID) pmkt.mu.Unlock() return nil } // UpdatePlugin обновляет плагин до новой версии. func (pmkt *PluginMarketplace) UpdatePlugin(pluginID, newVersion string) error { if err := pmkt.UninstallPlugin(pluginID); err != nil { return err } return pmkt.InstallPlugin(pluginID, newVersion) } // GetInstalledPlugins возвращает список установленных плагинов. func (pmkt *PluginMarketplace) GetInstalledPlugins() []string { pmkt.mu.RLock() defer pmkt.mu.RUnlock() result := make([]string, 0, len(pmkt.installed)) for id := range pmkt.installed { result = append(result, id) } return result } // ========== Plugin Struct ========== // PluginStatus представляет состояние плагина. type PluginStatus int32 const ( StatusLoaded PluginStatus = iota // Загружен, но не запущен StatusRunning // Запущен и работает StatusStopped // Остановлен StatusError // В состоянии ошибки ) // PluginEventLog представляет событие плагина с временными метками. type PluginEventLog struct { PluginName string `json:"plugin_name"` EventType string `json:"event_type"` Data interface{} `json:"data"` Timestamp int64 `json:"timestamp"` TimestampStr string `json:"timestamp_str"` DurationMs int64 `json:"duration_ms,omitempty"` Error string `json:"error,omitempty"` } // Plugin представляет загруженный Lua-плагин. // // Поле instrumented хранит ссылку на *InstrumentedLState, чтобы // избежать утечки при выгрузке плагина (ранее type assertion // interface{}(L).(*InstrumentedLState) всегда возвращал false, // поскольку L имеет тип *lua.LState). type Plugin struct { Name string FilePath string Status atomicInt32 LState *lua.LState instrumented *InstrumentedLState // Ссылка на обёртку для возврата в пул logger storage.LoggerInterface storage *storage.Storage mu sync.RWMutex loadedAt time.Time loadedAtMs int64 loadedAtStr string version string author string description string eventLog []PluginEventLog eventLogMu sync.RWMutex maxEventLog int closed bool isEngine bool engineName string // Реестр активных batch'ей, созданных через Lua-API. // Ключ — ID batch'а (uint64), значение — *storage.Batch. // Используется функциями begin_transaction/commit_transaction // в Lua-API: `begin_transaction` создаёт batch и регистрирует его // здесь, `commit_transaction` извлекает и коммитит. activeBatches sync.Map } // IsClosed возвращает статус закрытия плагина. func (p *Plugin) IsClosed() bool { p.mu.RLock() defer p.mu.RUnlock() return p.closed } // closePlugin закрывает плагин и освобождает ресурсы. func (p *Plugin) closePlugin() { p.mu.Lock() defer p.mu.Unlock() if p.closed { return } p.closed = true if p.LState != nil { p.LState.Close() } } // IsEnginePlugin возвращает true, если плагин является движком. func (p *Plugin) IsEnginePlugin() bool { return p.isEngine } // GetEngineName возвращает имя движка, если плагин является движком. func (p *Plugin) GetEngineName() string { return p.engineName } // RegisterBatch регистрирует batch в плагине. func (p *Plugin) RegisterBatch(id uint64, batch *storage.Batch) { if p == nil || batch == nil { return } p.activeBatches.Store(id, batch) } // UnregisterBatch снимает регистрацию batch. func (p *Plugin) UnregisterBatch(id uint64) { if p == nil { return } p.activeBatches.Delete(id) } // GetBatch возвращает batch по ID. func (p *Plugin) GetBatch(id uint64) (*storage.Batch, bool) { if p == nil { return nil, false } val, ok := p.activeBatches.Load(id) if !ok { return nil, false } batch, ok := val.(*storage.Batch) return batch, ok } // ========== PluginManager ========== // PluginManager управляет всеми загруженными плагинами. type PluginManager struct { plugins sync.Map logger storage.LoggerInterface storage *storage.Storage pluginsDir string eventBus chan PluginEvent eventBusClosed atomicBool enabled bool depResolver *DependencyResolver hotReload *HotReloadManager marketplace *PluginMarketplace statePool *LuaStatePool engineRegistry *EngineRegistry // nextBatchID — счётчик для ID batch'ей, созданных через Lua-API. nextBatchID atomicUint64 } // PluginEvent представляет событие от плагина. type PluginEvent struct { PluginName string EventType string Data interface{} Timestamp int64 TimestampStr string } // NewPluginManager создаёт новый менеджер плагинов. func NewPluginManager(pluginsDir string, logger storage.LoggerInterface, store *storage.Storage, enabled bool) *PluginManager { pm := &PluginManager{ logger: logger, storage: store, pluginsDir: pluginsDir, eventBus: make(chan PluginEvent, 1000), enabled: enabled, depResolver: NewDependencyResolver(), statePool: GetLuaStatePool(), } pm.engineRegistry = NewEngineRegistry(logger, pm) if !enabled { if logger != nil { logger.Info("Plugin system is disabled") } return pm } if err := os.MkdirAll(pluginsDir, 0755); err != nil { if logger != nil { logger.Error(fmt.Sprintf("Failed to create plugins directory: %v", err)) } } go pm.eventLoop() pm.hotReload = NewHotReloadManager(pm, logger) pm.hotReload.Start() go pm.autoLoadPlugins() if logger != nil { logger.Info(fmt.Sprintf("Plugin system initialized, plugins directory: %s", pluginsDir)) } return pm } // NewPluginManagerFromConfig создаёт новый менеджер плагинов из конфигурации. func NewPluginManagerFromConfig(cfg *PluginConfig, logger storage.LoggerInterface, store *storage.Storage) *PluginManager { if cfg == nil { return NewPluginManager("plugins", logger, store, false) } pm := NewPluginManager(cfg.ScriptDir, logger, store, cfg.Enabled) if cfg.Enabled && len(cfg.AllowList) > 0 { if pm.logger != nil { pm.logger.Info(fmt.Sprintf("Plugin allow list configured: %v", cfg.AllowList)) } } return pm } // GetEngineRegistry возвращает реестр движков. func (pm *PluginManager) GetEngineRegistry() *EngineRegistry { return pm.engineRegistry } // getPluginLState извлекает *lua.LState и *Plugin из значения sync.Map. func getPluginLState(plugin interface{}) (*lua.LState, *Plugin) { switch v := plugin.(type) { case *Plugin: v.mu.RLock() defer v.mu.RUnlock() return v.LState, v case *PluginWithDeps: v.Plugin.mu.RLock() defer v.Plugin.mu.RUnlock() return v.Plugin.LState, v.Plugin default: return nil, nil } } // getPluginInstrumented извлекает *InstrumentedLState из плагина. func getPluginInstrumented(plugin interface{}) *InstrumentedLState { switch v := plugin.(type) { case *Plugin: v.mu.RLock() defer v.mu.RUnlock() return v.instrumented case *PluginWithDeps: v.Plugin.mu.RLock() defer v.Plugin.mu.RUnlock() return v.Plugin.instrumented default: return nil } } // RegisterEngineFromLua регистрирует движок из Lua-плагина. func (pm *PluginManager) RegisterEngineFromLua(pluginName string, engineName string, factoryFn *lua.LFunction) error { if !pm.enabled { return fmt.Errorf("plugin system is disabled") } val, ok := pm.plugins.Load(pluginName) if !ok { return fmt.Errorf("plugin not found: %s", pluginName) } L, p := getPluginLState(val) if L == nil || p == nil || p.IsClosed() { return fmt.Errorf("plugin has no Lua state or is closed") } factory := func(config map[string]interface{}, logger storage.LoggerInterface) (CustomEngine, error) { return pm.createEngineFromLua(pluginName, engineName, factoryFn, config, logger) } pm.engineRegistry.RegisterEngine(engineName, factory) p.mu.Lock() p.isEngine = true p.engineName = engineName p.mu.Unlock() if pm.logger != nil { pm.logger.Info(fmt.Sprintf("Engine registered from plugin %s: %s", pluginName, engineName)) } return nil } // createEngineFromLua создаёт экземпляр Lua-движка. func (pm *PluginManager) createEngineFromLua(pluginName, engineName string, factoryFn *lua.LFunction, config map[string]interface{}, logger storage.LoggerInterface) (CustomEngine, error) { val, ok := pm.plugins.Load(pluginName) if !ok { return nil, fmt.Errorf("plugin not found: %s", pluginName) } L, p := getPluginLState(val) if L == nil || p == nil || p.IsClosed() { return nil, fmt.Errorf("plugin has no Lua state or is closed") } configTable := pm.goValueToLua(L, config).(*lua.LTable) if err := L.CallByParam(lua.P{ Fn: factoryFn, NRet: 1, Protect: true, }, configTable); err != nil { return nil, fmt.Errorf("failed to create engine instance: %v", err) } engineObj := L.Get(-1) L.Pop(1) if engineObj.Type() != lua.LTUserData { return nil, fmt.Errorf("factory function must return a userdata object") } engine := &LuaCustomEngine{ pluginName: pluginName, engineName: engineName, L: L, engineObj: engineObj.(*lua.LUserData), logger: logger, pm: pm, } return engine, nil } // ========== LuaCustomEngine - обёртка для Lua-движка ========== // LuaCustomEngine реализует CustomEngine для Lua-плагинов. type LuaCustomEngine struct { pluginName string engineName string L *lua.LState engineObj *lua.LUserData logger storage.LoggerInterface pm *PluginManager mu sync.RWMutex initialized bool closed bool } // Name возвращает имя движка. func (e *LuaCustomEngine) Name() string { return e.engineName } // Version возвращает версию движка. func (e *LuaCustomEngine) Version() string { return e.callStringMethod("version", nil) } // Initialize инициализирует движок с переданной конфигурацией. func (e *LuaCustomEngine) Initialize(config map[string]interface{}) error { e.mu.Lock() defer e.mu.Unlock() configTable := e.pm.goValueToLua(e.L, config).(*lua.LTable) err := e.callMethod("initialize", configTable) if err != nil { return err } e.initialized = true return nil } // Insert вставляет документ в движок. func (e *LuaCustomEngine) Insert(doc *storage.Document) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } docTable := e.pm.goValueToLua(e.L, doc.ToMap()).(*lua.LTable) return e.callMethod("insert", docTable) } // Find находит документ по ID. func (e *LuaCustomEngine) Find(id string) (*storage.Document, error) { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return nil, fmt.Errorf("engine not initialized or closed") } result := e.callMethodWithReturn("find", lua.LString(id)) if result == nil || result == lua.LNil { return nil, fmt.Errorf("document not found: %s", id) } if table, ok := result.(*lua.LTable); ok { return e.tableToDocument(table), nil } return nil, fmt.Errorf("invalid return type from find") } // Update обновляет документ. func (e *LuaCustomEngine) Update(id string, updates map[string]interface{}) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } updatesTable := e.pm.goValueToLua(e.L, updates).(*lua.LTable) return e.callMethod("update", lua.LString(id), updatesTable) } // Delete удаляет документ. func (e *LuaCustomEngine) Delete(id string) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } return e.callMethod("delete", lua.LString(id)) } // BatchInsert вставляет несколько документов. func (e *LuaCustomEngine) BatchInsert(docs []*storage.Document) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } docsTable := e.L.NewTable() for i, doc := range docs { docTable := e.pm.goValueToLua(e.L, doc.ToMap()).(*lua.LTable) docsTable.RawSetInt(i+1, docTable) } return e.callMethod("batch_insert", docsTable) } // BatchUpdate выполняет пакетное обновление. func (e *LuaCustomEngine) BatchUpdate(updates map[string]map[string]interface{}) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } updatesTable := e.pm.goValueToLua(e.L, updates).(*lua.LTable) return e.callMethod("batch_update", updatesTable) } // BatchDelete выполняет пакетное удаление. func (e *LuaCustomEngine) BatchDelete(ids []string) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } idsTable := e.L.NewTable() for i, id := range ids { idsTable.RawSetInt(i+1, lua.LString(id)) } return e.callMethod("batch_delete", idsTable) } // FindByFilter находит документы по фильтру. func (e *LuaCustomEngine) FindByFilter(filter func(*storage.Document) bool) ([]*storage.Document, error) { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return nil, fmt.Errorf("engine not initialized or closed") } result := e.callMethodWithReturn("find_all", nil) if result == nil || result == lua.LNil { return []*storage.Document{}, nil } docs := make([]*storage.Document, 0) if table, ok := result.(*lua.LTable); ok { table.ForEach(func(key, value lua.LValue) { if docTable, ok := value.(*lua.LTable); ok { doc := e.tableToDocument(docTable) if filter(doc) { docs = append(docs, doc) } } }) } return docs, nil } // FindByIndex находит документы по индексу. func (e *LuaCustomEngine) FindByIndex(indexName string, value interface{}) ([]*storage.Document, error) { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return nil, fmt.Errorf("engine not initialized or closed") } luaValue := e.pm.goValueToLua(e.L, value) result := e.callMethodWithReturn("find_by_index", lua.LString(indexName), luaValue) if result == nil || result == lua.LNil { return []*storage.Document{}, nil } docs := make([]*storage.Document, 0) if table, ok := result.(*lua.LTable); ok { table.ForEach(func(key, value lua.LValue) { if docTable, ok := value.(*lua.LTable); ok { docs = append(docs, e.tableToDocument(docTable)) } }) } return docs, nil } // CreateIndex создаёт индекс. func (e *LuaCustomEngine) CreateIndex(name string, fields []string, unique bool) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } fieldsTable := e.L.NewTable() for i, field := range fields { fieldsTable.RawSetInt(i+1, lua.LString(field)) } return e.callMethod("create_index", lua.LString(name), fieldsTable, lua.LBool(unique)) } // DropIndex удаляет индекс. func (e *LuaCustomEngine) DropIndex(name string) error { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return fmt.Errorf("engine not initialized or closed") } return e.callMethod("drop_index", lua.LString(name)) } // GetIndexes возвращает список индексов. func (e *LuaCustomEngine) GetIndexes() []string { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return []string{} } result := e.callMethodWithReturn("get_indexes", nil) if result == nil || result == lua.LNil { return []string{} } indexes := make([]string, 0) if table, ok := result.(*lua.LTable); ok { table.ForEach(func(key, value lua.LValue) { if str, ok := value.(lua.LString); ok { indexes = append(indexes, string(str)) } }) } return indexes } // Count возвращает количество документов. func (e *LuaCustomEngine) Count() int64 { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return 0 } result := e.callMethodWithReturn("count", nil) if result == nil || result == lua.LNil { return 0 } if num, ok := result.(lua.LNumber); ok { return int64(num) } return 0 } // Size возвращает размер данных в байтах. func (e *LuaCustomEngine) Size() int64 { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return 0 } result := e.callMethodWithReturn("size", nil) if result == nil || result == lua.LNil { return 0 } if num, ok := result.(lua.LNumber); ok { return int64(num) } return 0 } // GetStats возвращает статистику движка. func (e *LuaCustomEngine) GetStats() map[string]interface{} { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return make(map[string]interface{}) } result := e.callMethodWithReturn("get_stats", nil) if result == nil || result == lua.LNil { return make(map[string]interface{}) } if table, ok := result.(*lua.LTable); ok { return e.pm.luaTableToMap(table) } return make(map[string]interface{}) } // Close закрывает движок и освобождает ресурсы. func (e *LuaCustomEngine) Close() error { e.mu.Lock() defer e.mu.Unlock() if e.closed { return nil } if e.initialized { _ = e.callMethod("close", nil) } e.closed = true return nil } // OnDocumentInserted вызывается при вставке документа. func (e *LuaCustomEngine) OnDocumentInserted(doc *storage.Document) { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return } docTable := e.pm.goValueToLua(e.L, doc.ToMap()).(*lua.LTable) _ = e.callMethod("on_document_inserted", docTable) } // OnDocumentUpdated вызывается при обновлении документа. func (e *LuaCustomEngine) OnDocumentUpdated(oldDoc, newDoc *storage.Document) { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return } oldTable := e.pm.goValueToLua(e.L, oldDoc.ToMap()).(*lua.LTable) newTable := e.pm.goValueToLua(e.L, newDoc.ToMap()).(*lua.LTable) _ = e.callMethod("on_document_updated", oldTable, newTable) } // OnDocumentDeleted вызывается при удалении документа. func (e *LuaCustomEngine) OnDocumentDeleted(doc *storage.Document) { e.mu.RLock() defer e.mu.RUnlock() if !e.initialized || e.closed { return } docTable := e.pm.goValueToLua(e.L, doc.ToMap()).(*lua.LTable) _ = e.callMethod("on_document_deleted", docTable) } // ========== Вспомогательные методы LuaCustomEngine ========== // callMethod вызывает метод Lua-объекта без возвращаемого значения. func (e *LuaCustomEngine) callMethod(methodName string, args ...lua.LValue) error { method := e.L.GetField(e.engineObj, methodName) if method.Type() != lua.LTFunction { return fmt.Errorf("method %s not found in engine", methodName) } callArgs := make([]lua.LValue, 0, len(args)+1) callArgs = append(callArgs, e.engineObj) callArgs = append(callArgs, args...) if err := e.L.CallByParam(lua.P{ Fn: method, NRet: 0, Protect: true, }, callArgs...); err != nil { return fmt.Errorf("failed to call method %s: %v", methodName, err) } return nil } // callMethodWithReturn вызывает метод Lua-объекта с возвращаемым значением. func (e *LuaCustomEngine) callMethodWithReturn(methodName string, args ...lua.LValue) lua.LValue { method := e.L.GetField(e.engineObj, methodName) if method.Type() != lua.LTFunction { return lua.LNil } callArgs := make([]lua.LValue, 0, len(args)+1) callArgs = append(callArgs, e.engineObj) callArgs = append(callArgs, args...) if err := e.L.CallByParam(lua.P{ Fn: method, NRet: 1, Protect: true, }, callArgs...); err != nil { if e.logger != nil { e.logger.Error(fmt.Sprintf("Failed to call method %s: %v", methodName, err)) } return lua.LNil } result := e.L.Get(-1) e.L.Pop(1) return result } // callStringMethod вызывает метод и возвращает строковый результат. func (e *LuaCustomEngine) callStringMethod(methodName string, args ...lua.LValue) string { result := e.callMethodWithReturn(methodName, args...) if result == nil || result == lua.LNil { return "" } if str, ok := result.(lua.LString); ok { return string(str) } return "" } // tableToDocument преобразует Lua-таблицу в документ. func (e *LuaCustomEngine) tableToDocument(table *lua.LTable) *storage.Document { doc := storage.NewDocument() table.ForEach(func(key, value lua.LValue) { keyStr := key.String() if strings.HasPrefix(keyStr, "_") { return } goValue := e.pm.luaValueToGo(value) doc.SetField(keyStr, goValue) }) return doc } // ========== PluginManager методы ========== // autoLoadPlugins автоматически загружает все .lua файлы из директории плагинов. func (pm *PluginManager) autoLoadPlugins() { if !pm.enabled { return } pm.loadPluginsFromDir() ticker := time.NewTicker(10 * time.Second) defer ticker.Stop() for range ticker.C { pm.loadPluginsFromDir() } } // loadPluginsFromDir загружает все плагины из директории. func (pm *PluginManager) loadPluginsFromDir() { entries, err := os.ReadDir(pm.pluginsDir) if err != nil { if pm.logger != nil { pm.logger.Error(fmt.Sprintf("Failed to read plugins directory: %v", err)) } return } for _, entry := range entries { if entry.IsDir() { continue } name := entry.Name() if !strings.HasSuffix(name, ".lua") { continue } pluginName := strings.TrimSuffix(name, ".lua") if _, exists := pm.plugins.Load(pluginName); !exists { pluginPath := filepath.Join(pm.pluginsDir, name) if err := pm.LoadPluginWithSandbox(pluginName, pluginPath, SandboxLimits{ MaxCPUTime: DefaultCPULimit, MaxMemory: DefaultMemoryLimit, MaxExecutionTime: DefaultExecutionTimeout, MaxInstructions: 1000000, }); err != nil { if pm.logger != nil { pm.logger.Error(fmt.Sprintf("Failed to auto-load plugin %s: %v", pluginName, err)) } } else if pm.logger != nil { pm.logger.Info(fmt.Sprintf("Auto-loaded plugin: %s", pluginName)) } } } } // LoadPlugin загружает Lua-плагин из файла с ограничениями по умолчанию. func (pm *PluginManager) LoadPlugin(name, filePath string) error { return pm.LoadPluginWithSandbox(name, filePath, SandboxLimits{ MaxCPUTime: DefaultCPULimit, MaxMemory: DefaultMemoryLimit, MaxExecutionTime: DefaultExecutionTimeout, MaxInstructions: 1000000, }) } // LoadPluginWithSandbox загружает плагин с ограничениями, используя пул состояний. // // ИЗМЕНЕНО (2026-10-03): формат loaded_at_str теперь "02-01-2006 15:04:05.000". func (pm *PluginManager) LoadPluginWithSandbox(name, filePath string, limits SandboxLimits) error { if !pm.enabled { return fmt.Errorf("plugin system is disabled") } if pm.statePool.GetActiveCount() >= int32(MaxLuaStates) { return fmt.Errorf("max Lua states limit reached: %d", MaxLuaStates) } script, err := os.ReadFile(filePath) if err != nil { return fmt.Errorf("failed to read plugin file: %v", err) } manifest := pm.loadManifest(filePath) if manifest != nil && len(manifest.Dependencies) > 0 { if err := pm.checkDependencies(manifest.Dependencies); err != nil { return fmt.Errorf("dependency check failed: %v", err) } } L, err := pm.statePool.Acquire(limits) if err != nil { return fmt.Errorf("failed to acquire Lua state: %v", err) } // Регистрируем API до выполнения скрипта. // Для транзакционных функций нам нужна ссылка на Plugin, // но она создаётся только после DoString. Поэтому регистрируем // функции, которые будут искать плагин по имени через pm.plugins. // Это работает, потому что к моменту вызова Lua-функции плагин // уже будет зарегистрирован в pm.plugins. pm.registerDatabaseFunctions(L.LState) pm.registerTransactionFunctions(L.LState, name) pm.registerTriggerFunctions(L.LState) pm.registerTimestampFunctions(L.LState) pm.registerEngineFunctions(L.LState, name) // Выполняем скрипт с таймаутом. // Ограничение по времени работает через L.SetContext(ctx), // который был установлен в NewInstrumentedLState/Reset/SetLimits. done := make(chan error, 1) go func() { defer func() { if r := recover(); r != nil { done <- fmt.Errorf("panic during script execution: %v", r) } }() done <- L.DoString(string(script)) }() select { case err := <-done: if err != nil { pm.statePool.Release(L) return fmt.Errorf("failed to execute plugin script: %v", err) } case <-time.After(MaxPluginLoadTime): // Отменяем через контекст L.mu.RLock() cancel := L.cancel L.mu.RUnlock() if cancel != nil { cancel() } // Ждём завершения горутины, чтобы не было утечки select { case <-done: case <-time.After(1 * time.Second): } pm.statePool.Release(L) return fmt.Errorf("plugin load timeout") } version := pm.getPluginMetadata(L.LState, "version") author := pm.getPluginMetadata(L.LState, "author") description := pm.getPluginMetadata(L.LState, "description") isEngine := pm.getPluginMetadata(L.LState, "engine_type") != "" engineName := pm.getPluginMetadata(L.LState, "engine_type") now := time.Now() nowMs := now.UnixMilli() // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". nowStr := now.Format("02-01-2006 15:04:05.000") basePlugin := &Plugin{ Name: name, FilePath: filePath, LState: L.LState, instrumented: L, logger: pm.logger, storage: pm.storage, loadedAt: now, loadedAtMs: nowMs, loadedAtStr: nowStr, version: version, author: author, description: description, eventLog: make([]PluginEventLog, 0), maxEventLog: DefaultMaxEventLogSize, isEngine: isEngine, engineName: engineName, } basePlugin.Status.Store(int32(StatusLoaded)) var plugin interface{} if manifest != nil { pluginWithDeps := &PluginWithDeps{ Plugin: basePlugin, Manifest: manifest, Dependencies: make(map[string]*PluginWithDeps), Dependents: make(map[string]*PluginWithDeps), } plugin = pluginWithDeps pm.depResolver.RegisterPlugin(pluginWithDeps) } else { plugin = basePlugin } pm.plugins.Store(name, plugin) if pm.hotReload != nil { pm.hotReload.WatchPlugin(filePath) } if pm.logger != nil { engineInfo := "" if isEngine { engineInfo = fmt.Sprintf(" [ENGINE: %s]", engineName) } pm.logger.Info(fmt.Sprintf("Plugin loaded: %s v%s by %s - %s%s at %s", name, version, author, description, engineInfo, nowStr)) } auditTimeMs, auditTimeStr := storage.GetCurrentTimestamp() storage.LogAudit("PLUGIN_LOAD", "PLUGIN", name, map[string]interface{}{ "version": version, "author": author, "description": description, "loaded_at": nowMs, "loaded_at_str": nowStr, "audit_time": auditTimeMs, "audit_time_str": auditTimeStr, "limits_cpu_ms": limits.MaxCPUTime.Milliseconds(), "limits_memory_mb": limits.MaxMemory / 1024 / 1024, "active_lua_states": pm.statePool.GetActiveCount(), "is_engine": isEngine, "engine_type": engineName, }) if err := pm.callPluginFunction(plugin, "on_load"); err != nil { if pm.logger != nil { pm.logger.Warn(fmt.Sprintf("Plugin %s on_load error: %v", name, err)) } pm.addPluginEventLog(plugin, "on_load_error", nil, err.Error()) } else { pm.addPluginEventLog(plugin, "on_load", nil, "") } if isEngine && engineName != "" { if err := pm.registerEngineFromPlugin(plugin, engineName); err != nil { if pm.logger != nil { pm.logger.Warn(fmt.Sprintf("Failed to register engine from plugin %s: %v", name, err)) } } } return nil } // registerEngineFromPlugin регистрирует движок из загруженного плагина. func (pm *PluginManager) registerEngineFromPlugin(plugin interface{}, engineName string) error { L, p := getPluginLState(plugin) if L == nil || p == nil { return fmt.Errorf("plugin has no Lua state") } factoryFn := L.GetGlobal("create_engine") if factoryFn.Type() != lua.LTFunction { factoryFn = L.GetGlobal("new_engine") if factoryFn.Type() != lua.LTFunction { return fmt.Errorf("engine plugin must export create_engine or new_engine function") } } return pm.RegisterEngineFromLua(p.Name, engineName, factoryFn.(*lua.LFunction)) } // loadManifest загружает манифест плагина из JSON-файла. func (pm *PluginManager) loadManifest(pluginPath string) *PluginManifest { manifestPath := strings.TrimSuffix(pluginPath, ".lua") + ".json" data, err := os.ReadFile(manifestPath) if err != nil { return nil } var manifest PluginManifest if err := json.Unmarshal(data, &manifest); err != nil { return nil } return &manifest } // checkDependencies проверяет зависимости плагина. func (pm *PluginManager) checkDependencies(deps []PluginDependency) error { for _, dep := range deps { if dep.Optional { continue } val, ok := pm.plugins.Load(dep.Name) if !ok { return fmt.Errorf("missing dependency: %s", dep.Name) } if withDeps, ok := val.(*PluginWithDeps); ok && withDeps.Manifest != nil { if !pm.depResolver.CheckVersionCompatibility(dep, withDeps.Manifest.Version) { return fmt.Errorf("version incompatibility: %s requires %s but have %s", dep.Name, dep.Version, withDeps.Manifest.Version) } } } return nil } // registerTimestampFunctions регистрирует функции для работы с временными метками в Lua. // // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". func (pm *PluginManager) registerTimestampFunctions(L *lua.LState) { L.SetGlobal("get_current_timestamp_ms", L.NewFunction(func(L *lua.LState) int { L.Push(lua.LNumber(time.Now().UnixMilli())) return 1 })) L.SetGlobal("get_current_timestamp_str", L.NewFunction(func(L *lua.LState) int { L.Push(lua.LString(time.Now().Format("02-01-2006 15:04:05.000"))) return 1 })) L.SetGlobal("get_current_time", L.NewFunction(func(L *lua.LState) int { now := time.Now() table := L.NewTable() table.RawSetString("unix_ms", lua.LNumber(now.UnixMilli())) table.RawSetString("unix_sec", lua.LNumber(now.Unix())) table.RawSetString("string", lua.LString(now.Format("02-01-2006 15:04:05.000"))) table.RawSetString("year", lua.LNumber(now.Year())) table.RawSetString("month", lua.LNumber(int(now.Month()))) table.RawSetString("day", lua.LNumber(now.Day())) table.RawSetString("hour", lua.LNumber(now.Hour())) table.RawSetString("minute", lua.LNumber(now.Minute())) table.RawSetString("second", lua.LNumber(now.Second())) L.Push(table) return 1 })) } // registerDatabaseFunctions регистрирует функции доступа к СУБД в Lua. func (pm *PluginManager) registerDatabaseFunctions(L *lua.LState) { L.SetGlobal("get_database", L.NewFunction(func(L *lua.LState) int { if pm.storage == nil { L.Push(lua.LNil) L.Push(lua.LString("storage not available")) return 2 } dbName := L.CheckString(1) db, err := pm.storage.GetDatabase(dbName) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } ud := L.NewUserData() ud.Value = db L.SetMetatable(ud, L.GetTypeMetatable("database")) L.Push(ud) return 1 })) L.SetGlobal("get_collection", L.NewFunction(func(L *lua.LState) int { if pm.storage == nil { L.Push(lua.LNil) L.Push(lua.LString("storage not available")) return 2 } dbName := L.CheckString(1) collName := L.CheckString(2) db, err := pm.storage.GetDatabase(dbName) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } coll, err := db.GetCollection(collName) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } ud := L.NewUserData() ud.Value = coll L.SetMetatable(ud, L.GetTypeMetatable("collection")) L.Push(ud) return 1 })) L.SetGlobal("plugin_log", L.NewFunction(func(L *lua.LState) int { level := L.CheckString(1) message := L.CheckString(2) if pm.logger != nil { logMsg := fmt.Sprintf("[PLUGIN] %s: %s", level, message) switch level { case "debug": pm.logger.Debug(logMsg) case "info": pm.logger.Info(logMsg) case "warn": pm.logger.Warn(logMsg) case "error": pm.logger.Error(logMsg) default: pm.logger.Info(logMsg) } } return 0 })) L.SetGlobal("emit_event", L.NewFunction(func(L *lua.LState) int { // Проверяем, не закрыт ли eventBus, чтобы избежать // паники "send on closed channel". if pm.eventBusClosed.Load() { return 0 } eventType := L.CheckString(1) eventData := L.CheckAny(2) nowMs := time.Now().UnixMilli() // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". nowStr := time.Now().Format("02-01-2006 15:04:05.000") event := PluginEvent{ EventType: eventType, Data: pm.luaValueToGo(eventData), Timestamp: nowMs, TimestampStr: nowStr, } select { case pm.eventBus <- event: default: if pm.logger != nil { pm.logger.Warn("Plugin event bus full, event dropped") } } return 0 })) pm.setupDatabaseMetatable(L) pm.setupCollectionMetatable(L) } // registerEngineFunctions регистрирует функции для создания кастомных движков. func (pm *PluginManager) registerEngineFunctions(L *lua.LState, pluginName string) { L.SetGlobal("register_engine", L.NewFunction(func(L *lua.LState) int { engineName := L.CheckString(1) factoryFn := L.CheckFunction(2) if err := pm.RegisterEngineFromLua(pluginName, engineName, factoryFn); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) L.SetGlobal("get_engine", L.NewFunction(func(L *lua.LState) int { engineName := L.CheckString(1) engine := pm.engineRegistry.GetEngine(engineName) if engine == nil { L.Push(lua.LNil) L.Push(lua.LString(fmt.Sprintf("engine not found: %s", engineName))) return 2 } ud := L.NewUserData() ud.Value = engine L.SetMetatable(ud, L.GetTypeMetatable("custom_engine")) L.Push(ud) return 1 })) L.SetGlobal("list_engines", L.NewFunction(func(L *lua.LState) int { engines := pm.engineRegistry.ListEngines() table := L.NewTable() for i, name := range engines { table.RawSetInt(i+1, lua.LString(name)) } L.Push(table) return 1 })) L.SetGlobal("list_active_engines", L.NewFunction(func(L *lua.LState) int { engines := pm.engineRegistry.ListActiveEngines() table := L.NewTable() for i, name := range engines { table.RawSetInt(i+1, lua.LString(name)) } L.Push(table) return 1 })) // Настройка метатаблицы для кастомного движка engineMt := L.NewTypeMetatable("custom_engine") L.SetField(engineMt, "__index", L.NewFunction(func(L *lua.LState) int { engine := L.CheckUserData(1).Value.(CustomEngine) method := L.CheckString(2) switch method { case "name": L.Push(lua.LString(engine.Name())) case "version": L.Push(lua.LString(engine.Version())) case "count": L.Push(L.NewFunction(func(L *lua.LState) int { L.Push(lua.LNumber(engine.Count())) return 1 })) case "size": L.Push(L.NewFunction(func(L *lua.LState) int { L.Push(lua.LNumber(engine.Size())) return 1 })) case "insert": L.Push(L.NewFunction(func(L *lua.LState) int { docData := L.CheckTable(1) doc := storage.NewDocument() docData.ForEach(func(key, value lua.LValue) { if key.Type() == lua.LTString { doc.SetField(key.String(), pm.luaValueToGo(value)) } }) if err := engine.Insert(doc); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "find": L.Push(L.NewFunction(func(L *lua.LState) int { id := L.CheckString(1) doc, err := engine.Find(id) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } resultTable := pm.goValueToLua(L, doc.ToMap()).(*lua.LTable) L.Push(resultTable) return 1 })) case "update": L.Push(L.NewFunction(func(L *lua.LState) int { id := L.CheckString(1) updates := L.CheckTable(2) updatesMap := pm.luaTableToMap(updates) if err := engine.Update(id, updatesMap); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "delete": L.Push(L.NewFunction(func(L *lua.LState) int { id := L.CheckString(1) if err := engine.Delete(id); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "create_index": L.Push(L.NewFunction(func(L *lua.LState) int { name := L.CheckString(1) fieldsTable := L.CheckTable(2) unique := L.OptBool(3, false) fields := make([]string, 0) fieldsTable.ForEach(func(key, value lua.LValue) { if str, ok := value.(lua.LString); ok { fields = append(fields, string(str)) } }) if err := engine.CreateIndex(name, fields, unique); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "drop_index": L.Push(L.NewFunction(func(L *lua.LState) int { name := L.CheckString(1) if err := engine.DropIndex(name); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "get_stats": L.Push(L.NewFunction(func(L *lua.LState) int { stats := engine.GetStats() statsTable := pm.goValueToLua(L, stats).(*lua.LTable) L.Push(statsTable) return 1 })) default: L.Push(lua.LNil) } return 1 })) } // registerTransactionFunctions регистрирует функции для работы с batch-операциями // в Lua. Сохранены старые имена (begin_transaction и т.д.) для совместимости // с существующими Lua-плагинами, но семантически это batch-операции. // // ИЗМЕНЕНО (2026-10): storage.BeginTransaction и связанные функции были // удалены при переходе на batch. Теперь: // - begin_transaction создаёт *storage.Batch через storage.NewBatch // и регистрирует его в Plugin.activeBatches по ID. // - commit_transaction извлекает batch из Plugin.activeBatches и // вызывает batch.Commit(). // - abort_transaction просто снимает регистрацию batch (не коммитит). // - has_active_transaction проверяет, есть ли у плагина активные batch'и. // - get_current_transaction_id возвращает ID последнего созданного batch'а. // - get_active_transactions возвращает список активных batch'ей плагина. // // Параметр pluginName используется для поиска Plugin в pm.plugins. func (pm *PluginManager) registerTransactionFunctions(L *lua.LState, pluginName string) { // Вспомогательная функция для получения текущего плагина. getCurrentPlugin := func() *Plugin { val, ok := pm.plugins.Load(pluginName) if !ok { return nil } _, p := getPluginLState(val) return p } L.SetGlobal("begin_transaction", L.NewFunction(func(L *lua.LState) int { p := getCurrentPlugin() if p == nil { L.Push(lua.LNil) L.Push(lua.LString("plugin not registered")) return 2 } // Проверяем, что storage доступен. if pm.storage == nil { L.Push(lua.LNil) L.Push(lua.LString("storage not available")) return 2 } // Читаем параметры: database, collection. database := L.OptString(1, "") collection := L.OptString(2, "") batch := storage.NewBatch(database, collection) if batch == nil { L.Push(lua.LNil) L.Push(lua.LString("failed to create batch")) return 2 } batchID := pm.nextBatchID.Add(1) p.RegisterBatch(batchID, batch) ud := L.NewUserData() ud.Value = batch L.SetMetatable(ud, L.GetTypeMetatable("transaction")) L.SetField(ud, "id", lua.LNumber(batchID)) L.Push(ud) return 1 })) L.SetGlobal("commit_transaction", L.NewFunction(func(L *lua.LState) int { p := getCurrentPlugin() if p == nil { L.Push(lua.LString("plugin not registered")) return 1 } // Пробуем получить batch из userdata, если передан. var batchID uint64 if L.GetTop() >= 1 { ud := L.CheckUserData(1) if idVal := L.GetField(ud, "id"); idVal != lua.LNil { if num, ok := idVal.(lua.LNumber); ok { batchID = uint64(num) } } } // Если batchID == 0 — коммитим последний активный batch. // Для этого используем любой batch из activeBatches. if batchID == 0 { var found *storage.Batch var foundID uint64 p.activeBatches.Range(func(key, value interface{}) bool { if b, ok := value.(*storage.Batch); ok { found = b if id, ok := key.(uint64); ok { foundID = id } } return false }) if found == nil { L.Push(lua.LString("no active batch")) return 1 } batchID = foundID } batch, ok := p.GetBatch(batchID) if !ok { L.Push(lua.LString("batch not found")) return 1 } if err := batch.Commit(); err != nil { p.UnregisterBatch(batchID) L.Push(lua.LString(err.Error())) return 1 } // Commit уже снял регистрацию через defer, но на всякий случай. p.UnregisterBatch(batchID) L.Push(lua.LNil) return 1 })) L.SetGlobal("abort_transaction", L.NewFunction(func(L *lua.LState) int { p := getCurrentPlugin() if p == nil { L.Push(lua.LString("plugin not registered")) return 1 } var batchID uint64 if L.GetTop() >= 1 { ud := L.CheckUserData(1) if idVal := L.GetField(ud, "id"); idVal != lua.LNil { if num, ok := idVal.(lua.LNumber); ok { batchID = uint64(num) } } } if batchID == 0 { // Снимаем регистрацию со всех batch'ей плагина. p.activeBatches.Range(func(key, value interface{}) bool { if id, ok := key.(uint64); ok { p.UnregisterBatch(id) } return true }) } else { p.UnregisterBatch(batchID) } // batch просто забывается — его операции не применяются. // Регистрация в globalBatchManager.activeBatches снимается // через UnregisterBatch (мы вызываем его внутри). if bm := storage.GetBatchManager(); bm != nil { // Если batch всё ещё зарегистрирован в глобальном менеджере, // снимаем и там. if b, ok := p.GetBatch(batchID); ok && b != nil { bm.UnregisterBatch(b) } } L.Push(lua.LNil) return 1 })) L.SetGlobal("has_active_transaction", L.NewFunction(func(L *lua.LState) int { p := getCurrentPlugin() if p == nil { L.Push(lua.LBool(false)) return 1 } hasActive := false p.activeBatches.Range(func(key, value interface{}) bool { hasActive = true return false }) L.Push(lua.LBool(hasActive)) return 1 })) L.SetGlobal("get_current_transaction_id", L.NewFunction(func(L *lua.LState) int { p := getCurrentPlugin() if p == nil { L.Push(lua.LNil) return 1 } var lastID uint64 p.activeBatches.Range(func(key, value interface{}) bool { if id, ok := key.(uint64); ok { if id > lastID { lastID = id } } return true }) if lastID == 0 { L.Push(lua.LNil) } else { L.Push(lua.LString(fmt.Sprintf("%d", lastID))) } return 1 })) L.SetGlobal("get_active_transactions", L.NewFunction(func(L *lua.LState) int { p := getCurrentPlugin() table := L.NewTable() if p == nil { L.Push(table) return 1 } idx := 0 p.activeBatches.Range(func(key, value interface{}) bool { batch, ok := value.(*storage.Batch) if !ok { return true } idx++ txTable := L.NewTable() txTable.RawSetString("id", lua.LString(fmt.Sprintf("%d", batch.ID))) txTable.RawSetString("status", lua.LString("active")) txTable.RawSetString("start_time", lua.LNumber(batch.Timestamp)) txTable.RawSetString("operation_count", lua.LNumber(len(batch.Operations))) table.RawSetInt(idx, txTable) return true }) L.Push(table) return 1 })) // Настройка метатаблицы для batch'ей (для совместимости с "transaction"). mt := L.NewTypeMetatable("transaction") L.SetField(mt, "__index", L.NewFunction(func(L *lua.LState) int { ud := L.CheckUserData(1) batch, ok := ud.Value.(*storage.Batch) if !ok { L.Push(lua.LNil) return 1 } method := L.CheckString(2) switch method { case "get_id": L.Push(lua.LString(fmt.Sprintf("%d", batch.ID))) case "get_operation_count": L.Push(lua.LNumber(len(batch.Operations))) case "get_start_time": L.Push(lua.LNumber(batch.Timestamp)) case "get_status": L.Push(lua.LString("active")) case "add_insert": L.Push(L.NewFunction(func(L *lua.LState) int { docData := L.CheckTable(1) doc := storage.NewDocument() docData.ForEach(func(key, value lua.LValue) { if key.Type() == lua.LTString { doc.SetField(key.String(), pm.luaValueToGo(value)) } }) batch.AddInsert(doc) L.Push(lua.LNil) return 1 })) case "add_update": L.Push(L.NewFunction(func(L *lua.LState) int { docID := L.CheckString(1) updates := L.CheckTable(2) updatesMap := pm.luaTableToMap(updates) batch.AddUpdate(docID, updatesMap) L.Push(lua.LNil) return 1 })) case "add_delete": L.Push(L.NewFunction(func(L *lua.LState) int { docID := L.CheckString(1) batch.AddDelete(docID) L.Push(lua.LNil) return 1 })) case "add_restore": L.Push(L.NewFunction(func(L *lua.LState) int { docID := L.CheckString(1) batch.AddRestore(docID) L.Push(lua.LNil) return 1 })) case "commit": L.Push(L.NewFunction(func(L *lua.LState) int { if err := batch.Commit(); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) default: L.Push(lua.LNil) } return 1 })) } // registerTriggerFunctions регистрирует функции для работы с триггерами. func (pm *PluginManager) registerTriggerFunctions(L *lua.LState) { tm := storage.GetTriggerManager() L.SetGlobal("create_trigger", L.NewFunction(func(L *lua.LState) int { if tm == nil { L.Push(lua.LString("trigger manager not available")) return 1 } database := L.CheckString(1) collection := L.CheckString(2) name := L.CheckString(3) event := L.CheckString(4) config := L.CheckTable(5) configMap := make(map[string]interface{}) config.ForEach(func(key, value lua.LValue) { if key.Type() == lua.LTString { configMap[key.String()] = pm.luaValueToGo(value) } }) if err := tm.CreateTrigger(database, collection, name, storage.TriggerEvent(event), configMap); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) L.SetGlobal("drop_trigger", L.NewFunction(func(L *lua.LState) int { if tm == nil { L.Push(lua.LString("trigger manager not available")) return 1 } collection := L.CheckString(1) event := L.CheckString(2) name := L.CheckString(3) if err := tm.DropTrigger(collection, event, name); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) L.SetGlobal("enable_trigger", L.NewFunction(func(L *lua.LState) int { if tm == nil { L.Push(lua.LString("trigger manager not available")) return 1 } collection := L.CheckString(1) event := L.CheckString(2) name := L.CheckString(3) if err := tm.EnableTrigger(collection, event, name); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) L.SetGlobal("disable_trigger", L.NewFunction(func(L *lua.LState) int { if tm == nil { L.Push(lua.LString("trigger manager not available")) return 1 } collection := L.CheckString(1) event := L.CheckString(2) name := L.CheckString(3) if err := tm.DisableTrigger(collection, event, name); err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) L.SetGlobal("list_triggers", L.NewFunction(func(L *lua.LState) int { if tm == nil { L.Push(L.NewTable()) return 1 } collection := L.OptString(1, "") triggers := tm.ListTriggers(collection) table := L.NewTable() for i, trigger := range triggers { triggerTable := L.NewTable() triggerTable.RawSetString("name", lua.LString(trigger.Name)) triggerTable.RawSetString("collection", lua.LString(trigger.Collection)) triggerTable.RawSetString("event", lua.LString(string(trigger.Event))) triggerTable.RawSetString("action", lua.LString(string(trigger.Action))) triggerTable.RawSetString("enabled", lua.LBool(trigger.Enabled)) triggerTable.RawSetString("description", lua.LString(trigger.Description)) table.RawSetInt(i+1, triggerTable) } L.Push(table) return 1 })) } // setupDatabaseMetatable настраивает методы для объекта базы данных в Lua. func (pm *PluginManager) setupDatabaseMetatable(L *lua.LState) { mt := L.NewTypeMetatable("database") L.SetField(mt, "__index", L.NewFunction(func(L *lua.LState) int { db := L.CheckUserData(1).Value.(*storage.Database) method := L.CheckString(2) switch method { case "create_collection": L.Push(L.NewFunction(func(L *lua.LState) int { name := L.CheckString(1) err := db.CreateCollection(name) if err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "get_collection": L.Push(L.NewFunction(func(L *lua.LState) int { name := L.CheckString(1) coll, err := db.GetCollection(name) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } ud := L.NewUserData() ud.Value = coll L.SetMetatable(ud, L.GetTypeMetatable("collection")) L.Push(ud) return 1 })) case "drop_collection": L.Push(L.NewFunction(func(L *lua.LState) int { name := L.CheckString(1) err := db.DropCollection(name) if err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "list_collections": L.Push(L.NewFunction(func(L *lua.LState) int { collections := db.ListCollections() table := L.NewTable() for i, name := range collections { table.RawSetInt(i+1, lua.LString(name)) } L.Push(table) return 1 })) case "name": L.Push(lua.LString(db.Name())) default: L.Push(lua.LNil) } return 1 })) } // setupCollectionMetatable настраивает методы для объекта коллекции в Lua. func (pm *PluginManager) setupCollectionMetatable(L *lua.LState) { mt := L.NewTypeMetatable("collection") L.SetField(mt, "__index", L.NewFunction(func(L *lua.LState) int { coll := L.CheckUserData(1).Value.(*storage.Collection) method := L.CheckString(2) switch method { case "insert": L.Push(L.NewFunction(func(L *lua.LState) int { doc := L.CheckTable(1) fields := make(map[string]interface{}) doc.ForEach(func(key, value lua.LValue) { if key.Type() == lua.LTString { fields[key.String()] = pm.luaValueToGo(value) } }) newDoc := storage.NewDocument() for k, v := range fields { newDoc.SetField(k, v) } err := coll.Insert(newDoc) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } L.Push(lua.LString(newDoc.ID)) L.Push(lua.LNil) return 2 })) case "find": L.Push(L.NewFunction(func(L *lua.LState) int { id := L.CheckString(1) doc, err := coll.Find(id) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } table := L.NewTable() table.RawSetString("_id", lua.LString(doc.ID)) for k, v := range doc.GetFields() { table.RawSetString(k, pm.goValueToLua(L, v)) } L.Push(table) return 1 })) case "find_by_index": L.Push(L.NewFunction(func(L *lua.LState) int { indexName := L.CheckString(1) value := pm.luaValueToGo(L.CheckAny(2)) docs, err := coll.FindByIndex(indexName, value) if err != nil { L.Push(lua.LNil) L.Push(lua.LString(err.Error())) return 2 } table := L.NewTable() for i, doc := range docs { docTable := L.NewTable() docTable.RawSetString("_id", lua.LString(doc.ID)) for k, v := range doc.GetFields() { docTable.RawSetString(k, pm.goValueToLua(L, v)) } table.RawSetInt(i+1, docTable) } L.Push(table) return 1 })) case "update": L.Push(L.NewFunction(func(L *lua.LState) int { id := L.CheckString(1) updates := L.CheckTable(2) fields := make(map[string]interface{}) updates.ForEach(func(key, value lua.LValue) { if key.Type() == lua.LTString { fields[key.String()] = pm.luaValueToGo(value) } }) err := coll.Update(id, fields) if err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "delete": L.Push(L.NewFunction(func(L *lua.LState) int { id := L.CheckString(1) err := coll.Delete(id) if err != nil { L.Push(lua.LString(err.Error())) return 1 } L.Push(lua.LNil) return 1 })) case "count": L.Push(L.NewFunction(func(L *lua.LState) int { count := coll.Count() L.Push(lua.LNumber(count)) return 1 })) case "name": L.Push(lua.LString(coll.Name())) default: L.Push(lua.LNil) } return 1 })) } // luaValueToGo конвертирует Lua-значение в Go-значение. func (pm *PluginManager) luaValueToGo(val lua.LValue) interface{} { if val == nil || val == lua.LNil { return nil } switch v := val.(type) { case lua.LString: return string(v) case lua.LNumber: return float64(v) case lua.LBool: return bool(v) case *lua.LTable: return pm.luaTableToMap(v) default: return v.String() } } // luaTableToMap конвертирует Lua-таблицу в Go-карту. func (pm *PluginManager) luaTableToMap(table *lua.LTable) map[string]interface{} { result := make(map[string]interface{}) table.ForEach(func(key, value lua.LValue) { keyStr := "unknown" if key.Type() == lua.LTString { keyStr = key.String() } else if key.Type() == lua.LTNumber { keyStr = fmt.Sprintf("%d", int64(key.(lua.LNumber))) } result[keyStr] = pm.luaValueToGo(value) }) return result } // goValueToLua конвертирует Go-значение в Lua-значение. func (pm *PluginManager) goValueToLua(L *lua.LState, val interface{}) lua.LValue { if val == nil { return lua.LNil } switch v := val.(type) { case string: return lua.LString(v) case int: return lua.LNumber(float64(v)) case int64: return lua.LNumber(float64(v)) case float32: return lua.LNumber(float64(v)) case float64: return lua.LNumber(v) case bool: return lua.LBool(v) case map[string]interface{}: table := L.NewTable() for k, val := range v { table.RawSetString(k, pm.goValueToLua(L, val)) } return table case []interface{}: table := L.NewTable() for i, val := range v { table.RawSetInt(i+1, pm.goValueToLua(L, val)) } return table default: return lua.LString(fmt.Sprintf("%v", v)) } } // getPluginMetadata извлекает метаданные из загруженного Lua-скрипта. func (pm *PluginManager) getPluginMetadata(L *lua.LState, field string) string { val := L.GetGlobal(field) if str, ok := val.(lua.LString); ok { return string(str) } return "" } // callPluginFunction вызывает функцию плагина по имени с защитой от паники. func (pm *PluginManager) callPluginFunction(plugin interface{}, funcName string) error { L, p := getPluginLState(plugin) if L == nil || p == nil || p.IsClosed() { return fmt.Errorf("plugin has no Lua state or is closed") } fn := L.GetGlobal(funcName) if fn == lua.LNil { return nil // Функция не определена - не ошибка } var callErr error func() { defer func() { if r := recover(); r != nil { callErr = fmt.Errorf("panic during %s: %v", funcName, r) } }() callErr = L.CallByParam(lua.P{ Fn: fn, NRet: 0, Protect: true, }) }() if callErr != nil { return fmt.Errorf("failed to call %s: %v", funcName, callErr) } return nil } // addPluginEventLog добавляет событие в лог плагина. // // ИЗМЕНЕНО (2026-10-03): формат nowStr теперь "02-01-2006 15:04:05.000". func (pm *PluginManager) addPluginEventLog(plugin interface{}, eventType string, data interface{}, errMsg string) { var p *Plugin switch v := plugin.(type) { case *Plugin: p = v case *PluginWithDeps: p = v.Plugin default: return } p.eventLogMu.Lock() defer p.eventLogMu.Unlock() nowMs := time.Now().UnixMilli() // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". nowStr := time.Now().Format("02-01-2006 15:04:05.000") eventLog := PluginEventLog{ PluginName: p.Name, EventType: eventType, Data: data, Timestamp: nowMs, TimestampStr: nowStr, Error: errMsg, } p.eventLog = append(p.eventLog, eventLog) if len(p.eventLog) > p.maxEventLog { p.eventLog = p.eventLog[1:] // Ограничиваем размер лога } } // eventLoop обрабатывает события от плагинов. // Завершается при закрытии eventBus (флаг eventBusClosed и close канала). func (pm *PluginManager) eventLoop() { for event := range pm.eventBus { if pm.logger != nil { pm.logger.Debug(fmt.Sprintf("Plugin event [%s] at %s: %+v", event.EventType, event.TimestampStr, event.Data)) } pm.plugins.Range(func(key, value interface{}) bool { go pm.notifyPlugin(value, event) return true }) } } // notifyPlugin уведомляет конкретный плагин о событии. func (pm *PluginManager) notifyPlugin(plugin interface{}, event PluginEvent) { L, p := getPluginLState(plugin) if L == nil || p == nil || p.IsClosed() { return } fn := L.GetGlobal("on_event") if fn == lua.LNil { return } event.PluginName = p.Name eventTable := L.NewTable() eventTable.RawSetString("type", lua.LString(event.EventType)) eventTable.RawSetString("plugin_name", lua.LString(event.PluginName)) eventTable.RawSetString("timestamp", lua.LNumber(event.Timestamp)) eventTable.RawSetString("timestamp_str", lua.LString(event.TimestampStr)) eventTable.RawSetString("data", pm.goValueToLua(L, event.Data)) var callErr error func() { defer func() { if r := recover(); r != nil { callErr = fmt.Errorf("panic during on_event: %v", r) } }() callErr = L.CallByParam(lua.P{ Fn: fn, NRet: 0, Protect: true, }, eventTable) }() if callErr != nil && pm.logger != nil { pm.logger.Error(fmt.Sprintf("Plugin %s on_event error: %v", p.Name, callErr)) } } // ExecutePlugin выполняет пользовательскую функцию плагина. func (pm *PluginManager) ExecutePlugin(pluginName, funcName string, args ...interface{}) (interface{}, error) { if !pm.enabled { return nil, fmt.Errorf("plugin system is disabled") } val, ok := pm.plugins.Load(pluginName) if !ok { return nil, fmt.Errorf("plugin not found: %s", pluginName) } L, p := getPluginLState(val) if L == nil || p == nil { return nil, fmt.Errorf("unknown plugin type") } if PluginStatus(p.Status.Load()) != StatusRunning { return nil, fmt.Errorf("plugin %s is not running", pluginName) } if p.IsClosed() { return nil, fmt.Errorf("plugin %s is closed", pluginName) } fn := L.GetGlobal(funcName) if fn == lua.LNil { return nil, fmt.Errorf("function %s not found in plugin %s", funcName, pluginName) } luaArgs := make([]lua.LValue, len(args)) for i, arg := range args { luaArgs[i] = pm.goValueToLua(L, arg) } startTime := time.Now() var ret lua.LValue var callErr error func() { defer func() { if r := recover(); r != nil { callErr = fmt.Errorf("panic during execution: %v", r) } }() callErr = L.CallByParam(lua.P{ Fn: fn, NRet: 1, Protect: true, }, luaArgs...) if callErr == nil { ret = L.Get(-1) L.Pop(1) } }() duration := time.Since(startTime) if callErr != nil { pm.addPluginEventLog(val, funcName, args, callErr.Error()) return nil, fmt.Errorf("plugin execution failed: %v", callErr) } pm.addPluginEventLog(val, funcName, args, "") if pm.logger != nil { pm.logger.Debug(fmt.Sprintf("Plugin %s executed %s in %dms", pluginName, funcName, duration.Milliseconds())) } return pm.luaValueToGo(ret), nil } // UnloadPlugin выгружает плагин и возвращает состояние в пул. // // Используем поле plugin.instrumented вместо type assertion // interface{}(L).(*InstrumentedLState), который всегда возвращал false, // что приводило к утечке Lua-состояний и некорректному уменьшению active. // // ИЗМЕНЕНО (2026-10-03): формат unloaded_at_str — "02-01-2006 15:04:05.000". func (pm *PluginManager) UnloadPlugin(name string) error { if !pm.enabled { return fmt.Errorf("plugin system is disabled") } val, ok := pm.plugins.Load(name) if !ok { return fmt.Errorf("plugin not found: %s", name) } // Сначала отменяем все активные batch'и плагина, // чтобы не оставить их зарегистрированными в глобальном менеджере. var p *Plugin switch v := val.(type) { case *Plugin: p = v case *PluginWithDeps: p = v.Plugin } if p != nil { p.activeBatches.Range(func(key, value interface{}) bool { if b, ok := value.(*storage.Batch); ok && b != nil { if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(b) } } if id, ok := key.(uint64); ok { p.UnregisterBatch(id) } return true }) } if err := pm.callPluginFunction(val, "on_unload"); err != nil { if pm.logger != nil { pm.logger.Warn(fmt.Sprintf("Plugin %s on_unload error: %v", name, err)) } pm.addPluginEventLog(val, "on_unload_error", nil, err.Error()) } else { pm.addPluginEventLog(val, "on_unload", nil, "") } // Извлекаем InstrumentedLState из плагина inst := getPluginInstrumented(val) if p != nil { p.Status.Store(int32(StatusStopped)) p.closePlugin() } // Возвращаем состояние в пул (или закрываем, если пул закрыт) if inst != nil { pm.statePool.Release(inst) } auditTimeMs, auditTimeStr := storage.GetCurrentTimestamp() storage.LogAudit("PLUGIN_UNLOAD", "PLUGIN", name, map[string]interface{}{ "unloaded_at": time.Now().UnixMilli(), // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". "unloaded_at_str": time.Now().Format("02-01-2006 15:04:05.000"), "audit_time": auditTimeMs, "audit_time_str": auditTimeStr, "active_lua_states": pm.statePool.GetActiveCount(), }) pm.plugins.Delete(name) if pm.logger != nil { pm.logger.Info(fmt.Sprintf("Plugin unloaded: %s", name)) } return nil } // StartPlugin запускает плагин. // // ИЗМЕНЕНО (2026-10-03): формат started_at_str — "02-01-2006 15:04:05.000". func (pm *PluginManager) StartPlugin(name string) error { if !pm.enabled { return fmt.Errorf("plugin system is disabled") } val, ok := pm.plugins.Load(name) if !ok { return fmt.Errorf("plugin not found: %s", name) } var p *Plugin switch v := val.(type) { case *Plugin: p = v case *PluginWithDeps: p = v.Plugin default: return fmt.Errorf("unknown plugin type") } if p.IsClosed() { return fmt.Errorf("plugin is closed") } p.Status.Store(int32(StatusRunning)) if err := pm.callPluginFunction(val, "on_start"); err != nil { p.Status.Store(int32(StatusError)) pm.addPluginEventLog(val, "on_start_error", nil, err.Error()) return fmt.Errorf("failed to start plugin: %v", err) } pm.addPluginEventLog(val, "on_start", nil, "") auditTimeMs, auditTimeStr := storage.GetCurrentTimestamp() storage.LogAudit("PLUGIN_START", "PLUGIN", name, map[string]interface{}{ "started_at": time.Now().UnixMilli(), // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". "started_at_str": time.Now().Format("02-01-2006 15:04:05.000"), "audit_time": auditTimeMs, "audit_time_str": auditTimeStr, "active_lua_states": pm.statePool.GetActiveCount(), }) if pm.logger != nil { pm.logger.Info(fmt.Sprintf("Plugin started: %s", name)) } return nil } // StopPlugin останавливает плагин. // // ИЗМЕНЕНО (2026-10-03): формат stopped_at_str — "02-01-2006 15:04:05.000". func (pm *PluginManager) StopPlugin(name string) error { if !pm.enabled { return fmt.Errorf("plugin system is disabled") } val, ok := pm.plugins.Load(name) if !ok { return fmt.Errorf("plugin not found: %s", name) } if err := pm.callPluginFunction(val, "on_stop"); err != nil { if pm.logger != nil { pm.logger.Warn(fmt.Sprintf("Plugin %s on_stop error: %v", name, err)) } pm.addPluginEventLog(val, "on_stop_error", nil, err.Error()) } else { pm.addPluginEventLog(val, "on_stop", nil, "") } var p *Plugin switch v := val.(type) { case *Plugin: p = v case *PluginWithDeps: p = v.Plugin } if p != nil { p.Status.Store(int32(StatusStopped)) } auditTimeMs, auditTimeStr := storage.GetCurrentTimestamp() storage.LogAudit("PLUGIN_STOP", "PLUGIN", name, map[string]interface{}{ "stopped_at": time.Now().UnixMilli(), // ИЗМЕНЕНО (2026-10-03): формат "02-01-2006 15:04:05.000". "stopped_at_str": time.Now().Format("02-01-2006 15:04:05.000"), "audit_time": auditTimeMs, "audit_time_str": auditTimeStr, "active_lua_states": pm.statePool.GetActiveCount(), }) if pm.logger != nil { pm.logger.Info(fmt.Sprintf("Plugin stopped: %s", name)) } return nil } // ListPlugins возвращает список всех загруженных плагинов. func (pm *PluginManager) ListPlugins() []*Plugin { plugins := make([]*Plugin, 0) pm.plugins.Range(func(key, value interface{}) bool { switch v := value.(type) { case *Plugin: plugins = append(plugins, v) case *PluginWithDeps: plugins = append(plugins, v.Plugin) } return true }) return plugins } // GetPlugin возвращает плагин по имени. func (pm *PluginManager) GetPlugin(name string) (*Plugin, error) { val, ok := pm.plugins.Load(name) if !ok { return nil, fmt.Errorf("plugin not found: %s", name) } switch v := val.(type) { case *Plugin: return v, nil case *PluginWithDeps: return v.Plugin, nil default: return nil, fmt.Errorf("unknown plugin type") } } // GetPluginWithDeps возвращает плагин с зависимостями. func (pm *PluginManager) GetPluginWithDeps(name string) (*PluginWithDeps, error) { val, ok := pm.plugins.Load(name) if !ok { return nil, fmt.Errorf("plugin not found: %s", name) } if withDeps, ok := val.(*PluginWithDeps); ok { return withDeps, nil } return nil, fmt.Errorf("plugin %s does not support dependencies", name) } // LoadPluginsInOrder загружает плагины в правильном порядке зависимостей. func (pm *PluginManager) LoadPluginsInOrder() error { order, err := pm.depResolver.ResolveOrder() if err != nil { return err } for _, plugin := range order { if err := pm.StartPlugin(plugin.Name); err != nil { return fmt.Errorf("failed to start plugin %s: %v", plugin.Name, err) } } return nil } // GetHotReloadManager возвращает менеджер горячей перезагрузки. func (pm *PluginManager) GetHotReloadManager() *HotReloadManager { return pm.hotReload } // GetMarketplace возвращает маркетплейс плагинов. func (pm *PluginManager) GetMarketplace() *PluginMarketplace { return pm.marketplace } // EnablePluginMarketplace включает маркетплейс плагинов. func (pm *PluginManager) EnablePluginMarketplace(apiEndpoint string) { pm.marketplace = NewPluginMarketplace(apiEndpoint, pm, pm.logger) } // IsEnabled возвращает статус системы плагинов. func (pm *PluginManager) IsEnabled() bool { return pm.enabled } // GetPluginsDir возвращает директорию с плагинами. func (pm *PluginManager) GetPluginsDir() string { return pm.pluginsDir } // GetPluginEventLog возвращает лог событий плагина. func (p *Plugin) GetPluginEventLog() []PluginEventLog { p.eventLogMu.RLock() defer p.eventLogMu.RUnlock() result := make([]PluginEventLog, len(p.eventLog)) copy(result, p.eventLog) return result } // ClearPluginEventLog очищает лог событий плагина. func (p *Plugin) ClearPluginEventLog() { p.eventLogMu.Lock() defer p.eventLogMu.Unlock() p.eventLog = make([]PluginEventLog, 0) } // Version возвращает версию плагина. func (p *Plugin) Version() string { return p.version } // Author возвращает автора плагина. func (p *Plugin) Author() string { return p.author } // Description возвращает описание плагина. func (p *Plugin) Description() string { return p.description } // LoadedAt возвращает время загрузки плагина. func (p *Plugin) LoadedAt() time.Time { return p.loadedAt } // LoadedAtMs возвращает время загрузки плагина в миллисекундах. func (p *Plugin) LoadedAtMs() int64 { return p.loadedAtMs } // LoadedAtStr возвращает строковое представление времени загрузки плагина. func (p *Plugin) LoadedAtStr() string { return p.loadedAtStr } // GetStatus возвращает статус плагина. func (p *Plugin) GetStatus() PluginStatus { return PluginStatus(p.Status.Load()) } // GetLuaStatePoolStats возвращает статистику пула Lua-состояний. func (pm *PluginManager) GetLuaStatePoolStats() map[string]interface{} { return map[string]interface{}{ "active_count": pm.statePool.GetActiveCount(), "pool_size": pm.statePool.GetPoolSize(), "max_allowed": MaxLuaStates, "ttl_seconds": LuaStateTTL.Seconds(), } } // Close закрывает менеджер плагинов и все ресурсы. // // eventBus закрывается через флаг eventBusClosed, чтобы // избежать паники "send on closed channel" при отправке событий // из плагинов после закрытия. func (pm *PluginManager) Close() error { if pm.hotReload != nil { pm.hotReload.Stop() } if err := pm.engineRegistry.CloseAll(); err != nil { if pm.logger != nil { pm.logger.Error(fmt.Sprintf("Error closing engines: %v", err)) } } // Выгружаем все плагины pm.plugins.Range(func(key, value interface{}) bool { name := key.(string) _ = pm.UnloadPlugin(name) return true }) // Помечаем eventBus как закрытый и закрываем канал. // eventLoop завершится по range. if pm.eventBusClosed.CAS(false, true) { close(pm.eventBus) } pm.statePool.Close() return nil }