From 7804ce23380bda706236e816520bf34f2146c7ec Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=93=D1=80=D0=B8=D0=B3=D0=BE=D1=80=D0=B8=D0=B9=20=D0=A1?= =?UTF-8?q?=D0=B0=D1=84=D1=80=D0=BE=D0=BD=D0=BE=D0=B2?= Date: Sun, 4 Oct 2026 21:04:40 +0000 Subject: [PATCH] Update internal/migration/schema_migrator.go --- internal/migration/schema_migrator.go | 491 ++++++++++++++++++-------- 1 file changed, 348 insertions(+), 143 deletions(-) diff --git a/internal/migration/schema_migrator.go b/internal/migration/schema_migrator.go index 5d1a166..47c3766 100644 --- a/internal/migration/schema_migrator.go +++ b/internal/migration/schema_migrator.go @@ -9,16 +9,38 @@ */ // Файл: internal/migration/schema_migrator.go -// Назначение: Миграция схемы данных при обновлении версии СУБД -// Новая функциональность: Прозрачное обновление схемы данных (автоматическая миграция документов) +// Назначение: Миграция схемы данных при обновлении версии СУБД. +// Новая функциональность: прозрачное обновление схемы данных +// (автоматическая миграция документов). // -// ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции -// тип storage.Transaction был удалён. Все ссылки на *storage.Transaction -// заменены на *storage.Batch. Миграции схемы работают в основном с -// in-memory структурой SchemaDefinition, поэтому batch используется -// только как маркер атомарности: если пользовательская миграция пишет -// в СУБД, она может добавлять операции в batch, который коммитится -// в конце MigrateSchema атомарно через batch.Commit(). +// ИСПРАВЛЕНО (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 миграциях помечены "_". package migration @@ -27,6 +49,7 @@ import ( "fmt" "os" "path/filepath" + "regexp" "sort" "strings" "sync" @@ -36,14 +59,18 @@ import ( "futriis/internal/storage" ) -// SchemaDefinition определяет схему коллекции +// ============================================================================= +// СХЕМА +// ============================================================================= + +// SchemaDefinition определяет схему коллекции. type SchemaDefinition struct { Version string `json:"version"` Collections map[string]*CollectionSchema `json:"collections"` UpdatedAt int64 `json:"updated_at"` } -// CollectionSchema определяет схему одной коллекции +// CollectionSchema определяет схему одной коллекции. type CollectionSchema struct { Fields map[string]*FieldDefinition `json:"fields"` Required []string `json:"required"` @@ -51,7 +78,7 @@ type CollectionSchema struct { UpdatedAt int64 `json:"updated_at"` } -// FieldDefinition определяет поле в схеме +// FieldDefinition определяет поле в схеме. type FieldDefinition struct { Type string `json:"type"` // string, int, float, bool, array, object Required bool `json:"required"` @@ -60,13 +87,17 @@ type FieldDefinition struct { ValidateValue interface{} `json:"validate_value"` } -// IndexDefinition определяет индекс +// IndexDefinition определяет индекс. type IndexDefinition struct { Name string `json:"name"` Fields []string `json:"fields"` Unique bool `json:"unique"` } +// ============================================================================= +// МИГРАЦИИ +// ============================================================================= + // SchemaMigration представляет миграцию схемы. // // ИЗМЕНЕНО (2026-10): функции Up/Down теперь принимают *storage.Batch @@ -85,7 +116,7 @@ type SchemaMigration struct { AppliedAt int64 } -// SchemaMigrationRecord запись о применённой миграции схемы +// SchemaMigrationRecord — запись о применённой миграции схемы. type SchemaMigrationRecord struct { ID string `json:"id"` Version string `json:"version"` @@ -95,7 +126,7 @@ type SchemaMigrationRecord struct { Error string `json:"error,omitempty"` } -// DocumentMigrationStrategy стратегия миграции документа +// DocumentMigrationStrategy — стратегия миграции документа. type DocumentMigrationStrategy int const ( @@ -104,7 +135,11 @@ const ( StrategyLazy // Ленивая миграция при доступе ) -// SchemaMigrator управляет миграциями схемы данных +// ============================================================================= +// SCHEMA MIGRATOR +// ============================================================================= + +// SchemaMigrator управляет миграциями схемы данных. type SchemaMigrator struct { store *storage.Storage logger *log.Logger @@ -119,7 +154,7 @@ type SchemaMigrator struct { migrationStats *MigrationStats } -// MigrationStats статистика миграции +// MigrationStats — статистика миграции. type MigrationStats struct { TotalDocuments int64 MigratedDocuments int64 @@ -130,43 +165,48 @@ type MigrationStats struct { mu sync.RWMutex } -// NewSchemaMigrator создаёт новый мигратор схемы +// NewSchemaMigrator создаёт новый мигратор схемы. func NewSchemaMigrator(store *storage.Storage, logger *log.Logger, migrationDir string) *SchemaMigrator { sm := &SchemaMigrator{ - store: store, - logger: logger, - migrations: make(map[string]*SchemaMigration), - applied: make(map[string]*SchemaMigrationRecord), - migrationDir: migrationDir, - strategy: StrategyLazy, + store: store, + logger: logger, + migrations: make(map[string]*SchemaMigration), + applied: make(map[string]*SchemaMigrationRecord), + migrationDir: migrationDir, + strategy: StrategyLazy, migrationStats: &MigrationStats{}, schema: &SchemaDefinition{ Version: "1.0.0", Collections: make(map[string]*CollectionSchema), UpdatedAt: time.Now().UnixMilli(), }, + currentVersion: "1.0.0", } - // Создаём директорию для миграций (без паники при ошибке) + // Создаём директорию для миграций (без паники при ошибке). if err := os.MkdirAll(migrationDir, 0755); err != nil { if logger != nil { logger.Error(fmt.Sprintf("Failed to create migration directory %s: %v", migrationDir, err)) } } - // Загружаем существующую схему + // Загружаем существующую схему. sm.loadSchema() - // Загружаем применённые миграции + // Загружаем применённые миграции. sm.loadAppliedMigrations() - // Регистрируем встроенные миграции схемы + // Регистрируем встроенные миграции схемы. sm.registerBuiltinSchemaMigrations() return sm } -// loadSchema загружает определение схемы из файла +// ============================================================================= +// ЗАГРУЗКА / СОХРАНЕНИЕ +// ============================================================================= + +// loadSchema загружает определение схемы из файла. func (sm *SchemaMigrator) loadSchema() { path := filepath.Join(sm.migrationDir, "schema.json") @@ -201,11 +241,12 @@ func (sm *SchemaMigrator) loadSchema() { sm.mu.Unlock() if sm.logger != nil { - sm.logger.Info(fmt.Sprintf("Loaded schema version %s with %d collections", schema.Version, len(schema.Collections))) + sm.logger.Info(fmt.Sprintf("Loaded schema version %s with %d collections", + schema.Version, len(schema.Collections))) } } -// saveSchemaLocked сохраняет определение схемы атомарно (через временный файл + rename). +// saveSchemaLocked сохраняет определение схемы атомарно (временный файл + rename). // // ВАЖНО: метод НЕ берёт sm.mu, чтобы его можно было вызывать из функций, // уже удерживающих sm.mu.Lock(). Вызывающий сам отвечает за консистентность @@ -225,14 +266,7 @@ func (sm *SchemaMigrator) saveSchemaLocked() error { return os.Rename(tmpPath, path) } -// saveSchema — публичная обёртка, берущая RLock. -func (sm *SchemaMigrator) saveSchema() error { - sm.mu.RLock() - defer sm.mu.RUnlock() - return sm.saveSchemaLocked() -} - -// createDefaultSchema создаёт схему по умолчанию +// createDefaultSchema создаёт схему по умолчанию. func (sm *SchemaMigrator) createDefaultSchema() { sm.mu.Lock() defer sm.mu.Unlock() @@ -267,7 +301,7 @@ func (sm *SchemaMigrator) createDefaultSchema() { _ = sm.saveSchemaLocked() } -// registerBuiltinSchemaMigrations регистрирует встроенные миграции схемы +// registerBuiltinSchemaMigrations регистрирует встроенные миграции схемы. func (sm *SchemaMigrator) registerBuiltinSchemaMigrations() { sm.RegisterSchemaMigration(&SchemaMigration{ ID: "schema_001_add_timestamps", @@ -315,7 +349,7 @@ func (sm *SchemaMigrator) registerBuiltinSchemaMigrations() { }) } -// RegisterSchemaMigration регистрирует миграцию схемы +// RegisterSchemaMigration регистрирует миграцию схемы. func (sm *SchemaMigrator) RegisterSchemaMigration(m *SchemaMigration) { if m == nil || m.ID == "" { return @@ -325,7 +359,7 @@ func (sm *SchemaMigrator) RegisterSchemaMigration(m *SchemaMigration) { sm.migrations[m.ID] = m } -// loadAppliedMigrations загружает применённые миграции +// loadAppliedMigrations загружает применённые миграции. func (sm *SchemaMigrator) loadAppliedMigrations() { path := filepath.Join(sm.migrationDir, "schema_migrations.json") @@ -352,12 +386,22 @@ func (sm *SchemaMigrator) loadAppliedMigrations() { // saveAppliedMigrationsLocked сохраняет применённые миграции. // Не берёт sm.mu (вызывающий уже держит Lock). +// +// ИСПРАВЛЕНО: записи сортируются по AppliedAt, затем по ID — +// порядок в JSON стабильный и человекочитаемый. func (sm *SchemaMigrator) saveAppliedMigrationsLocked() error { records := make([]SchemaMigrationRecord, 0, len(sm.applied)) for _, record := range sm.applied { records = append(records, *record) } + 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 + }) + data, err := json.MarshalIndent(records, "", " ") if err != nil { return err @@ -371,18 +415,28 @@ func (sm *SchemaMigrator) saveAppliedMigrationsLocked() error { return os.Rename(tmpPath, path) } +// ============================================================================= +// МИГРАЦИЯ СХЕМЫ +// ============================================================================= + // MigrateSchema выполняет миграцию схемы до указанной версии. // -// ИЗМЕНЕНО (2026-10): вместо storage.BeginTransaction/CommitCurrentTransaction -// используется storage.NewBatch + batch.Commit. Batch создаётся для -// (db="_schema", coll="_migrations") — это виртуальная коллекция, куда -// миграции могут добавлять операции. Если миграция не пишет в СУБД, -// batch просто останется пустым и его Commit — no-op. +// ИЗМЕНЕНО (2026-10): используется storage.NewBatch + batch.Commit. +// Batch создаётся для (db="_system", coll="_schema_migrations") — +// это виртуальная коллекция, куда миграции могут добавлять операции. +// Если миграция не пишет в СУБД, batch просто останется пустым и +// его Commit — no-op. // -// ВАЖНО: этот метод НЕ удерживает sm.mu на всё время выполнения, -// чтобы избежать deadlock с saveSchema/saveAppliedMigrations и -// чтобы не блокировать чтение схемы на время долгих миграций. -// Лок берётся только для доступа к картам migrations/applied и к sm.schema. +// ИСПРАВЛЕНО (аудит): +// - добавлена nil-проверка sm.schema (паники больше не будет); +// - изменения применяются на диск ДО обновления in-memory. Если +// запись падает — sm.schema, sm.applied и sm.currentVersion +// откатываются к прежним значениям. +// +// ВАЖНО: метод НЕ удерживает sm.mu на всё время выполнения, чтобы +// избежать deadlock с saveSchema/saveAppliedMigrations и не блокировать +// чтение схемы на время долгих миграций. Лок берётся только для +// доступа к картам migrations/applied и к sm.schema. func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { if targetVersion == "" { return fmt.Errorf("target version is empty") @@ -392,9 +446,12 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { sm.logger.Info(fmt.Sprintf("Starting schema migration to version %s", targetVersion)) } - // Снимок текущего состояния под RLock + // Снимок текущего состояния под RLock. sm.mu.RLock() - currentVersion := sm.schema.Version + currentVersion := "1.0.0" + if sm.schema != nil { + currentVersion = sm.schema.Version + } allMigrations := make([]*SchemaMigration, 0, len(sm.migrations)) for _, m := range sm.migrations { allMigrations = append(allMigrations, m) @@ -412,12 +469,12 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { return nil } - // Сортируем по версии + // Сортируем по версии. sort.Slice(allMigrations, func(i, j int) bool { return sm.isVersionLess(allMigrations[i].Version, allMigrations[j].Version) }) - // Отбираем миграции для применения + // Отбираем миграции для применения. var toApply []*SchemaMigration for _, m := range allMigrations { if sm.isVersionGreater(m.Version, currentVersion) && !sm.isVersionGreater(m.Version, targetVersion) { @@ -443,7 +500,12 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { sm.mu.RUnlock() if schemaCopy == nil { - return fmt.Errorf("failed to copy schema") + // Если исходной схемы нет — создаём пустую с версией 1.0.0. + schemaCopy = &SchemaDefinition{ + Version: "1.0.0", + Collections: make(map[string]*CollectionSchema), + UpdatedAt: time.Now().UnixMilli(), + } } // Создаём batch для миграции. Batch — это группа операций, @@ -457,12 +519,9 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { return fmt.Errorf("failed to create batch") } - // Обработка паники с именованным возвратом err + // Обработка паники с именованным возвратом err. defer func() { if r := recover(); r != nil { - // Batch не коммитим — просто забываем о нём. - // Регистрация в activeBatches будет снята в Commit, - // но мы его не вызываем. Нужно явно снять регистрацию. if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(batch) } @@ -473,7 +532,7 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { } }() - // Новые applied-записи, которые закоммитим в sm.applied только при успехе + // Новые applied-записи, которые закоммитим в sm.applied только при успехе. newApplied := make(map[string]*SchemaMigrationRecord) for _, m := range toApply { @@ -484,7 +543,6 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { startTime := time.Now() if m.Up == nil { - // Снимаем регистрацию batch перед выходом с ошибкой. if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(batch) } @@ -521,29 +579,48 @@ func (sm *SchemaMigrator) MigrateSchema(targetVersion string) (err error) { return fmt.Errorf("failed to commit schema migration batch: %v", err) } - // Только теперь применяем изменения в память + // Только теперь применяем изменения: сначала пишем на диск, + // потом обновляем in-memory. Если запись падает — откатываем всё. sm.mu.Lock() + 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* работали с новым состоянием. sm.schema = schemaCopy + sm.currentVersion = targetVersion for id, rec := range newApplied { sm.applied[id] = rec } - sm.currentVersion = targetVersion - // Сохраняем на диск (под Lock, вызывая Locked-версии) + // 2. Пишем схему на диск. if err := sm.saveSchemaLocked(); err != nil { - sm.mu.Unlock() + sm.schema = oldSchema + sm.currentVersion = oldVersion + sm.applied = oldApplied return fmt.Errorf("failed to save schema: %v", err) } + + // 3. Пишем applied на диск. При ошибке откатываем схему и на диске тоже. if err := sm.saveAppliedMigrationsLocked(); err != nil { - sm.mu.Unlock() + 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)) + } return fmt.Errorf("failed to save applied migrations: %v", err) } - sm.mu.Unlock() if sm.logger != nil { sm.logger.Info(fmt.Sprintf("Schema migration to version %s completed successfully", targetVersion)) } - return nil } @@ -567,6 +644,10 @@ func (sm *SchemaMigrator) deepCopySchemaLocked() *SchemaDefinition { return © } +// ============================================================================= +// МИГРАЦИЯ ДОКУМЕНТОВ +// ============================================================================= + // MigrateDocument прозрачно мигрирует отдельный документ к актуальной схеме. // Возвращает копию документа и флаг changed (были ли изменения). func (sm *SchemaMigrator) MigrateDocument(doc *storage.Document, collectionName string) (*storage.Document, bool, error) { @@ -575,24 +656,31 @@ func (sm *SchemaMigrator) MigrateDocument(doc *storage.Document, collectionName } sm.mu.RLock() - collSchema, exists := sm.schema.Collections[collectionName] - if !exists { - collSchema = sm.schema.Collections["_default"] + 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 + } } - // Копируем схему, чтобы не держать RLock во время работы с документом + // Копируем схему, чтобы не держать RLock во время работы с документом. schemaCopy := copyCollectionSchema(collSchema) sm.mu.RUnlock() if schemaCopy == nil { - // Нет схемы — возвращаем документ без изменений + // Нет схемы — возвращаем документ без изменений. return doc.Clone(), false, nil } migratedDoc := doc.Clone() changed := false - // Добавляем недостающие поля с значениями по умолчанию + // Добавляем недостающие поля с значениями по умолчанию. for fieldName, fieldDef := range schemaCopy.Fields { + if fieldDef == nil { + continue + } if _, err := migratedDoc.GetField(fieldName); err != nil { if fieldDef.Default != nil { migratedDoc.SetField(fieldName, fieldDef.Default) @@ -604,8 +692,11 @@ func (sm *SchemaMigrator) MigrateDocument(doc *storage.Document, collectionName } } - // Проверяем типы полей и преобразуем при необходимости + // Проверяем типы полей и преобразуем при необходимости. for fieldName, fieldDef := range schemaCopy.Fields { + if fieldDef == nil { + continue + } if value, err := migratedDoc.GetField(fieldName); err == nil { converted, needsConversion := sm.convertFieldType(value, fieldDef.Type) if needsConversion { @@ -647,9 +738,11 @@ func copyCollectionSchema(src *CollectionSchema) *CollectionSchema { // MigrateCollectionDocuments прозрачно мигрирует все документы коллекции. // -// ИЗМЕНЕНО (2026-10): вместо поэлементного collection.Update используем -// batch для атомарности всей миграции коллекции. Это даёт атомарность -// на уровне коллекции: либо все документы мигрированы, либо ни один. +// ИСПРАВЛЕНО (аудит): +// - передаём version документа в Batch.AddUpdate, чтобы coll.Update +// увеличивал счётчик версии. Раньше миграция обновляла поля, но +// version документа не менялся. +// - ранний выход при пустой коллекции — не создаём лишний batch. func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collection) error { if collection == nil { return fmt.Errorf("nil collection") @@ -657,9 +750,22 @@ func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collect startTime := time.Now() + 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 + } + sm.migrationStats.mu.Lock() sm.migrationStats.StartTime = startTime.UnixMilli() - sm.migrationStats.TotalDocuments = collection.Count() + sm.migrationStats.TotalDocuments = int64(len(docs)) sm.migrationStats.MigratedDocuments = 0 sm.migrationStats.FailedDocuments = 0 sm.migrationStats.SkippedDocuments = 0 @@ -667,18 +773,16 @@ func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collect if sm.logger != nil { sm.logger.Info(fmt.Sprintf("Starting migration of collection %s with %d documents", - collection.Name(), collection.Count())) + collection.Name(), len(docs))) } - docs := collection.GetAllDocuments() - // Создаём один batch для всей коллекции — атомарность на уровне коллекции. batch := storage.NewBatch(collection.DBName(), collection.Name()) if batch == nil { return fmt.Errorf("failed to create batch") } - // Снимаем регистрацию batch при выходе в любом случае. + // Снимаем регистрацию batch при выходе, если не закоммитили. committed := false defer func() { if !committed { @@ -708,6 +812,11 @@ func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collect } updates := migratedDoc.GetFields() + // ИСПРАВЛЕНО: явно переносим version в updates, чтобы coll.Update + // увеличил счётчик версии документа. + if migratedDoc.Version > 0 { + updates["version"] = migratedDoc.Version + } batch.AddUpdate(doc.ID, updates) migratedCount++ } @@ -728,7 +837,7 @@ func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collect } committed = true } else { - // Batch пуст — просто снимаем регистрацию. + // Batch пуст — снимаем регистрацию вручную и отмечаем как "закоммичен". if bm := storage.GetBatchManager(); bm != nil { bm.UnregisterBatch(batch) } @@ -751,7 +860,7 @@ func (sm *SchemaMigrator) MigrateCollectionDocuments(collection *storage.Collect return nil } -// convertFieldType преобразует значение поля к нужному типу +// convertFieldType преобразует значение поля к нужному типу. func (sm *SchemaMigrator) convertFieldType(value interface{}, targetType string) (interface{}, bool) { switch targetType { case "string": @@ -803,6 +912,10 @@ func (sm *SchemaMigrator) convertFieldType(value interface{}, targetType string) return value, false } +// ============================================================================= +// ВЕРСИИ +// ============================================================================= + // getSortedSchemaMigrations возвращает отсортированный список миграций. // НЕ берёт sm.mu — вызывающий должен обеспечить консистентность // (например, вызвать под sm.mu.RLock или использовать снимок). @@ -819,12 +932,12 @@ func (sm *SchemaMigrator) getSortedSchemaMigrations() []*SchemaMigration { return migrations } -// isVersionLess сравнивает версии +// isVersionLess сравнивает версии. func (sm *SchemaMigrator) isVersionLess(v1, v2 string) bool { return compareVersions(v1, v2) < 0 } -// isVersionGreater сравнивает версии +// isVersionGreater сравнивает версии. func (sm *SchemaMigrator) isVersionGreater(v1, v2 string) bool { return compareVersions(v1, v2) > 0 } @@ -882,20 +995,26 @@ func compareVersions(v1, v2 string) int { return 0 } -// ========== Конкретные реализации миграций схемы ========== +// ============================================================================= +// КОНКРЕТНЫЕ РЕАЛИЗАЦИИ МИГРАЦИЙ СХЕМЫ +// ============================================================================= // // ИЗМЕНЕНО (2026-10): сигнатуры принимают *storage.Batch вместо // *storage.Transaction. Сами миграции работают только с in-memory // структурой SchemaDefinition, поэтому параметр batch не используется // в большинстве реализаций — он оставлен для совместимости с API и // для случаев, когда миграция хочет добавить операции в СУБД. +// Неиспользуемые параметры помечены "_". -// migrateSchemaAddTimestamps добавляет поля created_at и updated_at -func (sm *SchemaMigrator) migrateSchemaAddTimestamps(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaAddTimestamps добавляет поля created_at и updated_at. +func (sm *SchemaMigrator) migrateSchemaAddTimestamps(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } if collSchema.Fields == nil { collSchema.Fields = make(map[string]*FieldDefinition) } @@ -918,12 +1037,15 @@ func (sm *SchemaMigrator) migrateSchemaAddTimestamps(batch *storage.Batch, schem return nil } -// migrateSchemaRemoveTimestamps удаляет поля created_at и updated_at -func (sm *SchemaMigrator) migrateSchemaRemoveTimestamps(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaRemoveTimestamps удаляет поля created_at и updated_at. +func (sm *SchemaMigrator) migrateSchemaRemoveTimestamps(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } delete(collSchema.Fields, "created_at") delete(collSchema.Fields, "updated_at") collSchema.UpdatedAt = time.Now().UnixMilli() @@ -931,12 +1053,15 @@ func (sm *SchemaMigrator) migrateSchemaRemoveTimestamps(batch *storage.Batch, sc return nil } -// migrateSchemaAddSoftDelete добавляет поддержку мягкого удаления -func (sm *SchemaMigrator) migrateSchemaAddSoftDelete(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaAddSoftDelete добавляет поддержку мягкого удаления. +func (sm *SchemaMigrator) migrateSchemaAddSoftDelete(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } if collSchema.Fields == nil { collSchema.Fields = make(map[string]*FieldDefinition) } @@ -959,12 +1084,15 @@ func (sm *SchemaMigrator) migrateSchemaAddSoftDelete(batch *storage.Batch, schem return nil } -// migrateSchemaRemoveSoftDelete удаляет поддержку мягкого удаления -func (sm *SchemaMigrator) migrateSchemaRemoveSoftDelete(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaRemoveSoftDelete удаляет поддержку мягкого удаления. +func (sm *SchemaMigrator) migrateSchemaRemoveSoftDelete(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } delete(collSchema.Fields, "deleted_at") delete(collSchema.Fields, "deleted") collSchema.UpdatedAt = time.Now().UnixMilli() @@ -972,12 +1100,15 @@ func (sm *SchemaMigrator) migrateSchemaRemoveSoftDelete(batch *storage.Batch, sc return nil } -// migrateSchemaAddVersioning добавляет версионирование -func (sm *SchemaMigrator) migrateSchemaAddVersioning(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaAddVersioning добавляет версионирование. +func (sm *SchemaMigrator) migrateSchemaAddVersioning(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } if collSchema.Fields == nil { collSchema.Fields = make(map[string]*FieldDefinition) } @@ -993,24 +1124,30 @@ func (sm *SchemaMigrator) migrateSchemaAddVersioning(batch *storage.Batch, schem return nil } -// migrateSchemaRemoveVersioning удаляет версионирование -func (sm *SchemaMigrator) migrateSchemaRemoveVersioning(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaRemoveVersioning удаляет версионирование. +func (sm *SchemaMigrator) migrateSchemaRemoveVersioning(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } delete(collSchema.Fields, "_version") collSchema.UpdatedAt = time.Now().UnixMilli() } return nil } -// migrateSchemaAddIndexes добавляет индексы -func (sm *SchemaMigrator) migrateSchemaAddIndexes(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaAddIndexes добавляет индексы. +func (sm *SchemaMigrator) migrateSchemaAddIndexes(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for collectionName, collSchema := range schema.Collections { + if collSchema == nil { + continue + } defaultIndexes := []IndexDefinition{ {Name: "idx_created_at", Fields: []string{"created_at"}, Unique: false}, {Name: "idx_updated_at", Fields: []string{"updated_at"}, Unique: false}, @@ -1037,12 +1174,15 @@ func (sm *SchemaMigrator) migrateSchemaAddIndexes(batch *storage.Batch, schema * return nil } -// migrateSchemaRemoveIndexes удаляет индексы -func (sm *SchemaMigrator) migrateSchemaRemoveIndexes(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaRemoveIndexes удаляет индексы. +func (sm *SchemaMigrator) migrateSchemaRemoveIndexes(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } newIndexes := make([]IndexDefinition, 0) for _, idx := range collSchema.Indexes { if !strings.HasPrefix(idx.Name, "idx_") { @@ -1055,12 +1195,15 @@ func (sm *SchemaMigrator) migrateSchemaRemoveIndexes(batch *storage.Batch, schem return nil } -// migrateSchemaAddValidation добавляет валидацию полей -func (sm *SchemaMigrator) migrateSchemaAddValidation(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaAddValidation добавляет валидацию полей. +func (sm *SchemaMigrator) migrateSchemaAddValidation(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } if fieldDef, ok := collSchema.Fields["_id"]; ok && fieldDef != nil { fieldDef.Required = true } @@ -1069,12 +1212,15 @@ func (sm *SchemaMigrator) migrateSchemaAddValidation(batch *storage.Batch, schem return nil } -// migrateSchemaRemoveValidation удаляет валидацию -func (sm *SchemaMigrator) migrateSchemaRemoveValidation(batch *storage.Batch, schema *SchemaDefinition) error { +// migrateSchemaRemoveValidation удаляет валидацию. +func (sm *SchemaMigrator) migrateSchemaRemoveValidation(_ *storage.Batch, schema *SchemaDefinition) error { if schema == nil { return fmt.Errorf("nil schema") } for _, collSchema := range schema.Collections { + if collSchema == nil { + continue + } for _, fieldDef := range collSchema.Fields { if fieldDef == nil { continue @@ -1088,12 +1234,11 @@ func (sm *SchemaMigrator) migrateSchemaRemoveValidation(batch *storage.Batch, sc return nil } -// ========== Публичные методы для управления схемой ========== +// ============================================================================= +// ПУБЛИЧНЫЕ МЕТОДЫ ДЛЯ УПРАВЛЕНИЯ СХЕМОЙ +// ============================================================================= // AddCollectionSchema добавляет схему для новой коллекции. -// -// Изменение сохраняется в отдельной копии, чтобы избежать deadlock -// с saveSchemaLocked (вызывается под Lock). func (sm *SchemaMigrator) AddCollectionSchema(collectionName string, schema *CollectionSchema) error { if collectionName == "" { return fmt.Errorf("empty collection name") @@ -1111,6 +1256,13 @@ func (sm *SchemaMigrator) AddCollectionSchema(collectionName string, schema *Col sm.mu.Lock() defer sm.mu.Unlock() + 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) } @@ -1120,18 +1272,21 @@ func (sm *SchemaMigrator) AddCollectionSchema(collectionName string, schema *Col return sm.saveSchemaLocked() } -// GetCollectionSchema возвращает схему коллекции +// GetCollectionSchema возвращает схему коллекции. func (sm *SchemaMigrator) GetCollectionSchema(collectionName string) *CollectionSchema { sm.mu.RLock() defer sm.mu.RUnlock() + if sm.schema == nil { + return nil + } if schema, exists := sm.schema.Collections[collectionName]; exists { return schema } return sm.schema.Collections["_default"] } -// UpdateFieldSchema обновляет схему поля +// UpdateFieldSchema обновляет схему поля. func (sm *SchemaMigrator) UpdateFieldSchema(collectionName, fieldName string, fieldDef *FieldDefinition) error { if collectionName == "" || fieldName == "" || fieldDef == nil { return fmt.Errorf("invalid arguments") @@ -1140,8 +1295,19 @@ func (sm *SchemaMigrator) UpdateFieldSchema(collectionName, fieldName string, fi sm.mu.Lock() defer sm.mu.Unlock() + 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) + } + collSchema, exists := sm.schema.Collections[collectionName] - if !exists { + if !exists || collSchema == nil { collSchema = &CollectionSchema{ Fields: make(map[string]*FieldDefinition), Required: make([]string, 0), @@ -1161,7 +1327,7 @@ func (sm *SchemaMigrator) UpdateFieldSchema(collectionName, fieldName string, fi return sm.saveSchemaLocked() } -// GetSchema возвращает текущую схему (копию) +// GetSchema возвращает текущую схему (копию). func (sm *SchemaMigrator) GetSchema() *SchemaDefinition { sm.mu.RLock() defer sm.mu.RUnlock() @@ -1183,7 +1349,7 @@ func (sm *SchemaMigrator) GetSchema() *SchemaDefinition { return copySchema } -// GetSchemaVersion возвращает текущую версию схемы +// GetSchemaVersion возвращает текущую версию схемы. func (sm *SchemaMigrator) GetSchemaVersion() string { sm.mu.RLock() defer sm.mu.RUnlock() @@ -1193,7 +1359,7 @@ func (sm *SchemaMigrator) GetSchemaVersion() string { return sm.schema.Version } -// GetMigrationStats возвращает статистику миграции +// GetMigrationStats возвращает статистику миграции. func (sm *SchemaMigrator) GetMigrationStats() *MigrationStats { sm.migrationStats.mu.RLock() defer sm.migrationStats.mu.RUnlock() @@ -1209,7 +1375,7 @@ func (sm *SchemaMigrator) GetMigrationStats() *MigrationStats { return stats } -// SetMigrationStrategy устанавливает стратегию миграции +// SetMigrationStrategy устанавливает стратегию миграции. func (sm *SchemaMigrator) SetMigrationStrategy(strategy DocumentMigrationStrategy) { sm.mu.Lock() defer sm.mu.Unlock() @@ -1226,16 +1392,20 @@ func (sm *SchemaMigrator) SetMigrationStrategy(strategy DocumentMigrationStrateg } } -// ValidateDocumentAgainstSchema проверяет документ на соответствие схеме +// ValidateDocumentAgainstSchema проверяет документ на соответствие схеме. func (sm *SchemaMigrator) ValidateDocumentAgainstSchema(doc *storage.Document, collectionName string) error { if doc == nil { return fmt.Errorf("nil document") } sm.mu.RLock() - collSchema, exists := sm.schema.Collections[collectionName] - if !exists { - collSchema = sm.schema.Collections["_default"] + 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 + } } schemaCopy := copyCollectionSchema(collSchema) sm.mu.RUnlock() @@ -1269,7 +1439,7 @@ func (sm *SchemaMigrator) ValidateDocumentAgainstSchema(doc *storage.Document, c return nil } -// validateFieldType проверяет тип поля +// validateFieldType проверяет тип поля. func (sm *SchemaMigrator) validateFieldType(value interface{}, expectedType string) error { switch expectedType { case "string": @@ -1300,7 +1470,10 @@ func (sm *SchemaMigrator) validateFieldType(value interface{}, expectedType stri return nil } -// validateFieldValue проверяет значение поля по правилам валидации +// validateFieldValue проверяет значение поля по правилам валидации. +// +// ИСПРАВЛЕНО: для правила "regex" используется regexp.MatchString, +// а не strings.Contains (имя не соответствовало поведению). func (sm *SchemaMigrator) validateFieldValue(value interface{}, fieldDef *FieldDefinition) error { switch fieldDef.Validate { case "min": @@ -1331,7 +1504,11 @@ func (sm *SchemaMigrator) validateFieldValue(value interface{}, fieldDef *FieldD case "regex": if pattern, ok := fieldDef.ValidateValue.(string); ok { if str, ok := value.(string); ok { - if !strings.Contains(str, pattern) { + re, err := regexp.Compile(pattern) + if err != nil { + return fmt.Errorf("invalid regex pattern '%s': %v", pattern, err) + } + if !re.MatchString(str) { return fmt.Errorf("value '%s' does not match pattern '%s'", str, pattern) } } @@ -1357,22 +1534,41 @@ func toFloat64(v interface{}) (float64, bool) { return 0, false } -// GetSchemaStatus возвращает статус схемы +// ============================================================================= +// СТАТУС +// ============================================================================= + +// GetSchemaStatus возвращает статус схемы. +// +// ИСПРАВЛЕНО: добавлена nil-проверка sm.schema. Раньше метод +// паниковал, если schema не была загружена. func (sm *SchemaMigrator) GetSchemaStatus() *SchemaStatus { sm.mu.RLock() defer sm.mu.RUnlock() + 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 + } + status := &SchemaStatus{ - CurrentVersion: sm.schema.Version, - TotalCollections: len(sm.schema.Collections), + CurrentVersion: version, + TotalCollections: len(collections), TotalMigrations: len(sm.migrations), AppliedMigrations: len(sm.applied), PendingMigrations: len(sm.migrations) - len(sm.applied), - UpdatedAt: sm.schema.UpdatedAt, - Collections: make([]CollectionSchemaInfo, 0, len(sm.schema.Collections)), + UpdatedAt: updatedAt, + Collections: make([]CollectionSchemaInfo, 0, len(collections)), } - for name, collSchema := range sm.schema.Collections { + for name, collSchema := range collections { + if collSchema == nil { + continue + } status.Collections = append(status.Collections, CollectionSchemaInfo{ Name: name, FieldsCount: len(collSchema.Fields), @@ -1384,7 +1580,7 @@ func (sm *SchemaMigrator) GetSchemaStatus() *SchemaStatus { return status } -// SchemaStatus статус схемы +// SchemaStatus — статус схемы. type SchemaStatus struct { CurrentVersion string `json:"current_version"` TotalCollections int `json:"total_collections"` @@ -1395,7 +1591,7 @@ type SchemaStatus struct { Collections []CollectionSchemaInfo `json:"collections"` } -// CollectionSchemaInfo информация о схеме коллекции +// CollectionSchemaInfo — информация о схеме коллекции. type CollectionSchemaInfo struct { Name string `json:"name"` FieldsCount int `json:"fields_count"` @@ -1403,7 +1599,7 @@ type CollectionSchemaInfo struct { UpdatedAt int64 `json:"updated_at"` } -// MigrationStatusItem элемент статуса миграции (совместимость со старым кодом) +// MigrationStatusItem — элемент статуса миграции (совместимость со старым кодом). type MigrationStatusItem struct { ID string `json:"id"` Version string `json:"version"` @@ -1414,7 +1610,7 @@ type MigrationStatusItem struct { Success bool `json:"success"` } -// MigrationStatus статус миграций (совместимость со старым кодом) +// MigrationStatus — статус миграций (совместимость со старым кодом). type MigrationStatus struct { CurrentVersion string `json:"current_version"` TotalMigrations int `json:"total_migrations"` @@ -1423,13 +1619,20 @@ type MigrationStatus struct { Migrations []MigrationStatusItem `json:"migrations"` } -// GetStatus возвращает статус миграций (для совместимости со старым кодом) +// GetStatus возвращает статус миграций (для совместимости со старым кодом). +// +// ИСПРАВЛЕНО: добавлена nil-проверка sm.schema. func (sm *SchemaMigrator) GetStatus() *MigrationStatus { sm.mu.RLock() defer sm.mu.RUnlock() + version := "1.0.0" + if sm.schema != nil { + version = sm.schema.Version + } + status := &MigrationStatus{ - CurrentVersion: sm.schema.Version, + CurrentVersion: version, TotalMigrations: len(sm.migrations), AppliedMigrations: len(sm.applied), PendingMigrations: len(sm.migrations) - len(sm.applied), @@ -1456,7 +1659,9 @@ func (sm *SchemaMigrator) GetStatus() *MigrationStatus { return status } -// ========== Совместимость со старым кодом ========== +// ============================================================================= +// СОВМЕСТИМОСТЬ СО СТАРЫМ КОДОМ +// ============================================================================= // Migration (старая структура для совместимости). // @@ -1474,7 +1679,7 @@ type Migration struct { AppliedAt int64 } -// MigrationRecord (старая структура для совместимости) +// MigrationRecord (старая структура для совместимости). type MigrationRecord struct { ID string `json:"id"` Version string `json:"version"` @@ -1511,7 +1716,7 @@ func (sm *SchemaMigrator) RegisterMigration(m *Migration) { sm.RegisterSchemaMigration(schemaMig) } -// Migrate выполняет миграцию (старый интерфейс) +// Migrate выполняет миграцию (старый интерфейс). func (sm *SchemaMigrator) Migrate(targetVersion string) error { return sm.MigrateSchema(targetVersion) -} +} \ No newline at end of file