- PoW (Hashcash) authentication - X.509 CA for phone number certificates (+0 XXX YYY ZZZ) - QUIC transport (hysteria quic-go fork) with ACME TLS - Custom append-only DB engine (from NikoGram) - Call signaling (dial/ring/accept/reject/end) - E2EE media relay (X25519 + ChaCha20-Poly1305) - Brutal congestion control (from hysteria) - Media datagram relay (Opus/VP9/H264/H265) - Graceful shutdown
170 lines
3.4 KiB
Go
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
|
|
}
|