// package kstore_kv -- локальное быстрое key-value хранилище ядра. package kstore_kv import ( "os" "sync" "time" "github.com/dgraph-io/badger/v4" "gitp78su.ipnodns.ru/svi/kern/v4/d0" "gitp78su.ipnodns.ru/svi/kern/v4/d1" "gitp78su.ipnodns.ru/svi/kern/v4/d1/log_buf" "gitp78su.ipnodns.ru/svi/kern/v4/d2/kern_ctx" "gitp78su.ipnodns.ru/svi/kern/v4/d2/kern_ent/stream_name" "gitp78su.ipnodns.ru/svi/kern/v4/d2/kspec" "gitp78su.ipnodns.ru/svi/kern/v4/d2/lsi/store_ent" "gitp78su.ipnodns.ru/svi/kern/v4/d2/lsi/store_ent/rec_key" "gitp78su.ipnodns.ru/svi/kern/v4/d2/lsi/store_spec" ) var ( storeStreamName = stream_name.NewAStreamName("kstore_kv") // Имя потока для ожидателя потоков ) // kStoreKv -- локальное хранилище ядра. type kStoreKv struct { sync.RWMutex owner *d0.LOwner kCtx *kern_ctx.KernCtx lCtx *d1.LocalCtx log *d1.LogBuf wg kspec.IKernelWg storePath *d0.LTxt db *badger.DB isWork *d0.MSafeBool } var ( kernStore *kStoreKv // Глобальный объект block sync.Mutex ) // GetKernelStore -- возвращает новое локальное хранилище ядра. func GetKernelStore() store_spec.IKernelStoreKv { block.Lock() defer block.Unlock() if kernStore != nil { kernStore.log.Debug("GetKernelStore()") return kernStore } owner := d0.LetOwner(1, "isol_kern_store") kCtx := kern_ctx.GetKernCtx() param := &log_buf.LogBufParam{ Prefix_: d0.LetTxt("kStoreKv"), IsTerm_: d0.MutSafeBool(owner, true), } log := d1.NewLogBuf(param) sf := &kStoreKv{ owner: owner, kCtx: kCtx, lCtx: d1.NewLocalCtx(owner,kCtx.SelfCtx()), wg: kCtx.Wg(), isWork: d0.MutSafeBool(owner, true), log: log, } sf.open() kernStore = sf kCtx.Set(d0.LetTxt("kernStoreKV"), kernStore, "fast KV store on Badger") return kernStore } // Log -- возвращает локальный лог. func (sf *kStoreKv) Log() *d1.LogBuf { return sf.log } // Set -- устанавливает значение по ключу. func (sf *kStoreKv) Set(rec store_ent.IRecKvFix) *d0.Result[d0.Bool] { sf.Lock() defer sf.Unlock() key := rec.Key() sf.log.Debug("Set(): key='%v'", key.Get()) fnSet := func(txn *badger.Txn) error { err := txn.Set(key.Byte(), rec.Val().Byte()) return err } err := sf.db.Update(fnSet) if err != nil { err := d0.LetErr("Set(): key=%v, err=\n\t%w", key, err) sf.log.Err(err) return d0.ResErr[d0.Bool](err) } return d0.ResOk(d0.Bool(true)) } // Get -- возвращает значение по ключу. func (sf *kStoreKv) Get(key store_ent.IStoreKeyFix) *d0.Result[[]byte] { sf.RLock() defer sf.RUnlock() sf.log.Debug("Get(): key='%v'", key) var binVal []byte fnGet := func(txn *badger.Txn) error { item, err := txn.Get(key.Byte()) if err != nil { return err } binVal, err = item.ValueCopy(binVal) return err } err := sf.db.View(fnGet) if err != nil { err := d0.LetErr("Get(): key=%v, err=\n\t%v", key, err) sf.log.Err(err) return d0.ResErr[[]byte](err) } return d0.ResOk(binVal) } // ByPrefix -- фильтрует ключи по префиксу. func (sf *kStoreKv) ByPrefix(prefix store_ent.IStoreKeyFix) *d0.Result[[]store_ent.AStoreKey] { var ( binKey []byte lstKey = []store_ent.AStoreKey{} ) // fnValue := func(v []byte) error { // fmt.Printf("key=%s, value=%s\n", key, v) // return nil // } fnPrefix := func(txn *badger.Txn) error { it := txn.NewIterator(badger.DefaultIteratorOptions) defer it.Close() binPref := prefix.Byte() for it.Seek(binPref); it.ValidForPrefix(binPref); it.Next() { item := it.Item() binKey = item.Key() // err := item.Value(fnValue) // if err != nil { // return err // } lstKey = append(lstKey, rec_key.ARecKey(binKey)) } return nil } err := sf.db.View(fnPrefix) if err != nil { err := d0.LetErr("ByPrefix(): in find, err=\n\t%w", err) return d0.ResErr[[]store_ent.AStoreKey](err) } return d0.ResOk(lstKey) } // Delete -- удалить ключ из хранилища. func (sf *kStoreKv) Delete(key store_ent.IStoreKeyFix) *d0.Result[d0.Bool] { sf.Lock() defer sf.Unlock() sf.log.Debug("Delete(): key='%v'", key) fnDelete := func(txn *badger.Txn) error { err := txn.Delete(key.Byte()) return err } err := sf.db.Update(fnDelete) if err != nil { err := d0.LetErr("Delete(): key=%v, err=\n\t%w", key, err) sf.log.Err(err) return d0.ResErr[d0.Bool](err) } return d0.ResOk(d0.Bool(true)) } // Открывает базу при создании. func (sf *kStoreKv) open() { sf.Lock() defer sf.Unlock() sf.log.Debug("open()") strPath := os.Getenv("LOCAL_STORE_PATH") d0.If(strPath == "").Hassert("open(): env LOCAL_STORE_PATH not set") pwd, err := os.Getwd() d0.If(err != nil).Hassert("open(): in get PWD, err=\n\t%v", err) sf.storePath = d0.LetTxt(pwd + strPath + "/db_local") err = os.MkdirAll(sf.storePath.String(), 0750) d0.If(err != nil).Hassert("open(): in make dir %v, err=\n\t%v", sf.storePath, err) sf.db, err = badger.Open(badger.DefaultOptions(sf.storePath.String())) d0.If(err != nil).Hassert("open(): in open DB %v, err=\n\t%v", sf.storePath, err) sf.wg.Add(storeStreamName) sf.isWork.Set() go sf.close() go sf.clean() } // Выполняет периодическую сборку мусора в файле. func (sf *kStoreKv) clean() { chRun := make(chan d0.Num, 2) defer close(chRun) fnClean := func() { sf.Lock() defer sf.Unlock() _ = sf.db.RunValueLogGC(0.7) } chRun <- 1 for { select { case <-sf.kCtx.SelfCtx().Done(): // надо прекратить работу return case <-chRun: // Пора поработать fnClean() } time.Sleep(time.Second * 1) } } // Ожидает последнего потока под отдельной блокировкой. func (sf *kStoreKv) wait(chWait chan d0.Num) { for { time.Sleep(time.Millisecond * 5) if sf.wg.Len() <= 1 { break } } close(chWait) } // Ожидает закрытия контекста ядра, закрывает хранилище. func (sf *kStoreKv) close() { sf.kCtx.Wait() sf.Lock() defer sf.Unlock() if sf.isWork.IsNot() { return } chWait := make(chan d0.Num, 2) go sf.wait(chWait) <-chWait sf.isWork.Reset() err := sf.db.Close() d0.If(err != nil).Assert("kStoreKv.close(): in close DB, err=\n\t%v", err) sf.wg.Done(storeStreamName) sf.log.Debug("close(): done") }