diff --git a/cmd/backrest/backrest.go b/cmd/backrest/backrest.go index 8fc4f448..3a5d7b6a 100644 --- a/cmd/backrest/backrest.go +++ b/cmd/backrest/backrest.go @@ -25,7 +25,6 @@ import ( "github.com/garethgeorge/backrest/internal/logstore" "github.com/garethgeorge/backrest/internal/metric" "github.com/garethgeorge/backrest/internal/oplog" - "github.com/garethgeorge/backrest/internal/oplog/bboltstore" "github.com/garethgeorge/backrest/internal/oplog/sqlitestore" "github.com/garethgeorge/backrest/internal/orchestrator" "github.com/garethgeorge/backrest/internal/resticinstaller" @@ -87,7 +86,6 @@ func main() { if err != nil { zap.S().Fatalf("error creating oplog: %v", err) } - migrateBboltOplog(opstore) migratePopulateGuids(opstore, cfg) // Create rotating log storage @@ -271,52 +269,6 @@ func installLoggers() { zap.S().Infof("backrest version %v@%v, using log directory: %v", version, commit, logsDir) } -// migrateBboltOplog migrates the old bbolt oplog to the new sqlite oplog. -// It is careful to ensure that all migrations are applied before copying -// operations directly to the sqlite logstore. -func migrateBboltOplog(logstore oplog.OpStore) { - oldBboltOplogFile := path.Join(env.DataDir(), "oplog.boltdb") - if _, err := os.Stat(oldBboltOplogFile); err != nil { - return - } - - zap.S().Warnf("found old bbolt oplog file %q, migrating to sqlite", oldBboltOplogFile) - oldOpstore, err := bboltstore.NewBboltStore(oldBboltOplogFile) - if err != nil { - zap.S().Fatalf("error opening old bolt opstore: %v", oldBboltOplogFile, err) - } - oldOplog, err := oplog.NewOpLog(oldOpstore) - if err != nil { - zap.S().Fatalf("error opening old bolt oplog: %v", oldBboltOplogFile, err) - } - - var errs []error - var count int - if err := oldOplog.Query(oplog.Query{}, func(op *v1.Operation) error { - if err := logstore.Add(op); err != nil { - errs = append(errs, err) - zap.L().Warn("failed to migrate operation", zap.Error(err), zap.Any("operation", op)) - } else { - count++ - } - return nil - }); err != nil { - zap.S().Warnf("couldn't migrate all operations from the old bbolt oplog, if this recurs delete the file %q and restart", oldBboltOplogFile) - zap.S().Fatalf("error migrating old bbolt oplog: %v", err) - } - - if len(errs) > 0 { - zap.S().Errorf("encountered %d errors migrating old bbolt oplog, see logs for details.", len(errs), oldBboltOplogFile) - } - if err := oldOpstore.Close(); err != nil { - zap.S().Warnf("error closing old bbolt oplog: %v", err) - } - if err := os.Rename(oldBboltOplogFile, oldBboltOplogFile+".deprecated"); err != nil { - zap.S().Warnf("error removing old bbolt oplog: %v", err) - } - zap.S().Infof("migrated %d operations from old bbolt oplog to sqlite", count) -} - func migratePopulateGuids(logstore oplog.OpStore, cfg *v1.Config) { repoToGUID := make(map[string]string) for _, repo := range cfg.Repos { diff --git a/go.mod b/go.mod index 1ca4556d..3ffeb7bf 100644 --- a/go.mod +++ b/go.mod @@ -30,7 +30,7 @@ require ( github.com/natefinch/atomic v1.0.1 github.com/ncruces/zenity v0.10.14 github.com/prometheus/client_golang v1.22.0 - go.etcd.io/bbolt v1.4.0 + github.com/vearutop/statigz v1.5.0 go.uber.org/multierr v1.11.0 go.uber.org/zap v1.27.0 golang.org/x/crypto v0.38.0 @@ -61,7 +61,6 @@ require ( github.com/go-stack/stack v1.8.1 // indirect github.com/hashicorp/errwrap v1.1.0 // indirect github.com/josephspurrier/goversioninfo v1.5.0 // indirect - github.com/klauspost/compress v1.18.0 // indirect github.com/mattn/go-isatty v0.0.20 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/ncruces/go-strftime v0.1.9 // indirect @@ -71,7 +70,6 @@ require ( github.com/prometheus/procfs v0.16.1 // indirect github.com/randall77/makefat v0.0.0-20210315173500-7ddd0e42c844 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect - github.com/vearutop/statigz v1.5.0 // indirect go.opentelemetry.io/auto/sdk v1.1.0 // indirect go.opentelemetry.io/otel v1.35.0 // indirect go.opentelemetry.io/otel/metric v1.35.0 // indirect diff --git a/go.sum b/go.sum index f2b1d4e7..e8fb01f8 100644 --- a/go.sum +++ b/go.sum @@ -1,16 +1,16 @@ -al.essio.dev/pkg/shellescape v1.5.1 h1:86HrALUujYS/h+GtqoB26SBEdkWfmMI6FubjXlsXyho= -al.essio.dev/pkg/shellescape v1.5.1/go.mod h1:6sIqp7X2P6mThCQ7twERpZTuigpr6KbZWtls1U8I890= al.essio.dev/pkg/shellescape v1.6.0 h1:NxFcEqzFSEVCGN2yq7Huv/9hyCEGVa/TncnOOBBeXHA= al.essio.dev/pkg/shellescape v1.6.0/go.mod h1:6sIqp7X2P6mThCQ7twERpZTuigpr6KbZWtls1U8I890= -connectrpc.com/connect v1.17.0 h1:W0ZqMhtVzn9Zhn2yATuUokDLO5N+gIuBWMOnsQrfmZk= -connectrpc.com/connect v1.17.0/go.mod h1:0292hj1rnx8oFrStN7cB4jjVBeqs+Yx5yDIC2prWDO8= connectrpc.com/connect v1.18.1 h1:PAg7CjSAGvscaf6YZKUefjoih5Z/qYkyaTrBW8xvYPw= connectrpc.com/connect v1.18.1/go.mod h1:0292hj1rnx8oFrStN7cB4jjVBeqs+Yx5yDIC2prWDO8= github.com/akavel/rsrc v0.10.2 h1:Zxm8V5eI1hW4gGaYsJQUhxpjkENuG91ki8B4zCrvEsw= github.com/akavel/rsrc v0.10.2/go.mod h1:uLoCtb9J+EyAqh+26kdrTgmzRBFPGOolLWKpdxkKq+c= +github.com/andybalholm/brotli v1.1.1 h1:PR2pgnyFznKEugtsUo0xLdDop5SKXd5Qf5ysW+7XdTA= +github.com/andybalholm/brotli v1.1.1/go.mod h1:05ib4cKhjx3OQYUY22hTVd34Bc8upXjOLL2rKwwZBoA= github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/bool64/dev v0.2.39 h1:kP8DnMGlWXhGYJEZE/J0l/gVBdbuhoPGL+MJG4QbofE= +github.com/bool64/dev v0.2.39/go.mod h1:iJbh1y/HkunEPhgebWRNcs8wfGq7sjvJ6W5iabL8ACg= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/containrrr/shoutrrr v0.8.0 h1:mfG2ATzIS7NR2Ec6XL+xyoHzN97H8WPjir8aYzJUSec= @@ -66,15 +66,11 @@ github.com/go-task/slim-sprig v0.0.0-20230315185526-52ccab3ef572 h1:tfuBGBXKqDEe github.com/go-task/slim-sprig v0.0.0-20230315185526-52ccab3ef572/go.mod h1:9Pwr4B2jHnOSGXyyzV8ROjYa2ojvAY6HCGYYfMoC3Ls= github.com/gofrs/flock v0.12.1 h1:MTLVXXHf8ekldpJk3AKicLij9MdwOWkZ+a/jHHZby9E= github.com/gofrs/flock v0.12.1/go.mod h1:9zxTsyu5xtJ9DK+1tFZyibEV7y3uwDxPPfbxeeHCoD0= -github.com/golang-jwt/jwt/v5 v5.2.1 h1:OuVbFODueb089Lh128TAcimifWaLhJwVflnrgM17wHk= -github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/google/go-cmp v0.5.8/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= -github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= -github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/pprof v0.0.0-20240409012703-83162a5b38cd h1:gbpYu9NMq8jhDVbvlGkMFWCjLFlqqEZjEmObmhUy6Vo= @@ -94,12 +90,8 @@ github.com/hectane/go-acl v0.0.0-20230122075934-ca0b05cb1adb h1:PGufWXXDq9yaev6x github.com/hectane/go-acl v0.0.0-20230122075934-ca0b05cb1adb/go.mod h1:QiyDdbZLaJ/mZP4Zwc9g2QsfaEA4o7XvvgZegSci5/E= github.com/jarcoal/httpmock v1.3.0 h1:2RJ8GP0IIaWwcC9Fp2BmVi8Kog3v2Hn7VXM3fTd+nuc= github.com/jarcoal/httpmock v1.3.0/go.mod h1:3yb8rc4BI7TCBhFY8ng0gjuLKJNquuDNiPaZjnENuYg= -github.com/josephspurrier/goversioninfo v1.4.1 h1:5LvrkP+n0tg91J9yTkoVnt/QgNnrI1t4uSsWjIonrqY= -github.com/josephspurrier/goversioninfo v1.4.1/go.mod h1:JWzv5rKQr+MmW+LvM412ToT/IkYDZjaclF2pKDss8IY= github.com/josephspurrier/goversioninfo v1.5.0 h1:9TJtORoyf4YMoWSOo/cXFN9A/lB3PniJ91OxIH6e7Zg= github.com/josephspurrier/goversioninfo v1.5.0/go.mod h1:6MoTvFZ6GKJkzcdLnU5T/RGYUbHQbKpYeNP0AgQLd2o= -github.com/klauspost/compress v1.17.11 h1:In6xLpyWOi1+C7tXUUWv2ot1QvBjxevKAaI6IXrJmUc= -github.com/klauspost/compress v1.17.11/go.mod h1:pMDklpSncoRMuLFrf1W9Ss9KT+0rH90U12bZKk7uwG0= github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= @@ -109,11 +101,8 @@ github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0 github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/lxn/walk v0.0.0-20210112085537-c389da54e794/go.mod h1:E23UucZGqpuUANJooIbHWCufXvOcT6E7Stq81gU+CSQ= github.com/lxn/win v0.0.0-20210218163916-a377121e959e/go.mod h1:KxxjdtRkfNoYDCUP5ryK7XJJNTnpC8atvtmTheChOtk= -github.com/mattn/go-colorable v0.1.13 h1:fFA4WZxdEF4tXPZVKMLwD8oUnCTTo08duU7wxecdEvA= -github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg= github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= -github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= @@ -133,20 +122,12 @@ github.com/oxtoacart/bpool v0.0.0-20190530202638-03653db5a59c/go.mod h1:X07ZCGwU github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/prometheus/client_golang v1.20.5 h1:cxppBPuYhUnsO6yo/aoRol4L7q7UFfdm+bR9r+8l63Y= -github.com/prometheus/client_golang v1.20.5/go.mod h1:PIEt8X02hGcP8JWbeHyeZ53Y/jReSnHgO035n//V5WE= github.com/prometheus/client_golang v1.22.0 h1:rb93p9lokFEsctTys46VnV1kLCDpVZ0a/Y92Vm0Zc6Q= github.com/prometheus/client_golang v1.22.0/go.mod h1:R7ljNsLXhuQXYZYtw6GAE9AZg8Y7vEW5scdCXrWRXC0= -github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= -github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= -github.com/prometheus/common v0.60.1 h1:FUas6GcOw66yB/73KC+BOZoFJmbo/1pojoILArPAaSc= -github.com/prometheus/common v0.60.1/go.mod h1:h0LYf1R1deLSKtD4Vdg8gy4RuOvENW2J/h19V5NADQw= github.com/prometheus/common v0.63.0 h1:YR/EIY1o3mEFP/kZCD7iDMnLPlGyuU2Gb3HIcXnA98k= github.com/prometheus/common v0.63.0/go.mod h1:VVFF/fBIoToEnWRVkYoXEkq3R3paCoxG9PXP74SnV18= -github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= -github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg= github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is= github.com/randall77/makefat v0.0.0-20210315173500-7ddd0e42c844 h1:GranzK4hv1/pqTIhMTXt2X8MmMOuH3hMeUR0o9SP5yc= @@ -157,34 +138,26 @@ github.com/skratchdot/open-golang v0.0.0-20200116055534-eef842397966/go.mod h1:s github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= -github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= -github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg= -github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/vearutop/statigz v1.5.0 h1:FuWwZiT82yBw4xbWdWIawiP2XFTyEPhIo8upRxiKLqk= github.com/vearutop/statigz v1.5.0/go.mod h1:oHmjFf3izfCO804Di1ZjB666P3fAlVzJEx2k6jNt/Gk= github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= -go.etcd.io/bbolt v1.3.11 h1:yGEzV1wPz2yVCLsD8ZAiGHhHVlczyC9d1rP43/VCRJ0= -go.etcd.io/bbolt v1.3.11/go.mod h1:dksAq7YMXoljX0xu6VF5DMZGbhYYoLUalEiSySYAS4I= -go.etcd.io/bbolt v1.4.0 h1:TU77id3TnN/zKr7CO/uk+fBCwF2jGcMuw2B/FMAzYIk= -go.etcd.io/bbolt v1.4.0/go.mod h1:AsD+OCi/qPN1giOX1aiLAha3o1U8rAz65bvN4j0sRuk= go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= go.opentelemetry.io/otel v1.9.0/go.mod h1:np4EoPGzoPs3O67xUVNoPPcmSvsfOxNlNA4F4AC+0Eo= -go.opentelemetry.io/otel v1.32.0 h1:WnBN+Xjcteh0zdk01SVqV55d/m62NJLJdIyb4y/WO5U= -go.opentelemetry.io/otel v1.32.0/go.mod h1:00DCVSB0RQcnzlwyTfqtxSm+DRr9hpYrHjNGiBHVQIg= go.opentelemetry.io/otel v1.35.0 h1:xKWKPxrxB6OtMCbmMY021CqC45J+3Onta9MqjhnusiQ= go.opentelemetry.io/otel v1.35.0/go.mod h1:UEqy8Zp11hpkUrL73gSlELM0DupHoiq72dR+Zqel/+Y= -go.opentelemetry.io/otel/metric v1.32.0 h1:xV2umtmNcThh2/a/aCP+h64Xx5wsj8qqnkYZktzNa0M= -go.opentelemetry.io/otel/metric v1.32.0/go.mod h1:jH7CIbbK6SH2V2wE16W05BHCtIDzauciCRLoc/SyMv8= go.opentelemetry.io/otel/metric v1.35.0 h1:0znxYu2SNyuMSQT4Y9WDWej0VpcsxkuklLa4/siN90M= go.opentelemetry.io/otel/metric v1.35.0/go.mod h1:nKVFgxBZ2fReX6IlyW28MgZojkoAkJGaE8CpgeAU3oE= +go.opentelemetry.io/otel/sdk v1.34.0 h1:95zS4k/2GOy069d321O8jWgYsW3MzVV+KuSPKp7Wr1A= +go.opentelemetry.io/otel/sdk v1.34.0/go.mod h1:0e/pNiaMAqaykJGKbi+tSjWfNNHMTxoC9qANsCzbyxU= +go.opentelemetry.io/otel/sdk/metric v1.34.0 h1:5CeK9ujjbFVL5c1PhLuStg1wxA7vQv7ce1EK0Gyvahk= +go.opentelemetry.io/otel/sdk/metric v1.34.0/go.mod h1:jQ/r8Ze28zRKoNRdkjCZxfs6YvBTG1+YIqyFVFYec5w= go.opentelemetry.io/otel/trace v1.9.0/go.mod h1:2737Q0MuG8q1uILYm2YYVkAyLtOofiTNGg6VODnOiPo= -go.opentelemetry.io/otel/trace v1.32.0 h1:WIC9mYrXf8TmY/EXuULKc8hR17vE+Hjv2cssQDe03fM= -go.opentelemetry.io/otel/trace v1.32.0/go.mod h1:+i4rkvCraA+tG6AzwloGaCtkx53Fa+L+V8e9a7YvhT8= go.opentelemetry.io/otel/trace v1.35.0 h1:dPpEfJu1sDIqruz7BHFG3c7528f6ddfSWfFDVt/xgMs= go.opentelemetry.io/otel/trace v1.35.0/go.mod h1:WUk7DtFp1Aw2MkvqGdwiXYDZZNvA/1J8o6xRXLrIkyc= go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc= @@ -199,12 +172,8 @@ go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= -golang.org/x/crypto v0.29.0 h1:L5SG1JTTXupVV3n6sUqMTeWbjAyfPwoda2DLX8J8FrQ= -golang.org/x/crypto v0.29.0/go.mod h1:+F4F4N5hv6v38hfeYwTdx20oUvLLc+QfrE9Ax9HtgRg= golang.org/x/crypto v0.38.0 h1:jt+WWG8IZlBnVbomuhg2Mdq0+BBQaHbtqHEFEigjUV8= golang.org/x/crypto v0.38.0/go.mod h1:MvrbAqul58NNYPKnOra203SB9vpuZW0e+RRZV+Ggqjw= -golang.org/x/image v0.22.0 h1:UtK5yLUzilVrkjMAZAZ34DXGpASN8i8pj8g+O+yd10g= -golang.org/x/image v0.22.0/go.mod h1:9hPFhljd4zZ1GNSIZJ49sqbp45GKK9t6w+iXvGqZUz4= golang.org/x/image v0.27.0 h1:C8gA4oWU/tKkdCfYT6T2u4faJu3MeNS5O8UPWlPF61w= golang.org/x/image v0.27.0/go.mod h1:xbdrClrAUway1MUTEZDq9mz/UpRwYAkFFNUslZtcB+g= golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= @@ -215,14 +184,10 @@ golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4/go.mod h1:p54w0d4576C0XHj96bSt6lcn1PtDYWL6XObtHCRCNQM= -golang.org/x/net v0.31.0 h1:68CPQngjLL0r2AlUKiSxtQFKvzRVbnzLwMUn5SzcLHo= -golang.org/x/net v0.31.0/go.mod h1:P4fl1q7dY2hnZFxEk4pPSkDHF+QqjitcnDjUQyMM+pM= golang.org/x/net v0.40.0 h1:79Xs7wF06Gbdcg4kdCCIQArK11Z1hr5POQ6+fIYHNuY= golang.org/x/net v0.40.0/go.mod h1:y0hY0exeL2Pku80/zKK7tpntoX23cqL3Oa6njdgRtds= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.9.0 h1:fEo0HyrW1GIgZdpbhCRO0PkJajUS5H9IFUztCgEo2jQ= -golang.org/x/sync v0.9.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= golang.org/x/sync v0.14.0 h1:woo0S4Yywslg6hp4eUFjTVOyKt0RookbpAHG4c1HmhQ= golang.org/x/sync v0.14.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -232,18 +197,13 @@ golang.org/x/sys v0.0.0-20201018230417-eeed37f84f13/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.27.0 h1:wBqf8DvsY9Y/2P8gAfPDEYNuS30J4lPHJxXSb/nJZ+s= -golang.org/x/sys v0.27.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.33.0 h1:q3i8TbbEz+JRD9ywIRlyRAQbM0qF7hu24q3teo2hbuw= golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.20.0 h1:gK/Kv2otX8gz+wn7Rmb3vT96ZwuoxnQlY+HlJVj7Qug= -golang.org/x/text v0.20.0/go.mod h1:D4IsuqiFMhST5bX19pQ9ikHC2GsaKyk/oF+pn3ducp4= golang.org/x/text v0.25.0 h1:qVyWApTSYLk/drJRO5mDlNYskwQznZmkpV2c8q9zls4= golang.org/x/text v0.25.0/go.mod h1:WEdwpYrmk1qmdHvhkSTNPm3app7v4rsT8F2UD6+VHIA= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= @@ -255,20 +215,12 @@ golang.org/x/tools v0.26.0/go.mod h1:TPVVj70c7JJ3WCazhD8OdXcZg/og+b9+tH/KxylGwH0 golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -google.golang.org/genproto/googleapis/api v0.0.0-20241113202542-65e8d215514f h1:M65LEviCfuZTfrfzwwEoxVtgvfkFkBUbFnRbxCXuXhU= -google.golang.org/genproto/googleapis/api v0.0.0-20241113202542-65e8d215514f/go.mod h1:Yo94eF2nj7igQt+TiJ49KxjIH8ndLYPZMIRSiRcEbg0= google.golang.org/genproto/googleapis/api v0.0.0-20250505200425-f936aa4a68b2 h1:vPV0tzlsK6EzEDHNNH5sa7Hs9bd7iXR7B1tSiPepkV0= google.golang.org/genproto/googleapis/api v0.0.0-20250505200425-f936aa4a68b2/go.mod h1:pKLAc5OolXC3ViWGI62vvC0n10CpwAtRcTNCFwTKBEw= -google.golang.org/genproto/googleapis/rpc v0.0.0-20241113202542-65e8d215514f h1:C1QccEa9kUwvMgEUORqQD9S17QesQijxjZ84sO82mfo= -google.golang.org/genproto/googleapis/rpc v0.0.0-20241113202542-65e8d215514f/go.mod h1:GX3210XPVPUjJbTUbvwI8f2IpZDMZuPJWDzDuebbviI= google.golang.org/genproto/googleapis/rpc v0.0.0-20250505200425-f936aa4a68b2 h1:IqsN8hx+lWLqlN+Sc3DoMy/watjofWiU8sRFgQ8fhKM= google.golang.org/genproto/googleapis/rpc v0.0.0-20250505200425-f936aa4a68b2/go.mod h1:qQ0YXyHHx3XkvlzUtpXDkS29lDSafHMZBAZDc03LQ3A= -google.golang.org/grpc v1.68.0 h1:aHQeeJbo8zAkAa3pRzrVjZlbz6uSfeOXlJNQM0RAbz0= -google.golang.org/grpc v1.68.0/go.mod h1:fmSPC5AsjSBCK54MyHRx48kpOti1/jRfOlwEWywNjWA= google.golang.org/grpc v1.72.0 h1:S7UkcVa60b5AAQTaO6ZKamFp1zMZSU0fGDK2WZLbBnM= google.golang.org/grpc v1.72.0/go.mod h1:wH5Aktxcg25y1I3w7H69nHfXdOG3UiadoBtjh3izSDM= -google.golang.org/protobuf v1.35.2 h1:8Ar7bF+apOIoThw1EdZl0p1oWvMqTHmpA2fRTyZO8io= -google.golang.org/protobuf v1.35.2/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE= google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY= google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY= gopkg.in/Knetic/govaluate.v3 v3.0.0/go.mod h1:csKLBORsPbafmSCGTEh3U7Ozmsuq8ZSIlKk1bcqph0E= diff --git a/internal/oplog/bboltstore/bboltstore.go b/internal/oplog/bboltstore/bboltstore.go deleted file mode 100644 index 83524f80..00000000 --- a/internal/oplog/bboltstore/bboltstore.go +++ /dev/null @@ -1,458 +0,0 @@ -package bboltstore - -import ( - "errors" - "fmt" - "os" - "path" - "slices" - "testing" - "time" - - v1 "github.com/garethgeorge/backrest/gen/go/v1" - "github.com/garethgeorge/backrest/internal/oplog" - "github.com/garethgeorge/backrest/internal/oplog/bboltstore/indexutil" - "github.com/garethgeorge/backrest/internal/oplog/bboltstore/serializationutil" - "github.com/garethgeorge/backrest/internal/protoutil" - "go.etcd.io/bbolt" - bolt "go.etcd.io/bbolt" - "google.golang.org/protobuf/proto" -) - -type EventType int - -const ( - EventTypeUnknown = EventType(iota) - EventTypeOpCreated = EventType(iota) - EventTypeOpUpdated = EventType(iota) -) - -var ( - SystemBucket = []byte("oplog.system") // system stores metadata - OpLogBucket = []byte("oplog.log") // oplog stores existant operations. - RepoIndexBucket = []byte("oplog.repo_idx") // repo_index tracks IDs of operations affecting a given repo - PlanIndexBucket = []byte("oplog.plan_idx") // plan_index tracks IDs of operations affecting a given plan - FlowIdIndexBucket = []byte("oplog.flow_id_idx") // flow_id_index tracks IDs of operations affecting a given flow - InstanceIndexBucket = []byte("oplog.instance_idx") // instance_id_index tracks IDs of operations affecting a given instance - SnapshotIndexBucket = []byte("oplog.snapshot_idx") // snapshot_index tracks IDs of operations affecting a given snapshot -) - -// OpLog represents a log of operations performed. -// Operations are indexed by repo and plan. -type BboltStore struct { - db *bolt.DB -} - -var _ oplog.OpStore = &BboltStore{} - -func NewBboltStore(databasePath string) (*BboltStore, error) { - if err := os.MkdirAll(path.Dir(databasePath), 0700); err != nil { - return nil, fmt.Errorf("error creating database directory: %s", err) - } - - db, err := bolt.Open(databasePath, 0600, &bolt.Options{Timeout: 1 * time.Second}) - if err != nil { - return nil, fmt.Errorf("error opening database: %s", err) - } - - o := &BboltStore{ - db: db, - } - - if err := db.Update(func(tx *bolt.Tx) error { - // Create the buckets if they don't exist - for _, bucket := range [][]byte{ - SystemBucket, OpLogBucket, RepoIndexBucket, PlanIndexBucket, SnapshotIndexBucket, FlowIdIndexBucket, InstanceIndexBucket, - } { - if _, err := tx.CreateBucketIfNotExists(bucket); err != nil { - return fmt.Errorf("creating bucket %s: %s", string(bucket), err) - } - } - - return nil - }); err != nil { - return nil, err - } - - return o, nil -} - -func (o *BboltStore) Close() error { - return o.db.Close() -} - -func (o *BboltStore) Version() (int64, error) { - var version int64 - o.db.View(func(tx *bolt.Tx) error { - b := tx.Bucket(SystemBucket) - if b == nil { - return nil - } - var err error - version, err = serializationutil.Btoi(b.Get([]byte("version"))) - return err - }) - return version, nil -} - -func (o *BboltStore) SetVersion(version int64) error { - return o.db.Update(func(tx *bolt.Tx) error { - b, err := tx.CreateBucketIfNotExists(SystemBucket) - if err != nil { - return fmt.Errorf("creating system bucket: %w", err) - } - return b.Put([]byte("version"), serializationutil.Itob(version)) - }) -} - -// Add adds a generic operation to the operation log. -func (o *BboltStore) Add(ops ...*v1.Operation) error { - - return o.db.Update(func(tx *bolt.Tx) error { - for _, op := range ops { - err := o.addOperationHelper(tx, op) - if err != nil { - return err - } - } - return nil - }) -} - -func (o *BboltStore) Update(ops ...*v1.Operation) error { - for _, op := range ops { - if op.Id == 0 { - return errors.New("operation does not have an ID, OpLog.Update expects operation with an ID") - } - } - return o.db.Update(func(tx *bolt.Tx) error { - var err error - for _, op := range ops { - _, err = o.deleteOperationHelper(tx, op.Id) - if err != nil { - return fmt.Errorf("deleting existing value prior to update: %w", err) - } - if err := o.addOperationHelper(tx, op); err != nil { - return fmt.Errorf("adding updated value: %w", err) - } - } - return nil - }) -} - -func (o *BboltStore) Delete(ids ...int64) ([]*v1.Operation, error) { - removedOps := make([]*v1.Operation, 0, len(ids)) - err := o.db.Update(func(tx *bolt.Tx) error { - for _, id := range ids { - removed, err := o.deleteOperationHelper(tx, id) - if err != nil { - return fmt.Errorf("deleting operation %v: %w", id, err) - } - removedOps = append(removedOps, removed) - } - return nil - }) - return removedOps, err -} - -func (o *BboltStore) getOperationHelper(b *bolt.Bucket, id int64) (*v1.Operation, error) { - bytes := b.Get(serializationutil.Itob(id)) - if bytes == nil { - return nil, fmt.Errorf("opid %v: %w", id, oplog.ErrNotExist) - } - - var op v1.Operation - if err := proto.Unmarshal(bytes, &op); err != nil { - return nil, fmt.Errorf("error unmarshalling operation: %w", err) - } - - return &op, nil -} - -func (o *BboltStore) nextID(b *bolt.Bucket, unixTimeMs int64) (int64, error) { - seq, err := b.NextSequence() - if err != nil { - return 0, fmt.Errorf("next sequence: %w", err) - } - return int64(unixTimeMs<<20) | int64(seq&((1<<20)-1)), nil -} - -func (o *BboltStore) addOperationHelper(tx *bolt.Tx, op *v1.Operation) error { - b := tx.Bucket(OpLogBucket) - if op.Id == 0 { - var err error - op.Id, err = o.nextID(b, time.Now().UnixMilli()) - if err != nil { - return fmt.Errorf("create next operation ID: %w", err) - } - } - - if op.FlowId == 0 { - op.FlowId = op.Id - } - - if err := protoutil.ValidateOperation(op); err != nil { - return fmt.Errorf("validating operation: %w", err) - } - - bytes, err := proto.Marshal(op) - if err != nil { - return fmt.Errorf("error marshalling operation: %w", err) - } - - if err := b.Put(serializationutil.Itob(op.Id), bytes); err != nil { - return fmt.Errorf("error putting operation into bucket: %w", err) - } - - // Update always universal indices - if op.RepoId != "" { - if err := indexutil.IndexByteValue(tx.Bucket(RepoIndexBucket), []byte(op.RepoId), op.Id); err != nil { - return fmt.Errorf("error adding operation to repo index: %w", err) - } - } - - if op.PlanId != "" { - if err := indexutil.IndexByteValue(tx.Bucket(PlanIndexBucket), []byte(op.PlanId), op.Id); err != nil { - return fmt.Errorf("error adding operation to repo index: %w", err) - } - } - - if op.SnapshotId != "" { - if err := indexutil.IndexByteValue(tx.Bucket(SnapshotIndexBucket), []byte(op.SnapshotId), op.Id); err != nil { - return fmt.Errorf("error adding operation to snapshot index: %w", err) - } - } - - if op.FlowId != 0 { - if err := indexutil.IndexByteValue(tx.Bucket(FlowIdIndexBucket), serializationutil.Itob(op.FlowId), op.Id); err != nil { - return fmt.Errorf("error adding operation to flow index: %w", err) - } - } - - if op.InstanceId != "" { - if err := indexutil.IndexByteValue(tx.Bucket(InstanceIndexBucket), []byte(op.InstanceId), op.Id); err != nil { - return fmt.Errorf("error adding operation to instance index: %w", err) - } - } - - return nil -} - -func (o *BboltStore) deleteOperationHelper(tx *bolt.Tx, id int64) (*v1.Operation, error) { - b := tx.Bucket(OpLogBucket) - - prevValue, err := o.getOperationHelper(b, id) - if err != nil { - return nil, fmt.Errorf("getting operation %v: %w", id, err) - } - - if prevValue.PlanId != "" { - if err := indexutil.IndexRemoveByteValue(tx.Bucket(PlanIndexBucket), []byte(prevValue.PlanId), id); err != nil { - return nil, fmt.Errorf("removing operation %v from plan index: %w", id, err) - } - } - - if prevValue.RepoId != "" { - if err := indexutil.IndexRemoveByteValue(tx.Bucket(RepoIndexBucket), []byte(prevValue.RepoId), id); err != nil { - return nil, fmt.Errorf("removing operation %v from repo index: %w", id, err) - } - } - - if prevValue.SnapshotId != "" { - if err := indexutil.IndexRemoveByteValue(tx.Bucket(SnapshotIndexBucket), []byte(prevValue.SnapshotId), id); err != nil { - return nil, fmt.Errorf("removing operation %v from snapshot index: %w", id, err) - } - } - - if prevValue.FlowId != 0 { - if err := indexutil.IndexRemoveByteValue(tx.Bucket(FlowIdIndexBucket), serializationutil.Itob(prevValue.FlowId), id); err != nil { - return nil, fmt.Errorf("removing operation %v from flow index: %w", id, err) - } - } - - if prevValue.InstanceId != "" { - if err := indexutil.IndexRemoveByteValue(tx.Bucket(InstanceIndexBucket), []byte(prevValue.InstanceId), id); err != nil { - return nil, fmt.Errorf("removing operation %v from instance index: %w", id, err) - } - } - - if err := b.Delete(serializationutil.Itob(id)); err != nil { - return nil, fmt.Errorf("deleting operation %v from bucket: %w", id, err) - } - - return prevValue, nil -} - -func (o *BboltStore) Get(id int64) (*v1.Operation, error) { - var op *v1.Operation - if err := o.db.View(func(tx *bolt.Tx) error { - var err error - op, err = o.getOperationHelper(tx.Bucket(OpLogBucket), id) - return err - }); err != nil { - return nil, err - } - return op, nil -} - -// Query represents a query to the operation log. -type Query struct { - RepoId *string - PlanId *string - SnapshotId *string - FlowId *int64 - InstanceId *string - Ids []int64 -} - -func (o *BboltStore) Query(q oplog.Query, f func(*v1.Operation) error) error { - return o.queryHelper(q, func(tx *bbolt.Tx, op *v1.Operation) error { - return f(op) - }, true) -} - -func (o *BboltStore) QueryMetadata(q oplog.Query, f func(oplog.OpMetadata) error) error { - return errors.New("not implemented") -} - -func (o *BboltStore) Transform(q oplog.Query, f func(*v1.Operation) (*v1.Operation, error)) error { - return o.queryHelper(q, func(tx *bbolt.Tx, op *v1.Operation) error { - origId := op.Id - transformed, err := f(op) - if err != nil { - return err - } - if transformed == nil { - return nil - } - transformed.Modno = op.Modno + 1 - if _, err := o.deleteOperationHelper(tx, origId); err != nil { - return fmt.Errorf("deleting old operation: %w", err) - } - if err := o.addOperationHelper(tx, transformed); err != nil { - return fmt.Errorf("adding updated operation: %w", err) - } - return nil - }, false) -} - -func (o *BboltStore) queryHelper(query oplog.Query, do func(tx *bbolt.Tx, op *v1.Operation) error, isReadOnly bool) error { - helper := func(tx *bolt.Tx) error { - iterators := make([]indexutil.IndexIterator, 0, 5) - if query.PlanID != nil { - iterators = append(iterators, indexutil.IndexSearchByteValue(tx.Bucket(PlanIndexBucket), []byte(*query.PlanID))) - } - if query.SnapshotID != nil { - iterators = append(iterators, indexutil.IndexSearchByteValue(tx.Bucket(SnapshotIndexBucket), []byte(*query.SnapshotID))) - } - if query.FlowID != nil { - iterators = append(iterators, indexutil.IndexSearchByteValue(tx.Bucket(FlowIdIndexBucket), serializationutil.Itob(*query.FlowID))) - } - if query.InstanceID != nil { - iterators = append(iterators, indexutil.IndexSearchByteValue(tx.Bucket(InstanceIndexBucket), []byte(*query.InstanceID))) - } - - var ids []int64 - if len(iterators) == 0 && len(query.OpIDs) == 0 { - if query.Limit == 0 && query.Offset == 0 && !query.Reversed { - return o.forAll(tx, func(op *v1.Operation) error { - if query.Match(op) { - return do(tx, op) - } - return nil - }) - } else { - b := tx.Bucket(OpLogBucket) - c := b.Cursor() - for k, _ := c.First(); k != nil; k, _ = c.Next() { - if id, err := serializationutil.Btoi(k); err != nil { - continue // skip corrupt keys - } else { - ids = append(ids, id) - } - } - } - } else if len(iterators) > 0 { - ids = indexutil.CollectAll()(indexutil.NewJoinIterator(iterators...)) - } - ids = append(ids, query.OpIDs...) - - if query.Reversed { - slices.Reverse(ids) - } - if query.Offset > 0 { - if len(ids) <= query.Offset { - return nil - } - ids = ids[query.Offset:] - } - if query.Limit > 0 && len(ids) > query.Limit { - ids = ids[:query.Limit] - } - - return o.forOpsByIds(tx, ids, func(op *v1.Operation) error { - if query.Match(op) { - return do(tx, op) - } - return nil - }) - } - if isReadOnly { - return o.db.View(helper) - } else { - return o.db.Update(helper) - } -} - -func (o *BboltStore) forOpsByIds(tx *bolt.Tx, ids []int64, do func(*v1.Operation) error) error { - b := tx.Bucket(OpLogBucket) - for _, id := range ids { - op, err := o.getOperationHelper(b, id) - if err != nil { - return err - } - if err := do(op); err != nil { - if err == oplog.ErrStopIteration { - break - } - return err - } - } - return nil -} - -func (o *BboltStore) forAll(tx *bolt.Tx, do func(*v1.Operation) error) error { - b := tx.Bucket(OpLogBucket) - c := b.Cursor() - for k, v := c.First(); k != nil; k, v = c.Next() { - var op v1.Operation - if err := proto.Unmarshal(v, &op); err != nil { - return fmt.Errorf("error unmarshalling operation: %w", err) - } - if err := do(&op); err != nil { - if err == oplog.ErrStopIteration { - break - } - return err - } - } - return nil -} - -func (o *BboltStore) ResetForTest(t *testing.T) { - if err := o.db.Update(func(tx *bolt.Tx) error { - for _, bucket := range [][]byte{ - SystemBucket, OpLogBucket, RepoIndexBucket, PlanIndexBucket, SnapshotIndexBucket, FlowIdIndexBucket, InstanceIndexBucket, - } { - if err := tx.DeleteBucket(bucket); err != nil { - return fmt.Errorf("deleting bucket %s: %w", string(bucket), err) - } - if _, err := tx.CreateBucketIfNotExists(bucket); err != nil { - return fmt.Errorf("creating bucket %s: %w", string(bucket), err) - } - } - return nil - }); err != nil { - t.Fatalf("error resetting database: %s", err) - } -} diff --git a/internal/oplog/bboltstore/indexutil/indexutil.go b/internal/oplog/bboltstore/indexutil/indexutil.go deleted file mode 100644 index 66147f90..00000000 --- a/internal/oplog/bboltstore/indexutil/indexutil.go +++ /dev/null @@ -1,164 +0,0 @@ -package indexutil - -import ( - "bytes" - - "github.com/garethgeorge/backrest/internal/oplog/bboltstore/serializationutil" - bolt "go.etcd.io/bbolt" -) - -// IndexByteValue indexes a value and recordId tuple creating multimap from value to lists of associated recordIds. -func IndexByteValue(b *bolt.Bucket, value []byte, recordId int64) error { - key := serializationutil.BytesToKey(value) - key = append(key, serializationutil.Itob(recordId)...) - return b.Put(key, []byte{}) -} - -func IndexRemoveByteValue(b *bolt.Bucket, value []byte, recordId int64) error { - key := serializationutil.BytesToKey(value) - key = append(key, serializationutil.Itob(recordId)...) - return b.Delete(key) -} - -// IndexSearchByteValue searches the index given a value and returns an iterator over the associated recordIds. -func IndexSearchByteValue(b *bolt.Bucket, value []byte) IndexIterator { - return newSearchIterator(b, serializationutil.BytesToKey(value)) -} - -type IndexIterator interface { - Next() (int64, bool) -} - -type SeekableIndexIterator interface { - IndexIterator - Seek(int64) (int64, bool) // seek to the first recordId >= id and return it or return false. -} - -type IndexSearchIterator struct { - c *bolt.Cursor - k []byte - prefix []byte -} - -var _ SeekableIndexIterator = &IndexSearchIterator{} - -func newSearchIterator(b *bolt.Bucket, prefix []byte) IndexIterator { - c := b.Cursor() - k, _ := c.Seek(prefix) - return &IndexSearchIterator{ - c: c, - k: k, - prefix: prefix, - } -} - -func (i *IndexSearchIterator) Next() (int64, bool) { - if i.k == nil || !bytes.HasPrefix(i.k, i.prefix) { - return 0, false - } - id, err := serializationutil.Btoi(i.k[len(i.prefix):]) - if err != nil { - // this sholud never happen, if it does it indicates database corruption. - return 0, false - } - i.k, _ = i.c.Next() - return id, true -} - -func (i *IndexSearchIterator) Seek(id int64) (int64, bool) { - seekTo := []byte{} - seekTo = append(seekTo, i.prefix...) - seekTo = append(seekTo, serializationutil.Itob(id)...) - k, _ := i.c.Seek(seekTo) - if k == nil || !bytes.HasPrefix(k, i.prefix) { - return 0, false - } - id, err := serializationutil.Btoi(k[len(i.prefix):]) - if err != nil { - return 0, false - } - return id, true -} - -type JoinIterator struct { - iters []IndexIterator - seekables []SeekableIndexIterator -} - -func NewJoinIterator(iters ...IndexIterator) *JoinIterator { - seekables := make([]SeekableIndexIterator, 0, len(iters)) - for _, iter := range iters { - if seekable, ok := iter.(SeekableIndexIterator); ok { - seekables = append(seekables, seekable) - } else { - seekables = append(seekables, nil) - } - } - return &JoinIterator{ - iters: iters, - seekables: seekables, - } -} - -func (j *JoinIterator) Next() (int64, bool) { - if len(j.iters) == 0 { - return 0, false - } - - nexts := make([]int64, len(j.iters)) - for idx, iter := range j.iters { - id, ok := iter.Next() - if !ok { - return 0, false - } - nexts[idx] = id - } - - for { - var ok bool - maxIdx := 0 - allSame := true - for idx, id := range nexts { - if id > nexts[maxIdx] { - maxIdx = idx - } - if id != nexts[0] { - allSame = false - } - } - - if allSame { - return nexts[0], true - } - - for idx, id := range nexts { - if id == nexts[maxIdx] { - continue - } - - if j.seekables[idx] != nil { - nexts[idx], ok = j.seekables[idx].Seek(nexts[maxIdx]) - if !ok { - return 0, false - } - } else { - nexts[idx], ok = j.iters[idx].Next() - if !ok { - return 0, false - } - } - } - } -} - -type Collector func(IndexIterator) []int64 - -func CollectAll() Collector { - return func(iter IndexIterator) []int64 { - ids := make([]int64, 0, 100) - for id, ok := iter.Next(); ok; id, ok = iter.Next() { - ids = append(ids, id) - } - return ids - } -} diff --git a/internal/oplog/bboltstore/indexutil/indexutil_test.go b/internal/oplog/bboltstore/indexutil/indexutil_test.go deleted file mode 100644 index e90958d5..00000000 --- a/internal/oplog/bboltstore/indexutil/indexutil_test.go +++ /dev/null @@ -1,102 +0,0 @@ -package indexutil - -import ( - "fmt" - "reflect" - "testing" - - "go.etcd.io/bbolt" -) - -func TestIndexing(t *testing.T) { - db, err := bbolt.Open(t.TempDir()+"/test.boltdb", 0600, nil) - if err != nil { - t.Fatalf("error opening database: %s", err) - } - defer db.Close() - - if err := db.Update(func(tx *bbolt.Tx) error { - b, err := tx.CreateBucket([]byte("test")) - if err != nil { - return fmt.Errorf("error creating bucket: %s", err) - } - for id := 0; id < 100; id += 1 { - if err := IndexByteValue(b, []byte("document"), int64(id)); err != nil { - return err - } - } - return nil - }); err != nil { - t.Fatalf("db.Update error: %v", err) - } - - if err := db.View(func(tx *bbolt.Tx) error { - b := tx.Bucket([]byte("test")) - ids := CollectAll()(IndexSearchByteValue(b, []byte("document"))) - if len(ids) != 100 { - t.Errorf("want 100 ids, got %d", len(ids)) - } - ids = CollectAll()(IndexSearchByteValue(b, []byte("other"))) - if len(ids) != 0 { - t.Errorf("want 0 ids, got %d", len(ids)) - } - return nil - }); err != nil { - t.Fatalf("db.View error: %v", err) - } -} - -func TestIndexJoin(t *testing.T) { - // Arrange - db, err := bbolt.Open(t.TempDir()+"/test.boltdb", 0600, nil) - if err != nil { - t.Fatalf("error opening database: %s", err) - } - defer db.Close() - - if err := db.Update(func(tx *bbolt.Tx) error { - b, err := tx.CreateBucket([]byte("test")) - if err != nil { - return fmt.Errorf("error creating bucket: %s", err) - } - for id := 0; id < 150; id += 1 { - if err := IndexByteValue(b, []byte("document"), int64(id)); err != nil { - return err - } - } - - for id := 0; id < 100; id += 2 { - if err := IndexByteValue(b, []byte("other"), int64(id)); err != nil { - return err - } - } - - return nil - }); err != nil { - t.Fatalf("db.Update error: %v", err) - } - - if err := db.View(func(tx *bbolt.Tx) error { - // Act - b := tx.Bucket([]byte("test")) - ids := CollectAll()(NewJoinIterator(IndexSearchByteValue(b, []byte("document")), IndexSearchByteValue(b, []byte("other")))) - - // Assert - if len(ids) != 50 { - t.Errorf("want 50 ids, got %d", len(ids)) - } - - wantIds := []int64{} - for id := 0; id < 100; id += 2 { - wantIds = append(wantIds, int64(id)) - } - - if !reflect.DeepEqual(ids, wantIds) { - t.Errorf("want %v, got %v", wantIds, ids) - } - - return nil - }); err != nil { - t.Fatalf("db.View error: %v", err) - } -} diff --git a/internal/oplog/bboltstore/serializationutil/serializationutil.go b/internal/oplog/bboltstore/serializationutil/serializationutil.go deleted file mode 100644 index 397282bb..00000000 --- a/internal/oplog/bboltstore/serializationutil/serializationutil.go +++ /dev/null @@ -1,46 +0,0 @@ -package serializationutil - -import ( - "encoding/binary" - "errors" -) - -var ErrInvalidLength = errors.New("invalid length") - -func Itob(v int64) []byte { - b := make([]byte, 8) - binary.BigEndian.PutUint64(b, uint64(v)) - return b -} - -func Btoi(b []byte) (int64, error) { - if len(b) != 8 { - return 0, ErrInvalidLength - } - return int64(binary.BigEndian.Uint64(b)), nil -} - -func Stob(v string) []byte { - b := make([]byte, 0, len(v)+8) - b = append(b, Itob(int64(len(v)))...) - b = append(b, []byte(v)...) - return b -} - -func Btos(b []byte) (string, int64, error) { - if len(b) < 8 { - return "", 0, ErrInvalidLength - } - length, _ := Btoi(b[:8]) - if int64(len(b)) < 8+length { - return "", 0, ErrInvalidLength - } - return string(b[8 : 8+length]), 8 + length, nil -} - -func BytesToKey(b []byte) []byte { - key := make([]byte, 0, 8+len(b)) - key = append(key, Itob(int64(len(b)))...) - key = append(key, b...) - return key -} diff --git a/internal/oplog/bboltstore/serializationutil/serializationutil_test.go b/internal/oplog/bboltstore/serializationutil/serializationutil_test.go deleted file mode 100644 index 524a768e..00000000 --- a/internal/oplog/bboltstore/serializationutil/serializationutil_test.go +++ /dev/null @@ -1,23 +0,0 @@ -package serializationutil - -import "testing" - -func TestItoa(t *testing.T) { - nums := []int64{0, 1, 2, 3, 4, 1 << 32, int64(1) << 62} - for _, num := range nums { - b := Itob(num) - if v, _ := Btoi(b); v != num { - t.Errorf("itob/btoi failed for %d", num) - } - } -} - -func TestStob(t *testing.T) { - strs := []string{"", "a", "ab", "abc", "abcd", "abcde", "abcdef"} - for _, str := range strs { - b := Stob(str) - if val, _, _ := Btos(b); val != str { - t.Errorf("stob/btos failed for %s", str) - } - } -} diff --git a/internal/oplog/storetests/storecontract_test.go b/internal/oplog/storetests/storecontract_test.go index 10f12a8e..53232e17 100644 --- a/internal/oplog/storetests/storecontract_test.go +++ b/internal/oplog/storetests/storecontract_test.go @@ -7,7 +7,6 @@ import ( v1 "github.com/garethgeorge/backrest/gen/go/v1" "github.com/garethgeorge/backrest/internal/oplog" - "github.com/garethgeorge/backrest/internal/oplog/bboltstore" "github.com/garethgeorge/backrest/internal/oplog/memstore" "github.com/garethgeorge/backrest/internal/oplog/sqlitestore" "github.com/google/go-cmp/cmp" @@ -21,12 +20,6 @@ const ( ) func StoresForTest(t testing.TB) map[string]oplog.OpStore { - bboltstore, err := bboltstore.NewBboltStore(t.TempDir() + "/test.boltdb") - if err != nil { - t.Fatalf("error creating bbolt store: %s", err) - } - t.Cleanup(func() { bboltstore.Close() }) - sqlitestoreinst, err := sqlitestore.NewSqliteStore(t.TempDir() + "/test.sqlite") if err != nil { t.Fatalf("error creating sqlite store: %s", err) @@ -40,7 +33,6 @@ func StoresForTest(t testing.TB) map[string]oplog.OpStore { t.Cleanup(func() { sqlitememstore.Close() }) return map[string]oplog.OpStore{ - "bbolt": bboltstore, "memory": memstore.NewMemStore(), "sqlite": sqlitestoreinst, "sqlitemem": sqlitememstore, @@ -751,10 +743,6 @@ func TestQueryMetadata(t *testing.T) { t.Parallel() for name, store := range StoresForTest(t) { t.Run(name, func(t *testing.T) { - if name == "bbolt" { - t.Skip("bbolt does not support metadata") - } - log, err := oplog.NewOpLog(store) if err != nil { t.Fatalf("error creating oplog: %v", err)