initial db refactoring

This commit is contained in:
ston1th 2019-09-09 08:27:51 +02:00
commit e47a874711
5 changed files with 192 additions and 57 deletions

View file

@ -84,6 +84,19 @@ func (bs *BoltStore) Restore(r io.Reader) (err error) {
})
}
func (bs *BoltStore) Tx(writeable bool, f func(tx KVStore) error) error {
tx, err := db.db.Begin(writeable)
if err != nil {
return err
}
err = f(&BoltStoreTx{tx, bs})
if err != nil {
tx.Rollback()
return err
}
return tx.Commit()
}
func (bs *BoltStore) Get(key string, v interface{}) (err error) {
err = bs.db.View(func(tx *bolt.Tx) error {
b := tx.Bucket([]byte(defaultBoltBucket)).Get([]byte(key))
@ -120,12 +133,23 @@ func (bs *BoltStore) ForEach(f func(string, []byte) error) error {
}
func (bs *BoltStore) ForEachPrefix(prefix string, f func(string, []byte) error) error {
return bs.ForEach(func(k string, v []byte) error {
if trim, ok := hasTrimPrefix(k, prefix); ok {
return f(trim, v)
return bs.db.View(func(tx *bolt.Tx) error {
c := tx.Bucket([]byte(defaultBoltBucket)).Cursor()
p := []byte(prefix)
for k, v := c.Seek(p); k != nil && bytes.HasPrefix(k, p); k, v = c.Next() {
err := f(trimPrefix(k, p), v)
if err != nil {
return err
}
}
return nil
})
//return bs.ForEach(func(k string, v []byte) error {
// if trim, ok := hasTrimPrefix(k, prefix); ok {
// return f(trim, v)
// }
// return nil
//})
}
func (bs *BoltStore) Delete(key string) error {
@ -137,3 +161,52 @@ func (bs *BoltStore) Delete(key string) error {
func (bs *BoltStore) Close() error {
return bs.db.Close()
}
type BoltStoreTx struct {
tx *bolt.Tx
bs *BoltStore
}
func (tx *BoltStoreTx) Get(key string, v interface{}) (err error) {
b := tx.tx.Bucket([]byte(defaultBoltBucket)).Get([]byte(key))
if b == nil {
return ErrKeyNotFound
}
if v != nil {
return tx.bs.Unmarshal(b, v)
}
return nil
}
func (tx *BoltStoreTx) Set(key string, v interface{}) error {
var b []byte
if v != nil {
b, err = tx.bs.Marshal(v)
if err != nil {
return
}
}
return tx.tx.Bucket([]byte(defaultBoltBucket)).Put([]byte(key), b)
}
func (tx *BoltStoreTx) ForEach(f func(string, []byte) error) error {
return tx.tx.Bucket([]byte(defaultBoltBucket)).ForEach(func(k, v []byte) error {
return f(string(k), v)
})
}
func (tx *BoltStoreTx) ForEachPrefix(prefix string, f func(string, []byte) error) error {
c := tx.tx.Bucket([]byte(defaultBoltBucket)).Cursor()
p := []byte(prefix)
for k, v := c.Seek(p); k != nil && bytes.HasPrefix(k, p); k, v = c.Next() {
err := f(trimPrefix(k, p), v)
if err != nil {
return err
}
}
return nil
}
func (tx *BoltStoreTx) Delete(key string) error {
return tx.tx.Bucket([]byte(defaultBoltBucket)).Delete([]byte(key))
}