fix: apply oplog migrations correctly using new storage interface

This commit is contained in:
garethgeorge
2024-09-09 00:22:33 -07:00
parent 426af294cd
commit 491a6a6725
6 changed files with 35 additions and 23 deletions
+4 -1
View File
@@ -79,7 +79,10 @@ func main() {
}
defer opstore.Close()
oplog := oplog.NewOpLog(opstore)
oplog, err := oplog.NewOpLog(opstore)
if err != nil {
zap.S().Fatalf("error creating oplog: %v", err)
}
// Create rotating log storage
logStore, err := logwriter.NewLogManager(path.Join(env.DataDir(), "rotatinglogs"), 14) // 14 days of logs
+2 -7
View File
@@ -83,7 +83,7 @@ func (o *BboltStore) Close() error {
func (o *BboltStore) Version() (int64, error) {
var version int64
err := o.db.View(func(tx *bolt.Tx) error {
o.db.View(func(tx *bolt.Tx) error {
b := tx.Bucket(SystemBucket)
if b == nil {
return nil
@@ -92,7 +92,7 @@ func (o *BboltStore) Version() (int64, error) {
version, err = serializationutil.Btoi(b.Get([]byte("version")))
return err
})
return version, err
return version, nil
}
func (o *BboltStore) SetVersion(version int64) error {
@@ -107,11 +107,6 @@ func (o *BboltStore) SetVersion(version int64) error {
// Add adds a generic operation to the operation log.
func (o *BboltStore) Add(ops ...*v1.Operation) error {
for _, op := range ops {
if op.Id != 0 {
return errors.New("operation already has an ID, OpLog.Add is expected to set the ID")
}
}
return o.db.Update(func(tx *bolt.Tx) error {
for _, op := range ops {
@@ -21,7 +21,7 @@ func Btoi(b []byte) (int64, error) {
}
func Stob(v string) []byte {
var b []byte
b := make([]byte, 0, len(v)+8)
b = append(b, Itob(int64(len(v)))...)
b = append(b, []byte(v)...)
return b
@@ -39,7 +39,7 @@ func Btos(b []byte) (string, int64, error) {
}
func BytesToKey(b []byte) []byte {
var key []byte
key := make([]byte, 0, 8+len(b))
key = append(key, Itob(int64(len(b)))...)
key = append(key, b...)
return key
-4
View File
@@ -1,7 +1,6 @@
package memstore
import (
"errors"
"slices"
"sync"
@@ -113,9 +112,6 @@ func (m *MemStore) Add(op ...*v1.Operation) error {
defer m.mu.Unlock()
for _, o := range op {
if o.Id != 0 {
return errors.New("operation already has an ID, OpLog.Add is expected to set the ID")
}
m.nextID++
o.Id = m.nextID
if o.FlowId == 0 {
+1 -1
View File
@@ -11,7 +11,7 @@ import (
var migrations = []func(*OpLog) error{
migration001FlowID,
migration002InstanceID,
migrationNoop, // migration003Reset Validated,
migrationNoop,
migration002InstanceID, // re-run migration002InstanceID to fix improperly set instance IDs
}
+26 -8
View File
@@ -30,10 +30,16 @@ type OpLog struct {
subscribers []*Subscription
}
func NewOpLog(store OpStore) *OpLog {
return &OpLog{
func NewOpLog(store OpStore) (*OpLog, error) {
o := &OpLog{
store: store,
}
if err := ApplyMigrations(o); err != nil {
return nil, err
}
return o, nil
}
func (o *OpLog) curSubscribers() []*Subscription {
@@ -64,24 +70,36 @@ func (o *OpLog) Get(opID int64) (*v1.Operation, error) {
return o.store.Get(opID)
}
func (o *OpLog) Add(op ...*v1.Operation) error {
if err := o.store.Add(op...); err != nil {
func (o *OpLog) Add(ops ...*v1.Operation) error {
for _, o := range ops {
if o.Id != 0 {
return errors.New("operation already has an ID, OpLog.Add is expected to set the ID")
}
}
if err := o.store.Add(ops...); err != nil {
return err
}
for _, sub := range o.curSubscribers() {
(*sub)(op, OPERATION_ADDED)
(*sub)(ops, OPERATION_ADDED)
}
return nil
}
func (o *OpLog) Update(op ...*v1.Operation) error {
if err := o.store.Update(op...); err != nil {
func (o *OpLog) Update(ops ...*v1.Operation) error {
for _, o := range ops {
if o.Id == 0 {
return errors.New("operation does not have an ID, OpLog.Update is expected to have an ID")
}
}
if err := o.store.Update(ops...); err != nil {
return err
}
for _, sub := range o.curSubscribers() {
(*sub)(op, OPERATION_UPDATED)
(*sub)(ops, OPERATION_UPDATED)
}
return nil
}