2026-10-03 20:33:34 +00:00
|
|
|
|
/*
|
|
|
|
|
|
* 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/migration/schema_migrator.go
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Назначение: Миграция схемы данных при обновлении версии СУБД.
|
|
|
|
|
|
// Новая функциональность: прозрачное обновление схемы данных
|
|
|
|
|
|
// (автоматическая миграция документов).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
//
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// ИСПРАВЛЕНО (2026-10, переход ACID → Batch):
|
|
|
|
|
|
// После замены ACID-транзакций на batch-операции тип storage.Transaction
|
|
|
|
|
|
// был удалён. Все ссылки на *storage.Transaction заменены на
|
|
|
|
|
|
// *storage.Batch. Миграции схемы работают в основном с in-memory
|
|
|
|
|
|
// структурой SchemaDefinition, поэтому batch используется только как
|
|
|
|
|
|
// маркер атомарности: если пользовательская миграция пишет в СУБД,
|
|
|
|
|
|
// она добавляет операции в batch, который коммитится в конце
|
|
|
|
|
|
// MigrateSchema атомарно через batch.Commit().
|
|
|
|
|
|
//
|
|
|
|
|
|
// ИСПРАВЛЕНО (2026-10, аудит):
|
|
|
|
|
|
// 1. GetStatus, GetSchemaStatus, MigrateSchema и MigrateDocument больше
|
|
|
|
|
|
// не паникуют при sm.schema == nil. При отсутствии схемы
|
|
|
|
|
|
// используется пустая схема версии "1.0.0".
|
|
|
|
|
|
// 2. MigrateSchema применяет изменения на диск ДО обновления
|
|
|
|
|
|
// in-memory состояния. Если запись падает — sm.schema, sm.applied
|
|
|
|
|
|
// и sm.currentVersion откатываются к прежним значениям.
|
|
|
|
|
|
// Раньше in-memory менялось первым, и расхождение с диском
|
|
|
|
|
|
// оставалось видимым до рестарта.
|
|
|
|
|
|
// 3. MigrateCollectionDocuments теперь передаёт version документа в
|
|
|
|
|
|
// Batch.AddUpdate, чтобы coll.Update увеличивал версию. Раньше
|
|
|
|
|
|
// миграция обновляла поля, но version документа не менялся.
|
|
|
|
|
|
// 4. validateFieldValue для "regex" использует regexp.MatchString,
|
|
|
|
|
|
// а не strings.Contains (имя не соответствовало поведению).
|
|
|
|
|
|
// 5. saveAppliedMigrationsLocked сортирует записи по AppliedAt и ID —
|
|
|
|
|
|
// порядок в JSON стабильный (раньше был случайный, из map).
|
|
|
|
|
|
// 6. MigrateCollectionDocuments делает ранний выход при пустой
|
|
|
|
|
|
// коллекции — не создаёт и не регистрирует лишний batch.
|
|
|
|
|
|
// 7. Неиспользуемые параметры batch в Up/Down миграциях помечены "_".
|
2026-10-03 20:33:34 +00:00
|
|
|
|
|
|
|
|
|
|
package migration
|
|
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
|
"encoding/json"
|
|
|
|
|
|
"fmt"
|
|
|
|
|
|
"os"
|
|
|
|
|
|
"path/filepath"
|
2026-10-04 21:04:40 +00:00
|
|
|
|
"regexp"
|
2026-10-03 20:33:34 +00:00
|
|
|
|
"sort"
|
|
|
|
|
|
"strings"
|
|
|
|
|
|
"sync"
|
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
|
|
"futriis/internal/log"
|
|
|
|
|
|
"futriis/internal/storage"
|
|
|
|
|
|
)
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// СХЕМА
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
|
|
|
|
|
// SchemaDefinition определяет схему коллекции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type SchemaDefinition struct {
|
|
|
|
|
|
Version string `json:"version"`
|
|
|
|
|
|
Collections map[string]*CollectionSchema `json:"collections"`
|
|
|
|
|
|
UpdatedAt int64 `json:"updated_at"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// CollectionSchema определяет схему одной коллекции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type CollectionSchema struct {
|
|
|
|
|
|
Fields map[string]*FieldDefinition `json:"fields"`
|
|
|
|
|
|
Required []string `json:"required"`
|
|
|
|
|
|
Indexes []IndexDefinition `json:"indexes"`
|
|
|
|
|
|
UpdatedAt int64 `json:"updated_at"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// FieldDefinition определяет поле в схеме.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type FieldDefinition struct {
|
|
|
|
|
|
Type string `json:"type"` // string, int, float, bool, array, object
|
|
|
|
|
|
Required bool `json:"required"`
|
|
|
|
|
|
Default interface{} `json:"default"`
|
|
|
|
|
|
Validate string `json:"validate"` // regex, min, max, enum
|
|
|
|
|
|
ValidateValue interface{} `json:"validate_value"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// IndexDefinition определяет индекс.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type IndexDefinition struct {
|
|
|
|
|
|
Name string `json:"name"`
|
|
|
|
|
|
Fields []string `json:"fields"`
|
|
|
|
|
|
Unique bool `json:"unique"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// МИГРАЦИИ
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
// SchemaMigration представляет миграцию схемы.
|
|
|
|
|
|
//
|
|
|
|
|
|
// ИЗМЕНЕНО (2026-10): функции Up/Down теперь принимают *storage.Batch
|
|
|
|
|
|
// вместо *storage.Transaction. Batch — это группа операций, применяемых
|
|
|
|
|
|
// атомарно через WAL. Если миграции нужно изменить данные в СУБД, она
|
|
|
|
|
|
// добавляет операции в batch. Если миграция меняет только in-memory
|
|
|
|
|
|
// SchemaDefinition — batch не используется, но остаётся доступным для
|
|
|
|
|
|
// совместимости и возможного расширения.
|
|
|
|
|
|
type SchemaMigration struct {
|
|
|
|
|
|
ID string
|
|
|
|
|
|
Version string
|
|
|
|
|
|
Description string
|
|
|
|
|
|
Up func(batch *storage.Batch, schema *SchemaDefinition) error
|
|
|
|
|
|
Down func(batch *storage.Batch, schema *SchemaDefinition) error
|
|
|
|
|
|
CreatedAt int64
|
|
|
|
|
|
AppliedAt int64
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// SchemaMigrationRecord — запись о применённой миграции схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type SchemaMigrationRecord struct {
|
|
|
|
|
|
ID string `json:"id"`
|
|
|
|
|
|
Version string `json:"version"`
|
|
|
|
|
|
Description string `json:"description"`
|
|
|
|
|
|
AppliedAt int64 `json:"applied_at"`
|
|
|
|
|
|
Success bool `json:"success"`
|
|
|
|
|
|
Error string `json:"error,omitempty"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// DocumentMigrationStrategy — стратегия миграции документа.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type DocumentMigrationStrategy int
|
|
|
|
|
|
|
|
|
|
|
|
const (
|
|
|
|
|
|
StrategyInPlace DocumentMigrationStrategy = iota // Обновление на месте
|
|
|
|
|
|
StrategyCopyNew // Копирование в новую коллекцию
|
|
|
|
|
|
StrategyLazy // Ленивая миграция при доступе
|
|
|
|
|
|
)
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// SCHEMA MIGRATOR
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
|
|
|
|
|
// SchemaMigrator управляет миграциями схемы данных.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type SchemaMigrator struct {
|
|
|
|
|
|
store *storage.Storage
|
|
|
|
|
|
logger *log.Logger
|
|
|
|
|
|
migrations map[string]*SchemaMigration
|
|
|
|
|
|
applied map[string]*SchemaMigrationRecord
|
|
|
|
|
|
mu sync.RWMutex
|
|
|
|
|
|
migrationDir string
|
|
|
|
|
|
currentVersion string
|
|
|
|
|
|
targetVersion string
|
|
|
|
|
|
schema *SchemaDefinition
|
|
|
|
|
|
strategy DocumentMigrationStrategy
|
|
|
|
|
|
migrationStats *MigrationStats
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// MigrationStats — статистика миграции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type MigrationStats struct {
|
|
|
|
|
|
TotalDocuments int64
|
|
|
|
|
|
MigratedDocuments int64
|
|
|
|
|
|
FailedDocuments int64
|
|
|
|
|
|
SkippedDocuments int64
|
|
|
|
|
|
StartTime int64
|
|
|
|
|
|
EndTime int64
|
|
|
|
|
|
mu sync.RWMutex
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// NewSchemaMigrator создаёт новый мигратор схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func NewSchemaMigrator(store *storage.Storage, logger *log.Logger, migrationDir string) *SchemaMigrator {
|
|
|
|
|
|
sm := &SchemaMigrator{
|
2026-10-04 21:04:40 +00:00
|
|
|
|
store: store,
|
|
|
|
|
|
logger: logger,
|
|
|
|
|
|
migrations: make(map[string]*SchemaMigration),
|
|
|
|
|
|
applied: make(map[string]*SchemaMigrationRecord),
|
|
|
|
|
|
migrationDir: migrationDir,
|
|
|
|
|
|
strategy: StrategyLazy,
|
2026-10-03 20:33:34 +00:00
|
|
|
|
migrationStats: &MigrationStats{},
|
|
|
|
|
|
schema: &SchemaDefinition{
|
|
|
|
|
|
Version: "1.0.0",
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
},
|
2026-10-04 21:04:40 +00:00
|
|
|
|
currentVersion: "1.0.0",
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Создаём директорию для миграций (без паники при ошибке).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if err := os.MkdirAll(migrationDir, 0755); err != nil {
|
|
|
|
|
|
if logger != nil {
|
|
|
|
|
|
logger.Error(fmt.Sprintf("Failed to create migration directory %s: %v", migrationDir, err))
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Загружаем существующую схему.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.loadSchema()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Загружаем применённые миграции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.loadAppliedMigrations()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Регистрируем встроенные миграции схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.registerBuiltinSchemaMigrations()
|
|
|
|
|
|
|
|
|
|
|
|
return sm
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// ЗАГРУЗКА / СОХРАНЕНИЕ
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
|
|
|
|
|
// loadSchema загружает определение схемы из файла.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) loadSchema() {
|
|
|
|
|
|
path := filepath.Join(sm.migrationDir, "schema.json")
|
|
|
|
|
|
|
|
|
|
|
|
data, err := os.ReadFile(path)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Debug("No existing schema found, creating default schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.createDefaultSchema()
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
var schema SchemaDefinition
|
|
|
|
|
|
if err := json.Unmarshal(data, &schema); err != nil {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Error(fmt.Sprintf("Failed to unmarshal schema: %v", err))
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.createDefaultSchema()
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if schema.Collections == nil {
|
|
|
|
|
|
schema.Collections = make(map[string]*CollectionSchema)
|
|
|
|
|
|
}
|
|
|
|
|
|
if schema.Version == "" {
|
|
|
|
|
|
schema.Version = "1.0.0"
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
sm.schema = &schema
|
|
|
|
|
|
sm.currentVersion = schema.Version
|
|
|
|
|
|
sm.mu.Unlock()
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
sm.logger.Info(fmt.Sprintf("Loaded schema version %s with %d collections",
|
|
|
|
|
|
schema.Version, len(schema.Collections)))
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// saveSchemaLocked сохраняет определение схемы атомарно (временный файл + rename).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
//
|
|
|
|
|
|
// ВАЖНО: метод НЕ берёт sm.mu, чтобы его можно было вызывать из функций,
|
|
|
|
|
|
// уже удерживающих sm.mu.Lock(). Вызывающий сам отвечает за консистентность
|
|
|
|
|
|
// читаемого состояния sm.schema.
|
|
|
|
|
|
func (sm *SchemaMigrator) saveSchemaLocked() error {
|
|
|
|
|
|
data, err := json.MarshalIndent(sm.schema, "", " ")
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return err
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
path := filepath.Join(sm.migrationDir, "schema.json")
|
|
|
|
|
|
tmpPath := path + ".tmp"
|
|
|
|
|
|
|
|
|
|
|
|
if err := os.WriteFile(tmpPath, data, 0644); err != nil {
|
|
|
|
|
|
return err
|
|
|
|
|
|
}
|
|
|
|
|
|
return os.Rename(tmpPath, path)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// createDefaultSchema создаёт схему по умолчанию.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) createDefaultSchema() {
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
|
|
|
|
|
|
sm.schema = &SchemaDefinition{
|
|
|
|
|
|
Version: "1.0.0",
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.schema.Collections["_default"] = &CollectionSchema{
|
|
|
|
|
|
Fields: map[string]*FieldDefinition{
|
|
|
|
|
|
"_id": {
|
|
|
|
|
|
Type: "string",
|
|
|
|
|
|
Required: true,
|
|
|
|
|
|
},
|
|
|
|
|
|
"created_at": {
|
|
|
|
|
|
Type: "int64",
|
|
|
|
|
|
Default: int64(0),
|
|
|
|
|
|
},
|
|
|
|
|
|
"updated_at": {
|
|
|
|
|
|
Type: "int64",
|
|
|
|
|
|
Default: int64(0),
|
|
|
|
|
|
},
|
|
|
|
|
|
},
|
|
|
|
|
|
Required: []string{"_id"},
|
|
|
|
|
|
Indexes: []IndexDefinition{},
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.currentVersion = "1.0.0"
|
|
|
|
|
|
_ = sm.saveSchemaLocked()
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// registerBuiltinSchemaMigrations регистрирует встроенные миграции схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) registerBuiltinSchemaMigrations() {
|
|
|
|
|
|
sm.RegisterSchemaMigration(&SchemaMigration{
|
|
|
|
|
|
ID: "schema_001_add_timestamps",
|
|
|
|
|
|
Version: "1.1.0",
|
|
|
|
|
|
Description: "Add created_at and updated_at timestamps to all collections",
|
|
|
|
|
|
Up: sm.migrateSchemaAddTimestamps,
|
|
|
|
|
|
Down: sm.migrateSchemaRemoveTimestamps,
|
|
|
|
|
|
CreatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
|
|
sm.RegisterSchemaMigration(&SchemaMigration{
|
|
|
|
|
|
ID: "schema_002_add_soft_delete",
|
|
|
|
|
|
Version: "1.2.0",
|
|
|
|
|
|
Description: "Add soft delete support (deleted_at, deleted fields)",
|
|
|
|
|
|
Up: sm.migrateSchemaAddSoftDelete,
|
|
|
|
|
|
Down: sm.migrateSchemaRemoveSoftDelete,
|
|
|
|
|
|
CreatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
|
|
sm.RegisterSchemaMigration(&SchemaMigration{
|
|
|
|
|
|
ID: "schema_003_add_versioning",
|
|
|
|
|
|
Version: "1.3.0",
|
|
|
|
|
|
Description: "Add document versioning (_version field)",
|
|
|
|
|
|
Up: sm.migrateSchemaAddVersioning,
|
|
|
|
|
|
Down: sm.migrateSchemaRemoveVersioning,
|
|
|
|
|
|
CreatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
|
|
sm.RegisterSchemaMigration(&SchemaMigration{
|
|
|
|
|
|
ID: "schema_004_add_indexes",
|
|
|
|
|
|
Version: "2.0.0",
|
|
|
|
|
|
Description: "Add secondary indexes for performance",
|
|
|
|
|
|
Up: sm.migrateSchemaAddIndexes,
|
|
|
|
|
|
Down: sm.migrateSchemaRemoveIndexes,
|
|
|
|
|
|
CreatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
|
|
sm.RegisterSchemaMigration(&SchemaMigration{
|
|
|
|
|
|
ID: "schema_005_add_validation",
|
|
|
|
|
|
Version: "2.1.0",
|
|
|
|
|
|
Description: "Add field validation rules",
|
|
|
|
|
|
Up: sm.migrateSchemaAddValidation,
|
|
|
|
|
|
Down: sm.migrateSchemaRemoveValidation,
|
|
|
|
|
|
CreatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
})
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// RegisterSchemaMigration регистрирует миграцию схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) RegisterSchemaMigration(m *SchemaMigration) {
|
|
|
|
|
|
if m == nil || m.ID == "" {
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
sm.migrations[m.ID] = m
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// loadAppliedMigrations загружает применённые миграции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) loadAppliedMigrations() {
|
|
|
|
|
|
path := filepath.Join(sm.migrationDir, "schema_migrations.json")
|
|
|
|
|
|
|
|
|
|
|
|
data, err := os.ReadFile(path)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
var records []SchemaMigrationRecord
|
|
|
|
|
|
if err := json.Unmarshal(data, &records); err != nil {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Error(fmt.Sprintf("Failed to unmarshal schema migrations: %v", err))
|
|
|
|
|
|
}
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
|
|
|
|
|
|
for i := range records {
|
|
|
|
|
|
sm.applied[records[i].ID] = &records[i]
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// saveAppliedMigrationsLocked сохраняет применённые миграции.
|
|
|
|
|
|
// Не берёт sm.mu (вызывающий уже держит Lock).
|
2026-10-04 21:04:40 +00:00
|
|
|
|
//
|
|
|
|
|
|
// ИСПРАВЛЕНО: записи сортируются по AppliedAt, затем по ID —
|
|
|
|
|
|
// порядок в JSON стабильный и человекочитаемый.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) saveAppliedMigrationsLocked() error {
|
|
|
|
|
|
records := make([]SchemaMigrationRecord, 0, len(sm.applied))
|
|
|
|
|
|
for _, record := range sm.applied {
|
|
|
|
|
|
records = append(records, *record)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
sort.Slice(records, func(i, j int) bool {
|
|
|
|
|
|
if records[i].AppliedAt != records[j].AppliedAt {
|
|
|
|
|
|
return records[i].AppliedAt < records[j].AppliedAt
|
|
|
|
|
|
}
|
|
|
|
|
|
return records[i].ID < records[j].ID
|
|
|
|
|
|
})
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
data, err := json.MarshalIndent(records, "", " ")
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return err
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
path := filepath.Join(sm.migrationDir, "schema_migrations.json")
|
|
|
|
|
|
tmpPath := path + ".tmp"
|
|
|
|
|
|
if err := os.WriteFile(tmpPath, data, 0644); err != nil {
|
|
|
|
|
|
return err
|
|
|
|
|
|
}
|
|
|
|
|
|
return os.Rename(tmpPath, path)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// МИГРАЦИЯ СХЕМЫ
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
// MigrateSchema выполняет миграцию схемы до указанной версии.
|
|
|
|
|
|
//
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// ИЗМЕНЕНО (2026-10): используется storage.NewBatch + batch.Commit.
|
|
|
|
|
|
// Batch создаётся для (db="_system", coll="_schema_migrations") —
|
|
|
|
|
|
// это виртуальная коллекция, куда миграции могут добавлять операции.
|
|
|
|
|
|
// Если миграция не пишет в СУБД, batch просто останется пустым и
|
|
|
|
|
|
// его Commit — no-op.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
//
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// ИСПРАВЛЕНО (аудит):
|
|
|
|
|
|
// - добавлена nil-проверка sm.schema (паники больше не будет);
|
|
|
|
|
|
// - изменения применяются на диск ДО обновления in-memory. Если
|
|
|
|
|
|
// запись падает — sm.schema, sm.applied и sm.currentVersion
|
|
|
|
|
|
// откатываются к прежним значениям.
|
|
|
|
|
|
//
|
|
|
|
|
|
// ВАЖНО: метод НЕ удерживает sm.mu на всё время выполнения, чтобы
|
|
|
|
|
|
// избежать deadlock с saveSchema/saveAppliedMigrations и не блокировать
|
|
|
|
|
|
// чтение схемы на время долгих миграций. Лок берётся только для
|
|
|
|
|
|
// доступа к картам migrations/applied и к sm.schema.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) {
|
|
|
|
|
|
if targetVersion == "" {
|
|
|
|
|
|
return fmt.Errorf("target version is empty")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Starting schema migration to version %s", targetVersion))
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Снимок текущего состояния под RLock.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.mu.RLock()
|
2026-10-04 21:04:40 +00:00
|
|
|
|
currentVersion := "1.0.0"
|
|
|
|
|
|
if sm.schema != nil {
|
|
|
|
|
|
currentVersion = sm.schema.Version
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
allMigrations := make([]*SchemaMigration, 0, len(sm.migrations))
|
|
|
|
|
|
for _, m := range sm.migrations {
|
|
|
|
|
|
allMigrations = append(allMigrations, m)
|
|
|
|
|
|
}
|
|
|
|
|
|
appliedCopy := make(map[string]bool, len(sm.applied))
|
|
|
|
|
|
for id := range sm.applied {
|
|
|
|
|
|
appliedCopy[id] = true
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.mu.RUnlock()
|
|
|
|
|
|
|
|
|
|
|
|
if currentVersion == targetVersion {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info("Already at target version, no migration needed")
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Сортируем по версии.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sort.Slice(allMigrations, func(i, j int) bool {
|
|
|
|
|
|
return sm.isVersionLess(allMigrations[i].Version, allMigrations[j].Version)
|
|
|
|
|
|
})
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Отбираем миграции для применения.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
var toApply []*SchemaMigration
|
|
|
|
|
|
for _, m := range allMigrations {
|
|
|
|
|
|
if sm.isVersionGreater(m.Version, currentVersion) && !sm.isVersionGreater(m.Version, targetVersion) {
|
|
|
|
|
|
if !appliedCopy[m.ID] {
|
|
|
|
|
|
toApply = append(toApply, m)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Applying %d schema migrations", len(toApply)))
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Если targetVersion меньше currentVersion — это откат, который не поддерживается.
|
|
|
|
|
|
if len(toApply) == 0 && sm.isVersionGreater(currentVersion, targetVersion) {
|
|
|
|
|
|
return fmt.Errorf("downgrade from %s to %s is not supported", currentVersion, targetVersion)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Копируем схему, чтобы изменения были атомарны: применяем к копии,
|
|
|
|
|
|
// а в sm.schema записываем только после успешного коммита.
|
|
|
|
|
|
sm.mu.RLock()
|
|
|
|
|
|
schemaCopy := sm.deepCopySchemaLocked()
|
|
|
|
|
|
sm.mu.RUnlock()
|
|
|
|
|
|
|
|
|
|
|
|
if schemaCopy == nil {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Если исходной схемы нет — создаём пустую с версией 1.0.0.
|
|
|
|
|
|
schemaCopy = &SchemaDefinition{
|
|
|
|
|
|
Version: "1.0.0",
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Создаём batch для миграции. Batch — это группа операций,
|
|
|
|
|
|
// применяемых атомарно. Если ни одна миграция не пишет в СУБД,
|
|
|
|
|
|
// batch останется пустым, а его Commit — no-op.
|
|
|
|
|
|
//
|
|
|
|
|
|
// Используем специальную виртуальную коллекцию "_schema_migrations",
|
|
|
|
|
|
// чтобы отделить batch миграций от пользовательских batch'ей.
|
|
|
|
|
|
batch := storage.NewBatch("_system", "_schema_migrations")
|
|
|
|
|
|
if batch == nil {
|
|
|
|
|
|
return fmt.Errorf("failed to create batch")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Обработка паники с именованным возвратом err.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
defer func() {
|
|
|
|
|
|
if r := recover(); r != nil {
|
|
|
|
|
|
if bm := storage.GetBatchManager(); bm != nil {
|
|
|
|
|
|
bm.UnregisterBatch(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Error(fmt.Sprintf("Panic during schema migration: %v", r))
|
|
|
|
|
|
}
|
|
|
|
|
|
err = fmt.Errorf("panic during schema migration: %v", r)
|
|
|
|
|
|
}
|
|
|
|
|
|
}()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Новые applied-записи, которые закоммитим в sm.applied только при успехе.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
newApplied := make(map[string]*SchemaMigrationRecord)
|
|
|
|
|
|
|
|
|
|
|
|
for _, m := range toApply {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Applying schema migration: %s (%s)", m.ID, m.Description))
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
startTime := time.Now()
|
|
|
|
|
|
|
|
|
|
|
|
if m.Up == nil {
|
|
|
|
|
|
if bm := storage.GetBatchManager(); bm != nil {
|
|
|
|
|
|
bm.UnregisterBatch(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
return fmt.Errorf("migration %s has nil Up function", m.ID)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if err := m.Up(batch, schemaCopy); err != nil {
|
|
|
|
|
|
if bm := storage.GetBatchManager(); bm != nil {
|
|
|
|
|
|
bm.UnregisterBatch(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
return fmt.Errorf("schema migration %s failed: %v", m.ID, err)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
schemaCopy.Version = m.Version
|
|
|
|
|
|
schemaCopy.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
|
|
|
|
|
|
newApplied[m.ID] = &SchemaMigrationRecord{
|
|
|
|
|
|
ID: m.ID,
|
|
|
|
|
|
Version: m.Version,
|
|
|
|
|
|
Description: m.Description,
|
|
|
|
|
|
AppliedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
Success: true,
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Schema migration %s completed in %v", m.ID, time.Since(startTime)))
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Коммитим batch ДО записи на диск.
|
|
|
|
|
|
// Batch.Commit снимает регистрацию в activeBatches автоматически.
|
|
|
|
|
|
// Если batch пуст — Commit вернёт nil без записи в WAL.
|
|
|
|
|
|
if err := batch.Commit(); err != nil {
|
|
|
|
|
|
return fmt.Errorf("failed to commit schema migration batch: %v", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Только теперь применяем изменения: сначала пишем на диск,
|
|
|
|
|
|
// потом обновляем in-memory. Если запись падает — откатываем всё.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.mu.Lock()
|
2026-10-04 21:04:40 +00:00
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
|
|
|
|
|
|
oldSchema := sm.schema
|
|
|
|
|
|
oldVersion := sm.currentVersion
|
|
|
|
|
|
oldApplied := make(map[string]*SchemaMigrationRecord, len(sm.applied))
|
|
|
|
|
|
for k, v := range sm.applied {
|
|
|
|
|
|
oldApplied[k] = v
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// 1. Временно подменяем схему и applied, чтобы save* работали с новым состоянием.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.schema = schemaCopy
|
2026-10-04 21:04:40 +00:00
|
|
|
|
sm.currentVersion = targetVersion
|
2026-10-03 20:33:34 +00:00
|
|
|
|
for id, rec := range newApplied {
|
|
|
|
|
|
sm.applied[id] = rec
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// 2. Пишем схему на диск.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if err := sm.saveSchemaLocked(); err != nil {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
sm.schema = oldSchema
|
|
|
|
|
|
sm.currentVersion = oldVersion
|
|
|
|
|
|
sm.applied = oldApplied
|
2026-10-03 20:33:34 +00:00
|
|
|
|
return fmt.Errorf("failed to save schema: %v", err)
|
|
|
|
|
|
}
|
2026-10-04 21:04:40 +00:00
|
|
|
|
|
|
|
|
|
|
// 3. Пишем applied на диск. При ошибке откатываем схему и на диске тоже.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if err := sm.saveAppliedMigrationsLocked(); err != nil {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
sm.schema = oldSchema
|
|
|
|
|
|
sm.currentVersion = oldVersion
|
|
|
|
|
|
sm.applied = oldApplied
|
|
|
|
|
|
// Пытаемся восстановить schema.json на диске.
|
|
|
|
|
|
if restoreErr := sm.saveSchemaLocked(); restoreErr != nil && sm.logger != nil {
|
|
|
|
|
|
sm.logger.Error(fmt.Sprintf("Failed to restore schema after applied-migration save error: %v", restoreErr))
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
return fmt.Errorf("failed to save applied migrations: %v", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Schema migration to version %s completed successfully", targetVersion))
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// deepCopySchemaLocked создаёт глубокую копию sm.schema через JSON.
|
|
|
|
|
|
// Вызывающий должен удерживать sm.mu (RLock).
|
|
|
|
|
|
func (sm *SchemaMigrator) deepCopySchemaLocked() *SchemaDefinition {
|
|
|
|
|
|
if sm.schema == nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
data, err := json.Marshal(sm.schema)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
var copy SchemaDefinition
|
|
|
|
|
|
if err := json.Unmarshal(data, ©); err != nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
if copy.Collections == nil {
|
|
|
|
|
|
copy.Collections = make(map[string]*CollectionSchema)
|
|
|
|
|
|
}
|
|
|
|
|
|
return ©
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// МИГРАЦИЯ ДОКУМЕНТОВ
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
// MigrateDocument прозрачно мигрирует отдельный документ к актуальной схеме.
|
|
|
|
|
|
// Возвращает копию документа и флаг changed (были ли изменения).
|
|
|
|
|
|
func (sm *SchemaMigrator) MigrateDocument(doc *storage.Document, collectionName string) (*storage.Document, bool, error) {
|
|
|
|
|
|
if doc == nil {
|
|
|
|
|
|
return nil, false, fmt.Errorf("nil document")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.mu.RLock()
|
2026-10-04 21:04:40 +00:00
|
|
|
|
var collSchema *CollectionSchema
|
|
|
|
|
|
if sm.schema != nil {
|
|
|
|
|
|
if s, exists := sm.schema.Collections[collectionName]; exists {
|
|
|
|
|
|
collSchema = s
|
|
|
|
|
|
} else if s, exists := sm.schema.Collections["_default"]; exists {
|
|
|
|
|
|
collSchema = s
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Копируем схему, чтобы не держать RLock во время работы с документом.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
schemaCopy := copyCollectionSchema(collSchema)
|
|
|
|
|
|
sm.mu.RUnlock()
|
|
|
|
|
|
|
|
|
|
|
|
if schemaCopy == nil {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Нет схемы — возвращаем документ без изменений.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
return doc.Clone(), false, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
migratedDoc := doc.Clone()
|
|
|
|
|
|
changed := false
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Добавляем недостающие поля с значениями по умолчанию.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
for fieldName, fieldDef := range schemaCopy.Fields {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if fieldDef == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if _, err := migratedDoc.GetField(fieldName); err != nil {
|
|
|
|
|
|
if fieldDef.Default != nil {
|
|
|
|
|
|
migratedDoc.SetField(fieldName, fieldDef.Default)
|
|
|
|
|
|
changed = true
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Debug(fmt.Sprintf("Added default value for field %s in document %s", fieldName, doc.ID))
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Проверяем типы полей и преобразуем при необходимости.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
for fieldName, fieldDef := range schemaCopy.Fields {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if fieldDef == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if value, err := migratedDoc.GetField(fieldName); err == nil {
|
|
|
|
|
|
converted, needsConversion := sm.convertFieldType(value, fieldDef.Type)
|
|
|
|
|
|
if needsConversion {
|
|
|
|
|
|
migratedDoc.SetField(fieldName, converted)
|
|
|
|
|
|
changed = true
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Debug(fmt.Sprintf("Converted field %s type for document %s", fieldName, doc.ID))
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if changed {
|
|
|
|
|
|
migratedDoc.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
migratedDoc.Version++
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
return migratedDoc, changed, nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// copyCollectionSchema создаёт глубокую копию CollectionSchema через JSON.
|
|
|
|
|
|
func copyCollectionSchema(src *CollectionSchema) *CollectionSchema {
|
|
|
|
|
|
if src == nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
data, err := json.Marshal(src)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
var copy CollectionSchema
|
|
|
|
|
|
if err := json.Unmarshal(data, ©); err != nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
if copy.Fields == nil {
|
|
|
|
|
|
copy.Fields = make(map[string]*FieldDefinition)
|
|
|
|
|
|
}
|
|
|
|
|
|
return ©
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// MigrateCollectionDocuments прозрачно мигрирует все документы коллекции.
|
|
|
|
|
|
//
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// ИСПРАВЛЕНО (аудит):
|
|
|
|
|
|
// - передаём version документа в Batch.AddUpdate, чтобы coll.Update
|
|
|
|
|
|
// увеличивал счётчик версии. Раньше миграция обновляла поля, но
|
|
|
|
|
|
// version документа не менялся.
|
|
|
|
|
|
// - ранний выход при пустой коллекции — не создаём лишний batch.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collection) error {
|
|
|
|
|
|
if collection == nil {
|
|
|
|
|
|
return fmt.Errorf("nil collection")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
startTime := time.Now()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
docs := collection.GetAllDocuments()
|
|
|
|
|
|
if len(docs) == 0 {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Collection %s is empty, no migration needed", collection.Name()))
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.migrationStats.mu.Lock()
|
|
|
|
|
|
sm.migrationStats.TotalDocuments = 0
|
|
|
|
|
|
sm.migrationStats.StartTime = startTime.UnixMilli()
|
|
|
|
|
|
sm.migrationStats.EndTime = startTime.UnixMilli()
|
|
|
|
|
|
sm.migrationStats.mu.Unlock()
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.migrationStats.mu.Lock()
|
|
|
|
|
|
sm.migrationStats.StartTime = startTime.UnixMilli()
|
2026-10-04 21:04:40 +00:00
|
|
|
|
sm.migrationStats.TotalDocuments = int64(len(docs))
|
2026-10-03 20:33:34 +00:00
|
|
|
|
sm.migrationStats.MigratedDocuments = 0
|
|
|
|
|
|
sm.migrationStats.FailedDocuments = 0
|
|
|
|
|
|
sm.migrationStats.SkippedDocuments = 0
|
|
|
|
|
|
sm.migrationStats.mu.Unlock()
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Starting migration of collection %s with %d documents",
|
2026-10-04 21:04:40 +00:00
|
|
|
|
collection.Name(), len(docs)))
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Создаём один batch для всей коллекции — атомарность на уровне коллекции.
|
|
|
|
|
|
batch := storage.NewBatch(collection.DBName(), collection.Name())
|
|
|
|
|
|
if batch == nil {
|
|
|
|
|
|
return fmt.Errorf("failed to create batch")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Снимаем регистрацию batch при выходе, если не закоммитили.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
committed := false
|
|
|
|
|
|
defer func() {
|
|
|
|
|
|
if !committed {
|
|
|
|
|
|
if bm := storage.GetBatchManager(); bm != nil {
|
|
|
|
|
|
bm.UnregisterBatch(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}()
|
|
|
|
|
|
|
|
|
|
|
|
migratedCount := int64(0)
|
|
|
|
|
|
failedCount := int64(0)
|
|
|
|
|
|
skippedCount := int64(0)
|
|
|
|
|
|
|
|
|
|
|
|
for _, doc := range docs {
|
|
|
|
|
|
migratedDoc, changed, err := sm.MigrateDocument(doc, collection.Name())
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
failedCount++
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Error(fmt.Sprintf("Failed to migrate document %s: %v", doc.ID, err))
|
|
|
|
|
|
}
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if !changed {
|
|
|
|
|
|
skippedCount++
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
updates := migratedDoc.GetFields()
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// ИСПРАВЛЕНО: явно переносим version в updates, чтобы coll.Update
|
|
|
|
|
|
// увеличил счётчик версии документа.
|
|
|
|
|
|
if migratedDoc.Version > 0 {
|
|
|
|
|
|
updates["version"] = migratedDoc.Version
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
batch.AddUpdate(doc.ID, updates)
|
|
|
|
|
|
migratedCount++
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Коммитим batch атомарно.
|
|
|
|
|
|
if migratedCount > 0 {
|
|
|
|
|
|
if err := batch.Commit(); err != nil {
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Error(fmt.Sprintf("Failed to commit migration batch for %s: %v",
|
|
|
|
|
|
collection.Name(), err))
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.migrationStats.mu.Lock()
|
|
|
|
|
|
sm.migrationStats.FailedDocuments = failedCount + migratedCount
|
|
|
|
|
|
sm.migrationStats.SkippedDocuments = skippedCount
|
|
|
|
|
|
sm.migrationStats.EndTime = time.Now().UnixMilli()
|
|
|
|
|
|
sm.migrationStats.mu.Unlock()
|
|
|
|
|
|
return fmt.Errorf("failed to commit migration batch: %v", err)
|
|
|
|
|
|
}
|
|
|
|
|
|
committed = true
|
|
|
|
|
|
} else {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Batch пуст — снимаем регистрацию вручную и отмечаем как "закоммичен".
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if bm := storage.GetBatchManager(); bm != nil {
|
|
|
|
|
|
bm.UnregisterBatch(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
committed = true
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.migrationStats.mu.Lock()
|
|
|
|
|
|
sm.migrationStats.MigratedDocuments = migratedCount
|
|
|
|
|
|
sm.migrationStats.FailedDocuments = failedCount
|
|
|
|
|
|
sm.migrationStats.SkippedDocuments = skippedCount
|
|
|
|
|
|
sm.migrationStats.EndTime = time.Now().UnixMilli()
|
|
|
|
|
|
total := sm.migrationStats.TotalDocuments
|
|
|
|
|
|
sm.migrationStats.mu.Unlock()
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Migration of collection %s completed in %v: %d migrated, %d failed, %d skipped (total %d)",
|
|
|
|
|
|
collection.Name(), time.Since(startTime), migratedCount, failedCount, skippedCount, total))
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// convertFieldType преобразует значение поля к нужному типу.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) convertFieldType(value interface{}, targetType string) (interface{}, bool) {
|
|
|
|
|
|
switch targetType {
|
|
|
|
|
|
case "string":
|
|
|
|
|
|
if _, ok := value.(string); !ok {
|
|
|
|
|
|
return fmt.Sprintf("%v", value), true
|
|
|
|
|
|
}
|
|
|
|
|
|
case "int64":
|
|
|
|
|
|
switch v := value.(type) {
|
|
|
|
|
|
case int:
|
|
|
|
|
|
return int64(v), true
|
|
|
|
|
|
case int32:
|
|
|
|
|
|
return int64(v), true
|
|
|
|
|
|
case int64:
|
|
|
|
|
|
// уже нужный тип
|
|
|
|
|
|
case float64:
|
|
|
|
|
|
return int64(v), true
|
|
|
|
|
|
case string:
|
|
|
|
|
|
var i int64
|
|
|
|
|
|
if _, err := fmt.Sscanf(v, "%d", &i); err == nil {
|
|
|
|
|
|
return i, true
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
case "float64":
|
|
|
|
|
|
switch v := value.(type) {
|
|
|
|
|
|
case int:
|
|
|
|
|
|
return float64(v), true
|
|
|
|
|
|
case int32:
|
|
|
|
|
|
return float64(v), true
|
|
|
|
|
|
case int64:
|
|
|
|
|
|
return float64(v), true
|
|
|
|
|
|
case float32:
|
|
|
|
|
|
return float64(v), true
|
|
|
|
|
|
case float64:
|
|
|
|
|
|
// уже нужный тип
|
|
|
|
|
|
case string:
|
|
|
|
|
|
var f float64
|
|
|
|
|
|
if _, err := fmt.Sscanf(v, "%f", &f); err == nil {
|
|
|
|
|
|
return f, true
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
case "bool":
|
|
|
|
|
|
if _, ok := value.(bool); !ok {
|
|
|
|
|
|
if str, ok := value.(string); ok {
|
|
|
|
|
|
return strings.ToLower(str) == "true" || str == "1", true
|
|
|
|
|
|
}
|
|
|
|
|
|
return false, true
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
return value, false
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// ВЕРСИИ
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
// getSortedSchemaMigrations возвращает отсортированный список миграций.
|
|
|
|
|
|
// НЕ берёт sm.mu — вызывающий должен обеспечить консистентность
|
|
|
|
|
|
// (например, вызвать под sm.mu.RLock или использовать снимок).
|
|
|
|
|
|
func (sm *SchemaMigrator) getSortedSchemaMigrations() []*SchemaMigration {
|
|
|
|
|
|
migrations := make([]*SchemaMigration, 0, len(sm.migrations))
|
|
|
|
|
|
for _, m := range sm.migrations {
|
|
|
|
|
|
migrations = append(migrations, m)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sort.Slice(migrations, func(i, j int) bool {
|
|
|
|
|
|
return sm.isVersionLess(migrations[i].Version, migrations[j].Version)
|
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
|
|
return migrations
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// isVersionLess сравнивает версии.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) isVersionLess(v1, v2 string) bool {
|
|
|
|
|
|
return compareVersions(v1, v2) < 0
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// isVersionGreater сравнивает версии.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) isVersionGreater(v1, v2 string) bool {
|
|
|
|
|
|
return compareVersions(v1, v2) > 0
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// versionToInt преобразует версию в число (для совместимости).
|
|
|
|
|
|
// ВНИМАНИЕ: могут быть коллизии (1.2.1000 == 1.3.0). Для сравнения
|
|
|
|
|
|
// используется compareVersions.
|
|
|
|
|
|
func (sm *SchemaMigrator) versionToInt(version string) int64 {
|
|
|
|
|
|
major, minor, patch := parseVersion(version)
|
|
|
|
|
|
return int64(major)*1_000_000_000 + int64(minor)*1_000_000 + int64(patch)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// parseVersion разбирает строку версии вида "X.Y.Z" (или "X.Y", или "X").
|
|
|
|
|
|
// Отсутствующие компоненты считаются нулями. Невалидные компоненты — 0.
|
|
|
|
|
|
func parseVersion(version string) (int, int, int) {
|
|
|
|
|
|
parts := strings.Split(version, ".")
|
|
|
|
|
|
var major, minor, patch int
|
|
|
|
|
|
if len(parts) > 0 {
|
|
|
|
|
|
fmt.Sscanf(parts[0], "%d", &major)
|
|
|
|
|
|
}
|
|
|
|
|
|
if len(parts) > 1 {
|
|
|
|
|
|
fmt.Sscanf(parts[1], "%d", &minor)
|
|
|
|
|
|
}
|
|
|
|
|
|
if len(parts) > 2 {
|
|
|
|
|
|
fmt.Sscanf(parts[2], "%d", &patch)
|
|
|
|
|
|
}
|
|
|
|
|
|
return major, minor, patch
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// compareVersions сравнивает две версии покомпонентно.
|
|
|
|
|
|
// Возвращает -1, 0, 1. Отсутствующие компоненты считаются нулями,
|
|
|
|
|
|
// что устраняет проблему "1.2" vs "1.2.0" (они равны).
|
|
|
|
|
|
func compareVersions(v1, v2 string) int {
|
|
|
|
|
|
major1, minor1, patch1 := parseVersion(v1)
|
|
|
|
|
|
major2, minor2, patch2 := parseVersion(v2)
|
|
|
|
|
|
|
|
|
|
|
|
if major1 != major2 {
|
|
|
|
|
|
if major1 < major2 {
|
|
|
|
|
|
return -1
|
|
|
|
|
|
}
|
|
|
|
|
|
return 1
|
|
|
|
|
|
}
|
|
|
|
|
|
if minor1 != minor2 {
|
|
|
|
|
|
if minor1 < minor2 {
|
|
|
|
|
|
return -1
|
|
|
|
|
|
}
|
|
|
|
|
|
return 1
|
|
|
|
|
|
}
|
|
|
|
|
|
if patch1 != patch2 {
|
|
|
|
|
|
if patch1 < patch2 {
|
|
|
|
|
|
return -1
|
|
|
|
|
|
}
|
|
|
|
|
|
return 1
|
|
|
|
|
|
}
|
|
|
|
|
|
return 0
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// КОНКРЕТНЫЕ РЕАЛИЗАЦИИ МИГРАЦИЙ СХЕМЫ
|
|
|
|
|
|
// =============================================================================
|
2026-10-03 20:33:34 +00:00
|
|
|
|
//
|
|
|
|
|
|
// ИЗМЕНЕНО (2026-10): сигнатуры принимают *storage.Batch вместо
|
|
|
|
|
|
// *storage.Transaction. Сами миграции работают только с in-memory
|
|
|
|
|
|
// структурой SchemaDefinition, поэтому параметр batch не используется
|
|
|
|
|
|
// в большинстве реализаций — он оставлен для совместимости с API и
|
|
|
|
|
|
// для случаев, когда миграция хочет добавить операции в СУБД.
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Неиспользуемые параметры помечены "_".
|
2026-10-03 20:33:34 +00:00
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaAddTimestamps добавляет поля created_at и updated_at.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaAddTimestamps(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if collSchema.Fields == nil {
|
|
|
|
|
|
collSchema.Fields = make(map[string]*FieldDefinition)
|
|
|
|
|
|
}
|
|
|
|
|
|
if _, exists := collSchema.Fields["created_at"]; !exists {
|
|
|
|
|
|
collSchema.Fields["created_at"] = &FieldDefinition{
|
|
|
|
|
|
Type: "int64",
|
|
|
|
|
|
Required: false,
|
|
|
|
|
|
Default: int64(0),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if _, exists := collSchema.Fields["updated_at"]; !exists {
|
|
|
|
|
|
collSchema.Fields["updated_at"] = &FieldDefinition{
|
|
|
|
|
|
Type: "int64",
|
|
|
|
|
|
Required: false,
|
|
|
|
|
|
Default: int64(0),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaRemoveTimestamps удаляет поля created_at и updated_at.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaRemoveTimestamps(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
delete(collSchema.Fields, "created_at")
|
|
|
|
|
|
delete(collSchema.Fields, "updated_at")
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaAddSoftDelete добавляет поддержку мягкого удаления.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaAddSoftDelete(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if collSchema.Fields == nil {
|
|
|
|
|
|
collSchema.Fields = make(map[string]*FieldDefinition)
|
|
|
|
|
|
}
|
|
|
|
|
|
if _, exists := collSchema.Fields["deleted_at"]; !exists {
|
|
|
|
|
|
collSchema.Fields["deleted_at"] = &FieldDefinition{
|
|
|
|
|
|
Type: "int64",
|
|
|
|
|
|
Required: false,
|
|
|
|
|
|
Default: int64(0),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if _, exists := collSchema.Fields["deleted"]; !exists {
|
|
|
|
|
|
collSchema.Fields["deleted"] = &FieldDefinition{
|
|
|
|
|
|
Type: "bool",
|
|
|
|
|
|
Required: false,
|
|
|
|
|
|
Default: false,
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaRemoveSoftDelete удаляет поддержку мягкого удаления.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaRemoveSoftDelete(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
delete(collSchema.Fields, "deleted_at")
|
|
|
|
|
|
delete(collSchema.Fields, "deleted")
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaAddVersioning добавляет версионирование.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaAddVersioning(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if collSchema.Fields == nil {
|
|
|
|
|
|
collSchema.Fields = make(map[string]*FieldDefinition)
|
|
|
|
|
|
}
|
|
|
|
|
|
if _, exists := collSchema.Fields["_version"]; !exists {
|
|
|
|
|
|
collSchema.Fields["_version"] = &FieldDefinition{
|
|
|
|
|
|
Type: "int64",
|
|
|
|
|
|
Required: false,
|
|
|
|
|
|
Default: int64(1),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaRemoveVersioning удаляет версионирование.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaRemoveVersioning(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
delete(collSchema.Fields, "_version")
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaAddIndexes добавляет индексы.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaAddIndexes(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for collectionName, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
defaultIndexes := []IndexDefinition{
|
|
|
|
|
|
{Name: "idx_created_at", Fields: []string{"created_at"}, Unique: false},
|
|
|
|
|
|
{Name: "idx_updated_at", Fields: []string{"updated_at"}, Unique: false},
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
for _, idx := range defaultIndexes {
|
|
|
|
|
|
found := false
|
|
|
|
|
|
for _, existing := range collSchema.Indexes {
|
|
|
|
|
|
if existing.Name == idx.Name {
|
|
|
|
|
|
found = true
|
|
|
|
|
|
break
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if !found {
|
|
|
|
|
|
collSchema.Indexes = append(collSchema.Indexes, idx)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
sm.logger.Debug(fmt.Sprintf("Added indexes to collection %s", collectionName))
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaRemoveIndexes удаляет индексы.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaRemoveIndexes(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
newIndexes := make([]IndexDefinition, 0)
|
|
|
|
|
|
for _, idx := range collSchema.Indexes {
|
|
|
|
|
|
if !strings.HasPrefix(idx.Name, "idx_") {
|
|
|
|
|
|
newIndexes = append(newIndexes, idx)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.Indexes = newIndexes
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaAddValidation добавляет валидацию полей.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaAddValidation(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if fieldDef, ok := collSchema.Fields["_id"]; ok && fieldDef != nil {
|
|
|
|
|
|
fieldDef.Required = true
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// migrateSchemaRemoveValidation удаляет валидацию.
|
|
|
|
|
|
func (sm *SchemaMigrator) migrateSchemaRemoveValidation(_ *storage.Batch, schema *SchemaDefinition) error {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
for _, collSchema := range schema.Collections {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
for _, fieldDef := range collSchema.Fields {
|
|
|
|
|
|
if fieldDef == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
|
|
|
|
|
fieldDef.Required = false
|
|
|
|
|
|
fieldDef.Validate = ""
|
|
|
|
|
|
fieldDef.ValidateValue = nil
|
|
|
|
|
|
}
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// ПУБЛИЧНЫЕ МЕТОДЫ ДЛЯ УПРАВЛЕНИЯ СХЕМОЙ
|
|
|
|
|
|
// =============================================================================
|
2026-10-03 20:33:34 +00:00
|
|
|
|
|
|
|
|
|
|
// AddCollectionSchema добавляет схему для новой коллекции.
|
|
|
|
|
|
func (sm *SchemaMigrator) AddCollectionSchema(collectionName string, schema *CollectionSchema) error {
|
|
|
|
|
|
if collectionName == "" {
|
|
|
|
|
|
return fmt.Errorf("empty collection name")
|
|
|
|
|
|
}
|
|
|
|
|
|
if schema == nil {
|
|
|
|
|
|
return fmt.Errorf("nil schema")
|
|
|
|
|
|
}
|
|
|
|
|
|
if schema.Fields == nil {
|
|
|
|
|
|
schema.Fields = make(map[string]*FieldDefinition)
|
|
|
|
|
|
}
|
|
|
|
|
|
if schema.Indexes == nil {
|
|
|
|
|
|
schema.Indexes = make([]IndexDefinition, 0)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if sm.schema == nil {
|
|
|
|
|
|
sm.schema = &SchemaDefinition{
|
|
|
|
|
|
Version: "1.0.0",
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if sm.schema.Collections == nil {
|
|
|
|
|
|
sm.schema.Collections = make(map[string]*CollectionSchema)
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.schema.Collections[collectionName] = schema
|
|
|
|
|
|
sm.schema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
|
|
|
|
|
|
return sm.saveSchemaLocked()
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// GetCollectionSchema возвращает схему коллекции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) GetCollectionSchema(collectionName string) *CollectionSchema {
|
|
|
|
|
|
sm.mu.RLock()
|
|
|
|
|
|
defer sm.mu.RUnlock()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if sm.schema == nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
if schema, exists := sm.schema.Collections[collectionName]; exists {
|
|
|
|
|
|
return schema
|
|
|
|
|
|
}
|
|
|
|
|
|
return sm.schema.Collections["_default"]
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// UpdateFieldSchema обновляет схему поля.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) UpdateFieldSchema(collectionName, fieldName string, fieldDef *FieldDefinition) error {
|
|
|
|
|
|
if collectionName == "" || fieldName == "" || fieldDef == nil {
|
|
|
|
|
|
return fmt.Errorf("invalid arguments")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if sm.schema == nil {
|
|
|
|
|
|
sm.schema = &SchemaDefinition{
|
|
|
|
|
|
Version: "1.0.0",
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if sm.schema.Collections == nil {
|
|
|
|
|
|
sm.schema.Collections = make(map[string]*CollectionSchema)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
collSchema, exists := sm.schema.Collections[collectionName]
|
2026-10-04 21:04:40 +00:00
|
|
|
|
if !exists || collSchema == nil {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
collSchema = &CollectionSchema{
|
|
|
|
|
|
Fields: make(map[string]*FieldDefinition),
|
|
|
|
|
|
Required: make([]string, 0),
|
|
|
|
|
|
Indexes: make([]IndexDefinition, 0),
|
|
|
|
|
|
UpdatedAt: time.Now().UnixMilli(),
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.schema.Collections[collectionName] = collSchema
|
|
|
|
|
|
}
|
|
|
|
|
|
if collSchema.Fields == nil {
|
|
|
|
|
|
collSchema.Fields = make(map[string]*FieldDefinition)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
collSchema.Fields[fieldName] = fieldDef
|
|
|
|
|
|
collSchema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
sm.schema.UpdatedAt = time.Now().UnixMilli()
|
|
|
|
|
|
|
|
|
|
|
|
return sm.saveSchemaLocked()
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// GetSchema возвращает текущую схему (копию).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) GetSchema() *SchemaDefinition {
|
|
|
|
|
|
sm.mu.RLock()
|
|
|
|
|
|
defer sm.mu.RUnlock()
|
|
|
|
|
|
|
|
|
|
|
|
if sm.schema == nil {
|
|
|
|
|
|
return &SchemaDefinition{
|
|
|
|
|
|
Version: "1.0.0",
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
copySchema := sm.deepCopySchemaLocked()
|
|
|
|
|
|
if copySchema == nil {
|
|
|
|
|
|
return &SchemaDefinition{
|
|
|
|
|
|
Version: sm.schema.Version,
|
|
|
|
|
|
Collections: make(map[string]*CollectionSchema),
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
return copySchema
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// GetSchemaVersion возвращает текущую версию схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) GetSchemaVersion() string {
|
|
|
|
|
|
sm.mu.RLock()
|
|
|
|
|
|
defer sm.mu.RUnlock()
|
|
|
|
|
|
if sm.schema == nil {
|
|
|
|
|
|
return "1.0.0"
|
|
|
|
|
|
}
|
|
|
|
|
|
return sm.schema.Version
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// GetMigrationStats возвращает статистику миграции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) GetMigrationStats() *MigrationStats {
|
|
|
|
|
|
sm.migrationStats.mu.RLock()
|
|
|
|
|
|
defer sm.migrationStats.mu.RUnlock()
|
|
|
|
|
|
|
|
|
|
|
|
stats := &MigrationStats{
|
|
|
|
|
|
TotalDocuments: sm.migrationStats.TotalDocuments,
|
|
|
|
|
|
MigratedDocuments: sm.migrationStats.MigratedDocuments,
|
|
|
|
|
|
FailedDocuments: sm.migrationStats.FailedDocuments,
|
|
|
|
|
|
SkippedDocuments: sm.migrationStats.SkippedDocuments,
|
|
|
|
|
|
StartTime: sm.migrationStats.StartTime,
|
|
|
|
|
|
EndTime: sm.migrationStats.EndTime,
|
|
|
|
|
|
}
|
|
|
|
|
|
return stats
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// SetMigrationStrategy устанавливает стратегию миграции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) SetMigrationStrategy(strategy DocumentMigrationStrategy) {
|
|
|
|
|
|
sm.mu.Lock()
|
|
|
|
|
|
defer sm.mu.Unlock()
|
|
|
|
|
|
sm.strategy = strategy
|
|
|
|
|
|
|
|
|
|
|
|
if sm.logger != nil {
|
|
|
|
|
|
strategyName := "InPlace"
|
|
|
|
|
|
if strategy == StrategyCopyNew {
|
|
|
|
|
|
strategyName = "CopyNew"
|
|
|
|
|
|
} else if strategy == StrategyLazy {
|
|
|
|
|
|
strategyName = "Lazy"
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.logger.Info(fmt.Sprintf("Migration strategy set to %s", strategyName))
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// ValidateDocumentAgainstSchema проверяет документ на соответствие схеме.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) ValidateDocumentAgainstSchema(doc *storage.Document, collectionName string) error {
|
|
|
|
|
|
if doc == nil {
|
|
|
|
|
|
return fmt.Errorf("nil document")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
sm.mu.RLock()
|
2026-10-04 21:04:40 +00:00
|
|
|
|
var collSchema *CollectionSchema
|
|
|
|
|
|
if sm.schema != nil {
|
|
|
|
|
|
if s, exists := sm.schema.Collections[collectionName]; exists {
|
|
|
|
|
|
collSchema = s
|
|
|
|
|
|
} else if s, exists := sm.schema.Collections["_default"]; exists {
|
|
|
|
|
|
collSchema = s
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
|
|
|
|
|
schemaCopy := copyCollectionSchema(collSchema)
|
|
|
|
|
|
sm.mu.RUnlock()
|
|
|
|
|
|
|
|
|
|
|
|
if schemaCopy == nil {
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
for _, requiredField := range schemaCopy.Required {
|
|
|
|
|
|
if _, err := doc.GetField(requiredField); err != nil {
|
|
|
|
|
|
return fmt.Errorf("required field '%s' is missing", requiredField)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
for fieldName, fieldDef := range schemaCopy.Fields {
|
|
|
|
|
|
if fieldDef == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
|
|
|
|
|
if value, err := doc.GetField(fieldName); err == nil {
|
|
|
|
|
|
if err := sm.validateFieldType(value, fieldDef.Type); err != nil {
|
|
|
|
|
|
return fmt.Errorf("field '%s' validation failed: %v", fieldName, err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if fieldDef.Validate != "" {
|
|
|
|
|
|
if err := sm.validateFieldValue(value, fieldDef); err != nil {
|
|
|
|
|
|
return fmt.Errorf("field '%s' value validation failed: %v", fieldName, err)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// validateFieldType проверяет тип поля.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) validateFieldType(value interface{}, expectedType string) error {
|
|
|
|
|
|
switch expectedType {
|
|
|
|
|
|
case "string":
|
|
|
|
|
|
if _, ok := value.(string); !ok {
|
|
|
|
|
|
return fmt.Errorf("expected string, got %T", value)
|
|
|
|
|
|
}
|
|
|
|
|
|
case "int64":
|
|
|
|
|
|
switch value.(type) {
|
|
|
|
|
|
case int, int32, int64, float64:
|
|
|
|
|
|
// допустимые типы
|
|
|
|
|
|
default:
|
|
|
|
|
|
return fmt.Errorf("expected int64, got %T", value)
|
|
|
|
|
|
}
|
|
|
|
|
|
case "float64":
|
|
|
|
|
|
switch value.(type) {
|
|
|
|
|
|
case float64, float32, int, int32, int64:
|
|
|
|
|
|
// допустимые типы
|
|
|
|
|
|
default:
|
|
|
|
|
|
return fmt.Errorf("expected float64, got %T", value)
|
|
|
|
|
|
}
|
|
|
|
|
|
case "bool":
|
|
|
|
|
|
if _, ok := value.(bool); !ok {
|
|
|
|
|
|
return fmt.Errorf("expected bool, got %T", value)
|
|
|
|
|
|
}
|
|
|
|
|
|
case "array", "object":
|
|
|
|
|
|
// без строгой проверки
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// validateFieldValue проверяет значение поля по правилам валидации.
|
|
|
|
|
|
//
|
|
|
|
|
|
// ИСПРАВЛЕНО: для правила "regex" используется regexp.MatchString,
|
|
|
|
|
|
// а не strings.Contains (имя не соответствовало поведению).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) validateFieldValue(value interface{}, fieldDef *FieldDefinition) error {
|
|
|
|
|
|
switch fieldDef.Validate {
|
|
|
|
|
|
case "min":
|
|
|
|
|
|
if min, ok := toFloat64(fieldDef.ValidateValue); ok {
|
|
|
|
|
|
if v, ok := toFloat64(value); ok && v < min {
|
|
|
|
|
|
return fmt.Errorf("value %.2f is less than minimum %.2f", v, min)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
case "max":
|
|
|
|
|
|
if max, ok := toFloat64(fieldDef.ValidateValue); ok {
|
|
|
|
|
|
if v, ok := toFloat64(value); ok && v > max {
|
|
|
|
|
|
return fmt.Errorf("value %.2f exceeds maximum %.2f", v, max)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
case "enum":
|
|
|
|
|
|
if enumVals, ok := fieldDef.ValidateValue.([]interface{}); ok {
|
|
|
|
|
|
found := false
|
|
|
|
|
|
for _, ev := range enumVals {
|
|
|
|
|
|
if fmt.Sprintf("%v", value) == fmt.Sprintf("%v", ev) {
|
|
|
|
|
|
found = true
|
|
|
|
|
|
break
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if !found {
|
|
|
|
|
|
return fmt.Errorf("value %v not in enum list", value)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
case "regex":
|
|
|
|
|
|
if pattern, ok := fieldDef.ValidateValue.(string); ok {
|
|
|
|
|
|
if str, ok := value.(string); ok {
|
2026-10-04 21:04:40 +00:00
|
|
|
|
re, err := regexp.Compile(pattern)
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
|
return fmt.Errorf("invalid regex pattern '%s': %v", pattern, err)
|
|
|
|
|
|
}
|
|
|
|
|
|
if !re.MatchString(str) {
|
2026-10-03 20:33:34 +00:00
|
|
|
|
return fmt.Errorf("value '%s' does not match pattern '%s'", str, pattern)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
return nil
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// toFloat64 пытается привести значение к float64.
|
|
|
|
|
|
func toFloat64(v interface{}) (float64, bool) {
|
|
|
|
|
|
switch x := v.(type) {
|
|
|
|
|
|
case float64:
|
|
|
|
|
|
return x, true
|
|
|
|
|
|
case float32:
|
|
|
|
|
|
return float64(x), true
|
|
|
|
|
|
case int:
|
|
|
|
|
|
return float64(x), true
|
|
|
|
|
|
case int32:
|
|
|
|
|
|
return float64(x), true
|
|
|
|
|
|
case int64:
|
|
|
|
|
|
return float64(x), true
|
|
|
|
|
|
}
|
|
|
|
|
|
return 0, false
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// СТАТУС
|
|
|
|
|
|
// =============================================================================
|
|
|
|
|
|
|
|
|
|
|
|
// GetSchemaStatus возвращает статус схемы.
|
|
|
|
|
|
//
|
|
|
|
|
|
// ИСПРАВЛЕНО: добавлена nil-проверка sm.schema. Раньше метод
|
|
|
|
|
|
// паниковал, если schema не была загружена.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) GetSchemaStatus() *SchemaStatus {
|
|
|
|
|
|
sm.mu.RLock()
|
|
|
|
|
|
defer sm.mu.RUnlock()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
version := "1.0.0"
|
|
|
|
|
|
var collections map[string]*CollectionSchema
|
|
|
|
|
|
var updatedAt int64
|
|
|
|
|
|
if sm.schema != nil {
|
|
|
|
|
|
version = sm.schema.Version
|
|
|
|
|
|
collections = sm.schema.Collections
|
|
|
|
|
|
updatedAt = sm.schema.UpdatedAt
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
status := &SchemaStatus{
|
2026-10-04 21:04:40 +00:00
|
|
|
|
CurrentVersion: version,
|
|
|
|
|
|
TotalCollections: len(collections),
|
2026-10-03 20:33:34 +00:00
|
|
|
|
TotalMigrations: len(sm.migrations),
|
|
|
|
|
|
AppliedMigrations: len(sm.applied),
|
|
|
|
|
|
PendingMigrations: len(sm.migrations) - len(sm.applied),
|
2026-10-04 21:04:40 +00:00
|
|
|
|
UpdatedAt: updatedAt,
|
|
|
|
|
|
Collections: make([]CollectionSchemaInfo, 0, len(collections)),
|
2026-10-03 20:33:34 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
for name, collSchema := range collections {
|
|
|
|
|
|
if collSchema == nil {
|
|
|
|
|
|
continue
|
|
|
|
|
|
}
|
2026-10-03 20:33:34 +00:00
|
|
|
|
status.Collections = append(status.Collections, CollectionSchemaInfo{
|
|
|
|
|
|
Name: name,
|
|
|
|
|
|
FieldsCount: len(collSchema.Fields),
|
|
|
|
|
|
IndexesCount: len(collSchema.Indexes),
|
|
|
|
|
|
UpdatedAt: collSchema.UpdatedAt,
|
|
|
|
|
|
})
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
return status
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// SchemaStatus — статус схемы.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type SchemaStatus struct {
|
|
|
|
|
|
CurrentVersion string `json:"current_version"`
|
|
|
|
|
|
TotalCollections int `json:"total_collections"`
|
|
|
|
|
|
TotalMigrations int `json:"total_migrations"`
|
|
|
|
|
|
AppliedMigrations int `json:"applied_migrations"`
|
|
|
|
|
|
PendingMigrations int `json:"pending_migrations"`
|
|
|
|
|
|
UpdatedAt int64 `json:"updated_at"`
|
|
|
|
|
|
Collections []CollectionSchemaInfo `json:"collections"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// CollectionSchemaInfo — информация о схеме коллекции.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type CollectionSchemaInfo struct {
|
|
|
|
|
|
Name string `json:"name"`
|
|
|
|
|
|
FieldsCount int `json:"fields_count"`
|
|
|
|
|
|
IndexesCount int `json:"indexes_count"`
|
|
|
|
|
|
UpdatedAt int64 `json:"updated_at"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// MigrationStatusItem — элемент статуса миграции (совместимость со старым кодом).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type MigrationStatusItem struct {
|
|
|
|
|
|
ID string `json:"id"`
|
|
|
|
|
|
Version string `json:"version"`
|
|
|
|
|
|
Description string `json:"description"`
|
|
|
|
|
|
Applied bool `json:"applied"`
|
|
|
|
|
|
AppliedAt int64 `json:"applied_at,omitempty"`
|
|
|
|
|
|
CreatedAt int64 `json:"created_at"`
|
|
|
|
|
|
Success bool `json:"success"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// MigrationStatus — статус миграций (совместимость со старым кодом).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type MigrationStatus struct {
|
|
|
|
|
|
CurrentVersion string `json:"current_version"`
|
|
|
|
|
|
TotalMigrations int `json:"total_migrations"`
|
|
|
|
|
|
AppliedMigrations int `json:"applied_migrations"`
|
|
|
|
|
|
PendingMigrations int `json:"pending_migrations"`
|
|
|
|
|
|
Migrations []MigrationStatusItem `json:"migrations"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// GetStatus возвращает статус миграций (для совместимости со старым кодом).
|
|
|
|
|
|
//
|
|
|
|
|
|
// ИСПРАВЛЕНО: добавлена nil-проверка sm.schema.
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) GetStatus() *MigrationStatus {
|
|
|
|
|
|
sm.mu.RLock()
|
|
|
|
|
|
defer sm.mu.RUnlock()
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
version := "1.0.0"
|
|
|
|
|
|
if sm.schema != nil {
|
|
|
|
|
|
version = sm.schema.Version
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-03 20:33:34 +00:00
|
|
|
|
status := &MigrationStatus{
|
2026-10-04 21:04:40 +00:00
|
|
|
|
CurrentVersion: version,
|
2026-10-03 20:33:34 +00:00
|
|
|
|
TotalMigrations: len(sm.migrations),
|
|
|
|
|
|
AppliedMigrations: len(sm.applied),
|
|
|
|
|
|
PendingMigrations: len(sm.migrations) - len(sm.applied),
|
|
|
|
|
|
Migrations: make([]MigrationStatusItem, 0),
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
for _, m := range sm.getSortedSchemaMigrations() {
|
|
|
|
|
|
item := MigrationStatusItem{
|
|
|
|
|
|
ID: m.ID,
|
|
|
|
|
|
Version: m.Version,
|
|
|
|
|
|
Description: m.Description,
|
|
|
|
|
|
CreatedAt: m.CreatedAt,
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if record, ok := sm.applied[m.ID]; ok {
|
|
|
|
|
|
item.Applied = true
|
|
|
|
|
|
item.AppliedAt = record.AppliedAt
|
|
|
|
|
|
item.Success = record.Success
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
status.Migrations = append(status.Migrations, item)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
return status
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// =============================================================================
|
|
|
|
|
|
// СОВМЕСТИМОСТЬ СО СТАРЫМ КОДОМ
|
|
|
|
|
|
// =============================================================================
|
2026-10-03 20:33:34 +00:00
|
|
|
|
|
|
|
|
|
|
// Migration (старая структура для совместимости).
|
|
|
|
|
|
//
|
|
|
|
|
|
// ИЗМЕНЕНО (2026-10): сигнатуры Up/Down принимают *storage.Batch вместо
|
|
|
|
|
|
// *storage.Transaction. Batch — это группа операций, применяемых
|
|
|
|
|
|
// атомарно через WAL. Семантика та же: миграция может добавить
|
|
|
|
|
|
// операции в batch, а batch.Commit() применит их атомарно.
|
|
|
|
|
|
type Migration struct {
|
|
|
|
|
|
ID string
|
|
|
|
|
|
Version string
|
|
|
|
|
|
Description string
|
|
|
|
|
|
Up func(batch *storage.Batch) error
|
|
|
|
|
|
Down func(batch *storage.Batch) error
|
|
|
|
|
|
CreatedAt int64
|
|
|
|
|
|
AppliedAt int64
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// MigrationRecord (старая структура для совместимости).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
type MigrationRecord struct {
|
|
|
|
|
|
ID string `json:"id"`
|
|
|
|
|
|
Version string `json:"version"`
|
|
|
|
|
|
Description string `json:"description"`
|
|
|
|
|
|
AppliedAt int64 `json:"applied_at"`
|
|
|
|
|
|
Success bool `json:"success"`
|
|
|
|
|
|
Error string `json:"error,omitempty"`
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// RegisterMigration регистрирует миграцию (старый интерфейс).
|
|
|
|
|
|
// Безопасно обрабатывает nil Up/Down.
|
|
|
|
|
|
func (sm *SchemaMigrator) RegisterMigration(m *Migration) {
|
|
|
|
|
|
if m == nil {
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
schemaMig := &SchemaMigration{
|
|
|
|
|
|
ID: m.ID,
|
|
|
|
|
|
Version: m.Version,
|
|
|
|
|
|
Description: m.Description,
|
|
|
|
|
|
CreatedAt: m.CreatedAt,
|
|
|
|
|
|
}
|
|
|
|
|
|
if m.Up != nil {
|
|
|
|
|
|
up := m.Up
|
|
|
|
|
|
schemaMig.Up = func(batch *storage.Batch, schema *SchemaDefinition) error {
|
|
|
|
|
|
return up(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
if m.Down != nil {
|
|
|
|
|
|
down := m.Down
|
|
|
|
|
|
schemaMig.Down = func(batch *storage.Batch, schema *SchemaDefinition) error {
|
|
|
|
|
|
return down(batch)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
sm.RegisterSchemaMigration(schemaMig)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-10-04 21:04:40 +00:00
|
|
|
|
// Migrate выполняет миграцию (старый интерфейс).
|
2026-10-03 20:33:34 +00:00
|
|
|
|
func (sm *SchemaMigrator) Migrate(targetVersion string) error {
|
|
|
|
|
|
return sm.MigrateSchema(targetVersion)
|
2026-10-04 21:04:40 +00:00
|
|
|
|
}
|