Upload files to "internal/storage"
This commit is contained in:
1 parent
ab48d50988
commit
5e230b1a9c
1 file changed
+481
@@ -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)
|
||||
}
|
||||
Reference in new issue
Block a user