mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
moving hard to OID model
This commit is contained in:
+46
-60
@@ -6,31 +6,41 @@ package wstore
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/wavetermdev/thenextwave/pkg/shellexec"
|
||||
"github.com/wavetermdev/thenextwave/pkg/util/ds"
|
||||
"github.com/wavetermdev/thenextwave/pkg/waveobj"
|
||||
)
|
||||
|
||||
var WorkspaceMap = ds.NewSyncMap[*Workspace]()
|
||||
var TabMap = ds.NewSyncMap[*Tab]()
|
||||
var BlockMap = ds.NewSyncMap[*Block]()
|
||||
|
||||
func init() {
|
||||
waveobj.RegisterType[*Client]()
|
||||
waveobj.RegisterType[*Window]()
|
||||
waveobj.RegisterType[*Workspace]()
|
||||
waveobj.RegisterType[*Tab]()
|
||||
waveobj.RegisterType[*Block]()
|
||||
}
|
||||
|
||||
type Client struct {
|
||||
ClientId string `json:"clientid"`
|
||||
OID string `json:"oid"`
|
||||
Version int `json:"version"`
|
||||
MainWindowId string `json:"mainwindowid"`
|
||||
}
|
||||
|
||||
func (c Client) GetId() string {
|
||||
return c.ClientId
|
||||
func (*Client) GetOType() string {
|
||||
return "client"
|
||||
}
|
||||
|
||||
// stores the ui-context of the window
|
||||
// workspaceid, active tab, active block within each tab, window size, etc.
|
||||
type Window struct {
|
||||
WindowId string `json:"windowid"`
|
||||
OID string `json:"oid"`
|
||||
Version int `json:"version"`
|
||||
WorkspaceId string `json:"workspaceid"`
|
||||
ActiveTabId string `json:"activetabid"`
|
||||
ActiveBlockMap map[string]string `json:"activeblockmap"` // map from tabid to blockid
|
||||
@@ -39,42 +49,30 @@ type Window struct {
|
||||
LastFocusTs int64 `json:"lastfocusts"`
|
||||
}
|
||||
|
||||
func (w Window) GetId() string {
|
||||
return w.WindowId
|
||||
func (*Window) GetOType() string {
|
||||
return "window"
|
||||
}
|
||||
|
||||
type Workspace struct {
|
||||
Lock *sync.Mutex `json:"-"`
|
||||
WorkspaceId string `json:"workspaceid"`
|
||||
Name string `json:"name"`
|
||||
TabIds []string `json:"tabids"`
|
||||
OID string `json:"oid"`
|
||||
Version int `json:"version"`
|
||||
Name string `json:"name"`
|
||||
TabIds []string `json:"tabids"`
|
||||
}
|
||||
|
||||
func (ws Workspace) GetId() string {
|
||||
return ws.WorkspaceId
|
||||
}
|
||||
|
||||
func (ws *Workspace) WithLock(f func()) {
|
||||
ws.Lock.Lock()
|
||||
defer ws.Lock.Unlock()
|
||||
f()
|
||||
func (*Workspace) GetOType() string {
|
||||
return "workspace"
|
||||
}
|
||||
|
||||
type Tab struct {
|
||||
Lock *sync.Mutex `json:"-"`
|
||||
TabId string `json:"tabid"`
|
||||
Name string `json:"name"`
|
||||
BlockIds []string `json:"blockids"`
|
||||
OID string `json:"oid"`
|
||||
Version int `json:"version"`
|
||||
Name string `json:"name"`
|
||||
BlockIds []string `json:"blockids"`
|
||||
}
|
||||
|
||||
func (tab Tab) GetId() string {
|
||||
return tab.TabId
|
||||
}
|
||||
|
||||
func (tab *Tab) WithLock(f func()) {
|
||||
tab.Lock.Lock()
|
||||
defer tab.Lock.Unlock()
|
||||
f()
|
||||
func (*Tab) GetOType() string {
|
||||
return "tab"
|
||||
}
|
||||
|
||||
type FileDef struct {
|
||||
@@ -108,7 +106,8 @@ type WinSize struct {
|
||||
}
|
||||
|
||||
type Block struct {
|
||||
BlockId string `json:"blockid"`
|
||||
OID string `json:"oid"`
|
||||
Version int `json:"version"`
|
||||
BlockDef *BlockDef `json:"blockdef"`
|
||||
Controller string `json:"controller"`
|
||||
View string `json:"view"`
|
||||
@@ -116,45 +115,32 @@ type Block struct {
|
||||
RuntimeOpts *RuntimeOpts `json:"runtimeopts,omitempty"`
|
||||
}
|
||||
|
||||
func (b *Block) GetOType() string {
|
||||
func (*Block) GetOType() string {
|
||||
return "block"
|
||||
}
|
||||
|
||||
func (b Block) GetId() string {
|
||||
return b.BlockId
|
||||
}
|
||||
|
||||
// TODO remove
|
||||
func (b *Block) WithLock(f func()) {
|
||||
f()
|
||||
}
|
||||
|
||||
func CreateTab(workspaceId string, name string) (*Tab, error) {
|
||||
tab := &Tab{
|
||||
Lock: &sync.Mutex{},
|
||||
TabId: uuid.New().String(),
|
||||
OID: uuid.New().String(),
|
||||
Name: name,
|
||||
BlockIds: []string{},
|
||||
}
|
||||
TabMap.Set(tab.TabId, tab)
|
||||
TabMap.Set(tab.OID, tab)
|
||||
ws := WorkspaceMap.Get(workspaceId)
|
||||
if ws == nil {
|
||||
return nil, fmt.Errorf("workspace not found: %q", workspaceId)
|
||||
}
|
||||
ws.WithLock(func() {
|
||||
ws.TabIds = append(ws.TabIds, tab.TabId)
|
||||
})
|
||||
ws.TabIds = append(ws.TabIds, tab.OID)
|
||||
return tab, nil
|
||||
}
|
||||
|
||||
func CreateWorkspace() (*Workspace, error) {
|
||||
ws := &Workspace{
|
||||
Lock: &sync.Mutex{},
|
||||
WorkspaceId: uuid.New().String(),
|
||||
TabIds: []string{},
|
||||
OID: uuid.New().String(),
|
||||
TabIds: []string{},
|
||||
}
|
||||
WorkspaceMap.Set(ws.WorkspaceId, ws)
|
||||
_, err := CreateTab(ws.WorkspaceId, "Tab 1")
|
||||
WorkspaceMap.Set(ws.OID, ws)
|
||||
_, err := CreateTab(ws.OID, "Tab 1")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -164,7 +150,7 @@ func CreateWorkspace() (*Workspace, error) {
|
||||
func EnsureInitialData() error {
|
||||
ctx, cancelFn := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancelFn()
|
||||
clientCount, err := DBGetCount[Client](ctx)
|
||||
clientCount, err := DBGetCount[*Client](ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error getting client count: %w", err)
|
||||
}
|
||||
@@ -175,7 +161,7 @@ func EnsureInitialData() error {
|
||||
workspaceId := uuid.New().String()
|
||||
tabId := uuid.New().String()
|
||||
client := &Client{
|
||||
ClientId: uuid.New().String(),
|
||||
OID: uuid.New().String(),
|
||||
MainWindowId: windowId,
|
||||
}
|
||||
err = DBInsert(ctx, client)
|
||||
@@ -183,7 +169,7 @@ func EnsureInitialData() error {
|
||||
return fmt.Errorf("error inserting client: %w", err)
|
||||
}
|
||||
window := &Window{
|
||||
WindowId: windowId,
|
||||
OID: windowId,
|
||||
WorkspaceId: workspaceId,
|
||||
ActiveTabId: tabId,
|
||||
ActiveBlockMap: make(map[string]string),
|
||||
@@ -201,16 +187,16 @@ func EnsureInitialData() error {
|
||||
return fmt.Errorf("error inserting window: %w", err)
|
||||
}
|
||||
ws := &Workspace{
|
||||
WorkspaceId: workspaceId,
|
||||
Name: "default",
|
||||
TabIds: []string{tabId},
|
||||
OID: workspaceId,
|
||||
Name: "default",
|
||||
TabIds: []string{tabId},
|
||||
}
|
||||
err = DBInsert(ctx, ws)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error inserting workspace: %w", err)
|
||||
}
|
||||
tab := &Tab{
|
||||
TabId: uuid.New().String(),
|
||||
OID: uuid.New().String(),
|
||||
Name: "Tab 1",
|
||||
BlockIds: []string{},
|
||||
}
|
||||
|
||||
+75
-106
@@ -6,155 +6,124 @@ package wstore
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"reflect"
|
||||
|
||||
"github.com/wavetermdev/thenextwave/pkg/waveobj"
|
||||
)
|
||||
|
||||
const Table_Client = "db_client"
|
||||
const Table_Workspace = "db_workspace"
|
||||
const Table_Tab = "db_tab"
|
||||
const Table_Block = "db_block"
|
||||
const Table_Window = "db_window"
|
||||
|
||||
// can replace with struct tags in the future
|
||||
type ObjectWithId interface {
|
||||
GetId() string
|
||||
func waveObjTableName(w waveobj.WaveObj) string {
|
||||
return "db_" + w.GetOType()
|
||||
}
|
||||
|
||||
// can replace these with struct tags in the future
|
||||
var idColumnName = map[string]string{
|
||||
Table_Client: "clientid",
|
||||
Table_Workspace: "workspaceid",
|
||||
Table_Tab: "tabid",
|
||||
Table_Block: "blockid",
|
||||
Table_Window: "windowid",
|
||||
func tableNameGen[T waveobj.WaveObj]() string {
|
||||
var zeroObj T
|
||||
return "db_" + zeroObj.GetOType()
|
||||
}
|
||||
|
||||
var tableToType = map[string]reflect.Type{
|
||||
Table_Client: reflect.TypeOf(Client{}),
|
||||
Table_Workspace: reflect.TypeOf(Workspace{}),
|
||||
Table_Tab: reflect.TypeOf(Tab{}),
|
||||
Table_Block: reflect.TypeOf(Block{}),
|
||||
Table_Window: reflect.TypeOf(Window{}),
|
||||
}
|
||||
|
||||
var typeToTable map[reflect.Type]string
|
||||
|
||||
func init() {
|
||||
typeToTable = make(map[reflect.Type]string)
|
||||
for k, v := range tableToType {
|
||||
typeToTable[v] = k
|
||||
}
|
||||
}
|
||||
|
||||
func DBGetCount[T ObjectWithId](ctx context.Context) (int, error) {
|
||||
func DBGetCount[T waveobj.WaveObj](ctx context.Context) (int, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (int, error) {
|
||||
var valInstance T
|
||||
table := typeToTable[reflect.TypeOf(valInstance)]
|
||||
if table == "" {
|
||||
return 0, fmt.Errorf("unknown table type: %T", valInstance)
|
||||
}
|
||||
table := tableNameGen[T]()
|
||||
query := fmt.Sprintf("SELECT count(*) FROM %s", table)
|
||||
return tx.GetInt(query), nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBGetSingleton[T ObjectWithId](ctx context.Context) (*T, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (*T, error) {
|
||||
var rtn T
|
||||
query := fmt.Sprintf("SELECT data FROM %s LIMIT 1", typeToTable[reflect.TypeOf(rtn)])
|
||||
jsonData := tx.GetString(query)
|
||||
return TxReadJson[T](tx, jsonData), nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBGet[T ObjectWithId](ctx context.Context, id string) (*T, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (*T, error) {
|
||||
var rtn T
|
||||
table := typeToTable[reflect.TypeOf(rtn)]
|
||||
if table == "" {
|
||||
return nil, fmt.Errorf("unknown table type: %T", rtn)
|
||||
}
|
||||
query := fmt.Sprintf("SELECT data FROM %s WHERE %s = ?", table, idColumnName[table])
|
||||
jsonData := tx.GetString(query, id)
|
||||
return TxReadJson[T](tx, jsonData), nil
|
||||
})
|
||||
}
|
||||
|
||||
type idDataType struct {
|
||||
Id string
|
||||
Data string
|
||||
OId string
|
||||
Version int
|
||||
Data []byte
|
||||
}
|
||||
|
||||
func DBSelectMap[T ObjectWithId](ctx context.Context, ids []string) (map[string]*T, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (map[string]*T, error) {
|
||||
var valInstance T
|
||||
table := typeToTable[reflect.TypeOf(valInstance)]
|
||||
if table == "" {
|
||||
return nil, fmt.Errorf("unknown table type: %T", &valInstance)
|
||||
func DBGetSingleton[T waveobj.WaveObj](ctx context.Context) (T, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (T, error) {
|
||||
table := tableNameGen[T]()
|
||||
query := fmt.Sprintf("SELECT oid, version, data FROM %s LIMIT 1", table)
|
||||
var row idDataType
|
||||
tx.Get(&row, query)
|
||||
rtn, err := waveobj.FromJsonGen[T](row.Data)
|
||||
if err != nil {
|
||||
return rtn, err
|
||||
}
|
||||
waveobj.SetVersion(rtn, row.Version)
|
||||
return rtn, nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBGet[T waveobj.WaveObj](ctx context.Context, id string) (T, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (T, error) {
|
||||
table := tableNameGen[T]()
|
||||
query := fmt.Sprintf("SELECT oid, version, data FROM %s WHERE oid = ?", table)
|
||||
var row idDataType
|
||||
tx.Get(&row, query, id)
|
||||
rtn, err := waveobj.FromJsonGen[T](row.Data)
|
||||
if err != nil {
|
||||
return rtn, err
|
||||
}
|
||||
waveobj.SetVersion(rtn, row.Version)
|
||||
return rtn, nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBSelectMap[T waveobj.WaveObj](ctx context.Context, ids []string) (map[string]T, error) {
|
||||
return WithTxRtn(ctx, func(tx *TxWrap) (map[string]T, error) {
|
||||
table := tableNameGen[T]()
|
||||
var rows []idDataType
|
||||
query := fmt.Sprintf("SELECT %s, data FROM %s WHERE %s IN (SELECT value FROM json_each(?))", idColumnName[table], table, idColumnName[table])
|
||||
query := fmt.Sprintf("SELECT oid, version, data FROM %s WHERE oid IN (SELECT value FROM json_each(?))", table)
|
||||
tx.Select(&rows, query, ids)
|
||||
rtnMap := make(map[string]*T)
|
||||
rtnMap := make(map[string]T)
|
||||
for _, row := range rows {
|
||||
if row.Id == "" || row.Data == "" {
|
||||
if row.OId == "" || len(row.Data) == 0 {
|
||||
continue
|
||||
}
|
||||
r := TxReadJson[T](tx, row.Data)
|
||||
if r == nil {
|
||||
continue
|
||||
waveObj, err := waveobj.FromJsonGen[T](row.Data)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
rtnMap[(*r).GetId()] = r
|
||||
waveobj.SetVersion(waveObj, row.Version)
|
||||
rtnMap[row.OId] = waveObj
|
||||
}
|
||||
return rtnMap, nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBDelete[T ObjectWithId](ctx context.Context, id string) error {
|
||||
func DBDelete[T waveobj.WaveObj](ctx context.Context, id string) error {
|
||||
return WithTx(ctx, func(tx *TxWrap) error {
|
||||
var rtn T
|
||||
table := typeToTable[reflect.TypeOf(rtn)]
|
||||
if table == "" {
|
||||
return fmt.Errorf("unknown table type: %T", rtn)
|
||||
}
|
||||
query := fmt.Sprintf("DELETE FROM %s WHERE %s = ?", table, idColumnName[table])
|
||||
table := tableNameGen[T]()
|
||||
query := fmt.Sprintf("DELETE FROM %s WHERE oid = ?", table)
|
||||
tx.Exec(query, id)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBUpdate[T ObjectWithId](ctx context.Context, val *T) error {
|
||||
if val == nil {
|
||||
return fmt.Errorf("cannot update nil value")
|
||||
}
|
||||
if (*val).GetId() == "" {
|
||||
func DBUpdate(ctx context.Context, val waveobj.WaveObj) error {
|
||||
oid := waveobj.GetOID(val)
|
||||
if oid == "" {
|
||||
return fmt.Errorf("cannot update %T value with empty id", val)
|
||||
}
|
||||
jsonData, err := waveobj.ToJson(val)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return WithTx(ctx, func(tx *TxWrap) error {
|
||||
table := typeToTable[reflect.TypeOf(*val)]
|
||||
if table == "" {
|
||||
return fmt.Errorf("unknown table type: %T", *val)
|
||||
}
|
||||
query := fmt.Sprintf("UPDATE %s SET data = ? WHERE %s = ?", table, idColumnName[table])
|
||||
tx.Exec(query, TxJson(tx, val), (*val).GetId())
|
||||
table := waveObjTableName(val)
|
||||
query := fmt.Sprintf("UPDATE %s SET data = ?, version = version+1 WHERE oid = ?", table)
|
||||
tx.Exec(query, jsonData, oid)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func DBInsert[T ObjectWithId](ctx context.Context, val *T) error {
|
||||
if val == nil {
|
||||
return fmt.Errorf("cannot insert nil value")
|
||||
}
|
||||
if (*val).GetId() == "" {
|
||||
func DBInsert[T waveobj.WaveObj](ctx context.Context, val T) error {
|
||||
oid := waveobj.GetOID(val)
|
||||
if oid == "" {
|
||||
return fmt.Errorf("cannot insert %T value with empty id", val)
|
||||
}
|
||||
jsonData, err := waveobj.ToJson(val)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return WithTx(ctx, func(tx *TxWrap) error {
|
||||
table := typeToTable[reflect.TypeOf(*val)]
|
||||
if table == "" {
|
||||
return fmt.Errorf("unknown table type: %T", *val)
|
||||
}
|
||||
query := fmt.Sprintf("INSERT INTO %s (%s, data) VALUES (?, ?)", table, idColumnName[table])
|
||||
tx.Exec(query, (*val).GetId(), TxJson(tx, val))
|
||||
table := waveObjTableName(val)
|
||||
query := fmt.Sprintf("INSERT INTO %s (oid, version, data) VALUES (?, ?, ?)", table)
|
||||
tx.Exec(query, oid, 1, jsonData)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
@@ -5,7 +5,6 @@ package wstore
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"path"
|
||||
@@ -63,24 +62,3 @@ func WithTx(ctx context.Context, fn func(tx *TxWrap) error) error {
|
||||
func WithTxRtn[RT any](ctx context.Context, fn func(tx *TxWrap) (RT, error)) (RT, error) {
|
||||
return txwrap.WithTxRtn(ctx, globalDB, fn)
|
||||
}
|
||||
|
||||
func TxJson(tx *TxWrap, v any) string {
|
||||
barr, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
tx.SetErr(fmt.Errorf("json marshal (%T): %w", v, err))
|
||||
return ""
|
||||
}
|
||||
return string(barr)
|
||||
}
|
||||
|
||||
func TxReadJson[T any](tx *TxWrap, jsonData string) *T {
|
||||
if jsonData == "" {
|
||||
return nil
|
||||
}
|
||||
var v T
|
||||
err := json.Unmarshal([]byte(jsonData), &v)
|
||||
if err != nil {
|
||||
tx.SetErr(fmt.Errorf("json unmarshal (%T): %w", v, err))
|
||||
}
|
||||
return &v
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user