/* * 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 }