engine.mx raw
1 // Package store provides the Nostr event storage engine.
2 // It uses an append-only WAL for event data and sorted flat files
3 // for indexes. Single cooperative thread - no locks.
4 package store
5
6 import (
7 "git.smesh.lol/moxie/pkg/mxutil"
8 "bytes"
9 "fmt"
10 "math"
11 "os"
12 "path/filepath"
13 "runtime"
14 "sort"
15
16 "git.smesh.lol/nostr/pkg/event"
17 "git.smesh.lol/nostr/pkg/filter"
18 "git.smesh.lol/nostr/pkg/hex"
19 "git.smesh.lol/morly/pkg/store/checkpoint"
20 "git.smesh.lol/morly/pkg/store/index"
21 "git.smesh.lol/morly/pkg/store/serial"
22 "git.smesh.lol/morly/pkg/store/sorted"
23 "git.smesh.lol/morly/pkg/store/wal"
24 )
25
26 // Engine is the main storage engine.
27 type Engine struct {
28 dir string
29 w *wal.WAL
30
31 ckpt *checkpoint.Checkpoint
32 lastSer uint64
33 saveN int32
34 dirty int32 // records Put since the last sidecar write
35
36 nextPubSer uint64
37
38 eid *sorted.File
39 fpc *sorted.File
40 sei *sorted.File
41 ca *sorted.File
42 exp *sorted.File
43 kc *sorted.File
44 pc *sorted.File
45 kpc *sorted.File
46 tc *sorted.File
47 tkc *sorted.File
48 tpc *sorted.File
49 tkp *sorted.File
50 wrd *sorted.File
51 pks *sorted.File
52 spk *sorted.File
53
54 // Graph indexes.
55 epg *sorted.File
56 peg *sorted.File
57 eeg *sorted.File
58 gee *sorted.File
59 ppg *sorted.File
60 gpp *sorted.File
61 }
62
63 // Open opens or creates a store at dir.
64 func Open(dir string) (eng *Engine, derr error) {
65 // The engine and its index files outlive this call, and index files are
66 // self-mutating types whose sovereign arena is bound to the arena that
67 // holds the File. A frame arena here is released at return and takes
68 // those arenas, and every pointer into them, with it.
69 runtime.SovereignSetArena(runtime.RootArena())
70 defer runtime.SovereignRestoreArena(runtime.RootArena())
71 os.Stderr.Write("store.Open: enter dir=" | dir | "\n")
72 idxDir := filepath.Join(dir, "idx")
73 os.Stderr.Write("store.Open: MkdirAll " | idxDir | "\n")
74 if merr := os.MkdirAll(idxDir, 0755); merr != nil {
75 return nil, merr
76 }
77 os.Stderr.Write("store.Open: wal.Open\n")
78 w, err := wal.Open(filepath.Join(dir, "wal"))
79 if err != nil {
80 return nil, err
81 }
82 os.Stderr.Write("store.Open: creating engine\n")
83 e := &Engine{dir: dir, w: w, nextPubSer: 1}
84 os.Stderr.Write("store.Open: engine created\n")
85
86 open := func(name string, recLen, cmpLen int32) (*sorted.File, error) {
87 os.Stderr.Write("store.Open: opening index " | name | "\n")
88 f, oerr := sorted.Open(filepath.Join(idxDir, name|".dat"), recLen, cmpLen)
89 if oerr != nil {
90 os.Stderr.Write("store.Open: FAILED " | name | ": " | oerr.Error() | "\n")
91 } else {
92 os.Stderr.Write("store.Open: OK " | name | "\n")
93 }
94 return f, oerr
95 }
96 if e.eid, err = open("eid", index.EidKeyLen, index.EidCmpLen); err != nil {
97 return nil, err
98 }
99 if e.fpc, err = open("fpc", index.FpcKeyLen, index.FpcCmpLen); err != nil {
100 return nil, err
101 }
102 if e.sei, err = open("sei", index.SeiRecLen, index.SeiKeyLen); err != nil {
103 return nil, err
104 }
105 if e.ca, err = open("ca", index.CAKeyLen, index.CAKeyLen); err != nil {
106 return nil, err
107 }
108 if e.exp, err = open("exp", index.ExpKeyLen, index.ExpKeyLen); err != nil {
109 return nil, err
110 }
111 if e.kc, err = open("kc", index.KCKeyLen, index.KCKeyLen); err != nil {
112 return nil, err
113 }
114 if e.pc, err = open("pc", index.PCKeyLen, index.PCKeyLen); err != nil {
115 return nil, err
116 }
117 if e.kpc, err = open("kpc", index.KPCKeyLen, index.KPCKeyLen); err != nil {
118 return nil, err
119 }
120 if e.tc, err = open("tc", index.TCKeyLen, index.TCKeyLen); err != nil {
121 return nil, err
122 }
123 if e.tkc, err = open("tkc", index.TKCKeyLen, index.TKCKeyLen); err != nil {
124 return nil, err
125 }
126 if e.tpc, err = open("tpc", index.TPCKeyLen, index.TPCKeyLen); err != nil {
127 return nil, err
128 }
129 if e.tkp, err = open("tkp", index.TKPKeyLen, index.TKPKeyLen); err != nil {
130 return nil, err
131 }
132 if e.wrd, err = open("wrd", index.WrdKeyLen, index.WrdKeyLen); err != nil {
133 return nil, err
134 }
135 if e.pks, err = open("pks", index.PksKeyLen, index.PksCmpLen); err != nil {
136 return nil, err
137 }
138 if e.spk, err = open("spk", index.SpkRecLen, index.SpkKeyLen); err != nil {
139 return nil, err
140 }
141 if e.epg, err = open("epg", index.EpgKeyLen, index.EpgKeyLen); err != nil {
142 return nil, err
143 }
144 if e.peg, err = open("peg", index.PegKeyLen, index.PegKeyLen); err != nil {
145 return nil, err
146 }
147 if e.eeg, err = open("eeg", index.EegKeyLen, index.EegKeyLen); err != nil {
148 return nil, err
149 }
150 if e.gee, err = open("gee", index.GeeKeyLen, index.GeeKeyLen); err != nil {
151 return nil, err
152 }
153 if e.ppg, err = open("ppg", index.PpgKeyLen, index.PpgKeyLen); err != nil {
154 return nil, err
155 }
156 if e.gpp, err = open("gpp", index.GppKeyLen, index.GppKeyLen); err != nil {
157 return nil, err
158 }
159
160 // Open checkpoint.
161 os.Stderr.Write("store.Open: opening checkpoint\n")
162 ckpt, err := checkpoint.Open(filepath.Join(dir, "checkpoint.dat"))
163 if err != nil {
164 return nil, err
165 }
166 os.Stderr.Write("store.Open: checkpoint opened\n")
167 e.ckpt = ckpt
168 ckptSer := ckpt.Get()
169 os.Stderr.Write("store.Open: ckptSer obtained\n")
170
171 // Recovery: replay WAL entries after checkpoint.
172 if ckptSer == 0 {
173 fmt.Println("store: no checkpoint, rebuilding indexes from WAL...")
174 if rerr := e.rebuildIndexes(); rerr != nil {
175 return nil, fmt.Errorf("store: recovery failed: %w", rerr)
176 }
177 if ferr := e.Flush(); ferr != nil {
178 return nil, ferr
179 }
180 } else {
181 count := 0
182 ferr := e.w.ForEachFrom(ckptSer, func(ser uint64, data []byte) bool {
183 ev := event.New()
184 if uerr := ev.UnmarshalBinary(bytes.NewReader(data)); uerr != nil {
185 return true
186 }
187 e.indexEvent(ev, ser)
188 e.lastSer = ser
189 count++
190 return true
191 })
192 if ferr == wal.ErrCheckpointStale {
193 fmt.Println("store: checkpoint stale, full rebuild...")
194 for _, f := range e.allFiles() {
195 f.Clear()
196 }
197 e.nextPubSer = 1
198 if rerr := e.rebuildIndexes(); rerr != nil {
199 return nil, fmt.Errorf("store: recovery failed: %w", rerr)
200 }
201 if ferr2 := e.Flush(); ferr2 != nil {
202 return nil, ferr2
203 }
204 } else if ferr != nil {
205 return nil, fmt.Errorf("store: incremental recovery failed: %w", ferr)
206 } else if count > 0 {
207 fmt.Println("store: replayed", count, "events from checkpoint")
208 }
209 }
210
211 // Quick-flush + checkpoint after recovery.
212 if e.lastSer > 0 {
213 e.quickCheckpoint()
214 }
215
216 // Remove legacy clean marker if present.
217 os.Remove(filepath.Join(dir, ".clean"))
218
219 // Recover next pubkey serial from the spk index.
220 if rec, ok := e.spk.Last(); ok {
221 e.nextPubSer = serial.Get(rec[3:]) + 1
222 }
223 return e, nil
224 }
225
226 // Close flushes all indexes, checkpoints, and closes everything.
227 func (e *Engine) Close() (err error) {
228 fmt.Fprintf(os.Stderr, "store: closing, flushing indexes...\n")
229 if ferr := e.Flush(); ferr != nil {
230 fmt.Fprintf(os.Stderr, "store: flush on close: %v\n", ferr)
231 }
232 files := e.allFiles()
233 for _, f := range files {
234 f.SkipFlush()
235 if cerr := f.Close(); cerr != nil {
236 return cerr
237 }
238 }
239 fmt.Fprintf(os.Stderr, "store: indexes closed, syncing WAL...\n")
240 err := e.w.Close()
241 if err == nil {
242 fmt.Fprintf(os.Stderr, "store: shutdown complete\n")
243 }
244 return err
245 }
246
247 // Flush syncs the WAL, merge-flushes all indexes, and sets the checkpoint.
248 func (e *Engine) Flush() (err error) {
249 if serr := e.w.Sync(); serr != nil {
250 return serr
251 }
252 names := []string{
253 "eid", "fpc", "sei", "ca", "exp", "kc", "pc", "kpc",
254 "tc", "tkc", "tpc", "tkp", "wrd", "pks", "spk",
255 "epg", "peg", "eeg", "gee", "ppg", "gpp",
256 }
257 for i, f := range e.allFiles() {
258 if f.Dirty() {
259 os.Stderr.Write("store: flushing " | names[i] | "...\n")
260 }
261 if ferr := f.Flush(); ferr != nil {
262 return ferr
263 }
264 }
265 if e.lastSer > 0 {
266 if kerr := e.ckpt.Set(e.lastSer); kerr != nil {
267 return kerr
268 }
269 }
270 return nil
271 }
272
273 // quickCheckpoint does an incremental flush: WAL sync, QuickFlush all indexes,
274 // then set checkpoint to the minimum flushed serial across all indexes.
275 func (e *Engine) quickCheckpoint() {
276 e.writeSidecars(true)
277 e.saveN = 0
278 }
279
280 // writeSidecars moves the in-memory index buffers into the .buf sidecars so
281 // that a reader which only sees those files observes the records. Durability is
282 // a separate concern: fsync is what costs (18 files per checkpoint, ~60ms), so
283 // a query-facing flush writes without syncing and the periodic checkpoint, Flush
284 // and Close sync.
285 func (e *Engine) writeSidecars(sync bool) {
286 if !sync && e.dirty == 0 {
287 return
288 }
289 if sync {
290 e.w.Sync()
291 }
292 minSer := flushSidecars(e, sync)
293 if minSer < math.MaxUint64 {
294 e.ckpt.Set(minSer)
295 }
296 e.dirty = 0
297 }
298
299 // flushSidecars quick-flushes every index file, flushes the ones that asked to
300 // be merged and returns the minimum flushed serial (MaxUint64 when none).
301 //
302 // A free function: writeSidecars writes the receiver, so its own allocations
303 // land in the Engine's sovereign data arena and stay until compaction - and
304 // this is on every query's flush path, where the file list from allFiles and
305 // the merge set are pure scratch. Here they die with the frame.
306 func flushSidecars(e *Engine, sync bool) (minSer uint64) {
307 var merge []*sorted.File
308 // Declared outside the loops: a declaration in a loop body allocates on
309 // every iteration.
310 var needsMerge bool
311 var err error
312 for _, f := range e.allFiles() {
313 needsMerge, err = f.QuickFlush(e.lastSer, sync)
314 if err != nil {
315 fmt.Fprintf(os.Stderr, "store: quickflush error: %v\n", err)
316 }
317 if needsMerge {
318 merge = mxutil.Ensure(merge, 1)
319 merge = push(merge, f)
320 }
321 }
322 // Checkpoint = min(lastFlushedSer) across all indexes.
323 minSer = uint64(math.MaxUint64)
324 var s uint64
325 for _, f := range e.allFiles() {
326 s = f.LastFlushedSer()
327 if s > 0 && s < minSer {
328 minSer = s
329 }
330 }
331 for _, f := range merge {
332 f.Flush()
333 }
334 return minSer
335 }
336
337 // indexEvent adds all index entries for an event at the given serial.
338 func (e *Engine) indexEvent(ev *event.E, ser uint64) {
339 pubHash := serial.PubHash(ev.Pubkey)
340 idHash := serial.IdHash(ev.ID)
341
342 e.eid.Put(index.MakeEid(idHash, ser))
343 e.sei.Put(index.MakeSeiRec(ser, ev.ID))
344 e.fpc.Put(index.MakeFpc(ser, ev.ID, pubHash, ev.CreatedAt))
345 e.ca.Put(index.MakeCA(ev.CreatedAt, ser))
346 e.kc.Put(index.MakeKC(ev.Kind, ev.CreatedAt, ser))
347 e.pc.Put(index.MakePC(pubHash, ev.CreatedAt, ser))
348 e.kpc.Put(index.MakeKPC(ev.Kind, pubHash, ev.CreatedAt, ser))
349
350 if ev.Tags != nil {
351 for _, tg := range ev.Tags.T {
352 if tg.Len() < 2 {
353 continue
354 }
355 key := tg.Key()
356 if len(key) != 1 {
357 continue
358 }
359 tagKey := key[0]
360 val := tg.ValueHex()
361 if len(val) == 0 {
362 continue
363 }
364 valHash := serial.Ident(val)
365 e.tc.Put(index.MakeTC(tagKey, valHash, ev.CreatedAt, ser))
366 e.tkc.Put(index.MakeTKC(ev.Kind, tagKey, valHash, ev.CreatedAt, ser))
367 e.tpc.Put(index.MakeTPC(pubHash, tagKey, valHash, ev.CreatedAt, ser))
368 e.tkp.Put(index.MakeTKP(ev.Kind, pubHash, tagKey, valHash, ev.CreatedAt, ser))
369 }
370 }
371
372 if len(ev.Content) > 0 {
373 words := splitWords(ev.Content)
374 seen := map[string]bool{}
375 for _, w := range words {
376 if seen[string(w)] {
377 continue
378 }
379 seen[string(w)] = true
380 e.wrd.Put(index.MakeWrd(serial.Ident(w), ser))
381 }
382 }
383
384 e.populateGraphEdges(ev, ser)
385 }
386
387 // rebuildIndexes clears all indexes and rebuilds from WAL.
388 // replaySink re-indexes WAL entries during recovery. It is a method on a single
389 // pointer because that is the shape a callback reaches reliably: a closure that
390 // captures more than one word hands the callee an environment pointer that is
391 // not always valid, and the replay then faults on its first entry (reproduced
392 // with a 2000-entry WAL, not with 200).
393 type replaySink struct {
394 e *Engine
395 count int32
396 }
397
398 func (s *replaySink) add(ser uint64, data []byte) (ok bool) {
399 ev := event.New()
400 if uerr := ev.UnmarshalBinary(bytes.NewReader(data)); uerr != nil {
401 return true
402 }
403 s.e.indexEvent(ev, ser)
404 s.e.lastSer = ser
405 s.count++
406 return true
407 }
408
409 func (e *Engine) rebuildIndexes() (err error) {
410 for _, f := range e.allFiles() {
411 if cerr := f.Clear(); cerr != nil {
412 return cerr
413 }
414 }
415 e.nextPubSer = 1
416
417 sink := &replaySink{e: e}
418 err = e.w.ForEach(sink.add)
419 if err != nil {
420 return err
421 }
422 fmt.Println("store: rebuilt indexes for", sink.count, "events")
423 return nil
424 }
425
426 func (e *Engine) allFiles() (ss []*sorted.File) {
427 return []*sorted.File{
428 e.eid, e.fpc, e.sei, e.ca, e.exp, e.kc, e.pc, e.kpc,
429 e.tc, e.tkc, e.tpc, e.tkp, e.wrd, e.pks, e.spk,
430 e.epg, e.peg, e.eeg, e.gee, e.ppg, e.gpp,
431 }
432 }
433
434 // SaveEvent persists an event and all its indexes.
435 func (e *Engine) SaveEvent(ev *event.E) (err error) {
436 // Duplicate check via ID hash.
437 idHash := serial.IdHash(ev.ID)
438 if e.hasEvent(idHash, ev.ID) {
439 return fmt.Errorf("duplicate event")
440 }
441
442 // Write event bytes to WAL.
443 data := ev.MarshalBinaryToBytes(nil)
444 ser, err := e.w.Append(data)
445 if err != nil {
446 return err
447 }
448
449 e.indexEvent(ev, ser)
450 e.lastSer = ser
451 e.saveN++
452 e.dirty++
453 // Checkpoint in batches. quickCheckpoint fsyncs every index file, and doing
454 // that per accepted event held ingest to about 15 events/second (200 saves
455 // measured 12.8s), so a burst of a few hundred events stalled the relay and
456 // its clients timed out. A record is visible to queries as soon as it is
457 // Put - the scan merges the writer's in-memory buffer - and the WAL holds it
458 // until the next checkpoint, so the indexes only have to catch up
459 // periodically, on Close, and whenever a query asks (flushPending).
460 if e.saveN >= checkpointEvery {
461 e.quickCheckpoint()
462 }
463 return nil
464 }
465
466 // checkpointEvery is the number of accepted events between index checkpoints.
467 const checkpointEvery = 64
468
469 // MaxSerial returns the database's monotonic record counter: the serial of the
470 // most recent stored event, or 0 when nothing has been stored. Ephemeral events
471 // never reach the WAL, so they never consume a serial.
472 func (e *Engine) MaxSerial() (n uint64) { return e.lastSer }
473
474 // flushPending checkpoints pending writes so that another handle on the same
475 // files (or a reader that only sees the sidecars) observes them.
476 func (e *Engine) flushPending() {
477 e.writeSidecars(false)
478 }
479
480 // populateGraphEdges creates all graph index entries for an event.
481 func (e *Engine) populateGraphEdges(ev *event.E, ser uint64) {
482 authorSer := e.getOrCreatePubkeySerial(ev.Pubkey)
483 kind := ev.Kind
484
485 // epg/peg: author edge.
486 e.epg.Put(index.MakeEpg(ser, authorSer, kind, index.DirAuthor))
487 e.peg.Put(index.MakePeg(authorSer, kind, index.DirAuthor, ser))
488
489 if ev.Tags == nil {
490 return
491 }
492
493 for _, tg := range ev.Tags.T {
494 if tg.Len() < 2 {
495 continue
496 }
497 key := tg.Key()
498 if len(key) != 1 {
499 continue
500 }
501
502 switch key[0] {
503 case 'p':
504 valHex := tg.ValueHex()
505 if len(valHex) != 64 {
506 continue
507 }
508 pkBytes, err := hex.DecAppend(nil, valHex)
509 if err != nil || len(pkBytes) != 32 {
510 continue
511 }
512 pkSer := e.getOrCreatePubkeySerial(pkBytes)
513
514 // epg/peg: p-tag edges.
515 e.epg.Put(index.MakeEpg(ser, pkSer, kind, index.DirPTagOut))
516 e.peg.Put(index.MakePeg(pkSer, kind, index.DirPTagIn, ser))
517
518 // ppg/gpp: direct pubkey-to-pubkey edge (skip self-reference).
519 if pkSer != authorSer {
520 e.ppg.Put(index.MakePpg(authorSer, pkSer, kind, index.DirPKOut, ser))
521 e.gpp.Put(index.MakeGpp(pkSer, kind, index.DirPKIn, authorSer, ser))
522 }
523
524 case 'e':
525 valHex := tg.ValueHex()
526 if len(valHex) != 64 {
527 continue
528 }
529 evBytes, err := hex.DecAppend(nil, valHex)
530 if err != nil || len(evBytes) != 32 {
531 continue
532 }
533 // Look up target event serial - only create edge if target exists.
534 tgtSer, ok := e.getEventSerial(evBytes)
535 if !ok {
536 continue
537 }
538 // eeg/gee: e-tag edges.
539 e.eeg.Put(index.MakeEeg(ser, tgtSer, kind, index.DirETagOut))
540 e.gee.Put(index.MakeGee(tgtSer, kind, index.DirETagIn, ser))
541 }
542 }
543 }
544
545 // getEventSerial returns the WAL serial for a 32-byte event ID, if it exists.
546 func (e *Engine) getEventSerial(id []byte) (ser uint64, ok bool) {
547 idHash := serial.IdHash(id)
548 rec, hit := e.eid.Get(index.MakeEid(idHash, 0))
549 if !hit {
550 return 0, false
551 }
552 s := serial.Get(rec[3+serial.HashLen:])
553 eid, exists := e.getEventIDBySerial(s)
554 if !exists || !bytes.Equal(eid, id) {
555 return 0, false
556 }
557 return s, true
558 }
559
560 // SearchWord returns serials of events containing the given word.
561 func (e *Engine) SearchWord(word []byte) (ss []uint64) {
562 wh := serial.Ident(word)
563 return e.collectSerials(e.wrd,
564 index.MakeWrd(wh, 0),
565 index.MakeWrd(wh, serial.Max),
566 index.WrdKeyLen)
567 }
568
569 func splitWords(content []byte) (ss [][]byte) {
570 var words [][]byte
571 var word []byte
572 for _, b := range content {
573 if b >= 'A' && b <= 'Z' {
574 word = mxutil.Ensure(word, 1)
575 word = push(word, b+32)
576 } else if (b >= 'a' && b <= 'z') || (b >= '0' && b <= '9') {
577 word = mxutil.Ensure(word, 1)
578 word = push(word, b)
579 } else {
580 if len(word) >= 3 {
581 w := []byte{:len(word)}
582 copy(w, word)
583 words = mxutil.Ensure(words, 1)
584 words = push(words, w)
585 }
586 word = word[:0]
587 }
588 }
589 if len(word) >= 3 {
590 w := []byte{:len(word)}
591 copy(w, word)
592 words = mxutil.Ensure(words, 1)
593 words = push(words, w)
594 }
595 return words
596 }
597
598 // Search performs full-text search over the word index. Returns at most
599 // limit events (0 = no limit). Uses the same splitWords as the indexer,
600 // so word minimum length (3) and normalisation match.
601 func (e *Engine) Search(query []byte, limit int32) (ss []*event.E) {
602 return searchQuery(e, query, limit)
603 }
604
605 // searchQuery is a free function: the query's scratch - the split words, the
606 // per-word serial sets, the intersections and every SearchWord/GetBySerial
607 // return - dies with its frame. Allocated in Search itself they would land in
608 // the Engine's sovereign data arena, which a mutating method's allocations do,
609 // and stay there until compaction.
610 func searchQuery(e *Engine, query []byte, limit int32) (ss []*event.E) {
611 words := splitWords(query)
612 if len(words) == 0 {
613 return nil
614 }
615 var sets [][]uint64
616 for _, w := range words {
617 serials := e.SearchWord(w)
618 if len(serials) == 0 {
619 return nil
620 }
621 sets = mxutil.Ensure(sets, 1)
622 sets = push(sets, serials)
623 }
624 result := sets[0]
625 for i := 1; i < len(sets); i++ {
626 result = intersectU64(result, sets[i])
627 if len(result) == 0 {
628 return nil
629 }
630 }
631 var out []*event.E
632 for i := len(result) - 1; i >= 0; i-- {
633 ev, err := e.GetBySerial(result[i])
634 if err != nil || ev == nil {
635 continue
636 }
637 out = mxutil.Ensure(out, 1)
638 out = push(out, ev)
639 if limit > 0 && len(out) >= limit {
640 break
641 }
642 }
643 return out
644 }
645
646 func intersectU64(a, b []uint64) (ss []uint64) {
647 set := map[uint64]bool{}
648 for _, v := range a {
649 set[v] = true
650 }
651 var out []uint64
652 for _, v := range b {
653 if set[v] {
654 out = mxutil.Ensure(out, 1)
655 out = push(out, v)
656 }
657 }
658 return out
659 }
660
661 // GetBySerial reads an event from the WAL by its serial. The event is handed
662 // back as a deep copy: UnmarshalBinary fills receiver fields in the decoding
663 // frame, and a receiver write that is not part of a return leaves the decoded
664 // tag fields in that frame. Clone builds the value as a return value, so the
665 // relocation moves the whole graph into the caller's arena.
666 func (e *Engine) GetBySerial(ser uint64) (ev *event.E, err error) {
667 data, err := e.w.Read(ser)
668 if err != nil {
669 return nil, err
670 }
671 tmp := event.New()
672 if uerr := tmp.UnmarshalBinary(bytes.NewReader(data)); uerr != nil {
673 return nil, uerr
674 }
675 ev = tmp.Clone()
676 return
677 }
678
679 // QueryEvents returns events matching the filter, sorted newest-first.
680 func (e *Engine) QueryEvents(f *filter.F) (evs []*event.E, err error) {
681 // Make records written since the last checkpoint visible to every reader,
682 // not just this handle's own buffer. This is the only receiver write on the
683 // path, so it stays in the method: a free function receives e as an
684 // immutable borrow and cannot call a mutating method on it.
685 e.flushPending()
686 return queryEvents(e, f)
687 }
688
689 // queryEvents is a free function: the scan serials, the dedup map, every
690 // per-candidate GetBySerial return and the result list die with its frame.
691 // Written in QueryEvents itself they would land in the Engine's sovereign data
692 // arena - a method that writes the receiver allocates there - and stay until
693 // compaction, on the relay's hottest path.
694 func queryEvents(e *Engine, f *filter.F) (evs []*event.E, qerr error) {
695 if f.Ids != nil && f.Ids.Len() > 0 {
696 return e.queryByIDs(f)
697 }
698
699 var since, until int64
700 if f.Since != nil && f.Since.I64() != 0 {
701 since = f.Since.I64()
702 }
703 if f.Until != nil && f.Until.I64() != 0 {
704 until = f.Until.I64()
705 } else {
706 until = math.MaxInt64
707 }
708
709 hasKinds := f.Kinds != nil && f.Kinds.Len() > 0
710 hasAuthors := f.Authors != nil && f.Authors.Len() > 0
711 hasTags := f.Tags != nil && f.Tags.Len() > 0
712
713 var serials []uint64
714 switch {
715 case hasTags && hasKinds && hasAuthors:
716 serials = e.scanTKP(f, since, until)
717 case hasTags && hasKinds:
718 serials = e.scanTKC(f, since, until)
719 case hasTags && hasAuthors:
720 serials = e.scanTPC(f, since, until)
721 case hasTags:
722 serials = e.scanTC(f, since, until)
723 case hasKinds && hasAuthors:
724 serials = e.scanKPC(f, since, until)
725 case hasKinds:
726 serials = e.scanKC(f, since, until)
727 case hasAuthors:
728 serials = e.scanPC(f, since, until)
729 default:
730 serials = e.scanCA(since, until)
731 }
732
733 // Deduplicate.
734 seen := map[uint64]bool{}
735 deduped := serials[:0]
736 for _, s := range serials {
737 if !seen[s] {
738 seen[s] = true
739 deduped = mxutil.Ensure(deduped, 1)
740 deduped = push(deduped, s)
741 }
742 }
743
744 // Fetch and apply the full filter. Newest-first ordering is applied by the
745 // caller when it frames the events.
746 //
747 // Index order is ascending by (..., created_at), so walking the serials
748 // backwards visits the newest first. When the filter carries a limit that
749 // is the order the caller wants, and there is no reason to decode every
750 // older match: with the scan no longer capped, a limited REQ over a store
751 // holding thousands of events decoded all of them and the relay's clients
752 // timed out. Filters without a limit keep the forward walk.
753 var results []*event.E
754 var limit int32
755 if f.Limit != nil && *f.Limit > 0 {
756 if *f.Limit > 0x7fffffff {
757 limit = 0x7fffffff
758 } else {
759 limit = int32(*f.Limit)
760 }
761 }
762 start, end, step := int32(0), int32(len(deduped)), int32(1)
763 if limit > 0 {
764 start, end, step = int32(len(deduped))-1, -1, -1
765 }
766 for i := start; i != end; i += step {
767 ev, err := e.GetBySerial(deduped[i])
768 if err != nil {
769 continue
770 }
771 if !f.Matches(ev) {
772 continue
773 }
774 results = mxutil.Ensure(results, 1)
775 results = push(results, ev)
776 if limit > 0 && int32(len(results)) >= limit {
777 break
778 }
779 }
780
781 sort.Sort(&event.S{E: results}) // reverse chronological
782
783 // The result limit is applied by the caller when it frames the events
784 // (cfg.QueryResultLimit in the relay), so a store-side reslice is
785 // redundant. It is also the one construct here that miscompiles:
786 // results[:*f.Limit] returned an empty slice to the caller even with
787 // matches in hand.
788 return results, nil
789 }
790
791 // DeleteEvent removes an event and its primary index entries.
792 func (e *Engine) DeleteEvent(id []byte) (err error) {
793 idHash := serial.IdHash(id)
794
795 // Find the serial for this event. The eid key is prefix|idHash|serial, so
796 // the first record under the id-hash prefix is the event and the serial
797 // follows in the record. This is a direct lookup rather than a scan
798 // callback: a callback's captured state does not reliably reach the
799 // enclosing frame, which silently turned every delete into "event not
800 // found" and left replaceable events to accumulate.
801 rec, ok := e.eid.GetPrefix(index.MakeEidPrefix(idHash))
802 if !ok {
803 return fmt.Errorf("event not found")
804 }
805 ser := serial.Get(rec[3+serial.HashLen:])
806
807 // Read the event to get kind/pubkey/timestamp for index cleanup.
808 ev, err := e.GetBySerial(ser)
809 if err != nil {
810 return err
811 }
812
813 pubHash := serial.PubHash(ev.Pubkey)
814
815 // Delete primary index entries.
816 e.eid.Delete(index.MakeEid(idHash, ser))
817 e.ca.Delete(index.MakeCA(ev.CreatedAt, ser))
818 e.kc.Delete(index.MakeKC(ev.Kind, ev.CreatedAt, ser))
819 e.pc.Delete(index.MakePC(pubHash, ev.CreatedAt, ser))
820 e.kpc.Delete(index.MakeKPC(ev.Kind, pubHash, ev.CreatedAt, ser))
821
822 return nil
823 }
824
825 // ResolveBlob returns every stored event that references the blob hash through
826 // an `x` tag, newest first. The tc index already holds a record per
827 // single-letter tag, so this is the same key-range scan a {"#x": [...]} filter
828 // runs - no new index family, no on-disk change.
829 //
830 // The scan and the seen map are scratch, not object state, so they live in a
831 // free function that gets its own frame arena; the method itself only returns.
832 // This is the queryEvents/QueryEvents split for the same reason.
833 func (e *Engine) ResolveBlob(xHex []byte) (out []*event.E) {
834 return resolveBlob(e, xHex)
835 }
836
837 func resolveBlob(e *Engine, xHex []byte) (out []*event.E) {
838 if len(xHex) == 0 {
839 return nil
840 }
841 vh := serial.Ident(xHex)
842 serials := e.collectSerials(e.tc,
843 index.MakeTC('x', vh, 0, 0),
844 index.MakeTC('x', vh, math.MaxInt64, serial.Max),
845 index.TCKeyLen)
846 if len(serials) == 0 {
847 return nil
848 }
849 // Index order is ascending by (created_at, serial), so walking it
850 // backwards visits the newest event first. One event can carry the same x
851 // tag twice and contribute two records; dedup by serial.
852 seen := map[uint64]bool{}
853 for i := int32(len(serials)) - 1; i >= 0; i-- {
854 ser := serials[i]
855 if seen[ser] {
856 continue
857 }
858 seen[ser] = true
859 ev, err := e.GetBySerial(ser)
860 if err != nil {
861 continue
862 }
863 out = mxutil.Ensure(out, 1)
864 out = push(out, ev)
865 }
866 return
867 }
868
869 // GetByID retrieves an event by its 32-byte ID.
870 func (e *Engine) GetByID(id []byte) (ev *event.E, err error) {
871 // Point lookup, exactly as queryByIDs does it: the eid index compares only
872 // prefix|id-hash (EidCmpLen), so the 16-byte range
873 // Scan(MakeEid(hash,0), MakeEid(hash,Max)) this used was an empty range and
874 // GetByID answered "event not found" for every event that was stored.
875 // negentropy's FindHave calls this and dropped every event on the error,
876 // so a sync found nothing to send.
877 ser, ok := e.getEventSerial(id)
878 if !ok {
879 return nil, fmt.Errorf("event not found")
880 }
881 ev, err := e.GetBySerial(ser)
882 if err != nil {
883 return nil, err
884 }
885 return ev, nil
886 }
887
888 // --- internal helpers ---
889
890 func (e *Engine) hasEvent(idHash, fullID []byte) (ok bool) {
891 // Point lookup: the eid key is prefix|id-hash|serial and cmpLen stops at
892 // the hash, so Get returns the first record for this hash and the full ID
893 // check filters hash collisions. A range Scan here cost a Stat, a whole
894 // .buf sidecar read and a 256-record chunk buffer per event.
895 rec, found := e.eid.Get(index.MakeEid(idHash, 0))
896 if !found {
897 return false
898 }
899 ser := serial.Get(rec[3+serial.HashLen:])
900 eid, hasID := e.getEventIDBySerial(ser)
901 if !hasID {
902 return false
903 }
904 return bytes.Equal(eid, fullID)
905 }
906
907 func (e *Engine) getOrCreatePubkeySerial(pubkey []byte) (n uint64) {
908 pubHash := serial.PubHash(pubkey)
909 rec, exists := e.pks.Get(index.MakePks(pubHash, 0))
910 if exists {
911 return serial.Get(rec[3+serial.HashLen:])
912 }
913 ser := e.nextPubSer
914 e.nextPubSer++
915 e.pks.Put(index.MakePks(pubHash, ser))
916 e.spk.Put(index.MakeSpkRec(ser, pubkey))
917 return ser
918 }
919
920 func (e *Engine) getEventIDBySerial(ser uint64) (id []byte, ok bool) {
921 key := index.MakeSei(ser)
922 rec, ok := e.sei.Get(key)
923 if !ok {
924 return nil, false
925 }
926 return rec[index.SeiKeyLen:], true
927 }
928
929 func (e *Engine) queryByIDs(f *filter.F) (evs []*event.E, qerr error) {
930 var results []*event.E
931 for _, id := range f.Ids.T {
932 // Point lookup, not a range Scan: the eid index compares only
933 // prefix|id-hash (EidCmpLen), so two records for the same hash compare
934 // equal and Scan(MakeEid(hash,0), MakeEid(hash,Max)) is an empty
935 // range - every by-id query came back with no events. getEventSerial
936 // does the Get and verifies the full 32-byte id, which also filters
937 // hash collisions.
938 ser, ok := e.getEventSerial(id)
939 if !ok {
940 continue
941 }
942 if ev, err := e.GetBySerial(ser); err == nil {
943 results = mxutil.Ensure(results, 1)
944 results = push(results, ev)
945 }
946 }
947 return results, nil
948 }
949
950 func (e *Engine) collectSerials(idx *sorted.File, start, end []byte, keyLen int32) (ss []uint64) {
951 // No allocation inside the callback: its frame arena dies per call (see
952 // queryByIDs). A buffer plus a count keeps the scan allocation-free.
953 // The scan callback is a separate function: only a single captured pointer
954 // reliably reaches the enclosing frame, so writes to a captured local array
955 // were lost (the counter advanced, every value read back zero). Keep the
956 // collected state behind a pointer.
957 //
958 // The buffer is sized from the index, not fixed: a fixed sink silently
959 // stopped the scan at its array size, so once a store held more records
960 // than that, every kinds/authors/created-at query missed the records that
961 // sorted after the first batch - a freshly published event was invisible
962 // to every scan while the by-id lookup, whose range holds a handful of
963 // records, still found it.
964 total := idx.Count()
965 if total <= 0 {
966 return nil
967 }
968 if total > math.MaxInt32 {
969 total = math.MaxInt32
970 }
971 state := &serialSink{buf: []uint64{:int32(total)}, keyLen: keyLen}
972 idx.Scan(start, end, func(rec []byte) bool { return state.add(rec) })
973 n := state.n
974 ss = []uint64{:n}
975 for i := int32(0); i < n; i++ {
976 ss[i] = state.buf[i]
977 }
978 return ss
979 }
980
981 // serialSink collects serials from an index scan.
982 type serialSink struct {
983 buf []uint64
984 n int32
985 keyLen int32
986 }
987
988 func (s *serialSink) add(rec []byte) (ok bool) {
989 if s.n >= int32(len(s.buf)) {
990 return false
991 }
992 s.buf[s.n] = serial.Get(rec[s.keyLen-serial.Len:])
993 s.n++
994 return true
995 }
996
997 func (e *Engine) scanCA(since, until int64) (ss []uint64) {
998 return e.collectSerials(e.ca,
999 index.MakeCA(since, 0),
1000 index.MakeCA(until, serial.Max),
1001 index.CAKeyLen)
1002 }
1003
1004 func (e *Engine) scanKC(f *filter.F, since, until int64) (ss []uint64) {
1005 var out []uint64
1006 for _, k := range f.Kinds.K {
1007 got := e.collectSerials(e.kc,
1008 index.MakeKC(k.K, since, 0),
1009 index.MakeKC(k.K, until, serial.Max),
1010 index.KCKeyLen)
1011 out = out | got
1012 }
1013 return out
1014 }
1015
1016 func (e *Engine) scanPC(f *filter.F, since, until int64) (ss []uint64) {
1017 var out []uint64
1018 for _, author := range f.Authors.T {
1019 ph := serial.PubHash(author)
1020 out = out | e.collectSerials(e.pc,
1021 index.MakePC(ph, since, 0),
1022 index.MakePC(ph, until, serial.Max),
1023 index.PCKeyLen)
1024 }
1025 return out
1026 }
1027
1028 func (e *Engine) scanKPC(f *filter.F, since, until int64) (ss []uint64) {
1029 var out []uint64
1030 for _, k := range f.Kinds.K {
1031 for _, author := range f.Authors.T {
1032 ph := serial.PubHash(author)
1033 out = out | e.collectSerials(e.kpc,
1034 index.MakeKPC(k.K, ph, since, 0),
1035 index.MakeKPC(k.K, ph, until, serial.Max),
1036 index.KPCKeyLen)
1037 }
1038 }
1039 return out
1040 }
1041
1042 func (e *Engine) scanTC(f *filter.F, since, until int64) (ss []uint64) {
1043 var out []uint64
1044 for _, tg := range f.Tags.T {
1045 if tg.Len() < 2 || len(tg.Key()) != 1 {
1046 continue
1047 }
1048 tk := tg.Key()[0]
1049 for _, val := range tg.T[1:] {
1050 vh := serial.Ident(val)
1051 out = out | e.collectSerials(e.tc,
1052 index.MakeTC(tk, vh, since, 0),
1053 index.MakeTC(tk, vh, until, serial.Max),
1054 index.TCKeyLen)
1055 }
1056 }
1057 return out
1058 }
1059
1060 func (e *Engine) scanTKC(f *filter.F, since, until int64) (ss []uint64) {
1061 var out []uint64
1062 for _, tg := range f.Tags.T {
1063 if tg.Len() < 2 || len(tg.Key()) != 1 {
1064 continue
1065 }
1066 tk := tg.Key()[0]
1067 for _, val := range tg.T[1:] {
1068 vh := serial.Ident(val)
1069 for _, k := range f.Kinds.K {
1070 out = out | e.collectSerials(e.tkc,
1071 index.MakeTKC(k.K, tk, vh, since, 0),
1072 index.MakeTKC(k.K, tk, vh, until, serial.Max),
1073 index.TKCKeyLen)
1074 }
1075 }
1076 }
1077 return out
1078 }
1079
1080 func (e *Engine) scanTPC(f *filter.F, since, until int64) (ss []uint64) {
1081 var out []uint64
1082 for _, tg := range f.Tags.T {
1083 if tg.Len() < 2 || len(tg.Key()) != 1 {
1084 continue
1085 }
1086 tk := tg.Key()[0]
1087 for _, val := range tg.T[1:] {
1088 vh := serial.Ident(val)
1089 for _, author := range f.Authors.T {
1090 ph := serial.PubHash(author)
1091 out = out | e.collectSerials(e.tpc,
1092 index.MakeTPC(ph, tk, vh, since, 0),
1093 index.MakeTPC(ph, tk, vh, until, serial.Max),
1094 index.TPCKeyLen)
1095 }
1096 }
1097 }
1098 return out
1099 }
1100
1101 func (e *Engine) scanTKP(f *filter.F, since, until int64) (ss []uint64) {
1102 var out []uint64
1103 for _, tg := range f.Tags.T {
1104 if tg.Len() < 2 || len(tg.Key()) != 1 {
1105 continue
1106 }
1107 tk := tg.Key()[0]
1108 for _, val := range tg.T[1:] {
1109 vh := serial.Ident(val)
1110 for _, k := range f.Kinds.K {
1111 for _, author := range f.Authors.T {
1112 ph := serial.PubHash(author)
1113 out = out | e.collectSerials(e.tkp,
1114 index.MakeTKP(k.K, ph, tk, vh, since, 0),
1115 index.MakeTKP(k.K, ph, tk, vh, until, serial.Max),
1116 index.TKPKeyLen)
1117 }
1118 }
1119 }
1120 }
1121 return out
1122 }
1123
1124 // --- Graph traversal ---
1125
1126 // GetReferencingEvents returns serials of events that reference targetSer via e-tags.
1127 // Uses the gee (reverse event-event) index.
1128 func (e *Engine) GetReferencingEvents(targetSer uint64) (ss []uint64) {
1129 start := index.MakeGee(targetSer, 0, 0, 0)
1130 end := index.MakeGee(targetSer, 0xFFFF, 0xFF, serial.Max)
1131 var out []uint64
1132 e.gee.Scan(start, end, func(rec []byte) bool {
1133 out = mxutil.Ensure(out, 1)
1134 out = push(out, serial.Get(rec[11:]))
1135 return true
1136 })
1137 return out
1138 }
1139
1140 // GetETagTargets returns serials of events that srcSer references via e-tags.
1141 // Uses the eeg (forward event-event) index.
1142 func (e *Engine) GetETagTargets(srcSer uint64) (ss []uint64) {
1143 start := index.MakeEeg(srcSer, 0, 0, 0)
1144 end := index.MakeEeg(srcSer, serial.Max, 0xFFFF, 0xFF)
1145 var out []uint64
1146 e.eeg.Scan(start, end, func(rec []byte) bool {
1147 out = mxutil.Ensure(out, 1)
1148 out = push(out, serial.Get(rec[8:]))
1149 return true
1150 })
1151 return out
1152 }
1153
1154 // TraverseThread performs BFS traversal of thread structure via e-tags.
1155 // Returns all event IDs (32-byte binary) reachable from seedID within maxDepth.
1156 func (e *Engine) TraverseThread(seedID []byte, maxDepth int32, direction string) (ss [][]byte) {
1157 seedSer, ok := e.getEventSerial(seedID)
1158 if !ok {
1159 return nil
1160 }
1161 if direction == "" {
1162 direction = "both"
1163 }
1164 return traverseThread(e, seedSer, maxDepth, direction)
1165 }
1166
1167 // traverseThread is a free function: the breadth-first scratch - the visited
1168 // set, each frontier slice, and every GetReferencingEvents/GetETagTargets
1169 // return - dies with its frame instead of accumulating in the Engine's
1170 // sovereign data arena until the next compaction.
1171 func traverseThread(e *Engine, seedSer uint64, maxDepth int32, direction string) (allIDs [][]byte) {
1172 visited := map[uint64]bool{}
1173 visited[seedSer] = true
1174 frontier := []uint64{seedSer}
1175
1176 emptyStreak := 0
1177
1178 for depth := 1; depth <= maxDepth; depth++ {
1179 var next []uint64
1180 found := 0
1181
1182 for _, evSer := range frontier {
1183 if direction == "both" || direction == "inbound" {
1184 for _, refSer := range e.GetReferencingEvents(evSer) {
1185 if visited[refSer] {
1186 continue
1187 }
1188 visited[refSer] = true
1189 if eid, hasID := e.getEventIDBySerial(refSer); hasID {
1190 allIDs = mxutil.Ensure(allIDs, 1)
1191 allIDs = push(allIDs, eid)
1192 found++
1193 }
1194 next = mxutil.Ensure(next, 1)
1195 next = push(next, refSer)
1196 }
1197 }
1198 if direction == "both" || direction == "outbound" {
1199 for _, tgtSer := range e.GetETagTargets(evSer) {
1200 if visited[tgtSer] {
1201 continue
1202 }
1203 visited[tgtSer] = true
1204 if eid, hasID := e.getEventIDBySerial(tgtSer); hasID {
1205 allIDs = mxutil.Ensure(allIDs, 1)
1206 allIDs = push(allIDs, eid)
1207 found++
1208 }
1209 next = mxutil.Ensure(next, 1)
1210 next = push(next, tgtSer)
1211 }
1212 }
1213 }
1214
1215 if found == 0 {
1216 emptyStreak++
1217 if emptyStreak >= 2 {
1218 break
1219 }
1220 } else {
1221 emptyStreak = 0
1222 }
1223 frontier = next
1224 }
1225 return allIDs
1226 }
1227