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 }