Files
futriix/internal/storage/aof.go
T

482 lines
13 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/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)
}