chore: drop dependency on bbolt

This commit is contained in:
Gareth George
2025-05-05 22:45:07 -07:00
parent 48636558ee
commit 20dba34ce7
9 changed files with 10 additions and 913 deletions
-48
View File
@@ -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 {
+1 -3
View File
@@ -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
+9 -57
View File
@@ -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=
-458
View File
@@ -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)
}
}
@@ -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
}
}
@@ -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)
}
}
@@ -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
}
@@ -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)
}
}
}
@@ -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)