qcc/internal/store/engine.go
2026-06-30 12:52:09 +03:00

170 lines
3.4 KiB
Go

package store
import (
"fmt"
"os"
"path/filepath"
"sync"
)
type Engine struct {
mu sync.RWMutex
DataDir string
Tables map[string]*Table
frames map[string]*frameState
}
type frameState struct {
lastFrameID uint64
lastHash [32]byte
frameCount int
}
func NewEngine(dataDir string) (*Engine, error) {
dirs := []string{
dataDir,
filepath.Join(dataDir, "identities"),
filepath.Join(dataDir, "calls"),
filepath.Join(dataDir, "indexes"),
filepath.Join(dataDir, "acme"),
}
for _, d := range dirs {
if err := os.MkdirAll(d, 0700); err != nil {
return nil, fmt.Errorf("mkdir %s: %w", d, err)
}
}
e := &Engine{
DataDir: dataDir,
Tables: make(map[string]*Table),
frames: make(map[string]*frameState),
}
if err := e.initDB(); err != nil {
return nil, err
}
return e, nil
}
func (e *Engine) initDB() error {
path := filepath.Join(e.DataDir, "qcc.db")
if _, err := os.Stat(path); os.IsNotExist(err) {
f, err := os.Create(path)
if err != nil {
return err
}
defer f.Close()
header := make([]byte, 14)
copy(header[0:4], []byte("QCC\x00"))
header[4] = 1
header[5] = 0
_, err = f.Write(header)
return err
}
return nil
}
func (e *Engine) GetTable(name string) (*Table, error) {
e.mu.Lock()
defer e.mu.Unlock()
if t, ok := e.Tables[name]; ok {
return t, nil
}
fs, ok := e.frames[name]
if !ok {
fs = &frameState{}
e.frames[name] = fs
lastFrame, lastHash, count := e.scanFrames(name)
fs.lastFrameID = lastFrame
fs.lastHash = lastHash
fs.frameCount = count
}
tablePath := filepath.Join(e.DataDir, name)
t := &Table{
engine: e,
name: name,
path: tablePath,
lastHash: fs.lastHash,
nextID: fs.lastFrameID + 1,
frameCount: fs.frameCount,
index: make(map[uint64]IndexEntry),
}
e.Tables[name] = t
if err := t.loadLastIndex(); err != nil {
return nil, err
}
return t, nil
}
func (e *Engine) scanFrames(name string) (uint64, [32]byte, int) {
dir := filepath.Join(e.DataDir, name)
entries, err := os.ReadDir(dir)
if err != nil {
return 0, [32]byte{}, 0
}
var lastID uint64
var lastHash [32]byte
count := 0
var prevFrameID uint64
var prevHash [32]byte
for _, entry := range entries {
if filepath.Ext(entry.Name()) != ".frm" {
continue
}
data, err := os.ReadFile(filepath.Join(dir, entry.Name()))
if err != nil {
continue
}
if len(data) < 66 {
continue
}
storedChecksum := data[len(data)-32:]
computedChecksum := ComputeFrameChecksum(data[:len(data)-32])
if string(storedChecksum) != string(computedChecksum[:]) {
continue
}
f, err := UnmarshalFrame(data)
if err != nil {
continue
}
if count > 0 {
if f.Header.PrevFrameID != prevFrameID ||
string(f.Header.PrevHash[:]) != string(prevHash[:]) {
continue
}
}
if f.Header.FrameID > lastID {
lastID = f.Header.FrameID
lastHash = computedChecksum
}
prevFrameID = f.Header.FrameID
prevHash = computedChecksum
count++
}
return lastID, lastHash, count
}
func (e *Engine) UpdateFrameState(name string, frameID uint64, hash [32]byte) {
e.mu.Lock()
defer e.mu.Unlock()
if fs, ok := e.frames[name]; ok {
if frameID > fs.lastFrameID {
fs.lastFrameID = frameID
fs.lastHash = hash
fs.frameCount++
}
}
}
func (e *Engine) Close() error {
e.mu.Lock()
defer e.mu.Unlock()
for _, t := range e.Tables {
if err := t.flushIndex(); err != nil {
return err
}
}
return nil
}