// Package store provides the Nostr event storage engine. // It uses an append-only WAL for event data and sorted flat files // for indexes. Single cooperative thread - no locks. package store import ( "git.smesh.lol/moxie/pkg/mxutil" "bytes" "fmt" "math" "os" "path/filepath" "runtime" "sort" "git.smesh.lol/nostr/pkg/event" "git.smesh.lol/nostr/pkg/filter" "git.smesh.lol/nostr/pkg/hex" "git.smesh.lol/morly/pkg/store/checkpoint" "git.smesh.lol/morly/pkg/store/index" "git.smesh.lol/morly/pkg/store/serial" "git.smesh.lol/morly/pkg/store/sorted" "git.smesh.lol/morly/pkg/store/wal" ) // Engine is the main storage engine. type Engine struct { dir string w *wal.WAL ckpt *checkpoint.Checkpoint lastSer uint64 saveN int32 dirty int32 // records Put since the last sidecar write nextPubSer uint64 eid *sorted.File fpc *sorted.File sei *sorted.File ca *sorted.File exp *sorted.File kc *sorted.File pc *sorted.File kpc *sorted.File tc *sorted.File tkc *sorted.File tpc *sorted.File tkp *sorted.File wrd *sorted.File pks *sorted.File spk *sorted.File // Graph indexes. epg *sorted.File peg *sorted.File eeg *sorted.File gee *sorted.File ppg *sorted.File gpp *sorted.File } // Open opens or creates a store at dir. func Open(dir string) (eng *Engine, derr error) { // The engine and its index files outlive this call, and index files are // self-mutating types whose sovereign arena is bound to the arena that // holds the File. A frame arena here is released at return and takes // those arenas, and every pointer into them, with it. runtime.SovereignSetArena(runtime.RootArena()) defer runtime.SovereignRestoreArena(runtime.RootArena()) os.Stderr.Write("store.Open: enter dir=" | dir | "\n") idxDir := filepath.Join(dir, "idx") os.Stderr.Write("store.Open: MkdirAll " | idxDir | "\n") if merr := os.MkdirAll(idxDir, 0755); merr != nil { return nil, merr } os.Stderr.Write("store.Open: wal.Open\n") w, err := wal.Open(filepath.Join(dir, "wal")) if err != nil { return nil, err } os.Stderr.Write("store.Open: creating engine\n") e := &Engine{dir: dir, w: w, nextPubSer: 1} os.Stderr.Write("store.Open: engine created\n") open := func(name string, recLen, cmpLen int32) (*sorted.File, error) { os.Stderr.Write("store.Open: opening index " | name | "\n") f, oerr := sorted.Open(filepath.Join(idxDir, name|".dat"), recLen, cmpLen) if oerr != nil { os.Stderr.Write("store.Open: FAILED " | name | ": " | oerr.Error() | "\n") } else { os.Stderr.Write("store.Open: OK " | name | "\n") } return f, oerr } if e.eid, err = open("eid", index.EidKeyLen, index.EidCmpLen); err != nil { return nil, err } if e.fpc, err = open("fpc", index.FpcKeyLen, index.FpcCmpLen); err != nil { return nil, err } if e.sei, err = open("sei", index.SeiRecLen, index.SeiKeyLen); err != nil { return nil, err } if e.ca, err = open("ca", index.CAKeyLen, index.CAKeyLen); err != nil { return nil, err } if e.exp, err = open("exp", index.ExpKeyLen, index.ExpKeyLen); err != nil { return nil, err } if e.kc, err = open("kc", index.KCKeyLen, index.KCKeyLen); err != nil { return nil, err } if e.pc, err = open("pc", index.PCKeyLen, index.PCKeyLen); err != nil { return nil, err } if e.kpc, err = open("kpc", index.KPCKeyLen, index.KPCKeyLen); err != nil { return nil, err } if e.tc, err = open("tc", index.TCKeyLen, index.TCKeyLen); err != nil { return nil, err } if e.tkc, err = open("tkc", index.TKCKeyLen, index.TKCKeyLen); err != nil { return nil, err } if e.tpc, err = open("tpc", index.TPCKeyLen, index.TPCKeyLen); err != nil { return nil, err } if e.tkp, err = open("tkp", index.TKPKeyLen, index.TKPKeyLen); err != nil { return nil, err } if e.wrd, err = open("wrd", index.WrdKeyLen, index.WrdKeyLen); err != nil { return nil, err } if e.pks, err = open("pks", index.PksKeyLen, index.PksCmpLen); err != nil { return nil, err } if e.spk, err = open("spk", index.SpkRecLen, index.SpkKeyLen); err != nil { return nil, err } if e.epg, err = open("epg", index.EpgKeyLen, index.EpgKeyLen); err != nil { return nil, err } if e.peg, err = open("peg", index.PegKeyLen, index.PegKeyLen); err != nil { return nil, err } if e.eeg, err = open("eeg", index.EegKeyLen, index.EegKeyLen); err != nil { return nil, err } if e.gee, err = open("gee", index.GeeKeyLen, index.GeeKeyLen); err != nil { return nil, err } if e.ppg, err = open("ppg", index.PpgKeyLen, index.PpgKeyLen); err != nil { return nil, err } if e.gpp, err = open("gpp", index.GppKeyLen, index.GppKeyLen); err != nil { return nil, err } // Open checkpoint. os.Stderr.Write("store.Open: opening checkpoint\n") ckpt, err := checkpoint.Open(filepath.Join(dir, "checkpoint.dat")) if err != nil { return nil, err } os.Stderr.Write("store.Open: checkpoint opened\n") e.ckpt = ckpt ckptSer := ckpt.Get() os.Stderr.Write("store.Open: ckptSer obtained\n") // Recovery: replay WAL entries after checkpoint. if ckptSer == 0 { fmt.Println("store: no checkpoint, rebuilding indexes from WAL...") if rerr := e.rebuildIndexes(); rerr != nil { return nil, fmt.Errorf("store: recovery failed: %w", rerr) } if ferr := e.Flush(); ferr != nil { return nil, ferr } } else { count := 0 ferr := e.w.ForEachFrom(ckptSer, func(ser uint64, data []byte) bool { ev := event.New() if uerr := ev.UnmarshalBinary(bytes.NewReader(data)); uerr != nil { return true } e.indexEvent(ev, ser) e.lastSer = ser count++ return true }) if ferr == wal.ErrCheckpointStale { fmt.Println("store: checkpoint stale, full rebuild...") for _, f := range e.allFiles() { f.Clear() } e.nextPubSer = 1 if rerr := e.rebuildIndexes(); rerr != nil { return nil, fmt.Errorf("store: recovery failed: %w", rerr) } if ferr2 := e.Flush(); ferr2 != nil { return nil, ferr2 } } else if ferr != nil { return nil, fmt.Errorf("store: incremental recovery failed: %w", ferr) } else if count > 0 { fmt.Println("store: replayed", count, "events from checkpoint") } } // Quick-flush + checkpoint after recovery. if e.lastSer > 0 { e.quickCheckpoint() } // Remove legacy clean marker if present. os.Remove(filepath.Join(dir, ".clean")) // Recover next pubkey serial from the spk index. if rec, ok := e.spk.Last(); ok { e.nextPubSer = serial.Get(rec[3:]) + 1 } return e, nil } // Close flushes all indexes, checkpoints, and closes everything. func (e *Engine) Close() (err error) { fmt.Fprintf(os.Stderr, "store: closing, flushing indexes...\n") if ferr := e.Flush(); ferr != nil { fmt.Fprintf(os.Stderr, "store: flush on close: %v\n", ferr) } files := e.allFiles() for _, f := range files { f.SkipFlush() if cerr := f.Close(); cerr != nil { return cerr } } fmt.Fprintf(os.Stderr, "store: indexes closed, syncing WAL...\n") err := e.w.Close() if err == nil { fmt.Fprintf(os.Stderr, "store: shutdown complete\n") } return err } // Flush syncs the WAL, merge-flushes all indexes, and sets the checkpoint. func (e *Engine) Flush() (err error) { if serr := e.w.Sync(); serr != nil { return serr } names := []string{ "eid", "fpc", "sei", "ca", "exp", "kc", "pc", "kpc", "tc", "tkc", "tpc", "tkp", "wrd", "pks", "spk", "epg", "peg", "eeg", "gee", "ppg", "gpp", } for i, f := range e.allFiles() { if f.Dirty() { os.Stderr.Write("store: flushing " | names[i] | "...\n") } if ferr := f.Flush(); ferr != nil { return ferr } } if e.lastSer > 0 { if kerr := e.ckpt.Set(e.lastSer); kerr != nil { return kerr } } return nil } // quickCheckpoint does an incremental flush: WAL sync, QuickFlush all indexes, // then set checkpoint to the minimum flushed serial across all indexes. func (e *Engine) quickCheckpoint() { e.writeSidecars(true) e.saveN = 0 } // writeSidecars moves the in-memory index buffers into the .buf sidecars so // that a reader which only sees those files observes the records. Durability is // a separate concern: fsync is what costs (18 files per checkpoint, ~60ms), so // a query-facing flush writes without syncing and the periodic checkpoint, Flush // and Close sync. func (e *Engine) writeSidecars(sync bool) { if !sync && e.dirty == 0 { return } if sync { e.w.Sync() } minSer := flushSidecars(e, sync) if minSer < math.MaxUint64 { e.ckpt.Set(minSer) } e.dirty = 0 } // flushSidecars quick-flushes every index file, flushes the ones that asked to // be merged and returns the minimum flushed serial (MaxUint64 when none). // // A free function: writeSidecars writes the receiver, so its own allocations // land in the Engine's sovereign data arena and stay until compaction - and // this is on every query's flush path, where the file list from allFiles and // the merge set are pure scratch. Here they die with the frame. func flushSidecars(e *Engine, sync bool) (minSer uint64) { var merge []*sorted.File // Declared outside the loops: a declaration in a loop body allocates on // every iteration. var needsMerge bool var err error for _, f := range e.allFiles() { needsMerge, err = f.QuickFlush(e.lastSer, sync) if err != nil { fmt.Fprintf(os.Stderr, "store: quickflush error: %v\n", err) } if needsMerge { merge = mxutil.Ensure(merge, 1) merge = push(merge, f) } } // Checkpoint = min(lastFlushedSer) across all indexes. minSer = uint64(math.MaxUint64) var s uint64 for _, f := range e.allFiles() { s = f.LastFlushedSer() if s > 0 && s < minSer { minSer = s } } for _, f := range merge { f.Flush() } return minSer } // indexEvent adds all index entries for an event at the given serial. func (e *Engine) indexEvent(ev *event.E, ser uint64) { pubHash := serial.PubHash(ev.Pubkey) idHash := serial.IdHash(ev.ID) e.eid.Put(index.MakeEid(idHash, ser)) e.sei.Put(index.MakeSeiRec(ser, ev.ID)) e.fpc.Put(index.MakeFpc(ser, ev.ID, pubHash, ev.CreatedAt)) e.ca.Put(index.MakeCA(ev.CreatedAt, ser)) e.kc.Put(index.MakeKC(ev.Kind, ev.CreatedAt, ser)) e.pc.Put(index.MakePC(pubHash, ev.CreatedAt, ser)) e.kpc.Put(index.MakeKPC(ev.Kind, pubHash, ev.CreatedAt, ser)) if ev.Tags != nil { for _, tg := range ev.Tags.T { if tg.Len() < 2 { continue } key := tg.Key() if len(key) != 1 { continue } tagKey := key[0] val := tg.ValueHex() if len(val) == 0 { continue } valHash := serial.Ident(val) e.tc.Put(index.MakeTC(tagKey, valHash, ev.CreatedAt, ser)) e.tkc.Put(index.MakeTKC(ev.Kind, tagKey, valHash, ev.CreatedAt, ser)) e.tpc.Put(index.MakeTPC(pubHash, tagKey, valHash, ev.CreatedAt, ser)) e.tkp.Put(index.MakeTKP(ev.Kind, pubHash, tagKey, valHash, ev.CreatedAt, ser)) } } if len(ev.Content) > 0 { words := splitWords(ev.Content) seen := map[string]bool{} for _, w := range words { if seen[string(w)] { continue } seen[string(w)] = true e.wrd.Put(index.MakeWrd(serial.Ident(w), ser)) } } e.populateGraphEdges(ev, ser) } // rebuildIndexes clears all indexes and rebuilds from WAL. // replaySink re-indexes WAL entries during recovery. It is a method on a single // pointer because that is the shape a callback reaches reliably: a closure that // captures more than one word hands the callee an environment pointer that is // not always valid, and the replay then faults on its first entry (reproduced // with a 2000-entry WAL, not with 200). type replaySink struct { e *Engine count int32 } func (s *replaySink) add(ser uint64, data []byte) (ok bool) { ev := event.New() if uerr := ev.UnmarshalBinary(bytes.NewReader(data)); uerr != nil { return true } s.e.indexEvent(ev, ser) s.e.lastSer = ser s.count++ return true } func (e *Engine) rebuildIndexes() (err error) { for _, f := range e.allFiles() { if cerr := f.Clear(); cerr != nil { return cerr } } e.nextPubSer = 1 sink := &replaySink{e: e} err = e.w.ForEach(sink.add) if err != nil { return err } fmt.Println("store: rebuilt indexes for", sink.count, "events") return nil } func (e *Engine) allFiles() (ss []*sorted.File) { return []*sorted.File{ e.eid, e.fpc, e.sei, e.ca, e.exp, e.kc, e.pc, e.kpc, e.tc, e.tkc, e.tpc, e.tkp, e.wrd, e.pks, e.spk, e.epg, e.peg, e.eeg, e.gee, e.ppg, e.gpp, } } // SaveEvent persists an event and all its indexes. func (e *Engine) SaveEvent(ev *event.E) (err error) { // Duplicate check via ID hash. idHash := serial.IdHash(ev.ID) if e.hasEvent(idHash, ev.ID) { return fmt.Errorf("duplicate event") } // Write event bytes to WAL. data := ev.MarshalBinaryToBytes(nil) ser, err := e.w.Append(data) if err != nil { return err } e.indexEvent(ev, ser) e.lastSer = ser e.saveN++ e.dirty++ // Checkpoint in batches. quickCheckpoint fsyncs every index file, and doing // that per accepted event held ingest to about 15 events/second (200 saves // measured 12.8s), so a burst of a few hundred events stalled the relay and // its clients timed out. A record is visible to queries as soon as it is // Put - the scan merges the writer's in-memory buffer - and the WAL holds it // until the next checkpoint, so the indexes only have to catch up // periodically, on Close, and whenever a query asks (flushPending). if e.saveN >= checkpointEvery { e.quickCheckpoint() } return nil } // checkpointEvery is the number of accepted events between index checkpoints. const checkpointEvery = 64 // MaxSerial returns the database's monotonic record counter: the serial of the // most recent stored event, or 0 when nothing has been stored. Ephemeral events // never reach the WAL, so they never consume a serial. func (e *Engine) MaxSerial() (n uint64) { return e.lastSer } // flushPending checkpoints pending writes so that another handle on the same // files (or a reader that only sees the sidecars) observes them. func (e *Engine) flushPending() { e.writeSidecars(false) } // populateGraphEdges creates all graph index entries for an event. func (e *Engine) populateGraphEdges(ev *event.E, ser uint64) { authorSer := e.getOrCreatePubkeySerial(ev.Pubkey) kind := ev.Kind // epg/peg: author edge. e.epg.Put(index.MakeEpg(ser, authorSer, kind, index.DirAuthor)) e.peg.Put(index.MakePeg(authorSer, kind, index.DirAuthor, ser)) if ev.Tags == nil { return } for _, tg := range ev.Tags.T { if tg.Len() < 2 { continue } key := tg.Key() if len(key) != 1 { continue } switch key[0] { case 'p': valHex := tg.ValueHex() if len(valHex) != 64 { continue } pkBytes, err := hex.DecAppend(nil, valHex) if err != nil || len(pkBytes) != 32 { continue } pkSer := e.getOrCreatePubkeySerial(pkBytes) // epg/peg: p-tag edges. e.epg.Put(index.MakeEpg(ser, pkSer, kind, index.DirPTagOut)) e.peg.Put(index.MakePeg(pkSer, kind, index.DirPTagIn, ser)) // ppg/gpp: direct pubkey-to-pubkey edge (skip self-reference). if pkSer != authorSer { e.ppg.Put(index.MakePpg(authorSer, pkSer, kind, index.DirPKOut, ser)) e.gpp.Put(index.MakeGpp(pkSer, kind, index.DirPKIn, authorSer, ser)) } case 'e': valHex := tg.ValueHex() if len(valHex) != 64 { continue } evBytes, err := hex.DecAppend(nil, valHex) if err != nil || len(evBytes) != 32 { continue } // Look up target event serial - only create edge if target exists. tgtSer, ok := e.getEventSerial(evBytes) if !ok { continue } // eeg/gee: e-tag edges. e.eeg.Put(index.MakeEeg(ser, tgtSer, kind, index.DirETagOut)) e.gee.Put(index.MakeGee(tgtSer, kind, index.DirETagIn, ser)) } } } // getEventSerial returns the WAL serial for a 32-byte event ID, if it exists. func (e *Engine) getEventSerial(id []byte) (ser uint64, ok bool) { idHash := serial.IdHash(id) rec, hit := e.eid.Get(index.MakeEid(idHash, 0)) if !hit { return 0, false } s := serial.Get(rec[3+serial.HashLen:]) eid, exists := e.getEventIDBySerial(s) if !exists || !bytes.Equal(eid, id) { return 0, false } return s, true } // SearchWord returns serials of events containing the given word. func (e *Engine) SearchWord(word []byte) (ss []uint64) { wh := serial.Ident(word) return e.collectSerials(e.wrd, index.MakeWrd(wh, 0), index.MakeWrd(wh, serial.Max), index.WrdKeyLen) } func splitWords(content []byte) (ss [][]byte) { var words [][]byte var word []byte for _, b := range content { if b >= 'A' && b <= 'Z' { word = mxutil.Ensure(word, 1) word = push(word, b+32) } else if (b >= 'a' && b <= 'z') || (b >= '0' && b <= '9') { word = mxutil.Ensure(word, 1) word = push(word, b) } else { if len(word) >= 3 { w := []byte{:len(word)} copy(w, word) words = mxutil.Ensure(words, 1) words = push(words, w) } word = word[:0] } } if len(word) >= 3 { w := []byte{:len(word)} copy(w, word) words = mxutil.Ensure(words, 1) words = push(words, w) } return words } // Search performs full-text search over the word index. Returns at most // limit events (0 = no limit). Uses the same splitWords as the indexer, // so word minimum length (3) and normalisation match. func (e *Engine) Search(query []byte, limit int32) (ss []*event.E) { return searchQuery(e, query, limit) } // searchQuery is a free function: the query's scratch - the split words, the // per-word serial sets, the intersections and every SearchWord/GetBySerial // return - dies with its frame. Allocated in Search itself they would land in // the Engine's sovereign data arena, which a mutating method's allocations do, // and stay there until compaction. func searchQuery(e *Engine, query []byte, limit int32) (ss []*event.E) { words := splitWords(query) if len(words) == 0 { return nil } var sets [][]uint64 for _, w := range words { serials := e.SearchWord(w) if len(serials) == 0 { return nil } sets = mxutil.Ensure(sets, 1) sets = push(sets, serials) } result := sets[0] for i := 1; i < len(sets); i++ { result = intersectU64(result, sets[i]) if len(result) == 0 { return nil } } var out []*event.E for i := len(result) - 1; i >= 0; i-- { ev, err := e.GetBySerial(result[i]) if err != nil || ev == nil { continue } out = mxutil.Ensure(out, 1) out = push(out, ev) if limit > 0 && len(out) >= limit { break } } return out } func intersectU64(a, b []uint64) (ss []uint64) { set := map[uint64]bool{} for _, v := range a { set[v] = true } var out []uint64 for _, v := range b { if set[v] { out = mxutil.Ensure(out, 1) out = push(out, v) } } return out } // GetBySerial reads an event from the WAL by its serial. The event is handed // back as a deep copy: UnmarshalBinary fills receiver fields in the decoding // frame, and a receiver write that is not part of a return leaves the decoded // tag fields in that frame. Clone builds the value as a return value, so the // relocation moves the whole graph into the caller's arena. func (e *Engine) GetBySerial(ser uint64) (ev *event.E, err error) { data, err := e.w.Read(ser) if err != nil { return nil, err } tmp := event.New() if uerr := tmp.UnmarshalBinary(bytes.NewReader(data)); uerr != nil { return nil, uerr } ev = tmp.Clone() return } // QueryEvents returns events matching the filter, sorted newest-first. func (e *Engine) QueryEvents(f *filter.F) (evs []*event.E, err error) { // Make records written since the last checkpoint visible to every reader, // not just this handle's own buffer. This is the only receiver write on the // path, so it stays in the method: a free function receives e as an // immutable borrow and cannot call a mutating method on it. e.flushPending() return queryEvents(e, f) } // queryEvents is a free function: the scan serials, the dedup map, every // per-candidate GetBySerial return and the result list die with its frame. // Written in QueryEvents itself they would land in the Engine's sovereign data // arena - a method that writes the receiver allocates there - and stay until // compaction, on the relay's hottest path. func queryEvents(e *Engine, f *filter.F) (evs []*event.E, qerr error) { if f.Ids != nil && f.Ids.Len() > 0 { return e.queryByIDs(f) } var since, until int64 if f.Since != nil && f.Since.I64() != 0 { since = f.Since.I64() } if f.Until != nil && f.Until.I64() != 0 { until = f.Until.I64() } else { until = math.MaxInt64 } hasKinds := f.Kinds != nil && f.Kinds.Len() > 0 hasAuthors := f.Authors != nil && f.Authors.Len() > 0 hasTags := f.Tags != nil && f.Tags.Len() > 0 var serials []uint64 switch { case hasTags && hasKinds && hasAuthors: serials = e.scanTKP(f, since, until) case hasTags && hasKinds: serials = e.scanTKC(f, since, until) case hasTags && hasAuthors: serials = e.scanTPC(f, since, until) case hasTags: serials = e.scanTC(f, since, until) case hasKinds && hasAuthors: serials = e.scanKPC(f, since, until) case hasKinds: serials = e.scanKC(f, since, until) case hasAuthors: serials = e.scanPC(f, since, until) default: serials = e.scanCA(since, until) } // Deduplicate. seen := map[uint64]bool{} deduped := serials[:0] for _, s := range serials { if !seen[s] { seen[s] = true deduped = mxutil.Ensure(deduped, 1) deduped = push(deduped, s) } } // Fetch and apply the full filter. Newest-first ordering is applied by the // caller when it frames the events. // // Index order is ascending by (..., created_at), so walking the serials // backwards visits the newest first. When the filter carries a limit that // is the order the caller wants, and there is no reason to decode every // older match: with the scan no longer capped, a limited REQ over a store // holding thousands of events decoded all of them and the relay's clients // timed out. Filters without a limit keep the forward walk. var results []*event.E var limit int32 if f.Limit != nil && *f.Limit > 0 { if *f.Limit > 0x7fffffff { limit = 0x7fffffff } else { limit = int32(*f.Limit) } } start, end, step := int32(0), int32(len(deduped)), int32(1) if limit > 0 { start, end, step = int32(len(deduped))-1, -1, -1 } for i := start; i != end; i += step { ev, err := e.GetBySerial(deduped[i]) if err != nil { continue } if !f.Matches(ev) { continue } results = mxutil.Ensure(results, 1) results = push(results, ev) if limit > 0 && int32(len(results)) >= limit { break } } sort.Sort(&event.S{E: results}) // reverse chronological // The result limit is applied by the caller when it frames the events // (cfg.QueryResultLimit in the relay), so a store-side reslice is // redundant. It is also the one construct here that miscompiles: // results[:*f.Limit] returned an empty slice to the caller even with // matches in hand. return results, nil } // DeleteEvent removes an event and its primary index entries. func (e *Engine) DeleteEvent(id []byte) (err error) { idHash := serial.IdHash(id) // Find the serial for this event. The eid key is prefix|idHash|serial, so // the first record under the id-hash prefix is the event and the serial // follows in the record. This is a direct lookup rather than a scan // callback: a callback's captured state does not reliably reach the // enclosing frame, which silently turned every delete into "event not // found" and left replaceable events to accumulate. rec, ok := e.eid.GetPrefix(index.MakeEidPrefix(idHash)) if !ok { return fmt.Errorf("event not found") } ser := serial.Get(rec[3+serial.HashLen:]) // Read the event to get kind/pubkey/timestamp for index cleanup. ev, err := e.GetBySerial(ser) if err != nil { return err } pubHash := serial.PubHash(ev.Pubkey) // Delete primary index entries. e.eid.Delete(index.MakeEid(idHash, ser)) e.ca.Delete(index.MakeCA(ev.CreatedAt, ser)) e.kc.Delete(index.MakeKC(ev.Kind, ev.CreatedAt, ser)) e.pc.Delete(index.MakePC(pubHash, ev.CreatedAt, ser)) e.kpc.Delete(index.MakeKPC(ev.Kind, pubHash, ev.CreatedAt, ser)) return nil } // ResolveBlob returns every stored event that references the blob hash through // an `x` tag, newest first. The tc index already holds a record per // single-letter tag, so this is the same key-range scan a {"#x": [...]} filter // runs - no new index family, no on-disk change. // // The scan and the seen map are scratch, not object state, so they live in a // free function that gets its own frame arena; the method itself only returns. // This is the queryEvents/QueryEvents split for the same reason. func (e *Engine) ResolveBlob(xHex []byte) (out []*event.E) { return resolveBlob(e, xHex) } func resolveBlob(e *Engine, xHex []byte) (out []*event.E) { if len(xHex) == 0 { return nil } vh := serial.Ident(xHex) serials := e.collectSerials(e.tc, index.MakeTC('x', vh, 0, 0), index.MakeTC('x', vh, math.MaxInt64, serial.Max), index.TCKeyLen) if len(serials) == 0 { return nil } // Index order is ascending by (created_at, serial), so walking it // backwards visits the newest event first. One event can carry the same x // tag twice and contribute two records; dedup by serial. seen := map[uint64]bool{} for i := int32(len(serials)) - 1; i >= 0; i-- { ser := serials[i] if seen[ser] { continue } seen[ser] = true ev, err := e.GetBySerial(ser) if err != nil { continue } out = mxutil.Ensure(out, 1) out = push(out, ev) } return } // GetByID retrieves an event by its 32-byte ID. func (e *Engine) GetByID(id []byte) (ev *event.E, err error) { // Point lookup, exactly as queryByIDs does it: the eid index compares only // prefix|id-hash (EidCmpLen), so the 16-byte range // Scan(MakeEid(hash,0), MakeEid(hash,Max)) this used was an empty range and // GetByID answered "event not found" for every event that was stored. // negentropy's FindHave calls this and dropped every event on the error, // so a sync found nothing to send. ser, ok := e.getEventSerial(id) if !ok { return nil, fmt.Errorf("event not found") } ev, err := e.GetBySerial(ser) if err != nil { return nil, err } return ev, nil } // --- internal helpers --- func (e *Engine) hasEvent(idHash, fullID []byte) (ok bool) { // Point lookup: the eid key is prefix|id-hash|serial and cmpLen stops at // the hash, so Get returns the first record for this hash and the full ID // check filters hash collisions. A range Scan here cost a Stat, a whole // .buf sidecar read and a 256-record chunk buffer per event. rec, found := e.eid.Get(index.MakeEid(idHash, 0)) if !found { return false } ser := serial.Get(rec[3+serial.HashLen:]) eid, hasID := e.getEventIDBySerial(ser) if !hasID { return false } return bytes.Equal(eid, fullID) } func (e *Engine) getOrCreatePubkeySerial(pubkey []byte) (n uint64) { pubHash := serial.PubHash(pubkey) rec, exists := e.pks.Get(index.MakePks(pubHash, 0)) if exists { return serial.Get(rec[3+serial.HashLen:]) } ser := e.nextPubSer e.nextPubSer++ e.pks.Put(index.MakePks(pubHash, ser)) e.spk.Put(index.MakeSpkRec(ser, pubkey)) return ser } func (e *Engine) getEventIDBySerial(ser uint64) (id []byte, ok bool) { key := index.MakeSei(ser) rec, ok := e.sei.Get(key) if !ok { return nil, false } return rec[index.SeiKeyLen:], true } func (e *Engine) queryByIDs(f *filter.F) (evs []*event.E, qerr error) { var results []*event.E for _, id := range f.Ids.T { // Point lookup, not a range Scan: the eid index compares only // prefix|id-hash (EidCmpLen), so two records for the same hash compare // equal and Scan(MakeEid(hash,0), MakeEid(hash,Max)) is an empty // range - every by-id query came back with no events. getEventSerial // does the Get and verifies the full 32-byte id, which also filters // hash collisions. ser, ok := e.getEventSerial(id) if !ok { continue } if ev, err := e.GetBySerial(ser); err == nil { results = mxutil.Ensure(results, 1) results = push(results, ev) } } return results, nil } func (e *Engine) collectSerials(idx *sorted.File, start, end []byte, keyLen int32) (ss []uint64) { // No allocation inside the callback: its frame arena dies per call (see // queryByIDs). A buffer plus a count keeps the scan allocation-free. // The scan callback is a separate function: only a single captured pointer // reliably reaches the enclosing frame, so writes to a captured local array // were lost (the counter advanced, every value read back zero). Keep the // collected state behind a pointer. // // The buffer is sized from the index, not fixed: a fixed sink silently // stopped the scan at its array size, so once a store held more records // than that, every kinds/authors/created-at query missed the records that // sorted after the first batch - a freshly published event was invisible // to every scan while the by-id lookup, whose range holds a handful of // records, still found it. total := idx.Count() if total <= 0 { return nil } if total > math.MaxInt32 { total = math.MaxInt32 } state := &serialSink{buf: []uint64{:int32(total)}, keyLen: keyLen} idx.Scan(start, end, func(rec []byte) bool { return state.add(rec) }) n := state.n ss = []uint64{:n} for i := int32(0); i < n; i++ { ss[i] = state.buf[i] } return ss } // serialSink collects serials from an index scan. type serialSink struct { buf []uint64 n int32 keyLen int32 } func (s *serialSink) add(rec []byte) (ok bool) { if s.n >= int32(len(s.buf)) { return false } s.buf[s.n] = serial.Get(rec[s.keyLen-serial.Len:]) s.n++ return true } func (e *Engine) scanCA(since, until int64) (ss []uint64) { return e.collectSerials(e.ca, index.MakeCA(since, 0), index.MakeCA(until, serial.Max), index.CAKeyLen) } func (e *Engine) scanKC(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, k := range f.Kinds.K { got := e.collectSerials(e.kc, index.MakeKC(k.K, since, 0), index.MakeKC(k.K, until, serial.Max), index.KCKeyLen) out = out | got } return out } func (e *Engine) scanPC(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, author := range f.Authors.T { ph := serial.PubHash(author) out = out | e.collectSerials(e.pc, index.MakePC(ph, since, 0), index.MakePC(ph, until, serial.Max), index.PCKeyLen) } return out } func (e *Engine) scanKPC(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, k := range f.Kinds.K { for _, author := range f.Authors.T { ph := serial.PubHash(author) out = out | e.collectSerials(e.kpc, index.MakeKPC(k.K, ph, since, 0), index.MakeKPC(k.K, ph, until, serial.Max), index.KPCKeyLen) } } return out } func (e *Engine) scanTC(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, tg := range f.Tags.T { if tg.Len() < 2 || len(tg.Key()) != 1 { continue } tk := tg.Key()[0] for _, val := range tg.T[1:] { vh := serial.Ident(val) out = out | e.collectSerials(e.tc, index.MakeTC(tk, vh, since, 0), index.MakeTC(tk, vh, until, serial.Max), index.TCKeyLen) } } return out } func (e *Engine) scanTKC(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, tg := range f.Tags.T { if tg.Len() < 2 || len(tg.Key()) != 1 { continue } tk := tg.Key()[0] for _, val := range tg.T[1:] { vh := serial.Ident(val) for _, k := range f.Kinds.K { out = out | e.collectSerials(e.tkc, index.MakeTKC(k.K, tk, vh, since, 0), index.MakeTKC(k.K, tk, vh, until, serial.Max), index.TKCKeyLen) } } } return out } func (e *Engine) scanTPC(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, tg := range f.Tags.T { if tg.Len() < 2 || len(tg.Key()) != 1 { continue } tk := tg.Key()[0] for _, val := range tg.T[1:] { vh := serial.Ident(val) for _, author := range f.Authors.T { ph := serial.PubHash(author) out = out | e.collectSerials(e.tpc, index.MakeTPC(ph, tk, vh, since, 0), index.MakeTPC(ph, tk, vh, until, serial.Max), index.TPCKeyLen) } } } return out } func (e *Engine) scanTKP(f *filter.F, since, until int64) (ss []uint64) { var out []uint64 for _, tg := range f.Tags.T { if tg.Len() < 2 || len(tg.Key()) != 1 { continue } tk := tg.Key()[0] for _, val := range tg.T[1:] { vh := serial.Ident(val) for _, k := range f.Kinds.K { for _, author := range f.Authors.T { ph := serial.PubHash(author) out = out | e.collectSerials(e.tkp, index.MakeTKP(k.K, ph, tk, vh, since, 0), index.MakeTKP(k.K, ph, tk, vh, until, serial.Max), index.TKPKeyLen) } } } } return out } // --- Graph traversal --- // GetReferencingEvents returns serials of events that reference targetSer via e-tags. // Uses the gee (reverse event-event) index. func (e *Engine) GetReferencingEvents(targetSer uint64) (ss []uint64) { start := index.MakeGee(targetSer, 0, 0, 0) end := index.MakeGee(targetSer, 0xFFFF, 0xFF, serial.Max) var out []uint64 e.gee.Scan(start, end, func(rec []byte) bool { out = mxutil.Ensure(out, 1) out = push(out, serial.Get(rec[11:])) return true }) return out } // GetETagTargets returns serials of events that srcSer references via e-tags. // Uses the eeg (forward event-event) index. func (e *Engine) GetETagTargets(srcSer uint64) (ss []uint64) { start := index.MakeEeg(srcSer, 0, 0, 0) end := index.MakeEeg(srcSer, serial.Max, 0xFFFF, 0xFF) var out []uint64 e.eeg.Scan(start, end, func(rec []byte) bool { out = mxutil.Ensure(out, 1) out = push(out, serial.Get(rec[8:])) return true }) return out } // TraverseThread performs BFS traversal of thread structure via e-tags. // Returns all event IDs (32-byte binary) reachable from seedID within maxDepth. func (e *Engine) TraverseThread(seedID []byte, maxDepth int32, direction string) (ss [][]byte) { seedSer, ok := e.getEventSerial(seedID) if !ok { return nil } if direction == "" { direction = "both" } return traverseThread(e, seedSer, maxDepth, direction) } // traverseThread is a free function: the breadth-first scratch - the visited // set, each frontier slice, and every GetReferencingEvents/GetETagTargets // return - dies with its frame instead of accumulating in the Engine's // sovereign data arena until the next compaction. func traverseThread(e *Engine, seedSer uint64, maxDepth int32, direction string) (allIDs [][]byte) { visited := map[uint64]bool{} visited[seedSer] = true frontier := []uint64{seedSer} emptyStreak := 0 for depth := 1; depth <= maxDepth; depth++ { var next []uint64 found := 0 for _, evSer := range frontier { if direction == "both" || direction == "inbound" { for _, refSer := range e.GetReferencingEvents(evSer) { if visited[refSer] { continue } visited[refSer] = true if eid, hasID := e.getEventIDBySerial(refSer); hasID { allIDs = mxutil.Ensure(allIDs, 1) allIDs = push(allIDs, eid) found++ } next = mxutil.Ensure(next, 1) next = push(next, refSer) } } if direction == "both" || direction == "outbound" { for _, tgtSer := range e.GetETagTargets(evSer) { if visited[tgtSer] { continue } visited[tgtSer] = true if eid, hasID := e.getEventIDBySerial(tgtSer); hasID { allIDs = mxutil.Ensure(allIDs, 1) allIDs = push(allIDs, eid) found++ } next = mxutil.Ensure(next, 1) next = push(next, tgtSer) } } } if found == 0 { emptyStreak++ if emptyStreak >= 2 { break } } else { emptyStreak = 0 } frontier = next } return allIDs }