diff --git a/internal/cluster/crossdc_migrator.go b/internal/cluster/crossdc_migrator.go new file mode 100644 index 0000000..f49b41f --- /dev/null +++ b/internal/cluster/crossdc_migrator.go @@ -0,0 +1,1609 @@ +/* + * 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/cluster/crossdc_migrator.go +// Назначение: Кросс-датацентровая миграция данных +// Алгоритм: Асинхронная репликация с CDC (Change Data Capture) и очередью изменений + +package cluster + +import ( + "bytes" + "compress/gzip" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "net/http" + "os" + "path/filepath" + "sort" + "sync" + "sync/atomic" + "time" + + "futriis/internal/config" + "futriis/internal/log" + "futriis/internal/storage" +) + +// ============================================================================= +// ТИПЫ ДАННЫХ ДЛЯ МИГРАЦИИ +// ============================================================================= + +// MigrationStatus представляет статус миграции +type MigrationStatus string + +const ( + MigrationStatusIdle MigrationStatus = "idle" + MigrationStatusPreparing MigrationStatus = "preparing" + MigrationStatusMigrating MigrationStatus = "migrating" + MigrationStatusDeltaSync MigrationStatus = "delta_sync" + MigrationStatusValidating MigrationStatus = "validating" + MigrationStatusCompleted MigrationStatus = "completed" + MigrationStatusFailed MigrationStatus = "failed" + MigrationStatusPaused MigrationStatus = "paused" +) + +// MigrationMode представляет режим миграции +type MigrationMode string + +const ( + MigrationModeManual MigrationMode = "manual" // Полностью ручной + MigrationModeSemiAuto MigrationMode = "semi_auto" // Полуавтоматический + MigrationModeAuto MigrationMode = "auto" // Полностью автоматический +) + +// ChangeType представляет тип изменения +type ChangeType string + +const ( + ChangeInsert ChangeType = "insert" + ChangeUpdate ChangeType = "update" + ChangeDelete ChangeType = "delete" +) + +// ChangeRecord представляет запись об изменении +type ChangeRecord struct { + ID string `json:"id"` + Database string `json:"database"` + Collection string `json:"collection"` + DocumentID string `json:"document_id"` + ChangeType ChangeType `json:"change_type"` + Document map[string]interface{} `json:"document,omitempty"` + PreviousDoc map[string]interface{} `json:"previous_doc,omitempty"` + Timestamp int64 `json:"timestamp"` + Version uint64 `json:"version"` + Checksum string `json:"checksum"` + Applied bool `json:"applied"` + AppliedAt int64 `json:"applied_at,omitempty"` + LSN uint64 `json:"lsn,omitempty"` +} + +// MigrationTask представляет задачу миграции +type MigrationTask struct { + ID string `json:"id"` + SourceDC string `json:"source_dc"` + TargetDC string `json:"target_dc"` + Databases []string `json:"databases"` + Collections map[string][]string `json:"collections"` // database -> collections + Status MigrationStatus `json:"status"` + ProgressPercent float64 `json:"progress_percent"` + TotalDocuments int64 `json:"total_documents"` + MigratedDocs int64 `json:"migrated_docs"` + FailedDocs int64 `json:"failed_docs"` + SkippedDocs int64 `json:"skipped_docs"` + StartTime int64 `json:"start_time"` + EndTime int64 `json:"end_time"` + LastCheckpoint int64 `json:"last_checkpoint"` + CheckpointData map[string]interface{} `json:"checkpoint_data"` + Error string `json:"error,omitempty"` + CreatedAt int64 `json:"created_at"` + UpdatedAt int64 `json:"updated_at"` + mu sync.RWMutex `json:"-"` +} + +// MigrationStats представляет статистику миграции +type MigrationStats struct { + TotalChanges int64 `json:"total_changes"` + AppliedChanges int64 `json:"applied_changes"` + FailedChanges int64 `json:"failed_changes"` + SkippedChanges int64 `json:"skipped_changes"` + LatencyAvg int64 `json:"latency_avg_ms"` + LatencyMax int64 `json:"latency_max_ms"` + Throughput int64 `json:"throughput_docs_per_sec"` + mu sync.RWMutex +} + +// MigrationCheckpoint представляет чекпоинт миграции +type MigrationCheckpoint struct { + TaskID string `json:"task_id"` + LastLSN uint64 `json:"last_lsn"` + ProcessedDocs map[string]bool `json:"processed_docs"` + Timestamp int64 `json:"timestamp"` + Version uint64 `json:"version"` +} + +// ============================================================================= +// CHANGE QUEUE - ОЧЕРЕДЬ ИЗМЕНЕНИЙ С ПЕРСИСТЕНТНОСТЬЮ +// ============================================================================= + +// ChangeQueue представляет очередь изменений для миграции +type ChangeQueue struct { + changes []*ChangeRecord + mu sync.RWMutex + maxSize int + persisted bool + filePath string + lastLSN atomic.Uint64 +} + +// NewChangeQueue создаёт новую очередь изменений +func NewChangeQueue(maxSize int, filePath string) *ChangeQueue { + q := &ChangeQueue{ + changes: make([]*ChangeRecord, 0, maxSize), + maxSize: maxSize, + filePath: filePath, + persisted: false, + } + + // Загружаем сохранённую очередь + if filePath != "" { + q.load() + } + + return q +} + +// load загружает очередь из файла +func (cq *ChangeQueue) load() { + data, err := os.ReadFile(cq.filePath) + if err != nil { + return + } + + var records []*ChangeRecord + if err := json.Unmarshal(data, &records); err != nil { + return + } + + cq.mu.Lock() + defer cq.mu.Unlock() + cq.changes = records + + // Восстанавливаем последний LSN + for _, r := range records { + if r.LSN > cq.lastLSN.Load() { + cq.lastLSN.Store(r.LSN) + } + } +} + +// save сохраняет очередь в файл +func (cq *ChangeQueue) save() { + if cq.filePath == "" { + return + } + + cq.mu.RLock() + data, err := json.Marshal(cq.changes) + cq.mu.RUnlock() + + if err != nil { + return + } + + os.WriteFile(cq.filePath, data, 0644) +} + +// Add добавляет изменение в очередь +func (cq *ChangeQueue) Add(change *ChangeRecord) bool { + cq.mu.Lock() + defer cq.mu.Unlock() + + if len(cq.changes) >= cq.maxSize { + // Удаляем половину старых записей + cq.changes = cq.changes[len(cq.changes)/2:] + } + + // Устанавливаем LSN + change.LSN = cq.lastLSN.Add(1) + + cq.changes = append(cq.changes, change) + + // Асинхронное сохранение + if cq.filePath != "" { + go cq.save() + } + + return true +} + +// GetAll возвращает все изменения и очищает очередь +func (cq *ChangeQueue) GetAll() []*ChangeRecord { + cq.mu.Lock() + defer cq.mu.Unlock() + + changes := cq.changes + cq.changes = make([]*ChangeRecord, 0, cq.maxSize) + + if cq.filePath != "" { + os.WriteFile(cq.filePath, []byte("[]"), 0644) + } + + return changes +} + +// GetChangesSince возвращает изменения после указанного LSN +func (cq *ChangeQueue) GetChangesSince(lsn uint64) []*ChangeRecord { + cq.mu.RLock() + defer cq.mu.RUnlock() + + result := make([]*ChangeRecord, 0) + for _, change := range cq.changes { + if change.LSN > lsn { + result = append(result, change) + } + } + return result +} + +// GetChangesAfter возвращает изменения после указанного времени +func (cq *ChangeQueue) GetChangesAfter(timestamp int64) []*ChangeRecord { + cq.mu.RLock() + defer cq.mu.RUnlock() + + result := make([]*ChangeRecord, 0) + for _, change := range cq.changes { + if change.Timestamp > timestamp { + result = append(result, change) + } + } + return result +} + +// GetLastLSN возвращает последний LSN +func (cq *ChangeQueue) GetLastLSN() uint64 { + return cq.lastLSN.Load() +} + +// Clear очищает очередь +func (cq *ChangeQueue) Clear() { + cq.mu.Lock() + defer cq.mu.Unlock() + + cq.changes = make([]*ChangeRecord, 0, cq.maxSize) + if cq.filePath != "" { + os.WriteFile(cq.filePath, []byte("[]"), 0644) + } +} + +// Size возвращает размер очереди +func (cq *ChangeQueue) Size() int { + cq.mu.RLock() + defer cq.mu.RUnlock() + return len(cq.changes) +} + +// ============================================================================= +// ОСНОВНАЯ СТРУКТУРА МИГРАТОРА +// ============================================================================= + +// CrossDCMigrator управляет кросс-датацентровой миграцией данных +type CrossDCMigrator struct { + config *config.MigrationConfig + store *storage.Storage + logger *log.Logger + httpClient *http.Client + mu sync.RWMutex + tasks map[string]*MigrationTask + stats *MigrationStats + changeQueue *ChangeQueue + active atomic.Bool + stopChan chan struct{} + wg sync.WaitGroup + mode MigrationMode + currentTaskID string + checkpointPath string + resumedFrom *MigrationTask + deltaTicker *time.Ticker + changeSubscribers []chan *ChangeRecord + subscriberMu sync.RWMutex +} + +// NewCrossDCMigrator создаёт новый экземпляр мигратора +func NewCrossDCMigrator(cfg *config.MigrationConfig, store *storage.Storage, logger *log.Logger) *CrossDCMigrator { + if cfg == nil { + cfg = &config.MigrationConfig{ + Enabled: false, + Mode: "semi_auto", + Source: &config.DatacenterConfig{ + Name: "dc-primary", + Endpoint: "localhost:8080", + TimeoutSec: 30, + }, + Target: &config.DatacenterConfig{ + Name: "dc-secondary", + Endpoint: "localhost:8081", + TimeoutSec: 30, + }, + Settings: &config.MigrationSettings{ + BatchSize: 1000, + Workers: 4, + Compression: "snappy", + ResumeEnabled: true, + CheckpointIntervalSec: 30, + MaxRetries: 3, + RetryBackoffSec: 5, + }, + Delta: &config.DeltaSyncConfig{ + Enabled: true, + IntervalSec: 60, + MaxLagSec: 300, + }, + Validation: &config.ValidationConfig{ + Enabled: true, + SamplePercent: 10, + MaxErrors: 100, + }, + } + } + + cdm := &CrossDCMigrator{ + config: cfg, + store: store, + logger: logger, + tasks: make(map[string]*MigrationTask), + stats: &MigrationStats{}, + changeQueue: NewChangeQueue(100000, "futriis/migration/change_queue.json"), + stopChan: make(chan struct{}), + mode: MigrationMode(cfg.Mode), + checkpointPath: "futriis/migration/checkpoints", + changeSubscribers: make([]chan *ChangeRecord, 0), + httpClient: &http.Client{ + Timeout: time.Duration(cfg.Source.TimeoutSec) * time.Second, + }, + } + + // Создаём директорию для чекпоинтов + os.MkdirAll(cdm.checkpointPath, 0755) + os.MkdirAll("futriis/migration", 0755) + + // Загружаем сохранённые задачи + cdm.loadTasks() + + return cdm +} + +// Start запускает мигратор +func (cdm *CrossDCMigrator) Start() { + if !cdm.active.CompareAndSwap(false, true) { + return + } + + cdm.logger.Info("Cross-datacenter migrator started") + + // Восстанавливаем незавершённые миграции + cdm.resumePendingMigrations() + + // Запускаем обработку изменений + cdm.wg.Add(1) + go cdm.processChanges() + + // Запускаем дельта-синхронизацию если включена + if cdm.config.Delta.Enabled { + cdm.wg.Add(1) + go cdm.deltaSyncLoop() + } + + // Запускаем сохранение чекпоинтов + cdm.wg.Add(1) + go cdm.checkpointLoop() +} + +// Stop останавливает мигратор +func (cdm *CrossDCMigrator) Stop() { + if !cdm.active.Load() { + return + } + + close(cdm.stopChan) + cdm.wg.Wait() + cdm.active.Store(false) + + cdm.saveTasks() + cdm.logger.Info("Cross-datacenter migrator stopped") +} + +// ============================================================================= +// ОСНОВНЫЕ МЕТОДЫ МИГРАЦИИ +// ============================================================================= + +// StartMigration начинает миграцию данных +func (cdm *CrossDCMigrator) StartMigration(sourceDC, targetDC string, databases []string, collections map[string][]string) (*MigrationTask, error) { + if !cdm.config.Enabled { + return nil, fmt.Errorf("migration is disabled in configuration") + } + + cdm.mu.Lock() + defer cdm.mu.Unlock() + + // Проверяем, есть ли активная миграция + if cdm.currentTaskID != "" { + if task, ok := cdm.tasks[cdm.currentTaskID]; ok { + task.mu.RLock() + status := task.Status + task.mu.RUnlock() + if status == MigrationStatusMigrating || status == MigrationStatusDeltaSync { + return nil, fmt.Errorf("a migration is already in progress") + } + } + } + + taskID := fmt.Sprintf("mig_%d_%s", time.Now().UnixNano(), sourceDC) + + task := &MigrationTask{ + ID: taskID, + SourceDC: sourceDC, + TargetDC: targetDC, + Databases: databases, + Collections: collections, + Status: MigrationStatusPreparing, + ProgressPercent: 0, + StartTime: time.Now().UnixMilli(), + CheckpointData: make(map[string]interface{}), + CreatedAt: time.Now().UnixMilli(), + UpdatedAt: time.Now().UnixMilli(), + } + + cdm.tasks[taskID] = task + cdm.currentTaskID = taskID + + cdm.logger.Info(fmt.Sprintf("Migration task %s created: %s -> %s", taskID, sourceDC, targetDC)) + + // Запускаем миграцию асинхронно + cdm.wg.Add(1) + go cdm.runMigration(task) + + return task, nil +} + +// runMigration выполняет миграцию +func (cdm *CrossDCMigrator) runMigration(task *MigrationTask) { + defer cdm.wg.Done() + defer func() { + if r := recover(); r != nil { + cdm.logger.Error(fmt.Sprintf("Migration %s panicked: %v", task.ID, r)) + task.mu.Lock() + task.Status = MigrationStatusFailed + task.Error = fmt.Sprintf("panic: %v", r) + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + } + }() + + cdm.logger.Info(fmt.Sprintf("Starting migration %s", task.ID)) + + // 1. Подготовка - сбор метаданных + if err := cdm.prepareMigration(task); err != nil { + cdm.setTaskError(task, err) + return + } + + // 2. Основная миграция данных + if err := cdm.migrateData(task); err != nil { + cdm.setTaskError(task, err) + return + } + + // 3. Дельта-синхронизация (если включена) + if cdm.config.Delta.Enabled && cdm.mode != MigrationModeManual { + task.mu.Lock() + task.Status = MigrationStatusDeltaSync + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + if err := cdm.runDeltaSync(task); err != nil { + cdm.logger.Warn(fmt.Sprintf("Delta sync failed: %v", err)) + // Не считаем это фатальной ошибкой + } + } + + // 4. Валидация данных + if cdm.config.Validation.Enabled { + task.mu.Lock() + task.Status = MigrationStatusValidating + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + if err := cdm.validateMigration(task); err != nil { + cdm.logger.Warn(fmt.Sprintf("Validation failed: %v", err)) + // Продолжаем, но отмечаем проблемы + } + } + + // 5. Завершение + task.mu.Lock() + task.Status = MigrationStatusCompleted + task.EndTime = time.Now().UnixMilli() + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Info(fmt.Sprintf("Migration %s completed successfully", task.ID)) +} + +// prepareMigration подготавливает миграцию +func (cdm *CrossDCMigrator) prepareMigration(task *MigrationTask) error { + cdm.logger.Debug(fmt.Sprintf("Preparing migration %s", task.ID)) + + // Получаем список баз данных + databases := task.Databases + if len(databases) == 0 { + databases = cdm.store.ListDatabases() + } + + totalDocs := int64(0) + for _, dbName := range databases { + db, err := cdm.store.GetDatabase(dbName) + if err != nil { + continue + } + + collections := task.Collections[dbName] + if len(collections) == 0 { + collections = db.ListCollections() + } + + for _, collName := range collections { + coll, err := db.GetCollection(collName) + if err != nil { + continue + } + totalDocs += coll.Count() + } + } + + task.mu.Lock() + task.TotalDocuments = totalDocs + task.Status = MigrationStatusMigrating + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Info(fmt.Sprintf("Prepared migration %s: %d documents to migrate", task.ID, totalDocs)) + return nil +} + +// migrateData выполняет основную миграцию данных +func (cdm *CrossDCMigrator) migrateData(task *MigrationTask) error { + cdm.logger.Debug(fmt.Sprintf("Migrating data for task %s", task.ID)) + + databases := task.Databases + if len(databases) == 0 { + databases = cdm.store.ListDatabases() + } + + // Проверяем наличие чекпоинта для возобновления + var checkpoint *MigrationCheckpoint + if cdm.config.Settings.ResumeEnabled { + checkpoint = cdm.loadCheckpoint(task.ID) + if checkpoint != nil { + // Преобразуем map[string]bool в map[string]interface{} для совместимости + checkpointData := make(map[string]interface{}) + for k, v := range checkpoint.ProcessedDocs { + checkpointData[k] = v + } + task.mu.Lock() + task.CheckpointData = checkpointData + task.mu.Unlock() + cdm.logger.Info(fmt.Sprintf("Resuming migration %s from checkpoint (LSN: %d)", task.ID, checkpoint.LastLSN)) + } + } + + for _, dbName := range databases { + // Проверяем остановку + if !cdm.active.Load() { + return fmt.Errorf("migration stopped") + } + + db, err := cdm.store.GetDatabase(dbName) + if err != nil { + cdm.logger.Warn(fmt.Sprintf("Database %s not found: %v", dbName, err)) + continue + } + + collections := task.Collections[dbName] + if len(collections) == 0 { + collections = db.ListCollections() + } + + // Проверяем исключения + excludeMap := make(map[string]bool) + for _, excl := range cdm.config.Settings.ExcludeCollections { + excludeMap[excl] = true + } + + for _, collName := range collections { + if excludeMap[collName] { + continue + } + + // Проверяем остановку + if !cdm.active.Load() { + return fmt.Errorf("migration stopped") + } + + // Проверяем, была ли коллекция уже обработана (по чекпоинту) + checkpointKey := fmt.Sprintf("%s:%s", dbName, collName) + if checkpoint != nil { + if done, ok := checkpoint.ProcessedDocs[checkpointKey]; ok && done { + cdm.logger.Debug(fmt.Sprintf("Collection %s.%s already migrated, skipping", dbName, collName)) + continue + } + } + + if err := cdm.migrateCollection(task, dbName, collName, checkpoint); err != nil { + cdm.logger.Error(fmt.Sprintf("Failed to migrate collection %s.%s: %v", dbName, collName, err)) + // Продолжаем с другими коллекциями + } + + // Сохраняем чекпоинт для коллекции + if cdm.config.Settings.ResumeEnabled { + if checkpoint == nil { + checkpoint = &MigrationCheckpoint{ + TaskID: task.ID, + ProcessedDocs: make(map[string]bool), + Timestamp: time.Now().UnixMilli(), + Version: 1, + } + } + checkpoint.ProcessedDocs[checkpointKey] = true + checkpoint.Timestamp = time.Now().UnixMilli() + checkpoint.Version++ + cdm.saveCheckpoint(task.ID, checkpoint) + // Преобразуем для task.CheckpointData + checkpointData := make(map[string]interface{}) + for k, v := range checkpoint.ProcessedDocs { + checkpointData[k] = v + } + task.mu.Lock() + task.CheckpointData = checkpointData + task.mu.Unlock() + } + } + } + + return nil +} + +// migrateCollection мигрирует одну коллекцию +func (cdm *CrossDCMigrator) migrateCollection(task *MigrationTask, dbName, collName string, checkpoint *MigrationCheckpoint) error { + db, err := cdm.store.GetDatabase(dbName) + if err != nil { + return err + } + + coll, err := db.GetCollection(collName) + if err != nil { + return err + } + + // Получаем все документы + docs := coll.GetAllDocuments() + if len(docs) == 0 { + return nil + } + + batchSize := cdm.config.Settings.BatchSize + workers := cdm.config.Settings.Workers + + // Создаём каналы для параллельной обработки + docChan := make(chan *storage.Document, batchSize) + resultChan := make(chan error, workers) + + // Запускаем воркеры + var wg sync.WaitGroup + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for doc := range docChan { + err := cdm.sendDocumentToTarget(task, dbName, collName, doc) + if err != nil { + resultChan <- err + task.mu.Lock() + task.FailedDocs++ + task.mu.Unlock() + } else { + task.mu.Lock() + task.MigratedDocs++ + // Обновляем прогресс + if task.TotalDocuments > 0 { + task.ProgressPercent = float64(task.MigratedDocs) / float64(task.TotalDocuments) * 100 + } + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + } + } + }() + } + + // Отправляем документы в канал + sent := 0 + for _, doc := range docs { + // Проверяем остановку + if !cdm.active.Load() { + close(docChan) + wg.Wait() + return fmt.Errorf("migration stopped") + } + + // Проверяем, был ли документ уже обработан (по чекпоинту) + if checkpoint != nil { + if done, ok := checkpoint.ProcessedDocs[doc.ID]; ok && done { + task.mu.Lock() + task.SkippedDocs++ + task.mu.Unlock() + continue + } + } + + docChan <- doc + sent++ + + // Сохраняем чекпоинт каждые N документов + if sent%batchSize == 0 && cdm.config.Settings.ResumeEnabled { + if checkpoint == nil { + checkpoint = &MigrationCheckpoint{ + TaskID: task.ID, + ProcessedDocs: make(map[string]bool), + Timestamp: time.Now().UnixMilli(), + Version: 1, + } + } + checkpoint.ProcessedDocs[doc.ID] = true + checkpoint.Timestamp = time.Now().UnixMilli() + checkpoint.Version++ + cdm.saveCheckpoint(task.ID, checkpoint) + } + } + + close(docChan) + wg.Wait() + + // Проверяем ошибки + var lastErr error + select { + case err := <-resultChan: + lastErr = err + default: + } + + // Обновляем статистику + task.mu.Lock() + task.ProgressPercent = 100.0 + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + return lastErr +} + +// sendDocumentToTarget отправляет документ на целевой датацентр +func (cdm *CrossDCMigrator) sendDocumentToTarget(task *MigrationTask, dbName, collName string, doc *storage.Document) error { + // Формируем данные документа + docData := doc.GetFields() + + // Добавляем метаданные + docData["_meta"] = map[string]interface{}{ + "source_dc": task.SourceDC, + "source_db": dbName, + "source_collection": collName, + "migration_id": task.ID, + "migrated_at": time.Now().UnixMilli(), + } + + // Создаём запись изменения + change := &ChangeRecord{ + ID: fmt.Sprintf("%s_%s_%s", task.ID, doc.ID, time.Now().Format("20060102150405")), + Database: dbName, + Collection: collName, + DocumentID: doc.ID, + ChangeType: ChangeInsert, + Document: docData, + Timestamp: time.Now().UnixMilli(), + Version: doc.Version, + Applied: false, + } + + // Вычисляем контрольную сумму + data, _ := json.Marshal(docData) + hash := sha256.Sum256(data) + change.Checksum = hex.EncodeToString(hash[:]) + + // Добавляем в очередь изменений + cdm.changeQueue.Add(change) + + // Пытаемся отправить напрямую + targetEndpoint := cdm.config.Target.Endpoint + if targetEndpoint == "" { + return fmt.Errorf("target endpoint not configured") + } + + // Формируем запрос + reqData := map[string]interface{}{ + "database": dbName, + "collection": collName, + "document": docData, + "migration": map[string]interface{}{ + "id": task.ID, + "source": task.SourceDC, + }, + } + + jsonData, err := json.Marshal(reqData) + if err != nil { + return err + } + + // Сжимаем данные + compressedData, err := cdm.compressData(jsonData) + if err != nil { + return err + } + + // Отправляем HTTP запрос + url := fmt.Sprintf("%s/api/migration/receive", targetEndpoint) + req, err := http.NewRequest("POST", url, bytes.NewReader(compressedData)) + if err != nil { + return err + } + + req.Header.Set("Content-Type", "application/octet-stream") + req.Header.Set("Content-Encoding", cdm.config.Settings.Compression) + req.Header.Set("X-Migration-ID", task.ID) + + // Отправляем с повторными попытками + var lastErr error + for retry := 0; retry < cdm.config.Settings.MaxRetries; retry++ { + resp, err := cdm.httpClient.Do(req) + if err == nil { + resp.Body.Close() + if resp.StatusCode == http.StatusOK { + change.Applied = true + change.AppliedAt = time.Now().UnixMilli() + return nil + } + lastErr = fmt.Errorf("status: %s", resp.Status) + } else { + lastErr = err + } + + time.Sleep(time.Duration(cdm.config.Settings.RetryBackoffSec) * time.Second) + } + + return lastErr +} + +// ============================================================================= +// ДЕЛЬТА-СИНХРОНИЗАЦИЯ +// ============================================================================= + +// runDeltaSync выполняет дельта-синхронизацию +func (cdm *CrossDCMigrator) runDeltaSync(task *MigrationTask) error { + cdm.logger.Info(fmt.Sprintf("Starting delta sync for migration %s", task.ID)) + + // Получаем изменения после последнего чекпоинта + since := task.LastCheckpoint + if since == 0 { + since = task.StartTime + } + + // Собираем изменения из очереди + changes := cdm.changeQueue.GetChangesAfter(since) + + batchSize := cdm.config.Settings.BatchSize + for i := 0; i < len(changes); i += batchSize { + end := i + batchSize + if end > len(changes) { + end = len(changes) + } + + batch := changes[i:end] + if err := cdm.applyChangeBatch(task, batch); err != nil { + cdm.logger.Error(fmt.Sprintf("Failed to apply change batch: %v", err)) + return err + } + + // Обновляем прогресс + task.mu.Lock() + task.MigratedDocs += int64(len(batch)) + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + } + + task.mu.Lock() + task.LastCheckpoint = time.Now().UnixMilli() + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Info(fmt.Sprintf("Delta sync completed for migration %s: %d changes applied", task.ID, len(changes))) + return nil +} + +// applyChangeBatch применяет пакет изменений +func (cdm *CrossDCMigrator) applyChangeBatch(task *MigrationTask, changes []*ChangeRecord) error { + for _, change := range changes { + if !cdm.active.Load() { + return fmt.Errorf("migration stopped") + } + + // Применяем изменение в целевом датацентре + if err := cdm.applyChange(task, change); err != nil { + cdm.logger.Error(fmt.Sprintf("Failed to apply change %s: %v", change.ID, err)) + continue + } + } + return nil +} + +// applyChange применяет одно изменение +func (cdm *CrossDCMigrator) applyChange(task *MigrationTask, change *ChangeRecord) error { + // Формируем запрос к целевому датацентру + reqData := map[string]interface{}{ + "database": change.Database, + "collection": change.Collection, + "document_id": change.DocumentID, + "change_type": string(change.ChangeType), + "document": change.Document, + "previous_doc": change.PreviousDoc, + "migration_id": task.ID, + "timestamp": change.Timestamp, + } + + jsonData, err := json.Marshal(reqData) + if err != nil { + return err + } + + compressedData, err := cdm.compressData(jsonData) + if err != nil { + return err + } + + url := fmt.Sprintf("%s/api/migration/apply", cdm.config.Target.Endpoint) + req, err := http.NewRequest("POST", url, bytes.NewReader(compressedData)) + if err != nil { + return err + } + + req.Header.Set("Content-Type", "application/octet-stream") + req.Header.Set("Content-Encoding", cdm.config.Settings.Compression) + + var lastErr error + for retry := 0; retry < cdm.config.Settings.MaxRetries; retry++ { + resp, err := cdm.httpClient.Do(req) + if err == nil { + resp.Body.Close() + if resp.StatusCode == http.StatusOK { + change.Applied = true + change.AppliedAt = time.Now().UnixMilli() + cdm.stats.mu.Lock() + cdm.stats.AppliedChanges++ + cdm.stats.mu.Unlock() + return nil + } + lastErr = fmt.Errorf("status: %s", resp.Status) + } else { + lastErr = err + } + + time.Sleep(time.Duration(cdm.config.Settings.RetryBackoffSec) * time.Second) + } + + cdm.stats.mu.Lock() + cdm.stats.FailedChanges++ + cdm.stats.mu.Unlock() + + return lastErr +} + +// ============================================================================= +// ВАЛИДАЦИЯ +// ============================================================================= + +// validateMigration выполняет валидацию мигрированных данных +func (cdm *CrossDCMigrator) validateMigration(task *MigrationTask) error { + cdm.logger.Info(fmt.Sprintf("Validating migration %s", task.ID)) + + if cdm.config.Validation.SamplePercent <= 0 { + return nil + } + + databases := task.Databases + if len(databases) == 0 { + databases = cdm.store.ListDatabases() + } + + errors := 0 + + for _, dbName := range databases { + db, err := cdm.store.GetDatabase(dbName) + if err != nil { + continue + } + + collections := task.Collections[dbName] + if len(collections) == 0 { + collections = db.ListCollections() + } + + for _, collName := range collections { + coll, err := db.GetCollection(collName) + if err != nil { + continue + } + + docs := coll.GetAllDocuments() + sampleSize := len(docs) * cdm.config.Validation.SamplePercent / 100 + if sampleSize < 1 { + sampleSize = 1 + } + + for i := 0; i < sampleSize && i < len(docs); i++ { + doc := docs[i] + if err := cdm.validateDocument(task, dbName, collName, doc); err != nil { + errors++ + if errors >= cdm.config.Validation.MaxErrors { + return fmt.Errorf("validation stopped: too many errors (%d)", errors) + } + } + } + } + } + + cdm.logger.Info(fmt.Sprintf("Validation completed for %s: %d errors", task.ID, errors)) + return nil +} + +// validateDocument валидирует один документ +func (cdm *CrossDCMigrator) validateDocument(task *MigrationTask, dbName, collName string, doc *storage.Document) error { + docData := doc.GetFields() + data, _ := json.Marshal(docData) + hash := sha256.Sum256(data) + checksum := hex.EncodeToString(hash[:]) + + // Запрашиваем документ из целевого датацентра + url := fmt.Sprintf("%s/api/migration/validate?database=%s&collection=%s&document_id=%s", + cdm.config.Target.Endpoint, dbName, collName, doc.ID) + + resp, err := cdm.httpClient.Get(url) + if err != nil { + return fmt.Errorf("validation request failed: %v", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("target document not found: %s", doc.ID) + } + + var result map[string]interface{} + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return fmt.Errorf("failed to decode validation response: %v", err) + } + + targetChecksum, ok := result["checksum"].(string) + if !ok { + return fmt.Errorf("invalid checksum format") + } + + if targetChecksum != checksum { + return fmt.Errorf("checksum mismatch for %s: source=%s, target=%s", doc.ID, checksum, targetChecksum) + } + + return nil +} + +// ============================================================================= +// ЧЕКПОИНТЫ +// ============================================================================= + +// saveCheckpoint сохраняет чекпоинт миграции +func (cdm *CrossDCMigrator) saveCheckpoint(taskID string, checkpoint *MigrationCheckpoint) { + if !cdm.config.Settings.ResumeEnabled { + return + } + + cdm.mu.Lock() + defer cdm.mu.Unlock() + + path := filepath.Join(cdm.checkpointPath, fmt.Sprintf("%s.json", taskID)) + + checkpoint.Timestamp = time.Now().UnixMilli() + checkpoint.Version++ + + jsonData, err := json.MarshalIndent(checkpoint, "", " ") + if err != nil { + cdm.logger.Warn(fmt.Sprintf("Failed to marshal checkpoint: %v", err)) + return + } + + if err := os.WriteFile(path, jsonData, 0644); err != nil { + cdm.logger.Warn(fmt.Sprintf("Failed to save checkpoint: %v", err)) + return + } + + cdm.logger.Debug(fmt.Sprintf("Checkpoint saved for task %s (LSN: %d, docs: %d)", + taskID, checkpoint.LastLSN, len(checkpoint.ProcessedDocs))) +} + +// loadCheckpoint загружает чекпоинт миграции +func (cdm *CrossDCMigrator) loadCheckpoint(taskID string) *MigrationCheckpoint { + path := filepath.Join(cdm.checkpointPath, fmt.Sprintf("%s.json", taskID)) + + data, err := os.ReadFile(path) + if err != nil { + if !os.IsNotExist(err) { + cdm.logger.Warn(fmt.Sprintf("Failed to read checkpoint: %v", err)) + } + return nil + } + + var checkpoint MigrationCheckpoint + if err := json.Unmarshal(data, &checkpoint); err != nil { + cdm.logger.Warn(fmt.Sprintf("Failed to unmarshal checkpoint: %v", err)) + return nil + } + + cdm.logger.Debug(fmt.Sprintf("Checkpoint loaded for task %s (LSN: %d, docs: %d)", + taskID, checkpoint.LastLSN, len(checkpoint.ProcessedDocs))) + return &checkpoint +} + +// ============================================================================= +// ДЕЛЬТА-СИНХРОНИЗАЦИЯ В ЦИКЛЕ +// ============================================================================= + +// deltaSyncLoop выполняет периодическую дельта-синхронизацию +func (cdm *CrossDCMigrator) deltaSyncLoop() { + defer cdm.wg.Done() + + interval := time.Duration(cdm.config.Delta.IntervalSec) * time.Second + ticker := time.NewTicker(interval) + defer ticker.Stop() + + for { + select { + case <-cdm.stopChan: + return + case <-ticker.C: + if !cdm.active.Load() { + continue + } + cdm.runPeriodicDeltaSync() + } + } +} + +// runPeriodicDeltaSync выполняет периодическую дельта-синхронизацию +func (cdm *CrossDCMigrator) runPeriodicDeltaSync() { + cdm.mu.RLock() + taskID := cdm.currentTaskID + cdm.mu.RUnlock() + + if taskID == "" { + return + } + + task, ok := cdm.tasks[taskID] + if !ok { + return + } + + task.mu.RLock() + status := task.Status + task.mu.RUnlock() + + if status != MigrationStatusCompleted && status != MigrationStatusDeltaSync { + return + } + + cdm.logger.Debug("Running periodic delta sync") + + if err := cdm.runDeltaSync(task); err != nil { + cdm.logger.Warn(fmt.Sprintf("Periodic delta sync failed: %v", err)) + } +} + +// ============================================================================= +// ОБРАБОТКА ИЗМЕНЕНИЙ В ОЧЕРЕДИ +// ============================================================================= + +// processChanges обрабатывает изменения из очереди +func (cdm *CrossDCMigrator) processChanges() { + defer cdm.wg.Done() + + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + + for { + select { + case <-cdm.stopChan: + return + case <-ticker.C: + if !cdm.active.Load() { + continue + } + cdm.flushChangeQueue() + } + } +} + +// flushChangeQueue сбрасывает очередь изменений +func (cdm *CrossDCMigrator) flushChangeQueue() { + changes := cdm.changeQueue.GetAll() + if len(changes) == 0 { + return + } + + cdm.mu.RLock() + taskID := cdm.currentTaskID + cdm.mu.RUnlock() + + if taskID == "" { + cdm.logger.Warn("No active migration, dropping changes") + cdm.changeQueue.Clear() + return + } + + task, ok := cdm.tasks[taskID] + if !ok { + cdm.logger.Warn("Migration task not found, dropping changes") + cdm.changeQueue.Clear() + return + } + + // Применяем изменения + for _, change := range changes { + if !cdm.active.Load() { + return + } + + if err := cdm.applyChange(task, change); err != nil { + cdm.logger.Warn(fmt.Sprintf("Failed to apply change: %v", err)) + } + } + + cdm.logger.Debug(fmt.Sprintf("Flushed %d changes to target", len(changes))) +} + +// ============================================================================= +// ВСПОМОГАТЕЛЬНЫЕ МЕТОДЫ +// ============================================================================= + +// compressData сжимает данные +func (cdm *CrossDCMigrator) compressData(data []byte) ([]byte, error) { + var buf bytes.Buffer + + switch cdm.config.Settings.Compression { + case "gzip", "": + w := gzip.NewWriter(&buf) + defer w.Close() + if _, err := w.Write(data); err != nil { + return nil, err + } + if err := w.Flush(); err != nil { + return nil, err + } + default: + return data, nil + } + + return buf.Bytes(), nil +} + +// setTaskError устанавливает ошибку задачи +func (cdm *CrossDCMigrator) setTaskError(task *MigrationTask, err error) { + task.mu.Lock() + task.Status = MigrationStatusFailed + task.Error = err.Error() + task.EndTime = time.Now().UnixMilli() + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Error(fmt.Sprintf("Migration %s failed: %v", task.ID, err)) +} + +// loadTasks загружает сохранённые задачи +func (cdm *CrossDCMigrator) loadTasks() { + path := "futriis/migration/tasks.json" + data, err := os.ReadFile(path) + if err != nil { + return + } + + var tasks map[string]*MigrationTask + if err := json.Unmarshal(data, &tasks); err != nil { + return + } + + cdm.mu.Lock() + for id, task := range tasks { + cdm.tasks[id] = task + } + cdm.mu.Unlock() + + cdm.logger.Debug(fmt.Sprintf("Loaded %d migration tasks", len(tasks))) +} + +// saveTasks сохраняет задачи +func (cdm *CrossDCMigrator) saveTasks() { + cdm.mu.RLock() + tasks := make(map[string]*MigrationTask) + for id, task := range cdm.tasks { + tasks[id] = task + } + cdm.mu.RUnlock() + + data, err := json.MarshalIndent(tasks, "", " ") + if err != nil { + cdm.logger.Warn(fmt.Sprintf("Failed to marshal tasks: %v", err)) + return + } + + if err := os.WriteFile("futriis/migration/tasks.json", data, 0644); err != nil { + cdm.logger.Warn(fmt.Sprintf("Failed to save tasks: %v", err)) + } +} + +// resumePendingMigrations восстанавливает незавершённые миграции +func (cdm *CrossDCMigrator) resumePendingMigrations() { + cdm.mu.RLock() + var pending []*MigrationTask + for _, task := range cdm.tasks { + task.mu.RLock() + status := task.Status + task.mu.RUnlock() + if status == MigrationStatusMigrating || status == MigrationStatusDeltaSync { + pending = append(pending, task) + } + } + cdm.mu.RUnlock() + + for _, task := range pending { + if cdm.config.Settings.ResumeEnabled { + cdm.logger.Info(fmt.Sprintf("Resuming migration %s", task.ID)) + cdm.wg.Add(1) + go cdm.runMigration(task) + } + } +} + +// checkpointLoop периодически сохраняет чекпоинты +func (cdm *CrossDCMigrator) checkpointLoop() { + defer cdm.wg.Done() + + interval := time.Duration(cdm.config.Settings.CheckpointIntervalSec) * time.Second + ticker := time.NewTicker(interval) + defer ticker.Stop() + + for { + select { + case <-cdm.stopChan: + return + case <-ticker.C: + if !cdm.active.Load() { + continue + } + cdm.mu.RLock() + taskID := cdm.currentTaskID + cdm.mu.RUnlock() + + if taskID == "" { + continue + } + + task, ok := cdm.tasks[taskID] + if !ok { + continue + } + + task.mu.RLock() + status := task.Status + checkpointData := task.CheckpointData + task.mu.RUnlock() + + if status != MigrationStatusMigrating && status != MigrationStatusDeltaSync { + continue + } + + // Сохраняем чекпоинт + checkpoint := &MigrationCheckpoint{ + TaskID: taskID, + LastLSN: cdm.changeQueue.GetLastLSN(), + ProcessedDocs: make(map[string]bool), + Timestamp: time.Now().UnixMilli(), + Version: 1, + } + + // Конвертируем checkpointData в map[string]bool + for k, v := range checkpointData { + if b, ok := v.(bool); ok { + checkpoint.ProcessedDocs[k] = b + } + } + + cdm.saveCheckpoint(taskID, checkpoint) + } + } +} + +// ============================================================================= +// ПУБЛИЧНЫЕ МЕТОДЫ ДЛЯ REPL +// ============================================================================= + +// GetMigrationStatus возвращает статус миграции +func (cdm *CrossDCMigrator) GetMigrationStatus(taskID string) (*MigrationTask, error) { + cdm.mu.RLock() + defer cdm.mu.RUnlock() + + task, ok := cdm.tasks[taskID] + if !ok { + return nil, fmt.Errorf("task %s not found", taskID) + } + + return task, nil +} + +// GetCurrentTaskID возвращает ID текущей задачи +func (cdm *CrossDCMigrator) GetCurrentTaskID() string { + cdm.mu.RLock() + defer cdm.mu.RUnlock() + return cdm.currentTaskID +} + +// ListTasks возвращает список всех задач +func (cdm *CrossDCMigrator) ListTasks() []*MigrationTask { + cdm.mu.RLock() + defer cdm.mu.RUnlock() + + tasks := make([]*MigrationTask, 0, len(cdm.tasks)) + for _, task := range cdm.tasks { + tasks = append(tasks, task) + } + + // Сортируем по времени создания (новые сверху) + sort.Slice(tasks, func(i, j int) bool { + return tasks[i].CreatedAt > tasks[j].CreatedAt + }) + + return tasks +} + +// PauseMigration приостанавливает миграцию +func (cdm *CrossDCMigrator) PauseMigration(taskID string) error { + cdm.mu.RLock() + task, ok := cdm.tasks[taskID] + cdm.mu.RUnlock() + + if !ok { + return fmt.Errorf("task %s not found", taskID) + } + + task.mu.Lock() + if task.Status != MigrationStatusMigrating && task.Status != MigrationStatusDeltaSync { + task.mu.Unlock() + return fmt.Errorf("task %s is not in a running state", taskID) + } + task.Status = MigrationStatusPaused + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Info(fmt.Sprintf("Migration %s paused", taskID)) + return nil +} + +// ResumeMigration возобновляет миграцию +func (cdm *CrossDCMigrator) ResumeMigration(taskID string) error { + cdm.mu.RLock() + task, ok := cdm.tasks[taskID] + cdm.mu.RUnlock() + + if !ok { + return fmt.Errorf("task %s not found", taskID) + } + + task.mu.Lock() + if task.Status != MigrationStatusPaused { + task.mu.Unlock() + return fmt.Errorf("task %s is not paused", taskID) + } + task.Status = MigrationStatusMigrating + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Info(fmt.Sprintf("Migration %s resumed", taskID)) + cdm.wg.Add(1) + go cdm.runMigration(task) + + return nil +} + +// CancelMigration отменяет миграцию +func (cdm *CrossDCMigrator) CancelMigration(taskID string) error { + cdm.mu.RLock() + task, ok := cdm.tasks[taskID] + cdm.mu.RUnlock() + + if !ok { + return fmt.Errorf("task %s not found", taskID) + } + + task.mu.Lock() + task.Status = MigrationStatusFailed + task.Error = "cancelled by user" + task.EndTime = time.Now().UnixMilli() + task.UpdatedAt = time.Now().UnixMilli() + task.mu.Unlock() + + cdm.logger.Info(fmt.Sprintf("Migration %s cancelled", taskID)) + return nil +} + +// GetMigrationStats возвращает статистику миграции +func (cdm *CrossDCMigrator) GetMigrationStats() *MigrationStats { + cdm.stats.mu.RLock() + defer cdm.stats.mu.RUnlock() + return &MigrationStats{ + TotalChanges: cdm.stats.TotalChanges, + AppliedChanges: cdm.stats.AppliedChanges, + FailedChanges: cdm.stats.FailedChanges, + SkippedChanges: cdm.stats.SkippedChanges, + LatencyAvg: cdm.stats.LatencyAvg, + LatencyMax: cdm.stats.LatencyMax, + Throughput: cdm.stats.Throughput, + } +} + +// SubscribeChanges подписывается на изменения +func (cdm *CrossDCMigrator) SubscribeChanges() <-chan *ChangeRecord { + ch := make(chan *ChangeRecord, 1000) + cdm.subscriberMu.Lock() + cdm.changeSubscribers = append(cdm.changeSubscribers, ch) + cdm.subscriberMu.Unlock() + return ch +} + +// UnsubscribeChanges отписывается от изменений +func (cdm *CrossDCMigrator) UnsubscribeChanges(ch <-chan *ChangeRecord) { + cdm.subscriberMu.Lock() + defer cdm.subscriberMu.Unlock() + + for i, sub := range cdm.changeSubscribers { + if sub == ch { + cdm.changeSubscribers = append(cdm.changeSubscribers[:i], cdm.changeSubscribers[i+1:]...) + close(sub) + break + } + } +} + +// GetQueueStats возвращает статистику очереди +func (cdm *CrossDCMigrator) GetQueueStats() map[string]interface{} { + return map[string]interface{}{ + "size": cdm.changeQueue.Size(), + "last_lsn": cdm.changeQueue.GetLastLSN(), + "max_size": 100000, + "file_path": cdm.changeQueue.filePath, + } +} + +// GetConfig возвращает конфигурацию мигратора (публичный метод для доступа из REPL) +func (cdm *CrossDCMigrator) GetConfig() *config.MigrationConfig { + return cdm.config +}