/* * 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/storage/aof.go // Назначение: Append-Only File (AOF) для мутаций, выполненных вне batch. // // ЗАЧЕМ НУЖЕН AOF: // WAL покрывает только batch-операции (групповые мутации). // Прямые вызовы coll.Insert / coll.Update / coll.Delete идут мимо WAL. // AOF логирует каждую отдельную мутацию и позволяет: // 1. Восстановить состояние после краха без потери данных. // 2. Вести полный audit trail. // 3. Использовать для репликации. // // ФОРМАТ: // Каждая запись — JSON-строка, заканчивающаяся '\n'. // Поля: type, database, collection, document_id, data, timestamp, crc. // Файлы разбиваются по размеру (по умолчанию 64 МБ). // При старте AOFManager сканирует директорию и загружает все записи. package storage import ( "bufio" "encoding/json" "fmt" "os" "path/filepath" "sort" "strings" "sync" "time" ) // ============================================================================= // КОНСТАНТЫ // ============================================================================= const ( // Максимальный размер одного AOF-файла. AOFMaxFileSize = 64 * 1024 * 1024 // Префикс и суффикс AOF-файлов. AOFFilePrefix = "aof_" AOFFileSuffix = ".log" // Размер буфера для bufio.Writer. AOFBufferSize = 64 * 1024 ) // ============================================================================= // ТИПЫ // ============================================================================= // AOFOperation — одна операция в AOF. type AOFOperation struct { Type string `json:"type"` // insert, update, delete, restore Database string `json:"database"` Collection string `json:"collection"` DocumentID string `json:"document_id"` Data map[string]interface{} `json:"data,omitempty"` Timestamp int64 `json:"timestamp"` CRC uint32 `json:"crc,omitempty"` } // AOFManager управляет append-only file. type AOFManager struct { mu sync.Mutex dir string currentFile *os.File writer *bufio.Writer currentSize int64 fileIndex uint32 fsyncEnabled bool logger LoggerInterface closed bool } // NewAOFManager создаёт AOF-менеджер. func NewAOFManager(dir string, fsyncEnabled bool, logger LoggerInterface) (*AOFManager, error) { if err := os.MkdirAll(dir, 0755); err != nil { return nil, fmt.Errorf("failed to create AOF directory: %v", err) } aof := &AOFManager{ dir: dir, fsyncEnabled: fsyncEnabled, logger: logger, } // Находим максимальный индекс файла. files, err := aof.listFiles() if err != nil { return nil, err } if len(files) > 0 { lastFile := files[len(files)-1] var idx uint32 if _, err := fmt.Sscanf(filepath.Base(lastFile), AOFFilePrefix+"%d"+AOFFileSuffix, &idx); err == nil { aof.fileIndex = idx } } // Открываем текущий файл на дозапись. if err := aof.openCurrentFile(); err != nil { return nil, err } return aof, nil } // SetLogger устанавливает логгер. func (aof *AOFManager) SetLogger(logger LoggerInterface) { aof.mu.Lock() defer aof.mu.Unlock() aof.logger = logger } // listFiles возвращает отсортированный список AOF-файлов. func (aof *AOFManager) listFiles() ([]string, error) { pattern := filepath.Join(aof.dir, AOFFilePrefix+"*"+AOFFileSuffix) files, err := filepath.Glob(pattern) if err != nil { return nil, err } sort.Strings(files) return files, nil } // openCurrentFile открывает текущий файл для дозаписи. func (aof *AOFManager) openCurrentFile() error { path := filepath.Join(aof.dir, fmt.Sprintf("%s%d%s", AOFFilePrefix, aof.fileIndex, AOFFileSuffix)) f, err := os.OpenFile(path, os.O_CREATE|os.O_APPEND|os.O_RDWR, 0644) if err != nil { return fmt.Errorf("failed to open AOF file: %v", err) } stat, err := f.Stat() if err != nil { f.Close() return err } aof.currentFile = f aof.writer = bufio.NewWriterSize(f, AOFBufferSize) aof.currentSize = stat.Size() return nil } // rotateIfNeeded проверяет размер и ротирует файл. func (aof *AOFManager) rotateIfNeeded() error { if aof.currentSize < AOFMaxFileSize { return nil } if err := aof.writer.Flush(); err != nil { return err } if aof.fsyncEnabled { RealFsync(aof.currentFile) } if err := aof.currentFile.Close(); err != nil { return err } aof.fileIndex++ return aof.openCurrentFile() } // Append записывает одну операцию в AOF. func (aof *AOFManager) Append(op *AOFOperation) error { aof.mu.Lock() defer aof.mu.Unlock() if aof.closed { return fmt.Errorf("AOF is closed") } op.Timestamp = time.Now().UnixMilli() // Вычисляем CRC. data, err := json.Marshal(op) if err != nil { return err } op.CRC = crc32(data) data, err = json.Marshal(op) if err != nil { return err } data = append(data, '\n') if _, err := aof.writer.Write(data); err != nil { return err } aof.currentSize += int64(len(data)) if err := aof.writer.Flush(); err != nil { return err } if aof.fsyncEnabled { RealFsync(aof.currentFile) } return aof.rotateIfNeeded() } // AppendBatch записывает все операции batch в AOF. func (aof *AOFManager) AppendBatch(batch *Batch) error { for _, op := range batch.Operations { aofOp := &AOFOperation{ Type: op.Type, Database: op.Database, Collection: op.Collection, DocumentID: op.DocumentID, Data: op.Data, } if err := aof.Append(aofOp); err != nil { return err } } return nil } // ReadAll читает все операции из всех AOF-файлов. func (aof *AOFManager) ReadAll() ([]*AOFOperation, error) { aof.mu.Lock() defer aof.mu.Unlock() files, err := aof.listFiles() if err != nil { return nil, err } ops := make([]*AOFOperation, 0) for _, file := range files { fileOps, err := aof.readFile(file) if err != nil { if aof.logger != nil { aof.logger.Warn(fmt.Sprintf("Failed to read AOF file %s: %v", file, err)) } continue } ops = append(ops, fileOps...) } return ops, nil } // readFile читает операции из одного файла. func (aof *AOFManager) readFile(path string) ([]*AOFOperation, error) { f, err := os.Open(path) if err != nil { return nil, err } defer f.Close() ops := make([]*AOFOperation, 0) scanner := bufio.NewScanner(f) scanner.Buffer(make([]byte, 0, 64*1024), 16*1024*1024) for scanner.Scan() { line := scanner.Bytes() if len(line) == 0 { continue } var op AOFOperation if err := json.Unmarshal(line, &op); err != nil { continue } // Проверяем CRC. expectedCRC := op.CRC op.CRC = 0 data, _ := json.Marshal(&op) calculatedCRC := crc32(data) op.CRC = expectedCRC if calculatedCRC != expectedCRC { continue } ops = append(ops, &op) } return ops, scanner.Err() } // Replay воспроизводит все операции AOF, применяя их к хранилищу. // // Используется при старте для восстановления состояния после краха. func (aof *AOFManager) Replay(storage *Storage) error { ops, err := aof.ReadAll() if err != nil { return err } if len(ops) == 0 { return nil } if aof.logger != nil { aof.logger.Info(fmt.Sprintf("Replaying %d AOF operations", len(ops))) } for _, op := range ops { db, err := storage.GetDatabase(op.Database) if err != nil { // База данных могла быть создана позже — пропускаем. continue } coll, err := db.GetCollection(op.Collection) if err != nil { continue } switch op.Type { case "insert": doc := NewDocumentWithID(op.DocumentID) for k, v := range op.Data { doc.SetField(k, v) } doc.CreatedAt = op.Timestamp doc.UpdatedAt = op.Timestamp coll.Insert(doc) case "update": coll.Update(op.DocumentID, op.Data) case "delete": coll.Delete(op.DocumentID) case "restore": coll.RestoreDeleted(op.DocumentID) } } return nil } // Close закрывает AOF-менеджер. func (aof *AOFManager) Close() error { aof.mu.Lock() defer aof.mu.Unlock() if aof.closed { return nil } aof.closed = true if aof.writer != nil { if err := aof.writer.Flush(); err != nil { return err } } if aof.currentFile != nil { if aof.fsyncEnabled { RealFsync(aof.currentFile) } return aof.currentFile.Close() } return nil } // GetStats возвращает статистику AOF. func (aof *AOFManager) GetStats() map[string]interface{} { aof.mu.Lock() defer aof.mu.Unlock() files, _ := aof.listFiles() var totalSize int64 for _, f := range files { if stat, err := os.Stat(f); err == nil { totalSize += stat.Size() } } return map[string]interface{}{ "file_count": len(files), "total_size": totalSize, "current_file": fmt.Sprintf("%s%d%s", AOFFilePrefix, aof.fileIndex, AOFFileSuffix), "current_size": aof.currentSize, "fsync_enabled": aof.fsyncEnabled, "dir": aof.dir, } } // ============================================================================= // ИНТЕГРАЦИЯ С COLLECTION // ============================================================================= // LogAOFInsert логирует вставку документа. func LogAOFInsert(dbName, collName string, doc *Document) { if globalBatchManager == nil || globalBatchManager.aof == nil { return } op := &AOFOperation{ Type: "insert", Database: dbName, Collection: collName, DocumentID: doc.ID, Data: doc.GetFields(), } if err := globalBatchManager.aof.Append(op); err != nil { if globalBatchManager.logger != nil { globalBatchManager.logger.Warn(fmt.Sprintf("AOF insert failed: %v", err)) } } } // LogAOFUpdate логирует обновление документа. func LogAOFUpdate(dbName, collName, docID string, updates map[string]interface{}) { if globalBatchManager == nil || globalBatchManager.aof == nil { return } op := &AOFOperation{ Type: "update", Database: dbName, Collection: collName, DocumentID: docID, Data: updates, } if err := globalBatchManager.aof.Append(op); err != nil { if globalBatchManager.logger != nil { globalBatchManager.logger.Warn(fmt.Sprintf("AOF update failed: %v", err)) } } } // LogAOFDelete логирует удаление документа. func LogAOFDelete(dbName, collName, docID string) { if globalBatchManager == nil || globalBatchManager.aof == nil { return } op := &AOFOperation{ Type: "delete", Database: dbName, Collection: collName, DocumentID: docID, } if err := globalBatchManager.aof.Append(op); err != nil { if globalBatchManager.logger != nil { globalBatchManager.logger.Warn(fmt.Sprintf("AOF delete failed: %v", err)) } } } // LogAOFRestore логирует восстановление документа. func LogAOFRestore(dbName, collName, docID string) { if globalBatchManager == nil || globalBatchManager.aof == nil { return } op := &AOFOperation{ Type: "restore", Database: dbName, Collection: collName, DocumentID: docID, } if err := globalBatchManager.aof.Append(op); err != nil { if globalBatchManager.logger != nil { globalBatchManager.logger.Warn(fmt.Sprintf("AOF restore failed: %v", err)) } } } // GetAOFStats возвращает статистику AOF. func GetAOFStats() map[string]interface{} { if globalBatchManager == nil || globalBatchManager.aof == nil { return map[string]interface{}{"error": "AOF not initialized"} } return globalBatchManager.aof.GetStats() } // ReplayAOF воспроизводит AOF при старте. func ReplayAOF(storage *Storage) error { if globalBatchManager == nil || globalBatchManager.aof == nil { return nil } return globalBatchManager.aof.Replay(storage) } // ============================================================================= // ВСПОМОГАТЕЛЬНЫЕ // ============================================================================= // trimAOFExtension убирает суффикс из имени файла. func trimAOFExtension(name string) string { return strings.TrimSuffix(name, AOFFileSuffix) }