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