Files
futriix/internal/cluster/crossdc_migrator.go

1610 lines
51 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/*
* 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
}