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