mirror of
https://github.com/aptly-dev/aptly.git
synced 2026-07-26 13:47:40 +00:00
etcd: implement transactions
- use temporary db for lookups in transactions - use batch implementation to commit transaction
This commit is contained in:
@@ -132,7 +132,8 @@ func (s *EtcDDBSuite) TestTransactionCommit(c *C) {
|
|||||||
transaction.Put(key2, value2)
|
transaction.Put(key2, value2)
|
||||||
v, err := s.db.Get(key)
|
v, err := s.db.Get(key)
|
||||||
c.Check(v, DeepEquals, value)
|
c.Check(v, DeepEquals, value)
|
||||||
transaction.Delete(key)
|
err = transaction.Delete(key)
|
||||||
|
c.Assert(err, IsNil)
|
||||||
|
|
||||||
_, err = transaction.Get(key2)
|
_, err = transaction.Get(key2)
|
||||||
c.Assert(err, IsNil)
|
c.Assert(err, IsNil)
|
||||||
@@ -152,5 +153,5 @@ func (s *EtcDDBSuite) TestTransactionCommit(c *C) {
|
|||||||
c.Check(v2, DeepEquals, value2)
|
c.Check(v2, DeepEquals, value2)
|
||||||
|
|
||||||
_, err = transaction.Get(key)
|
_, err = transaction.Get(key)
|
||||||
c.Assert(err, IsNil)
|
c.Assert(err, NotNil)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -40,6 +40,7 @@ func (s *EtcDStorage) Get(key []byte) (value []byte, err error) {
|
|||||||
}
|
}
|
||||||
for _, kv := range getResp.Kvs {
|
for _, kv := range getResp.Kvs {
|
||||||
value = kv.Value
|
value = kv.Value
|
||||||
|
break
|
||||||
}
|
}
|
||||||
if len(value) == 0 {
|
if len(value) == 0 {
|
||||||
err = database.ErrNotFound
|
err = database.ErrNotFound
|
||||||
@@ -169,12 +170,11 @@ func (s *EtcDStorage) CreateBatch() database.Batch {
|
|||||||
|
|
||||||
// OpenTransaction creates new transaction.
|
// OpenTransaction creates new transaction.
|
||||||
func (s *EtcDStorage) OpenTransaction() (database.Transaction, error) {
|
func (s *EtcDStorage) OpenTransaction() (database.Transaction, error) {
|
||||||
cli, err := internalOpen(s.url)
|
tmpdb, err := s.CreateTemporary()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
kvc := clientv3.NewKV(cli)
|
return &transaction{s: s, tmpdb: tmpdb}, nil
|
||||||
return &transaction{t: kvc}, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// CompactDB does nothing for etcd
|
// CompactDB does nothing for etcd
|
||||||
|
|||||||
@@ -3,51 +3,71 @@ package etcddb
|
|||||||
import (
|
import (
|
||||||
"github.com/aptly-dev/aptly/database"
|
"github.com/aptly-dev/aptly/database"
|
||||||
clientv3 "go.etcd.io/etcd/client/v3"
|
clientv3 "go.etcd.io/etcd/client/v3"
|
||||||
"go.etcd.io/etcd/client/v3/clientv3util"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type transaction struct {
|
type transaction struct {
|
||||||
t clientv3.KV
|
s *EtcDStorage
|
||||||
|
tmpdb database.Storage
|
||||||
|
ops []clientv3.Op
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get implements database.Reader interface.
|
// Get implements database.Reader interface.
|
||||||
func (t *transaction) Get(key []byte) ([]byte, error) {
|
func (t *transaction) Get(key []byte) (value []byte, err error) {
|
||||||
getResp, err := t.t.Get(Ctx, string(key))
|
value, err = t.tmpdb.Get(key)
|
||||||
|
// if not found, search main db
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
value, err = t.s.Get(key)
|
||||||
}
|
}
|
||||||
|
return
|
||||||
var value []byte
|
|
||||||
for _, kv := range getResp.Kvs {
|
|
||||||
valc := make([]byte, len(kv.Value))
|
|
||||||
copy(valc, kv.Value)
|
|
||||||
value = valc
|
|
||||||
}
|
|
||||||
|
|
||||||
return value, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Put implements database.Writer interface.
|
// Put implements database.Writer interface.
|
||||||
func (t *transaction) Put(key, value []byte) (err error) {
|
func (t *transaction) Put(key, value []byte) (err error) {
|
||||||
_, err = t.t.Txn(Ctx).
|
err = t.tmpdb.Put(key, value)
|
||||||
If().Then(clientv3.OpPut(string(key), string(value))).Commit()
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
t.ops = append(t.ops, clientv3.OpPut(string(key), string(value)))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// Delete implements database.Writer interface.
|
// Delete implements database.Writer interface.
|
||||||
func (t *transaction) Delete(key []byte) (err error) {
|
func (t *transaction) Delete(key []byte) (err error) {
|
||||||
_, err = t.t.Txn(Ctx).
|
err = t.tmpdb.Delete(key)
|
||||||
If(clientv3util.KeyExists(string(key))).
|
if err != nil {
|
||||||
Then(clientv3.OpDelete(string(key))).Commit()
|
return
|
||||||
|
}
|
||||||
|
t.ops = append(t.ops, clientv3.OpDelete(string(key)))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *transaction) Commit() (err error) {
|
func (t *transaction) Commit() (err error) {
|
||||||
|
kv := clientv3.NewKV(t.s.db)
|
||||||
|
|
||||||
|
batchSize := 128
|
||||||
|
for i := 0; i < len(t.ops); i += batchSize {
|
||||||
|
txn := kv.Txn(Ctx)
|
||||||
|
end := i + batchSize
|
||||||
|
if end > len(t.ops) {
|
||||||
|
end = len(t.ops)
|
||||||
|
}
|
||||||
|
|
||||||
|
batch := t.ops[i:end]
|
||||||
|
txn.Then(batch...)
|
||||||
|
_, err = txn.Commit()
|
||||||
|
if err != nil {
|
||||||
|
panic(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
t.ops = []clientv3.Op{}
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// Discard is safe to call after Commit(), it would be no-op
|
// Discard is safe to call after Commit(), it would be no-op
|
||||||
func (t *transaction) Discard() {
|
func (t *transaction) Discard() {
|
||||||
|
t.ops = []clientv3.Op{}
|
||||||
|
t.tmpdb.Drop()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user