diff --git a/internal/storage/aof.go b/internal/storage/aof.go new file mode 100644 index 0000000..3b37104 --- /dev/null +++ b/internal/storage/aof.go @@ -0,0 +1,481 @@ +/* + * 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) +}