(cherry picked from commit bbaa7dd0f1)
This commit is contained in:
Matthew Jaffee 2022-12-13 10:32:30 -06:00 committed by Joe Friedrich
parent ed7c6d419e
commit ff3595d759
10 changed files with 630 additions and 494 deletions

View file

@ -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()
}
//////////////////////////////////////////////////////////////

View file

@ -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)

View file

@ -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 {

View file

@ -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
// }

View file

@ -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)
// }

149
dax/storage/encoding.go Normal file
View file

@ -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
}

View file

@ -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"`
// }

View file

@ -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"`
// }

View file

@ -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

View file

@ -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