mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
Block store (#578)
Contains the implementation of the block store In this pr is a simple way to send and receive data through a database I have implemented the base functionality as well as quite a few tests to make sure that everything works There are a few methods that have yet to be implemented, but theoretically they should be implemented as calls to the other functions, ie append should just be a call to WriteAt This doesn't affect anything yet so it can safely be merged whenever. I don't want this pr to stagnate like file view, so I'm happy to write multiple prs for this
This commit is contained in:
@@ -5,6 +5,7 @@ go 1.22
|
||||
toolchain go1.22.0
|
||||
|
||||
require (
|
||||
github.com/alecthomas/units v0.0.0-20231202071711-9a357b53e9c9
|
||||
github.com/alessio/shellescape v1.4.1
|
||||
github.com/armon/circbuf v0.0.0-20190214190532-5111143e8da2
|
||||
github.com/creack/pty v1.1.18
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
github.com/alecthomas/units v0.0.0-20231202071711-9a357b53e9c9 h1:ez/4by2iGztzR4L0zgAOR8lTQK9VlyBVVd7G4omaOQs=
|
||||
github.com/alecthomas/units v0.0.0-20231202071711-9a357b53e9c9/go.mod h1:OMCwj8VM1Kc9e19TLln2VL61YJF0x1XFtfdL4JdbSyE=
|
||||
github.com/alessio/shellescape v1.4.1 h1:V7yhSDDn8LP4lc4jS8pFkt0zCnzVJlG5JXy9BVKJUX0=
|
||||
github.com/alessio/shellescape v1.4.1/go.mod h1:PZAiSCk0LJaZkiCSkPv8qIobYglO3FPpyFjDCtHLS30=
|
||||
github.com/armon/circbuf v0.0.0-20190214190532-5111143e8da2 h1:7Ip0wMmLHLRJdrloDxZfhMm0xrLXZS8+COSu2bXmEQs=
|
||||
@@ -55,6 +57,7 @@ github.com/sawka/txwrap v0.1.2 h1:v8xS0Z1LE7/6vMZA81PYihI+0TSR6Zm1MalzzBIuXKc=
|
||||
github.com/sawka/txwrap v0.1.2/go.mod h1:T3nlw2gVpuolo6/XEetvBbk1oMXnY978YmBFy1UyHvw=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
|
||||
github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
|
||||
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
|
||||
github.com/wavetermdev/ssh_config v0.0.0-20240306041034-17e2087ebde2 h1:onqZrJVap1sm15AiIGTfWzdr6cEF0KdtddeuuOVhzyY=
|
||||
@@ -71,6 +74,8 @@ golang.org/x/sys v0.15.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/term v0.15.0 h1:y/Oo/a/q3IXu26lQgl04j/gjuBDOBlx7X6Om1j2CPW4=
|
||||
golang.org/x/term v0.15.0/go.mod h1:BDl952bC7+uMoWR75FIrCDx79TPU9oHkTZ9yRbYOrX0=
|
||||
golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
mvdan.cc/sh/v3 v3.7.0 h1:lSTjdP/1xsddtaKfGg7Myu7DnlHItd3/M2tomOcNNBg=
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,242 @@
|
||||
package blockstore
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"path"
|
||||
"sync"
|
||||
|
||||
"github.com/jmoiron/sqlx"
|
||||
_ "github.com/mattn/go-sqlite3"
|
||||
"github.com/sawka/txwrap"
|
||||
"github.com/wavetermdev/waveterm/wavesrv/pkg/dbutil"
|
||||
"github.com/wavetermdev/waveterm/wavesrv/pkg/scbase"
|
||||
)
|
||||
|
||||
const DBFileName = "blockstore.db"
|
||||
|
||||
type SingleConnDBGetter struct {
|
||||
SingleConnLock *sync.Mutex
|
||||
}
|
||||
|
||||
var dbWrap *SingleConnDBGetter
|
||||
|
||||
type TxWrap = txwrap.TxWrap
|
||||
|
||||
func InitDBState() {
|
||||
dbWrap = &SingleConnDBGetter{SingleConnLock: &sync.Mutex{}}
|
||||
}
|
||||
|
||||
func (dbg *SingleConnDBGetter) GetDB(ctx context.Context) (*sqlx.DB, error) {
|
||||
db, err := GetDB(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
dbg.SingleConnLock.Lock()
|
||||
return db, nil
|
||||
}
|
||||
|
||||
func (dbg *SingleConnDBGetter) ReleaseDB(db *sqlx.DB) {
|
||||
dbg.SingleConnLock.Unlock()
|
||||
}
|
||||
|
||||
func WithTx(ctx context.Context, fn func(tx *TxWrap) error) error {
|
||||
return txwrap.DBGWithTx(ctx, dbWrap, fn)
|
||||
}
|
||||
|
||||
func WithTxRtn[RT any](ctx context.Context, fn func(tx *TxWrap) (RT, error)) (RT, error) {
|
||||
var rtn RT
|
||||
txErr := WithTx(ctx, func(tx *TxWrap) error {
|
||||
temp, err := fn(tx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rtn = temp
|
||||
return nil
|
||||
})
|
||||
return rtn, txErr
|
||||
}
|
||||
|
||||
var globalDBLock = &sync.Mutex{}
|
||||
var globalDB *sqlx.DB
|
||||
var globalDBErr error
|
||||
|
||||
func GetDBName() string {
|
||||
scHome := scbase.GetWaveHomeDir()
|
||||
return path.Join(scHome, DBFileName)
|
||||
}
|
||||
|
||||
func GetDB(ctx context.Context) (*sqlx.DB, error) {
|
||||
if txwrap.IsTxWrapContext(ctx) {
|
||||
return nil, fmt.Errorf("cannot call GetDB from within a running transaction")
|
||||
}
|
||||
globalDBLock.Lock()
|
||||
defer globalDBLock.Unlock()
|
||||
if globalDB == nil && globalDBErr == nil {
|
||||
dbName := GetDBName()
|
||||
globalDB, globalDBErr = sqlx.Open("sqlite3", fmt.Sprintf("file:%s?cache=shared&mode=rwc&_journal_mode=WAL&_busy_timeout=5000", dbName))
|
||||
if globalDBErr != nil {
|
||||
globalDBErr = fmt.Errorf("opening db[%s]: %w", dbName, globalDBErr)
|
||||
log.Printf("[db] error: %v\n", globalDBErr)
|
||||
} else {
|
||||
log.Printf("[db] successfully opened db %s\n", dbName)
|
||||
}
|
||||
}
|
||||
return globalDB, globalDBErr
|
||||
}
|
||||
|
||||
func CloseDB() {
|
||||
globalDBLock.Lock()
|
||||
defer globalDBLock.Unlock()
|
||||
if globalDB == nil {
|
||||
return
|
||||
}
|
||||
err := globalDB.Close()
|
||||
if err != nil {
|
||||
log.Printf("[db] error closing database: %v\n", err)
|
||||
}
|
||||
globalDB = nil
|
||||
}
|
||||
|
||||
func (f *FileInfo) ToMap() map[string]interface{} {
|
||||
rtn := make(map[string]interface{})
|
||||
log.Printf("fileInfo ToMap is unimplemented!")
|
||||
return rtn
|
||||
}
|
||||
|
||||
func (fInfo *FileInfo) FromMap(m map[string]interface{}) bool {
|
||||
fileOpts := FileOptsType{}
|
||||
dbutil.QuickSetBool(&fileOpts.Circular, m, "circular")
|
||||
dbutil.QuickSetInt64(&fileOpts.MaxSize, m, "maxsize")
|
||||
|
||||
var metaJson []byte
|
||||
dbutil.QuickSetBytes(&metaJson, m, "meta")
|
||||
var fileMeta FileMeta
|
||||
err := json.Unmarshal(metaJson, &fileMeta)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
dbutil.QuickSetStr(&fInfo.BlockId, m, "blockid")
|
||||
dbutil.QuickSetStr(&fInfo.Name, m, "name")
|
||||
dbutil.QuickSetInt64(&fInfo.Size, m, "size")
|
||||
dbutil.QuickSetInt64(&fInfo.CreatedTs, m, "createdts")
|
||||
dbutil.QuickSetInt64(&fInfo.ModTs, m, "modts")
|
||||
fInfo.Opts = fileOpts
|
||||
fInfo.Meta = fileMeta
|
||||
return true
|
||||
}
|
||||
|
||||
func GetFileInfo(ctx context.Context, blockId string, name string) (*FileInfo, error) {
|
||||
fInfoArr, txErr := WithTxRtn(ctx, func(tx *TxWrap) ([]*FileInfo, error) {
|
||||
var rtn []*FileInfo
|
||||
query := `SELECT * FROM block_file WHERE name = 'file-1'`
|
||||
marr := tx.SelectMaps(query)
|
||||
for _, m := range marr {
|
||||
rtn = append(rtn, dbutil.FromMap[*FileInfo](m))
|
||||
}
|
||||
return rtn, nil
|
||||
})
|
||||
if txErr != nil {
|
||||
return nil, fmt.Errorf("GetFileInfo database error: %v", txErr)
|
||||
}
|
||||
if len(fInfoArr) > 1 {
|
||||
return nil, fmt.Errorf("GetFileInfo duplicate files in database")
|
||||
}
|
||||
if len(fInfoArr) == 0 {
|
||||
return nil, fmt.Errorf("GetFileInfo: File not found")
|
||||
}
|
||||
fInfo := fInfoArr[0]
|
||||
return fInfo, nil
|
||||
}
|
||||
|
||||
func GetCacheFromDB(ctx context.Context, blockId string, name string, off int64, length int64, cacheNum int64) (*[]byte, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (*[]byte, error) {
|
||||
var cacheData *[]byte
|
||||
query := `SELECT substr(data,?,?) FROM block_data WHERE blockid = ? AND name = ? and partidx = ?`
|
||||
tx.Get(&cacheData, query, off, length+1, blockId, name, cacheNum)
|
||||
if cacheData == nil {
|
||||
cacheData = &[]byte{}
|
||||
}
|
||||
return cacheData, nil
|
||||
})
|
||||
}
|
||||
|
||||
func DeleteFileFromDB(ctx context.Context, blockId string, name string) error {
|
||||
txErr := WithTx(ctx, func(tx *TxWrap) error {
|
||||
query := `DELETE from block_file where blockid = ? AND name = ?`
|
||||
tx.Exec(query, blockId, name)
|
||||
return nil
|
||||
})
|
||||
if txErr != nil {
|
||||
return txErr
|
||||
}
|
||||
txErr = WithTx(ctx, func(tx *TxWrap) error {
|
||||
query := `DELETE from block_data where blockid = ? AND name = ?`
|
||||
tx.Exec(query, blockId, name)
|
||||
return nil
|
||||
})
|
||||
if txErr != nil {
|
||||
return txErr
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func DeleteBlockFromDB(ctx context.Context, blockId string) error {
|
||||
txErr := WithTx(ctx, func(tx *TxWrap) error {
|
||||
query := `DELETE from block_file where blockid = ?`
|
||||
tx.Exec(query, blockId)
|
||||
return nil
|
||||
})
|
||||
if txErr != nil {
|
||||
return txErr
|
||||
}
|
||||
txErr = WithTx(ctx, func(tx *TxWrap) error {
|
||||
query := `DELETE from block_data where blockid = ?`
|
||||
tx.Exec(query, blockId)
|
||||
return nil
|
||||
})
|
||||
if txErr != nil {
|
||||
return txErr
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func GetAllFilesInDBForBlockId(ctx context.Context, blockId string) ([]*FileInfo, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) ([]*FileInfo, error) {
|
||||
var rtn []*FileInfo
|
||||
query := `SELECT * FROM block_file where blockid = ?`
|
||||
marr := tx.SelectMaps(query, blockId)
|
||||
for _, m := range marr {
|
||||
rtn = append(rtn, dbutil.FromMap[*FileInfo](m))
|
||||
}
|
||||
return rtn, nil
|
||||
})
|
||||
}
|
||||
|
||||
func GetAllFilesInDB(ctx context.Context) ([]*FileInfo, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) ([]*FileInfo, error) {
|
||||
var rtn []*FileInfo
|
||||
query := `SELECT * FROM block_file`
|
||||
marr := tx.SelectMaps(query)
|
||||
for _, m := range marr {
|
||||
rtn = append(rtn, dbutil.FromMap[*FileInfo](m))
|
||||
}
|
||||
return rtn, nil
|
||||
})
|
||||
}
|
||||
|
||||
func GetAllBlockIdsInDB(ctx context.Context) ([]string, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) ([]string, error) {
|
||||
var rtn []string
|
||||
query := `SELECT DISTINCT blockid FROM block_file`
|
||||
marr := tx.SelectMaps(query)
|
||||
for _, m := range marr {
|
||||
var blockId string
|
||||
dbutil.QuickSetStr(&blockId, m, "blockid")
|
||||
rtn = append(rtn, blockId)
|
||||
}
|
||||
return rtn, nil
|
||||
})
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,21 @@
|
||||
CREATE TABLE schema_migrations (version uint64,dirty bool);
|
||||
CREATE UNIQUE INDEX version_unique ON schema_migrations (version);
|
||||
CREATE TABLE block_file (
|
||||
blockid varchar(36) NOT NULL,
|
||||
name varchar(200) NOT NULL,
|
||||
maxsize bigint NOT NULL,
|
||||
circular boolean NOT NULL,
|
||||
size bigint NOT NULL,
|
||||
createdts bigint NOT NULL,
|
||||
modts bigint NOT NULL,
|
||||
meta json NOT NULL,
|
||||
PRIMARY KEY (blockid, name)
|
||||
);
|
||||
|
||||
CREATE TABLE block_data (
|
||||
blockid varchar(36) NOT NULL,
|
||||
name varchar(200) NOT NULL,
|
||||
partidx int NOT NULL,
|
||||
data blob NOT NULL,
|
||||
PRIMARY KEY(blockid, name, partidx)
|
||||
);
|
||||
Reference in New Issue
Block a user