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