memory.go raw
1 package memory
2
3 import (
4 "bytes"
5 "encoding/json"
6 "time"
7
8 "git.mleku.dev/mleku/dendrite/pkg/ratio"
9 "git.mleku.dev/mleku/dendrite/pkg/spore"
10 "github.com/dgraph-io/badger/v4"
11 "github.com/dgraph-io/badger/v4/options"
12 )
13
14 // DB is the persistent cross-generation memory store.
15 type DB struct {
16 db *badger.DB
17 }
18
19 // Open creates or opens a memory database at the given directory.
20 func Open(dir string) (*DB, error) {
21 opts := badger.DefaultOptions(dir)
22 opts.CompactL0OnClose = true
23 opts.LmaxCompaction = true
24 opts.Compression = options.None
25 opts.Logger = nil // silent
26 opts.MemTableSize = 16 << 20 // 16MB memtable (default 64MB)
27 opts.ValueLogFileSize = 64 << 20 // 64MB vlog files (default 1GB)
28 opts.NumMemtables = 2 // reduce from default 5
29 opts.NumLevelZeroTables = 2 // reduce from default 5
30 opts.NumLevelZeroTablesStall = 5 // reduce from default 15
31 db, err := badger.Open(opts)
32 if err != nil {
33 return nil, err
34 }
35 return &DB{db: db}, nil
36 }
37
38 // Close closes the database.
39 func (d *DB) Close() error {
40 if d.db == nil {
41 return nil
42 }
43 return d.db.Close()
44 }
45
46 // genMeta is the JSON value stored in generation records.
47 type genMeta struct {
48 Timestamp int64 `json:"ts"`
49 ParentHash string `json:"parent,omitempty"`
50 InstanceID uint32 `json:"inst"`
51 ConfigHash string `json:"cfg,omitempty"`
52 }
53
54 // RecordGeneration writes generation metadata.
55 func (d *DB) RecordGeneration(gen uint32, parentHash string, instanceID uint32, ts time.Time) error {
56 meta := genMeta{
57 Timestamp: ts.UnixMilli(),
58 ParentHash: parentHash,
59 InstanceID: instanceID,
60 }
61 val, err := json.Marshal(meta)
62 if err != nil {
63 return err
64 }
65 return d.db.Update(func(txn *badger.Txn) error {
66 return txn.Set(GenKey(gen), val)
67 })
68 }
69
70 // BondRecord is a single bond event: an element type bonded at a site.
71 type BondRecord struct {
72 Tag string
73 SiteID uint32
74 }
75
76 // RecordBonds writes bond events for a generation.
77 // Uses nil values — the key IS the index.
78 func (d *DB) RecordBonds(gen uint32, bonds []BondRecord) error {
79 if len(bonds) == 0 {
80 return nil
81 }
82 return d.db.Update(func(txn *badger.Txn) error {
83 for _, b := range bonds {
84 h := TagHash(b.Tag)
85 key := BndKey(h, gen, b.SiteID)
86 if err := txn.Set(key, nil); err != nil {
87 return err
88 }
89 }
90 return nil
91 })
92 }
93
94 // RecordMissing writes missing site records (negative space) for a generation.
95 func (d *DB) RecordMissing(gen uint32, missing []spore.TagCount) error {
96 if len(missing) == 0 {
97 return nil
98 }
99 return d.db.Update(func(txn *badger.Txn) error {
100 for _, m := range missing {
101 h := TagHash(m.Tag)
102 key := MisKey(h, gen, uint32(m.Count))
103 if err := txn.Set(key, nil); err != nil {
104 return err
105 }
106 }
107 return nil
108 })
109 }
110
111 // RecordTypeSig writes type signature snapshot for a generation.
112 func (d *DB) RecordTypeSig(gen uint32, typeSig []spore.TagCount) error {
113 if len(typeSig) == 0 {
114 return nil
115 }
116 return d.db.Update(func(txn *badger.Txn) error {
117 for _, tc := range typeSig {
118 h := TagHash(tc.Tag)
119 key := TypKey(h, uint32(tc.Count), gen)
120 if err := txn.Set(key, nil); err != nil {
121 return err
122 }
123 }
124 return nil
125 })
126 }
127
128 // RecordConnectivity writes per-type connectivity for a generation.
129 func (d *DB) RecordConnectivity(gen uint32, connectivity []spore.TagRatio) error {
130 if len(connectivity) == 0 {
131 return nil
132 }
133 return d.db.Update(func(txn *badger.Txn) error {
134 for _, c := range connectivity {
135 h := TagHash(c.Tag)
136 key := ConKey(h, gen)
137 val := EncodeConValue(c.Value.Num, c.Value.Denom)
138 if err := txn.Set(key, val); err != nil {
139 return err
140 }
141 }
142 return nil
143 })
144 }
145
146 // RecordHealth writes a health snapshot for a generation.
147 func (d *DB) RecordHealth(gen uint32, occupied, total uint32, avgLockIn ratio.Ratio) error {
148 key := HltKey(gen)
149 val := EncodeHltValue(occupied, total, avgLockIn.Num, avgLockIn.Denom)
150 return d.db.Update(func(txn *badger.Txn) error {
151 return txn.Set(key, val)
152 })
153 }
154
155 // RecordLockInDist writes lock-in depth distribution for a generation.
156 // buckets maps quantized bucket (0-255) → count.
157 func (d *DB) RecordLockInDist(gen uint32, buckets map[byte]uint32) error {
158 if len(buckets) == 0 {
159 return nil
160 }
161 return d.db.Update(func(txn *badger.Txn) error {
162 for bucket, count := range buckets {
163 key := LckKey(bucket, gen)
164 val := EncodeU32Value(count)
165 if err := txn.Set(key, val); err != nil {
166 return err
167 }
168 }
169 return nil
170 })
171 }
172
173 // RecordHexagramOps writes hexagram operation counts for a generation.
174 // ops maps operation code (0-7) → count.
175 func (d *DB) RecordHexagramOps(gen uint32, ops map[byte]uint32) error {
176 if len(ops) == 0 {
177 return nil
178 }
179 return d.db.Update(func(txn *badger.Txn) error {
180 for op, count := range ops {
181 key := HexKey(op, gen)
182 val := EncodeU32Value(count)
183 if err := txn.Set(key, val); err != nil {
184 return err
185 }
186 }
187 return nil
188 })
189 }
190
191 // RecordMindsicle stores a mindsicle (frozen lattice) for a generation.
192 func (d *DB) RecordMindsicle(gen uint32, data []byte) error {
193 key := MndKey(gen)
194 return d.db.Update(func(txn *badger.Txn) error {
195 return txn.Set(key, data)
196 })
197 }
198
199 // LoadMindsicle retrieves a mindsicle for a specific generation.
200 func (d *DB) LoadMindsicle(gen uint32) ([]byte, error) {
201 key := MndKey(gen)
202 var val []byte
203 err := d.db.View(func(txn *badger.Txn) error {
204 item, err := txn.Get(key)
205 if err != nil {
206 return err
207 }
208 val, err = item.ValueCopy(nil)
209 return err
210 })
211 return val, err
212 }
213
214 // LatestMindsicle returns the highest generation for which a mindsicle exists,
215 // along with its data. Returns (0, nil, ErrKeyNotFound) if none.
216 func (d *DB) LatestMindsicle() (uint32, []byte, error) {
217 var gen uint32
218 var data []byte
219 err := d.db.View(func(txn *badger.Txn) error {
220 opts := badger.DefaultIteratorOptions
221 opts.Prefix = PrefixMnd[:]
222 opts.Reverse = true
223 it := txn.NewIterator(opts)
224 defer it.Close()
225 // Seek to the end of the mnd prefix range.
226 it.Seek(PrefixEnd(PrefixMnd[:]))
227 if !it.Valid() {
228 // Try seeking to just the prefix for the first item.
229 it.Rewind()
230 if !it.Valid() {
231 return badger.ErrKeyNotFound
232 }
233 }
234 item := it.Item()
235 key := item.Key()
236 if len(key) < 7 || key[0] != PrefixMnd[0] || key[1] != PrefixMnd[1] || key[2] != PrefixMnd[2] {
237 return badger.ErrKeyNotFound
238 }
239 gen = DecodeGen(key)
240 var err error
241 data, err = item.ValueCopy(nil)
242 return err
243 })
244 return gen, data, err
245 }
246
247 // RecordEWMAState stores serialized EWMA detector state for a generation.
248 func (d *DB) RecordEWMAState(gen uint32, state []byte) error {
249 key := EwmKey(gen)
250 return d.db.Update(func(txn *badger.Txn) error {
251 return txn.Set(key, state)
252 })
253 }
254
255 // LoadLatestEWMAState returns the highest generation for which EWMA state
256 // exists, along with its data. Returns (0, nil, ErrKeyNotFound) if none.
257 func (d *DB) LoadLatestEWMAState() (uint32, []byte, error) {
258 var gen uint32
259 var data []byte
260 err := d.db.View(func(txn *badger.Txn) error {
261 opts := badger.DefaultIteratorOptions
262 opts.Prefix = PrefixEwm[:]
263 opts.Reverse = true
264 it := txn.NewIterator(opts)
265 defer it.Close()
266 it.Seek(PrefixEnd(PrefixEwm[:]))
267 if !it.Valid() {
268 it.Rewind()
269 if !it.Valid() {
270 return badger.ErrKeyNotFound
271 }
272 }
273 item := it.Item()
274 key := item.Key()
275 if len(key) < 7 || key[0] != PrefixEwm[0] || key[1] != PrefixEwm[1] || key[2] != PrefixEwm[2] {
276 return badger.ErrKeyNotFound
277 }
278 gen = DecodeGen(key)
279 var err error
280 data, err = item.ValueCopy(nil)
281 return err
282 })
283 return gen, data, err
284 }
285
286 // RecordADSR writes the ADSR phase distribution for a generation.
287 func (d *DB) RecordADSR(gen uint32, counts [4]uint32) error {
288 key := AdrKey(gen)
289 val := EncodeAdrValue(counts)
290 return d.db.Update(func(txn *badger.Txn) error {
291 return txn.Set(key, val)
292 })
293 }
294
295 // RecordFitness writes fitness scores for a generation.
296 func (d *DB) RecordFitness(gen uint32, source, binary, behav, overall ratio.Ratio) error {
297 return d.db.Update(func(txn *badger.Txn) error {
298 dims := []struct {
299 dim byte
300 r ratio.Ratio
301 }{
302 {DimSource, source},
303 {DimBinary, binary},
304 {DimBehav, behav},
305 {DimOverall, overall},
306 }
307 for _, d := range dims {
308 key := FitKey(d.dim, gen)
309 val := EncodeFitValue(d.r.Num, d.r.Denom)
310 if err := txn.Set(key, val); err != nil {
311 return err
312 }
313 }
314 return nil
315 })
316 }
317
318 // WalkerCheckpoint is the JSON value stored in the singleton wlk key.
319 // Contains everything needed to resume the walker at the exact position.
320 type WalkerCheckpoint struct {
321 Epoch uint32 `json:"epoch"`
322 Seed uint64 `json:"seed"`
323 Position int `json:"position"`
324 GenNum uint32 `json:"gen_num"`
325 Files []string `json:"files"`
326 Root string `json:"root"`
327 }
328
329 // RecordFileScore accumulates raw and accreted counts for a file.
330 // Uses read-modify-write within a single transaction.
331 func (d *DB) RecordFileScore(filePath string, rawDelta, accretedDelta int64) error {
332 h := TagHash(filePath)
333 key := FscKey(h)
334 return d.db.Update(func(txn *badger.Txn) error {
335 var prevAccreted, prevRaw int64
336 item, err := txn.Get(key)
337 if err == nil {
338 _ = item.Value(func(val []byte) error {
339 prevAccreted, prevRaw = DecodeFitValue(val)
340 return nil
341 })
342 }
343 return txn.Set(key, EncodeFitValue(prevAccreted+accretedDelta, prevRaw+rawDelta))
344 })
345 }
346
347 // LoadFileScores loads all per-file accretion scores.
348 // Returns a map of file path hash → [accreted, raw].
349 func (d *DB) LoadFileScores() (map[[8]byte][2]int64, error) {
350 scores := make(map[[8]byte][2]int64)
351 prefix := PrefixFsc[:]
352 err := d.db.View(func(txn *badger.Txn) error {
353 opts := badger.DefaultIteratorOptions
354 opts.Prefix = prefix
355 it := txn.NewIterator(opts)
356 defer it.Close()
357 for it.Seek(prefix); it.Valid(); it.Next() {
358 item := it.Item()
359 key := item.Key()
360 if !bytes.HasPrefix(key, prefix) {
361 break
362 }
363 if len(key) < 11 {
364 continue
365 }
366 var h [8]byte
367 copy(h[:], key[3:11])
368 _ = item.Value(func(val []byte) error {
369 accreted, raw := DecodeFitValue(val)
370 scores[h] = [2]int64{accreted, raw}
371 return nil
372 })
373 }
374 return nil
375 })
376 return scores, err
377 }
378
379 // RecordWalkerCheckpoint persists the walker state for exact resume.
380 func (d *DB) RecordWalkerCheckpoint(cp WalkerCheckpoint) error {
381 val, err := json.Marshal(cp)
382 if err != nil {
383 return err
384 }
385 return d.db.Update(func(txn *badger.Txn) error {
386 return txn.Set(WlkKey(), val)
387 })
388 }
389
390 // LoadWalkerCheckpoint loads the persisted walker checkpoint.
391 func (d *DB) LoadWalkerCheckpoint() (WalkerCheckpoint, error) {
392 var cp WalkerCheckpoint
393 err := d.db.View(func(txn *badger.Txn) error {
394 item, err := txn.Get(WlkKey())
395 if err != nil {
396 return err
397 }
398 return item.Value(func(val []byte) error {
399 return json.Unmarshal(val, &cp)
400 })
401 })
402 return cp, err
403 }
404
405 // RecordOracleState stores the current oracle state (JSON), overwriting any previous.
406 // Uses a singleton key — only the latest state matters; per-gen history is in ohx keys.
407 func (d *DB) RecordOracleState(data []byte) error {
408 key := OrcKey()
409 return d.db.Update(func(txn *badger.Txn) error {
410 return txn.Set(key, data)
411 })
412 }
413
414 // LoadOracleState returns the current oracle state, or ErrKeyNotFound if none exists.
415 func (d *DB) LoadOracleState() ([]byte, error) {
416 var data []byte
417 err := d.db.View(func(txn *badger.Txn) error {
418 item, err := txn.Get(OrcKey())
419 if err != nil {
420 return err
421 }
422 data, err = item.ValueCopy(nil)
423 return err
424 })
425 return data, err
426 }
427
428 // RecordOracleReading stores an individual reading by sequence and generation.
429 func (d *DB) RecordOracleReading(seq, gen uint32, data []byte) error {
430 key := OhxKey(seq, gen)
431 return d.db.Update(func(txn *badger.Txn) error {
432 return txn.Set(key, data)
433 })
434 }
435
436 // LoadOracleHistory loads the last N oracle readings, ordered by sequence.
437 func (d *DB) LoadOracleHistory(lastN int) ([][]byte, error) {
438 var results [][]byte
439 err := d.db.View(func(txn *badger.Txn) error {
440 opts := badger.DefaultIteratorOptions
441 opts.Prefix = PrefixOhx[:]
442 opts.Reverse = true
443 it := txn.NewIterator(opts)
444 defer it.Close()
445 it.Seek(PrefixEnd(PrefixOhx[:]))
446 count := 0
447 for ; it.Valid() && count < lastN; it.Next() {
448 item := it.Item()
449 key := item.Key()
450 if !bytes.HasPrefix(key, PrefixOhx[:]) {
451 break
452 }
453 val, err := item.ValueCopy(nil)
454 if err != nil {
455 return err
456 }
457 results = append(results, val)
458 count++
459 }
460 return nil
461 })
462 // Reverse to get chronological order (earliest first).
463 for i, j := 0, len(results)-1; i < j; i, j = i+1, j-1 {
464 results[i], results[j] = results[j], results[i]
465 }
466 return results, err
467 }
468
469 // --- Recognition methods ---
470
471 // RecordProfile stores a recognition profile snapshot for a generation.
472 func (d *DB) RecordProfile(gen uint32, data []byte) error {
473 key := PrfKey(gen)
474 return d.db.Update(func(txn *badger.Txn) error {
475 return txn.Set(key, data)
476 })
477 }
478
479 // LoadLatestProfile returns the highest generation for which a profile exists.
480 func (d *DB) LoadLatestProfile() (uint32, []byte, error) {
481 var gen uint32
482 var data []byte
483 err := d.db.View(func(txn *badger.Txn) error {
484 opts := badger.DefaultIteratorOptions
485 opts.Prefix = PrefixPrf[:]
486 opts.Reverse = true
487 it := txn.NewIterator(opts)
488 defer it.Close()
489 it.Seek(PrefixEnd(PrefixPrf[:]))
490 if !it.Valid() {
491 it.Rewind()
492 if !it.Valid() {
493 return badger.ErrKeyNotFound
494 }
495 }
496 item := it.Item()
497 key := item.Key()
498 if len(key) < 7 || key[0] != PrefixPrf[0] || key[1] != PrefixPrf[1] || key[2] != PrefixPrf[2] {
499 return badger.ErrKeyNotFound
500 }
501 gen = DecodeGen(key)
502 var err error
503 data, err = item.ValueCopy(nil)
504 return err
505 })
506 return gen, data, err
507 }
508
509 // RecordConvergence stores convergence state for a generation.
510 func (d *DB) RecordConvergence(gen uint32, data []byte) error {
511 key := CvgKey(gen)
512 return d.db.Update(func(txn *badger.Txn) error {
513 return txn.Set(key, data)
514 })
515 }
516
517 // RecordModelSpore stores a model fingerprint spore by name.
518 func (d *DB) RecordModelSpore(modelName string, data []byte) error {
519 h := TagHash(modelName)
520 key := MdlKey(h)
521 // Store the name alongside the data so we can recover it.
522 val := append([]byte(modelName+"\x00"), data...)
523 return d.db.Update(func(txn *badger.Txn) error {
524 return txn.Set(key, val)
525 })
526 }
527
528 // LoadModelSpore retrieves a model fingerprint spore by name.
529 func (d *DB) LoadModelSpore(modelName string) ([]byte, error) {
530 h := TagHash(modelName)
531 key := MdlKey(h)
532 var data []byte
533 err := d.db.View(func(txn *badger.Txn) error {
534 item, err := txn.Get(key)
535 if err != nil {
536 return err
537 }
538 val, err := item.ValueCopy(nil)
539 if err != nil {
540 return err
541 }
542 // Skip the name prefix.
543 idx := bytes.IndexByte(val, 0)
544 if idx >= 0 && idx < len(val)-1 {
545 data = val[idx+1:]
546 } else {
547 data = val
548 }
549 return nil
550 })
551 return data, err
552 }
553
554 // ListModelSpores returns the names of all stored model fingerprints.
555 func (d *DB) ListModelSpores() ([]string, error) {
556 var names []string
557 prefix := PrefixMdl[:]
558 err := d.db.View(func(txn *badger.Txn) error {
559 opts := badger.DefaultIteratorOptions
560 opts.Prefix = prefix
561 it := txn.NewIterator(opts)
562 defer it.Close()
563 for it.Seek(prefix); it.Valid(); it.Next() {
564 item := it.Item()
565 key := item.Key()
566 if !bytes.HasPrefix(key, prefix) {
567 break
568 }
569 _ = item.Value(func(val []byte) error {
570 idx := bytes.IndexByte(val, 0)
571 if idx > 0 {
572 names = append(names, string(val[:idx]))
573 }
574 return nil
575 })
576 }
577 return nil
578 })
579 return names, err
580 }
581
582 // RecordDirectiveResult stores the execution result of an oracle directive.
583 func (d *DB) RecordDirectiveResult(seq uint32, hash [8]byte, data []byte) error {
584 key := OsrKey(seq, hash)
585 return d.db.Update(func(txn *badger.Txn) error {
586 return txn.Set(key, data)
587 })
588 }
589