diff --git a/api_directive.go b/api_directive.go index b4948e562..10928fff5 100644 --- a/api_directive.go +++ b/api_directive.go @@ -102,19 +102,19 @@ type directiveJobTableKeys struct { directiveJobType idx *Index tkey dax.TableKey - partition dax.VersionedPartition + partition dax.PartitionNum } type directiveJobFieldKeys struct { directiveJobType tkey dax.TableKey - field dax.VersionedField + field dax.FieldName } type directiveJobShards struct { directiveJobType tkey dax.TableKey - shard dax.VersionedShard + shard dax.ShardNum } // directiveWorker is a worker in a worker pool which handles portions of a @@ -342,7 +342,7 @@ func (api *API) pushJobsTableKeys(ctx context.Context, jobs chan<- directiveJobT jobs <- directiveJobTableKeys{ idx: idx, tkey: tkey, - partition: partition, + partition: partition.Num, } } } @@ -352,6 +352,8 @@ func (api *API) loadTableKeys(ctx context.Context, idx *Index, tkey dax.TableKey qtid := tkey.QualifiedTableID() mgr := api.serverlessStorage.GetTableKeyManager(qtid, partition) + + // load latest snapshot rc, err := mgr.LoadLatestSnapshot() if err != nil { return errors.Wrap(err, "loading table key snapshot") @@ -359,28 +361,18 @@ func (api *API) loadTableKeys(ctx context.Context, idx *Index, tkey dax.TableKey defer rc.Close() if err := api.TranslateIndexDB(ctx, string(tkey), int(partition), rc); err != nil { return errors.Wrap(err, "restoring table keys") - } - if err := func() error { - store := idx.TranslateStore(int(partition)) - reader, err := mgr.LoadWriteLog() + // define write log loading in a function since we have to do it + // before and after locking + loadWriteLog := func() error { + writelog, err := mgr.LoadWriteLog() if err != nil { - return errors.Wrap(err, "loading write log table keys") - } - reader := api.writeLogReader.TableKeyReader(ctx, qtid, partition.Num, partition.Version) - if err := reader.Open(); err != nil { - // TODO: this log can be confusing because on a create - // table, there is no log file yet, so an error is expected. - // Instead of swallowing this error, we need to check the - // error code and handle it differently. This means the - // writelogger will need to return an error indicating that - // the log file does not exist, but that that is expected. - // log.Printf("could not open log file for table: %s, partition: %d: version: %d, err: %s", table, partition.Num, partition.Version, err) - return nil + return errors.Wrap(err, "getting write log reader for table keys") } + reader := storage.NewTableKeyReader(qtid, partition, writelog) defer reader.Close() - + store := idx.TranslateStore(int(partition)) for msg, err := reader.Read(); err != io.EOF; msg, err = reader.Read() { if err != nil { return errors.Wrap(err, "reading from log reader") @@ -391,18 +383,21 @@ func (api *API) loadTableKeys(ctx context.Context, idx *Index, tkey dax.TableKey } } } - return nil - }(); err != nil { + } + // 1st write log load + if err := loadWriteLog(); err != nil { return err } - // Set the table/partition/version in the holder. - if err := api.holder.versionStore.AddPartitions(ctx, qtid, partition); err != nil { - return errors.Wrap(err, "adding partition to sharder") + // acquire lock on this partition's keys + if err := mgr.Lock(); err != nil { + return errors.Wrap(err, "locking table key partition") } - return nil + // reload writelog in case of changes between last load and + // lock. The manager object takes care of only loading new data. + return loadWriteLog() } func (api *API) pushJobsFieldKeys(ctx context.Context, jobs chan<- directiveJobType, fromD, toD *dax.Directive) { @@ -414,53 +409,44 @@ func (api *API) pushJobsFieldKeys(ctx context.Context, jobs chan<- directiveJobT for _, field := range fields { jobs <- directiveJobFieldKeys{ tkey: tkey, - field: field, + field: field.Name, } } } } -func (api *API) loadFieldKeys(ctx context.Context, tkey dax.TableKey, field dax.VersionedField) error { +func (api *API) loadFieldKeys(ctx context.Context, tkey dax.TableKey, field dax.FieldName) error { qtid := tkey.QualifiedTableID() - // Load the previous snapshot. Version 0 doesn't have a snapshot - // file; it only has log entries. - if field.Version > 0 { - // Load field snapshot: version - 1 - previousVersion := field.Version - 1 - rc, err := api.snapshotReadWriter.ReadFieldKeys(ctx, qtid, field.Name, previousVersion) - if err != nil { - return errors.Wrap(err, "reading field keys snapshot") - } - defer rc.Close() + mgr := api.serverlessStorage.GetFieldKeyManager(qtid, field) - if err := api.TranslateFieldDB(ctx, string(tkey), string(field.Name), rc); err != nil { - return errors.Wrap(err, "restoring field keys") - } + // load latest snapshot + rc, err := mgr.LoadLatestSnapshot() + if err != nil { + return errors.Wrap(err, "loading field key snapshot") + } + defer rc.Close() + if err := api.TranslateFieldDB(ctx, string(tkey), string(field), rc); err != nil { + return errors.Wrap(err, "restoring field keys") } - if err := func() error { + // define write log loading in a function since we have to do it + // before and after locking + loadWriteLog := func() error { + writelog, err := mgr.LoadWriteLog() + if err != nil { + return errors.Wrap(err, "getting write log reader for field keys") + } + reader := storage.NewFieldKeyReader(qtid, field, writelog) + defer reader.Close() // Get field in order to find the translate store. - fld := api.holder.Field(string(tkey), string(field.Name)) + fld := api.holder.Field(string(tkey), string(field)) if fld == nil { - log.Printf("field not found in holder: %s", field.Name) + log.Printf("field not found in holder: %s", field) return nil } store := fld.TranslateStore() - reader := api.writeLogReader.FieldKeyReader(ctx, qtid, field.Name, field.Version) - if err := reader.Open(); err != nil { - // TODO: this log can be confusing because on a create - // table, there is no log file yet, so an error is expected. - // Instead of swallowing this error, we need to check the - // error code and handle it differently. This means the - // writelogger will need to return an error indicating that - // the log file does not exist, but that that is expected. - // log.Printf("could not open log file for table: %s, field: %s: version: %d, err: %s", table, field.Name, field.Version, err) - return nil - } - defer reader.Close() - for msg, err := reader.Read(); err != io.EOF; msg, err = reader.Read() { if err != nil { return errors.Wrap(err, "reading from log reader") @@ -471,18 +457,21 @@ func (api *API) loadFieldKeys(ctx context.Context, tkey dax.TableKey, field dax. } } } - return nil - }(); err != nil { + } + // 1st write log load + if err := loadWriteLog(); err != nil { return err } - // Set the table/field/version in the holder. - if err := api.holder.versionStore.AddFields(ctx, qtid, field); err != nil { - return errors.Wrap(err, "adding field to sharder") + // acquire lock on this partition's keys + if err := mgr.Lock(); err != nil { + return errors.Wrap(err, "locking field key partition") } - return nil + // reload writelog in case of changes between last load and + // lock. The manager object takes care of only loading new data. + return loadWriteLog() } func (api *API) pushJobsShards(ctx context.Context, jobs chan<- directiveJobType, fromD, toD *dax.Directive) { @@ -497,7 +486,7 @@ func (api *API) pushJobsShards(ctx context.Context, jobs chan<- directiveJobType for _, shard := range shards { jobs <- directiveJobShards{ tkey: tkey, - shard: shard, + shard: shard.Num, } } } @@ -506,39 +495,25 @@ func (api *API) pushJobsShards(ctx context.Context, jobs chan<- directiveJobType func (api *API) loadShard(ctx context.Context, tkey dax.TableKey, shard dax.ShardNum) error { qtid := tkey.QualifiedTableID() - partition := disco.ShardToShardPartition(string(tkey), uint64(shard.Num), disco.DefaultPartitionN) - partitionNum := dax.PartitionNum(partition) + partition := dax.PartitionNum(disco.ShardToShardPartition(string(tkey), uint64(shard), disco.DefaultPartitionN)) - // Load the previous snapshot. Version 0 doesn't have a snapshot - // file; it only has log entries. - if shard.Version > 0 { - // Load shard snapshot: version - 1 - previousVersion := shard.Version - 1 - rc, err := api.snapshotReadWriter.ReadShardData(ctx, qtid, partitionNum, shard.Num, previousVersion) - if err != nil { - return errors.Wrap(err, "reading shard data snapshot") - } - - if err := api.RestoreShard(ctx, string(tkey), uint64(shard.Num), rc); err != nil { - return errors.Wrap(err, "restoring shard data") - } + mgr := api.serverlessStorage.GetShardManager(qtid, partition, shard) + rc, err := mgr.LoadLatestSnapshot() + if err != nil { + return errors.Wrap(err, "reading latest snapshot for shard") + } + if err := api.RestoreShard(ctx, string(tkey), uint64(shard), rc); err != nil { + return errors.Wrap(err, "restoring shard data") } - // WriteLog reader. - if err := func() error { - reader := api.writeLogReader.ShardReader(ctx, qtid, partitionNum, shard.Num, shard.Version) - if err := reader.Open(); err != nil { - // TODO: this log can be confusing because on a create - // table, there is no log file yet, so an error is expected. - // Instead of swallowing this error, we need to check the - // error code and handle it differently. This means the - // writelogger will need to return an error indicating that - // the log file does not exist, but that that is expected. - // log.Printf("could not open log file for table: %s, partition: %d: version: %d, shard: %d, err: %s", table, partition, shard.Version, shard.Num, err) - return nil + // define write log loading in a func because we do it twice. + loadWriteLog := func() error { + writelog, err := mgr.LoadWriteLog() + if err != nil { + return errors.Wrap(err, "") } + reader := storage.NewShardReader(qtid, partition, shard, writelog) defer reader.Close() - for logMsg, err := reader.Read(); err != io.EOF; logMsg, err = reader.Read() { if err != nil { return errors.Wrap(err, "reading from log reader") @@ -629,18 +604,21 @@ func (api *API) loadShard(ctx context.Context, tkey dax.TableKey, shard dax.Shar } } } - return nil - }(); err != nil { + } + // 1st write log load + if err := loadWriteLog(); err != nil { return err } - // Set the table/shard/version in the holder. - if err := api.holder.versionStore.AddShards(ctx, qtid, shard); err != nil { - return errors.Wrap(err, "adding shard to sharder") + // acquire lock on this partition's keys + if err := mgr.Lock(); err != nil { + return errors.Wrap(err, "locking field key partition") } - return nil + // reload writelog in case of changes between last load and + // lock. The manager object takes care of only loading new data. + return loadWriteLog() } ////////////////////////////////////////////////////////////// diff --git a/dax/computer/interfaces.go b/dax/computer/interfaces.go index b6a423026..7e00e59fc 100644 --- a/dax/computer/interfaces.go +++ b/dax/computer/interfaces.go @@ -43,6 +43,7 @@ type SnapshotService interface { // One must not call LoadWriteLog until after calling // LoadLatestSnapshot. One must not call Append, IncrementWLVersion, // or Snapshot until after successfully calling Lock. +// TODO(jaffee), this doesn't need to be an interface. Remove. type ServerlessStorage interface { // LoadLatestSnapshot loads the latest available snapshot in the snapshot store. LoadLatestSnapshot() (data io.ReadCloser, err error) diff --git a/dax/computer/service/computer.go b/dax/computer/service/computer.go index b41b05123..776f9807f 100644 --- a/dax/computer/service/computer.go +++ b/dax/computer/service/computer.go @@ -157,7 +157,8 @@ func newCommand(addr dax.Address, cfg CommandConfig) *fbserver.Command { var writeLoggerImpl computer.WriteLogService if cfg.ComputerConfig.WriteLogger != "" { - writeLoggerImpl = writeloggerclient.New(dax.Address(cfg.ComputerConfig.WriteLogger)) + panic("running separate writelogger is currently unsupported") + // writeLoggerImpl = writeloggerclient.New(dax.Address(cfg.ComputerConfig.WriteLogger)) } else if wlSvc != nil { writeLoggerImpl = wlSvc } else { @@ -166,7 +167,8 @@ func newCommand(addr dax.Address, cfg CommandConfig) *fbserver.Command { var snapshotterImpl computer.SnapshotService if cfg.ComputerConfig.Snapshotter != "" { - snapshotterImpl = snapshotterclient.New(dax.Address(cfg.ComputerConfig.Snapshotter)) + panic("running separate snapshotter is currently unsupported") + // snapshotterImpl = snapshotterclient.New(dax.Address(cfg.ComputerConfig.Snapshotter)) } else if ssSvc != nil { snapshotterImpl = ssSvc } else { diff --git a/dax/snapshotter/client/client.go b/dax/snapshotter/client/client.go index ac8e5a72a..7e0b1e71c 100644 --- a/dax/snapshotter/client/client.go +++ b/dax/snapshotter/client/client.go @@ -1,109 +1,111 @@ // Package client contains an http implementation of the WriteLogger client. package client -import ( - "bytes" - "encoding/json" - "fmt" - "io" - "net/http" - "net/url" +// import ( +// "bytes" +// "encoding/json" +// "fmt" +// "io" +// "net/http" +// "net/url" "github.com/featurebasedb/featurebase/v3/dax" snapshotterhttp "github.com/featurebasedb/featurebase/v3/dax/snapshotter/http" "github.com/featurebasedb/featurebase/v3/errors" ) -const defaultScheme = "http" +// const defaultScheme = "http" -// Snapshotter is a client for the Snapshotter API methods. -type Snapshotter struct { - address dax.Address -} +// // TODO(jaffee): remove this? -func New(address dax.Address) *Snapshotter { - return &Snapshotter{ - address: address, - } -} +// // Snapshotter is a client for the Snapshotter API methods. +// type Snapshotter struct { +// address dax.Address +// } -func (s *Snapshotter) Write(bucket string, key string, version int, rc io.ReadCloser) error { - url := fmt.Sprintf("%s/snapshotter/write-snapshot?bucket=%s&key=%s&version=%d", - s.address.WithScheme(defaultScheme), - url.QueryEscape(bucket), - url.QueryEscape(key), - version, - ) +// func New(address dax.Address) *Snapshotter { +// return &Snapshotter{ +// address: address, +// } +// } - // Post the request. - resp, err := http.Post(url, "", rc) - if err != nil { - return errors.Wrap(err, "posting write-snapshot") - } - defer resp.Body.Close() +// func (s *Snapshotter) Write(bucket string, key string, version int, rc io.ReadCloser) error { +// url := fmt.Sprintf("%s/snapshotter/write-snapshot?bucket=%s&key=%s&version=%d", +// s.address.WithScheme(defaultScheme), +// url.QueryEscape(bucket), +// url.QueryEscape(key), +// version, +// ) - if resp.StatusCode != http.StatusOK { - b, _ := io.ReadAll(resp.Body) - return errors.Errorf("status code: %d: %s", resp.StatusCode, b) - } +// // Post the request. +// resp, err := http.Post(url, "", rc) +// if err != nil { +// return errors.Wrap(err, "posting write-snapshot") +// } +// defer resp.Body.Close() - var wsr snapshotterhttp.WriteSnapshotResponse - if err := json.NewDecoder(resp.Body).Decode(&wsr); err != nil { - return errors.Wrap(err, "reading response body") - } +// if resp.StatusCode != http.StatusOK { +// b, _ := io.ReadAll(resp.Body) +// return errors.Errorf("status code: %d: %s", resp.StatusCode, b) +// } - return nil -} +// var wsr snapshotterhttp.WriteSnapshotResponse +// if err := json.NewDecoder(resp.Body).Decode(&wsr); err != nil { +// return errors.Wrap(err, "reading response body") +// } -// WriteTo is exactly the same as Write, except that it takes an io.WriteTo -// instead of an io.ReadCloser. This needs to be cleaned up so that we're only -// using one or the other. -func (s *Snapshotter) WriteTo(bucket string, key string, version int, wrTo io.WriterTo) error { - url := fmt.Sprintf("%s/snapshotter/write-snapshot?bucket=%s&key=%s&version=%d", - s.address.WithScheme(defaultScheme), - url.QueryEscape(bucket), - url.QueryEscape(key), - version, - ) +// return nil +// } - buf := &bytes.Buffer{} - if _, err := wrTo.WriteTo(buf); err != nil { - return errors.Wrap(err, "writing to buffer") - } +// // WriteTo is exactly the same as Write, except that it takes an io.WriteTo +// // instead of an io.ReadCloser. This needs to be cleaned up so that we're only +// // using one or the other. +// func (s *Snapshotter) WriteTo(bucket string, key string, version int, wrTo io.WriterTo) error { +// url := fmt.Sprintf("%s/snapshotter/write-snapshot?bucket=%s&key=%s&version=%d", +// s.address.WithScheme(defaultScheme), +// url.QueryEscape(bucket), +// url.QueryEscape(key), +// version, +// ) - // Post the request. - resp, err := http.Post(url, "", buf) - if err != nil { - return errors.Wrap(err, "posting write-snapshot") - } - defer resp.Body.Close() +// buf := &bytes.Buffer{} +// if _, err := wrTo.WriteTo(buf); err != nil { +// return errors.Wrap(err, "writing to buffer") +// } - if resp.StatusCode != http.StatusOK { - b, _ := io.ReadAll(resp.Body) - return errors.Errorf("status code: %d: %s", resp.StatusCode, b) - } +// // Post the request. +// resp, err := http.Post(url, "", buf) +// if err != nil { +// return errors.Wrap(err, "posting write-snapshot") +// } +// defer resp.Body.Close() - var wsr snapshotterhttp.WriteSnapshotResponse - if err := json.NewDecoder(resp.Body).Decode(&wsr); err != nil { - return errors.Wrap(err, "reading response body") - } +// if resp.StatusCode != http.StatusOK { +// b, _ := io.ReadAll(resp.Body) +// return errors.Errorf("status code: %d: %s", resp.StatusCode, b) +// } - return nil -} +// var wsr snapshotterhttp.WriteSnapshotResponse +// if err := json.NewDecoder(resp.Body).Decode(&wsr); err != nil { +// return errors.Wrap(err, "reading response body") +// } -func (s *Snapshotter) Read(bucket string, key string, version int) (io.ReadCloser, error) { - url := fmt.Sprintf("%s/snapshotter/read-snapshot?bucket=%s&key=%s&version=%d", - s.address.WithScheme(defaultScheme), - url.QueryEscape(bucket), - url.QueryEscape(key), - version, - ) +// return nil +// } - // Get the request. - resp, err := http.Get(url) - if err != nil { - return nil, errors.Wrap(err, "getting read-snapshot") - } +// func (s *Snapshotter) Read(bucket string, key string, version int) (io.ReadCloser, error) { +// url := fmt.Sprintf("%s/snapshotter/read-snapshot?bucket=%s&key=%s&version=%d", +// s.address.WithScheme(defaultScheme), +// url.QueryEscape(bucket), +// url.QueryEscape(key), +// version, +// ) - return resp.Body, nil -} +// // Get the request. +// resp, err := http.Get(url) +// if err != nil { +// return nil, errors.Wrap(err, "getting read-snapshot") +// } + +// return resp.Body, nil +// } diff --git a/dax/snapshotter/http/handler.go b/dax/snapshotter/http/handler.go index b559b2845..804830d1a 100644 --- a/dax/snapshotter/http/handler.go +++ b/dax/snapshotter/http/handler.go @@ -1,120 +1,120 @@ package http -import ( - "encoding/json" - "io" - "net/http" - "strconv" +// import ( +// "encoding/json" +// "io" +// "net/http" +// "strconv" "github.com/gorilla/mux" "github.com/featurebasedb/featurebase/v3/dax/snapshotter" "github.com/featurebasedb/featurebase/v3/rbf" ) -func Handler(s *snapshotter.Snapshotter) http.Handler { - svr := &server{ - snapshotter: s, - } +// func Handler(s *snapshotter.Snapshotter) http.Handler { +// svr := &server{ +// snapshotter: s, +// } - router := mux.NewRouter() - router.HandleFunc("/health", svr.getHealth).Methods("GET").Name("GetHealth") - router.HandleFunc("/write-snapshot", svr.postWriteSnapshot).Methods("POST").Name("PostWriteSnapshot") - router.HandleFunc("/read-snapshot", svr.getReadSnapshot).Methods("GET").Name("GetReadSnapshot") - return router -} +// router := mux.NewRouter() +// router.HandleFunc("/health", svr.getHealth).Methods("GET").Name("GetHealth") +// router.HandleFunc("/write-snapshot", svr.postWriteSnapshot).Methods("POST").Name("PostWriteSnapshot") +// router.HandleFunc("/read-snapshot", svr.getReadSnapshot).Methods("GET").Name("GetReadSnapshot") +// return router +// } -type server struct { - snapshotter *snapshotter.Snapshotter -} +// type server struct { +// snapshotter *snapshotter.Snapshotter +// } -// GET /health -func (s *server) getHealth(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusOK) -} +// // GET /health +// func (s *server) getHealth(w http.ResponseWriter, r *http.Request) { +// w.WriteHeader(http.StatusOK) +// } -// POST /write-snapshot -func (s *server) postWriteSnapshot(w http.ResponseWriter, r *http.Request) { - bucket := r.URL.Query().Get("bucket") - if bucket == "" { - http.Error(w, "bucket required", http.StatusBadRequest) - return - } +// // POST /write-snapshot +// func (s *server) postWriteSnapshot(w http.ResponseWriter, r *http.Request) { +// bucket := r.URL.Query().Get("bucket") +// if bucket == "" { +// http.Error(w, "bucket required", http.StatusBadRequest) +// return +// } - key := r.URL.Query().Get("key") - if key == "" { - http.Error(w, "key required", http.StatusBadRequest) - return - } +// key := r.URL.Query().Get("key") +// if key == "" { +// http.Error(w, "key required", http.StatusBadRequest) +// return +// } - versionArg := r.URL.Query().Get("version") - versionInt64, err := strconv.ParseInt(versionArg, 10, 64) - if err != nil { - http.Error(w, "bad shard", http.StatusBadRequest) - return - } - version := int(versionInt64) +// versionArg := r.URL.Query().Get("version") +// versionInt64, err := strconv.ParseInt(versionArg, 10, 64) +// if err != nil { +// http.Error(w, "bad shard", http.StatusBadRequest) +// return +// } +// version := int(versionInt64) - body := r.Body - defer body.Close() +// body := r.Body +// defer body.Close() - if err := s.snapshotter.Write(bucket, key, version, body); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } +// if err := s.snapshotter.Write(bucket, key, version, body); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } - resp := &WriteSnapshotResponse{} +// resp := &WriteSnapshotResponse{} - if err := json.NewEncoder(w).Encode(resp); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } -} +// if err := json.NewEncoder(w).Encode(resp); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } +// } -type WriteSnapshotResponse struct{} +// type WriteSnapshotResponse struct{} -// GET /read-snapshot -func (s *server) getReadSnapshot(w http.ResponseWriter, r *http.Request) { - bucket := r.URL.Query().Get("bucket") - if bucket == "" { - http.Error(w, "bucket required", http.StatusBadRequest) - return - } +// // GET /read-snapshot +// func (s *server) getReadSnapshot(w http.ResponseWriter, r *http.Request) { +// bucket := r.URL.Query().Get("bucket") +// if bucket == "" { +// http.Error(w, "bucket required", http.StatusBadRequest) +// return +// } - key := r.URL.Query().Get("key") - if key == "" { - http.Error(w, "key required", http.StatusBadRequest) - return - } +// key := r.URL.Query().Get("key") +// if key == "" { +// http.Error(w, "key required", http.StatusBadRequest) +// return +// } - versionArg := r.URL.Query().Get("version") - versionInt64, err := strconv.ParseInt(versionArg, 10, 64) - if err != nil { - http.Error(w, "bad shard", http.StatusBadRequest) - return - } - version := int(versionInt64) +// versionArg := r.URL.Query().Get("version") +// versionInt64, err := strconv.ParseInt(versionArg, 10, 64) +// if err != nil { +// http.Error(w, "bad shard", http.StatusBadRequest) +// return +// } +// version := int(versionInt64) - rc, err := s.snapshotter.Read(bucket, key, version) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } - defer rc.Close() +// rc, err := s.snapshotter.Read(bucket, key, version) +// if err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } +// defer rc.Close() - // TODO: is rbf.PageSize a problem here for non-RBF snapshots (i.e. keys)? - // Copy data to response body. - if _, err := io.CopyBuffer(&passthroughWriter{w}, rc, make([]byte, rbf.PageSize)); err != nil { - http.Error(w, err.Error(), http.StatusInternalServerError) - return - } -} +// // TODO: is rbf.PageSize a problem here for non-RBF snapshots (i.e. keys)? +// // Copy data to response body. +// if _, err := io.CopyBuffer(&passthroughWriter{w}, rc, make([]byte, rbf.PageSize)); err != nil { +// http.Error(w, err.Error(), http.StatusInternalServerError) +// return +// } +// } -// passthroughWriter is used to remove non-Writer interfaces from an io.Writer. -// For example, a writer that implements io.ReaderFrom can change io.Copy() behavior. -type passthroughWriter struct { - w io.Writer -} +// // passthroughWriter is used to remove non-Writer interfaces from an io.Writer. +// // For example, a writer that implements io.ReaderFrom can change io.Copy() behavior. +// type passthroughWriter struct { +// w io.Writer +// } -func (w *passthroughWriter) Write(p []byte) (int, error) { - return w.w.Write(p) -} +// func (w *passthroughWriter) Write(p []byte) (int, error) { +// return w.w.Write(p) +// } diff --git a/dax/storage/encoding.go b/dax/storage/encoding.go new file mode 100644 index 000000000..d372e6e31 --- /dev/null +++ b/dax/storage/encoding.go @@ -0,0 +1,149 @@ +package storage + +import ( + "bufio" + "encoding/json" + "io" + + "github.com/molecula/featurebase/v3/dax" + "github.com/molecula/featurebase/v3/dax/computer" +) + +// TODO this needs to be genericized and moved + +type TableKeyReader struct { + table dax.TableKey + partition dax.PartitionNum + scanner *bufio.Scanner + closer io.Closer +} + +func NewTableKeyReader(qtid dax.QualifiedTableID, partition dax.PartitionNum, writelog io.ReadCloser) *TableKeyReader { + r := &TableKeyReader{ + table: qtid.Key(), + partition: partition, + scanner: bufio.NewScanner(writelog), + closer: writelog, + } + + return r +} + +func (r *TableKeyReader) Read() (computer.PartitionKeyMap, error) { + if r.scanner == nil { + return computer.PartitionKeyMap{}, io.EOF + } + + var b []byte + var out computer.PartitionKeyMap + + if r.scanner.Scan() { + b = r.scanner.Bytes() + if err := json.Unmarshal(b, &out); err != nil { + return out, err + } + return out, nil + } + if err := r.scanner.Err(); err != nil { + return out, err + } + + return out, io.EOF +} + +func (r *TableKeyReader) Close() error { + if r.closer != nil { + return r.closer.Close() + } + return nil +} + +type FieldKeyReader struct { + table dax.TableKey + field dax.FieldName + scanner *bufio.Scanner + closer io.Closer +} + +func NewFieldKeyReader(qtid dax.QualifiedTableID, field dax.FieldName, writelog io.ReadCloser) *FieldKeyReader { + r := &FieldKeyReader{ + table: qtid.Key(), + field: field, + scanner: bufio.NewScanner(writelog), + closer: writelog, + } + + return r +} + +func (r *FieldKeyReader) Read() (computer.FieldKeyMap, error) { + if r.scanner == nil { + return computer.FieldKeyMap{}, io.EOF + } + + var b []byte + var out computer.FieldKeyMap + + if r.scanner.Scan() { + b = r.scanner.Bytes() + if err := json.Unmarshal(b, &out); err != nil { + return out, err + } + return out, nil + } + if err := r.scanner.Err(); err != nil { + return out, err + } + + return out, io.EOF +} + +func (r *FieldKeyReader) Close() error { + if r.closer != nil { + return r.closer.Close() + } + return nil +} + +type ShardReader struct { + table dax.TableKey + partition dax.PartitionNum + shard dax.ShardNum + version int + scanner *bufio.Scanner + closer io.Closer +} + +func NewShardReader(qtid dax.QualifiedTableID, partition dax.PartitionNum, shard dax.ShardNum, writelog io.ReadCloser) *ShardReader { + r := &ShardReader{ + table: qtid.Key(), + partition: partition, + shard: shard, + scanner: bufio.NewScanner(writelog), + closer: writelog, + } + + return r +} + +func (r *ShardReader) Read() (computer.LogMessage, error) { + if r.scanner == nil { + return nil, io.EOF + } + + if r.scanner.Scan() { + return computer.UnmarshalLogMessage(r.scanner.Bytes()) + } + if err := r.scanner.Err(); err != nil { + return nil, err + } + + return nil, io.EOF +} + +func (r *ShardReader) Close() error { + if r.closer != nil { + return r.closer.Close() + } + return nil +} diff --git a/dax/writelogger/client/client.go b/dax/writelogger/client/client.go index 73ac31338..202d6164d 100644 --- a/dax/writelogger/client/client.go +++ b/dax/writelogger/client/client.go @@ -1,145 +1,147 @@ // Package client contains an http implementation of the WriteLogger client. package client -import ( - "bytes" - "encoding/json" - "fmt" - "io" - "net/http" +// import ( +// "bytes" +// "encoding/json" +// "fmt" +// "io" +// "net/http" "github.com/featurebasedb/featurebase/v3/dax" "github.com/featurebasedb/featurebase/v3/errors" ) -const defaultScheme = "http" +// // TODO(jaffee): remove this? -// WriteLogger is a client for the WriteLogger API methods. -type WriteLogger struct { - address dax.Address -} +// const defaultScheme = "http" -func New(address dax.Address) *WriteLogger { - return &WriteLogger{ - address: address, - } -} +// // WriteLogger is a client for the WriteLogger API methods. +// type WriteLogger struct { +// address dax.Address +// } -func (w *WriteLogger) AppendMessage(bucket string, key string, version int, msg []byte) error { - url := fmt.Sprintf("%s/writelogger/append-message", w.address.WithScheme(defaultScheme)) +// func New(address dax.Address) *WriteLogger { +// return &WriteLogger{ +// address: address, +// } +// } - req := &AppendMessageRequest{ - Bucket: bucket, - Key: key, - Version: version, - Message: msg, - } +// func (w *WriteLogger) AppendMessage(bucket string, key string, version int, msg []byte) error { +// url := fmt.Sprintf("%s/writelogger/append-message", w.address.WithScheme(defaultScheme)) - // Encode the request. - postBody, err := json.Marshal(req) - if err != nil { - return errors.Wrap(err, "marshalling post request") - } - requestBody := bytes.NewBuffer(postBody) +// req := &AppendMessageRequest{ +// Bucket: bucket, +// Key: key, +// Version: version, +// Message: msg, +// } - // Post the request. - resp, err := http.Post(url, "application/json", requestBody) - if err != nil { - return errors.Wrap(err, "posting append-message request") - } - defer resp.Body.Close() +// // Encode the request. +// postBody, err := json.Marshal(req) +// if err != nil { +// return errors.Wrap(err, "marshalling post request") +// } +// requestBody := bytes.NewBuffer(postBody) - if resp.StatusCode != http.StatusOK { - b, _ := io.ReadAll(resp.Body) - return errors.Errorf("status code: %d: %s", resp.StatusCode, b) - } +// // Post the request. +// resp, err := http.Post(url, "application/json", requestBody) +// if err != nil { +// return errors.Wrap(err, "posting append-message request") +// } +// defer resp.Body.Close() - var isr *AppendMessageResponse - if err := json.NewDecoder(resp.Body).Decode(&isr); err != nil { - return errors.Wrap(err, "reading response body") - } +// if resp.StatusCode != http.StatusOK { +// b, _ := io.ReadAll(resp.Body) +// return errors.Errorf("status code: %d: %s", resp.StatusCode, b) +// } - return nil -} +// var isr *AppendMessageResponse +// if err := json.NewDecoder(resp.Body).Decode(&isr); err != nil { +// return errors.Wrap(err, "reading response body") +// } -type AppendMessageRequest struct { - Bucket string `json:"bucket"` - Key string `json:"key"` - Version int `json:"version"` - Message []byte `json:"message"` -} -type AppendMessageResponse struct{} +// return nil +// } -func (w *WriteLogger) LogReader(bucket string, key string, version int) (io.Reader, io.Closer, error) { - url := fmt.Sprintf("%s/writelogger/log-reader", w.address.WithScheme(defaultScheme)) +// type AppendMessageRequest struct { +// Bucket string `json:"bucket"` +// Key string `json:"key"` +// Version int `json:"version"` +// Message []byte `json:"message"` +// } +// type AppendMessageResponse struct{} - req := &LogReaderRequest{ - Bucket: bucket, - Version: version, - Key: key, - } +// func (w *WriteLogger) LogReader(bucket string, key string, version int) (io.Reader, io.Closer, error) { +// url := fmt.Sprintf("%s/writelogger/log-reader", w.address.WithScheme(defaultScheme)) - // Encode the request. - postBody, err := json.Marshal(req) - if err != nil { - return nil, nil, errors.Wrap(err, "marshalling post request") - } - requestBody := bytes.NewBuffer(postBody) +// req := &LogReaderRequest{ +// Bucket: bucket, +// Version: version, +// Key: key, +// } - // Post the request. - resp, err := http.Post(url, "application/json", requestBody) - if err != nil { - return nil, nil, errors.Wrap(err, "posting log-reader request") - } +// // Encode the request. +// postBody, err := json.Marshal(req) +// if err != nil { +// return nil, nil, errors.Wrap(err, "marshalling post request") +// } +// requestBody := bytes.NewBuffer(postBody) - if resp.StatusCode != http.StatusOK { - b, _ := io.ReadAll(resp.Body) - defer resp.Body.Close() - return nil, nil, errors.Errorf("status code: %d: %s", resp.StatusCode, b) - } +// // Post the request. +// resp, err := http.Post(url, "application/json", requestBody) +// if err != nil { +// return nil, nil, errors.Wrap(err, "posting log-reader request") +// } - return resp.Body, resp.Body, nil -} +// if resp.StatusCode != http.StatusOK { +// b, _ := io.ReadAll(resp.Body) +// defer resp.Body.Close() +// return nil, nil, errors.Errorf("status code: %d: %s", resp.StatusCode, b) +// } -type LogReaderRequest struct { - Bucket string `json:"bucket"` - Version int `json:"version"` - Key string `json:"key"` -} +// return resp.Body, resp.Body, nil +// } -func (w *WriteLogger) DeleteLog(bucket string, key string, version int) error { - url := fmt.Sprintf("%s/writelogger/delete-log", w.address.WithScheme(defaultScheme)) +// type LogReaderRequest struct { +// Bucket string `json:"bucket"` +// Version int `json:"version"` +// Key string `json:"key"` +// } - req := &DeleteLogRequest{ - Bucket: bucket, - Version: version, - Key: key, - } +// func (w *WriteLogger) DeleteLog(bucket string, key string, version int) error { +// url := fmt.Sprintf("%s/writelogger/delete-log", w.address.WithScheme(defaultScheme)) - // Encode the request. - postBody, err := json.Marshal(req) - if err != nil { - return errors.Wrap(err, "marshalling post request") - } - requestBody := bytes.NewBuffer(postBody) +// req := &DeleteLogRequest{ +// Bucket: bucket, +// Version: version, +// Key: key, +// } - // Post the request. - resp, err := http.Post(url, "application/json", requestBody) - if err != nil { - return errors.Wrap(err, "posting log-reader request") - } +// // Encode the request. +// postBody, err := json.Marshal(req) +// if err != nil { +// return errors.Wrap(err, "marshalling post request") +// } +// requestBody := bytes.NewBuffer(postBody) - if resp.StatusCode != http.StatusOK { - b, _ := io.ReadAll(resp.Body) - defer resp.Body.Close() - return errors.Errorf("status code: %d: %s", resp.StatusCode, b) - } +// // Post the request. +// resp, err := http.Post(url, "application/json", requestBody) +// if err != nil { +// return errors.Wrap(err, "posting log-reader request") +// } - return nil -} +// if resp.StatusCode != http.StatusOK { +// b, _ := io.ReadAll(resp.Body) +// defer resp.Body.Close() +// return errors.Errorf("status code: %d: %s", resp.StatusCode, b) +// } -type DeleteLogRequest struct { - Bucket string `json:"bucket"` - Version int `json:"version"` - Key string `json:"key"` -} +// return nil +// } + +// type DeleteLogRequest struct { +// Bucket string `json:"bucket"` +// Version int `json:"version"` +// Key string `json:"key"` +// } diff --git a/dax/writelogger/http/handler.go b/dax/writelogger/http/handler.go index 6366a7f26..0fa0c5de2 100644 --- a/dax/writelogger/http/handler.go +++ b/dax/writelogger/http/handler.go @@ -1,121 +1,121 @@ package http -import ( - "encoding/json" - "io" - "net/http" +// import ( +// "encoding/json" +// "io" +// "net/http" "github.com/gorilla/mux" "github.com/featurebasedb/featurebase/v3/dax/writelogger" "github.com/featurebasedb/featurebase/v3/logger" ) -func Handler(w *writelogger.WriteLogger, logger logger.Logger) http.Handler { - svr := &server{ - writeLogger: w, - logger: logger, - } +// func Handler(w *writelogger.WriteLogger, logger logger.Logger) http.Handler { +// svr := &server{ +// writeLogger: w, +// logger: logger, +// } - router := mux.NewRouter() - router.HandleFunc("/health", svr.getHealth).Methods("GET").Name("GetHealth") - router.HandleFunc("/append-message", svr.postAppendMessage).Methods("POST").Name("PostAppendMessage") - router.HandleFunc("/log-reader", svr.postLogReader).Methods("POST").Name("PostLogReader") - router.HandleFunc("/delete-log", svr.postDeleteLog).Methods("POST").Name("PostDeleteLog") - return router -} +// router := mux.NewRouter() +// router.HandleFunc("/health", svr.getHealth).Methods("GET").Name("GetHealth") +// router.HandleFunc("/append-message", svr.postAppendMessage).Methods("POST").Name("PostAppendMessage") +// router.HandleFunc("/log-reader", svr.postLogReader).Methods("POST").Name("PostLogReader") +// router.HandleFunc("/delete-log", svr.postDeleteLog).Methods("POST").Name("PostDeleteLog") +// return router +// } -type server struct { - writeLogger *writelogger.WriteLogger - logger logger.Logger -} +// type server struct { +// writeLogger *writelogger.WriteLogger +// logger logger.Logger +// } -// GET /health -func (s *server) getHealth(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusOK) -} +// // GET /health +// func (s *server) getHealth(w http.ResponseWriter, r *http.Request) { +// w.WriteHeader(http.StatusOK) +// } -// POST /append-message -func (s *server) postAppendMessage(w http.ResponseWriter, r *http.Request) { - body := r.Body - defer body.Close() +// // POST /append-message +// func (s *server) postAppendMessage(w http.ResponseWriter, r *http.Request) { +// body := r.Body +// defer body.Close() - req := AppendMessageRequest{} - if err := json.NewDecoder(body).Decode(&req); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } +// req := AppendMessageRequest{} +// if err := json.NewDecoder(body).Decode(&req); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } - err := s.writeLogger.AppendMessage(req.Bucket, req.Key, req.Version, req.Message) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } +// err := s.writeLogger.AppendMessage(req.Bucket, req.Key, req.Version, req.Message) +// if err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } - resp := AppendMessageResponse{} - if err := json.NewEncoder(w).Encode(resp); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } -} +// resp := AppendMessageResponse{} +// if err := json.NewEncoder(w).Encode(resp); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } +// } -type AppendMessageRequest struct { - Bucket string `json:"bucket"` - Key string `json:"key"` - Version int `json:"version"` - Message []byte `json:"message"` -} +// type AppendMessageRequest struct { +// Bucket string `json:"bucket"` +// Key string `json:"key"` +// Version int `json:"version"` +// Message []byte `json:"message"` +// } -type AppendMessageResponse struct{} +// type AppendMessageResponse struct{} -// POST /log-reader -func (s *server) postLogReader(w http.ResponseWriter, r *http.Request) { - body := r.Body - defer body.Close() +// // POST /log-reader +// func (s *server) postLogReader(w http.ResponseWriter, r *http.Request) { +// body := r.Body +// defer body.Close() - req := LogReaderRequest{} - if err := json.NewDecoder(body).Decode(&req); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } +// req := LogReaderRequest{} +// if err := json.NewDecoder(body).Decode(&req); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } - reader, closer, err := s.writeLogger.LogReader(req.Bucket, req.Key, req.Version) - if err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } +// reader, closer, err := s.writeLogger.LogReader(req.Bucket, req.Key, req.Version) +// if err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } - defer closer.Close() +// defer closer.Close() - if _, err := io.Copy(w, reader); err != nil { - s.logger.Printf("error streaming log data: %s", err) - } -} +// if _, err := io.Copy(w, reader); err != nil { +// s.logger.Printf("error streaming log data: %s", err) +// } +// } -type LogReaderRequest struct { - Bucket string `json:"bucket"` - Version int `json:"version"` - Key string `json:"key"` -} +// type LogReaderRequest struct { +// Bucket string `json:"bucket"` +// Version int `json:"version"` +// Key string `json:"key"` +// } -// POST /delete-log -func (s *server) postDeleteLog(w http.ResponseWriter, r *http.Request) { - body := r.Body - defer body.Close() +// // POST /delete-log +// func (s *server) postDeleteLog(w http.ResponseWriter, r *http.Request) { +// body := r.Body +// defer body.Close() - req := DeleteLogRequest{} - if err := json.NewDecoder(body).Decode(&req); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } +// req := DeleteLogRequest{} +// if err := json.NewDecoder(body).Decode(&req); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } - if err := s.writeLogger.DeleteLog(req.Bucket, req.Key, req.Version); err != nil { - http.Error(w, err.Error(), http.StatusBadRequest) - return - } -} +// if err := s.writeLogger.DeleteLog(req.Bucket, req.Key, req.Version); err != nil { +// http.Error(w, err.Error(), http.StatusBadRequest) +// return +// } +// } -type DeleteLogRequest struct { - Bucket string `json:"bucket"` - Version int `json:"version"` - Key string `json:"key"` -} +// type DeleteLogRequest struct { +// Bucket string `json:"bucket"` +// Version int `json:"version"` +// Key string `json:"key"` +// } diff --git a/dax/writelogger/writelogger.go b/dax/writelogger/writelogger.go index ae7a8957a..22a8cde18 100644 --- a/dax/writelogger/writelogger.go +++ b/dax/writelogger/writelogger.go @@ -16,9 +16,9 @@ import ( ) type WriteLogger struct { - mu sync.RWMutex + dataDir string - dataDir string + mu sync.RWMutex logFiles map[string]*os.File lockFiles map[string]*os.File @@ -128,8 +128,6 @@ func (w *WriteLogger) lockFile(bucket, key string) (string, string) { func (w *WriteLogger) Lock(bucket, key string) error { lockDir, lockFile := w.lockFile(bucket, key) - fmt.Println("lock dir:", lockDir) - fmt.Println("lock fil:", lockFile) if err := os.MkdirAll(lockDir, 0777); err != nil { return errors.Wrapf(err, "lock dir %s", lockDir) @@ -138,6 +136,8 @@ func (w *WriteLogger) Lock(bucket, key string) error { if err != nil { return errors.Wrapf(err, "opening lock file: %s", lockFile) } + w.mu.Lock() + defer w.mu.Unlock() w.lockFiles[lockFile] = f // fd, err = syscall.Open(lockFile, syscall.O_RDWR|syscall.O_CREAT, 0644) // if err != nil { @@ -151,6 +151,8 @@ func (w *WriteLogger) Lock(bucket, key string) error { } func (w *WriteLogger) Unlock(bucket, key string) error { + w.mu.Lock() + defer w.mu.Unlock() // TODO(jaffee) since the file isn't guaranteed to be removed if // the process is killed, we should actually use flock instead of // EXCL file creation. Problem with that is it makes testing @@ -209,7 +211,7 @@ func (w *WriteLogger) logFileByKey(key string) (*os.File, error) { if err != nil { return nil, errors.Wrapf(err, "opening file: %s", filePath) } - fmt.Printf("opened %s for key %s\n", filePath, key) + w.logFiles[key] = f return f, nil diff --git a/dax/writelogger/writelogger_test.go b/dax/writelogger/writelogger_test.go index ddcf17745..62d1d4667 100644 --- a/dax/writelogger/writelogger_test.go +++ b/dax/writelogger/writelogger_test.go @@ -50,11 +50,11 @@ func TestWriteLogger(t *testing.T) { assert.NoError(t, err) // Read the message. - reader, closer, err := wl.LogReader(bucket(table, partition), key, version) + readcloser, err := wl.LogReader(bucket(table, partition), key, version) assert.NoError(t, err) - defer closer.Close() + defer readcloser.Close() - buf, err := io.ReadAll(reader) + buf, err := io.ReadAll(readcloser) assert.NoError(t, err) var out payload